From 8d9ce2adb7a77cb811dd8960899ba83a550ec439 Mon Sep 17 00:00:00 2001 From: Nikita Ashikhmin Date: Tue, 7 Jul 2026 15:57:15 +0400 Subject: [PATCH] Expose Codex agent message phases in ACP chunks --- src/CodexAcpServer.ts | 6 +- src/CodexEventHandler.ts | 17 +++- src/ContentChunks.ts | 29 +++++-- src/ResponseItemHistoryFallback.ts | 9 +- .../agent-message-events.test.ts | 85 +++++++++++++++++++ .../data/agent-message-phases.json | 42 +++++++++ .../response-item-history-fallback.test.ts | 28 ++++++ 7 files changed, 201 insertions(+), 15 deletions(-) create mode 100644 src/__tests__/CodexACPAgent/agent-message-events.test.ts create mode 100644 src/__tests__/CodexACPAgent/data/agent-message-phases.json diff --git a/src/CodexAcpServer.ts b/src/CodexAcpServer.ts index 6454ca86..22d5e703 100644 --- a/src/CodexAcpServer.ts +++ b/src/CodexAcpServer.ts @@ -69,6 +69,7 @@ import packageJson from "../package.json"; import {isJetBrains2026_1Client} from "./JBUtils"; import {resolveTerminalOutputMode, type TerminalOutputMode} from "./TerminalOutputMode"; import { + createCodexMessagePhaseMeta, createAgentTextMessageChunk, createAgentTextThoughtChunk, createUserMessageChunk, @@ -983,12 +984,15 @@ export class CodexAcpServer { case "subAgentActivity": case "sleep": return []; - case "agentMessage": + case "agentMessage": { + const meta = createCodexMessagePhaseMeta(item.phase); return [{ sessionUpdate: "agent_message_chunk", messageId: item.id, content: { type: "text", text: item.text }, + ...(meta ? { _meta: meta } : {}), }]; + } case "reasoning": return this.createReasoningUpdates(item); case "fileChange": diff --git a/src/CodexEventHandler.ts b/src/CodexEventHandler.ts index 5ed8ad79..c2a4d672 100644 --- a/src/CodexEventHandler.ts +++ b/src/CodexEventHandler.ts @@ -56,6 +56,7 @@ import { import { stripShellPrefix } from "./CommandUtils"; import {createTerminalOutputMeta, type TerminalOutputMode} from "./TerminalOutputMode"; import { + createCodexMessagePhaseMeta, createAgentTextMessageChunk, createAgentTextThoughtChunk, } from "./ContentChunks"; @@ -74,6 +75,7 @@ export class CodexEventHandler { private readonly seenReasoningDeltaItemIds = new Set(); private readonly terminalCommandIds = new Set(); private readonly terminalCommandOutputIds = new Set(); + private readonly agentMessagePhases = new Map(); constructor(connection: AcpClientConnection, sessionState: SessionState) { this.connection = connection; @@ -230,7 +232,8 @@ export class CodexEventHandler { } private async createTextEvent(event: AgentMessageDeltaNotification): Promise { - return createAgentTextMessageChunk(event.delta, event.itemId); + const phase = this.agentMessagePhases.get(event.itemId) ?? null; + return createAgentTextMessageChunk(event.delta, event.itemId, createCodexMessagePhaseMeta(phase)); } private async createConfigWarningEvent(event: ConfigWarningNotification): Promise { @@ -331,11 +334,13 @@ export class CodexEventHandler { return createImageGenerationStartUpdate(event.item); case "collabAgentToolCall": return createCollabAgentToolCallUpdate(event.item); + case "agentMessage": + this.rememberAgentMessagePhase(event.item); + return null; case "subAgentActivity": case "sleep": case "userMessage": case "hookPrompt": - case "agentMessage": case "reasoning": case "enteredReviewMode": case "exitedReviewMode": @@ -383,6 +388,9 @@ export class CodexEventHandler { return createWebSearchCompleteUpdate(event.item); case "collabAgentToolCall": return createCollabAgentToolCallCompleteUpdate(event.item); + case "agentMessage": + this.rememberAgentMessagePhase(event.item); + return null; case "exitedReviewMode": return this.createExitedReviewModeEvent(event.item); case "contextCompaction": @@ -392,7 +400,6 @@ export class CodexEventHandler { case "sleep": case "userMessage": case "hookPrompt": - case "agentMessage": case "enteredReviewMode": case "plan": return null; @@ -400,6 +407,10 @@ export class CodexEventHandler { } } + private rememberAgentMessagePhase(item: ThreadItem & { type: "agentMessage" }): void { + this.agentMessagePhases.set(item.id, item.phase); + } + private createCompletedReasoningEvent(item: ThreadItem & { type: "reasoning" }): UpdateSessionEvent | null { const parts = item.summary.length > 0 ? item.summary : item.content; const text = parts.filter(part => part.length > 0).join("\n\n"); diff --git a/src/ContentChunks.ts b/src/ContentChunks.ts index b111d347..2bef82e8 100644 --- a/src/ContentChunks.ts +++ b/src/ContentChunks.ts @@ -1,52 +1,67 @@ import type {ContentBlock} from "@agentclientprotocol/sdk"; import type {UpdateSessionEvent} from "./ACPSessionConnection"; -export function createUserMessageChunk(content: ContentBlock, messageId?: string): UpdateSessionEvent { +type AcpMeta = Record; + +export function createCodexMessagePhaseMeta(phase: string | null | undefined): AcpMeta | undefined { + if (!phase) { + return undefined; + } + return { codex: { phase } }; +} + +export function createUserMessageChunk(content: ContentBlock, messageId?: string, meta?: AcpMeta): UpdateSessionEvent { if (messageId) { return { sessionUpdate: "user_message_chunk", messageId, content, + ...(meta ? { _meta: meta } : {}), }; } return { sessionUpdate: "user_message_chunk", content, + ...(meta ? { _meta: meta } : {}), }; } -export function createAgentMessageChunk(content: ContentBlock, messageId?: string): UpdateSessionEvent { +export function createAgentMessageChunk(content: ContentBlock, messageId?: string, meta?: AcpMeta): UpdateSessionEvent { if (messageId) { return { sessionUpdate: "agent_message_chunk", messageId, content, + ...(meta ? { _meta: meta } : {}), }; } return { sessionUpdate: "agent_message_chunk", content, + ...(meta ? { _meta: meta } : {}), }; } -export function createAgentThoughtChunk(content: ContentBlock, messageId?: string): UpdateSessionEvent { +export function createAgentThoughtChunk(content: ContentBlock, messageId?: string, meta?: AcpMeta): UpdateSessionEvent { if (messageId) { return { sessionUpdate: "agent_thought_chunk", messageId, content, + ...(meta ? { _meta: meta } : {}), }; } return { sessionUpdate: "agent_thought_chunk", content, + ...(meta ? { _meta: meta } : {}), }; } -export function createAgentTextMessageChunk(text: string, messageId?: string): UpdateSessionEvent { - return createAgentMessageChunk({type: "text", text}, messageId); +export function createAgentTextMessageChunk(text: string, messageId?: string, meta?: AcpMeta): UpdateSessionEvent { + return createAgentMessageChunk({type: "text", text}, messageId, meta); } -export function createAgentTextThoughtChunk(text: string, messageId?: string): UpdateSessionEvent { - return createAgentThoughtChunk({type: "text", text}, messageId); +export function createAgentTextThoughtChunk(text: string, messageId?: string, meta?: AcpMeta): UpdateSessionEvent { + return createAgentThoughtChunk({type: "text", text}, messageId, meta); } diff --git a/src/ResponseItemHistoryFallback.ts b/src/ResponseItemHistoryFallback.ts index 11e49d68..0730455f 100644 --- a/src/ResponseItemHistoryFallback.ts +++ b/src/ResponseItemHistoryFallback.ts @@ -6,6 +6,7 @@ import { stripShellPrefix } from "./CommandUtils"; import type { CommandAction, Thread, ThreadItem } from "./app-server/v2"; import { createCommandActionEvent } from "./CodexToolCallMapper"; import { createTerminalOutputMeta, type TerminalOutputMode } from "./TerminalOutputMode"; +import { createAgentMessageChunk, createCodexMessagePhaseMeta } from "./ContentChunks"; type JsonRecord = Record; type AcpToolCallEvent = Extract; @@ -234,10 +235,10 @@ function createMessageUpdates(item: JsonRecord): UpdateSessionEvent[] { return []; } - return contentBlocksFromResponseContent(item["content"]).map((content) => ({ - sessionUpdate: "agent_message_chunk", - content, - })); + const phase = stringValue(item["phase"]); + return contentBlocksFromResponseContent(item["content"]).map((content) => ( + createAgentMessageChunk(content, undefined, createCodexMessagePhaseMeta(phase)) + )); } function createEventMsgUpdates(record: JsonRecord): UpdateSessionEvent[] | null { diff --git a/src/__tests__/CodexACPAgent/agent-message-events.test.ts b/src/__tests__/CodexACPAgent/agent-message-events.test.ts new file mode 100644 index 00000000..e4147d5f --- /dev/null +++ b/src/__tests__/CodexACPAgent/agent-message-events.test.ts @@ -0,0 +1,85 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import type { ServerNotification } from "../../app-server"; +import type { SessionState } from "../../CodexAcpServer"; +import { AgentMode } from "../../AgentMode"; +import { + createCodexMockTestFixture, + createTestSessionState, + setupPromptAndSendNotifications, + type CodexMockTestFixture +} from "../acp-test-utils"; + +describe("CodexEventHandler - agent message events", () => { + let mockFixture: CodexMockTestFixture; + const sessionId = "test-session-id"; + + beforeEach(() => { + mockFixture = createCodexMockTestFixture(); + vi.clearAllMocks(); + }); + + const sessionState: SessionState = createTestSessionState({ + sessionId, + currentModelId: "model-id[effort]", + agentMode: AgentMode.DEFAULT_AGENT_MODE + }); + + it("includes Codex message phase metadata on streamed agent message chunks", async () => { + const notifications: ServerNotification[] = [ + { + method: "item/started", + params: { + threadId: sessionId, + turnId: "turn-1", + startedAtMs: 0, + item: { + type: "agentMessage", + id: "commentary-message", + text: "", + phase: "commentary", + memoryCitation: null, + }, + }, + }, + { + method: "item/agentMessage/delta", + params: { + threadId: sessionId, + turnId: "turn-1", + itemId: "commentary-message", + delta: "Checking the relevant event mapping.", + }, + }, + { + method: "item/started", + params: { + threadId: sessionId, + turnId: "turn-1", + startedAtMs: 10, + item: { + type: "agentMessage", + id: "final-message", + text: "", + phase: "final_answer", + memoryCitation: null, + }, + }, + }, + { + method: "item/agentMessage/delta", + params: { + threadId: sessionId, + turnId: "turn-1", + itemId: "final-message", + delta: "Yes, here is the answer.", + }, + }, + ]; + + await setupPromptAndSendNotifications(mockFixture, sessionId, sessionState, notifications); + + await expect(mockFixture.getAcpConnectionDump([])).toMatchFileSnapshot( + "data/agent-message-phases.json" + ); + }); +}); diff --git a/src/__tests__/CodexACPAgent/data/agent-message-phases.json b/src/__tests__/CodexACPAgent/data/agent-message-phases.json new file mode 100644 index 00000000..8f1a4b1e --- /dev/null +++ b/src/__tests__/CodexACPAgent/data/agent-message-phases.json @@ -0,0 +1,42 @@ +{ + "method": "sessionUpdate", + "args": [ + { + "sessionId": "test-session-id", + "update": { + "sessionUpdate": "agent_message_chunk", + "messageId": "commentary-message", + "content": { + "type": "text", + "text": "Checking the relevant event mapping." + }, + "_meta": { + "codex": { + "phase": "commentary" + } + } + } + } + ] +} +{ + "method": "sessionUpdate", + "args": [ + { + "sessionId": "test-session-id", + "update": { + "sessionUpdate": "agent_message_chunk", + "messageId": "final-message", + "content": { + "type": "text", + "text": "Yes, here is the answer." + }, + "_meta": { + "codex": { + "phase": "final_answer" + } + } + } + } + ] +} \ No newline at end of file diff --git a/src/__tests__/CodexACPAgent/response-item-history-fallback.test.ts b/src/__tests__/CodexACPAgent/response-item-history-fallback.test.ts index 72301794..aaeef097 100644 --- a/src/__tests__/CodexACPAgent/response-item-history-fallback.test.ts +++ b/src/__tests__/CodexACPAgent/response-item-history-fallback.test.ts @@ -55,6 +55,26 @@ describe("ResponseItemHistoryFallback", () => { expect(thoughtTexts(updates)).toEqual(["Need to inspect the directory."]); }); + it("preserves assistant message phase metadata from response items", () => { + const updates = parseResponseItemHistoryFallback(jsonl([ + { + type: "response_item", + payload: { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "Final answer text." }], + phase: "final_answer", + }, + }, + functionCall("call-missing", "ls"), + functionCallOutput("call-missing", "Chunk ID: missing\nProcess exited with code 0\nOutput:\nREADME.md\n"), + ]), "terminal_output"); + + expect(agentMessageMetas(updates)).toEqual([ + { codex: { phase: "final_answer" } }, + ]); + }); + it("marks exec command outputs without exit footers failed when they report command errors", () => { const updates = parseResponseItemHistoryFallback(jsonl([ functionCall("call-read-failed", "cat missing.txt"), @@ -133,3 +153,11 @@ function thoughtTexts(updates: UpdateSessionEvent[] | null): string[] { )) .flatMap((update) => update.content.type === "text" ? [update.content.text] : []); } + +function agentMessageMetas(updates: UpdateSessionEvent[] | null): unknown[] { + return (updates ?? []) + .filter((update): update is Extract => ( + update.sessionUpdate === "agent_message_chunk" + )) + .map((update) => update._meta); +}