fix(agent): synchronize conversation preparation state

This commit is contained in:
yu
2026-08-05 10:40:10 +08:00
parent 46f29f70ca
commit 9bccd0ff1a
13 changed files with 511 additions and 124 deletions
+112 -45
View File
@@ -4,7 +4,7 @@ import path from "node:path";
import express, { type NextFunction, type Request, type Response } from "express";
import { runClaudeTurn } from "../agent/claude.js";
import { archiveCodexThread, CodexSkillLookupError, configureCodexSkill, generateCodexSkillDraft, interruptCodexTurn, listCodexModels, listCodexSkills, listCodexThreads, readCodexThread, resolveCodexApproval, resolveCodexSkill, resumeCodexThread, runCodexTurn, startCodexThread, summarizeCodexThread } from "../agent/codex.js";
import { archiveCodexThread, CodexSkillLookupError, configureCodexSkill, generateCodexSkillDraft, interruptCodexTurn, isRecoverableThreadError, listCodexModels, listCodexSkills, listCodexThreads, readCodexThread, resolveCodexApproval, resolveCodexSkill, resumeCodexThread, runCodexTurn, startCodexThread, summarizeCodexThread } from "../agent/codex.js";
import type { CodexReasoningEffort, CodexSkillSelector } from "../agent/codex-protocol.js";
import type { AgentAttachment, AgentPermissionMode } from "../agent/types.js";
import { AGENT_PROTOCOL_VERSION, CanvasSession } from "../canvas/session.js";
@@ -20,8 +20,9 @@ export function startHttpServer() {
config.url = `http://127.0.0.1:${port}`;
saveConfig(config);
const session = new CanvasSession();
const skillStore = new SkillStore(ensureSiteWorkspace(config).workspacePath);
const initialWorkspace = ensureSiteWorkspace(config);
const session = new CanvasSession(initialWorkspace.activeThreadId || "");
const skillStore = new SkillStore(initialWorkspace.workspacePath);
/** 将 Agent 事件广播到所属线程或全部网页。 */
const emit = (type: string, payload: unknown) => {
const value = payload && typeof payload === "object" && !Array.isArray(payload) ? payload as Record<string, unknown> : { value: payload };
@@ -29,6 +30,10 @@ export function startHttpServer() {
session.emitAll(type, value);
return;
}
if (type === "agent_bootstrap" && value.phase === "preheat") {
if (value.type === "mcp.startup") session.updateConversationMcp(String(value.name || ""), startupStatus(value.status), String(value.error || "") || null, String(value.failureReason || "") || null);
if (value.type === "mcp.complete") session.completeConversationMcpInventory(mcpInventory(value.services));
}
const scope = session.codexBusy ? session.codexEventScope : { threadId: "", turnId: "", sourceClientId: "" };
const threadId = String(value.threadId || value.thread_id || scope.threadId || ensureSiteWorkspace(config).activeThreadId || "");
const turnId = String(value.turnId || value.turn_id || scope.turnId || "");
@@ -43,10 +48,11 @@ export function startHttpServer() {
threadId ? session.emitThread(type, threadId, data) : session.emitAll(type, data);
};
/** 保存并广播当前站点工作空间的活跃线程。 */
const setActiveThread = (activeThreadId: string, payload: Record<string, unknown> = {}) => {
const setActiveThread = (activeThreadId: string, payload: Record<string, unknown> = {}, preserveConversation = false) => {
const workspace = updateSiteWorkspace(config, { activeThreadId: activeThreadId || undefined });
if (!preserveConversation) session.activateConversation(activeThreadId, String(payload.sourceClientId || "") || undefined);
if (!session.codexBusy && session.codexThreadId !== activeThreadId) session.setCodexState({ threadId: activeThreadId, turnId: "" });
session.emitThread("workspace_changed", activeThreadId, { ...payload, activeThreadId });
session.emitThread("workspace_changed", activeThreadId, { ...payload, activeThreadId, conversation: session.conversationStateSnapshot });
return workspace;
};
let draftThreadStart: ReturnType<typeof startCodexThread> | null = null;
@@ -54,19 +60,45 @@ export function startHttpServer() {
const prepareDraftThread = (clientId: string, permission: AgentPermissionMode) => {
if (draftThreadStart) return draftThreadStart;
const workspace = ensureSiteWorkspace(config);
emit("agent_bootstrap", { type: "codex.preparing", sourceClientId: clientId });
const start = startCodexThread(emit, workspace.workspacePath, permission);
draftThreadStart = start;
void start.then((thread) => {
if (draftThreadStart !== start) return;
draftThreadStart = null;
const threadId = String((thread as Record<string, unknown>).id || "");
if (threadId && !ensureSiteWorkspace(config).activeThreadId) setActiveThread(threadId, { emptyThread: true, draftThread: true, sourceClientId: clientId });
}).catch((error) => {
if (draftThreadStart === start) draftThreadStart = null;
emit("agent_bootstrap", { type: "codex.prepare_failed", sourceClientId: clientId, error: error instanceof Error ? error.message : String(error) });
});
return start;
let prepared!: ReturnType<typeof startCodexThread>;
prepared = (async () => {
emit("agent_bootstrap", { type: "codex.preparing", sourceClientId: clientId });
try {
const thread = await startCodexThread(emit, workspace.workspacePath, permission, true);
if (draftThreadStart !== prepared) return thread;
const threadId = String((thread as Record<string, unknown>).id || "");
if (threadId && !ensureSiteWorkspace(config).activeThreadId) {
session.completeConversationPreparation(threadId);
setActiveThread(threadId, { emptyThread: true, draftThread: true, sourceClientId: clientId }, true);
}
return thread;
} catch (error) {
if (draftThreadStart === prepared) {
const text = error instanceof Error ? error.message : String(error);
session.failConversationPreparation(text);
emit("agent_bootstrap", { type: "codex.prepare_failed", sourceClientId: clientId, error: text });
}
throw error;
} finally {
if (draftThreadStart === prepared) draftThreadStart = null;
}
})();
draftThreadStart = prepared;
return prepared;
};
/** 恢复已有线程并等待完整 MCP 清单,供启动恢复和手动切换共用。 */
const prepareExistingThread = async (threadId: string, clientId = "", permission: AgentPermissionMode = "request") => {
const workspace = ensureSiteWorkspace(config);
session.beginConversation({ conversationId: threadId, threadId, sourceClientId: clientId || undefined });
emit("agent_bootstrap", { type: "codex.preparing", threadId, sourceClientId: clientId || undefined });
const result = await resumeCodexThread(emit, threadId, workspace.workspacePath, permission, true);
session.completeConversationPreparation(threadId);
return result;
};
const failPreparedConversation = (error: unknown, threadId: string, clientId = "") => {
const text = error instanceof Error ? error.message : String(error);
session.failConversationPreparation(text);
emit("agent_bootstrap", { type: "codex.prepare_failed", threadId, sourceClientId: clientId || undefined, error: text });
};
const app = express();
app.disable("x-powered-by");
@@ -133,7 +165,7 @@ export function startHttpServer() {
app.post("/api/tools", route(async (req, res) => res.json({ ok: true, result: await session.callTool(req.body?.name, req.body?.input || {}) })));
app.get("/agent/codex/workspace", (_req, res) => {
const workspace = ensureSiteWorkspace(config);
res.json({ ok: true, workspace });
res.json({ ok: true, workspace, conversation: session.conversationStateSnapshot });
});
app.get("/agent/codex/models", route(async (_req, res) => res.json({ ok: true, ...(await listCodexModels(emit)) })));
app.get("/agent/codex/skills", route(async (req, res) => {
@@ -204,25 +236,26 @@ export function startHttpServer() {
app.get("/agent/codex/threads", route(async (req, res) => {
const workspace = ensureSiteWorkspace(config);
const result = await listCodexThreads(emit, { cwd: workspace.workspacePath, searchTerm: String(req.query.searchTerm || "") });
res.json({ ok: true, workspace, ...result });
res.json({ ok: true, workspace, conversation: session.conversationStateSnapshot, ...result });
}));
app.post("/agent/codex/threads/new", codexMutation(async (req, res) => {
const workspace = ensureSiteWorkspace(config);
const thread = await startCodexThread(emit, workspace.workspacePath, permissionMode(req.body?.permissionMode));
const activeThreadId = String((thread as Record<string, unknown>).id || "");
const nextWorkspace = setActiveThread(activeThreadId, { emptyThread: true, sourceClientId: String(req.body?.clientId || "") });
res.json({ ok: true, workspace: nextWorkspace, thread: summarizeCodexThread(thread), messages: [] });
}));
app.post("/agent/codex/threads/reset", codexMutation((req, res) => {
const clientId = String(req.body?.clientId || "");
const workspace = setActiveThread("", { emptyThread: true, draftThread: true, sourceClientId: clientId });
void prepareDraftThread(clientId, permissionMode(req.body?.permissionMode));
res.json({ ok: true, workspace });
session.beginConversation({ sourceClientId: clientId });
setActiveThread("", { emptyThread: true, draftThread: true, sourceClientId: clientId }, true);
const thread = await prepareDraftThread(clientId, permissionMode(req.body?.permissionMode));
res.json({ ok: true, workspace: ensureSiteWorkspace(config), conversation: session.conversationStateSnapshot, thread: summarizeCodexThread(thread), messages: [] });
}));
app.post("/agent/codex/threads/reset", codexMutation(async (req, res) => {
const clientId = String(req.body?.clientId || "");
session.beginConversation({ sourceClientId: clientId });
setActiveThread("", { emptyThread: true, draftThread: true, sourceClientId: clientId }, true);
await prepareDraftThread(clientId, permissionMode(req.body?.permissionMode));
res.json({ ok: true, workspace: ensureSiteWorkspace(config), conversation: session.conversationStateSnapshot });
}));
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)) });
res.json({ ok: true, workspace, conversation: session.conversationStateSnapshot, ...(await readCodexThread(emit, threadId, workspace.workspacePath)) });
}));
app.post("/agent/codex/history/ack", (req, res) => {
const threadId = String(req.body?.threadId || "");
@@ -231,18 +264,23 @@ export function startHttpServer() {
res.json({ ok: true });
});
app.post("/agent/codex/threads/:threadId/resume", codexMutation(async (req, res) => {
const workspace = ensureSiteWorkspace(config);
const threadId = routeParam(req.params.threadId);
const result = await resumeCodexThread(emit, threadId, workspace.workspacePath, permissionMode(req.body?.permissionMode));
const nextWorkspace = setActiveThread(threadId, { sourceClientId: String(req.body?.clientId || "") });
res.json({ ok: true, workspace: nextWorkspace, ...result });
const clientId = String(req.body?.clientId || "");
try {
const result = await prepareExistingThread(threadId, clientId, permissionMode(req.body?.permissionMode));
const nextWorkspace = setActiveThread(threadId, { sourceClientId: clientId }, true);
res.json({ ok: true, workspace: nextWorkspace, conversation: session.conversationStateSnapshot, ...result });
} catch (error) {
failPreparedConversation(error, threadId, clientId);
throw error;
}
}));
app.post("/agent/codex/threads/:threadId/delete", codexMutation(async (req, res) => {
const workspace = ensureSiteWorkspace(config);
const threadId = routeParam(req.params.threadId);
await archiveCodexThread(emit, threadId, workspace.workspacePath);
setActiveThread(workspace.activeThreadId === threadId ? "" : workspace.activeThreadId || "", { sourceClientId: String(req.body?.clientId || "") });
res.json({ ok: true });
const nextWorkspace = setActiveThread(workspace.activeThreadId === threadId ? "" : workspace.activeThreadId || "", { sourceClientId: String(req.body?.clientId || "") });
res.json({ ok: true, workspace: nextWorkspace, conversation: session.conversationStateSnapshot });
}));
app.post("/agent/codex/turn", codexMutation(async (req, res) => {
const attachments = Array.isArray(req.body?.attachments) ? (req.body.attachments as AgentAttachment[]) : [];
@@ -253,7 +291,15 @@ export function startHttpServer() {
if (!clientId || !session.hasClient(clientId)) return res.status(409).json({ ok: false, error: "发起任务的网页已断开,请重新连接后再试" });
const requestedThreadId = String(req.body?.threadId || "");
const activeThreadId = workspace.activeThreadId || "";
if (requestedThreadId !== activeThreadId) return res.status(409).json({ ok: false, error: "当前会话已在其他页面切换,请同步后重试" });
const conversation = session.conversationStateSnapshot;
const requestedConversationId = String(req.body?.conversationId || "");
const expectedRevision = Number(req.body?.expectedRevision || 0);
if (requestedThreadId !== activeThreadId || conversation.threadId !== activeThreadId || (requestedConversationId && requestedConversationId !== conversation.conversationId) || (expectedRevision && expectedRevision !== conversation.revision)) {
return res.status(409).json({ ok: false, code: "CONVERSATION_STALE", error: "当前会话已切换,已同步最新状态,请确认后重试", state: conversation });
}
if (!activeThreadId || !["ready", "warning"].includes(conversation.status)) {
return res.status(409).json({ ok: false, code: "CONVERSATION_NOT_READY", error: "Codex 对话仍在初始化,请等待 MCP 加载完成", state: conversation });
}
const model = String(req.body?.model || "") || undefined;
const effort = reasoningEffort(req.body?.effort);
const skill = req.body?.skill === undefined ? undefined : await resolveCodexSkill(emit, workspace.workspacePath, skillSelector(req.body.skill), true);
@@ -262,15 +308,10 @@ export function startHttpServer() {
let threadId = activeThreadId;
logger.info("Codex turn accepted", { threadId: req.body?.threadId, model: model || "default", reasoningEffort: effort || "default", promptLength: prompt.length, attachmentCount: attachments.length });
session.bindClient(clientId);
session.markConversationRunning(threadId);
session.setCodexState({ busy: true, threadId, turnId: "" });
try {
let turnId = "";
if (!threadId) {
const thread = await prepareDraftThread(clientId, permissionMode(req.body?.permissionMode));
threadId = String((thread as Record<string, unknown>).id || "");
setActiveThread(threadId, { emptyThread: true, sourceClientId: clientId });
}
session.setCodexState({ busy: true, threadId, turnId: "" });
const attachmentRefs = session.setTurnAttachments(clientId, attachments);
session.emitThread("chat_message", threadId, {
sourceClientId: clientId,
@@ -307,6 +348,7 @@ export function startHttpServer() {
threadId = actualThreadId;
setActiveThread(threadId, { emptyThread: true, sourceClientId: clientId });
}
session.markConversationRunning(threadId);
session.setCodexState({ busy: true, threadId, turnId: "" });
if (threadChanged) {
session.emitThread("chat_message", threadId, {
@@ -333,12 +375,14 @@ export function startHttpServer() {
session.clearTurnAttachments(clientId);
if (clientId) session.releaseClient(clientId);
session.setCodexState({ busy: false, threadId, turnId });
session.finishConversationRun(threadId);
},
});
res.json({ ok: true, threadId });
} catch (error) {
session.releaseClient(clientId);
session.setCodexState({ busy: false, threadId, turnId: "" });
session.finishConversationRun(threadId);
throw error;
}
}));
@@ -346,7 +390,7 @@ export function startHttpServer() {
/** 将 Codex 写操作串行化,避免多窗口在异步请求期间交叉修改会话。 */
function codexMutation(handler: (req: Request, res: Response) => unknown | Promise<unknown>) {
return route(async (req, res) => {
if (!session.beginCodexMutation()) return res.status(409).json({ ok: false, error: "Codex 正在运行或正在切换会话,请稍后重试" });
if (!session.beginCodexMutation()) return res.status(409).json({ ok: false, code: "CONVERSATION_BUSY", error: "Codex 正在运行或正在切换会话,请稍后重试", state: session.conversationStateSnapshot });
try {
return await handler(req, res);
} finally {
@@ -385,6 +429,15 @@ export function startHttpServer() {
console.log("Remove manually added MCP: codex mcp remove infinite-canvas");
if (logger.enabled) console.log(`Debug log: ${logger.filePath}`);
logger.info("Canvas Agent started", { url: config.url, workspace: ensureSiteWorkspace(config).workspacePath, debugLog: logger.filePath });
const activeThreadId = initialWorkspace.activeThreadId || "";
if (activeThreadId && session.beginCodexMutation()) {
void prepareExistingThread(activeThreadId).catch(async (error) => {
if (!isRecoverableThreadError(error)) return failPreparedConversation(error, activeThreadId);
session.beginConversation();
setActiveThread("", { emptyThread: true, draftThread: true }, true);
await prepareDraftThread("", "request");
}).finally(() => session.endCodexMutation()).catch(() => undefined);
}
});
}
@@ -406,6 +459,20 @@ function reasoningEffort(value: unknown): CodexReasoningEffort | undefined {
return value === "minimal" || value === "low" || value === "medium" || value === "high" || value === "xhigh" || value === "max" || value === "ultra" ? value : undefined;
}
function startupStatus(value: unknown): "starting" | "ready" | "failed" | "cancelled" {
return value === "starting" || value === "ready" || value === "failed" ? value : "cancelled";
}
function mcpInventory(value: unknown) {
if (!Array.isArray(value)) return [];
return value.flatMap((item) => {
if (!item || typeof item !== "object" || Array.isArray(item)) return [];
const server = item as Record<string, unknown>;
const name = String(server.name || "");
return name ? [{ name, authStatus: String(server.authStatus || "") || undefined }] : [];
});
}
/** 读取浏览器提交的 Skill 选择器;真实路径随后必须通过原生列表校验。 */
function skillSelector(value: unknown): CodexSkillSelector {
const selector = value && typeof value === "object" && !Array.isArray(value) ? value as Record<string, unknown> : {};