diff --git a/src/main/codex/codex-structured-conversation-stop.test.ts b/src/main/codex/codex-structured-conversation-stop.test.ts index 5d8a9c5627f..62ecaf8f2ea 100644 --- a/src/main/codex/codex-structured-conversation-stop.test.ts +++ b/src/main/codex/codex-structured-conversation-stop.test.ts @@ -52,7 +52,9 @@ describe('a Codex Stop that names no turn', () => { detail: { text: 'expected active turn id turn-journal but found turn-1', audience: 'person' - } + }, + // An invalid-request refusal: the named turn is not the one Codex is running. + turnNotRunning: true } }) expect(rig.interrupts().map((call) => call.params?.turnId)).toEqual(['turn-journal']) diff --git a/src/main/codex/codex-structured-session-adapter-fixture.ts b/src/main/codex/codex-structured-session-adapter-fixture.ts index 45c97f00066..f2b4e102c35 100644 --- a/src/main/codex/codex-structured-session-adapter-fixture.ts +++ b/src/main/codex/codex-structured-session-adapter-fixture.ts @@ -13,6 +13,7 @@ import { type CodexStructuredLaunch, type CodexStructuredSessionEvent } from './codex-structured-session-adapter' +import type { CodexStructuredSessionAdapterDeps } from './codex-structured-session-state' export const THREAD_ID = 'thread-abc' @@ -104,7 +105,9 @@ export function answerWithOpenedTurn( export function adapterFor( codex: ReturnType, launch: Partial = {}, - events: CodexStructuredSessionEvent[] = [] + events: CodexStructuredSessionEvent[] = [], + /** Host wiring the runtime adds, such as the late dispatch settlement. */ + deps: Partial = {} ): CodexStructuredSessionAdapter { let acquisitionGeneration = 0 return new CodexStructuredSessionAdapter({ @@ -120,7 +123,8 @@ export function adapterFor( openConnection: codex.openConnection, readProcessStartTime: async () => 1_700_000_000_000, now: () => 1_700_000_000_500, - mintAcquisitionGeneration: () => `generation-${++acquisitionGeneration}` + mintAcquisitionGeneration: () => `generation-${++acquisitionGeneration}`, + ...deps }) } diff --git a/src/main/codex/codex-structured-turn-cancellation.ts b/src/main/codex/codex-structured-turn-cancellation.ts index 99af63f9d6f..4330ebed0fb 100644 --- a/src/main/codex/codex-structured-turn-cancellation.ts +++ b/src/main/codex/codex-structured-turn-cancellation.ts @@ -5,10 +5,21 @@ import type { AgentSessionCancelOutcome } from '../native-chat/agent-session-wir import type { CodexSession } from './codex-structured-session-state' import type { CodexJournalTranslationAdmission } from './codex-structured-journal-contracts' +/** + * Codex's interrupt handler answers -32600 when the named turn is not its running one (no turn, + * another turn, thread not loaded) and -32603 when it could not submit the interrupt. A request it + * could not parse, or one before initialize, also answers -32600: read as not running, it keeps the + * child, as before. + */ +function isCodexTurnNotRunningRefusal(error: unknown): boolean { + return isCodexAppServerRequestError(error) && error.code === -32600 +} + /** * Codex answers a turn's interrupt only as that turn ends, so the answer is what confirms the - * Stop. The interrupt is the whole Stop: Codex kills the turn's one-shot commands itself and keeps - * its background terminals running until the thread ends. + * Stop. An answered interrupt is the whole Stop: Codex kills the turn's one-shot commands itself + * and keeps its background terminals running until the thread ends. One Codex could not carry out, + * or never answered, leaves the host to end the child (`performCancel`). */ export async function interruptCodexTurn(input: { session: CodexSession @@ -29,7 +40,13 @@ export async function interruptCodexTurn(input: { throw error } const detail = providerDiagnosticOf(error) - return { cancelled: false, refusal: detail ? { detail } : {} } + return { + cancelled: false, + refusal: { + ...(detail ? { detail } : {}), + ...(isCodexTurnNotRunningRefusal(error) ? { turnNotRunning: true } : {}) + } + } } const promptAdmission = input.onConfirmed?.() if (promptAdmission && !promptAdmission.accepted) { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-stop.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-stop.ts index 0d73ec723ce..4709234ad5c 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-stop.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-stop.ts @@ -1,11 +1,20 @@ // How a provider's Stop, and a prompt card's own Cancel, end what they end. Every member is -// optional: a provider that declares none keeps its child after a Stop, and its card's Cancel -// interrupts the turn holding the card. +// optional: a provider that declares none keeps its child after a Stop it took, and its card's +// Cancel interrupts the turn holding the card. import type { AgentJournalApprovalItem, AgentJournalQuestionItem } from '../../../shared/agent-session-journal-types' +import type { ProviderDiagnostic } from '../../../shared/agent-session-failure' + +/** `refusal`: the provider answered the Stop and declined it, in its own words when it gave any. + * `turnNotRunning`: its refusal is the kind it gives for a turn not running there, so the Stop keeps + * the child; absent, it could not interrupt that turn, which may run on. */ +export type AgentSessionCancelOutcome = { + cancelled: boolean + refusal?: { detail?: ProviderDiagnostic; turnNotRunning?: true } +} /** Where a card's Cancel goes: a dismissal (`dismissPrompt`), or the chat's Stop. */ export type AgentSessionPromptCancelRoute = { kind: 'dismiss' } | { kind: 'stop' } @@ -13,7 +22,8 @@ export type AgentSessionPromptCancelRoute = { kind: 'dismiss' } | { kind: 'stop' export type StructuredAgentSessionAdapterStop = { /** A Stop ends this provider's child after `cancelTurn`, whatever it answered, unless it named a * turn that is no longer live and the cancel answered that it did not take it; the next send - * resumes the conversation. Absent or false keeps the child after a Stop. */ + * resumes the conversation. Absent or false keeps the child after a Stop the provider took; a + * failed interrupt still ends a child running the turn the Stop meant. */ stopEndsSession?(sessionId: string): boolean /** What a Stop that ends the session waits on before it ends the child: resolves once the provider * has nothing in flight, a send it has not answered included, or when its grace, counted from 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 91f6d945c50..91ac8a7598c 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 @@ -32,12 +32,13 @@ import type { AgentSessionThreadGoalChange } from '../../../shared/agent-session-wire' import type { AgentSessionRefusalReason } from '../../../shared/agent-session-wire-refusals' -import type { - ProviderDiagnostic, - SubmissionRejectionFact -} from '../../../shared/agent-session-failure' +import type { SubmissionRejectionFact } from '../../../shared/agent-session-failure' import type { StructuredAgentSessionStopCause } from './structured-agent-session-stop-cause' -import type { StructuredAgentSessionAdapterStop } from './structured-agent-session-adapter-stop' +import type { + AgentSessionCancelOutcome, + StructuredAgentSessionAdapterStop +} from './structured-agent-session-adapter-stop' +export type { AgentSessionCancelOutcome } from './structured-agent-session-adapter-stop' export type { StructuredAgentSessionChildEndCause, StructuredAgentSessionStopCause @@ -242,12 +243,6 @@ export type StructuredAgentSessionSetOptionInput = { fence: number } -/** `refusal`: the provider answered the Stop and declined it, in its own words when it gave any. */ -export type AgentSessionCancelOutcome = { - cancelled: boolean - refusal?: { detail?: ProviderDiagnostic } -} - export type StructuredAgentSessionAdapter = StructuredAgentSessionAdapterStop & { /** Provider-aware capability check for hosts that route more than one adapter. */ supportsCreate?(location: AgentSessionExecutionLocation, agent: string): boolean diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-chat-stop.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-chat-stop.ts index 70707019295..b6224728b96 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-chat-stop.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-chat-stop.ts @@ -1,7 +1,9 @@ // The chat's Stop, however a client reached it: the Stop button or a question card's Cancel. One // body and one order: withdraw what is queued, record where the Stop took effect, interrupt, then -// end the child in the next step on the session's lane. The body is reachable only through -// `mutateWithChatStop`, which queues that step in the same synchronous call as the mutation. +// end the child: in this step when the interrupt failed with the turn still running, else in the +// next step on the session's lane for a provider whose Stop ends its session. The body is +// reachable only through `mutateWithChatStop`, which queues that step in the same synchronous call +// as the mutation. import { agentSessionFailureWords } from '../../../shared/agent-session-failure-words' import { agentSessionFailureFact } from '../../../shared/agent-session-failure' @@ -83,6 +85,9 @@ export function mutateWithChatStop( clientOperationId: envelope.clientOperationId, ...named, stopChild: () => context.stopAgent(sessionId), + onStopChildError: (error) => context.deps.onEventSinkError?.({ sessionId, error }), + // The host drops its child only once the exit is proven, and nothing else runs meanwhile. + childReleased: () => context.sessions.get(sessionId)?.child !== child, endSession: (owed) => { windDown = owed }, diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-codex-stop-row.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-codex-stop-row.test.ts index 9103393c26c..e93b44a6b08 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-codex-stop-row.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-codex-stop-row.test.ts @@ -6,8 +6,10 @@ import { openTestJournalHostDatabase } from '../agent-session-journal/journal-ho import { mkdtemp, rm } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' -import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { afterEach, beforeEach, describe, expect, it, vi, type MockInstance } from 'vitest' import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope' +import { CodexAppServerRequestError } from '../../codex/codex-app-server-connection' +import { CodexAppServerTimeoutError } from '../../codex/codex-app-server-session' import { THREAD_ID as THREAD, adapterFor, @@ -15,6 +17,7 @@ import { } from '../../codex/codex-structured-session-adapter-fixture' import { codexTurnLifecycleFake } from '../../codex/codex-turn-lifecycle-fake' import { openTestAgentSessionRecordStore } from '../../runtime/agent-session-record-store-test-harness' +import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter' import { StructuredAgentSessionHost } from './structured-agent-session-host' import { HOST_TEST_NOW as NOW, @@ -31,13 +34,15 @@ const CALLER = { callerKey: 'client-1' } let root: string let host: StructuredAgentSessionHost let turns: ReturnType +let codex: ReturnType +let notify: (method: string, params: unknown) => void +let disposeSession: MockInstance> beforeEach(async () => { root = await mkdtemp(join(tmpdir(), 'orca-codex-stop-row-')) resetHostTestOperationIds() - const codex = fakeCodex() - const notify = (method: string, params: unknown): void => - codex.connections.at(-1)?.handlers.onNotification?.(method, params) + codex = fakeCodex() + notify = (method, params) => codex.connections.at(-1)?.handlers.onNotification?.(method, params) turns = codexTurnLifecycleFake(THREAD, () => notify) codex.routes['turn/start'] = turns.routes['turn/start'] codex.routes['turn/interrupt'] = () => { @@ -51,10 +56,20 @@ beforeEach(async () => { return {} } const store = await openTestAgentSessionRecordStore(root) + // The runtime's wiring: an echo accepts its send, and an exit reaches the host. + const adapter = adapterFor(codex, {}, [], { + onDispatchSettledLate: (settlement) => void host.settleLateDispatch(settlement), + onEvent: (event) => { + if (event.type === 'ended' && 'cause' in event && event.cause === 'unexpected-exit') { + void host.handleAdapterEvent(event) + } + } + }) + disposeSession = vi.spyOn(adapter, 'disposeSession') host = new StructuredAgentSessionHost({ logger: createStructuredAgentSessionLogger(), store, - adapter: Object.assign(adapterFor(codex), { supportsCreate: () => true }), + adapter: Object.assign(adapter, { supportsCreate: () => true }), journalDatabase: openTestJournalHostDatabase(root), claimKeyId: 'key-1', mintSpawnToken: () => 'spawn-1', @@ -103,9 +118,13 @@ function stop(turnId?: string) { } async function runningTurn(): Promise { - expect(await send('count to 40')).toMatchObject({ ok: true }) + const sent = await send('count to 40') + if (!sent.ok) { + throw new Error(JSON.stringify(sent.refusal)) + } await vi.waitFor(() => expect(turns.turnId).toBe('turn-1')) turns.start() + turns.echo(sent.value.clientMessageId) await vi.waitFor(async () => expect((await journalRows()).turns).toEqual(['running'])) } @@ -113,10 +132,59 @@ async function journalRows() { const items = (await host.journalSnapshot(SESSION)).items return { statuses: items.flatMap((item) => (item.body.kind === 'status' ? [item.body.text] : [])), - turns: items.flatMap((item) => (item.body.kind === 'turn' ? [item.body.state] : [])) + turns: items.flatMap((item) => (item.body.kind === 'turn' ? [item.body.state] : [])), + outcomes: items.flatMap((item) => (item.body.kind === 'turn' ? [item.body.outcome] : [])) } } +/** Codex's -32600 states the named turn is not running; its -32603 is an interrupt it could not + * submit, the turn still running. */ +function refused(message: string, code: -32600 | -32603 = -32600): CodexAppServerRequestError { + return new CodexAppServerRequestError( + 'turn/interrupt', + code, + `codex app-server turn/interrupt failed: ${message}`, + message + ) +} + +function interruptFailure(failure: 'internal error' | 'unanswered'): Error { + return failure === 'internal error' + ? refused('failed to interrupt turn: channel closed', -32603) + : new CodexAppServerTimeoutError('codex app-server turn/interrupt exceeded 30000ms') +} + +/** The journal's writes wait a moment, so frames Orca received are not yet in the journal. */ +function holdJournalWrites(): void { + const journal = host['sessions'].get(SESSION)!.journal + const released = new Promise((resolve) => setTimeout(resolve, 20)) + const appendItem = journal.appendItem.bind(journal) + const appendLifecycleBatch = journal.appendLifecycleBatch.bind(journal) + vi.spyOn(journal, 'appendItem').mockImplementation(async (...args) => { + await released + return appendItem(...args) + }) + vi.spyOn(journal, 'appendLifecycleBatch').mockImplementation(async (...args) => { + await released + return appendLifecycleBatch(...args) + }) +} + +/** Codex picked the follow-up's turn and answered the send, and has not started it. */ +async function followUpUnopened(): Promise { + const sent = await send('and then this') + if (!sent.ok) { + throw new Error(JSON.stringify(sent.refusal)) + } + await vi.waitFor(() => expect(turns.turnId).toBe('turn-2')) + await host.flushStreamedEvents(SESSION) +} + +/** Whether the Stop ended the child: the host's stop, which proves the exit, with the user's cause. */ +function childEndedByStop(): boolean { + return disposeSession.mock.calls.some(([, cause]) => cause === 'user-stop') +} + describe('a Codex Stop that Codex answered', () => { it.each([ ['names no turn', undefined], @@ -132,6 +200,198 @@ describe('a Codex Stop that Codex answered', () => { expect((await journalRows()).statuses).toEqual(['Cancellation requested.']) expect(stopped).toMatchObject({ ok: true, value: { cancelled: true } }) + // Codex's background terminals live in its child, so a Stop it took keeps the child. + expect(childEndedByStop()).toBe(false) + expect(codex.connections.at(-1)?.closed).toBe(false) } ) }) + +describe('a Codex Stop whose interrupt failed', () => { + it.each([ + ['Codex could not submit it, naming no turn', undefined, 'internal error'], + ['Codex could not submit it, naming its turn', 'turn-1', 'internal error'], + ['unanswered, naming no turn', undefined, 'unanswered'], + ['unanswered, naming its turn', 'turn-1', 'unanswered'] + ] as const)( + "ends the child still running the turn, as the user's cancellation, when %s", + async (_case, turnId, failure) => { + await runningTurn() + codex.routes['turn/interrupt'] = () => { + throw interruptFailure(failure) + } + + const stopped = await stop(turnId) + await host.flushStreamedEvents(SESSION) + + expect(stopped).toMatchObject({ ok: true, value: { cancelled: true } }) + expect(disposeSession).toHaveBeenCalledExactlyOnceWith(SESSION, 'user-stop') + expect(codex.connections.at(-1)?.closed).toBe(true) + const rows = await journalRows() + expect(rows.turns).toEqual(['interrupted']) + expect(rows.outcomes).toEqual(['cancellation']) + expect(rows.statuses).toEqual(['Cancellation requested.']) + } + ) + + it('says Codex did not stop when a Stop naming the running turn could not prove the exit', async () => { + await runningTurn() + codex.routes['turn/interrupt'] = () => { + throw interruptFailure('internal error') + } + disposeSession.mockResolvedValueOnce(false) + + const stopped = await stop('turn-1') + await host.flushStreamedEvents(SESSION) + + expect(stopped).toMatchObject({ ok: true, value: { cancelled: false } }) + expect(disposeSession).toHaveBeenCalledExactlyOnceWith(SESSION, 'user-stop') + expect((await journalRows()).statuses).toEqual([ + "Codex didn't stop: failed to interrupt turn: channel closed." + ]) + }) + + it.each([ + ['naming no turn', undefined], + ['naming its turn', 'turn-1'] + ] as const)( + 'keeps the child when Codex has moved on to another turn, %s', + async (_case, turnId) => { + await runningTurn() + // Codex ended the turn and opened another before the interrupt reached it. + codex.routes['turn/interrupt'] = () => { + turns.end('completed') + notify('turn/started', { threadId: THREAD, turn: { id: 'turn-2', status: 'inProgress' } }) + throw refused('expected active turn id turn-1 but found turn-2') + } + + const stopped = await stop(turnId) + await host.flushStreamedEvents(SESSION) + + expect(stopped).toMatchObject({ ok: true, value: { cancelled: false } }) + expect(childEndedByStop()).toBe(false) + expect(codex.connections.at(-1)?.closed).toBe(false) + expect((await journalRows()).turns).toEqual(['completed', 'running']) + } + ) + + // Codex marks the turn ended before it writes the frame, so its refusal can come first. + it.each([ + ['naming no turn', undefined], + ['naming its turn', 'turn-1'] + ] as const)( + "keeps the child when Codex's refusal arrives before the turn's end, %s", + async (_case, turnId) => { + await runningTurn() + codex.routes['turn/interrupt'] = () => { + setTimeout(() => turns.end('completed'), 2) + throw refused('no active turn to interrupt') + } + + const stopped = await stop(turnId) + await vi.waitFor(async () => expect((await journalRows()).turns).toEqual(['completed'])) + await host.flushStreamedEvents(SESSION) + + expect(stopped).toMatchObject({ ok: true, value: { cancelled: false } }) + expect(childEndedByStop()).toBe(false) + expect(codex.connections.at(-1)?.closed).toBe(false) + expect((await journalRows()).outcomes).toEqual(['success']) + } + ) + + it.each([ + ['naming no turn', undefined], + ['naming its turn', 'turn-1'] + ] as const)( + 'keeps the child when the turn ended before the interrupt reached it, %s', + async (_case, turnId) => { + await runningTurn() + codex.routes['turn/interrupt'] = (params) => { + turns.end('completed') + return turns.routes['turn/interrupt'](params) + } + + const stopped = await stop(turnId) + await host.flushStreamedEvents(SESSION) + + expect(stopped).toMatchObject({ ok: true, value: { cancelled: false } }) + expect(childEndedByStop()).toBe(false) + expect(codex.connections.at(-1)?.closed).toBe(false) + expect((await journalRows()).turns).toEqual(['completed']) + } + ) + + it('keeps the child at rest, and writes no row, when a Stop names a turn that already ended', async () => { + await runningTurn() + turns.end('completed') + await vi.waitFor(async () => expect((await journalRows()).turns).toEqual(['completed'])) + codex.routes['turn/interrupt'] = turns.routes['turn/interrupt'] + + const stopped = await stop('turn-1') + await host.flushStreamedEvents(SESSION) + + expect(stopped).toMatchObject({ ok: true, value: { cancelled: false } }) + expect(childEndedByStop()).toBe(false) + expect(codex.connections.at(-1)?.closed).toBe(false) + expect((await journalRows()).statuses).toEqual([]) + }) + + // A phone names the turn it shows; a follow-up from elsewhere is not that Stop's to end. + it.each([ + ['Codex refused it', 'refused'], + ['the interrupt went unanswered', 'unanswered'] + ] as const)( + 'keeps the child when a Stop names a turn that ended and a follow-up has not opened, %s', + async (_case, failure) => { + await runningTurn() + turns.end('completed') + await vi.waitFor(async () => expect((await journalRows()).turns).toEqual(['completed'])) + await followUpUnopened() + codex.routes['turn/interrupt'] = + failure === 'refused' + ? turns.routes['turn/interrupt'] + : () => { + throw interruptFailure('unanswered') + } + + const stopped = await stop('turn-1') + await host.flushStreamedEvents(SESSION) + + expect(stopped).toMatchObject({ ok: true, value: { cancelled: false } }) + expect(childEndedByStop()).toBe(false) + expect(codex.connections.at(-1)?.closed).toBe(false) + } + ) + + it("decides on Codex's frames received before the interrupt failed, not on the journal's last write", async () => { + await runningTurn() + codex.routes['turn/interrupt'] = () => { + holdJournalWrites() + turns.end('completed') + notify('turn/started', { threadId: THREAD, turn: { id: 'turn-2', status: 'inProgress' } }) + throw interruptFailure('internal error') + } + + const stopped = await stop('turn-1') + await host.flushStreamedEvents(SESSION) + + expect(stopped).toMatchObject({ ok: true, value: { cancelled: false } }) + expect(childEndedByStop()).toBe(false) + expect((await journalRows()).turns).toEqual(['completed', 'running']) + }) + + it('leaves a child that exited during the interrupt to its exit', async () => { + await runningTurn() + codex.routes['turn/interrupt'] = () => { + const exited = new Error('codex app-server exited') + codex.connections.at(-1)?.handlers.onExit?.(exited) + throw exited + } + + const stopped = await stop() + await host.flushStreamedEvents(SESSION) + + expect(stopped).toMatchObject({ ok: true, value: { cancelled: false } }) + expect(childEndedByStop()).toBe(false) + }) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-stop.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-stop.test.ts index 0f7c471ed0a..a5af4b851ba 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-stop.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-stop.test.ts @@ -37,8 +37,12 @@ let dispatch: Mock let cancelTurn: Mock let awaitStarted: Mock> let closeSession: Mock> +let acknowledgeSessionRelease: Mock< + NonNullable +> /** Codex's answer by default: its Stop keeps the child. */ let stopEndsSession: boolean +let onEventSinkError: Mock<(failure: { sessionId: string; error: unknown }) => void> let events: StructuredAgentSessionEventSink | undefined function eventually(assertion: () => void | Promise): Promise { @@ -53,7 +57,9 @@ beforeEach(async () => { cancelTurn = vi.fn(async () => ({ cancelled: true })) awaitStarted = vi.fn(async () => undefined) closeSession = vi.fn(async () => true) + acknowledgeSessionRelease = vi.fn() stopEndsSession = false + onEventSinkError = vi.fn() store = await openTestAgentSessionRecordStore(root) host = new StructuredAgentSessionHost({ logger: createStructuredAgentSessionLogger(), @@ -81,6 +87,7 @@ beforeEach(async () => { dispatch, awaitStarted, closeSession, + acknowledgeSessionRelease, releaseAcquisition: vi.fn(async () => true), cancelTurn, stopEndsSession: () => stopEndsSession, @@ -90,6 +97,7 @@ beforeEach(async () => { journalDatabase: openTestJournalHostDatabase(root), claimKeyId: 'key-1', mintSpawnToken: () => 'spawn-1', + onEventSinkError, now: () => NOW }) expect(await host.attach(CALLER, hostTestAttachParams(null))).toMatchObject({ ok: true }) @@ -221,28 +229,92 @@ describe('a Stop that names no turn', () => { expect(dispatch).not.toHaveBeenCalled() }) - it('says the agent did not stop, in its words, when it refused', async () => { + it('ends the child when the provider could not interrupt the message still in flight', async () => { const { id, result } = send('hello') await result await eventually(async () => expect((await submission(id))?.handedOverAt).toBeDefined()) cancelTurn.mockResolvedValueOnce({ cancelled: false, - refusal: { detail: { text: 'no active turn to interrupt', audience: 'person' } } + refusal: { detail: { text: 'failed to interrupt turn', audience: 'person' } } + }) + + expect(await stop()).toMatchObject({ ok: true, value: { cancelled: true } }) + + expect(closeSession).toHaveBeenCalledExactlyOnceWith(SESSION, 'user-stop') + expect(await statusRows()).toEqual(['Cancellation requested.']) + }) + + it('says the agent did not stop, in its words, when it refused because its turn is not running', async () => { + const { id, result } = send('hello') + await result + await eventually(async () => expect((await submission(id))?.handedOverAt).toBeDefined()) + cancelTurn.mockResolvedValueOnce({ + cancelled: false, + refusal: { + detail: { text: 'no active turn to interrupt', audience: 'person' }, + turnNotRunning: true + } }) expect(await stop()).toMatchObject({ ok: true, value: { cancelled: false } }) - expect(cancelTurn).toHaveBeenCalledOnce() + expect(closeSession).not.toHaveBeenCalled() expect(await statusRows()).toEqual(["Codex didn't stop: no active turn to interrupt."]) }) - it('says the Stop is unconfirmed, not that nothing ran, when the provider never answered it', async () => { + it('ends the child when the provider never answered the interrupt', async () => { const { id, result } = send('hello') await result await eventually(async () => expect((await submission(id))?.handedOverAt).toBeDefined()) // Codex answers an interrupt as the turn ends, so a turn that never ends leaves it unanswered. cancelTurn.mockRejectedValueOnce(new Error('codex app-server turn/interrupt exceeded 30000ms')) + expect(await stop()).toMatchObject({ ok: true, value: { cancelled: true } }) + + expect(closeSession).toHaveBeenCalledExactlyOnceWith(SESSION, 'user-stop') + expect(await statusRows()).toEqual(['Cancellation requested.']) + }) + + it('says the agent did not stop, in its words, when it could not interrupt and the child end is unproven', async () => { + const { id, result } = send('hello') + await result + await eventually(async () => expect((await submission(id))?.handedOverAt).toBeDefined()) + cancelTurn.mockResolvedValueOnce({ + cancelled: false, + refusal: { detail: { text: 'failed to interrupt turn', audience: 'person' } } + }) + closeSession.mockResolvedValueOnce(false) + + expect(await stop()).toMatchObject({ ok: true, value: { cancelled: false } }) + + expect(closeSession).toHaveBeenCalledExactlyOnceWith(SESSION, 'user-stop') + expect(onEventSinkError).toHaveBeenCalledWith(expect.objectContaining({ sessionId: SESSION })) + expect(await statusRows()).toEqual(["Codex didn't stop: failed to interrupt turn."]) + }) + + it('reads as requested when the child exit was proven and only a later cleanup step failed', async () => { + const { id, result } = send('hello') + await result + await eventually(async () => expect((await submission(id))?.handedOverAt).toBeDefined()) + cancelTurn.mockRejectedValueOnce(new Error('codex app-server turn/interrupt exceeded 30000ms')) + acknowledgeSessionRelease.mockImplementationOnce(() => { + throw new Error('route release failed') + }) + + expect(await stop()).toMatchObject({ ok: true, value: { cancelled: true } }) + + expect(closeSession).toHaveBeenCalledExactlyOnceWith(SESSION, 'user-stop') + expect(onEventSinkError).toHaveBeenCalledWith(expect.objectContaining({ sessionId: SESSION })) + expect(await statusRows()).toEqual(['Cancellation requested.']) + }) + + it('says the Stop is unconfirmed, not that nothing ran, when neither the interrupt nor the child end is proven', async () => { + const { id, result } = send('hello') + await result + await eventually(async () => expect((await submission(id))?.handedOverAt).toBeDefined()) + cancelTurn.mockRejectedValueOnce(new Error('codex app-server turn/interrupt exceeded 30000ms')) + closeSession.mockResolvedValueOnce(false) + expect(await stop()).toMatchObject({ ok: true, value: { cancelled: false } }) expect(await statusRows()).toEqual(['Cancellation was not confirmed.']) @@ -337,6 +409,18 @@ describe('a Stop that names its turn, as an older client sends it', () => { expect(await statusRows()).toEqual([]) }) + it('keeps the child, and writes no row, when the provider could not interrupt it with the conversation at rest', async () => { + cancelTurn.mockResolvedValueOnce({ + cancelled: false, + refusal: { detail: { text: 'failed to interrupt turn', audience: 'person' } } + }) + + expect(await stop('turn-1')).toMatchObject({ ok: true, value: { cancelled: false } }) + + expect(closeSession).not.toHaveBeenCalled() + expect(await statusRows()).toEqual([]) + }) + async function queueOnHost(): Promise<{ id: string; release: () => void }> { const started = Promise.withResolvers() awaitStarted.mockImplementationOnce(() => started.promise) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts index d4adf43853b..57aa06299fc 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts @@ -39,6 +39,37 @@ export async function isMainAgentWorkingOnceFlushed( ) } +/** + * After an interrupt that failed: whether the session still runs what the Stop was sent for. One at + * rest, or running a different turn, is not the Stop's to end. A child that exited reads at rest: + * its exit ends its turn, and the host's own exit handling waits behind this step. + */ +async function stillRunsStoppedTurn( + ctx: Pick, + stoppedTurnId: string | null +): Promise { + if (!(await isMainAgentWorkingOnceFlushed(ctx))) { + return false + } + // Working with no turn open after the Stop's turn is a later send whose turn has not opened. + return stoppedTurnId === null || ctx.journal.activeTurnId() === stoppedTurnId +} + +/** The row for a Stop the provider declined, in its words when it gave any. */ +function stopRefusedNote( + ctx: Pick, + refusal: AgentSessionCancelOutcome['refusal'] +): AgentJournalStatusItem { + const detail = refusal?.detail + return { + kind: 'status', + ...agentSessionFailureWords(agentSessionFailureFact('stopRefused', detail ? { detail } : {}), { + ...ctx.failureTextContext, + surface: 'row' + }) + } +} + /** What a Stop that ends the provider's session leaves its next serialized step: whether the * provider took the interrupt, so its wind-down is worth waiting on, and when the interrupt went * out. No turn id: that step runs right behind the Stop, so no later turn can slip in between. */ @@ -75,8 +106,14 @@ export async function performCancel( scope?: 'background-tasks' taskId?: string prompt?: { itemId: string; expectedRevision: number } - /** Ends the provider child, for a running command the provider did not take the Stop on. */ + /** Ends the provider child, for a running command the provider did not take the Stop on, or a + * turn whose interrupt failed. */ stopChild?: () => Promise + /** A child end after a failed interrupt that threw. */ + onStopChildError?: (error: unknown) => void + /** After that throw: whether the host let go of the child, its exit proven before a later + * cleanup step failed. */ + childReleased?: () => boolean /** Hands the child's end to the Stop's next serialized step, for a provider whose Stop ends * its session. */ endSession?: (windDown: StructuredAgentSessionStopWindDown) => void @@ -119,6 +156,9 @@ export async function performCancel( const stoppedAt = Date.now() // The provider's own answer; unset when its cancel threw, leaving the effect unknown. let taken: boolean | undefined + // The provider could not interrupt the turn, or its cancel threw: the turn may run on. + let interruptFailed = false + let refusal: AgentSessionCancelOutcome['refusal'] try { const dispatchStatus = latestJournalDispatchObservation(ctx.journal, ctx.fence) const outcome: AgentSessionCancelOutcome = stoppedBefore @@ -145,6 +185,8 @@ export async function performCancel( }) taken = outcome.cancelled cancelled = outcome.cancelled + refusal = outcome.refusal + interruptFailed = refusal !== undefined && refusal.turnNotRunning !== true if (!cancelled && input.withdrewQueued && !(await isMainAgentWorkingOnceFlushed(ctx))) { // A Stop that withdrew what was queued and left nothing working ended what it was sent for, // named or not. The journal judges it: providers differ on refusing a turn that has ended. @@ -154,19 +196,13 @@ export async function performCancel( note = null } else if (!cancelled && input.turnId === undefined) { // Sent only while the chat reads working, so a Stop that ended nothing must say why. - const detail = outcome.refusal?.detail - note = { - kind: 'status', - ...agentSessionFailureWords( - agentSessionFailureFact('stopRefused', detail ? { detail } : {}), - { ...ctx.failureTextContext, surface: 'row' } - ) - } + note = stopRefusedNote(ctx, refusal) } } catch (error) { if (input.prompt) { throw error } + interruptFailed = true // The adapter's error is Orca's; the row says only that the stop is unconfirmed. note = { kind: 'status', @@ -191,8 +227,31 @@ export async function performCancel( await input.stopChild?.() cancelled = true note = { kind: 'status', text: STOP_NOTE_CANCELLATION_REQUESTED } - } - if (!cancelled && taken === false && input.turnId !== undefined) { + } else if ( + !cancelled && + interruptFailed && + input.stopChild && + // An unnamed Stop meant the turn the journal showed when it was sent. + (await stillRunsStoppedTurn(ctx, input.turnId ?? liveTurnId)) + ) { + // The interrupt failed and the turn runs on: only the child's end stops it. + let ended: boolean + try { + await input.stopChild() + ended = true + } catch (error) { + input.onStopChildError?.(error) + // Unless the exit was proven, the child may still run the turn and the failed row stays true. + ended = input.childReleased?.() === true + } + if (ended) { + cancelled = true + note = { kind: 'status', text: STOP_NOTE_CANCELLATION_REQUESTED } + } else if (taken !== undefined) { + // The turn was just read running, so a named Stop says it was refused rather than nothing. + note = stopRefusedNote(ctx, refusal) + } + } else if (!cancelled && taken === false && input.turnId !== undefined) { // Nothing was left of the turn it named and nothing else ended: a Stop that ends nothing writes no row. note = null }