fix(agent): synchronize shared Codex sessions

This commit is contained in:
yu
2026-07-17 18:30:10 +08:00
parent 062e4569aa
commit 5e1fd7a825
8 changed files with 347 additions and 127 deletions
+53 -29
View File
@@ -11,11 +11,12 @@ import type { AgentAttachment, AgentEmit } from "./types.js";
type Json = Record<string, unknown>;
type AgentEvent = Json & { type: string; usage?: unknown };
type PendingRequest = { resolve: (value: unknown) => void; reject: (error: Error) => void };
type CodexRunOptions = { threadId?: string; cwd?: string };
type CodexRunOptions = { threadId?: string; cwd?: string; appEmit?: AgentEmit; onStart?: () => void; onThread?: (threadId: string) => void; onTurn?: (turnId: string) => void; onFinish?: () => void };
type AgentHistoryMessage = { id: string; role: "user" | "assistant" | "tool" | "error"; title?: string; text: string; detail?: unknown; streamId?: string };
let codexQueue: Promise<unknown> = Promise.resolve();
let codexApp: CodexAppClient | null = null;
let codexAppStart: Promise<CodexAppClient> | null = null;
let codexThreadId = "";
const canvasAgentMcp = canvasAgentMcpCommand();
const require = createRequire(import.meta.url);
@@ -30,52 +31,56 @@ export async function runCodexTurn(prompt: string, emit: AgentEmit, attachments:
await codexQueue;
}
export function interruptCodexTurn() {
if (!codexApp) return false;
export function interruptCodexTurn(threadId?: string) {
if (!codexApp || (threadId && threadId !== codexThreadId)) return false;
return codexApp.interruptCurrentTurn();
}
async function runCodexTurnNow(prompt: string, emit: AgentEmit, attachments: AgentAttachment[], options: CodexRunOptions) {
let files: string[] = [];
try {
options.onStart?.();
files = await writeAttachmentFiles(attachments);
codexApp ||= await CodexAppClient.start(emit);
let threadId = await ensureCodexThread(codexApp, options, emit);
const app = await getCodexApp(options.appEmit || emit);
let threadId = await ensureCodexThread(app, options, emit);
options.onThread?.(threadId);
try {
await codexApp.startTurn(threadId, prompt, files);
await app.startTurn(threadId, prompt, files, options.onTurn);
} catch (error) {
if (!isRecoverableThreadError(error)) throw error;
emit("agent_log", { text: `Codex thread unavailable, starting a new thread: ${errorMessage(error)}` });
codexThreadId = "";
threadId = await ensureCodexThread(codexApp, { cwd: options.cwd }, emit);
await codexApp.startTurn(threadId, prompt, files);
threadId = await ensureCodexThread(app, { cwd: options.cwd }, emit);
options.onThread?.(threadId);
await app.startTurn(threadId, prompt, files, options.onTurn);
}
} catch (error) {
emit("agent_error", { message: errorMessage(error) });
} finally {
options.onFinish?.();
await Promise.all(files.map((file) => fs.unlink(file).catch(() => undefined)));
}
}
export async function startCodexThread(emit: AgentEmit, cwd?: string) {
codexApp ||= await CodexAppClient.start(emit);
const thread = await codexApp.startThread(cwd);
const app = await getCodexApp(emit);
const thread = await app.startThread(cwd);
codexThreadId = String(field(thread, "id") || "");
return thread;
}
export async function resumeCodexThread(emit: AgentEmit, threadId: string, cwd?: string) {
codexApp ||= await CodexAppClient.start(emit);
const app = await getCodexApp(emit);
await loadCodexThread(emit, threadId, cwd, false);
const thread = await codexApp.resumeThread(threadId, cwd);
const thread = await app.resumeThread(threadId, cwd);
assertThreadWorkspace(thread, cwd);
codexThreadId = String(field(thread, "id") || threadId);
return { thread, messages: threadMessages(thread) };
}
export async function listCodexThreads(emit: AgentEmit, options: { cwd: string; searchTerm?: string; limit?: number }) {
codexApp ||= await CodexAppClient.start(emit);
const result = await codexApp.listThreads({
const app = await getCodexApp(emit);
const result = await app.listThreads({
limit: options.limit || 40,
sortKey: "updated_at",
sortDirection: "desc",
@@ -97,9 +102,9 @@ export async function verifyCodexThreadWorkspace(emit: AgentEmit, threadId: stri
}
export async function archiveCodexThread(emit: AgentEmit, threadId: string, cwd?: string) {
codexApp ||= await CodexAppClient.start(emit);
const app = await getCodexApp(emit);
await loadCodexThread(emit, threadId, cwd, false);
await codexApp.archiveThread(threadId);
await app.archiveThread(threadId);
}
export function runClaudeTurn(prompt: string, emit: AgentEmit) {
@@ -131,7 +136,7 @@ async function ensureCodexThread(app: CodexAppClient, options: CodexRunOptions,
return codexThreadId;
}
function isRecoverableThreadError(error: unknown) {
export function isRecoverableThreadError(error: unknown) {
return /thread not loaded|no rollout found/i.test(errorMessage(error));
}
@@ -192,10 +197,11 @@ class CodexAppClient {
return this.request("thread/archive", { threadId });
}
async startTurn(threadId: string, prompt: string, images: string[]) {
async startTurn(threadId: string, prompt: string, images: string[], onTurn?: (turnId: string) => void) {
const result = await this.request("turn/start", { threadId, input: codexInput(prompt, images), approvalPolicy: "never" });
const turnId = String(field(field(result, "turn"), "id") || "");
if (!turnId) throw new Error("Codex app-server 没有返回 turn id");
onTurn?.(turnId);
const completed = this.completedTurns.get(turnId);
if (this.completedTurns.has(turnId)) {
this.completedTurns.delete(turnId);
@@ -267,9 +273,9 @@ class CodexAppClient {
} else if (turnId) {
this.completedTurns.set(turnId, error ? new Error(String(field(error, "message") || "Codex turn failed")) : null);
}
this.emit("agent_event", { agent: "codex", type: "stream.summary", delta_count: this.deltaCount });
this.emit("agent_event", { agent: "codex", type: "stream.summary", delta_count: this.deltaCount, ...codexEventScope(params) });
this.deltaCount = 0;
this.emit("agent_done", { agent: "codex", usage: event.usage });
this.emit("agent_done", { agent: "codex", usage: event.usage, ...codexEventScope(params) });
}
}
@@ -278,7 +284,7 @@ class CodexAppClient {
const text = `${this.textByItem.get(id) || ""}${String(field(params, "delta") || "")}`;
this.deltaCount += 1;
this.textByItem.set(id, text);
this.emit("agent_event", { agent: "codex", type: "item.updated", item: { id, type: "agent_message", text } });
this.emit("agent_event", { agent: "codex", type: "item.updated", item: { id, type: "agent_message", text }, ...codexEventScope(params) });
}
private answerServerRequest(message: Json) {
@@ -321,23 +327,41 @@ function codexInput(prompt: string, images: string[]) {
}
function normalizeCodexNotification(method: string, params: Json): AgentEvent | null {
if (method === "thread/started") return { type: "thread.started", thread_id: field(field(params, "thread"), "id") };
if (method === "turn/started") return { type: "turn.started" };
if (method === "turn/completed") return { type: "turn.completed", usage: null };
if (method === "item/started") return { type: "item.started", item: normalizeItem(field(params, "item")) };
if (method === "item/completed") return { type: "item.completed", item: normalizeItem(field(params, "item")) };
if (method === "error") return { type: "error", message: field(params, "message") };
const scope = codexEventScope(params);
if (method === "thread/started") return { type: "thread.started", ...scope };
if (method === "turn/started") return { type: "turn.started", ...scope };
if (method === "turn/completed") return { type: "turn.completed", usage: null, ...scope };
if (method === "item/started") return { type: "item.started", item: normalizeItem(field(params, "item")), ...scope };
if (method === "item/completed") return { type: "item.completed", item: normalizeItem(field(params, "item")), ...scope };
if (method === "error") return { type: "error", message: field(params, "message"), ...scope };
return null;
}
function codexEventScope(params: Json) {
const threadId = String(field(params, "threadId") || field(field(params, "thread"), "id") || "");
const turnId = String(field(params, "turnId") || field(field(params, "turn"), "id") || "");
return { ...(threadId ? { thread_id: threadId } : {}), ...(turnId ? { turn_id: turnId } : {}) };
}
async function loadCodexThread(emit: AgentEmit, threadId: string, cwd: string | undefined, includeTurns: boolean) {
codexApp ||= await CodexAppClient.start(emit);
const result = await codexApp.readThread(threadId, includeTurns);
const app = await getCodexApp(emit);
const result = await app.readThread(threadId, includeTurns);
const thread = field(result, "thread") || {};
assertThreadWorkspace(thread, cwd);
return thread;
}
async function getCodexApp(emit: AgentEmit) {
if (codexApp) return codexApp;
codexAppStart ||= CodexAppClient.start(emit);
try {
codexApp = await codexAppStart;
return codexApp;
} finally {
codexAppStart = null;
}
}
function assertThreadWorkspace(thread: unknown, cwd?: string) {
if (!cwd || threadInWorkspace(thread, cwd)) return;
throw new Error("该 Codex 会话不属于当前画布工作空间");
+72
View File
@@ -132,6 +132,78 @@ test("closing a client rejects its pending tool requests", async () => {
assert.match(outcome, /断开/);
});
test("shared thread events are broadcast with the active thread id", (t) => {
const session = new CanvasSession();
const first = connect(session, "first");
const second = connect(session, "second");
t.after(() => {
first.close();
second.close();
});
session.emitThread("workspace_changed", "thread-2", { activeThreadId: "thread-2" });
assert.deepEqual(first.event("workspace_changed"), { activeThreadId: "thread-2", threadId: "thread-2" });
assert.deepEqual(second.event("workspace_changed"), { activeThreadId: "thread-2", threadId: "thread-2" });
});
test("new clients receive the current Codex state and later updates", (t) => {
const session = new CanvasSession();
session.setCodexState({ busy: true, threadId: "thread-2", turnId: "turn-1" });
const client = connect(session, "first");
t.after(() => client.close());
assert.deepEqual(field(client.event("hello"), "codex"), { busy: true, threadId: "thread-2", turnId: "turn-1" });
session.setCodexState({ busy: false });
assert.deepEqual(client.event("codex_state"), { busy: false, threadId: "thread-2", turnId: "turn-1" });
});
test("a bound client remains the tool target while focus changes", async (t) => {
const session = new CanvasSession();
const first = connect(session, "first");
const second = connect(session, "second");
t.after(() => {
first.close();
second.close();
});
session.updateState(snapshot("canvas-first"), "first");
session.updateState(snapshot("canvas-second"), "second");
session.bindClient("first");
session.activateClient("second");
assert.equal(field(await session.callTool("canvas_get_state", {}), "projectId"), "canvas-first");
const result = session.callTool("canvas_create_text_node", { text: "bound" });
const call = first.event("tool_call");
assert.equal(second.event("tool_call"), undefined);
session.resolveResult("first", { requestId: String(field(call, "requestId")), result: { ok: true } });
assert.deepEqual(await result, { ok: true });
session.releaseClient("first");
assert.equal(field(await session.callTool("canvas_get_state", {}), "projectId"), "canvas-second");
});
test("closing the bound client falls back to the active client", async (t) => {
const session = new CanvasSession();
const first = connect(session, "first");
const second = connect(session, "second");
t.after(() => {
first.close();
second.close();
});
session.updateState(snapshot("canvas-first"), "first");
session.updateState(snapshot("canvas-second"), "second");
session.bindClient("first");
session.activateClient("second");
first.close();
assert.equal(field(await session.callTool("canvas_get_state", {}), "projectId"), "canvas-second");
const result = session.callTool("canvas_create_text_node", { text: "fallback" });
const call = second.event("tool_call");
session.resolveResult("second", { requestId: String(field(call, "requestId")), result: { ok: true } });
assert.deepEqual(await result, { ok: true });
});
function connect(session: CanvasSession, clientId: string) {
const response = new FakeSseResponse();
session.openEvents(new URL(`http://127.0.0.1/events?clientId=${clientId}`), response as unknown as ServerResponse);
+34 -4
View File
@@ -6,6 +6,7 @@ import { compactCanvasState, compactNode, isToolName, nextCanvasX, parseToolInpu
import type { CanvasNode, CanvasNodeType, CanvasSnapshot } from "./types.js";
type PendingRequest = { clientId: string; resolve: (value: unknown) => void; reject: (error: Error) => void };
export type CodexState = { busy: boolean; threadId: string; turnId: string };
const SITE_TOOLS = new Set<ToolName>([
"site_navigate",
@@ -26,14 +27,29 @@ export class CanvasSession {
private pending = new Map<string, PendingRequest>();
private canvasStates = new Map<string, CanvasSnapshot>();
private activeClientId = "";
private boundClientId = "";
private focusSequence = 0;
private codexState: CodexState = { busy: false, threadId: "", turnId: "" };
private get canvasState() {
return this.canvasStates.get(this.activeClientId) || null;
return this.canvasStates.get(this.targetClientId) || null;
}
private get targetClientId() {
return this.boundClientId || this.activeClientId;
}
health() {
return { ok: true, hasCanvas: Boolean(this.canvasState), clients: this.clients.size };
return { ok: true, hasCanvas: Boolean(this.canvasState), clients: this.clients.size, codexBusy: this.codexState.busy };
}
get codexBusy() {
return this.codexState.busy;
}
setCodexState(patch: Partial<CodexState>) {
this.codexState = { ...this.codexState, ...patch };
this.emitAll("codex_state", this.codexState);
}
openEvents(url: URL, res: ServerResponse) {
@@ -48,7 +64,7 @@ export class CanvasSession {
this.clientFocusOrder.set(clientId, ++this.focusSequence);
}
}
sendEvent(res, "hello", { ok: true, clientId });
sendEvent(res, "hello", { ok: true, clientId, codex: this.codexState });
const timer = setInterval(() => sendEvent(res, "ping", { time: Date.now() }), 15000);
res.on("close", () => {
clearInterval(timer);
@@ -56,6 +72,7 @@ export class CanvasSession {
this.clients.delete(clientId);
this.clientFocusOrder.delete(clientId);
this.canvasStates.delete(clientId);
if (this.boundClientId === clientId) this.boundClientId = "";
this.pending.forEach((item, requestId) => {
if (item.clientId !== clientId) return;
this.pending.delete(requestId);
@@ -77,6 +94,15 @@ export class CanvasSession {
this.clientFocusOrder.set(clientId, ++this.focusSequence);
}
bindClient(clientId: string) {
if (!this.clients.has(clientId)) throw new Error("当前网页未连接");
this.boundClientId = clientId;
}
releaseClient(clientId: string) {
if (this.boundClientId === clientId) this.boundClientId = "";
}
resolveResult(clientId: string, body: { requestId?: string; error?: string; result?: unknown }) {
const item = body.requestId ? this.pending.get(body.requestId) : null;
if (!item || !body.requestId || item.clientId !== clientId) return false;
@@ -89,6 +115,10 @@ export class CanvasSession {
this.clients.forEach((client) => sendEvent(client, type, payload));
}
emitThread(type: string, threadId: string, payload: Record<string, unknown> = {}) {
this.emitAll(type, { ...payload, threadId });
}
async callTool(name: unknown, rawInput: unknown) {
if (!isToolName(name)) throw new Error(`未知工具:${String(name)}`);
let tool: ToolName = name;
@@ -200,7 +230,7 @@ export class CanvasSession {
private async requestCanvasTool(name: ToolName, input: Record<string, unknown>) {
const requestId = crypto.randomUUID();
const clientId = this.activeClientId;
const clientId = this.targetClientId;
const client = this.clients.get(clientId);
if (!client) throw new Error("当前没有已连接画布");
sendEvent(client, "tool_call", { requestId, name, input });
+81 -20
View File
@@ -2,7 +2,7 @@ import express, { type NextFunction, type Request, type Response } from "express
import { DEFAULT_PORT, ensureSiteWorkspace, loadConfig, saveConfig, updateSiteWorkspace, type CanvasAgentConfig } from "./config.js";
import { CanvasSession } from "./canvas-session.js";
import { archiveCodexThread, interruptCodexTurn, listCodexThreads, readCodexThread, resumeCodexThread, runClaudeTurn, runCodexTurn, startCodexThread, summarizeCodexThread, verifyCodexThreadWorkspace, withAgentPrompt } from "./agents.js";
import { archiveCodexThread, interruptCodexTurn, isRecoverableThreadError, listCodexThreads, readCodexThread, resumeCodexThread, runClaudeTurn, runCodexTurn, startCodexThread, summarizeCodexThread, verifyCodexThreadWorkspace, withAgentPrompt } from "./agents.js";
import type { AgentAttachment } from "./types.js";
export function startHttpServer() {
@@ -12,7 +12,16 @@ export function startHttpServer() {
saveConfig(config);
const session = new CanvasSession();
const emit = (type: string, payload: unknown) => session.emitAll(type, payload);
const emit = (type: string, payload: unknown) => {
const data = payload && typeof payload === "object" && !Array.isArray(payload) ? payload as Record<string, unknown> : { value: payload };
const threadId = String(data.threadId || data.thread_id || ensureSiteWorkspace(config).activeThreadId || "");
threadId ? session.emitThread(type, threadId, data) : session.emitAll(type, data);
};
const setActiveThread = (activeThreadId: string, payload: Record<string, unknown> = {}) => {
const workspace = updateSiteWorkspace(config, { activeThreadId: activeThreadId || undefined });
session.emitThread("workspace_changed", activeThreadId, { ...payload, activeThreadId });
return workspace;
};
const app = express();
app.disable("x-powered-by");
app.use(express.json({ limit: "30mb" }));
@@ -52,48 +61,100 @@ export function startHttpServer() {
res.json({ ok: true, workspace, ...result });
}));
app.post("/agent/codex/threads/new", route(async (_req, res) => {
if (session.codexBusy) return res.status(409).json({ ok: false, error: "Codex 正在运行,请等待当前任务完成" });
const workspace = ensureSiteWorkspace(config);
const thread = await startCodexThread(emit, workspace.workspacePath);
const activeThreadId = String((thread as Record<string, unknown>).id || "");
updateSiteWorkspace(config, { activeThreadId });
res.json({ ok: true, workspace: { ...workspace, activeThreadId }, thread: summarizeCodexThread(thread), messages: [] });
const nextWorkspace = setActiveThread(activeThreadId, { emptyThread: true });
res.json({ ok: true, workspace: nextWorkspace, thread: summarizeCodexThread(thread), messages: [] });
}));
app.get("/agent/codex/threads/:threadId", route(async (req, res) => {
const workspace = ensureSiteWorkspace(config);
const threadId = routeParam(req.params.threadId);
res.json({ ok: true, workspace, ...(await readCodexThread(emit, threadId, workspace.workspacePath)) });
try {
res.json({ ok: true, workspace, ...(await readCodexThread(emit, threadId, workspace.workspacePath)) });
} catch (error) {
if (workspace.activeThreadId !== threadId || !isRecoverableThreadError(error)) throw error;
res.json({ ok: true, workspace, thread: { id: threadId, preview: "", cwd: workspace.workspacePath }, messages: [] });
}
}));
app.post("/agent/codex/threads/:threadId/resume", route(async (req, res) => {
if (session.codexBusy) return res.status(409).json({ ok: false, error: "Codex 正在运行,请等待当前任务完成" });
const workspace = ensureSiteWorkspace(config);
const threadId = routeParam(req.params.threadId);
const result = await resumeCodexThread(emit, threadId, workspace.workspacePath);
updateSiteWorkspace(config, { activeThreadId: threadId });
res.json({ ok: true, workspace: { ...workspace, activeThreadId: threadId }, ...result });
const nextWorkspace = setActiveThread(threadId);
res.json({ ok: true, workspace: nextWorkspace, ...result });
}));
app.post("/agent/codex/threads/:threadId/delete", route(async (req, res) => {
if (session.codexBusy) return res.status(409).json({ ok: false, error: "Codex 正在运行,请等待当前任务完成" });
const workspace = ensureSiteWorkspace(config);
const threadId = routeParam(req.params.threadId);
await archiveCodexThread(emit, threadId, workspace.workspacePath);
if (workspace.activeThreadId === threadId) updateSiteWorkspace(config, { activeThreadId: undefined });
setActiveThread(workspace.activeThreadId === threadId ? "" : workspace.activeThreadId || "");
res.json({ ok: true });
}));
app.post("/agent/codex/turn", route(async (req, res) => {
if (session.codexBusy) return res.status(409).json({ ok: false, error: "Codex 正在运行,请等待当前任务完成" });
const attachments = Array.isArray(req.body?.attachments) ? (req.body.attachments as AgentAttachment[]) : [];
const workspace = ensureSiteWorkspace(config);
let threadId = String(req.body?.threadId || workspace.activeThreadId || "");
if (!threadId) {
const thread = await startCodexThread(emit, workspace.workspacePath);
threadId = String((thread as Record<string, unknown>).id || "");
updateSiteWorkspace(config, { activeThreadId: threadId });
} else if (threadId !== workspace.activeThreadId) {
await verifyCodexThreadWorkspace(emit, threadId, workspace.workspacePath);
updateSiteWorkspace(config, { activeThreadId: threadId });
const prompt = String(req.body?.prompt || "");
if (!prompt.trim()) return res.status(400).json({ ok: false, error: "请输入任务内容" });
const clientId = String(req.body?.clientId || "");
session.setCodexState({ busy: true, threadId: String(req.body?.threadId || workspace.activeThreadId || ""), turnId: "" });
try {
let threadId = String(req.body?.threadId || workspace.activeThreadId || "");
let turnId = "";
if (!threadId) {
const thread = await startCodexThread(emit, workspace.workspacePath);
threadId = String((thread as Record<string, unknown>).id || "");
setActiveThread(threadId, { emptyThread: true });
} else if (threadId !== workspace.activeThreadId) {
await verifyCodexThreadWorkspace(emit, threadId, workspace.workspacePath);
setActiveThread(threadId);
}
const chatMessage = {
sourceClientId: clientId,
message: { id: String(req.body?.messageId || Date.now()), role: "user", text: String(req.body?.messageText || prompt || `发送了 ${attachments.length} 张图片`) },
};
let chatThreadId = "";
const turnEmit = (type: string, payload: unknown) => {
const data = payload && typeof payload === "object" && !Array.isArray(payload) ? payload as Record<string, unknown> : { value: payload };
session.emitThread(type, threadId, data);
};
void runCodexTurn(withAgentPrompt(prompt), turnEmit, attachments, {
threadId,
cwd: workspace.workspacePath,
appEmit: emit,
onStart: clientId ? () => session.bindClient(clientId) : undefined,
onThread: (actualThreadId) => {
if (actualThreadId !== threadId) {
threadId = actualThreadId;
setActiveThread(threadId, { emptyThread: true });
}
session.setCodexState({ busy: true, threadId, turnId: "" });
if (chatThreadId !== threadId) {
chatThreadId = threadId;
session.emitThread("chat_message", threadId, chatMessage);
}
},
onTurn: (actualTurnId) => {
turnId = actualTurnId;
session.setCodexState({ busy: true, threadId, turnId });
},
onFinish: () => {
if (clientId) session.releaseClient(clientId);
session.setCodexState({ busy: false, threadId, turnId });
},
});
res.json({ ok: true, threadId });
} catch (error) {
session.setCodexState({ busy: false, threadId: String(req.body?.threadId || workspace.activeThreadId || ""), turnId: "" });
throw error;
}
void runCodexTurn(withAgentPrompt(String(req.body?.prompt || "")), emit, attachments, { threadId, cwd: workspace.workspacePath });
res.json({ ok: true, threadId });
}));
app.post("/agent/codex/interrupt", (_req, res) => {
const ok = interruptCodexTurn();
app.post("/agent/codex/interrupt", (req, res) => {
const ok = interruptCodexTurn(String(req.body?.threadId || ""));
res.json({ ok });
});
app.post("/agent/claude/turn", (req, res) => {