From 6c9cd3263c8a38e995e2a997b9d9daa3638e25be Mon Sep 17 00:00:00 2001 From: Leeeon233 Date: Thu, 30 Jul 2026 03:52:49 +0000 Subject: [PATCH] feat: fork sessions by native turn id Model: gpt-5.6-sol --- AGENTS.md | 7 ++- src/AcpExtensions.ts | 17 +++--- src/CodexAcpClient.ts | 45 +--------------- src/CodexAcpServer.ts | 32 ++++------- src/CodexEventHandler.ts | 7 ++- src/ContentChunks.ts | 10 ++++ .../data/agent-message-phases.json | 6 +++ .../data/follow-up-no-duplicates.json | 15 ++++++ .../data/load-session-history.json | 5 ++ ...ession-response-item-history-fallback.json | 5 ++ .../CodexACPAgent/data/multiple-sessions.json | 10 ++++ .../CodexACPAgent/data/output-acp-events.json | 15 ++++++ .../CodexACPAgent/initialize.test.ts | 3 +- .../CodexACPAgent/session-fork.test.ts | 53 +------------------ .../CodexACPAgent/thread-goal-events.test.ts | 5 ++ src/index.ts | 8 --- 16 files changed, 101 insertions(+), 142 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index dc075fdb..1f0eb082 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -35,10 +35,9 @@ - App-server events: prefer `thread/*`, `turn/*`, and `item/*` event surfaces; avoid the deprecated `codex/event/*` API (planned removal). Keep implementations aligned with generated types in `src/app-server` (including `v2` exports). - Steer uses app-server `turn/steer` on the tracked active turn. Correlate `clientUserMessageId` and acknowledge only the matching `item/completed(userMessage)`; never emulate steer with a second `turn/start`. - Session fork uses app-server `thread/fork` and installs the returned child as an independent ACP - session. The temporary `_meta.lody.forkAtMessage` extension carries a standard ACP `messageId`; - resolve the containing Codex turn with `thread/read` inside this adapter before setting - `thread/fork.lastTurnId`. Never expose Codex turn IDs as ACP message IDs or emulate fork by - replaying source history. + session. Agent message updates expose their Codex turn id as `_meta.lody.turnId`; + `_meta.lody.forkAtTurn.turnId` is passed directly to `thread/fork.lastTurnId`. Do not maintain + a message-id mapping or emulate fork by replaying source history. - Codex reasoning summaries can echo trailing empty HTML comments from model instructions. Keep that provider-specific cleanup in `src/ReasoningText.ts` across live deltas and history replay; do not filter assistant text, raw reasoning, or HTML globally in the client renderer. diff --git a/src/AcpExtensions.ts b/src/AcpExtensions.ts index 5ba7dfd6..ddcdd780 100644 --- a/src/AcpExtensions.ts +++ b/src/AcpExtensions.ts @@ -12,19 +12,16 @@ export const ACP_EXT_SESSION_USAGE_UPDATE_METHOD = "_acp_ext:session_usage_updat export const ACP_EXT_SESSION_RATE_LIMITS_METHOD = "_acp_ext:session_rate_limits"; export const ACP_EXT_CODEX_PROPOSED_PLAN_METHOD = "_acp_ext:codex_proposed_plan"; export const CODEX_STEER_APPLIED_METHOD = "_codex/steerApplied"; -export const LODY_FORK_MESSAGE_BEFORE_ACTIVE_TURN_METHOD = - "_lody/session/fork-message-before-active-turn"; - -export function getLodyForkMessageId(meta: unknown): string | null { +export function getLodyForkTurnId(meta: unknown): string | null { if (typeof meta !== "object" || meta === null) return null; const lody = (meta as Record)["lody"]; if (typeof lody !== "object" || lody === null) return null; - const forkAtMessage = (lody as Record)["forkAtMessage"]; - if (typeof forkAtMessage !== "object" || forkAtMessage === null) return null; - const version = (forkAtMessage as Record)["version"]; - const messageId = (forkAtMessage as Record)["messageId"]; - return version === 1 && typeof messageId === "string" && messageId.length > 0 - ? messageId + const forkAtTurn = (lody as Record)["forkAtTurn"]; + if (typeof forkAtTurn !== "object" || forkAtTurn === null) return null; + const version = (forkAtTurn as Record)["version"]; + const turnId = (forkAtTurn as Record)["turnId"]; + return version === 1 && typeof turnId === "string" && turnId.length > 0 + ? turnId : null; } diff --git a/src/CodexAcpClient.ts b/src/CodexAcpClient.ts index 80df60bd..ccd8a187 100644 --- a/src/CodexAcpClient.ts +++ b/src/CodexAcpClient.ts @@ -23,7 +23,7 @@ import {AgentMode} from "./AgentMode"; import path from "node:path"; import {logger} from "./Logger"; import {sanitizeMcpServerName} from "./McpServerName"; -import {getLodyForkMessageId} from "./AcpExtensions"; +import {getLodyForkTurnId} from "./AcpExtensions"; import type { AccountLoginCompletedNotification, AccountUpdatedNotification, @@ -356,10 +356,7 @@ export class CodexAcpClient { ): Promise { const additionalDirectories = readAdditionalDirectories(request.cwd, request.additionalDirectories, request._meta); await this.refreshSkills(request.cwd, additionalDirectories); - const forkMessageId = getLodyForkMessageId(request._meta); - const forkTurnId = forkMessageId - ? await this.resolveForkTurnIdForMessage(request.sessionId, forkMessageId) - : null; + const forkTurnId = getLodyForkTurnId(request._meta); const response = await this.codexClient.threadFork({ config: await this.createSessionConfig(request.cwd, additionalDirectories, request.mcpServers ?? []), @@ -382,44 +379,6 @@ export class CodexAcpClient { }; } - async resolveForkTurnIdForMessage(threadId: string, messageId: string): Promise { - const response = await this.codexClient.threadRead({ - threadId, - includeTurns: true, - }); - const turn = response.thread.turns.find((candidate) => - candidate.items.some((item) => item.type === "agentMessage" && item.id === messageId) - ); - if (!turn) { - throw RequestError.invalidRequest("ACP message is not a forkable Codex turn boundary"); - } - return turn.id; - } - - async findMessageBeforeTurn(threadId: string, activeTurnId: string): Promise { - const response = await this.codexClient.threadRead({ - threadId, - includeTurns: true, - }); - const turns = response.thread.turns; - const activeIndex = turns.findIndex((turn) => turn.id === activeTurnId); - if (activeIndex < 0) { - throw RequestError.invalidRequest("Active turn changed before its preceding message was captured"); - } - for (let index = activeIndex - 1; index >= 0; index--) { - const turn = turns[index]; - if (turn && turn.status !== "inProgress") { - const agentMessage = [...turn.items] - .reverse() - .find((item) => item.type === "agentMessage"); - if (agentMessage) { - return agentMessage.id; - } - } - } - throw RequestError.invalidRequest("No completed assistant message exists before the active turn"); - } - async loadSession(request: acp.LoadSessionRequest, onSubscribed?: () => void): Promise { const additionalDirectories = readAdditionalDirectories(request.cwd, request.additionalDirectories, request._meta); await this.refreshSkills(request.cwd, additionalDirectories); diff --git a/src/CodexAcpServer.ts b/src/CodexAcpServer.ts index fc18f307..a937e671 100644 --- a/src/CodexAcpServer.ts +++ b/src/CodexAcpServer.ts @@ -47,7 +47,6 @@ import { getCodexSteerId, isExtMethodRequest, LEGACY_SET_SESSION_MODEL_METHOD, - LODY_FORK_MESSAGE_BEFORE_ACTIVE_TURN_METHOD, } from "./AcpExtensions"; import { createCollabAgentToolCallUpdate, @@ -81,7 +80,7 @@ import {isJetBrains2026_1Client} from "./JBUtils"; import {resolveTerminalOutputMode, type TerminalOutputMode} from "./TerminalOutputMode"; import {sanitizeReasoningParts} from "./ReasoningText"; import { - createCodexMessagePhaseMeta, + createCodexAgentMessageMeta, createAgentTextMessageChunk, createAgentTextThoughtChunk, createUserMessageChunk, @@ -268,7 +267,7 @@ export class CodexAcpServer { steer: CODEX_STEER_CAPABILITY, }, lody: { - forkAtMessage: {version: 1, beforeActiveTurn: true}, + forkAtTurn: {version: 1}, }, }, }, @@ -669,22 +668,6 @@ export class CodexAcpServer { }; } - async resolveMessageBeforeActiveTurn(params: {sessionId: string}): Promise<{messageId: string}> { - const session = this.sessions.get(params.sessionId); - const activeTurnId = session?.currentTurnId; - if (!activeTurnId) { - throw RequestError.invalidRequest("Session has no active turn"); - } - const messageId = await this.runWithProcessCheck(() => - this.codexAcpClient.findMessageBeforeTurn(params.sessionId, activeTurnId) - ); - logger.log("Resolved ACP message before active turn", { - sessionId: params.sessionId, - method: LODY_FORK_MESSAGE_BEFORE_ACTIVE_TURN_METHOD, - }); - return {messageId}; - } - async listSessions(params: acp.ListSessionsRequest): Promise { logger.log("Listing sessions...", {cwd: params.cwd, cursor: params.cursor}); await this.checkAuthorization(); @@ -1205,7 +1188,7 @@ export class CodexAcpServer { const threadUpdates: UpdateSessionEvent[] = []; for (const turn of thread.turns) { for (const item of turn.items) { - const updates = await this.createHistoryUpdates(item, sessionState); + const updates = await this.createHistoryUpdates(item, sessionState, turn.id); threadUpdates.push(...updates); } } @@ -1279,7 +1262,11 @@ export class CodexAcpServer { return normalized.length > 0 ? normalized : null; } - private async createHistoryUpdates(item: ThreadItem, sessionState: SessionState): Promise { + private async createHistoryUpdates( + item: ThreadItem, + sessionState: SessionState, + turnId: string, + ): Promise { switch (item.type) { case "userMessage": return this.createUserMessageUpdates(item); @@ -1288,12 +1275,11 @@ export class CodexAcpServer { case "sleep": return []; case "agentMessage": { - const meta = createCodexMessagePhaseMeta(item.phase); return [{ sessionUpdate: "agent_message_chunk", messageId: item.id, content: { type: "text", text: item.text }, - ...(meta ? { _meta: meta } : {}), + _meta: createCodexAgentMessageMeta(item.phase, turnId), }]; } case "reasoning": diff --git a/src/CodexEventHandler.ts b/src/CodexEventHandler.ts index 381da1f5..f08f6083 100644 --- a/src/CodexEventHandler.ts +++ b/src/CodexEventHandler.ts @@ -69,6 +69,7 @@ import { stripShellPrefix } from "./CommandUtils"; import {createTerminalOutputMeta, type TerminalOutputMode} from "./TerminalOutputMode"; import {ReasoningSummaryFilter, sanitizeReasoningParts} from "./ReasoningText"; import { + createCodexAgentMessageMeta, createCodexMessagePhaseMeta, createAgentTextMessageChunk, createAgentTextThoughtChunk, @@ -364,7 +365,11 @@ export class CodexEventHandler { private async createTextEvent(event: AgentMessageDeltaNotification): Promise { const phase = this.agentMessagePhases.get(event.itemId) ?? null; - return createAgentTextMessageChunk(event.delta, event.itemId, createCodexMessagePhaseMeta(phase)); + return createAgentTextMessageChunk( + event.delta, + event.itemId, + createCodexAgentMessageMeta(phase, event.turnId), + ); } private async createConfigWarningEvent(event: ConfigWarningNotification): Promise { diff --git a/src/ContentChunks.ts b/src/ContentChunks.ts index 2bef82e8..232efacf 100644 --- a/src/ContentChunks.ts +++ b/src/ContentChunks.ts @@ -10,6 +10,16 @@ export function createCodexMessagePhaseMeta(phase: string | null | undefined): A return { codex: { phase } }; } +export function createCodexAgentMessageMeta( + phase: string | null | undefined, + turnId: string, +): AcpMeta { + return { + ...(phase ? {codex: {phase}} : {}), + lody: {turnId}, + }; +} + export function createUserMessageChunk(content: ContentBlock, messageId?: string, meta?: AcpMeta): UpdateSessionEvent { if (messageId) { return { diff --git a/src/__tests__/CodexACPAgent/data/agent-message-phases.json b/src/__tests__/CodexACPAgent/data/agent-message-phases.json index 8f1a4b1e..a7c3a0d3 100644 --- a/src/__tests__/CodexACPAgent/data/agent-message-phases.json +++ b/src/__tests__/CodexACPAgent/data/agent-message-phases.json @@ -13,6 +13,9 @@ "_meta": { "codex": { "phase": "commentary" + }, + "lody": { + "turnId": "turn-1" } } } @@ -34,6 +37,9 @@ "_meta": { "codex": { "phase": "final_answer" + }, + "lody": { + "turnId": "turn-1" } } } diff --git a/src/__tests__/CodexACPAgent/data/follow-up-no-duplicates.json b/src/__tests__/CodexACPAgent/data/follow-up-no-duplicates.json index d992b603..1e7b2b19 100644 --- a/src/__tests__/CodexACPAgent/data/follow-up-no-duplicates.json +++ b/src/__tests__/CodexACPAgent/data/follow-up-no-duplicates.json @@ -9,6 +9,11 @@ "content": { "type": "text", "text": "He" + }, + "_meta": { + "lody": { + "turnId": "string" + } } } } @@ -25,6 +30,11 @@ "content": { "type": "text", "text": "ll" + }, + "_meta": { + "lody": { + "turnId": "string" + } } } } @@ -41,6 +51,11 @@ "content": { "type": "text", "text": "o!" + }, + "_meta": { + "lody": { + "turnId": "string" + } } } } diff --git a/src/__tests__/CodexACPAgent/data/load-session-history.json b/src/__tests__/CodexACPAgent/data/load-session-history.json index 43236816..3f7ceb1d 100644 --- a/src/__tests__/CodexACPAgent/data/load-session-history.json +++ b/src/__tests__/CodexACPAgent/data/load-session-history.json @@ -53,6 +53,11 @@ "content": { "type": "text", "text": "Hello!" + }, + "_meta": { + "lody": { + "turnId": "turn-1" + } } } } diff --git a/src/__tests__/CodexACPAgent/data/load-session-response-item-history-fallback.json b/src/__tests__/CodexACPAgent/data/load-session-response-item-history-fallback.json index 64671c33..2598bffe 100644 --- a/src/__tests__/CodexACPAgent/data/load-session-response-item-history-fallback.json +++ b/src/__tests__/CodexACPAgent/data/load-session-response-item-history-fallback.json @@ -197,6 +197,11 @@ "content": { "type": "text", "text": "The directory contains README.md and src." + }, + "_meta": { + "lody": { + "turnId": "turn-1" + } } } } diff --git a/src/__tests__/CodexACPAgent/data/multiple-sessions.json b/src/__tests__/CodexACPAgent/data/multiple-sessions.json index 4c884302..123d6254 100644 --- a/src/__tests__/CodexACPAgent/data/multiple-sessions.json +++ b/src/__tests__/CodexACPAgent/data/multiple-sessions.json @@ -9,6 +9,11 @@ "content": { "type": "text", "text": "Hello-1" + }, + "_meta": { + "lody": { + "turnId": "string" + } } } } @@ -25,6 +30,11 @@ "content": { "type": "text", "text": "Hello-2" + }, + "_meta": { + "lody": { + "turnId": "string" + } } } } diff --git a/src/__tests__/CodexACPAgent/data/output-acp-events.json b/src/__tests__/CodexACPAgent/data/output-acp-events.json index d992b603..1e7b2b19 100644 --- a/src/__tests__/CodexACPAgent/data/output-acp-events.json +++ b/src/__tests__/CodexACPAgent/data/output-acp-events.json @@ -9,6 +9,11 @@ "content": { "type": "text", "text": "He" + }, + "_meta": { + "lody": { + "turnId": "string" + } } } } @@ -25,6 +30,11 @@ "content": { "type": "text", "text": "ll" + }, + "_meta": { + "lody": { + "turnId": "string" + } } } } @@ -41,6 +51,11 @@ "content": { "type": "text", "text": "o!" + }, + "_meta": { + "lody": { + "turnId": "string" + } } } } diff --git a/src/__tests__/CodexACPAgent/initialize.test.ts b/src/__tests__/CodexACPAgent/initialize.test.ts index 1c15c11c..fa77aa68 100644 --- a/src/__tests__/CodexACPAgent/initialize.test.ts +++ b/src/__tests__/CodexACPAgent/initialize.test.ts @@ -70,9 +70,8 @@ describe('CodexACPAgent - initialize', () => { }, }, lody: { - forkAtMessage: { + forkAtTurn: { version: 1, - beforeActiveTurn: true, }, }, }, diff --git a/src/__tests__/CodexACPAgent/session-fork.test.ts b/src/__tests__/CodexACPAgent/session-fork.test.ts index a424fd92..e8b4089d 100644 --- a/src/__tests__/CodexACPAgent/session-fork.test.ts +++ b/src/__tests__/CodexACPAgent/session-fork.test.ts @@ -21,14 +21,6 @@ describe("ACP session fork", () => { vi.spyOn(codexAppServerClient, "skillsExtraRootsSet").mockResolvedValue(undefined); vi.spyOn(codexAppServerClient, "listSkills").mockResolvedValue({data: []}); vi.spyOn(codexAppServerClient, "configRead").mockResolvedValue({config: {}} as never); - const threadReadSpy = vi.spyOn(codexAppServerClient, "threadRead").mockResolvedValue({ - thread: { - turns: [{ - id: "completed-turn-id", - items: [{type: "agentMessage", id: "assistant-message-id"}], - }], - }, - } as never); const threadForkSpy = vi.spyOn(codexAppServerClient, "threadFork").mockResolvedValue({ thread: {id: "child-session-id"}, model: model.id, @@ -49,9 +41,9 @@ describe("ACP session fork", () => { mcpServers: [mcpServer], _meta: { lody: { - forkAtMessage: { + forkAtTurn: { version: 1, - messageId: "assistant-message-id", + turnId: "completed-turn-id", }, }, }, @@ -89,47 +81,6 @@ describe("ACP session fork", () => { }, }, }); - expect(threadReadSpy).toHaveBeenCalledWith({ - threadId: "source-session-id", - includeTurns: true, - }); - }); - - it("resolves the last terminal Codex turn before the active turn", async () => { - const fixture = createCodexMockTestFixture(); - const codexAcpClient = fixture.getCodexAcpClient(); - const codexAppServerClient = fixture.getCodexAppServerClient(); - vi.spyOn(codexAppServerClient, "threadRead").mockResolvedValue({ - thread: { - turns: [ - { - id: "completed-turn", - status: "completed", - items: [{type: "agentMessage", id: "assistant-message-id"}], - }, - {id: "active-turn", status: "inProgress", items: []}, - ], - }, - } as never); - - await expect( - codexAcpClient.findMessageBeforeTurn("source-session-id", "active-turn"), - ).resolves.toBe("assistant-message-id"); - }); - - it("rejects when the active Codex turn changed during capture", async () => { - const fixture = createCodexMockTestFixture(); - const codexAcpClient = fixture.getCodexAcpClient(); - const codexAppServerClient = fixture.getCodexAppServerClient(); - vi.spyOn(codexAppServerClient, "threadRead").mockResolvedValue({ - thread: { - turns: [{id: "completed-turn", status: "completed"}], - }, - } as never); - - await expect( - codexAcpClient.findMessageBeforeTurn("source-session-id", "stale-active-turn"), - ).rejects.toThrow("Invalid request"); }); it("installs the fork as an independent promptable ACP session", async () => { diff --git a/src/__tests__/CodexACPAgent/thread-goal-events.test.ts b/src/__tests__/CodexACPAgent/thread-goal-events.test.ts index fc690994..74fb6b97 100644 --- a/src/__tests__/CodexACPAgent/thread-goal-events.test.ts +++ b/src/__tests__/CodexACPAgent/thread-goal-events.test.ts @@ -178,6 +178,11 @@ describe("CodexEventHandler - thread goal events", () => { type: "text", text: "Because they kept losing interest in `any`.", }, + _meta: { + lody: { + turnId: "turn-1", + }, + }, }); expect(events[1]!.args[0].update).toEqual({ sessionUpdate: "session_info_update", diff --git a/src/index.ts b/src/index.ts index cc5bc28c..b886c6b7 100644 --- a/src/index.ts +++ b/src/index.ts @@ -14,7 +14,6 @@ import {runLoginCommand} from "./login"; import {runCodexCli} from "./CodexCli"; import { LEGACY_SET_SESSION_MODEL_METHOD, - LODY_FORK_MESSAGE_BEFORE_ACTIVE_TURN_METHOD, } from "./AcpExtensions"; const emptyExtensionParamsParser = z.preprocess( @@ -27,10 +26,6 @@ const legacySetSessionModelParamsParser = z.object({ modelId: z.string(), }).passthrough(); -const activeTurnForkMessageParamsParser = z.object({ - sessionId: z.string().min(1), -}); - if (process.argv.includes("--version")) { console.log(`${packageJson.name} ${packageJson.version}`); process.exit(0); @@ -140,8 +135,5 @@ function startAcpServer() { .onRequest("authentication/status", emptyExtensionParamsParser, (ctx) => getAgent().extMethod("authentication/status", ctx.params)) .onRequest("authentication/logout", emptyExtensionParamsParser, (ctx) => getAgent().extMethod("authentication/logout", ctx.params)) .onRequest(LEGACY_SET_SESSION_MODEL_METHOD, legacySetSessionModelParamsParser, (ctx) => getAgent().extMethod(LEGACY_SET_SESSION_MODEL_METHOD, ctx.params)) - .onRequest(LODY_FORK_MESSAGE_BEFORE_ACTIVE_TURN_METHOD, activeTurnForkMessageParamsParser, (ctx) => - getAgent().resolveMessageBeforeActiveTurn(ctx.params) - ) .connect(acpJsonStream); }