Files
pikasTech-HWLAB/internal/cloud/codex-stdio-session-turn-state.ts
T
2026-06-01 00:24:20 +08:00

580 lines
20 KiB
TypeScript

import { CODEX_STDIO_FIRST_TOKEN_PROGRESS_MS, CODEX_STDIO_NO_PROGRESS_DIAGNOSIS_MS } from "./codex-stdio-session.ts";
import {
appServerActivityTimeoutError,
appServerCommandStatus,
appServerHardTimeoutError,
appServerTerminalStatus,
commandExecutionCommandText,
commandExecutionErrorText,
commandExecutionOutputText,
extractAppServerRecord,
extractAppServerString,
firstNonEmpty,
optionalId,
redactText,
boundToolOutput,
tailText
} from "./codex-stdio-session-helpers.ts";
export function createAppServerTurnState({ traceRecorder, session } = {}) {
let threadId = optionalId(session?.threadId);
let turnId = null;
let assistantText = "";
let finalResponse = "";
let terminal = null;
let terminalResolve;
const terminalPromise = new Promise((resolve) => {
terminalResolve = resolve;
});
let lastActivityAt = Date.now();
let lastActivityLabel = "turn-state:created";
let lastWaitingFor = "app-server-turn";
let lastRoutedEvent = null;
let lastClientRequest = null;
let lastDiagnosis = null;
let lastProgressDiagnosisAt = 0;
let lastProgressDiagnosisKey = null;
function setThreadId(value) {
threadId = optionalId(value) ?? threadId;
}
function setTurnId(value) {
turnId = optionalId(value) ?? turnId;
}
function activitySnapshot(referenceNow = Date.now()) {
return {
lastActivityAt: new Date(lastActivityAt).toISOString(),
lastActivityLabel,
idleMs: Math.max(0, referenceNow - lastActivityAt),
waitingFor: lastWaitingFor
};
}
function observeActivity({ label = "app-server:activity", waitingFor = null } = {}) {
lastActivityAt = Date.now();
lastActivityLabel = label;
if (waitingFor) lastWaitingFor = waitingFor;
return activitySnapshot(lastActivityAt);
}
function appendTrace(event = {}) {
const appended = traceRecorder?.append(event);
lastRoutedEvent = compactTraceEvent(appended ?? event);
observeActivity({
label: appended?.label ?? event.label ?? event.type ?? "trace:event",
waitingFor: appended?.waitingFor ?? event.waitingFor ?? null
});
return appended;
}
function appendProgressTrace(event = {}) {
return traceRecorder?.append(event);
}
function hasAssistantOutput() {
return Boolean(firstNonEmpty(finalResponse, assistantText));
}
function appendFirstTokenProgressTrace(referenceNow = Date.now()) {
const activity = activitySnapshot(referenceNow);
if (hasAssistantOutput() || !turnId || terminal) return null;
if (!["assistant-message", "turn/completed"].includes(activity.waitingFor)) return null;
return appendProgressTrace({
type: "turn",
status: "running",
label: "turn:waiting:first_assistant_token",
message: "Codex app-server is waiting for the first assistant token; this progress trace does not reset the backend idle timer.",
progressOnly: true,
idleMs: activity.idleMs,
lastActivityAt: activity.lastActivityAt,
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
waitingFor: "assistant-message"
});
}
function appendNoProgressDiagnosisTrace(referenceNow = Date.now(), diagnosticMs = CODEX_STDIO_NO_PROGRESS_DIAGNOSIS_MS) {
if (terminal) return null;
const activity = activitySnapshot(referenceNow);
if (activity.idleMs < diagnosticMs) return null;
const diagnosis = diagnoseNoProgress("running-no-progress", activity);
const key = `${diagnosis.code}:${activity.waitingFor ?? "unknown"}`;
if (lastProgressDiagnosisKey === key && referenceNow - lastProgressDiagnosisAt < diagnosticMs) return null;
lastProgressDiagnosisKey = key;
lastProgressDiagnosisAt = referenceNow;
return appendProgressTrace({
type: "turn",
status: "running",
label: "turn:no_progress:diagnosed",
message: diagnosis.summary,
progressOnly: true,
diagnosis,
diagnosisCode: diagnosis.code,
idleMs: activity.idleMs,
lastActivityAt: activity.lastActivityAt,
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
waitingFor: activity.waitingFor
});
}
function finish(status, error = null, extra = {}) {
if (terminal) return;
terminal = {
terminalStatus: appServerTerminalStatus(status),
terminalError: error ? redactText(error) : null,
threadId,
turnId,
...activitySnapshot(),
...extra
};
terminalResolve(terminal);
}
function handle(message) {
const method = typeof message?.method === "string" ? message.method : "unknown";
observeActivity({
label: `app-server:${method}`,
waitingFor: appServerWaitingForMethod(method)
});
const params = extractAppServerRecord(message?.params);
const item = extractAppServerRecord(params?.item);
const turn = extractAppServerRecord(params?.turn);
if (method === "thread/started") {
setThreadId(extractAppServerString(extractAppServerRecord(params?.thread), "id"));
appendTrace({
type: "thread",
status: "completed",
label: "thread:started",
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
waitingFor: "turn/start"
});
return;
}
if (method === "turn/started") {
setTurnId(extractAppServerString(turn, "id"));
appendTrace({
type: "turn",
status: "started",
label: "turn:started",
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
waitingFor: "assistant-message"
});
return;
}
if (method === "item/agentMessage/delta") {
const delta = String(params?.delta ?? "");
assistantText += delta;
if (delta) {
traceRecorder?.appendAssistantDelta?.({
itemId: optionalId(params?.itemId),
chunk: delta,
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
waitingFor: "turn/completed"
});
}
return;
}
if (method === "item/completed" && item?.type === "agentMessage") {
const completedText = typeof item.text === "string" ? item.text : "";
if (completedText) finalResponse = completedText;
appendTrace({
type: "assistant_message",
status: "completed",
label: "assistant:item_completed",
itemId: optionalId(item.id),
message: redactText(completedText),
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
waitingFor: "turn/completed"
});
return;
}
if (method === "client/request/handled") {
const decision = extractAppServerString(params, "decision") ?? "handled";
const handledMethod = extractAppServerString(params, "method") ?? "client/request";
const pending = decision === "received";
lastClientRequest = {
method: handledMethod,
decision,
status: pending ? "pending" : "handled",
itemId: optionalId(params?.itemId ?? params?.targetItemId),
approvalId: optionalId(params?.approvalId),
observedAt: new Date().toISOString()
};
appendTrace({
type: "client_request",
status: pending ? "started" : decision === "unsupported" || decision === "denied" ? "degraded" : "completed",
label: `client/request:${decision}`,
toolName: handledMethod,
itemId: optionalId(params?.itemId ?? params?.targetItemId),
message: `Handled app-server client request ${handledMethod} with decision=${decision}.`,
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
waitingFor: pending ? "client/request" : "turn/completed"
});
return;
}
if (method === "item/started" && item?.type === "commandExecution") {
appendCommandExecutionTrace(item, {
status: "started",
label: "item/commandExecution:started"
});
return;
}
if (method === "item/completed" && item?.type === "commandExecution") {
appendCommandExecutionTrace(item, {
status: appServerCommandStatus(item),
label: "item/commandExecution:completed"
});
return;
}
if (method === "item/commandExecution/outputDelta" || method === "item/reasoning/summaryTextDelta" || method === "item/reasoning/textDelta") {
appendTrace({
type: method.startsWith("item/reasoning/") ? "reasoning" : "tool_call",
status: "output_chunk",
label: method,
outputSummary: String(params?.delta ?? "").slice(0, 400),
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
waitingFor: "turn/completed"
});
return;
}
if (method === "error") {
const error = extractAppServerRecord(params?.error);
const message = typeof error?.message === "string" ? error.message : "Codex app-server error";
const willRetry = params?.willRetry === true;
setThreadId(extractAppServerString(params, "threadId"));
setTurnId(extractAppServerString(params, "turnId"));
appendTrace({
type: willRetry ? "provider_retry" : "error",
status: willRetry ? "retrying" : "failed",
label: willRetry ? "app-server:retrying" : "app-server:error",
errorCode: willRetry ? "codex_stdio_provider_retry" : "codex_stdio_failed",
message: redactText(message),
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
waitingFor: willRetry ? "turn/completed" : "user-retry",
terminal: !willRetry
});
if (!willRetry) finish("failed", message);
return;
}
if (method === "turn/completed") {
const status = appServerTerminalStatus(turn?.status);
const error = extractAppServerRecord(turn?.error);
const message = typeof error?.message === "string" ? error.message : null;
appendTrace({
type: "turn",
status,
label: `turn:completed:${status ?? "unknown"}`,
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
errorCode: status === "completed" ? null : "codex_stdio_failed",
message: message ? redactText(message) : undefined,
terminal: status !== "completed"
});
finish(status, message);
}
}
function appendCommandExecutionTrace(item, { status, label }) {
const output = commandExecutionOutputText(item);
const command = commandExecutionCommandText(item);
const boundedCommand = boundToolOutput(redactText(command), 1600);
const boundedOutput = output ? boundToolOutput(redactText(tailText(output, 2000)), 2000) : { text: "", truncated: false };
appendTrace({
type: "tool_call",
stage: "tool_call",
status,
label,
itemId: optionalId(item.id),
toolName: "commandExecution",
command: boundedCommand.text,
commandBytes: Buffer.byteLength(command, "utf8"),
commandTruncated: boundedCommand.truncated,
exitCode: Number.isInteger(item.exitCode) ? item.exitCode : undefined,
durationMs: typeof item.durationMs === "number" ? item.durationMs : undefined,
outputBytes: Buffer.byteLength(output, "utf8"),
stdoutSummary: boundedOutput.text || undefined,
stderrSummary: commandExecutionErrorText(item),
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
waitingFor: status === "completed" ? "turn/completed" : "commandExecution/completed"
});
}
function diagnoseNoProgress(kind, activity = activitySnapshot()) {
const base = {
kind,
waitingFor: activity.waitingFor,
lastActivityAt: activity.lastActivityAt,
idleMs: activity.idleMs,
lastActivityLabel,
lastRoutedEvent,
threadId,
turnId
};
if (lastClientRequest && lastClientRequest.status !== "handled") {
return {
...base,
code: "app_server_client_request_stalled",
layer: "client-request",
summary: `Codex app-server client request ${lastClientRequest.method} did not reach a handled decision.`,
clientRequest: lastClientRequest
};
}
if (activity.waitingFor === "commandExecution/completed") {
return {
...base,
code: "app_server_command_execution_stalled",
layer: "tool-execution",
summary: "Codex app-server is waiting for commandExecution/completed; the tool command or tool transport did not return a terminal item."
};
}
if (lastRoutedEvent?.label === "item/commandExecution:completed" && activity.waitingFor === "turn/completed") {
return {
...base,
code: "app_server_tool_result_ack_stalled",
layer: "codex-kernel",
summary: "The commandExecution item completed and was routed to the turn, but Codex app-server did not emit the following turn/completed."
};
}
if (hasAssistantOutput() && activity.waitingFor === "turn/completed") {
return {
...base,
code: "app_server_assistant_finalization_stalled",
layer: "codex-kernel",
summary: "Assistant output was observed, but Codex app-server did not emit turn/completed."
};
}
if (activity.waitingFor === "assistant-message" || activity.waitingFor === "turn/completed") {
return {
...base,
code: "codex_kernel_stalled",
layer: "codex-kernel",
summary: "Codex app-server accepted the turn, but no further routed model/tool notification arrived before the no-activity timeout."
};
}
return {
...base,
code: "app_server_notification_stalled",
layer: "app-server-transport",
summary: `No routed app-server notification arrived while waiting for ${activity.waitingFor ?? "unknown"}.`
};
}
async function wait(timeoutMs, closedPromise, { hardTimeoutMs = null, diagnosticMs = CODEX_STDIO_NO_PROGRESS_DIAGNOSIS_MS } = {}) {
let activityTimer;
let hardTimer;
let progressTimer;
const timeoutPromise = new Promise((_, reject) => {
const checkActivity = () => {
const now = Date.now();
const activity = activitySnapshot(now);
if (activity.idleMs >= timeoutMs) {
const diagnosis = diagnoseNoProgress("activity-timeout", activity);
lastDiagnosis = diagnosis;
const error = appServerActivityTimeoutError({
timeoutMs,
idleMs: activity.idleMs,
lastActivityAt: activity.lastActivityAt,
waitingFor: activity.waitingFor,
threadId,
turnId,
partialAssistant: hasAssistantOutput(),
diagnosis
});
appendTrace({
type: "timeout",
status: "failed",
label: "timeout:no_activity",
errorCode: error.code,
diagnosis,
diagnosisCode: diagnosis.code,
message: error.message,
timeoutMs,
idleMs: activity.idleMs,
lastActivityAt: activity.lastActivityAt,
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
waitingFor: activity.waitingFor,
terminal: true
});
reject(error);
return;
}
activityTimer = setTimeout(checkActivity, Math.max(1, timeoutMs - activity.idleMs));
};
activityTimer = setTimeout(checkActivity, timeoutMs);
});
const progressPromise = new Promise(() => {
const tick = () => {
appendFirstTokenProgressTrace();
appendNoProgressDiagnosisTrace(Date.now(), diagnosticMs);
progressTimer = setTimeout(tick, CODEX_STDIO_FIRST_TOKEN_PROGRESS_MS);
};
progressTimer = setTimeout(tick, CODEX_STDIO_FIRST_TOKEN_PROGRESS_MS);
});
const hardTimeoutPromise = hardTimeoutMs
? new Promise((_, reject) => {
hardTimer = setTimeout(() => {
const activity = activitySnapshot();
const diagnosis = diagnoseNoProgress("hard-timeout", activity);
lastDiagnosis = diagnosis;
const error = appServerHardTimeoutError({
hardTimeoutMs,
timeoutMs,
idleMs: activity.idleMs,
lastActivityAt: activity.lastActivityAt,
waitingFor: activity.waitingFor,
threadId,
turnId,
diagnosis
});
appendTrace({
type: "timeout",
status: "failed",
label: "timeout:hard_cap",
errorCode: error.code,
diagnosis,
diagnosisCode: diagnosis.code,
message: error.message,
timeoutMs,
hardTimeoutMs,
idleMs: activity.idleMs,
lastActivityAt: activity.lastActivityAt,
sessionId: session?.sessionId,
sessionStatus: session?.status,
turn: session?.turn,
threadId,
turnId,
waitingFor: activity.waitingFor,
terminal: true
});
reject(error);
}, hardTimeoutMs);
})
: null;
const candidates = [terminalPromise];
if (closedPromise && typeof closedPromise.then === "function") {
candidates.push(closedPromise.then((exit) => {
if (terminal) return terminal;
return {
terminalStatus: "failed",
terminalError: hasAssistantOutput()
? "Codex app-server transport closed after partial assistant output but before turn/completed."
: "Codex app-server transport closed before turn/completed.",
threadId,
turnId,
...activitySnapshot(),
diagnosis: diagnoseNoProgress("transport-closed", activitySnapshot()),
appServerExit: exit
};
}));
}
try {
return await Promise.race([...candidates, timeoutPromise, hardTimeoutPromise, progressPromise].filter(Boolean));
} finally {
clearTimeout(activityTimer);
clearTimeout(hardTimer);
clearTimeout(progressTimer);
}
}
function snapshot() {
const activity = activitySnapshot();
return {
threadId,
turnId,
assistantText: redactText(assistantText),
finalResponse: redactText(firstNonEmpty(finalResponse, assistantText)),
terminalStatus: terminal?.terminalStatus ?? null,
terminalError: terminal?.terminalError ?? null,
lastActivityAt: terminal?.lastActivityAt ?? activity.lastActivityAt,
idleMs: terminal?.idleMs ?? activity.idleMs,
waitingFor: terminal?.waitingFor ?? activity.waitingFor,
diagnosis: terminal?.diagnosis ?? lastDiagnosis
};
}
return {
appendTrace,
handle,
setThreadId,
setTurnId,
wait,
snapshot
};
}
function compactTraceEvent(event = null) {
if (!event || typeof event !== "object") return null;
return {
type: event.type ?? null,
status: event.status ?? null,
label: event.label ?? null,
stage: event.stage ?? null,
toolName: event.toolName ?? null,
itemId: event.itemId ?? null,
waitingFor: event.waitingFor ?? null,
exitCode: Number.isInteger(event.exitCode) ? event.exitCode : null,
durationMs: Number.isFinite(event.durationMs) ? event.durationMs : null
};
}
function appServerWaitingForMethod(method) {
if (method === "client/request/handled") return "turn/completed";
if (method === "thread/started") return "turn/start";
if (method === "turn/started") return "assistant-message";
if (method === "item/agentMessage/delta") return "turn/completed";
if (method === "item/completed") return "turn/completed";
if (method === "item/commandExecution/outputDelta") return "turn/completed";
if (method === "item/reasoning/summaryTextDelta" || method === "item/reasoning/textDelta") return "turn/completed";
if (method === "error") return "turn/completed";
if (method === "turn/completed") return "turn/completed";
return "app-server-notification";
}