feat(canvas): restructure source directory for Canvas Agent, enhancing code maintainability

This commit is contained in:
HouYunFei
2026-07-29 10:47:06 +08:00
parent 92fd0ce129
commit 40c47dd7ff
20 changed files with 3095 additions and 777 deletions
+48
View File
@@ -0,0 +1,48 @@
import { spawn } from "node:child_process";
import { AGENT_PROMPT } from "../config.js";
import { errorMessage } from "../utils/value.js";
import type { AgentEmit } from "./types.js";
/** 使用 Claude CLI 执行一次带 Canvas Agent 工具的任务。 */
export function runClaudeTurn(prompt: string, emit: AgentEmit) {
const fullPrompt = withAgentPrompt(prompt);
if (!fullPrompt) return;
const child = spawnAgent("claude", ["-p", "--output-format", "stream-json", "--verbose", "--include-partial-messages", "--allowedTools", "mcp__infinite-canvas__*", fullPrompt], emit);
if (child) pipeJsonLines(child, emit, "claude");
}
/** 为 Claude CLI 请求拼接 Canvas Agent 指令。 */
function withAgentPrompt(prompt: string) {
return prompt.trim() ? `${AGENT_PROMPT}\n\n用户请求:${prompt}` : "";
}
/** 将 Claude CLI 的 JSON Lines 输出转换为 Agent 事件。 */
function pipeJsonLines(child: ReturnType<typeof spawn>, emit: AgentEmit, agent: string) {
let out = "";
child.stdout?.on("data", (chunk) => {
out += chunk.toString();
const lines = out.split(/\r?\n/);
out = lines.pop() || "";
lines.filter(Boolean).forEach((line) => {
try {
emit("agent_event", { agent, ...JSON.parse(line) });
} catch {
emit("agent_event", { agent, type: "raw", text: line });
}
});
});
child.stderr?.on("data", (chunk) => emit("agent_log", { text: chunk.toString() }));
child.on("error", (error) => emit("agent_error", { message: error.message }));
child.on("close", (code) => emit("agent_done", { agent, code }));
}
/** 启动外部 Agent CLI,并将同步启动异常转换为事件。 */
function spawnAgent(name: string, args: string[], emit: AgentEmit) {
try {
return spawn(name, args, { stdio: ["ignore", "pipe", "pipe"], shell: process.platform === "win32", windowsHide: true });
} catch (error) {
emit("agent_error", { message: errorMessage(error) });
return null;
}
}
+308
View File
@@ -0,0 +1,308 @@
import { spawn, type ChildProcess } from "node:child_process";
import { createRequire } from "node:module";
import path from "node:path";
import { fileURLToPath } from "node:url";
import { VERSION } from "../config.js";
import { logger } from "../utils/logger.js";
import { field, type JsonRecord } from "../utils/value.js";
import type { AgentEmit } from "./types.js";
type AgentEvent = JsonRecord & { type: string; usage?: unknown };
type PendingRequest = { resolve: (value: unknown) => void; reject: (error: Error) => void };
const canvasAgentMcp = canvasAgentMcpCommand();
const require = createRequire(import.meta.url);
/** 封装 Codex app-server 的 JSON-RPC 通信与事件转换。 */
export class CodexAppClient {
private nextId = 1;
private buffer = "";
private currentThreadId = "";
private textByItem = new Map<string, string>();
private lastUsage: unknown = null;
private pending = new Map<number, PendingRequest>();
private activeTurns = new Map<string, PendingRequest>();
private completedTurns = new Map<string, Error | null>();
/** 保存 app-server 子进程和事件出口。 */
private constructor(private child: ChildProcess, private emit: AgentEmit) {}
/** 启动并初始化 Codex app-server。 */
static async start(emit: AgentEmit, onExit: () => void) {
logger.info("Starting Codex app-server", { executable: process.execPath, codex: codexBin() });
const child = spawn(process.execPath, [codexBin(), "app-server", "--stdio"], { stdio: ["pipe", "pipe", "pipe"], windowsHide: true });
const client = new CodexAppClient(child, emit);
child.stdout?.on("data", (chunk) => client.read(chunk.toString()));
child.stderr?.on("data", (chunk) => {
const text = chunk.toString();
logger.warn("Codex app-server stderr", { text });
emit("agent_log", { text });
});
child.on("error", (error) => {
logger.error("Codex app-server process error", error);
emit("agent_error", { message: error.message });
});
child.on("exit", (code) => {
logger.warn("Codex app-server exited", { code });
client.failAll(`Codex app-server exited: ${code ?? 0}`);
onExit();
emit("agent_log", { text: `Codex app-server exited: ${code ?? 0}` });
});
await client.request("initialize", { clientInfo: { name: "canvas-agent", title: "Infinite Canvas Agent", version: VERSION }, capabilities: { experimentalApi: true, requestAttestation: false } });
client.notify("initialized");
return client;
}
/** 创建新的 Codex 线程。 */
async startThread(cwd?: string) {
const result = await this.request("thread/start", { approvalPolicy: "never", sandbox: "workspace-write", config: codexConfig(), ...(cwd ? { cwd } : {}), threadSource: "user" });
const thread = field(result, "thread") as JsonRecord | undefined;
const id = String(field(thread, "id") || "");
if (!id) throw new Error("Codex app-server 没有返回 thread id");
return thread || {};
}
/** 恢复已有 Codex 线程。 */
async resumeThread(threadId: string, cwd?: string) {
const result = await this.request("thread/resume", { threadId, approvalPolicy: "never", sandbox: "workspace-write", config: codexConfig(), ...(cwd ? { cwd } : {}) });
const thread = field(result, "thread") as JsonRecord | undefined;
const id = String(field(thread, "id") || "");
if (!id) throw new Error("Codex app-server 没有返回 thread id");
return thread || {};
}
/** 查询 Codex 线程列表。 */
listThreads(params: JsonRecord) {
return this.request("thread/list", params);
}
/** 读取指定 Codex 线程。 */
readThread(threadId: string, includeTurns = true) {
return this.request("thread/read", { threadId, includeTurns });
}
/** 归档指定 Codex 线程。 */
archiveThread(threadId: string) {
return this.request("thread/archive", { threadId });
}
/** 启动一个 Codex turn 并等待完成通知。 */
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");
this.currentThreadId = threadId;
onTurn?.(turnId);
const completed = this.completedTurns.get(turnId);
if (this.completedTurns.has(turnId)) {
this.completedTurns.delete(turnId);
if (completed) throw completed;
return;
}
await new Promise((resolve, reject) => this.activeTurns.set(turnId, { resolve, reject }));
}
/** 中断当前正在运行的 Codex turn。 */
interruptCurrentTurn() {
if (this.activeTurns.size === 0) return false;
try {
logger.warn("Interrupting active Codex turn", { threadId: this.currentThreadId, activeTurns: this.activeTurns.size });
this.child.kill("SIGINT");
return true;
} catch {
return false;
}
}
/** 发送 JSON-RPC 请求并保存待处理 Promise。 */
private request(method: string, params: unknown) {
const id = this.nextId++;
this.write({ id, method, params });
return new Promise((resolve, reject) => this.pending.set(id, { resolve, reject }));
}
/** 发送无需响应的 JSON-RPC 通知。 */
private notify(method: string, params?: unknown) {
this.write(params === undefined ? { method } : { method, params });
}
/** 将 JSON-RPC 消息写入 app-server 标准输入。 */
private write(value: unknown) {
const method = String(field(value, "method") || "");
const params = field(value, "params");
if (method) logger.debug(`Codex ${method}`, { id: field(value, "id"), threadId: field(params, "threadId") });
this.child.stdin?.write(`${JSON.stringify(value)}\n`);
}
/** 按行解析 app-server 标准输出。 */
private read(chunk: string) {
this.buffer += chunk;
const lines = this.buffer.split(/\r?\n/);
this.buffer = lines.pop() || "";
lines.filter(Boolean).forEach((line) => {
try {
this.handle(JSON.parse(line) as JsonRecord);
} catch (error) {
logger.warn("Invalid Codex app-server output", { error, line });
this.emit("agent_log", { text: line });
}
});
}
/** 分派单条 JSON-RPC 响应、请求或通知。 */
private handle(message: JsonRecord) {
const id = Number(message.id);
if (message.error && this.pending.has(id)) {
const error = String(field(message.error, "message") || "Codex request failed");
logger.warn("Codex request failed", { id, error });
return this.reject(id, error);
}
if (this.pending.has(id)) return this.resolve(id, message.result);
if (typeof message.method === "string" && "id" in message) return this.answerServerRequest(message);
if (typeof message.method === "string") this.handleNotification(message.method, (message.params || {}) as JsonRecord);
}
/** 转换并广播 app-server 通知。 */
private handleNotification(method: string, params: JsonRecord) {
if (method === "item/agentMessage/delta") return this.emitDelta(params);
if (method === "thread/tokenUsage/updated") {
this.lastUsage = normalizeUsage(params);
this.emit("agent_event", { agent: "codex", type: "usage.updated", usage: this.lastUsage, ...codexEventScope(params) });
return;
}
const event = normalizeCodexNotification(method, params);
if (!event) return;
if (event.type === "item.completed") {
const item = field(event, "item") as JsonRecord | undefined;
const id = String(field(item, "id") || "");
const streamedText = this.textByItem.get(id);
if (item?.type === "agent_message" && streamedText && !item.text) item.text = streamedText;
if (id) this.textByItem.delete(id);
}
if (event.type === "turn.completed") event.usage = this.lastUsage;
this.emit("agent_event", { agent: "codex", ...event });
if (event.type === "turn.completed") {
const turnId = String(field(params, "turnId") || field(field(params, "turn"), "id") || "");
const pending = this.activeTurns.get(turnId);
const error = field(field(params, "turn"), "error");
if (pending) {
this.activeTurns.delete(turnId);
error ? pending.reject(new Error(String(field(error, "message") || "Codex turn failed"))) : pending.resolve(event);
} else if (turnId) {
this.completedTurns.set(turnId, error ? new Error(String(field(error, "message") || "Codex turn failed")) : null);
}
if (this.activeTurns.size === 0) this.currentThreadId = "";
this.emit("agent_done", { agent: "codex", usage: event.usage, ...codexEventScope(params) });
}
}
/** 合并并广播 Agent 文本增量。 */
private emitDelta(params: JsonRecord) {
const id = String(field(params, "itemId") || "");
const text = `${this.textByItem.get(id) || ""}${String(field(params, "delta") || "")}`;
this.textByItem.set(id, text);
this.emit("agent_event", { agent: "codex", type: "item.updated", item: { id, type: "agent_message", text }, ...codexEventScope(params) });
}
/** 自动回复 app-server 发起的授权或交互请求。 */
private answerServerRequest(message: JsonRecord) {
const method = String(message.method);
const result = method === "mcpServer/elicitation/request" ? { action: "accept", content: {}, _meta: null } : { decision: "decline" };
this.write({ id: message.id, result });
this.emit("agent_event", { agent: "codex", type: "server.request", method, params: message.params, result });
}
/** 完成指定 JSON-RPC 请求。 */
private resolve(id: number, result: unknown) {
const pending = this.pending.get(id);
if (pending) (this.pending.delete(id), pending.resolve(result));
}
/** 拒绝指定 JSON-RPC 请求。 */
private reject(id: number, message: string) {
const pending = this.pending.get(id);
if (pending) (this.pending.delete(id), pending.reject(new Error(message)));
}
/** 拒绝进程退出时仍未完成的请求与 turn。 */
private failAll(message: string) {
[...this.pending.values(), ...this.activeTurns.values()].forEach((item) => item.reject(new Error(message)));
this.pending.clear();
this.activeTurns.clear();
this.currentThreadId = "";
}
}
/** 生成 Codex 调用 Canvas Agent MCP 的启动命令。 */
function canvasAgentMcpCommand() {
const current = process.argv.find((arg) => /index\.(t|j)s$/.test(arg)) || "";
const entry = path.resolve(current || fileURLToPath(new URL("../index.js", import.meta.url)));
const tsx = path.join(path.dirname(entry), "..", "node_modules", "tsx", "dist", "cli.mjs");
return entry.endsWith(".ts") ? { command: process.execPath, args: [tsx, entry, "mcp"] } : { command: process.execPath, args: [entry, "mcp"] };
}
/** 生成 Codex app-server 使用的 MCP 配置。 */
function codexConfig() {
return { mcp_servers: { "infinite-canvas": { command: canvasAgentMcp.command, args: canvasAgentMcp.args, default_tools_approval_mode: "approve", startup_timeout_sec: 20, tool_timeout_sec: 90 } } };
}
/** 将文本和本地图片转换为 Codex turn 输入。 */
function codexInput(prompt: string, images: string[]) {
return [{ type: "text", text: prompt, text_elements: [] }, ...images.map((file) => ({ type: "localImage", path: file }))];
}
/** 将 app-server 通知转换为前端使用的 Agent 事件。 */
function normalizeCodexNotification(method: string, params: JsonRecord): AgentEvent | null {
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, duration_ms: field(field(params, "turn"), "durationMs"), ...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;
}
/** 提取 Codex 事件所属的线程和 turn。 */
function codexEventScope(params: JsonRecord) {
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 } : {}) };
}
/** 统一 app-server item 的类型和参数格式。 */
function normalizeItem(item: unknown) {
const value = item && typeof item === "object" ? { ...(item as JsonRecord) } : {};
if (value.type === "agentMessage") value.type = "agent_message";
if (value.type === "mcpToolCall") value.type = "mcp_tool_call";
if (value.type === "agent_message" && typeof value.id === "string") value.text = String(value.text || "");
if ("arguments" in value) value.arguments = parseMaybeJson(value.arguments);
return value;
}
/** 将 Codex token usage 转换为前端字段。 */
function normalizeUsage(params: JsonRecord) {
const last = field(field(params, "tokenUsage"), "last") as JsonRecord | undefined;
return {
input_tokens: field(last, "inputTokens"),
cached_input_tokens: field(last, "cachedInputTokens"),
output_tokens: field(last, "outputTokens"),
reasoning_output_tokens: field(last, "reasoningOutputTokens"),
};
}
/** 尝试将字符串解析为 JSON,失败时保留原值。 */
function parseMaybeJson(value: unknown) {
if (typeof value !== "string") return value;
try {
return JSON.parse(value);
} catch {
return value;
}
}
/** 定位当前依赖中 Codex CLI 的执行文件。 */
function codexBin() {
return path.join(path.dirname(require.resolve("@openai/codex/package.json")), "bin", "codex.js");
}
+100
View File
@@ -0,0 +1,100 @@
import { field } from "../utils/value.js";
type AgentHistoryMessage = { id: string; role: "user" | "assistant" | "tool" | "error"; title?: string; text: string; detail?: unknown; streamId?: string };
/** 将 Codex 线程转换为列表展示所需的摘要。 */
export function summarizeCodexThread(thread: unknown) {
return {
id: String(field(thread, "id") || ""),
sessionId: String(field(thread, "sessionId") || ""),
preview: displayUserText(String(field(thread, "preview") || "")),
name: stringOrNull(field(thread, "name")),
cwd: String(field(thread, "cwd") || ""),
status: String(field(thread, "status") || ""),
source: field(thread, "source"),
threadSource: field(thread, "threadSource"),
createdAt: Number(field(thread, "createdAt") || 0),
updatedAt: Number(field(thread, "updatedAt") || 0),
};
}
/** 将 Codex turn items 转换为网页聊天历史。 */
export function threadMessages(thread: unknown): AgentHistoryMessage[] {
const turns = arrayValue(field(thread, "turns"));
const messages: AgentHistoryMessage[] = [];
turns.forEach((turn, turnIndex) => {
arrayValue(field(turn, "items")).forEach((item, itemIndex) => {
const type = String(field(item, "type") || "");
const id = String(field(item, "id") || `${turnIndex}-${itemIndex}`);
if (type === "userMessage") {
const text = displayUserText(userInputText(field(item, "content")));
if (text) messages.push({ id, role: "user", text });
}
if (type === "agentMessage") {
const text = String(field(item, "text") || "").trim();
if (text) messages.push({ id, role: "assistant", title: "Codex", text });
}
if (type === "mcpToolCall") {
const tool = String(field(item, "tool") || "工具调用");
const error = field(field(item, "error"), "message");
messages.push({ id, role: error ? "error" : "tool", title: toolName(tool), text: error ? String(error) : `${toolName(tool)} ${String(field(item, "status") || "完成")}`, detail: item });
}
if (type === "commandExecution") {
const command = String(field(item, "command") || "").trim();
if (command) messages.push({ id, role: "tool", title: "命令", text: command, detail: { cwd: field(item, "cwd"), status: field(item, "status"), exitCode: field(item, "exitCode") } });
}
if (type === "fileChange") messages.push({ id, role: "tool", title: "文件变更", text: "Codex 修改了文件", detail: item });
});
});
return messages.filter((item) => item.text).slice(-120);
}
/** 提取用户输入条目中的文本与附件占位信息。 */
function userInputText(content: unknown) {
return arrayValue(content)
.map((item) => {
const type = String(field(item, "type") || "");
if (type === "text") return String(field(item, "text") || "");
if (type === "image" || type === "localImage") return "图片附件";
if (type === "mention") return `@${String(field(item, "name") || "文件")}`;
return "";
})
.filter(Boolean)
.join("\n");
}
/** 移除用户消息中由旧流程拼接的 Agent 前置提示词。 */
function displayUserText(text: string) {
const value = text.trim();
const marker = "用户请求:";
const index = value.lastIndexOf(marker);
return (index >= 0 ? value.slice(index + marker.length) : value).trim();
}
/** 将未知值转换为数组。 */
function arrayValue(value: unknown) {
return Array.isArray(value) ? value : [];
}
/** 将非空字符串保留为字符串,否则返回 null。 */
function stringOrNull(value: unknown) {
return typeof value === "string" && value.trim() ? value : null;
}
/** 将 MCP 工具名称转换为聊天记录中的中文标题。 */
function toolName(name: string) {
if (name === "canvas_apply_ops") return "画布操作";
if (name === "canvas_get_state") return "读取画布";
if (name === "canvas_get_selection") return "读取选区";
if (name === "canvas_export_snapshot") return "导出快照";
if (name === "canvas_create_attachment_nodes") return "添加附件图片";
if (name === "canvas_create_text_node") return "创建文本";
if (name === "canvas_create_image_prompt_flow") return "创建生图流程";
if (name === "canvas_create_generation_flow") return "创建生成流程";
if (name === "canvas_generate_text") return "生成文本";
if (name === "canvas_generate_image") return "生成图片";
if (name === "canvas_generate_video") return "生成视频";
if (name === "canvas_generate_audio") return "生成音频";
if (name === "canvas_run_generation") return "触发生成";
return name;
}
+195
View File
@@ -0,0 +1,195 @@
import fs from "node:fs/promises";
import os from "node:os";
import path from "node:path";
import { logger } from "../utils/logger.js";
import { errorMessage, field } from "../utils/value.js";
import { CodexAppClient } from "./codex-client.js";
import { summarizeCodexThread, threadMessages } from "./codex-history.js";
import type { AgentAttachment, AgentEmit } from "./types.js";
type CodexRunOptions = { threadId?: string; cwd?: string; appEmit?: AgentEmit; onStart?: () => void; onThread?: (threadId: string) => void; onTurn?: (turnId: string) => void; onFinish?: () => void };
let codexQueue: Promise<unknown> = Promise.resolve();
let codexApp: CodexAppClient | null = null;
let codexAppStart: Promise<CodexAppClient> | null = null;
let codexThreadId = "";
export { summarizeCodexThread } from "./codex-history.js";
/** 将 Codex turn 加入串行队列并等待执行完成。 */
export async function runCodexTurn(prompt: string, emit: AgentEmit, attachments: AgentAttachment[] = [], options: CodexRunOptions = {}) {
if (!prompt.trim()) return;
codexQueue = codexQueue.catch(() => undefined).then(() => runCodexTurnNow(prompt, emit, attachments, options));
await codexQueue;
}
/** 中断当前线程正在执行的 Codex turn。 */
export function interruptCodexTurn(threadId?: string) {
if (!codexApp || (threadId && threadId !== codexThreadId)) return false;
return codexApp.interruptCurrentTurn();
}
/** 创建新的 Codex 线程并记录当前线程 ID。 */
export async function startCodexThread(emit: AgentEmit, cwd?: string) {
const app = await getCodexApp(emit);
const thread = await app.startThread(cwd);
codexThreadId = String(field(thread, "id") || "");
return thread;
}
/** 恢复指定 Codex 线程并返回聊天历史。 */
export async function resumeCodexThread(emit: AgentEmit, threadId: string, cwd?: string) {
const app = await getCodexApp(emit);
await loadCodexThread(emit, threadId, cwd, false);
const thread = await app.resumeThread(threadId, cwd);
assertThreadWorkspace(thread, cwd);
codexThreadId = String(field(thread, "id") || threadId);
return { thread, messages: threadMessages(thread) };
}
/** 查询当前工作空间中的 Codex 线程。 */
export async function listCodexThreads(emit: AgentEmit, options: { cwd: string; searchTerm?: string; limit?: number }) {
const app = await getCodexApp(emit);
const result = await app.listThreads({
limit: options.limit || 40,
sortKey: "updated_at",
sortDirection: "desc",
sourceKinds: ["cli", "vscode", "appServer", "exec"],
cwd: options.cwd,
...(options.searchTerm ? { searchTerm: options.searchTerm } : {}),
});
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 };
}
/** 读取指定 Codex 线程及其聊天历史。 */
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) };
}
/** 确认指定 Codex 线程属于当前工作空间。 */
export async function verifyCodexThreadWorkspace(emit: AgentEmit, threadId: string, cwd: string) {
await loadCodexThread(emit, threadId, cwd, false);
}
/** 归档指定 Codex 线程。 */
export async function archiveCodexThread(emit: AgentEmit, threadId: string, cwd?: string) {
const app = await getCodexApp(emit);
await loadCodexThread(emit, threadId, cwd, false);
await app.archiveThread(threadId);
}
/** 判断线程异常是否允许自动新建线程后重试。 */
export function isRecoverableThreadError(error: unknown) {
return /thread not loaded|no rollout found/i.test(errorMessage(error));
}
/** 执行一次 Codex turn,并负责附件临时文件和线程恢复。 */
async function runCodexTurnNow(prompt: string, emit: AgentEmit, attachments: AgentAttachment[], options: CodexRunOptions) {
let files: string[] = [];
try {
options.onStart?.();
files = await writeAttachmentFiles(attachments);
const app = await getCodexApp(options.appEmit || emit);
let threadId = await ensureCodexThread(app, options, emit);
options.onThread?.(threadId);
try {
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(app, { cwd: options.cwd }, emit);
options.onThread?.(threadId);
await app.startTurn(threadId, prompt, files, options.onTurn);
}
} catch (error) {
logger.error("Codex turn failed", error);
emit("agent_error", { message: errorMessage(error) });
} finally {
options.onFinish?.();
await Promise.all(files.map((file) => fs.unlink(file).catch(() => undefined)));
}
}
/** 恢复请求线程或创建新的 Codex 线程。 */
async function ensureCodexThread(app: CodexAppClient, options: CodexRunOptions, emit: AgentEmit) {
if (options.threadId) {
if (options.threadId === codexThreadId) return codexThreadId;
try {
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;
} catch (error) {
if (!isRecoverableThreadError(error)) throw error;
emit("agent_log", { text: `Codex thread unavailable, starting a new thread: ${errorMessage(error)}` });
}
}
if (!codexThreadId) {
const thread = await app.startThread(options.cwd);
codexThreadId = String(field(thread, "id") || "");
}
return codexThreadId;
}
/** 从 app-server 读取线程并校验工作空间。 */
async function loadCodexThread(emit: AgentEmit, threadId: string, cwd: string | undefined, includeTurns: boolean) {
const app = await getCodexApp(emit);
const result = await app.readThread(threadId, includeTurns);
const thread = field(result, "thread") || {};
assertThreadWorkspace(thread, cwd);
return thread;
}
/** 获取已启动的 Codex app-server 客户端。 */
async function getCodexApp(emit: AgentEmit) {
if (codexApp) return codexApp;
codexAppStart ||= CodexAppClient.start(emit, () => {
codexApp = null;
codexThreadId = "";
});
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 会话不属于当前画布工作空间");
}
/** 判断线程工作目录是否与当前工作空间一致。 */
function threadInWorkspace(thread: unknown, cwd: string) {
const threadCwd = String(field(thread, "cwd") || "");
return Boolean(threadCwd && path.resolve(threadCwd) === path.resolve(cwd));
}
/** 将图片附件写入临时文件供 Codex 读取。 */
async function writeAttachmentFiles(attachments: AgentAttachment[]) {
return await Promise.all(attachments.filter((item) => item.dataUrl?.startsWith("data:image/")).map(writeAttachmentFile));
}
/** 将单个 Data URL 图片附件写入临时文件。 */
async function writeAttachmentFile(item: AgentAttachment) {
const [, meta = "", data = ""] = item.dataUrl?.match(/^data:([^;]+);base64,(.+)$/) || [];
if (!data) throw new Error(`图片附件无效:${item.name || "未命名图片"}`);
const file = path.join(os.tmpdir(), `infinite-canvas-${Date.now()}-${Math.random().toString(16).slice(2)}.${imageExt(meta || item.type)}`);
await fs.writeFile(file, Buffer.from(data, "base64"));
return file;
}
/** 根据图片 MIME 类型返回临时文件扩展名。 */
function imageExt(type = "") {
if (type.includes("png")) return "png";
if (type.includes("webp")) return "webp";
return "jpg";
}
+5
View File
@@ -0,0 +1,5 @@
/** Agent 向网页广播事件的函数类型。 */
export type AgentEmit = (type: string, payload: unknown) => void;
/** 用户随当前 Agent 消息上传的附件。 */
export type AgentAttachment = { id?: string; name?: string; type?: string; size?: number; width?: number; height?: number; dataUrl?: string };