mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 00:02:10 +00:00
fix(native-chat): a failed Codex turn keeps its failure and "Worked for" (#23514)
* fix(native-chat): a Codex turn's first terminal settlement is final Codex follows a turn-ending `error` with a failed `turn/completed` for the same turn. The error settled the turn and dropped its start and attributed send; the completion then re-settled it from nothing, so the record lost its start time and flipped from completed/failure to interrupted/failure. A failed turn lost its "Worked for" and read as interrupted. Both ends now go through one settlement that refuses a turn already in the recent-turns window, so every later end (error then completion, a duplicate completion) adds nothing and overwrites nothing. * test(native-chat): the settles-once fixture is an ordinary failed turn The fixture was labelled and worded as a failed compaction; this change covers the ordinary-turn path, so its frames now read as one. * test(native-chat): a failed Codex turn keeps its record when its two ends coalesce in the queue
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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<typeof recorder>) {
|
||||
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
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -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<string, RecentTurn>()
|
||||
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,
|
||||
|
||||
@@ -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<void>)[] = []
|
||||
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<string, unknown>, 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)
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user