feat(agent): enhance Codex thread management with workspace verification and improved thread loading

This commit is contained in:
HouYunFei
2026-06-16 13:22:15 +08:00
parent 7f0199da59
commit b4515ed85f
11 changed files with 54 additions and 561 deletions
+38 -14
View File
@@ -53,34 +53,39 @@ export async function startCodexThread(emit: AgentEmit, cwd?: string) {
export async function resumeCodexThread(emit: AgentEmit, threadId: string, cwd?: string) {
codexApp ||= await CodexAppClient.start(emit);
await loadCodexThread(emit, threadId, cwd, false);
const thread = await codexApp.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; all?: boolean; searchTerm?: string; limit?: number }) {
export async function listCodexThreads(emit: AgentEmit, options: { cwd: string; searchTerm?: string; limit?: number }) {
codexApp ||= await CodexAppClient.start(emit);
const result = await codexApp.listThreads({
limit: options.limit || 40,
sortKey: "updated_at",
sortDirection: "desc",
sourceKinds: ["cli", "vscode", "appServer", "exec"],
...(options.all ? {} : { cwd: options.cwd }),
cwd: options.cwd,
...(options.searchTerm ? { searchTerm: options.searchTerm } : {}),
});
const data = Array.isArray(field(result, "data")) ? (field(result, "data") as unknown[]).map(summarizeCodexThread) : [];
const data = Array.isArray(field(result, "data")) ? (field(result, "data") as unknown[]).map(summarizeCodexThread).filter((thread) => threadInWorkspace(thread, options.cwd)) : [];
return { data, nextCursor: field(result, "nextCursor") || null, backwardsCursor: field(result, "backwardsCursor") || null };
}
export async function readCodexThread(emit: AgentEmit, threadId: string) {
codexApp ||= await CodexAppClient.start(emit);
const result = await codexApp.readThread(threadId);
const thread = field(result, "thread") || {};
export async function readCodexThread(emit: AgentEmit, threadId: string, cwd?: string) {
const thread = await loadCodexThread(emit, threadId, cwd, true);
return { thread: summarizeCodexThread(thread), messages: threadMessages(thread) };
}
export async function archiveCodexThread(emit: AgentEmit, threadId: string) {
export async function verifyCodexThreadWorkspace(emit: AgentEmit, threadId: string, cwd: string) {
await loadCodexThread(emit, threadId, cwd, false);
}
export async function archiveCodexThread(emit: AgentEmit, threadId: string, cwd?: string) {
codexApp ||= await CodexAppClient.start(emit);
await loadCodexThread(emit, threadId, cwd, false);
await codexApp.archiveThread(threadId);
}
@@ -93,10 +98,11 @@ export function runClaudeTurn(prompt: string, emit: AgentEmit) {
async function ensureCodexThread(app: CodexAppClient, options: CodexRunOptions) {
if (options.threadId) {
if (codexThreadId !== options.threadId) {
const thread = await app.resumeThread(options.threadId, options.cwd);
codexThreadId = String(field(thread, "id") || options.threadId);
}
const result = await app.readThread(options.threadId, false);
assertThreadWorkspace(field(result, "thread") || {}, options.cwd);
const thread = await app.resumeThread(options.threadId, options.cwd);
assertThreadWorkspace(thread, options.cwd);
codexThreadId = String(field(thread, "id") || options.threadId);
return codexThreadId;
}
if (!codexThreadId) {
@@ -155,8 +161,8 @@ class CodexAppClient {
return this.request("thread/list", params);
}
readThread(threadId: string) {
return this.request("thread/read", { threadId, includeTurns: true });
readThread(threadId: string, includeTurns = true) {
return this.request("thread/read", { threadId, includeTurns });
}
archiveThread(threadId: string) {
@@ -291,6 +297,24 @@ function normalizeCodexNotification(method: string, params: Json): AgentEvent |
return null;
}
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 thread = field(result, "thread") || {};
assertThreadWorkspace(thread, cwd);
return thread;
}
function assertThreadWorkspace(thread: unknown, cwd?: string) {
if (!cwd || threadInWorkspace(thread, cwd)) return;
throw new Error("该 Codex 会话不属于当前画布工作空间");
}
function threadInWorkspace(thread: unknown, cwd: string) {
const threadCwd = String(field(thread, "cwd") || "");
return Boolean(threadCwd && path.resolve(threadCwd) === path.resolve(cwd));
}
function normalizeItem(item: unknown) {
const value = item && typeof item === "object" ? { ...(item as Json) } : {};
if (value.type === "agentMessage") value.type = "agent_message";
+6 -4
View File
@@ -2,7 +2,7 @@ import express, { type NextFunction, type Request, type Response } from "express
import { DEFAULT_PORT, ensureCanvasWorkspace, loadConfig, saveConfig, updateCanvasWorkspace, type CanvasAgentConfig } from "./config.js";
import { CanvasSession } from "./canvas-session.js";
import { archiveCodexThread, listCodexThreads, readCodexThread, resumeCodexThread, runClaudeTurn, runCodexTurn, startCodexThread, summarizeCodexThread, withAgentPrompt } from "./agents.js";
import { archiveCodexThread, listCodexThreads, readCodexThread, resumeCodexThread, runClaudeTurn, runCodexTurn, startCodexThread, summarizeCodexThread, verifyCodexThreadWorkspace, withAgentPrompt } from "./agents.js";
import type { AgentAttachment } from "./types.js";
export function startHttpServer() {
@@ -44,7 +44,7 @@ export function startHttpServer() {
});
app.get("/agent/codex/threads", route(async (req, res) => {
const workspace = ensureCanvasWorkspace(config, String(req.query.canvasId || ""));
const result = await listCodexThreads(emit, { cwd: workspace.workspacePath, all: req.query.all === "1", searchTerm: String(req.query.searchTerm || "") });
const result = await listCodexThreads(emit, { cwd: workspace.workspacePath, searchTerm: String(req.query.searchTerm || "") });
res.json({ ok: true, workspace, ...result });
}));
app.post("/agent/codex/threads/new", route(async (req, res) => {
@@ -55,8 +55,9 @@ export function startHttpServer() {
res.json({ ok: true, workspace: { ...workspace, activeThreadId }, thread: summarizeCodexThread(thread), messages: [] });
}));
app.get("/agent/codex/threads/:threadId", route(async (req, res) => {
const workspace = ensureCanvasWorkspace(config, String(req.query.canvasId || ""));
const threadId = routeParam(req.params.threadId);
res.json({ ok: true, ...(await readCodexThread(emit, threadId)) });
res.json({ ok: true, workspace, ...(await readCodexThread(emit, threadId, workspace.workspacePath)) });
}));
app.post("/agent/codex/threads/:threadId/resume", route(async (req, res) => {
const workspace = ensureCanvasWorkspace(config, String(req.body?.canvasId || ""));
@@ -68,7 +69,7 @@ export function startHttpServer() {
app.post("/agent/codex/threads/:threadId/delete", route(async (req, res) => {
const workspace = ensureCanvasWorkspace(config, String(req.body?.canvasId || ""));
const threadId = routeParam(req.params.threadId);
await archiveCodexThread(emit, threadId);
await archiveCodexThread(emit, threadId, workspace.workspacePath);
if (workspace.activeThreadId === threadId) updateCanvasWorkspace(config, workspace.canvasId, { activeThreadId: undefined });
res.json({ ok: true });
}));
@@ -81,6 +82,7 @@ export function startHttpServer() {
threadId = String((thread as Record<string, unknown>).id || "");
updateCanvasWorkspace(config, workspace.canvasId, { activeThreadId: threadId });
} else if (threadId !== workspace.activeThreadId) {
await verifyCodexThreadWorkspace(emit, threadId, workspace.workspacePath);
updateCanvasWorkspace(config, workspace.canvasId, { activeThreadId: threadId });
}
void runCodexTurn(withAgentPrompt(String(req.body?.prompt || "")), emit, attachments, { threadId, cwd: workspace.workspacePath });