From 9e225d7db98261eb784661186a74a68133d54554 Mon Sep 17 00:00:00 2001 From: Mohamed Boudra Date: Fri, 6 Feb 2026 12:42:18 +0700 Subject: [PATCH] Fix Codex persisted timeline parsing and session persistence --- .../src/server/agent/agent-manager.test.ts | 121 +++++ .../server/src/server/agent/agent-manager.ts | 23 +- .../agent/providers/codex-app-server-agent.ts | 49 +- .../server/agent/providers/codex-mcp-agent.ts | 462 ++++++++++-------- .../providers/codex-rollout-parsing.test.ts | 73 +++ 5 files changed, 519 insertions(+), 209 deletions(-) diff --git a/packages/server/src/server/agent/agent-manager.test.ts b/packages/server/src/server/agent/agent-manager.test.ts index dc65cbd1f..067e3faf3 100644 --- a/packages/server/src/server/agent/agent-manager.test.ts +++ b/packages/server/src/server/agent/agent-manager.test.ts @@ -494,4 +494,125 @@ describe("AgentManager", () => { const updatedAgent = manager.getAgent(snapshot.id); expect(updatedAgent?.currentModeId).toBe("acceptEdits"); }); + + test("close during in-flight stream does not clear persistence sessionId", async () => { + const workdir = mkdtempSync(join(tmpdir(), "agent-manager-test-")); + const storagePath = join(workdir, "agents"); + const storage = new AgentStorage(storagePath, logger); + + class CloseRaceSession implements AgentSession { + readonly provider = "codex" as const; + readonly capabilities = TEST_CAPABILITIES; + readonly id = randomUUID(); + private threadId: string | null = this.id; + private releaseStream: (() => void) | null = null; + private closed = false; + + async run(): Promise { + return { sessionId: this.id, finalText: "", timeline: [] }; + } + + async *stream(): AsyncGenerator { + yield { type: "turn_started", provider: this.provider }; + if (!this.closed) { + await new Promise((resolve) => { + this.releaseStream = resolve; + }); + } + yield { type: "turn_canceled", provider: this.provider, reason: "closed" }; + } + + async *streamHistory(): AsyncGenerator {} + + async getRuntimeInfo() { + return { + provider: this.provider, + sessionId: this.threadId, + model: null, + modeId: null, + }; + } + + async getAvailableModes() { + return []; + } + + async getCurrentMode() { + return null; + } + + async setMode(): Promise {} + + getPendingPermissions() { + return []; + } + + async respondToPermission(): Promise {} + + describePersistence() { + if (!this.threadId) { + return null; + } + return { provider: this.provider, sessionId: this.threadId }; + } + + async interrupt(): Promise {} + + async close(): Promise { + this.closed = true; + this.threadId = null; + this.releaseStream?.(); + } + } + + class CloseRaceClient implements AgentClient { + readonly provider = "codex" as const; + readonly capabilities = TEST_CAPABILITIES; + + async isAvailable(): Promise { + return true; + } + + async createSession(): Promise { + return new CloseRaceSession(); + } + + async resumeSession(): Promise { + return new CloseRaceSession(); + } + } + + const manager = new AgentManager({ + clients: { + codex: new CloseRaceClient(), + }, + registry: storage, + logger, + idFactory: () => "close-race-agent", + }); + + const snapshot = await manager.createAgent({ + provider: "codex", + cwd: workdir, + }); + + const stream = manager.streamAgent(snapshot.id, "hello"); + await stream.next(); + + await manager.closeAgent(snapshot.id); + + // Drain stream finalizer path after close(). + while (true) { + const next = await stream.next(); + if (next.done) { + break; + } + } + + await manager.flush(); + await storage.flush(); + + const persisted = await storage.get(snapshot.id); + expect(persisted?.persistence?.sessionId).toBe(snapshot.persistence?.sessionId); + }); }); diff --git a/packages/server/src/server/agent/agent-manager.ts b/packages/server/src/server/agent/agent-manager.ts index 14f56d386..f3578ba9b 100644 --- a/packages/server/src/server/agent/agent-manager.ts +++ b/packages/server/src/server/agent/agent-manager.ts @@ -601,13 +601,17 @@ export class AgentManager { mutableAgent.pendingRun = null; mutableAgent.lifecycle = error ? "error" : "idle"; mutableAgent.lastError = error; - mutableAgent.persistence = attachPersistenceCwd( + const persistenceHandle = mutableAgent.session.describePersistence() ?? - (mutableAgent.runtimeInfo?.sessionId - ? { provider: mutableAgent.provider, sessionId: mutableAgent.runtimeInfo.sessionId } - : null), - mutableAgent.cwd - ); + (mutableAgent.runtimeInfo?.sessionId + ? { provider: mutableAgent.provider, sessionId: mutableAgent.runtimeInfo.sessionId } + : null); + if (persistenceHandle) { + mutableAgent.persistence = attachPersistenceCwd( + persistenceHandle, + mutableAgent.cwd + ); + } this.emitState(mutableAgent); }; @@ -1113,7 +1117,12 @@ export class AgentManager { case "thread_started": // Update persistence with the new session ID from the provider. // persistence.sessionId is the single source of truth for session identity. - agent.persistence = attachPersistenceCwd(agent.session.describePersistence(), agent.cwd); + { + const handle = agent.session.describePersistence(); + if (handle) { + agent.persistence = attachPersistenceCwd(handle, agent.cwd); + } + } break; case "timeline": this.recordTimeline(agent, event.item); diff --git a/packages/server/src/server/agent/providers/codex-app-server-agent.ts b/packages/server/src/server/agent/providers/codex-app-server-agent.ts index 9d34d3270..abe0d5a60 100644 --- a/packages/server/src/server/agent/providers/codex-app-server-agent.ts +++ b/packages/server/src/server/agent/providers/codex-app-server-agent.ts @@ -33,6 +33,7 @@ import { Dirent } from "node:fs"; import os from "node:os"; import path from "node:path"; import readline from "node:readline"; +import { loadCodexPersistedTimeline } from "./codex-mcp-agent.js"; const DEFAULT_TIMEOUT_MS = 14 * 24 * 60 * 60 * 1000; @@ -1086,22 +1087,40 @@ class CodexAppServerAgentSession implements AgentSession { private async loadPersistedHistory(): Promise { if (!this.client || !this.currentThreadId) return; try { + let rolloutTimeline: AgentTimelineItem[] = []; + try { + rolloutTimeline = await loadCodexPersistedTimeline( + this.currentThreadId, + undefined, + this.logger + ); + } catch { + rolloutTimeline = []; + } + const response = (await this.client.request("thread/read", { threadId: this.currentThreadId, includeTurns: true, })) as { thread?: { turns?: Array<{ items?: any[] }> } }; const thread = response?.thread; - if (!thread || !Array.isArray(thread.turns)) return; - const timeline: AgentTimelineItem[] = []; - for (const turn of thread.turns) { - const items = Array.isArray(turn.items) ? turn.items : []; - for (const item of items) { - const timelineItem = threadItemToTimeline(item, { cwd: this.config.cwd ?? null }); - if (timelineItem) { - timeline.push(timelineItem); + const threadTimeline: AgentTimelineItem[] = []; + if (thread && Array.isArray(thread.turns)) { + for (const turn of thread.turns) { + const items = Array.isArray(turn.items) ? turn.items : []; + for (const item of items) { + const timelineItem = threadItemToTimeline(item, { + cwd: this.config.cwd ?? null, + }); + if (timelineItem) { + threadTimeline.push(timelineItem); + } } } } + + const timeline = + rolloutTimeline.length > 0 ? rolloutTimeline : threadTimeline; + if (timeline.length > 0) { this.persistedHistory = timeline; this.historyPending = true; @@ -1775,19 +1794,27 @@ export class CodexAppServerAgentClient implements AgentClient { const title = thread.preview ?? null; let timeline: AgentTimelineItem[] = []; try { + const rolloutTimeline = await loadCodexPersistedTimeline( + threadId, + undefined, + this.logger + ); const read = (await client.request("thread/read", { threadId, includeTurns: true, })) as { thread?: { turns?: Array<{ items?: any[] }> } }; const turns = read.thread?.turns ?? []; - const items: AgentTimelineItem[] = []; + const itemsFromThreadRead: AgentTimelineItem[] = []; for (const turn of turns) { for (const item of turn.items ?? []) { const timelineItem = threadItemToTimeline(item, { cwd }); - if (timelineItem) items.push(timelineItem); + if (timelineItem) itemsFromThreadRead.push(timelineItem); } } - timeline = items; + timeline = + rolloutTimeline.length > 0 + ? rolloutTimeline + : itemsFromThreadRead; } catch { timeline = []; } diff --git a/packages/server/src/server/agent/providers/codex-mcp-agent.ts b/packages/server/src/server/agent/providers/codex-mcp-agent.ts index aa907135a..dd7fba32b 100644 --- a/packages/server/src/server/agent/providers/codex-mcp-agent.ts +++ b/packages/server/src/server/agent/providers/codex-mcp-agent.ts @@ -5097,10 +5097,20 @@ async function findRolloutFile( return null; } -const RolloutContentItemSchema = z.object({ - type: z.string(), - text: z.string(), -}).passthrough(); +const RolloutContentItemSchema = z + .union([ + z.object({ type: z.literal("input_text"), text: z.string() }), + z.object({ type: z.literal("output_text"), text: z.string() }), + z.object({ type: z.literal("reasoning_text"), text: z.string() }), + z.object({ type: z.literal("text"), text: z.string() }), + z + .object({ + type: z.string(), + text: z.string().optional(), + message: z.string().optional(), + }) + .passthrough(), + ]); const RolloutContentArraySchema = z.array(RolloutContentItemSchema); @@ -5116,77 +5126,131 @@ function extractContentTextByType(content: unknown, itemType: string): string { .trim(); } -const RolloutResponsePayloadSchema = z.object({ - type: z.string().optional(), - role: z.string().optional(), +const RolloutResponseMessagePayloadSchema = z.object({ + type: z.literal("message"), + role: z.enum(["user", "assistant"]).optional(), + content: z.unknown().optional(), +}); + +const RolloutResponseReasoningPayloadSchema = z.object({ + type: z.literal("reasoning"), content: z.unknown().optional(), - name: z.string().optional(), - call_id: z.string().optional(), - arguments: z.string().optional(), - output: z.string().optional(), summary: z.array(z.object({ text: z.string().optional() })).optional(), text: z.string().optional(), }); -const RolloutEventPayloadSchema = z.object({ - type: z.string().optional(), +const RolloutResponseFunctionCallPayloadSchema = z.object({ + type: z.literal("function_call"), + name: z.string().optional(), + call_id: z.string().optional(), + arguments: z.string().optional(), +}); + +const RolloutResponseCustomToolCallPayloadSchema = z.object({ + type: z.literal("custom_tool_call"), + name: z.string().optional(), + call_id: z.string().optional(), + arguments: z.string().optional(), +}); + +const RolloutResponseFunctionCallOutputPayloadSchema = z.object({ + type: z.literal("function_call_output"), + call_id: z.string().optional(), + output: z.string().optional(), +}); + +const RolloutEventAgentReasoningPayloadSchema = z.object({ + type: z.literal("agent_reasoning"), text: z.string().optional(), - message: z.string().optional(), }); -const RolloutEntrySchema = z.object({ - type: z.enum(["response_item", "event_msg"]), - payload: z.unknown().optional(), +const RolloutEventAgentMessagePayloadSchema = z.object({ + type: z.literal("agent_message"), + message: z + .union([ + z.string(), + z + .object({ + role: z.string().optional(), + message: z.string().optional(), + text: z.string().optional(), + }) + .passthrough(), + ]) + .optional(), }); -type RolloutEntry = z.infer; -type RolloutResponsePayload = z.infer; +const RolloutEventUserMessagePayloadSchema = z.object({ + type: z.literal("user_message"), + message: z + .union([ + z.string(), + z + .object({ + role: z.string().optional(), + message: z.string().optional(), + text: z.string().optional(), + }) + .passthrough(), + ]) + .optional(), +}); -function parseRolloutEntryFromLine(line: string): RolloutEntry | null { - if (!line) { - return null; - } - try { - const parsed = JSON.parse(line); - const result = RolloutEntrySchema.safeParse(parsed); - if (result.success) { - return result.data; - } - if ( - parsed && - typeof parsed === "object" && - typeof (parsed as { output?: unknown }).output === "string" - ) { - return parseRolloutEntryFromLine((parsed as { output: string }).output); - } - } catch { - return null; - } - return null; -} +type RolloutResponseReasoningPayload = z.infer< + typeof RolloutResponseReasoningPayloadSchema +>; +type ParsedRolloutRecord = + | { kind: "timeline"; item: AgentTimelineItem } + | { kind: "call"; name: string; callId?: string; input?: unknown } + | { kind: "output"; callId: string; output: string } + | { kind: "ignore" }; + +const RolloutMessageContentSchema = z + .union([ + z.string().transform((content) => content.trim()), + z.array( + z + .object({ + text: z.string().optional(), + message: z.string().optional(), + }) + .passthrough() + ).transform((content) => + content + .map((block) => block.text ?? block.message ?? "") + .map((text) => text.trim()) + .filter(Boolean) + .join("\n") + .trim() + ), + ]); function extractMessageText(content: unknown): string { - if (!Array.isArray(content)) { + const parsed = RolloutMessageContentSchema.safeParse(content); + if (!parsed.success) { return ""; } - const parts: string[] = []; - for (const block of content) { - if (!block || typeof block !== "object") { - continue; - } - const record = block as Record; - const text = typeof record.text === "string" ? record.text : undefined; - if (text && text.trim()) { - parts.push(text.trim()); - continue; - } - const message = - typeof record.message === "string" ? record.message : undefined; - if (message && message.trim()) { - parts.push(message.trim()); - } + return parsed.data; +} + +const RolloutEventMessageTextSchema = z + .union([ + z.string(), + z + .object({ + message: z.string().optional(), + text: z.string().optional(), + }) + .passthrough() + .transform((message) => message.message ?? message.text ?? ""), + ]); + +function extractEventMessageText(message: unknown): string { + const parsed = RolloutEventMessageTextSchema.safeParse(message); + if (!parsed.success) { + return ""; } - return parts.join("\n").trim(); + return parsed.data; } function isSyntheticRolloutUserMessage(text: string): boolean { @@ -5207,10 +5271,10 @@ function isSyntheticRolloutUserMessage(text: string): boolean { return false; } -function extractReasoningText(payload: RolloutResponsePayload): string { - if (Array.isArray(payload?.summary)) { +function extractReasoningText(payload: RolloutResponseReasoningPayload): string { + if (Array.isArray(payload.summary)) { const text = payload.summary - .map((item) => (item && typeof item.text === "string" ? item.text : "")) + .map((item) => item.text ?? "") .filter(Boolean) .join("\n") .trim(); @@ -5218,17 +5282,117 @@ function extractReasoningText(payload: RolloutResponsePayload): string { return text; } } - // Handle content array with reasoning_text items const contentText = extractContentTextByType(payload.content, "reasoning_text"); if (contentText) { return contentText; } - if (typeof payload?.text === "string") { + if (typeof payload.text === "string") { return payload.text; } return ""; } +function parseJsonLikeString(value: string): unknown { + try { + return JSON.parse(value); + } catch { + return value; + } +} + +const FunctionCallInputNormalizationSchema = z + .union([ + z.object({ cmd: z.string() }).transform((input) => ({ name: "Bash", input: { command: input.cmd } })), + z + .object({ command: z.array(z.string()) }) + .transform((input) => ({ name: "Bash", input: { command: input.command[2] ?? "" } })), + z.unknown().transform((input) => ({ name: "unknown", input })), + ]); + +const RolloutResponseRecordSchema = z + .union([ + RolloutResponseMessagePayloadSchema.transform((payload): ParsedRolloutRecord => { + const text = extractMessageText(payload.content); + const itemType = payload.role === "assistant" ? "assistant_message" : "user_message"; + const shouldEmit = text.length > 0 && (itemType !== "user_message" || !isSyntheticRolloutUserMessage(text)); + return shouldEmit + ? { kind: "timeline", item: { type: itemType, text } } + : { kind: "ignore" }; + }), + RolloutResponseReasoningPayloadSchema.transform((payload): ParsedRolloutRecord => { + const text = extractReasoningText(payload); + return text.length > 0 + ? { kind: "timeline", item: { type: "reasoning", text } } + : { kind: "ignore" }; + }), + z + .union([ + RolloutResponseFunctionCallPayloadSchema, + RolloutResponseCustomToolCallPayloadSchema, + ]) + .transform((payload): ParsedRolloutRecord => { + const rawName = payload.name ?? "unknown"; + const parsedArguments = payload.arguments ? parseJsonLikeString(payload.arguments) : undefined; + const normalized = + rawName === "exec_command" || rawName === "shell" + ? FunctionCallInputNormalizationSchema.parse(parsedArguments) + : { name: rawName, input: parsedArguments }; + const skip = rawName === "write_stdin"; + return skip + ? { kind: "ignore" } + : { + kind: "call", + name: normalized.name, + callId: payload.call_id, + input: normalized.input, + }; + }), + RolloutResponseFunctionCallOutputPayloadSchema.transform( + (payload): ParsedRolloutRecord => + payload.call_id && payload.output + ? { kind: "output", callId: payload.call_id, output: payload.output } + : { kind: "ignore" } + ), + z.unknown().transform((): ParsedRolloutRecord => ({ kind: "ignore" })), + ]); + +const RolloutEventRecordSchema = z + .union([ + RolloutEventAgentReasoningPayloadSchema.transform( + (payload): ParsedRolloutRecord => + payload.text + ? { kind: "timeline", item: { type: "reasoning", text: payload.text } } + : { kind: "ignore" } + ), + RolloutEventAgentMessagePayloadSchema.transform((payload): ParsedRolloutRecord => { + const text = extractEventMessageText(payload.message); + return text.length > 0 + ? { kind: "timeline", item: { type: "assistant_message", text } } + : { kind: "ignore" }; + }), + RolloutEventUserMessagePayloadSchema.transform((payload): ParsedRolloutRecord => { + const text = extractEventMessageText(payload.message); + const shouldEmit = text.length > 0 && !isSyntheticRolloutUserMessage(text); + return shouldEmit + ? { kind: "timeline", item: { type: "user_message", text } } + : { kind: "ignore" }; + }), + z.unknown().transform((): ParsedRolloutRecord => ({ kind: "ignore" })), + ]); + +const RolloutRecordSchema = z + .object({ + type: z.enum(["response_item", "event_msg"]), + payload: z.unknown().optional(), + item: z.unknown().optional(), + msg: z.unknown().optional(), + }) + .transform((entry) => + entry.type === "response_item" + ? RolloutResponseRecordSchema.parse(entry.payload ?? entry.item) + : RolloutEventRecordSchema.parse(entry.payload ?? entry.msg) + ); + function parseJsonRolloutTimeline( parsed: unknown ): AgentTimelineItem[] | null { @@ -5244,25 +5408,28 @@ function parseJsonRolloutTimeline( if (!entry || typeof entry !== "object") { continue; } - const record = entry as Record; - const type = record.type; - if (type === "message") { - const role = record.role; - const text = extractMessageText(record.content); - if (!text || typeof role !== "string") { + const messagePayloadResult = + RolloutResponseMessagePayloadSchema.safeParse(entry); + if (messagePayloadResult.success) { + const payload = messagePayloadResult.data; + const text = extractMessageText(payload.content); + if (!text) { continue; } - if (role === "assistant") { + if (payload.role === "assistant") { timeline.push({ type: "assistant_message", text }); - } else if (role === "user") { + } else if (payload.role === "user") { if (!isSyntheticRolloutUserMessage(text)) { timeline.push({ type: "user_message", text }); } } continue; } - if (type === "reasoning") { - const text = extractReasoningText(record as RolloutResponsePayload); + + const reasoningPayloadResult = + RolloutResponseReasoningPayloadSchema.safeParse(entry); + if (reasoningPayloadResult.success) { + const text = extractReasoningText(reasoningPayloadResult.data); if (text) { timeline.push({ type: "reasoning", text }); } @@ -5294,126 +5461,39 @@ export async function parseRolloutFile( .map((line) => line.trim()) .filter(Boolean); - // First pass: collect function_call_output entries by call_id - const outputsByCallId = new Map(); - for (const line of lines) { - const entry = parseRolloutEntryFromLine(line); - if (!entry || entry.type !== "response_item") continue; - const payloadResult = RolloutResponsePayloadSchema.safeParse(entry.payload); - if (!payloadResult.success) continue; - const payload = payloadResult.data; - if (payload.type === "function_call_output" && payload.call_id && payload.output) { - outputsByCallId.set(payload.call_id, payload.output); - } - } - - // Second pass: build timeline - const timeline: AgentTimelineItem[] = []; - - for (const line of lines) { - const entry = parseRolloutEntryFromLine(line); - if (!entry) continue; - - if (entry.type === "response_item") { - const payloadResult = RolloutResponsePayloadSchema.safeParse(entry.payload); - if (!payloadResult.success) continue; - const payload = payloadResult.data; - - switch (payload.type) { - case "message": { - const text = extractMessageText(payload.content); - if (text) { - if (payload.role === "assistant") { - timeline.push({ type: "assistant_message", text }); - } else if (payload.role === "user") { - if (!isSyntheticRolloutUserMessage(text)) { - timeline.push({ type: "user_message", text }); - } - } - } - break; - } - case "reasoning": { - const text = extractReasoningText(payload); - if (text) { - timeline.push({ type: "reasoning", text }); - } - break; - } - case "function_call": - case "custom_tool_call": { - const rawName = payload.name ?? "unknown"; - const callId = payload.call_id; - - // Skip internal polling calls - if (rawName === "write_stdin") { - break; - } - - let input: unknown; - if (payload.arguments) { - try { - input = JSON.parse(payload.arguments); - } catch { - input = payload.arguments; - } - } - - // Map exec_command and shell to Bash with normalized input - let name = rawName; - if (rawName === "exec_command" && input && typeof input === "object") { - const execInput = input as { cmd?: string }; - if (execInput.cmd) { - name = "Bash"; - input = { command: execInput.cmd }; - } - } else if (rawName === "shell" && input && typeof input === "object") { - // Older format: { command: ["bash", "-lc", "actual cmd"], workdir: "..." } - const shellInput = input as { command?: string[] }; - if (Array.isArray(shellInput.command) && shellInput.command.length >= 3) { - name = "Bash"; - // command[2] is the actual shell command after "bash -lc" - input = { command: shellInput.command[2] }; - } - } - - // Attach output if available - const output = callId ? outputsByCallId.get(callId) : undefined; - - timeline.push({ - type: "tool_call", - name, - callId, - status: "completed", - input, - ...(output ? { output } : {}), - }); - break; - } - case "function_call_output": - // Already processed in first pass - break; - default: - break; + const parsedRecords = lines + .map((line) => { + try { + return JSON.parse(line); + } catch { + return null; } - } else if (entry.type === "event_msg") { - const payloadResult = RolloutEventPayloadSchema.safeParse(entry.payload); - if (!payloadResult.success) continue; - const payload = payloadResult.data; + }) + .filter((record): record is unknown => record !== null) + .map((record) => RolloutRecordSchema.safeParse(record)) + .filter((result): result is { success: true; data: ParsedRolloutRecord } => result.success) + .map((result) => result.data); - if (payload.type === "agent_reasoning" && payload.text) { - timeline.push({ type: "reasoning", text: payload.text }); - } else if (payload.type === "agent_message" && payload.message) { - timeline.push({ type: "assistant_message", text: payload.message }); - } else if (payload.type === "user_message" && payload.message) { - if (!isSyntheticRolloutUserMessage(payload.message)) { - timeline.push({ type: "user_message", text: payload.message }); - } - } - } - } + const outputsByCallId = parsedRecords + .filter((record): record is Extract => record.kind === "output") + .reduce((map, record) => map.set(record.callId, record.output), new Map()); - return timeline; + return parsedRecords.flatMap((record): AgentTimelineItem[] => + record.kind === "timeline" + ? [record.item] + : record.kind === "call" + ? [ + { + type: "tool_call", + name: record.name, + callId: record.callId, + status: "completed", + input: record.input, + output: record.callId ? outputsByCallId.get(record.callId) : undefined, + }, + ] + : [] + ); } type CodexPersistedTimelineOptions = { diff --git a/packages/server/src/server/agent/providers/codex-rollout-parsing.test.ts b/packages/server/src/server/agent/providers/codex-rollout-parsing.test.ts index dfa050e16..68a70c7d1 100644 --- a/packages/server/src/server/agent/providers/codex-rollout-parsing.test.ts +++ b/packages/server/src/server/agent/providers/codex-rollout-parsing.test.ts @@ -216,6 +216,79 @@ describe("codex rollout parsing", () => { text: "Let me think about this.", }); }); + + test("parses legacy response_item shape using item", async () => { + const rolloutPath = join(tmpDir, "rollout.jsonl"); + const lines = [ + JSON.stringify({ + timestamp: "2026-01-22T07:08:54.378Z", + type: "response_item", + item: { + type: "function_call", + name: "exec_command", + arguments: '{"cmd":"echo hello"}', + call_id: "call_legacy_1", + }, + }), + JSON.stringify({ + timestamp: "2026-01-22T07:09:01.785Z", + type: "response_item", + item: { + type: "function_call_output", + call_id: "call_legacy_1", + output: "hello", + }, + }), + ]; + writeFileSync(rolloutPath, lines.join("\n") + "\n"); + + const timeline = await parseRolloutFile(rolloutPath); + const toolCalls = timeline.filter((i) => i.type === "tool_call"); + expect(toolCalls.length).toBe(1); + expect(toolCalls[0]).toMatchObject({ + type: "tool_call", + name: "Bash", + callId: "call_legacy_1", + input: { command: "echo hello" }, + output: "hello", + }); + }); + + test("parses legacy event_msg shape using msg", async () => { + const rolloutPath = join(tmpDir, "rollout.jsonl"); + const lines = [ + JSON.stringify({ + timestamp: "2026-01-22T07:08:54.378Z", + type: "event_msg", + msg: { + type: "agent_reasoning", + text: "thinking", + }, + }), + JSON.stringify({ + timestamp: "2026-01-22T07:08:54.378Z", + type: "event_msg", + msg: { + type: "agent_message", + message: { role: "assistant", message: "done" }, + }, + }), + JSON.stringify({ + timestamp: "2026-01-22T07:08:54.378Z", + type: "event_msg", + msg: { + type: "user_message", + message: { role: "user", message: "question" }, + }, + }), + ]; + writeFileSync(rolloutPath, lines.join("\n") + "\n"); + + const timeline = await parseRolloutFile(rolloutPath); + expect(timeline).toContainEqual({ type: "reasoning", text: "thinking" }); + expect(timeline).toContainEqual({ type: "assistant_message", text: "done" }); + expect(timeline).toContainEqual({ type: "user_message", text: "question" }); + }); }); describe("complex conversation", () => {