diff --git a/src/main/codex/codex-structured-journal-translation-turn-boundaries.ts b/src/main/codex/codex-structured-journal-translation-turn-boundaries.ts index d4881de1897..2c958858816 100644 --- a/src/main/codex/codex-structured-journal-translation-turn-boundaries.ts +++ b/src/main/codex/codex-structured-journal-translation-turn-boundaries.ts @@ -37,6 +37,14 @@ type TurnBoundaryEvent = { dispatchSequenceAtReceipt?: number } +type TurnTerminal = { + state: 'completed' | 'interrupted' + completedAt: number + /** Null when Codex named no verdict, or when the host inferred this end itself. */ + outcome?: AgentJournalTurnOutcome | null + durationMs?: number | null +} + /** Opens and settles the durable lifecycle row for each primary-thread turn. */ export class CodexJournalTurnBoundaries { private readonly recentTurns = new CodexJournalRecentTurns() @@ -146,46 +154,12 @@ export class CodexJournalTurnBoundaries { // turn boundary is no evidence contact was lost. Only `settleSession` may // write `unverifiable`. const status = readCodexTurnStatus(event.params) - const turnLifecycle = - event.threadId === this.deps.primaryThreadId() - ? this.settled(event.threadId, turnId, { - state: codexTurnLifecycleState(status), - outcome: codexTurnOutcome(status), - completedAt: this.receiptTime(event), - durationMs: readCodexTurnDurationMs(event.params) - }) - : null - const requestOrigin = this.deps.activeTurns.requestOrigin(event.threadId, turnId) - const latestDispatchSequence = this.deps.activeTurns.latestDispatchSequence( - event.threadId, - turnId - ) - const admission = settleCodexJournalTurn({ - sink: this.deps.sink, - sessionId: event.sessionId, - threadId: event.threadId, - turnId, - turnLifecycle, - streams: this.deps.items.streams, - activeItems: this.deps.items.activeItems, - pendingPrompts: this.deps.pendingPrompts, - ...(this.deps.clearPromptTurn ? { clearPromptTurn: this.deps.clearPromptTurn } : {}), - linkageFor: this.deps.linkageFor + return this.end(event, turnId, { + state: codexTurnLifecycleState(status), + outcome: codexTurnOutcome(status), + completedAt: this.receiptTime(event), + durationMs: readCodexTurnDurationMs(event.params) }) - if (admission.accepted) { - if (turnLifecycle) { - this.recentTurns.remember( - event.threadId, - turnLifecycle, - requestOrigin, - latestDispatchSequence - ) - } - this.deps.items.ordinals.forgetTurn(event.threadId, turnId) - this.deps.activeTurns.forget(event.threadId, turnId) - this.deps.resetActivity(event.threadId) - } - return admission } /** @@ -204,18 +178,34 @@ export class CodexJournalTurnBoundaries { return suppressionAdmission } const turnId = readCodexTurnId(event.params) ?? this.deps.activeTurns.current(event.threadId) - // An error naming an already-settled turn is not a second end: its terminal - // row holds the start and duration this one could not reconstruct. + // An error ends only a turn this host saw open; its `turn/completed` settles any other. if (!turnId || !this.deps.activeTurns.isActive(event.threadId, turnId)) { return CODEX_JOURNAL_ADMITTED } + return this.end(event, turnId, { + state: 'completed', + outcome: 'failure', + completedAt: this.receiptTime(event) + }) + } + + /** + * A turn's first terminal settlement is final. Codex follows a turn-ending + * `error` with a failed `turn/completed` for the same turn, and by then the + * start and attributed send this row carries are forgotten, so a second end + * could only overwrite the record with less. + */ + private end( + event: TurnBoundaryEvent, + turnId: string, + terminal: TurnTerminal + ): CodexJournalTranslationAdmission { + if (this.recentTurns.has(event.threadId, turnId)) { + return CODEX_JOURNAL_ADMITTED + } const turnLifecycle = event.threadId === this.deps.primaryThreadId() - ? this.settled(event.threadId, turnId, { - state: 'completed', - outcome: 'failure', - completedAt: this.receiptTime(event) - }) + ? this.settled(event.threadId, turnId, terminal) : null const requestOrigin = this.deps.activeTurns.requestOrigin(event.threadId, turnId) const latestDispatchSequence = this.deps.activeTurns.latestDispatchSequence( @@ -252,17 +242,7 @@ export class CodexJournalTurnBoundaries { /** Terminal lifecycle for a remembered turn; `startedAt` is absent when the start was never seen. * The verdict travels as one record so a caller cannot supply the state and drop the outcome. */ - settled( - threadId: string, - turnId: string, - terminal: { - state: 'completed' | 'interrupted' - completedAt: number - /** Null when Codex named no verdict, or when the host inferred this end itself. */ - outcome?: AgentJournalTurnOutcome | null - durationMs?: number | null - } - ): AgentJournalTurnLifecycle { + settled(threadId: string, turnId: string, terminal: TurnTerminal): AgentJournalTurnLifecycle { const startedAt = this.deps.activeTurns.startedAt(threadId, turnId) // Carried forward from the exact echoed send that was attributed to this turn. const requestOrigin = this.deps.activeTurns.requestOrigin(threadId, turnId) diff --git a/src/main/codex/codex-structured-journal-translation-turn-settles-once.test.ts b/src/main/codex/codex-structured-journal-translation-turn-settles-once.test.ts new file mode 100644 index 00000000000..a5c02bd3bfe --- /dev/null +++ b/src/main/codex/codex-structured-journal-translation-turn-settles-once.test.ts @@ -0,0 +1,210 @@ +import { describe, expect, it } from 'vitest' +import type { + AgentJournalItemBody, + AgentJournalItemIdentity +} from '../../shared/agent-session-journal-types' +import { + agentJournalItemKey, + agentJournalSubmissionKey +} from '../../shared/agent-session-journal-item-key' +import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' +import { createCodexJournalTranslator } from './codex-structured-journal-translation' +import type { CodexStructuredSessionEvent } from './codex-structured-session-adapter' + +const SESSION_ID = 'session-1' +const THREAD_ID = 'thread-abc' +const TURN_ID = 'turn-1' +const NEXT_TURN_ID = 'turn-2' +const CLIENT_MESSAGE_ID = 'client-1' + +type Row = { key: string; body: AgentJournalItemBody } + +function recorder() { + const rows: Row[] = [] + const sink: StructuredAgentSessionEventSink = { + appendItem: (identity: AgentJournalItemIdentity, body) => + rows.push({ key: agentJournalItemKey(identity), body }), + appendTombstone: () => {}, + publish: () => {} + } + return { sink, rows } +} + +function notification(method: string, params: unknown, observedAt: number) { + return { + type: 'notification', + sessionId: SESSION_ID, + threadId: THREAD_ID, + method, + params, + observedAt + } satisfies CodexStructuredSessionEvent +} + +function translator(tap: ReturnType) { + return createCodexJournalTranslator({ + sink: tap.sink, + sessionId: SESSION_ID, + primaryThreadId: () => THREAD_ID, + dispatchRequestOrigin: () => ({ requestedAt: 900, sequence: 0 }) + }) +} + +/** Every terminal write for a turn, in order: what a reader could observe. */ +function terminalWrites(rows: readonly Row[], turnId: string) { + return rows + .map((row) => row.body) + .filter((body) => body.kind === 'turn' && body.turnId === turnId && body.state !== 'running') +} + +/** The body the journal reducer keeps for a turn's lifecycle row. */ +function settledRecord(rows: readonly Row[], turnId: string) { + return rows + .map((row) => row.body) + .findLast((body) => body.kind === 'turn' && body.turnId === turnId) +} + +/** Codex's frames for a turn that fails: the error, then the failed completion. */ +function runFailedTurn(handle: (event: CodexStructuredSessionEvent) => unknown) { + handle(notification('turn/started', { turn: { id: TURN_ID } }, 1_000)) + handle( + notification( + 'item/started', + { + turnId: TURN_ID, + turn: { id: TURN_ID }, + item: { type: 'userMessage', id: 'user-1', clientId: CLIENT_MESSAGE_ID } + }, + 1_100 + ) + ) + handle( + notification( + 'item/completed', + { + turnId: TURN_ID, + item: { type: 'agentMessage', id: 'agent-1', text: 'Checking the build' } + }, + 1_500 + ) + ) + handle( + notification( + 'error', + { + threadId: THREAD_ID, + turnId: TURN_ID, + willRetry: false, + error: { message: 'stream disconnected before completion' } + }, + 2_000 + ) + ) + handle( + notification( + 'turn/completed', + { turn: { id: TURN_ID, status: 'failed', durationMs: 1_100 } }, + 2_100 + ) + ) +} + +describe('a Codex turn settles once', () => { + it('keeps the failure the error settled when the failed completion follows it', () => { + const tap = recorder() + const codex = translator(tap) + + runFailedTurn((event) => codex.handle(event)) + + expect(terminalWrites(tap.rows, TURN_ID)).toHaveLength(1) + expect(settledRecord(tap.rows, TURN_ID)).toEqual({ + kind: 'turn', + turnId: TURN_ID, + state: 'completed', + outcome: 'failure', + userItemId: agentJournalSubmissionKey(CLIENT_MESSAGE_ID), + startedAt: 1_000, + requestedAt: 900, + completedAt: 2_000 + }) + expect( + tap.rows.filter((row) => row.body.kind === 'status' && row.body.tone === 'error') + ).toHaveLength(1) + }) + + it('ignores a duplicate completion for a turn it already settled', () => { + const tap = recorder() + const codex = translator(tap) + + codex.handle(notification('turn/started', { turn: { id: TURN_ID } }, 1_000)) + codex.handle( + notification( + 'turn/completed', + { turn: { id: TURN_ID, status: 'completed', durationMs: 900 } }, + 2_000 + ) + ) + codex.handle(notification('turn/completed', { turn: { id: TURN_ID, status: 'failed' } }, 3_000)) + + expect(terminalWrites(tap.rows, TURN_ID)).toHaveLength(1) + expect(settledRecord(tap.rows, TURN_ID)).toMatchObject({ + state: 'completed', + outcome: 'success', + startedAt: 1_000, + completedAt: 2_000, + durationMs: 900 + }) + }) + + it('settles an ordinary turn exactly as before', () => { + const tap = recorder() + const codex = translator(tap) + + codex.handle(notification('turn/started', { turn: { id: TURN_ID } }, 1_000)) + codex.handle( + notification( + 'turn/completed', + { turn: { id: TURN_ID, status: 'completed', durationMs: 3_250 } }, + 4_500 + ) + ) + + expect(terminalWrites(tap.rows, TURN_ID)).toEqual([ + { + kind: 'turn', + turnId: TURN_ID, + state: 'completed', + outcome: 'success', + userItemId: `codex:${THREAD_ID}:${TURN_ID}:0`, + startedAt: 1_000, + completedAt: 4_500, + durationMs: 3_250 + } + ]) + }) + + it('settles the next turn on its own after a failed one', () => { + const tap = recorder() + const codex = translator(tap) + + runFailedTurn((event) => codex.handle(event)) + const failed = settledRecord(tap.rows, TURN_ID) + codex.handle(notification('turn/started', { turn: { id: NEXT_TURN_ID } }, 3_000)) + codex.handle( + notification( + 'turn/completed', + { turn: { id: NEXT_TURN_ID, status: 'completed', durationMs: 1_000 } }, + 4_000 + ) + ) + + expect(settledRecord(tap.rows, TURN_ID)).toEqual(failed) + expect(terminalWrites(tap.rows, NEXT_TURN_ID)).toHaveLength(1) + expect(settledRecord(tap.rows, NEXT_TURN_ID)).toMatchObject({ + state: 'completed', + outcome: 'success', + startedAt: 3_000, + completedAt: 4_000 + }) + }) +}) diff --git a/src/main/codex/codex-structured-journal-translation-turn-state.ts b/src/main/codex/codex-structured-journal-translation-turn-state.ts index ec1cff041db..832f57f3096 100644 --- a/src/main/codex/codex-structured-journal-translation-turn-state.ts +++ b/src/main/codex/codex-structured-journal-translation-turn-state.ts @@ -167,7 +167,8 @@ type RecentTurn = { bytes: number } -/** Bounded terminal lifecycle window for exact echoes that arrive after completion. */ +/** Bounded terminal lifecycle window: exact echoes that arrive after completion + * revise it, and a later end for a turn in it is not a second settlement. */ export class CodexJournalRecentTurns { private readonly turns = new Map() private retainedBytes = 0 @@ -231,6 +232,10 @@ export class CodexJournalRecentTurns { } } + has(threadId: string, turnId: string): boolean { + return this.turns.has(this.turnKey(threadId, turnId)) + } + requestOriginRevision( threadId: string, turnId: string, diff --git a/src/main/codex/codex-structured-journal-turn-settles-once-journal.test.ts b/src/main/codex/codex-structured-journal-turn-settles-once-journal.test.ts new file mode 100644 index 00000000000..6459b59755c --- /dev/null +++ b/src/main/codex/codex-structured-journal-turn-settles-once-journal.test.ts @@ -0,0 +1,111 @@ +// A failed Codex turn through the real path its record takes: +// translator → deferred sink queue → on-disk journal → the shared turn-timing reader. + +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it } from 'vitest' +import { readAgentJournalTurn } from '../../shared/agent-session-turn-record' +import { + completedStructuredAgentTurnSeconds, + selectStructuredAgentRunningTurnTiming +} from '../../shared/structured-agent-session-turn-timing' +import { openAgentSessionJournal } from '../native-chat/agent-session-journal/journal-store-factory' +import { createDeferredStructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' +import { createCodexJournalTranslator } from './codex-structured-journal-translation' + +const SESSION = 'session-codex-failed-turn' +const THREAD = 'thread-abc' +const TURN = 'turn-1' + +const cleanups: (() => Promise)[] = [] +afterEach(async () => { + for (const cleanup of cleanups.splice(0)) { + await cleanup() + } +}) + +async function session() { + const root = await mkdtemp(join(tmpdir(), 'orca-codex-failed-turn-')) + const journal = await openAgentSessionJournal({ + identity: { + sessionId: SESSION, + workspaceId: 'workspace-1', + hostId: 'local', + agent: 'codex', + providerHandle: { kind: 'codex', threadId: THREAD } + }, + journalDir: root, + now: () => 1_000 + }) + const deferred = createDeferredStructuredAgentSessionEventSink() + deferred.bind({ journal, fence: 1, publish: () => {} }) + cleanups.push(async () => { + deferred.close() + await journal.close() + await rm(root, { recursive: true, force: true }) + }) + const translator = createCodexJournalTranslator({ + sink: deferred.sink, + sessionId: SESSION, + primaryThreadId: () => THREAD, + schedule: (run) => { + run() + return () => {} + } + }) + const on = (method: string, params: Record, observedAt: number) => + translator.handle({ + type: 'notification', + sessionId: SESSION, + threadId: THREAD, + method, + params: { threadId: THREAD, ...params }, + observedAt + }) + return { + on, + drained: () => deferred.drained(), + items: async () => { + await deferred.drained() + return journal.snapshot().items + } + } +} + +describe('a failed Codex turn in the journal', () => { + it('keeps the failure and its duration when the failed completion lands while the error is still queued', async () => { + const { on, drained, items } = await session() + on('turn/started', { turn: { id: TURN } }, 1_000) + on( + 'item/completed', + { turnId: TURN, item: { type: 'agentMessage', id: 'agent-1', text: 'Checking the build' } }, + 1_500 + ) + await drained() + + // Codex writes both frames back to back; the error's status row is still being + // written when the failed completion arrives, so the error's settlement is queued. + on( + 'error', + { turnId: TURN, willRetry: false, error: { message: 'stream disconnected' } }, + 3_000 + ) + on('turn/completed', { turn: { id: TURN, status: 'failed', durationMs: 1_900 } }, 3_100) + + const rows = await items() + const turn = rows + .map((item) => readAgentJournalTurn(item.body)) + .findLast((record) => record?.turnId === TURN) + expect(turn).toMatchObject({ + state: 'completed', + outcome: 'failure', + startedAt: 1_000, + completedAt: 3_000 + }) + // The duration "Worked for" shows under the turn's message. + expect( + completedStructuredAgentTurnSeconds(selectStructuredAgentRunningTurnTiming(rows, TURN)) + ).toBe(2) + }) +})