fix(native-chat): a Codex Stop that Codex refuses or never answers ends the Codex process (#24334)

* fix(codex): a Stop whose interrupt Codex refused or never answered ends the session

A Codex Stop is the interrupt alone, so its background terminals keep running. When Codex refused
the interrupt, or never answered it, the turn kept running and the chat said "Codex didn't stop"
or "Cancellation was not confirmed." with no way to stop it short of closing the chat.

Now, after an interrupt that failed, the Stop re-reads the conversation once Codex's frames already
received have landed, and ends the child through the host's usual stop (proof of exit, then the
lease release) when it still runs what the Stop was sent for. It skips a conversation at rest, a
child whose exit already ended its turn, and one Codex has moved on to a different turn. The turn
then reads as the user's cancellation. If the child's exit is not proven, the failure is reported
and the Stop keeps its "didn't stop" or "not confirmed" row, which is still true.

A Stop Codex took keeps the child, unchanged.

* fix(codex): end the child only for an interrupt whose effect is unknown, on the turn the Stop meant

- Codex's invalid-request refusal (-32600: no active turn, another turn active, thread not
  loaded) states the named turn is not running, and can arrive before that turn's end frame.
  The Codex adapter now marks it `turnNotRunning`, and such a refusal never ends the child.
  Only an internal-error refusal (-32603), an unanswered interrupt, or a thrown cancel does.
- A Stop that named a turn, or an unnamed Stop that read one, ends the child only when the
  journal, after draining received frames, still shows that same turn. A session working on
  a follow-up whose turn has not opened is no longer ended by a Stop of the finished turn.
- When the child's exit was proven and only a later cleanup step failed, the Stop reads as
  requested; the failure is still reported.
- `AgentSessionCancelOutcome` moves beside the adapter's other Stop members (re-exported), which
  keeps the adapter file within its line budget.

* test(codex): pin that a failed interrupt decides on the drained journal

The frames Codex sent before the interrupt failed are held in the sink, so a read that skips the drain sees the stopped turn still running.

* test(codex): a refusal naming a turn Codex is not running carries turnNotRunning

* test(native-chat): fail the route release synchronously, as the adapter's acknowledgement is

* docs(codex): a -32600 Codex could not parse also reads as not running and keeps the child

* fix(native-chat): a named Stop whose child end is unproven says Codex didn't stop, not that the turn finished

The new branch has just read the turn running, so the named Stop's 'already finished' row was false.
This commit is contained in:
Brennan Benson
2026-10-01 11:27:53 -07:00
committed by GitHub
parent 34a582bd39
commit beccdec74c
9 changed files with 480 additions and 44 deletions
@@ -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'])
@@ -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<typeof fakeCodex>,
launch: Partial<CodexStructuredLaunch> = {},
events: CodexStructuredSessionEvent[] = []
events: CodexStructuredSessionEvent[] = [],
/** Host wiring the runtime adds, such as the late dispatch settlement. */
deps: Partial<CodexStructuredSessionAdapterDeps> = {}
): 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
})
}
@@ -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) {
@@ -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
@@ -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
@@ -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<TValue>(
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
},
@@ -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<typeof codexTurnLifecycleFake>
let codex: ReturnType<typeof fakeCodex>
let notify: (method: string, params: unknown) => void
let disposeSession: MockInstance<NonNullable<StructuredAgentSessionAdapter['disposeSession']>>
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<void> {
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<void> {
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)
})
})
@@ -37,8 +37,12 @@ let dispatch: Mock<StructuredAgentSessionAdapter['dispatch']>
let cancelTurn: Mock<StructuredAgentSessionAdapter['cancelTurn']>
let awaitStarted: Mock<NonNullable<StructuredAgentSessionAdapter['awaitStarted']>>
let closeSession: Mock<NonNullable<StructuredAgentSessionAdapter['closeSession']>>
let acknowledgeSessionRelease: Mock<
NonNullable<StructuredAgentSessionAdapter['acknowledgeSessionRelease']>
>
/** 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<void>): Promise<void> {
@@ -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<undefined>()
awaitStarted.mockImplementationOnce(() => started.promise)
@@ -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<AgentSessionTurnContext, 'journal' | 'fence' | 'flushStreamedEvents'>,
stoppedTurnId: string | null
): Promise<boolean> {
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<AgentSessionTurnContext, 'failureTextContext'>,
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<void>
/** 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
}