diff --git a/src/main/native-chat/agent-session-journal/journal-stop-turn-end.ts b/src/main/native-chat/agent-session-journal/journal-stop-turn-end.ts index b0975d0d607..821ceb78138 100644 --- a/src/main/native-chat/agent-session-journal/journal-stop-turn-end.ts +++ b/src/main/native-chat/agent-session-journal/journal-stop-turn-end.ts @@ -26,10 +26,7 @@ function stopIsAPersons(reason: JournalStopEvent['reason']): boolean { } } -type TurnEndState = Pick< - JournalReducerState, - 'items' | 'queuePauseMarks' | 'latestPersonTurnSequence' -> +type TurnEndState = Pick /** Whether `stop`, a person's, makes the end of turn `turnId` theirs: it named that turn, or named * none and stopped the turn item `itemId` opened. */ @@ -48,20 +45,26 @@ function stopIsTurnCancellation( } /** Pressed before any turn showed, a Stop stopped the first turn opened after it, and no later - * one: unless a send a person made since was accepted, whose turn that is. `itemId` null: a turn - * not yet opened. */ + * one: unless a send journaled since, of any origin, was not refused, whose turn that is. A Stop + * whose stopped send never opens a turn so binds nothing. `itemId` null: a turn not yet opened. */ function turnlessStopStopped( state: TurnEndState, stop: JournalLatestStop, itemId: string | null ): boolean { const createdAt = itemId === null ? null : (state.items.get(itemId)?.sequence ?? null) - if ( - (createdAt !== null && createdAt <= stop.sequence) || - state.latestPersonTurnSequence >= stop.sequence - ) { + if (createdAt !== null && createdAt <= stop.sequence) { return false } + for (const submission of state.submissions.values()) { + if ( + submission.dispatchState !== 'rejected' && + submission.acceptedSequence !== undefined && + submission.acceptedSequence > stop.sequence + ) { + return false + } + } for (const [otherId, item] of state.items) { if ( otherId !== itemId && @@ -103,19 +106,17 @@ export function personStopDecidesTurn( turnId: string | null, endedAt?: number ): boolean { - if (turnId !== null) { - const itemId = [...state.items].find( - ([, item]) => readAgentJournalTurn(item.body)?.turnId === turnId - )?.[0] - return stopEndsTurnAsCancellation(state, turnId, itemId ?? null, endedAt) - } const stop = state.queuePauseMarks.latestStop - return ( - stop !== null && - stop.event.turnId === undefined && - stopIsAPersons(stop.event.reason) && - turnlessStopStopped(state, stop, null) - ) + if (stop === null || !stopIsAPersons(stop.event.reason)) { + return false + } + if (turnId === null) { + return stop.event.turnId === undefined && turnlessStopStopped(state, stop, null) + } + const itemId = [...state.items].find( + ([, item]) => readAgentJournalTurn(item.body)?.turnId === turnId + )?.[0] + return stopEndsTurnAsCancellation(state, turnId, itemId ?? null, endedAt) } /** diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-stop-event-binding.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-stop-event-binding.test.ts new file mode 100644 index 00000000000..c7cef759d54 --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-stop-event-binding.test.ts @@ -0,0 +1,128 @@ +// Which turn a person's Stop that named no turn binds: only the one its stopped send opens. A Stop of +// a start that never landed stopped a send that opens no turn, and a send journaled after the Stop +// opens its own; neither is the Stop's, whatever sent it. + +import { afterEach, describe, expect, it } from 'vitest' +import { + AGENT_JOURNAL_THREAD_SCOPE, + type AgentJournalItemIdentity +} from '../../../shared/agent-session-journal-types' +import { readAgentJournalTurn } from '../../../shared/agent-session-turn-record' +import type { JournalStopEvent } from '../agent-session-journal/journal-row-schema' +import { settleStaleStructuredAgentSessionState } from './structured-agent-session-dead-generation-settlement' +import { HOST_TEST_SESSION } from './structured-agent-session-host-test-data' +import { + createQueuedMessageTestRig, + eventually, + type QueuedMessageTestRig +} from './structured-agent-session-queued-message-rig.test-fixture' + +let rig: QueuedMessageTestRig + +afterEach(() => rig.dispose()) + +const MAIL_TURN: AgentJournalItemIdentity = { + provider: 'codex', + threadId: 'thread-1', + turnId: 'turn-mail', + ordinal: 999 +} + +function journal() { + const open = rig.host.collaboratorsForTests().sessions.get(HOST_TEST_SESSION)?.journal + if (!open) { + throw new Error('expected the conversation open') + } + return open +} + +function stopEvents(): JournalStopEvent[] { + const since = journal().readSince({ epoch: journal().epoch, sequence: 0 }) + if (!since.ok) { + throw new Error(`expected rows, got reset ${since.reset}`) + } + return since.rows.flatMap((row) => + row.kind === 'tombstone' && row.stopEvent ? [row.stopEvent] : [] + ) +} + +function childPhase() { + return rig.host.collaboratorsForTests().sessions.get(HOST_TEST_SESSION)?.child?.phase +} + +function fence(): number { + return rig.store.getRecord(HOST_TEST_SESSION)?.lease.runtimeFence ?? 1 +} + +function mailTurn() { + return journal() + .snapshot() + .items.map((item) => readAgentJournalTurn(item.body)) + .find((turn) => turn?.turnId === 'turn-mail') +} + +/** A person's Stop of a start that never landed, whose send opens no turn; then orchestration mail + * starts a new child, which lands and runs the mail's turn. */ +async function mailTurnAfterStopOfStart(): Promise { + rig = await createQueuedMessageTestRig({ starting: true, restartable: true }) + let release: () => void = () => undefined + rig.awaitStarted.mockImplementation( + () => new Promise((resolve) => (release = () => resolve(undefined))) + ) + rig.send('work on this') + await eventually(() => expect(childPhase()).toBe('starting')) + expect(await rig.stop()).toMatchObject({ ok: true }) + release() + expect(stopEvents()).toEqual([expect.objectContaining({ reason: 'user-stop' })]) + expect(stopEvents()[0]).not.toHaveProperty('turnId') + await eventually(() => expect(childPhase()).toBeUndefined()) + rig.awaitStarted.mockImplementation(async () => undefined) + await rig.send('mail for the worker', undefined, { internal: true }).result + await eventually(() => expect(rig.dispatch).toHaveBeenCalled()) + await journal().appendItem( + MAIL_TURN, + { kind: 'turn', turnId: 'turn-mail', state: 'running', startedAt: Date.now() }, + { fence: fence(), turnScope: AGENT_JOURNAL_THREAD_SCOPE } + ) +} + +describe('a Stop of a start that never landed binds no later turn', () => { + it("writes the host's event when it evicts the mail turn, which reads as news", async () => { + await mailTurnAfterStopOfStart() + let atClose: JournalStopEvent[] = [] + rig.closeSession.mockImplementationOnce(async () => { + atClose = stopEvents() + return true + }) + + await rig.host.close(HOST_TEST_SESSION, 'evict') + + expect(atClose.map((event) => event.reason)).toEqual(['user-stop', 'evict']) + await rig.host.journalSnapshot(HOST_TEST_SESSION) + expect(mailTurn()).toMatchObject({ state: 'interrupted' }) + expect(mailTurn()).not.toHaveProperty('outcome') + }) + + it('settles a crash of the mail turn on relaunch as news', async () => { + await mailTurnAfterStopOfStart() + const owner = fence() + rig.crashRestartHostProcess() + await rig.host.journalSnapshot(HOST_TEST_SESSION) + + await settleStaleStructuredAgentSessionState({ + journal: journal(), + sessionId: HOST_TEST_SESSION, + fence: owner + 1, + acquisitionGeneration: 'generation-2', + deathEvidence: { + kind: 'exit-observed', + detail: 'the relaunch proved the old child gone', + observedAt: Date.now() + 60_000, + ownerFence: owner + } + }) + + expect(mailTurn()).toMatchObject({ state: 'interrupted' }) + expect(mailTurn()).not.toHaveProperty('outcome') + }) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-stop-event-entries.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-stop-event-entries.test.ts index 26f191a7f90..bccc020f7cb 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-stop-event-entries.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-stop-event-entries.test.ts @@ -60,28 +60,6 @@ async function runningTurn(turnId = 'turn-1'): Promise { return working } -/** The stopped send's turn opens after its turnless Stop and ends cut: the turn that Stop bound. */ -async function stoppedTurnEnded(): Promise { - const fence = rig.store.getRecord(HOST_TEST_SESSION)?.lease.runtimeFence ?? 1 - const scope = { fence, turnScope: AGENT_JOURNAL_THREAD_SCOPE } - const identity = { - provider: 'codex' as const, - threadId: 'thread-1', - turnId: 'turn-stopped', - ordinal: 998 - } - await journal().appendItem( - identity, - { kind: 'turn', turnId: 'turn-stopped', state: 'running', startedAt: 1 }, - scope - ) - await journal().appendItem( - identity, - { kind: 'turn', turnId: 'turn-stopped', state: 'interrupted', completedAt: Date.now() }, - scope - ) -} - async function queuedDraft(text: string): Promise { const queued = await rig.send(text, 'queue-if-active').result if (!queued.ok || !('queued' in queued.value)) { @@ -257,10 +235,11 @@ describe("a person's Stop pause and the Stop events after it", () => { await queuedDraft('queued behind the turn') expect(await rig.stop()).toMatchObject({ ok: true }) await rig.settleAccepted(working, 'stopped') - await stoppedTurnEnded() expect(await rig.queuePause()).toEqual({ reason: 'stopped' }) // Orchestration mail starts a turn the host sent, which lifts nothing. - await rig.send('mail for the lead', undefined, { internal: true }).result + const mail = rig.send('mail for the lead', undefined, { internal: true }) + await mail.result + await rig.settleAccepted(mail.id, 'mail') await journal().appendItem( { provider: 'codex', threadId: 'thread-1', turnId: 'turn-mail', ordinal: 999 }, { kind: 'turn', turnId: 'turn-mail', state: 'running', startedAt: 1 }, @@ -285,7 +264,6 @@ describe("a person's Stop pause and the Stop events after it", () => { const held = await queuedDraft('queued behind the turn') expect(await rig.stop()).toMatchObject({ ok: true }) await rig.settleAccepted(working, 'stopped') - await stoppedTurnEnded() // The agent at rest goes, writing nothing; mail then starts a new child, which never lands. await idleSweep().tick() expect(await rig.queuePause()).toEqual({ reason: 'stopped' })