diff --git a/src/main/codex/codex-structured-dispatch-echo.test.ts b/src/main/codex/codex-structured-dispatch-echo.test.ts index c7eb43bf96f..1eeb8186ad9 100644 --- a/src/main/codex/codex-structured-dispatch-echo.test.ts +++ b/src/main/codex/codex-structured-dispatch-echo.test.ts @@ -133,9 +133,9 @@ describe('the turn Codex answered a send into but has not opened', () => { it('is the turn the latest armed send was answered into', () => { const echoes = createCodexDispatchEchoes() echoes.arm('client-1') - echoes.bindTurn('client-1', 'thread-1', 'turn-1') + echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start') echoes.arm('client-2') - echoes.bindTurn('client-2', 'thread-1', 'turn-2') + echoes.bindTurn('client-2', 'thread-1', 'turn-2', 'start') expect(echoes.answeredUnopenedTurn('thread-1', NONE_OPEN)).toBe('turn-2') }) @@ -143,7 +143,7 @@ describe('the turn Codex answered a send into but has not opened', () => { it('is none once Codex opened that turn', () => { const echoes = createCodexDispatchEchoes() echoes.arm('client-1') - echoes.bindTurn('client-1', 'thread-1', 'turn-1') + echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start') expect(echoes.answeredUnopenedTurn('thread-1', new Set(['turn-1']))).toBeNull() }) @@ -151,7 +151,7 @@ describe('the turn Codex answered a send into but has not opened', () => { it('is none once that turn ended, even with its send still armed for an echo', () => { const echoes = createCodexDispatchEchoes() echoes.arm('client-1') - echoes.bindTurn('client-1', 'thread-1', 'turn-1') + echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start') // A completed end leaves its unechoed send armed. expect(echoes.endTurn('thread-1', 'turn-1', { status: 'completed' })).toEqual([]) @@ -161,9 +161,9 @@ describe('the turn Codex answered a send into but has not opened', () => { it('skips a turn a wait left unopened, and still names an earlier one', () => { const echoes = createCodexDispatchEchoes() echoes.arm('client-1') - echoes.bindTurn('client-1', 'thread-1', 'turn-1') + echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start') echoes.arm('client-2') - echoes.bindTurn('client-2', 'thread-1', 'turn-2') + echoes.bindTurn('client-2', 'thread-1', 'turn-2', 'start') echoes.leftUnopened('thread-1', 'turn-2') @@ -174,7 +174,7 @@ describe('the turn Codex answered a send into but has not opened', () => { const echoes = createCodexDispatchEchoes() echoes.arm('client-1') echoes.arm('client-2') - echoes.bindTurn('client-2', 'thread-2', 'turn-2') + echoes.bindTurn('client-2', 'thread-2', 'turn-2', 'start') expect(echoes.answeredUnopenedTurn('thread-1', NONE_OPEN)).toBeNull() }) diff --git a/src/main/codex/codex-structured-dispatch-echo.ts b/src/main/codex/codex-structured-dispatch-echo.ts index 775cac5bca4..e72cb3c63a9 100644 --- a/src/main/codex/codex-structured-dispatch-echo.ts +++ b/src/main/codex/codex-structured-dispatch-echo.ts @@ -1,5 +1,8 @@ import type { ProviderDiagnostic } from '../../shared/agent-session-failure' -import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types' +import type { + AgentJournalItemIdentity, + AgentJournalTurnJoin +} from '../../shared/agent-session-journal-types' /** Sends awaiting their echo. One bound to a turn that ended without taking it settles from that * end; any other whose echo never arrives is retired by the journal's recovery on exit. */ @@ -35,10 +38,15 @@ export type CodexDispatchEchoes = { /** Drops an armed send whose write never reached the provider. */ disarm: (clientMessageId: string) => void /** - * Binds a send to the turn Codex answered it into. Returns that turn's end when the answer is - * read after it; a send that end settles is no longer armed. + * Binds a send to the turn Codex answered it into, and how it joined that turn. Returns that + * turn's end when the answer is read after it; a send that end settles is no longer armed. */ - bindTurn: (clientMessageId: string, threadId: string, turnId: string) => CodexTurnEnd | null + bindTurn: ( + clientMessageId: string, + threadId: string, + turnId: string, + via: AgentJournalTurnJoin + ) => CodexTurnEnd | null /** The turn the latest armed send was answered into that is neither in `openTurnIds`, ended, nor * left unopened through a wait: one Codex has picked for the send but not opened. */ answeredUnopenedTurn: (threadId: string, openTurnIds: ReadonlySet) => string | null @@ -49,7 +57,11 @@ export type CodexDispatchEchoes = { * completed, which echoes its pending input first, so one it never echoed waits for recovery. An * interrupt withdraws an un-echoed send, steered or the turn's own input: neither reached history. */ - endTurn: (threadId: string, turnId: string, end: CodexTurnEnd) => string[] + endTurn: ( + threadId: string, + turnId: string, + end: CodexTurnEnd + ) => { clientMessageId: string; via: AgentJournalTurnJoin }[] /** Submission origin for this exact send, retained until its echo settles it. */ requestOrigin: (clientMessageId: string) => CodexDispatchRequestOrigin | null /** Highest causal sequence assigned to a dispatch in this session. */ @@ -61,7 +73,11 @@ export type CodexDispatchEchoes = { export function createCodexDispatchEchoes(): CodexDispatchEchoes { const armed = new Map< string, - { requestedAt: number | null; sequence: number; turn?: { threadId: string; turnId: string } } + { + requestedAt: number | null + sequence: number + turn?: { threadId: string; turnId: string; via: AgentJournalTurnJoin } + } >() const endedTurns = new Map() const unopenedTurns = new Set() @@ -85,12 +101,12 @@ export function createCodexDispatchEchoes(): CodexDispatchEchoes { }, settle: (clientMessageId) => armed.delete(clientMessageId), disarm: (clientMessageId) => void armed.delete(clientMessageId), - bindTurn: (clientMessageId, threadId, turnId) => { + bindTurn: (clientMessageId, threadId, turnId, via) => { const entry = armed.get(clientMessageId) if (!entry) { return null } - entry.turn = { threadId, turnId } + entry.turn = { threadId, turnId, via } const end = endedTurns.get(turnKey(threadId, turnId)) ?? null if (end && settles(end)) { armed.delete(clientMessageId) @@ -132,10 +148,10 @@ export function createCodexDispatchEchoes(): CodexDispatchEchoes { } const settled = [...armed].flatMap(([clientMessageId, entry]) => entry.turn && turnKey(entry.turn.threadId, entry.turn.turnId) === turn - ? [clientMessageId] + ? [{ clientMessageId, via: entry.turn.via }] : [] ) - for (const clientMessageId of settled) { + for (const { clientMessageId } of settled) { armed.delete(clientMessageId) } return settled diff --git a/src/main/codex/codex-structured-session-state.ts b/src/main/codex/codex-structured-session-state.ts index 277a2719673..4aa6485b1bf 100644 --- a/src/main/codex/codex-structured-session-state.ts +++ b/src/main/codex/codex-structured-session-state.ts @@ -1,4 +1,5 @@ import type { + AgentJournalAnsweredTurnIdentity, AgentJournalItemIdentity, AgentSessionJournalIdentity } from '../../shared/agent-session-journal-types' @@ -86,7 +87,7 @@ export type CodexStructuredSessionAdapterDeps = { | { providerIdentity: AgentJournalItemIdentity } | ({ state: 'rejected' - answeredInTurn: AgentJournalItemIdentity + answeredInTurn: AgentJournalAnsweredTurnIdentity } & AgentJournalDispatchRejection) ) ) => void diff --git a/src/main/codex/codex-structured-turn-end-settlement.test.ts b/src/main/codex/codex-structured-turn-end-settlement.test.ts index ccb9fa7ba45..2b6542a4333 100644 --- a/src/main/codex/codex-structured-turn-end-settlement.test.ts +++ b/src/main/codex/codex-structured-turn-end-settlement.test.ts @@ -317,7 +317,10 @@ describe('the turn a withdrawn Codex send names', () => { expect(rig.turnRecords.length).toBeGreaterThan(0) expect(new Set(rig.turnRecords.map(agentJournalItemKey)).size).toBe(1) expect(rig.settlements).toEqual([ - expect.objectContaining({ clientMessageId: 'client-1', answeredInTurn: rig.turnRecords[0] }) + expect.objectContaining({ + clientMessageId: 'client-1', + answeredInTurn: { turn: rig.turnRecords[0], via: 'start' } + }) ]) }) @@ -329,8 +332,14 @@ describe('the turn a withdrawn Codex send names', () => { rig.turns.end('interrupted') expect(rig.settlements.map((settlement) => [settlement.clientMessageId, settlement])).toEqual([ - ['client-1', expect.objectContaining({ answeredInTurn: rig.turnRecords[0] })], - ['client-2', expect.objectContaining({ answeredInTurn: rig.turnRecords[0] })] + [ + 'client-1', + expect.objectContaining({ answeredInTurn: { turn: rig.turnRecords[0], via: 'start' } }) + ], + [ + 'client-2', + expect.objectContaining({ answeredInTurn: { turn: rig.turnRecords[0], via: 'steer' } }) + ] ]) }) @@ -345,7 +354,7 @@ describe('the turn a withdrawn Codex send names', () => { await expect(sending).resolves.toMatchObject({ state: 'rejected', - answeredInTurn: rig.turnRecords[0] + answeredInTurn: { turn: rig.turnRecords[0], via: 'start' } }) }) @@ -369,9 +378,11 @@ describe('a send bound to a turn', () => { it('dies with the settlement its turn end makes', () => { const echoes = createCodexDispatchEchoes() echoes.arm('client-1') - echoes.bindTurn('client-1', 'thread-1', 'turn-1') + echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start') - expect(echoes.endTurn('thread-1', 'turn-1', { status: 'interrupted' })).toEqual(['client-1']) + expect(echoes.endTurn('thread-1', 'turn-1', { status: 'interrupted' })).toEqual([ + { clientMessageId: 'client-1', via: 'start' } + ]) expect(echoes.size).toBe(0) expect(echoes.settle('client-1')).toBe(false) }) @@ -379,20 +390,20 @@ describe('a send bound to a turn', () => { it('dies with its child, which forgets recorded turn ends too', () => { const echoes = createCodexDispatchEchoes() echoes.arm('client-1') - echoes.bindTurn('client-1', 'thread-1', 'turn-1') + echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start') echoes.endTurn('thread-2', 'turn-2', { status: 'interrupted' }) echoes.clear() expect(echoes.size).toBe(0) echoes.arm('client-2') - expect(echoes.bindTurn('client-2', 'thread-2', 'turn-2')).toBeNull() + expect(echoes.bindTurn('client-2', 'thread-2', 'turn-2', 'start')).toBeNull() }) it('is matched by thread as well as turn id', () => { const echoes = createCodexDispatchEchoes() echoes.arm('client-1') - echoes.bindTurn('client-1', 'thread-1', 'turn-1') + echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start') expect(echoes.endTurn('thread-2', 'turn-1', { status: 'interrupted' })).toEqual([]) expect(echoes.size).toBe(1) diff --git a/src/main/codex/codex-structured-turn-end-settlement.ts b/src/main/codex/codex-structured-turn-end-settlement.ts index 931f4d41a70..376d9adb6e1 100644 --- a/src/main/codex/codex-structured-turn-end-settlement.ts +++ b/src/main/codex/codex-structured-turn-end-settlement.ts @@ -17,7 +17,7 @@ import { agentSessionFailureWords, type AgentJournalDispatchRejection } from '../../shared/agent-session-failure-words' -import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types' +import type { AgentJournalAnsweredTurnIdentity } from '../../shared/agent-session-journal-types' import type { CodexTurnEnd } from './codex-structured-dispatch-echo' import { codexTurnLifecycleIdentity } from './codex-structured-journal-translation-turns' import type { CodexSession } from './codex-structured-session-state' @@ -43,8 +43,8 @@ export function codexDispatchRejection( export type CodexTurnEndSettlement = { clientMessageId: string state: 'rejected' - /** The turn Codex answered the send into, whose end settled it. */ - answeredInTurn: AgentJournalItemIdentity + /** The turn Codex answered the send into, whose end settled it, and how the send joined it. */ + answeredInTurn: AgentJournalAnsweredTurnIdentity } & AgentJournalDispatchRejection function errorDetail(params: unknown): ProviderDiagnostic | undefined { @@ -97,10 +97,14 @@ export function settleCodexSendsInEndedTurn( return } const rejection = codexTurnEndRejection(end) - const answeredInTurn = codexTurnLifecycleIdentity(frame.sessionId, turnId) - for (const clientMessageId of session.dispatchEchoes.endTurn(session.threadId, turnId, end)) { + const turn = codexTurnLifecycleIdentity(frame.sessionId, turnId) + for (const { clientMessageId, via } of session.dispatchEchoes.endTurn( + session.threadId, + turnId, + end + )) { if (rejection) { - settle({ clientMessageId, state: 'rejected', answeredInTurn, ...rejection }) + settle({ clientMessageId, state: 'rejected', answeredInTurn: { turn, via }, ...rejection }) } } } diff --git a/src/main/codex/codex-structured-turn-start.ts b/src/main/codex/codex-structured-turn-start.ts index 0aace8ee54b..271def87607 100644 --- a/src/main/codex/codex-structured-turn-start.ts +++ b/src/main/codex/codex-structured-turn-start.ts @@ -1,5 +1,8 @@ import { agentSessionFailureFact, providerDiagnosticOf } from '../../shared/agent-session-failure' -import type { AgentJournalMessageItem } from '../../shared/agent-session-journal-types' +import type { + AgentJournalMessageItem, + AgentJournalTurnJoin +} from '../../shared/agent-session-journal-types' import type { NativeChatBlock } from '../../shared/native-chat-types' import type { AgentSessionDispatchOutcome } from '../native-chat/agent-session-wire/structured-agent-session-adapter' import { @@ -108,7 +111,7 @@ async function steerCodexTurn( host: CodexTurnHost, expectedTurnId: string, input: { clientMessageId: string; body: AgentJournalMessageItem; timeoutMs?: number } -): Promise<{ turnId: string } | null> { +): Promise<{ turnId: string; via: 'steer' } | null> { try { const answer = await host.connection.request( 'turn/steer', @@ -120,7 +123,7 @@ async function steerCodexTurn( }, { timeoutMs: input.timeoutMs } ) - return { turnId: readCodexTurnId(answer) ?? expectedTurnId } + return { turnId: readCodexTurnId(answer) ?? expectedTurnId, via: 'steer' } } catch (error) { if (isCodexAppServerRequestError(error) || isCodexAppServerUnsupportedError(error)) { return null @@ -142,7 +145,7 @@ export async function startCodexTurn( requestedAt?: number timeoutMs?: number } -): Promise<{ turnId: string | null } | false> { +): Promise<{ turnId: string | null; via: AgentJournalTurnJoin } | false> { // Armed before the write: the echo and `turn/started` can both land while the // response is in flight, and the start must snapshot this send in its frontier. if (!host.dispatchEchoes.arm(input.clientMessageId, input.requestedAt)) { @@ -168,7 +171,7 @@ export async function startCodexTurn( }, { timeoutMs: input.timeoutMs } ) - return { turnId: readCodexTurnId(answer) } + return { turnId: readCodexTurnId(answer), via: 'start' } } /** @@ -187,7 +190,7 @@ export async function dispatchCodexTurn( }, timeoutMs: number | undefined ): Promise { - let answer: { turnId: string | null } | false + let answer: { turnId: string | null; via: AgentJournalTurnJoin } | false try { answer = await startCodexTurn(session, { ...input, timeoutMs }) } catch (error) { @@ -211,13 +214,21 @@ export async function dispatchCodexTurn( } // An answer read after the turn it names already ended is settled by that end. const endedFirst = answer.turnId - ? session.dispatchEchoes.bindTurn(input.clientMessageId, session.threadId, answer.turnId) + ? session.dispatchEchoes.bindTurn( + input.clientMessageId, + session.threadId, + answer.turnId, + answer.via + ) : null const rejection = endedFirst ? codexTurnEndRejection(endedFirst) : null return rejection && answer.turnId ? { state: 'rejected', - answeredInTurn: codexTurnLifecycleIdentity(input.sessionId, answer.turnId), + answeredInTurn: { + turn: codexTurnLifecycleIdentity(input.sessionId, answer.turnId), + via: answer.via + }, ...rejection } : { state: 'admitted' } diff --git a/src/main/native-chat/agent-session-journal/journal-dispatch-reducer.ts b/src/main/native-chat/agent-session-journal/journal-dispatch-reducer.ts index c2f7d45327a..b19fa2f1ce9 100644 --- a/src/main/native-chat/agent-session-journal/journal-dispatch-reducer.ts +++ b/src/main/native-chat/agent-session-journal/journal-dispatch-reducer.ts @@ -6,6 +6,7 @@ import { type UnreadAgentSessionFailureFact } from '../../../shared/agent-session-failure' import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key' +import type { AgentJournalAnsweredTurn } from '../../../shared/agent-session-journal-types' import { journalDispatchRowApplies } from './journal-dispatch-settlement' import type { JournalReducerState } from './journal-reducer' import { @@ -34,14 +35,11 @@ export function applyJournalDispatchRow( } else { delete submission.rejection } - if ( - row.state === 'rejected' && - typeof row.answeredInTurnItemId === 'string' && - row.answeredInTurnItemId - ) { - submission.answeredInTurnItemId = row.answeredInTurnItemId + const answeredInTurn = row.state === 'rejected' ? readAnsweredTurn(row.answeredInTurn) : undefined + if (answeredInTurn) { + submission.answeredInTurn = answeredInTurn } else { - delete submission.answeredInTurnItemId + delete submission.answeredInTurn } submission.resolvedAt = row.state === 'pending' ? null : row.ts if (row.state === 'pending') { @@ -70,6 +68,19 @@ export function applyJournalDispatchRow( }) } +/** A stored answered turn; one malformed, or naming a way of joining this build does not know, is + * dropped, never the row. */ +function readAnsweredTurn(value: unknown): AgentJournalAnsweredTurn | undefined { + if (typeof value !== 'object' || value === null) { + return undefined + } + const turnItemId = 'turnItemId' in value ? value.turnItemId : undefined + const via = 'via' in value ? value.via : undefined + return typeof turnItemId === 'string' && turnItemId && (via === 'start' || via === 'steer') + ? { turnItemId, via } + : undefined +} + /** A stored rejection fact, read where it can be placed; a kind it cannot place is kept as * written, so the classifier still knows a fact was there without this build claiming what it * says. Shared with the queued-draft table, whose returned card mirrors its submission. */ diff --git a/src/main/native-chat/agent-session-journal/journal-row-builders.ts b/src/main/native-chat/agent-session-journal/journal-row-builders.ts index 6873db0fd46..4484f26a93f 100644 --- a/src/main/native-chat/agent-session-journal/journal-row-builders.ts +++ b/src/main/native-chat/agent-session-journal/journal-row-builders.ts @@ -111,7 +111,12 @@ export function journalDispatchRowBuilder( reason: boundedDispatchReason(input), ...(input.state === 'rejected' ? { rejection: input.rejection } : {}), ...(input.state === 'rejected' && input.answeredInTurn - ? { answeredInTurnItemId: agentJournalItemKey(input.answeredInTurn) } + ? { + answeredInTurn: { + turnItemId: agentJournalItemKey(input.answeredInTurn.turn), + via: input.answeredInTurn.via + } + } : {}), ...journalRowBase(state().epoch, seq, input.fence, ts), ...(input.recovered ? { recovered: input.recovered } : {}), diff --git a/src/main/native-chat/agent-session-journal/journal-row-schema.ts b/src/main/native-chat/agent-session-journal/journal-row-schema.ts index 194290542d9..4957bd846be 100644 --- a/src/main/native-chat/agent-session-journal/journal-row-schema.ts +++ b/src/main/native-chat/agent-session-journal/journal-row-schema.ts @@ -16,6 +16,7 @@ import type { AgentSessionFailureFact } from '../../../shared/agent-session-failure' import { AGENT_SESSION_JOURNAL_SCHEMA_VERSION, + type AgentJournalAnsweredTurn, type AgentJournalDispatchState, type AgentJournalItemBody, type AgentJournalMessageItem, @@ -150,10 +151,10 @@ export type JournalDispatchRow = JournalRowBase & { /** On `rejected`: why, typed. Older readers keep the key and ignore it; a malformed one is * dropped when read, never the row. */ rejection?: AgentSessionFailureFact - /** On `rejected`: the turn record a Codex send was answered into (`turn/start` or `turn/steer`) - * when that turn's end settled the send. Absent on every other row. Older readers keep the key - * and ignore it. */ - answeredInTurnItemId?: string + /** On `rejected`: the turn a Codex send was answered into, and how it joined it, when that + * turn's end settled the send. Absent on every other row. Older readers keep the key and ignore + * it; a malformed one is dropped when read, never the row. */ + answeredInTurn?: AgentJournalAnsweredTurn } /** An item mutation may name its own producer, because one batch can CREATE diff --git a/src/main/native-chat/agent-session-journal/journal-store-contracts.ts b/src/main/native-chat/agent-session-journal/journal-store-contracts.ts index b2d2914b318..4232e805bbe 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-contracts.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-contracts.ts @@ -1,5 +1,6 @@ import type { AgentJournalDispatchRejection } from '../../../shared/agent-session-failure-words' import type { + AgentJournalAnsweredTurnIdentity, AgentJournalCursor, AgentJournalItemBody, AgentJournalItemIdentity, @@ -42,7 +43,7 @@ export type ResolveDispatchInput = { * `agentSessionFailureWords`, never written by hand. */ | ({ state: 'rejected' - answeredInTurn?: AgentJournalItemIdentity + answeredInTurn?: AgentJournalAnsweredTurnIdentity } & AgentJournalDispatchRejection) | { state: 'unknown'; reason?: string | null } ) diff --git a/src/main/native-chat/agent-session-journal/journal-submission-positions.test.ts b/src/main/native-chat/agent-session-journal/journal-submission-positions.test.ts index 31cfee3a25c..2ad8953ea03 100644 --- a/src/main/native-chat/agent-session-journal/journal-submission-positions.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-submission-positions.test.ts @@ -94,7 +94,7 @@ async function sendWithdrawnByItsTurnEnd() { clientMessageId: 'send-1', state: 'rejected', ...agentSessionFailureWords(agentSessionFailureFact('cancelled'), { surface: 'rejection' }), - answeredInTurn: TURN, + answeredInTurn: { turn: TURN, via: 'start' }, fence: 1 }) return journal @@ -107,23 +107,25 @@ describe('the turn a rejected submission was answered into', () => { expect(journal.submission('send-1')).toMatchObject({ dispatchState: 'rejected', - answeredInTurnItemId: turnItemId + answeredInTurn: { turnItemId, via: 'start' } }) expect(readAgentSessionHydrationPage(journal).submissions).toEqual([ - expect.objectContaining({ answeredInTurnItemId: turnItemId }) + expect.objectContaining({ answeredInTurn: { turnItemId, via: 'start' } }) ]) await journals.closeAll() const replayed = await journals.open({ identity: IDENTITY, stateDirectory: root! }) - expect(replayed.submission('send-1')).toMatchObject({ answeredInTurnItemId: turnItemId }) + expect(replayed.submission('send-1')).toMatchObject({ + answeredInTurn: { turnItemId, via: 'start' } + }) }) it('is absent on a take-back that names no turn, as on rows from older hosts', async () => { const { journal } = await sendHandedOverThenWithdrawn() expect(journal.submission('send-1')).toMatchObject({ dispatchState: 'rejected' }) - expect(journal.submission('send-1')).not.toHaveProperty('answeredInTurnItemId') + expect(journal.submission('send-1')).not.toHaveProperty('answeredInTurn') expect(readAgentSessionHydrationPage(journal).submissions[0]).not.toHaveProperty( - 'answeredInTurnItemId' + 'answeredInTurn' ) }) }) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts index 2aacb4d9633..d08bf45b8d3 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts @@ -12,6 +12,7 @@ import type { import type { AgentSessionBackgroundTaskStops } from '../../../shared/agent-child-work-stop-targets' import type { + AgentJournalAnsweredTurnIdentity, AgentJournalItemIdentity, AgentJournalItemBody, AgentJournalMessageItem, @@ -180,7 +181,7 @@ export type AgentSessionDispatchOutcome = * provider answered the send into, which ended before the answer was read. */ | ({ state: 'rejected' - answeredInTurn?: AgentJournalItemIdentity + answeredInTurn?: AgentJournalAnsweredTurnIdentity } & AgentJournalDispatchRejection) /** The call did not settle. Never re-send on the user's behalf. */ | { state: 'unknown'; reason: string } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-late-dispatch.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-late-dispatch.ts index 4384d5fd561..3d38494fea9 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-late-dispatch.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-late-dispatch.ts @@ -1,4 +1,7 @@ -import type { AgentJournalItemIdentity } from '../../../shared/agent-session-journal-types' +import type { + AgentJournalAnsweredTurnIdentity, + AgentJournalItemIdentity +} from '../../../shared/agent-session-journal-types' import type { AgentJournalDispatchRejection } from '../../../shared/agent-session-failure-words' import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations' import { structuredAgentSessionConversationFence } from './structured-agent-session-provider-child' @@ -14,7 +17,7 @@ export async function settleStructuredAgentSessionLateDispatch( | { providerIdentity: AgentJournalItemIdentity } | ({ state: 'rejected' - answeredInTurn?: AgentJournalItemIdentity + answeredInTurn?: AgentJournalAnsweredTurnIdentity } & AgentJournalDispatchRejection) | { state: 'unknown'; reason: string } ) diff --git a/src/main/runtime/structured-agent-session-codex-turn-end-settlement.test.ts b/src/main/runtime/structured-agent-session-codex-turn-end-settlement.test.ts index ccc922a70e5..0cd0570b83d 100644 --- a/src/main/runtime/structured-agent-session-codex-turn-end-settlement.test.ts +++ b/src/main/runtime/structured-agent-session-codex-turn-end-settlement.test.ts @@ -317,12 +317,12 @@ describe('the turn a withdrawn Codex send was answered into', () => { item.body.kind === 'turn' ? [item.itemId] : [] ), named: snapshot.submissions.find((entry) => entry.clientMessageId === clientMessageId) - ?.answeredInTurnItemId, - onPage: onPage?.answeredInTurnItemId + ?.answeredInTurn, + onPage: onPage?.answeredInTurn } } - it("is that turn's record, for the send that opened it and for one steered into it", async () => { + it("is that turn's record, started by the send that opened it and steered by a later one", async () => { const opening = await send('look around') await vi.waitFor(() => expect(answers).toBe(1)) turns.start() @@ -336,13 +336,13 @@ describe('the turn a withdrawn Codex send was answered into', () => { const { turnRecords } = await answeredInto(opening) expect(turnRecords).toHaveLength(1) - expect(await answeredInto(opening)).toMatchObject({ - named: turnRecords[0], - onPage: turnRecords[0] - }) - expect(await answeredInto(steered)).toMatchObject({ - named: turnRecords[0], - onPage: turnRecords[0] + const started = { turnItemId: turnRecords[0], via: 'start' } + const steeredIn = { turnItemId: turnRecords[0], via: 'steer' } + expect(await answeredInto(opening)).toEqual({ turnRecords, named: started, onPage: started }) + expect(await answeredInto(steered)).toEqual({ + turnRecords, + named: steeredIn, + onPage: steeredIn }) }) @@ -373,7 +373,7 @@ describe('the turn a withdrawn Codex send was answered into', () => { ) const { turnRecords, named } = await answeredInto(sent) expect(turnRecords).toHaveLength(1) - expect(named).toBe(turnRecords[0]) + expect(named).toEqual({ turnItemId: turnRecords[0], via: 'start' }) }) }) diff --git a/src/shared/agent-session-journal-schemas.ts b/src/shared/agent-session-journal-schemas.ts index 00d88b74de5..7a21d5816a2 100644 --- a/src/shared/agent-session-journal-schemas.ts +++ b/src/shared/agent-session-journal-schemas.ts @@ -319,7 +319,8 @@ export const AgentJournalSubmissionSchema = z.object({ submittedAt: z.number(), resolvedAt: z.number().nullable(), submittedSequence: z.number().int().optional(), - answeredInTurnItemId: z.string().min(1).optional(), + // Open like `dispatchState`: a way of joining a newer host names must not drop the submission. + answeredInTurn: z.object({ turnItemId: z.string().min(1), via: z.string().min(1) }).optional(), recovered: z.literal(true).optional(), handoverRecorded: z.literal(true).optional(), handedOverAt: z.number().optional(), diff --git a/src/shared/agent-session-journal-types.ts b/src/shared/agent-session-journal-types.ts index 4e21ee5f2e8..b6b33950de4 100644 --- a/src/shared/agent-session-journal-types.ts +++ b/src/shared/agent-session-journal-types.ts @@ -391,6 +391,18 @@ export type AgentJournalRenderItem = AgentJournalProducerLinkage & { // ─── Submissions ──────────────────────────────────────────────────────────── +/** The turn a send was answered into: its record's item id, and how the send joined it. `start`: + * the provider answered the send's start request with that turn; `steer`: Orca steered it into + * that running turn. Known limit: a start the provider silently folds into a running turn reads + * as `start`. Open: a newer host may name another way, which a reader leaves unclaimed. */ +export type AgentJournalAnsweredTurn = { turnItemId: string; via: AgentJournalTurnJoin } +export type AgentJournalTurnJoin = 'start' | 'steer' +/** The same, as a writer names it: the turn record's identity, keyed when the row is written. */ +export type AgentJournalAnsweredTurnIdentity = { + turn: AgentJournalItemIdentity + via: AgentJournalTurnJoin +} + export const AGENT_JOURNAL_DISPATCH_STATES = ['pending', 'accepted', 'rejected', 'unknown'] as const export type AgentJournalDispatchState = (typeof AGENT_JOURNAL_DISPATCH_STATES)[number] @@ -415,10 +427,10 @@ export type AgentJournalSubmission = { * never stored. Needed because a rejected send's own row moves to its rejection, which erases * where it was sent. Absent from hosts that predate it. */ submittedSequence?: number - /** On `rejected`: the item id of the turn record a Codex send was answered into, when that turn - * ended without taking it. Absent on every other send: accepted ones (the echo places them), - * Claude, queued take-backs, restart recovery, and rows from hosts that predate it. */ - answeredInTurnItemId?: string + /** On `rejected`: the turn a Codex send was answered into, when that turn ended without taking + * it. Absent on every other send: accepted ones (the echo places them), Claude, queued + * take-backs, restart recovery, and rows from hosts that predate it. */ + answeredInTurn?: AgentJournalAnsweredTurn /** Set when crash reconciliation resolved the dispatch, not the provider. A live * `unknown` is a send still outstanding; a recovered one outlived its writer. */ recovered?: true diff --git a/tests/e2e/cross-version-wire/submission-positions-downgrade.unit.test.ts b/tests/e2e/cross-version-wire/submission-positions-downgrade.unit.test.ts index 46ddf9715de..5265cfb1af9 100644 --- a/tests/e2e/cross-version-wire/submission-positions-downgrade.unit.test.ts +++ b/tests/e2e/cross-version-wire/submission-positions-downgrade.unit.test.ts @@ -40,7 +40,7 @@ function withoutPosition(page: AgentSessionHistoryPage): AgentSessionHistoryPage return { ...page, submissions: page.submissions.map( - ({ submittedSequence: _submitted, answeredInTurnItemId: _turn, ...rest }) => rest + ({ submittedSequence: _submitted, answeredInTurn: _turn, ...rest }) => rest ) } } @@ -78,10 +78,13 @@ test('a released client folds and draws a page whose submissions carry their jou state: 'rejected', ...agentSessionFailureWords(agentSessionFailureFact('cancelled'), { surface: 'rejection' }), answeredInTurn: { - provider: 'legacy', - agent: 'codex', - sessionId: IDENTITY.sessionId, - recordId: 'turn-lifecycle:turn-1' + turn: { + provider: 'legacy', + agent: 'codex', + sessionId: IDENTITY.sessionId, + recordId: 'turn-lifecycle:turn-1' + }, + via: 'start' }, fence: 1 }) @@ -93,7 +96,9 @@ test('a released client folds and draws a page whose submissions carry their jou }) ) expect(page.submissions.find((entry) => entry.clientMessageId === 'turn-ended')).toEqual( - expect.objectContaining({ answeredInTurnItemId: expect.any(String) }) + expect.objectContaining({ + answeredInTurn: { turnItemId: expect.any(String), via: 'start' } + }) ) const checkout = await materializeReleaseCheckout(BASELINE_REF) @@ -117,7 +122,7 @@ test('a released client folds and draws a page whose submissions carry their jou const state = reduce(empty, { type: 'history-page', page: from }) return { submissions: state.submissions.map( - ({ submittedSequence: _s, answeredInTurnItemId: _t, ...rest }) => rest + ({ submittedSequence: _s, answeredInTurn: _t, ...rest }) => rest ), messages: project(state.items, [], state.submissions) }