fix(agent): unify live and historical conversation state
This commit is contained in:
@@ -7,17 +7,27 @@ import stripAnsi from "strip-ansi";
|
||||
import { VERSION } from "../config.js";
|
||||
import { logger } from "../utils/logger.js";
|
||||
import { field, type JsonRecord } from "../utils/value.js";
|
||||
import { codexEventHistory, type CodexEventHistory } from "./codex-event-history.js";
|
||||
import type { CodexNotificationParams, CodexPlanUpdate, CodexReasoningEffort, CodexRequestMethod, CodexRequestParams, CodexRequestResult, CodexTurnInput } from "./codex-protocol.js";
|
||||
import type { AgentEmit, AgentPermissionMode } from "./types.js";
|
||||
|
||||
type AgentEvent = JsonRecord & { type: string; usage?: unknown };
|
||||
type PendingRequest = { resolve: (value: unknown) => void; reject: (error: Error) => void };
|
||||
type ItemDeltaParams = { threadId: string; turnId: string; itemId: string; delta: string };
|
||||
type ActiveTurn = PendingRequest & { threadId: string; turnId: string; prompt: string };
|
||||
type ItemDeltaParams = { threadId: string; turnId: string; itemId: string; delta: string; summaryIndex?: number };
|
||||
type PendingDelta = { delta: string; itemType: string; params: ItemDeltaParams; timer: ReturnType<typeof setTimeout> };
|
||||
type ApprovalRequest = { id: number; method: string; params: JsonRecord; decision?: string };
|
||||
type PendingTurnStart = { threadId: string; prompt: string; turnId?: string; onTurn?: (turnId: string) => void };
|
||||
|
||||
const canvasAgentMcp = canvasAgentMcpCommand();
|
||||
const require = createRequire(import.meta.url);
|
||||
const STREAM_UPDATE_INTERVAL_MS = 40;
|
||||
const supplementalItemTypes = new Set(["agent_message", "reasoning", "plan", "mcp_tool_call", "command_execution", "file_change", "dynamic_tool_call", "collab_tool_call", "web_search", "image_view", "image_generation", "context_compaction"]);
|
||||
|
||||
/** 表示错误已经通过 app-server 终态或进程事件通知过网页。 */
|
||||
export class CodexReportedError extends Error {
|
||||
override name = "CodexReportedError";
|
||||
}
|
||||
|
||||
/** 封装 Codex app-server 的 JSON-RPC 通信与事件转换。 */
|
||||
export class CodexAppClient {
|
||||
@@ -25,23 +35,37 @@ export class CodexAppClient {
|
||||
private buffer = "";
|
||||
private currentThreadId = "";
|
||||
private currentTurnId = "";
|
||||
private pendingTurnStart?: PendingTurnStart;
|
||||
private startedTurnKeys = new Set<string>();
|
||||
private textByItem = new Map<string, string>();
|
||||
private reasoningTextByItem = new Map<string, Map<number, string>>();
|
||||
private lastUsage: unknown = null;
|
||||
private pending = new Map<number, PendingRequest>();
|
||||
private activeTurns = new Map<string, PendingRequest>();
|
||||
private activeTurns = new Map<string, ActiveTurn>();
|
||||
private completedTurns = new Map<string, Error | null>();
|
||||
private pendingDeltas = new Map<string, PendingDelta>();
|
||||
private startedItems = new Map<string, JsonRecord>();
|
||||
private itemSequences = new Map<string, number>();
|
||||
private nextItemSequences = new Map<string, number>();
|
||||
private plansByTurn = new Map<string, CodexPlanUpdate>();
|
||||
private approvalRequests = new Map<string, { id: number; method: string; params: JsonRecord }>();
|
||||
private approvalRequests = new Map<string, ApprovalRequest>();
|
||||
private finalizingTurns = new Map<string, Promise<void>>();
|
||||
private failing = false;
|
||||
|
||||
/** 保存 app-server 子进程和事件出口。 */
|
||||
private constructor(private child: ChildProcess, private emit: AgentEmit) {}
|
||||
private constructor(private child: ChildProcess, private emit: AgentEmit, private eventHistory: Pick<CodexEventHistory, "record" | "recordTurn"> = codexEventHistory) {}
|
||||
|
||||
/** 启动并初始化 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);
|
||||
let stopped = false;
|
||||
const stop = () => {
|
||||
if (stopped) return;
|
||||
stopped = true;
|
||||
onExit();
|
||||
};
|
||||
child.stdout?.on("data", (chunk) => client.read(chunk.toString()));
|
||||
child.stderr?.on("data", (chunk) => {
|
||||
const text = stripAnsi(chunk.toString()).replace(/^\d{4}-\d{2}-\d{2}T[\d:.]+Z\s+/, "");
|
||||
@@ -51,11 +75,13 @@ export class CodexAppClient {
|
||||
child.on("error", (error) => {
|
||||
logger.error("Codex app-server process error", error);
|
||||
emit("agent_error", { message: error.message });
|
||||
client.failAll(error.message, true);
|
||||
stop();
|
||||
});
|
||||
child.on("exit", (code) => {
|
||||
logger.warn("Codex app-server exited", { code });
|
||||
client.failAll(`Codex app-server exited: ${code ?? 0}`);
|
||||
onExit();
|
||||
stop();
|
||||
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 } });
|
||||
@@ -104,35 +130,49 @@ export class CodexAppClient {
|
||||
|
||||
/** 清理已归档线程的任务计划缓存。 */
|
||||
clearPlanUpdates(threadId: string) {
|
||||
this.plansByTurn.forEach((item, turnId) => {
|
||||
if (item.threadId === threadId) this.plansByTurn.delete(turnId);
|
||||
this.plansByTurn.forEach((item, key) => {
|
||||
if (item.threadId === threadId) this.plansByTurn.delete(key);
|
||||
});
|
||||
}
|
||||
|
||||
/** 启动一个 Codex turn 并等待完成通知。 */
|
||||
async startTurn(threadId: string, prompt: string, images: string[], permissionMode: AgentPermissionMode, model?: string, effort?: CodexReasoningEffort, onTurn?: (turnId: string) => void) {
|
||||
this.currentThreadId = threadId;
|
||||
const { turn } = await this.request("turn/start", { threadId, input: codexInput(prompt, images), ...turnSettings(permissionMode), ...(model ? { model } : {}), ...(effort ? { effort } : {}) });
|
||||
const turnId = turn.id;
|
||||
if (!turnId) throw new Error("Codex app-server 没有返回 turn id");
|
||||
this.currentTurnId = turnId;
|
||||
onTurn?.(turnId);
|
||||
const completed = this.completedTurns.get(turnId);
|
||||
if (this.completedTurns.has(turnId)) {
|
||||
this.completedTurns.delete(turnId);
|
||||
this.currentThreadId = "";
|
||||
this.currentTurnId = "";
|
||||
if (completed) throw completed;
|
||||
return;
|
||||
this.currentTurnId = "";
|
||||
this.lastUsage = null;
|
||||
const pendingStart: PendingTurnStart = { threadId, prompt, onTurn };
|
||||
this.pendingTurnStart = pendingStart;
|
||||
try {
|
||||
const { turn } = await this.request("turn/start", { threadId, input: codexInput(prompt, images), ...turnSettings(permissionMode), ...(model ? { model } : {}), ...(effort ? { effort } : {}) });
|
||||
const turnId = turn.id;
|
||||
if (!turnId) throw new Error("Codex app-server 没有返回 turn id");
|
||||
pendingStart.turnId = turnId;
|
||||
this.currentTurnId = turnId;
|
||||
this.notifyTurnStarted(threadId, turnId, pendingStart);
|
||||
const turnKey = turnCacheKey(threadId, turnId);
|
||||
const completed = this.completedTurns.get(turnKey);
|
||||
if (this.completedTurns.has(turnKey)) {
|
||||
this.completedTurns.delete(turnKey);
|
||||
this.currentThreadId = "";
|
||||
this.currentTurnId = "";
|
||||
if (completed) throw completed;
|
||||
return;
|
||||
}
|
||||
await new Promise((resolve, reject) => this.activeTurns.set(turnKey, { resolve, reject, threadId, turnId, prompt }));
|
||||
} catch (error) {
|
||||
if (!this.currentTurnId) this.currentThreadId = "";
|
||||
throw error;
|
||||
} finally {
|
||||
if (this.pendingTurnStart === pendingStart) this.pendingTurnStart = undefined;
|
||||
if (pendingStart.turnId) this.startedTurnKeys.delete(turnCacheKey(threadId, pendingStart.turnId));
|
||||
}
|
||||
await new Promise((resolve, reject) => this.activeTurns.set(turnId, { resolve, reject }));
|
||||
}
|
||||
|
||||
/** 中断当前正在运行的 Codex turn。 */
|
||||
async interruptCurrentTurn() {
|
||||
/** 中断当前正在运行且属于指定线程的 Codex turn。 */
|
||||
async interruptCurrentTurn(requestedThreadId?: string) {
|
||||
const threadId = this.currentThreadId;
|
||||
const turnId = this.currentTurnId;
|
||||
if (!threadId || !turnId) return false;
|
||||
if (!threadId || !turnId || (requestedThreadId && requestedThreadId !== threadId)) return false;
|
||||
try {
|
||||
logger.warn("Interrupting active Codex turn", { threadId, turnId });
|
||||
await this.request("turn/interrupt", { threadId, turnId });
|
||||
@@ -147,10 +187,12 @@ export class CodexAppClient {
|
||||
resolveApproval(requestId: string, decision: string) {
|
||||
const request = this.approvalRequests.get(requestId);
|
||||
if (!request) return false;
|
||||
this.approvalRequests.delete(requestId);
|
||||
if (request.decision) return true;
|
||||
request.decision = decision;
|
||||
const permissions = field(request.params, "permissions") || field(request.params, "requestedPermissions");
|
||||
const accepted = decision === "accept" || decision === "acceptForSession";
|
||||
const result = request.method === "item/permissions/requestApproval"
|
||||
? { permissions: decision === "decline" ? {} : permissions || {}, scope: decision === "acceptForSession" ? "session" : "turn" }
|
||||
? { permissions: accepted ? permissions || {} : {}, scope: decision === "acceptForSession" ? "session" : "turn" }
|
||||
: { decision };
|
||||
this.write({ id: request.id, result });
|
||||
return true;
|
||||
@@ -209,14 +251,26 @@ export class CodexAppClient {
|
||||
private handleNotification(method: string, params: JsonRecord) {
|
||||
if (method === "serverRequest/resolved") {
|
||||
const requestId = String(field(params, "requestId") || "");
|
||||
if (requestId) this.approvalRequests.delete(requestId);
|
||||
this.emit("codex_approval_resolved", { requestId, ...params });
|
||||
const request = requestId ? this.approvalRequests.get(requestId) : undefined;
|
||||
if (request) {
|
||||
this.approvalRequests.delete(requestId);
|
||||
this.emit("codex_approval_resolved", { ...request.params, ...params, requestId, decision: request.decision });
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (!field(params, "threadId") && this.currentThreadId && (method === "turn/started" || method === "turn/completed" || method === "turn/plan/updated")) params = { ...params, threadId: this.currentThreadId };
|
||||
const turnEvent = method.startsWith("turn/") || method.startsWith("item/") || method === "thread/tokenUsage/updated" || method === "error";
|
||||
if (turnEvent) {
|
||||
const threadId = String(field(params, "threadId") || this.currentThreadId);
|
||||
const turnId = String(field(params, "turnId") || field(field(params, "turn"), "id") || this.currentTurnId);
|
||||
if (method === "turn/started") {
|
||||
this.currentThreadId = threadId;
|
||||
this.currentTurnId = turnId;
|
||||
this.notifyTurnStarted(threadId, turnId);
|
||||
}
|
||||
params = { ...params, ...(threadId ? { threadId } : {}), ...(turnId ? { turnId } : {}) };
|
||||
}
|
||||
if (method === "item/agentMessage/delta") {
|
||||
const value = params as unknown as CodexNotificationParams<"item/agentMessage/delta">;
|
||||
this.textByItem.set(value.itemId, `${this.textByItem.get(value.itemId) || ""}${value.delta}`);
|
||||
return this.emitDelta("agent_message", value);
|
||||
}
|
||||
if (method === "item/plan/delta") return this.emitDelta("plan", params as unknown as CodexNotificationParams<"item/plan/delta">);
|
||||
@@ -225,7 +279,7 @@ export class CodexAppClient {
|
||||
if (method === "turn/plan/updated") {
|
||||
const value = params as unknown as CodexNotificationParams<"turn/plan/updated">;
|
||||
const update: CodexPlanUpdate = { ...value, threadId: value.threadId || "" };
|
||||
if (update.threadId && update.turnId) this.plansByTurn.set(update.turnId, update);
|
||||
if (update.threadId && update.turnId) this.plansByTurn.set(turnCacheKey(update.threadId, update.turnId), update);
|
||||
params = update as unknown as JsonRecord;
|
||||
}
|
||||
if (method === "thread/tokenUsage/updated") {
|
||||
@@ -235,66 +289,162 @@ export class CodexAppClient {
|
||||
}
|
||||
const event = normalizeCodexNotification(method, params);
|
||||
if (!event) return;
|
||||
const eventScope = codexEventScope(params);
|
||||
if (event.type === "item.started" || event.type === "item.completed") {
|
||||
const item = field(event, "item") as JsonRecord | undefined;
|
||||
const id = String(field(item, "id") || "");
|
||||
const threadId = String(field(event, "thread_id") || this.currentThreadId);
|
||||
const turnId = String(field(event, "turn_id") || this.currentTurnId);
|
||||
const key = itemCacheKey({ threadId, turnId, itemId: id });
|
||||
if (id && threadId && turnId && event.type === "item.started") {
|
||||
this.assignItemSequence(threadId, turnId, id);
|
||||
this.startedItems.set(key, item || {});
|
||||
}
|
||||
if (id && event.type === "item.completed") {
|
||||
const started = this.startedItems.get(key);
|
||||
if (started) event.item = mergeDefinedRecord(started, item || {});
|
||||
this.startedItems.delete(key);
|
||||
}
|
||||
}
|
||||
if (event.type === "item.completed") {
|
||||
const item = field(event, "item") as JsonRecord | undefined;
|
||||
const id = String(field(item, "id") || "");
|
||||
this.flushDelta(id);
|
||||
const streamedText = this.textByItem.get(id);
|
||||
if (item?.type === "agent_message" && streamedText && !item.text) item.text = streamedText;
|
||||
if (id) this.textByItem.delete(id);
|
||||
const key = itemCacheKey({ threadId: String(field(event, "thread_id") || this.currentThreadId), turnId: String(field(event, "turn_id") || this.currentTurnId), itemId: id });
|
||||
this.flushDelta(key);
|
||||
const streamedText = this.textByItem.get(key);
|
||||
if ((item?.type === "agent_message" || item?.type === "plan") && streamedText && !String(item.text || "").trim()) item.text = streamedText;
|
||||
if (item?.type === "reasoning" && streamedText && !readableCodexText(field(item, "summary"))) item.summary = streamedText;
|
||||
if (item?.type === "command_execution" && streamedText && !String(item.aggregatedOutput || "").trim()) item.aggregatedOutput = streamedText;
|
||||
const threadId = String(field(event, "thread_id") || this.currentThreadId);
|
||||
const turnId = String(field(event, "turn_id") || this.currentTurnId);
|
||||
if (id && threadId && turnId && item && supplementalItemTypes.has(String(item.type || ""))) {
|
||||
const sequence = this.assignItemSequence(threadId, turnId, id);
|
||||
void this.eventHistory.record({ threadId, turnId, itemId: id, sequence, item }).catch((error) => logger.warn("Failed to persist Codex event history", { threadId, turnId, itemId: id, error }));
|
||||
}
|
||||
if (id) {
|
||||
this.textByItem.delete(key);
|
||||
this.reasoningTextByItem.delete(key);
|
||||
}
|
||||
}
|
||||
let turnPersistence: Promise<void> | undefined;
|
||||
if (event.type === "turn.completed") {
|
||||
const turn = field(params, "turn");
|
||||
const turnId = String(field(turn, "id") || field(params, "turnId") || "");
|
||||
const plan = this.plansByTurn.get(turnId);
|
||||
if (plan) this.plansByTurn.set(turnId, { ...plan, turnStatus: String(field(turn, "status") || "completed") });
|
||||
const threadId = String(field(params, "threadId") || field(event, "thread_id") || "");
|
||||
const planKey = turnCacheKey(threadId, turnId);
|
||||
const plan = this.plansByTurn.get(planKey);
|
||||
if (plan) this.plansByTurn.set(planKey, { ...plan, turnStatus: String(field(turn, "status") || "completed") });
|
||||
if (threadId && turnId) {
|
||||
const input = this.pendingTurnStart?.threadId === threadId ? this.pendingTurnStart.prompt : "";
|
||||
const turnRecord = { ...(turn && typeof turn === "object" && !Array.isArray(turn) ? turn as JsonRecord : { id: turnId, status: field(turn, "status") || "completed" }), ...(input ? { input } : {}) };
|
||||
turnPersistence = this.eventHistory.recordTurn({ threadId, turnId, turn: turnRecord }).catch((error) => logger.warn("Failed to persist Codex turn history", { threadId, turnId, error }));
|
||||
this.finalizingTurns.set(planKey, turnPersistence);
|
||||
}
|
||||
this.finishTurnDeltas(threadId, turnId);
|
||||
}
|
||||
if (event.type === "turn.completed") event.usage = this.lastUsage;
|
||||
this.emit("agent_event", { agent: "codex", ...event });
|
||||
if (event.type === "turn.completed") {
|
||||
const turn = (params as unknown as CodexNotificationParams<"turn/completed">).turn;
|
||||
const turnId = turn.id;
|
||||
const pending = this.activeTurns.get(turnId);
|
||||
const error = turn.error;
|
||||
if (pending) {
|
||||
this.activeTurns.delete(turnId);
|
||||
error ? pending.reject(new Error(error.message || "Codex turn failed")) : pending.resolve(event);
|
||||
} else if (turnId) {
|
||||
this.completedTurns.set(turnId, error ? new Error(error.message || "Codex turn failed") : null);
|
||||
}
|
||||
if (turnId === this.currentTurnId) {
|
||||
this.currentThreadId = "";
|
||||
this.currentTurnId = "";
|
||||
}
|
||||
this.emit("agent_done", { agent: "codex", usage: event.usage, ...codexEventScope(params) });
|
||||
event.usage = this.lastUsage;
|
||||
const complete = () => this.completeTurn(event, params, eventScope);
|
||||
if (turnPersistence) void turnPersistence.then(complete);
|
||||
else complete();
|
||||
return;
|
||||
}
|
||||
this.emit("agent_event", { agent: "codex", ...event });
|
||||
}
|
||||
|
||||
/** 补充历史落盘后再广播 turn 终态,确保界面完成状态可跨 Agent 重启恢复。 */
|
||||
private completeTurn(event: AgentEvent, params: JsonRecord, eventScope: ReturnType<typeof codexEventScope>) {
|
||||
const turn = (params as unknown as CodexNotificationParams<"turn/completed">).turn;
|
||||
const turnId = turn.id;
|
||||
const turnKey = turnCacheKey(String(field(params, "threadId") || field(event, "thread_id") || ""), turnId);
|
||||
this.finalizingTurns.delete(turnKey);
|
||||
this.emit("agent_event", { agent: "codex", ...event });
|
||||
const pending = this.activeTurns.get(turnKey);
|
||||
const error = turn.error;
|
||||
const failure = error ? new CodexReportedError(error.message || "Codex turn failed") : null;
|
||||
if (pending) {
|
||||
this.activeTurns.delete(turnKey);
|
||||
failure ? pending.reject(failure) : pending.resolve(event);
|
||||
} else if (turnId) {
|
||||
this.completedTurns.set(turnKey, failure);
|
||||
}
|
||||
if (turnId === this.currentTurnId) {
|
||||
this.currentThreadId = "";
|
||||
this.currentTurnId = "";
|
||||
}
|
||||
this.emit("agent_done", { agent: "codex", usage: event.usage, ...eventScope });
|
||||
}
|
||||
|
||||
/** 合并并广播 Agent 文本或执行输出增量。 */
|
||||
private emitDelta(itemType: string, params: ItemDeltaParams) {
|
||||
const id = params.itemId;
|
||||
const pending = this.pendingDeltas.get(id);
|
||||
params = { ...params, threadId: params.threadId || this.currentThreadId, turnId: params.turnId || this.currentTurnId };
|
||||
const key = itemCacheKey(params);
|
||||
this.textByItem.set(key, itemType === "reasoning" ? this.appendReasoningText(key, params) : `${this.textByItem.get(key) || ""}${params.delta}`);
|
||||
const pending = this.pendingDeltas.get(key);
|
||||
if (pending) {
|
||||
pending.delta += params.delta;
|
||||
pending.itemType = itemType;
|
||||
pending.params = params;
|
||||
return;
|
||||
}
|
||||
this.pendingDeltas.set(id, {
|
||||
this.pendingDeltas.set(key, {
|
||||
delta: params.delta,
|
||||
itemType,
|
||||
params,
|
||||
timer: setTimeout(() => this.flushDelta(id), STREAM_UPDATE_INTERVAL_MS),
|
||||
timer: setTimeout(() => this.flushDelta(key), STREAM_UPDATE_INTERVAL_MS),
|
||||
});
|
||||
}
|
||||
|
||||
/** 合并短时间内的文本增量,减少 SSE 传输和前端渲染次数。 */
|
||||
private flushDelta(id: string) {
|
||||
const pending = this.pendingDeltas.get(id);
|
||||
private flushDelta(key: string) {
|
||||
const pending = this.pendingDeltas.get(key);
|
||||
if (!pending) return;
|
||||
clearTimeout(pending.timer);
|
||||
this.pendingDeltas.delete(id);
|
||||
if (pending.delta) this.emit("agent_event", { agent: "codex", type: "item.updated", item: { id, type: pending.itemType, delta: pending.delta }, ...codexEventScope(pending.params as unknown as JsonRecord) });
|
||||
this.pendingDeltas.delete(key);
|
||||
if (pending.delta) this.emit("agent_event", { agent: "codex", type: "item.updated", item: { id: pending.params.itemId, type: pending.itemType, delta: pending.delta }, ...codexEventScope(pending.params as unknown as JsonRecord) });
|
||||
}
|
||||
|
||||
/** 按 Codex summaryIndex 保存 reasoning 分段,顺序与线程历史一致。 */
|
||||
private appendReasoningText(key: string, params: ItemDeltaParams) {
|
||||
const segments = this.reasoningTextByItem.get(key) || new Map<number, string>();
|
||||
this.reasoningTextByItem.set(key, segments);
|
||||
return appendReasoningDelta(segments, params.summaryIndex, params.delta);
|
||||
}
|
||||
|
||||
/** turn 结束时发送最后一批增量并清理未收到 item.completed 的缓存。 */
|
||||
private finishTurnDeltas(threadId: string, turnId: string) {
|
||||
const prefix = `${turnCacheKey(threadId, turnId)}\0`;
|
||||
[...this.pendingDeltas.keys()].filter((key) => key.startsWith(prefix)).forEach((key) => this.flushDelta(key));
|
||||
[...this.textByItem.keys()].filter((key) => key.startsWith(prefix)).forEach((key) => this.textByItem.delete(key));
|
||||
[...this.reasoningTextByItem.keys()].filter((key) => key.startsWith(prefix)).forEach((key) => this.reasoningTextByItem.delete(key));
|
||||
[...this.startedItems.keys()].filter((key) => key.startsWith(prefix)).forEach((key) => this.startedItems.delete(key));
|
||||
[...this.itemSequences.keys()].filter((key) => key.startsWith(prefix)).forEach((key) => this.itemSequences.delete(key));
|
||||
this.nextItemSequences.delete(turnCacheKey(threadId, turnId));
|
||||
}
|
||||
|
||||
/** 为一个 turn 内的 item 固定开始顺序,完成通知只更新内容。 */
|
||||
private assignItemSequence(threadId: string, turnId: string, itemId: string) {
|
||||
const key = itemCacheKey({ threadId, turnId, itemId });
|
||||
const existing = this.itemSequences.get(key);
|
||||
if (existing !== undefined) return existing;
|
||||
const turnKey = turnCacheKey(threadId, turnId);
|
||||
const sequence = (this.nextItemSequences.get(turnKey) || 0) + 1;
|
||||
this.nextItemSequences.set(turnKey, sequence);
|
||||
this.itemSequences.set(key, sequence);
|
||||
return sequence;
|
||||
}
|
||||
|
||||
/** 在通知或 turn/start 响应到达时回调一次 turn 启动状态。 */
|
||||
private notifyTurnStarted(threadId: string, turnId: string, fallback?: PendingTurnStart) {
|
||||
if (!threadId || !turnId) return;
|
||||
const pending = this.pendingTurnStart;
|
||||
const registration = pending?.threadId === threadId && (!pending.turnId || pending.turnId === turnId) ? pending : fallback;
|
||||
if (registration && !registration.turnId) registration.turnId = turnId;
|
||||
const key = turnCacheKey(threadId, turnId);
|
||||
if (this.startedTurnKeys.has(key)) return;
|
||||
if (!registration) return;
|
||||
this.startedTurnKeys.add(key);
|
||||
registration.onTurn?.(turnId);
|
||||
}
|
||||
|
||||
/** 自动回复 app-server 发起的授权或交互请求。 */
|
||||
@@ -325,19 +475,79 @@ export class CodexAppClient {
|
||||
}
|
||||
|
||||
/** 拒绝进程退出时仍未完成的请求与 turn。 */
|
||||
private failAll(message: string) {
|
||||
[...this.pending.values(), ...this.activeTurns.values()].forEach((item) => item.reject(new Error(message)));
|
||||
private failAll(message: string, reported = false) {
|
||||
if (this.failing) return;
|
||||
this.failing = true;
|
||||
this.approvalRequests.forEach((request, requestId) => this.emit("codex_approval_resolved", { ...request.params, requestId, decision: request.decision || "cancel" }));
|
||||
const failedTurns = new Map<string, { threadId: string; turnId: string; prompt: string }>();
|
||||
this.activeTurns.forEach(({ threadId, turnId, prompt }, key) => {
|
||||
if (!this.finalizingTurns.has(key)) failedTurns.set(key, { threadId, turnId, prompt });
|
||||
});
|
||||
const pendingStart = this.pendingTurnStart;
|
||||
if (pendingStart?.turnId) {
|
||||
const key = turnCacheKey(pendingStart.threadId, pendingStart.turnId);
|
||||
if (!this.finalizingTurns.has(key) && !failedTurns.has(key)) failedTurns.set(key, { threadId: pendingStart.threadId, turnId: pendingStart.turnId, prompt: pendingStart.prompt });
|
||||
}
|
||||
const finalizing = [...this.finalizingTurns.values()];
|
||||
const persistence = [...failedTurns.values()].map(({ threadId, turnId, prompt }) => {
|
||||
const turn = { id: turnId, status: "failed", error: { message }, ...(prompt ? { input: prompt } : {}) };
|
||||
return this.eventHistory.recordTurn({ threadId, turnId, turn }).catch((historyError) => logger.warn("Failed to persist Codex turn failure", { threadId, turnId, error: historyError }));
|
||||
});
|
||||
const error = reported || failedTurns.size || finalizing.length ? new CodexReportedError(message) : new Error(message);
|
||||
this.pendingDeltas.forEach((item) => clearTimeout(item.timer));
|
||||
this.pending.clear();
|
||||
this.activeTurns.clear();
|
||||
this.pendingDeltas.clear();
|
||||
this.textByItem.clear();
|
||||
this.reasoningTextByItem.clear();
|
||||
this.startedItems.clear();
|
||||
this.itemSequences.clear();
|
||||
this.nextItemSequences.clear();
|
||||
this.plansByTurn.clear();
|
||||
this.completedTurns.clear();
|
||||
this.approvalRequests.clear();
|
||||
this.pendingTurnStart = undefined;
|
||||
this.startedTurnKeys.clear();
|
||||
this.lastUsage = null;
|
||||
this.currentThreadId = "";
|
||||
this.currentTurnId = "";
|
||||
void Promise.all([...finalizing, ...persistence]).then(() => {
|
||||
failedTurns.forEach(({ threadId, turnId, prompt }) => {
|
||||
const turn = { id: turnId, status: "failed", error: { message }, ...(prompt ? { input: prompt } : {}) };
|
||||
this.emit("agent_event", { agent: "codex", type: "turn.completed", status: "failed", error: { message }, thread_id: threadId, turn_id: turnId, turn });
|
||||
this.emit("agent_done", { agent: "codex", status: "failed", error: { message }, thread_id: threadId, turn_id: turnId });
|
||||
});
|
||||
this.pending.forEach((item) => item.reject(error));
|
||||
this.activeTurns.forEach((item) => item.reject(error));
|
||||
this.pending.clear();
|
||||
this.activeTurns.clear();
|
||||
this.finalizingTurns.clear();
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/** 将 Codex 的字符串、摘要数组或文本对象转换为展示文本。 */
|
||||
function readableCodexText(value: unknown): string {
|
||||
if (typeof value === "string") return value.trim();
|
||||
if (Array.isArray(value)) return value.map(readableCodexText).filter(Boolean).join("\n");
|
||||
if (!value || typeof value !== "object") return "";
|
||||
return readableCodexText(field(value, "text"));
|
||||
}
|
||||
|
||||
/** 合并单个 reasoning 分段增量并按 summaryIndex 输出。 */
|
||||
export function appendReasoningDelta(segments: Map<number, string>, summaryIndex: number | undefined, delta: string) {
|
||||
const index = Number.isInteger(summaryIndex) ? Number(summaryIndex) : 0;
|
||||
segments.set(index, `${segments.get(index) || ""}${delta}`);
|
||||
return [...segments.entries()].sort(([left], [right]) => left - right).map(([, text]) => text.trim()).filter(Boolean).join("\n");
|
||||
}
|
||||
|
||||
/** 生成仅供进程内缓存使用的完整 turn 与 item 作用域键。 */
|
||||
function itemCacheKey(scope: { threadId: string; turnId: string; itemId: string }) {
|
||||
return `${turnCacheKey(scope.threadId, scope.turnId)}\0${scope.itemId}`;
|
||||
}
|
||||
|
||||
function turnCacheKey(threadId: string, turnId: string) {
|
||||
return `${threadId}\0${turnId}`;
|
||||
}
|
||||
|
||||
/** 生成 Codex 调用 Canvas Agent MCP 的启动命令。 */
|
||||
function canvasAgentMcpCommand() {
|
||||
const current = process.argv.find((arg) => /index\.(t|j)s$/.test(arg)) || "";
|
||||
@@ -395,7 +605,7 @@ function normalizeItem(item: unknown) {
|
||||
if (value.type === "commandExecution") value.type = "command_execution";
|
||||
if (value.type === "fileChange") value.type = "file_change";
|
||||
if (value.type === "dynamicToolCall") value.type = "dynamic_tool_call";
|
||||
if (value.type === "collabToolCall") value.type = "collab_tool_call";
|
||||
if (value.type === "collabToolCall" || value.type === "collabAgentToolCall") value.type = "collab_tool_call";
|
||||
if (value.type === "webSearch") value.type = "web_search";
|
||||
if (value.type === "imageView") value.type = "image_view";
|
||||
if (value.type === "imageGeneration") value.type = "image_generation";
|
||||
@@ -405,6 +615,14 @@ function normalizeItem(item: unknown) {
|
||||
return value;
|
||||
}
|
||||
|
||||
function mergeDefinedRecord(started: JsonRecord, completed: JsonRecord) {
|
||||
const merged = { ...started };
|
||||
Object.entries(completed).forEach(([key, value]) => {
|
||||
if (value !== undefined) merged[key] = value;
|
||||
});
|
||||
return merged;
|
||||
}
|
||||
|
||||
/** 将 Codex token usage 转换为前端字段。 */
|
||||
function normalizeUsage(params: CodexNotificationParams<"thread/tokenUsage/updated">) {
|
||||
const last = params.tokenUsage.last;
|
||||
|
||||
Reference in New Issue
Block a user