fix(claude): a Stop ends Claude's process, and the next message resumes the conversation (#24235)

* fix(claude): a Stop ends Claude's process, and the next message resumes the conversation

Claude's Stop now ends its child after the interrupt, whatever Claude answered, once the stopped turn ends or a 3 s grace runs out. The grace lets Claude's own result move the resume point past that turn, so the next send resumes after it. Background work ends with the child. Codex keeps its child.

* test(native-chat): the router answers stopEndsSession for the session's live owner

* test(native-chat): type the Stop test's close mock as the adapter's optional method

* fix(native-chat): a Claude Stop answers on its interrupt, and the next step on the chat's lane ends the child

The Stop now answers as soon as the interrupt step finishes. Ending Claude's
process runs as a second task on the session's serialized lane, queued in the
same tick as the Stop, so a message sent meanwhile reaches only the resumed
child. That step waits for the stopped turn to end (its result, the CLI's idle,
the child's exit or its close), then ends the child; a failure there is
reported and never fails the Stop, and the wind-down it leaves owed is retried.

The interrupt's answer and the turn's end share one 3 s grace counted from when
the interrupt goes out, so an interrupt Claude never answers ends the child in
about 3 s instead of after the 30 s control deadline. The wait is derived per
call from Claude's open turn: the armed latch and the close's user-stop branch
are gone, and the close keeps its whole deadline for proving the exit.

A Stop naming a turn that has since ended (the phone names the turn it last
saw) now ends the child too when Claude took it, as it does when Claude
interrupts a handed-over follow-up whose turn has not opened.

* test(native-chat): pin a Claude Stop's race, failure and timing cases on the shipping adapter

- Nothing queued while the Stop's first step waits on the interrupt runs before
  the child ends.
- A message sent with the pre-Stop fence during the Stop is accepted with no
  notice, and reaches only the resumed child, after the old child's close.
- An interrupt Claude never answers ends the child within the grace.
- A close that cannot prove the exit leaves the Stop answered with its row; the
  failure goes to the host's error hook.
- A Stop naming the turn that just ended ends the child when Claude interrupts
  the handed-over follow-up.
- A second Stop pressed while the first ends the child stays quiet.

* fix(native-chat): a Stop naming an ended turn ends Claude's child when its interrupt fails or goes unanswered

The phone names the turn it last saw. When Claude interrupted a handed-over
follow-up for that Stop and the answer failed or never came, the child stayed.
Only a provider that answers it did not take the Stop keeps the child now.
Also pins a queue-if-active send made while the Stop ends the child.

* fix(native-chat): a Claude card's Cancel denies an approval and stops on a question, never a bare interrupt

A permission card's Cancel interrupted the turn and kept Claude running, the
old Stop on a second control: a refused interrupt let the turn run on and
background work survived. Now the host asks the provider how a card's Cancel
is answered. For Claude, an approval's Cancel sends the same reply as its Deny
option, so the turn goes on without that tool; a question's Cancel runs the
chat's Stop, which ends the child, and the next send resumes the conversation.
A provider that gives no answer keeps today's path, so Codex is unchanged.

The host decides, so desktops and phones of any version keep sending the same
cancel and get the new behaviour; nothing on the wire changes.

With no Claude caller left, the prompt-cancel interrupt is gone: the bound
claim, the cancellation observation, the admission of the prompt's
cancellation, the prompt's turn binding, and withdrawing a refused stop (every
Claude interrupt now precedes the child's end, so the recorded stop stands).

* fix(native-chat): a Claude approval card's Stop option is the chat's Stop

The approval card's "Stop" option answered Claude with a deny that interrupts
the turn and kept Claude running, the last control on a Claude chat that
interrupted without ending the child. The host now routes it, like the card's
own Cancel, through the provider: for Claude that option is the chat's Stop.
It runs the same Stop body as the Stop button, in the same order (withdraw,
the Stop takes effect, interrupt), and the step that ends the child is queued
in the same synchronous call. The body now lives in one place,
structured-agent-session-chat-stop.ts, shared by the Stop button, a question
card's Cancel and this option.

The host decides, so an old card or client that sends option `cancel` gets the
real Stop; nothing on the wire changes. The interrupting deny reply is gone.

* test(native-chat): the option-route test's answer carries a whole journal resolution

* fix(native-chat): Claude cards drop their Stop option, and a plan card's Cancel asks Claude to wait

- Approval and plan cards no longer offer "Stop". Neither reference offers a
  Stop while an approval is up: the user denies, then stops. An older card's or
  client's `cancel` answer is a plain deny, and the respond path is main's again.
- A plan card's X / Esc no longer answers "Keep planning", which sent Claude
  straight back to revise and re-propose. It dismisses the plan: a deny that
  tells Claude to end its turn and wait for the user. No interrupt; the child
  stays. A tool approval's Cancel stays the Deny reply, and a question's Cancel
  stays the chat's Stop.
- The provider's routing takes the pending card itself
  (`routePromptCancel({ sessionId, prompt })`), so a plan is told from a tool
  approval, instead of a kind plus an optional option id that meant Cancel when
  absent. A routed answer may be one the card does not offer, like the
  dismissal.
- The chat's Stop is reachable only through `mutateWithChatStop`, which owns the
  mutation call and queues the step that ends the child in the same call, so
  no caller can run the Stop and forget the child's end.

* refactor(native-chat): cancelPlan no longer takes a stopChild nothing passes

The Stop's own body ends the child; the plan's default run serves only a card's
interrupt and the background-task stop, which never end it.

* fix(native-chat): a question card's Cancel settles the card in the Stop's first step, and answers Claude when nothing stops

A question card's Cancel runs the chat's Stop. The card stayed pending until
the child ended (up to about 3 s), so it could still be pressed and the second
press was refused once the chat rested. When the Stop found nothing in flight
(Claude asking after its main turn ended), the card stayed pending for good and
Claude's request went unanswered; the phone froze the card with every button
disabled.

Now the card is recorded as cancelled by the user in the Stop's own step, with
the receipt "Cancelled on <device>". A Stop that ends the session leaves
Claude's request to end with the child, so no reply races the interrupt; one
that ends nothing declines the request itself. The provider forgets the card
once the host records it, so neither Claude's own cancel nor the child's end
writes over the user's cancel. The card's Stop stops whatever the chat has in
flight, like the Stop button, instead of a turn the card names that may have
ended.

The same dismissal, a new provider member beside the cancel route, now answers
a plan card's Cancel: recorded as cancelled by the user, and Claude told to end
its turn and wait. That replaces resolving it as an option the card does not
offer.

* test(native-chat): the question-card Stop test has Claude cancel its held request, as it does once interrupted

* fix(native-chat): a Claude Stop before the echo reads Interrupted, and a question card's Stop stays in its own turn

- A Stop pressed after Send but before Claude echoed it read as a normal finish
  ("Worked for 0s", a green done check): with no turn row to wait on, the child
  was ended at once, and the echo that would have opened the stopped turn
  landed on a retired send. The wait before the child ends is now keyed on what
  Claude has in flight (an open turn or a send it has not answered), not on a
  journal turn id, within the same 3 s from the interrupt and woken by the same
  settles; a settle that leaves something in flight waits on. The echo opens
  the turn, the aborted result ends it Interrupted. A Claude that says nothing
  is ended when the grace runs out, the unanswered send doubt as before.
- A question card's Cancel takes the chat's Stop only when the card belongs to
  the turn running now, judged after draining the provider's lifecycle. A card
  a finished turn raised, such as a background agent's, is dismissed instead:
  Claude's request is declined and nothing stops.
- Claude forgets a dismissed card before the host records it, so Claude's own
  cancel landing during that write cannot replace the user's "Cancelled on
  <device>".

* test(native-chat): Stop tests read background work from the host's child records, and guard the optional dismissal

Main's child records replaced the adapter's background-task callback: the
background-work test now reads the host's child record (live before the Stop,
settled after the child ends) and the session's agent status (done, Interrupted,
nothing left for Monitoring).

* fix(native-chat): a tool card's Cancel reads Cancelled, and a Stop's background work reads stopped

- A tool approval card's X / Esc already sent Claude the plain deny, but the
  card read "Deny · Answered on <device>", the same as pressing Deny. It is
  now a dismissal like a plan card's: recorded as cancelled by the user, with
  the same deny reply, no interrupt and the child kept. Every card's Cancel
  now routes to a dismissal or the chat's Stop, so the option route is gone.
- When a Stop ended Claude's child, a background agent or shell it was running
  settled as an unknown ending (a neutral dot), while the same task's own stop
  reads Interrupted. A close Orca asks for (a Stop, a rest, a quit) now ends
  what still runs as stopped; only a session that dies on its own leaves the
  ending unknown. The close reads no stop cause. A rest never meets live child
  work: the idle sweep keeps a chat with any.
- If the host fails to record a dismissed card, Claude gets the card back, and
  a withdrawal Claude made meanwhile closes it, as before the host took it.

* refactor(native-chat): the session mutation path gets its own module, breaking the chat-stop import cycle

chat-stop.ts imported mutateStructuredAgentSession from host-mutations.ts,
which imports mutateWithChatStop back. The mutation context and the one
admit-then-serialize path now live in structured-agent-session-mutation-context.ts,
which both import; host-mutations re-exports the context type for its other
readers.

* fix(claude): a close that saw a descendant survive leaves background work's ending unknown

A Stop, rest or quit ended still-live background tasks as stopped even when the
close found Claude's root gone but a descendant still running. Only a close
that proved the whole tree gone stops them now; otherwise they settle unknown,
as for an exit of the session's own. The stopped ending is marked as Orca's,
so a frame of the child's own still replaces it.

* docs(native-chat): comments describe the Stop as it now works

What re-drives a failed child end, when the child-work decoder stops what is
live, what a Stop's wait re-reads on a result, and why a card's Stop is
unnamed.
This commit is contained in:
Brennan Benson
2026-09-30 23:39:40 -07:00
committed by GitHub
parent b245b3e019
commit 4bd55c07ea
44 changed files with 2148 additions and 859 deletions
@@ -193,6 +193,16 @@ describe('Claude child-work decoder', () => {
expect(decoder.drain(500)).toHaveLength(256)
})
it('ends what still runs as stopped when Orca ends the session, and only then', () => {
const decoder = decoderWith(backgroundAgent)
decoder.stopLive()
decoder.clear()
expect(decoder.drain(500)).toEqual([
expect.objectContaining({ type: 'ended', outcome: 'cancelled' }),
{ type: 'session-ended', observedAt: 500 }
])
})
it('marks the end of the provider session', () => {
const decoder = decoderWith(backgroundAgent)
decoder.clear()
+11 -3
View File
@@ -3,9 +3,9 @@
// A child is live from its `task_started` until its own terminal `task_updated` or
// `task_notification`; nothing else ends it. A roster (`background_tasks_changed`), a turn ending
// or a spawn call returning is the parent's view of the child, not the child's, and the CLI sends
// every child its own terminal frame, so none of them settles one. When the session ends, the
// host settles whatever is still live. Edges wait here until the frame is journaled, then take
// the host clock.
// every child its own terminal frame, so none of them settles one. When Orca ends the session and
// proves its tree gone, what is still live is stopped (`stopLive`); any other end leaves it for the
// host to settle as unknown. Edges wait here until the frame is journaled, then take the host clock.
import type {
AgentChildWorkKind,
@@ -124,6 +124,14 @@ export class ClaudeChildWorkDecoder {
}
}
/** Orca ended the session and proved its process tree gone: what still ran is stopped. The
* ending is Orca's, not the child's, so a frame of the child's own still replaces it. */
stopLive(): void {
for (const id of this.live.keys()) {
this.end(id, 'stopped', { basis: 'stop-acknowledged' })
}
}
/** The provider session is gone: the host settles what it still holds live. */
clear(): void {
this.live.clear()
@@ -11,15 +11,13 @@ import type { ClaudeCommandStart } from './claude-command-turn'
export type ClaudeJournalTranslator = {
handle: (event: ClaudeStructuredSessionEvent) => void
journalPrompts: Pick<ClaudeJournalPrompts, 'cancel' | 'resolve'>
journalPrompts: Pick<ClaudeJournalPrompts, 'resolve' | 'handOver' | 'cancel'>
/** The open turn's provider id — the same id its journal row carries, and the one
* a client's Stop names. Sole owner: no reader keeps a copy to disagree with. */
readonly currentTurnId: string | null
/** Orca is stopping this turn; its error end reads as the user's cancellation when they asked.
* False when the turn is no longer open. */
recordTurnStop: (turnId: string, cause: StructuredAgentSessionStopCause) => boolean
/** The provider refused the stop. */
withdrawTurnStop: (turnId: string) => void
/** The open turn's id while it is a conversation command's. */
readonly commandTurnId: string | null
/** Makes the host's command turn the open one until the command's result ends it. */
-7
View File
@@ -103,13 +103,6 @@ export class ClaudeOpenTurn {
return true
}
/** The provider refused the stop, so the turn goes on as if none was sent. */
withdrawStop(turnId: string): void {
if (this.current?.turnId === turnId) {
this.sentStop = null
}
}
/** Whether a turn is open inside a provider request cycle that has already
* done work — the state in which the CLI folds an arriving send into it. A
* cycle's first send is its opener, never a fold. */
+9 -72
View File
@@ -27,7 +27,6 @@ export type ClaudePendingPrompt = ClaudePromptPresentation & {
suggestions: PermissionUpdate[]
questionIds: readonly string[]
settle: ClaudePromptSettle
turnId?: string | null
}
export type ClaudePromptRegistration = ClaudePromptPresentation & {
@@ -37,12 +36,6 @@ export type ClaudePromptRegistration = ClaudePromptPresentation & {
input: Record<string, unknown>
suggestions: PermissionUpdate[]
settle: ClaudePromptSettle
turnId?: string | null
}
type PromptBinding = {
address: string
turnId: string | null
}
export type ClaudePromptClaim = {
@@ -50,11 +43,6 @@ export type ClaudePromptClaim = {
readonly found: { prompt: ClaudePendingPrompt }
}
type ClaudePromptCancellationObservation = {
promise: Promise<void>
resolve: () => void
}
export function isClaudePromptRecord(value: unknown): value is Record<string, unknown> {
return typeof value === 'object' && value !== null && !Array.isArray(value)
}
@@ -78,12 +66,9 @@ function questionId(question: Record<string, unknown>, index: number): string {
/** Session-local callback ownership; none of this state is reconstructed from the transcript. */
export class ClaudePromptRegistry {
private readonly prompts = new Map<string, ClaudePendingPrompt>()
private readonly journalBindings = new Map<string, PromptBinding>()
/** Journal item id to the prompt key it shows. */
private readonly journalBindings = new Map<string, string>()
private readonly claims = new Map<ClaudePendingPrompt, ClaudePromptClaim>()
private readonly cancellationObservations = new WeakMap<
ClaudePendingPrompt,
ClaudePromptCancellationObservation
>()
register(registration: ClaudePromptRegistration): ClaudePendingPrompt | null {
const toolUseId = readClaudePromptString(registration.toolUseId)
@@ -109,8 +94,7 @@ export class ClaudePromptRegistry {
...(registration.matchedAskRule ? { matchedAskRule: registration.matchedAskRule } : {}),
...(registration.subject ? { subject: registration.subject } : {}),
questionIds: questions.map(questionId),
settle: registration.settle,
turnId: registration.turnId ?? null
settle: registration.settle
}
this.prompts.set(prompt.promptKey, prompt)
return prompt
@@ -121,23 +105,17 @@ export class ClaudePromptRegistry {
if (!this.prompts.has(prompt.promptKey)) {
return false
}
const observation = this.cancellationObservations.get(prompt)
this.forget(prompt)
observation?.resolve()
return true
}
bindJournalItemId(journalItemId: string, promptKey: string, turnId: string | null = null): void {
const prompt = this.prompts.get(promptKey)
this.journalBindings.set(journalItemId, {
address: promptKey,
turnId: turnId ?? prompt?.turnId ?? null
})
bindJournalItemId(journalItemId: string, promptKey: string): void {
this.journalBindings.set(journalItemId, promptKey)
}
find(itemId: string): { prompt: ClaudePendingPrompt } | null {
const binding = this.journalBindings.get(itemId)
const prompt = this.prompts.get(binding?.address ?? itemId)
const promptKey = this.journalBindings.get(itemId)
const prompt = this.prompts.get(promptKey ?? itemId)
return prompt ? { prompt } : null
}
@@ -151,17 +129,6 @@ export class ClaudePromptRegistry {
return claim
}
claimBound(itemId: string, turnId: string): ClaudePromptClaim | null {
const binding = this.journalBindings.get(itemId)
const prompt = binding ? this.prompts.get(binding.address) : undefined
if (!binding || !prompt || binding.turnId !== turnId || this.claims.has(prompt)) {
return null
}
const claim = { itemId, found: { prompt } }
this.claims.set(prompt, claim)
return claim
}
ownsClaim(claim: ClaudePromptClaim): boolean {
return (
this.claims.get(claim.found.prompt) === claim &&
@@ -169,39 +136,12 @@ export class ClaudePromptRegistry {
)
}
ownsBoundClaim(claim: ClaudePromptClaim, itemId: string, turnId: string): boolean {
const binding = this.journalBindings.get(itemId)
return (
claim.itemId === itemId &&
this.claims.get(claim.found.prompt) === claim &&
binding?.address === claim.found.prompt.promptKey &&
binding.turnId === turnId &&
this.prompts.get(binding.address) === claim.found.prompt
)
}
releaseClaim(claim: ClaudePromptClaim): void {
if (this.claims.get(claim.found.prompt) === claim) {
this.claims.delete(claim.found.prompt)
}
}
observeCancellation(claim: ClaudePromptClaim): Promise<void> | null {
if (!this.ownsClaim(claim)) {
return null
}
let observation = this.cancellationObservations.get(claim.found.prompt)
if (!observation) {
let resolve = (): void => {}
const promise = new Promise<void>((settled) => {
resolve = settled
})
observation = { promise, resolve }
this.cancellationObservations.set(claim.found.prompt, observation)
}
return observation.promise
}
cancel(requestId: string): ClaudePendingPrompt | null {
const prompt = this.prompts.get(requestId) ?? null
if (prompt) {
@@ -213,8 +153,8 @@ export class ClaudePromptRegistry {
forget(prompt: ClaudePendingPrompt): void {
this.claims.delete(prompt)
this.prompts.delete(prompt.promptKey)
for (const [itemId, binding] of this.journalBindings) {
if (binding.address === prompt.promptKey) {
for (const [itemId, promptKey] of this.journalBindings) {
if (promptKey === prompt.promptKey) {
this.journalBindings.delete(itemId)
}
}
@@ -225,9 +165,6 @@ export class ClaudePromptRegistry {
this.prompts.clear()
this.journalBindings.clear()
this.claims.clear()
for (const prompt of pending) {
this.cancellationObservations.get(prompt)?.resolve()
}
return pending
}
}
@@ -0,0 +1,92 @@
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import {
awaitClaudeRequestEnd,
CLAUDE_STOP_GRACE_MS,
claudeStoppedRequestEndWait,
settleClaudeTurnEndWaiters
} from './claude-request-end-wait'
import type { ClaudeSession } from './claude-structured-session-state'
type InFlight = { translator: { currentTurnId: string | null }; dispatchWaiters: unknown[] }
// Only what Claude has in flight is read: the wait derives everything else from it. `state` is the
// same object, for the test to move Claude along.
function claudeWith(
currentTurnId: string | null = 'turn-1',
unanswered = 0
): { claude: ClaudeSession; state: InFlight } {
const state: InFlight = {
translator: { currentTurnId },
dispatchWaiters: Array.from({ length: unanswered }, () => ({}))
}
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the wait reads only `translator.currentTurnId` and `dispatchWaiters`, and keys its waiters by the session's identity.
return { claude: state as unknown as ClaudeSession, state }
}
function session(currentTurnId: string | null = 'turn-1'): ClaudeSession {
return claudeWith(currentTurnId).claude
}
async function settledWithin(promise: Promise<void>, ms: number): Promise<boolean> {
let settled = false
void promise.then(() => {
settled = true
})
await vi.advanceTimersByTimeAsync(ms)
return settled
}
describe("a Stop's wait for Claude to wind down what it has in flight", () => {
beforeEach(() => {
vi.useFakeTimers()
})
afterEach(() => {
vi.useRealTimers()
})
it('waits out the whole time for a turn Claude never ends', async () => {
const waiting = awaitClaudeRequestEnd(session(), CLAUDE_STOP_GRACE_MS)
expect(await settledWithin(waiting, CLAUDE_STOP_GRACE_MS - 1)).toBe(false)
expect(await settledWithin(waiting, 1)).toBe(true)
})
it('ends as soon as the stopped turn ends', async () => {
const { claude, state } = claudeWith()
const waiting = awaitClaudeRequestEnd(claude, CLAUDE_STOP_GRACE_MS)
state.translator.currentTurnId = null
settleClaudeTurnEndWaiters(claude)
expect(await settledWithin(waiting, 0)).toBe(true)
})
it('waits on a send Claude has not echoed, through its turn, to that turn’s end', async () => {
const { claude, state } = claudeWith(null, 1)
const waiting = awaitClaudeRequestEnd(claude, CLAUDE_STOP_GRACE_MS)
// The echo answers the send and opens its turn; a settle in between leaves it in flight.
state.dispatchWaiters.length = 0
state.translator.currentTurnId = 'turn-1'
settleClaudeTurnEndWaiters(claude)
expect(await settledWithin(waiting, 0)).toBe(false)
state.translator.currentTurnId = null
settleClaudeTurnEndWaiters(claude)
expect(await settledWithin(waiting, 0)).toBe(true)
})
it('holds nothing when Claude has nothing in flight', async () => {
expect(await settledWithin(awaitClaudeRequestEnd(session(null), 1_000), 0)).toBe(true)
})
it('waits only what the interrupt left of the grace, and nothing for a session that is gone', async () => {
const sessions = new Map([['session-1', session()]])
const wait = claudeStoppedRequestEndWait(sessions)
const stoppedAt = Date.now()
await vi.advanceTimersByTimeAsync(CLAUDE_STOP_GRACE_MS - 1_000)
const waiting = wait('session-1', stoppedAt)
expect(await settledWithin(waiting, 999)).toBe(false)
expect(await settledWithin(waiting, 1)).toBe(true)
expect(await settledWithin(wait('session-2', Date.now()), 0)).toBe(true)
})
})
@@ -0,0 +1,81 @@
// A Stop ends Claude's child. Before it does, Claude gets what is left of the grace to wind down what
// it has in flight through its own path: an unechoed send echoes and opens its turn, the stopped turn
// ends with its own result, and Claude writes both to its transcript. Keyed on Claude's own request,
// not on a journal turn: a Stop pressed before the echo has no turn row yet.
import type { ClaudeSession } from './claude-structured-session-state'
/** One budget from the Stop's interrupt: for Claude to answer it and wind its request down. */
export const CLAUDE_STOP_GRACE_MS = 3_000
const waiters = new WeakMap<ClaudeSession, Set<() => void>>()
/** Claude has a turn open, or a send it has not answered yet. */
function claudeRequestInFlight(session: ClaudeSession): boolean {
return (session.translator?.currentTurnId ?? null) !== null || session.dispatchWaiters.length > 0
}
function nextSettle(session: ClaudeSession): { settled: Promise<void>; forget: () => void } {
let waiting = waiters.get(session)
if (!waiting) {
waiting = new Set()
waiters.set(session, waiting)
}
const set = waiting
let wake!: () => void
const settled = new Promise<void>((resolve) => {
wake = resolve
})
set.add(wake)
return { settled, forget: () => set.delete(wake) }
}
/** Resolves once Claude has nothing in flight, re-read at each settle, or after `ms`. */
export async function awaitClaudeRequestEnd(session: ClaudeSession, ms: number): Promise<void> {
let elapsed = false
let timer: ReturnType<typeof setTimeout> | undefined
const budget = new Promise<void>((resolve) => {
timer = setTimeout(() => {
elapsed = true
resolve()
}, ms)
timer.unref?.()
})
try {
// A settle that leaves something in flight, such as a subagent's result, waits on.
while (!elapsed && claudeRequestInFlight(session)) {
const next = nextSettle(session)
try {
await Promise.race([next.settled, budget])
} finally {
next.forget()
}
}
} finally {
clearTimeout(timer)
}
}
/** The adapter's wait: whatever of the grace the Stop's interrupt left, counted from `stoppedAt`. */
export function claudeStoppedRequestEndWait(
sessions: Map<string, ClaudeSession>
): (sessionId: string, stoppedAt: number) => Promise<void> {
return async (sessionId, stoppedAt) => {
const session = sessions.get(sessionId)
if (session) {
await awaitClaudeRequestEnd(
session,
Math.max(0, stoppedAt + CLAUDE_STOP_GRACE_MS - Date.now())
)
}
}
}
/** Something Claude had in flight settled: a result, the CLI's idle, the child's exit or its close. */
export function settleClaudeTurnEndWaiters(session: ClaudeSession): void {
const waiting = waiters.get(session)
waiters.delete(session)
for (const wake of waiting ?? []) {
wake()
}
}
@@ -87,7 +87,7 @@ describe('Claude child work from captured frame orders', () => {
})
})
it('settles what still runs when the session ends, and keeps every record', async () => {
it('stops what still runs when Orca ends the session, and keeps every record', async () => {
const run = await producer(hostWithParent())
const [running] = until(MOVED_TO_BACKGROUND, 51_275)
run.replay(running)
@@ -99,8 +99,37 @@ describe('Claude child work from captured frame orders', () => {
outcome
}))
).toEqual([
{ description: 'Run 45s sleep command', membership: 'settled', outcome: 'unknown' },
{ description: 'Sleep 45 seconds then print 1', membership: 'settled', outcome: 'unknown' }
{ description: 'Run 45s sleep command', membership: 'settled', outcome: 'cancelled' },
{ description: 'Sleep 45 seconds then print 1', membership: 'settled', outcome: 'cancelled' }
])
})
it('leaves how they ended unknown when the close sees a descendant survive', async () => {
const run = await producer(hostWithParent())
const [running] = until(MOVED_TO_BACKGROUND, 51_275)
run.replay(running)
// The root exited, but Orca's own tree check saw a process of the session's still running.
const connection = run.claude.connections[0]!
connection.exitVerdict = { root: 'exited', tree: 'live' }
connection.close = async () => false
await expect(run.adapter.closeSession('session-1', 'user-stop')).rejects.toMatchObject({
name: 'AgentSessionAcquisitionRootExitObservedError'
})
expect(run.records().map(({ membership, outcome }) => ({ membership, outcome }))).toEqual([
{ membership: 'settled', outcome: 'unknown' },
{ membership: 'settled', outcome: 'unknown' }
])
})
it('leaves how they ended unknown when the session dies on its own', async () => {
const run = await producer(hostWithParent())
const [running] = until(MOVED_TO_BACKGROUND, 51_275)
run.replay(running)
run.claude.connections[0]!.handlers.onExit?.(new Error('claude exited'))
await run.adapter.drainObservedExits()
expect(run.records().map(({ membership, outcome }) => ({ membership, outcome }))).toEqual([
{ membership: 'settled', outcome: 'unknown' },
{ membership: 'settled', outcome: 'unknown' }
])
})
@@ -195,7 +195,7 @@ describe('Claude structured child-work producer', () => {
step.check?.()
}
await adapter.closeSession('session-1')
// The session is gone: what still ran settles unreported, and every record stays.
// Orca ended the session: what still ran is stopped with it, and every record stays.
const summary = records().map((record) => ({
description: record.description,
membership: record.membership,
@@ -210,7 +210,7 @@ describe('Claude structured child-work producer', () => {
generation: 1
},
{ description: 'npm test', membership: 'settled', outcome: 'cancelled', generation: 1 },
{ description: 'Audit the build', membership: 'settled', outcome: 'unknown', generation: 2 }
{ description: 'Audit the build', membership: 'settled', outcome: 'cancelled', generation: 2 }
])
expect(adapter.backgroundTaskState('session-1')).toBeUndefined()
})
@@ -197,7 +197,7 @@ describe('cancelClaudeTurn', () => {
})
describe('answerClaudePrompt', () => {
it('resolves cancellation observation when teardown clears the prompt registry', async () => {
it('forgets a pending prompt and its claim when teardown clears the registry', () => {
const prompts = new ClaudePromptRegistry()
const settle = vi.fn()
const prompt = prompts.register({
@@ -213,19 +213,8 @@ describe('answerClaudePrompt', () => {
if (!claim) {
throw new Error('expected prompt claim')
}
const observed = prompts.observeCancellation(claim)
if (!observed) {
throw new Error('expected cancellation observation')
}
let observedCancellation = false
void observed.then(() => {
observedCancellation = true
})
expect(prompts.clear()).toEqual([prompt])
await Promise.resolve()
expect(observedCancellation).toBe(true)
expect(prompts.find('journal-clear')).toBeNull()
expect(prompts.ownsClaim(claim)).toBe(false)
expect(settle).not.toHaveBeenCalled()
@@ -249,12 +238,12 @@ describe('answerClaudePrompt', () => {
handle: vi.fn(),
openTurnInLiveProviderCycle: false,
journalPrompts: {
cancel: vi.fn(() => ({ accepted: true as const })),
resolve: resolvePrompt
resolve: resolvePrompt,
handOver: () => () => {},
cancel: () => ({ accepted: true })
},
currentTurnId: null,
recordTurnStop: () => true,
withdrawTurnStop: () => {},
commandTurnId: null,
beginCommand: vi.fn(),
forgetCommand: vi.fn(),
@@ -58,12 +58,9 @@ export async function cancelClaudeTurn(
}
return { cancelled: true }
} catch (error) {
// The CLI refused. Any other error leaves the interrupt's effect unknown. Either way the Stop
// ends the child next, so the stop recorded on the turn stands.
if (error instanceof ClaudeControlRequestError) {
// The CLI refused, so the turn runs on and its own end means what it says. Any other error
// leaves the interrupt's effect unknown, and the stop the user asked for stands.
if (stopped) {
session.translator?.withdrawTurnStop(stopped.turnId)
}
return { cancelled: false }
}
throw error
@@ -1,5 +1,6 @@
// A user's Stop inside a live Claude chat interrupts the turn and keeps the session. The turn's end
// then comes from the CLI's result frame, which CLIs before 2.1.91 send with no terminal_reason.
// A user's Stop inside a live Claude chat interrupts the turn before the host ends its child. The
// turn's end then comes from the CLI's result frame, which CLIs before 2.1.91 send with no
// terminal_reason.
import { describe, expect, it, vi } from 'vitest'
import type { AgentJournalItemBody } from '../../shared/agent-session-journal-types'
@@ -155,11 +156,12 @@ describe("a user's Stop inside a live Claude chat", () => {
expect(settled(bodies, nextTurnId)).toMatchObject({ state: 'completed', outcome: 'failure' })
})
// The Stop ends the child next, so the turn it was asked for reads Interrupted however it ends.
it.each([
['naming the turn', true],
['naming no turn', false]
] as const)(
'keeps a failure the turn reaches after the CLI refused the interrupt, %s',
'reads a turn the CLI refused to interrupt as the user stopping it, %s',
async (_label, named) => {
const claude = fakeClaude({
routes: {
@@ -175,8 +177,8 @@ describe("a user's Stop inside a live Claude chat", () => {
).resolves.toEqual({ cancelled: false })
connection.handlers.onMessage?.(CUT_SHORT)
expect(settled(bodies, turnId)).toMatchObject({ state: 'completed', outcome: 'failure' })
expect(providerRows(bodies)).toHaveLength(1)
expect(settled(bodies, turnId)).toMatchObject({ outcome: 'cancellation' })
expect(providerRows(bodies)).toHaveLength(0)
}
)
})
@@ -10,7 +10,6 @@ export type ClaudePermissionCallbackDeps = {
sessionId: string
prompts: ClaudePromptRegistry
emit: (event: ClaudeStructuredSessionEvent) => void
currentTurnId?: () => string | null
}
function denySafeResult(toolUseId: string | undefined): PermissionResult {
@@ -52,8 +51,7 @@ export function buildClaudePermissionCallbacks(deps: ClaudePermissionCallbackDep
toolUseId: options.toolUseID,
input,
suggestions: options.suggestions ?? [],
settle,
turnId: deps.currentTurnId?.() ?? null
settle
})
if (!prompt) {
settle(denySafeResult(options.toolUseID))
@@ -202,6 +202,18 @@ export class ClaudeJournalPrompts {
this.deletePrompt(promptKey)
}
/** The host records the card itself, so nothing here writes it any more. The returned undo hands
* it back when that record fails, so Claude's own withdrawal can still close it. */
handOver(promptKey: string): () => void {
const entry = this.items.get(promptKey)
this.deletePrompt(promptKey)
return () => {
if (entry && !this.items.has(promptKey)) {
this.items.set(promptKey, { items: entry.items, cancellationPending: false })
}
}
}
clear(): void {
this.items.clear()
this.pendingCancellationTotal = 0
@@ -283,7 +283,6 @@ export function createClaudeJournalTranslator(
return turn.id
},
recordTurnStop: (turnId, cause) => turn.recordStop(turnId, cause),
withdrawTurnStop: (turnId) => turn.withdrawStop(turnId),
get commandTurnId() {
return turn.command ? turn.id : null
},
@@ -53,8 +53,7 @@ describe('Claude structured approval presentation', () => {
options: [
{ id: 'allow', label: 'Allow' },
{ id: 'allowForSession', label: 'Allow for this session' },
{ id: 'deny', label: 'Deny' },
{ id: 'cancel', label: 'Stop' }
{ id: 'deny', label: 'Deny' }
]
})
})
@@ -75,12 +74,33 @@ describe('Claude structured approval presentation', () => {
detail: '# Release\n\n- Run tests',
options: [
{ id: 'allow', label: 'Approve plan' },
{ id: 'deny', label: 'Keep planning' },
{ id: 'cancel', label: 'Stop' }
{ id: 'deny', label: 'Keep planning' }
]
})
})
it('answers a dismissal as a plain deny, and a dismissed plan by asking Claude to wait', () => {
const tool = buildClaudePromptReply(approvalPrompt({ command: 'ls' }), {
kind: 'option',
optionId: 'cancel'
})
const plan = buildClaudePromptReply(
approvalPrompt({ plan: '# Release' }, { subject: { kind: 'plan', text: '# Release' } }),
{ kind: 'option', optionId: 'cancel' }
)
expect(tool).toEqual({
behavior: 'deny',
message: 'User denied this action.',
toolUseID: expect.any(String)
})
expect(plan).toMatchObject({
behavior: 'deny',
message: expect.stringMatching(/wait for them/)
})
expect(plan).not.toHaveProperty('interrupt')
})
it('uses a plan-specific fallback title when the harness omits one', () => {
const item = claudeApprovalItem(
approvalPrompt({ plan: '# Release' }, { subject: { kind: 'plan', text: '# Release' } })
@@ -9,23 +9,19 @@ import { formatToolInput, truncateToolDetail } from '../../shared/native-chat-to
import { boundJournalPromptBody } from '../native-chat/agent-session-journal/journal-prompt-body-bounds'
import { claudeRecord, claudeText } from './claude-structured-item-translation'
import {
CLAUDE_APPROVAL_DECISIONS,
encodeClaudeQuestionOptionId,
type ClaudeApprovalDecision,
type ClaudePendingPrompt
} from './claude-structured-prompt-replies'
const APPROVAL_LABELS: Record<ClaudeApprovalDecision, string> = {
allow: 'Allow',
allowForSession: 'Allow for this session',
deny: 'Deny',
cancel: 'Stop'
}
const APPROVAL_OPTIONS: readonly AgentJournalPromptOption[] = [
{ id: 'allow', label: 'Allow' },
{ id: 'allowForSession', label: 'Allow for this session' },
{ id: 'deny', label: 'Deny' }
]
const PLAN_APPROVAL_OPTIONS: readonly AgentJournalPromptOption[] = [
{ id: 'allow', label: 'Approve plan' },
{ id: 'deny', label: 'Keep planning' },
{ id: 'cancel', label: 'Stop' }
{ id: 'deny', label: 'Keep planning' }
]
const PENDING = {
@@ -60,12 +56,9 @@ export function claudeApprovalItem(prompt: ClaudePendingPrompt): AgentJournalApp
...(prompt.matchedAskRule ? { matchedAskRule: prompt.matchedAskRule } : {}),
...(prompt.subject ? { subject: prompt.subject } : {}),
detail: detail || null,
options: planSubject
? PLAN_APPROVAL_OPTIONS.map((option) => ({ ...option }))
: CLAUDE_APPROVAL_DECISIONS.map((decision) => ({
id: decision,
label: APPROVAL_LABELS[decision]
})),
options: (planSubject ? PLAN_APPROVAL_OPTIONS : APPROVAL_OPTIONS).map((option) => ({
...option
})),
resolution: { ...PENDING }
})
}
@@ -2,18 +2,15 @@ import { AGENT_JOURNAL_THREAD_SCOPE } from '../../shared/agent-session-journal-t
import { describe, expect, it, vi } from 'vitest'
import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key'
import type { AgentJournalItemBody } from '../../shared/agent-session-journal-types'
import { readAgentJournalTurn } from '../../shared/agent-session-turn-record'
import type {
StructuredAgentSessionAppendOptions,
StructuredAgentSessionEventSink
} from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { ClaudeControlRequestError } from './claude-stream-json-connection'
import { ClaudeJournalPrompts } from './claude-structured-journal-prompts'
import { claudeQuestionItems } from './claude-structured-prompt-items'
import type { ClaudePendingPrompt } from './claude-structured-prompt-replies'
import type { ClaudeStructuredSessionEvent } from './claude-structured-session-state'
import {
PROVIDER_SESSION_ID,
USER_MESSAGE,
acquired,
adapterFor,
@@ -131,16 +128,6 @@ describe('Claude live prompt ownership', () => {
})
await vi.waitFor(() => expect(commitStarted).toHaveBeenCalledOnce())
await expect(
adapter.cancelTurn({
sessionId: 'session-1',
turnId: 'turn-1',
fence: 7,
prompt: { itemId: 'journal-prompt' }
})
).resolves.toEqual({ cancelled: false })
expect(claude.connections[0]?.calls.some((call) => call.subtype === 'interrupt')).toBe(false)
commitGate.resolve()
await answer
await expect(answered.promise).resolves.toMatchObject({
@@ -149,196 +136,80 @@ describe('Claude live prompt ownership', () => {
})
})
it('lets prompt cancellation win and waits for SDK abort cleanup', async () => {
const interruptGate = deferred()
const controller = new AbortController()
const claude = fakeClaude({
replayUuid: 'turn-1',
routes: { interrupt: () => interruptGate.promise }
it("keeps the user's dismissal when Claude cancels the request while the host records it", async () => {
const claude = fakeClaude({ replayUuid: 'turn-1' })
const recorded = lifecycleRecorder()
const adapter = adapterFor(claude)
await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9',
events: recorded.sink
})
const adapter = await acquired(claude)
await startTurn(adapter)
const connection = claude.connections[0]
if (!connection) {
throw new Error('expected Claude connection')
const request = new AbortController()
invokeCanUseTool(claude.connections[0]!, 'AskUserQuestion', 'permission-1', 'tool-1', {
input: { questions: [{ question: 'Which branch?', options: [{ label: 'main' }] }] },
signal: request.signal
})
const [card] = [...recorded.bodies].find(([, body]) => body.kind === 'question') ?? []
if (!card) {
throw new Error('expected the question card')
}
const answered = invokeCanUseTool(connection, 'Bash', 'permission-1', 'tool-1', {
input: { command: 'git status' },
signal: controller.signal
})
adapter.bindPromptItemId('session-1', 'journal-prompt', 'permission-1')
let cancellationSettled = false
const cancellation = adapter
.cancelTurn({
sessionId: 'session-1',
turnId: 'turn-1',
fence: 7,
prompt: { itemId: 'journal-prompt' }
})
.finally(() => {
cancellationSettled = true
})
await vi.waitFor(() => expect(claude.connections[0]?.calls.at(-1)?.subtype).toBe('interrupt'))
const commit = vi.fn(async () => undefined)
await expect(
adapter.answerPrompt({
sessionId: 'session-1',
itemId: 'journal-prompt',
kind: 'approval',
response: { kind: 'option', optionId: 'allow' },
fence: 7,
commit
})
).rejects.toThrow(/no longer waiting/)
expect(commit).not.toHaveBeenCalled()
interruptGate.resolve()
await Promise.resolve()
expect(cancellationSettled).toBe(false)
expect(answered.settled()).toBe(false)
controller.abort()
await expect(cancellation).resolves.toEqual({ cancelled: true })
await expect(answered.promise).resolves.toBeNull()
await expect(
adapter.answerPrompt({
sessionId: 'session-1',
itemId: 'journal-prompt',
kind: 'approval',
response: { kind: 'option', optionId: 'allow' },
fence: 7,
commit
})
).rejects.toThrow(/no longer waiting/)
expect(controller.signal.aborted).toBe(true)
expect(commit).not.toHaveBeenCalled()
})
it('cancels an owned prompt after another dispatch queues behind its turn', async () => {
const controller = new AbortController()
let queuedUuid = ''
const claude = fakeClaude({
replayUuids: ['turn-1', null],
capabilities: ['interrupt_cancel_queued_v1'],
routes: {
interrupt: () => {
controller.abort()
return { still_queued: [], cancelled: [queuedUuid] }
}
}
})
const lateSettlements: unknown[] = []
const adapter = await acquired(claude, {}, [], (settlement) => lateSettlements.push(settlement))
await startTurn(adapter)
const connection = claude.connections[0]
if (!connection) {
throw new Error('expected Claude connection')
const dismiss = adapter.dismissPrompt
if (!dismiss) {
throw new Error('expected Claude to dismiss a card')
}
const answered = invokeCanUseTool(connection, 'Bash', 'permission-queued', 'tool-queued', {
input: { command: 'git status' },
signal: controller.signal
})
adapter.bindPromptItemId('session-1', 'journal-prompt', 'permission-queued')
await expect(
adapter.dispatch({
sessionId: 'session-1',
clientMessageId: 'queued-message',
body: USER_MESSAGE,
fence: 7
})
).resolves.toEqual({ state: 'admitted' })
const sentUuid = connection.sent.at(-1)?.uuid
if (typeof sentUuid !== 'string') {
throw new Error('expected queued dispatch uuid')
}
queuedUuid = sentUuid
await expect(
adapter.cancelTurn({
sessionId: 'session-1',
turnId: 'turn-1',
fence: 7,
prompt: { itemId: 'journal-prompt' }
})
).resolves.toEqual({ cancelled: true })
await expect(answered.promise).resolves.toBeNull()
expect(connection.calls).toContainEqual({
subtype: 'interrupt',
params: { cancelQueued: true }
})
expect(lateSettlements).toContainEqual({
await dismiss({
sessionId: 'session-1',
clientMessageId: 'queued-message',
state: 'rejected',
reason: 'provider_cancelled_before_start',
rejection: { kind: 'cancelled' }
itemId: card,
fence: 7,
answer: false,
// Claude's own cancel lands mid-commit, as an interrupt's can.
commit: async () => request.abort()
})
expect(recorded.bodies.get(card)).toMatchObject({ resolution: { state: 'pending' } })
})
it('does not interrupt a queued turn when the CLI cannot cancel queued messages', async () => {
const claude = fakeClaude({ replayUuids: ['turn-1', null] })
const adapter = await acquired(claude)
await startTurn(adapter)
const connection = claude.connections[0]
if (!connection) {
throw new Error('expected Claude connection')
}
const controller = new AbortController()
const answered = invokeCanUseTool(connection, 'Bash', 'permission-legacy', 'tool-legacy', {
input: { command: 'git status' },
signal: controller.signal
it('hands the card back when the host fails to record it, so a withdrawal still closes it', async () => {
const claude = fakeClaude({ replayUuid: 'turn-1' })
const recorded = lifecycleRecorder()
const adapter = adapterFor(claude)
await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9',
events: recorded.sink
})
adapter.bindPromptItemId('session-1', 'journal-prompt', 'permission-legacy')
await expect(
adapter.dispatch({
sessionId: 'session-1',
clientMessageId: 'queued-message',
body: USER_MESSAGE,
fence: 7
})
).resolves.toEqual({ state: 'admitted' })
await expect(
adapter.cancelTurn({
sessionId: 'session-1',
turnId: 'turn-1',
fence: 7,
prompt: { itemId: 'journal-prompt' }
})
).resolves.toEqual({ cancelled: false })
expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(false)
controller.abort()
await expect(answered.promise).resolves.toBeNull()
})
it('does not interrupt a newer active turn through a stale prompt callback', async () => {
const claude = fakeClaude({ replayUuids: ['turn-1', 'turn-2'] })
const adapter = await acquired(claude)
await startTurn(adapter)
const connection = claude.connections[0]
if (!connection) {
throw new Error('expected Claude connection')
}
const controller = new AbortController()
const answered = invokeCanUseTool(connection, 'Bash', 'permission-stale', 'tool-stale', {
input: { command: 'git status' },
signal: controller.signal
const request = new AbortController()
invokeCanUseTool(claude.connections[0]!, 'AskUserQuestion', 'permission-1', 'tool-1', {
input: { questions: [{ question: 'Which branch?', options: [{ label: 'main' }] }] },
signal: request.signal
})
adapter.bindPromptItemId('session-1', 'journal-prompt', 'permission-stale')
await startTurn(adapter, 'turn-2')
const [card] = [...recorded.bodies].find(([, body]) => body.kind === 'question') ?? []
const dismiss = adapter.dismissPrompt
if (!card || !dismiss) {
throw new Error('expected the question card and a dismissal')
}
await expect(
adapter.cancelTurn({
dismiss({
sessionId: 'session-1',
turnId: 'turn-1',
itemId: card,
fence: 7,
prompt: { itemId: 'journal-prompt' }
answer: false,
// Claude withdraws the request, then the host's record fails.
commit: async () => {
request.abort()
throw new Error('journal write failed')
}
})
).resolves.toEqual({ cancelled: false })
expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(false)
expect(answered.settled()).toBe(false)
controller.abort()
await expect(answered.promise).resolves.toBeNull()
).rejects.toThrow('journal write failed')
expect(recorded.bodies.get(card)).toMatchObject({ resolution: { state: 'cancelled' } })
})
it('drops resolved prompt bodies instead of retaining them for the session lifetime', () => {
@@ -370,141 +241,6 @@ describe('Claude live prompt ownership', () => {
expect(prompts.size).toBe(0)
})
it('releases the callback claim after a failed interrupt', async () => {
const claude = fakeClaude({
replayUuid: 'turn-1',
routes: {
interrupt: () => {
throw new ClaudeControlRequestError('interrupt', 'not running')
}
}
})
const adapter = await acquired(claude)
await startTurn(adapter)
const connection = claude.connections[0]
if (!connection) {
throw new Error('expected Claude connection')
}
const answered = invokeCanUseTool(connection, 'Bash', 'permission-1', 'tool-1', {
input: { command: 'git status' }
})
adapter.bindPromptItemId('session-1', 'journal-prompt', 'permission-1')
await expect(
adapter.cancelTurn({
sessionId: 'session-1',
turnId: 'turn-1',
fence: 7,
prompt: { itemId: 'journal-prompt' }
})
).resolves.toEqual({ cancelled: false })
await adapter.answerPrompt({
sessionId: 'session-1',
itemId: 'journal-prompt',
kind: 'approval',
response: { kind: 'option', optionId: 'allow' },
fence: 7,
commit: async () => undefined
})
await expect(answered.promise).resolves.toMatchObject({
behavior: 'allow',
toolUseID: 'tool-1'
})
})
it('enqueues terminal prompt state before a confirmed cancellation resolves', async () => {
const controller = new AbortController()
const claude = fakeClaude({
replayUuid: 'turn-1',
routes: { interrupt: () => controller.abort() }
})
const recorded = lifecycleRecorder()
const adapter = adapterFor(claude)
await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9',
events: recorded.sink
})
await startTurn(adapter)
const connection = claude.connections[0]
if (!connection) {
throw new Error('expected Claude connection')
}
const answered = invokeCanUseTool(connection, 'Bash', 'permission-1', 'tool-1', {
input: { command: 'git status' },
signal: controller.signal
})
adapter.bindPromptItemId('session-1', 'journal-prompt', 'permission-1')
const promptItemId = [...recorded.bodies].find(([, body]) => body.kind === 'approval')?.[0]
const cancellation = adapter
.cancelTurn({
sessionId: 'session-1',
turnId: 'turn-1',
fence: 7,
prompt: { itemId: 'journal-prompt' }
})
.then((result) => {
recorded.order.push('resolved')
return result
})
await expect(cancellation).resolves.toEqual({ cancelled: true })
await expect(answered.promise).resolves.toBeNull()
if (!promptItemId) {
throw new Error('expected a recorded prompt item')
}
expect(recorded.order).toEqual(['prompt-lifecycle', 'resolved'])
expect(
[...recorded.bodies.values()].some(
(body) =>
(body.kind === 'approval' || body.kind === 'question') &&
body.resolution.state === 'pending'
)
).toBe(false)
expect(recorded.bodies.get(promptItemId)).toMatchObject({
resolution: { state: 'cancelled' }
})
expect(
[...recorded.bodies.values()].some(
(body) => readAgentJournalTurn(body)?.state === 'interrupted'
)
).toBe(false)
connection.handlers.onMessage?.({
type: 'result',
subtype: 'error_during_execution',
uuid: 'result-1',
session_id: PROVIDER_SESSION_ID,
is_error: true,
terminal_reason: 'aborted_tools',
errors: [],
duration_ms: 654
})
expect([...recorded.bodies.values()].find((body) => readAgentJournalTurn(body))).toMatchObject({
state: 'interrupted',
durationMs: 654
})
connection.handlers.onMessage?.({
type: 'result',
subtype: 'success',
uuid: 'result-duplicate',
session_id: PROVIDER_SESSION_ID,
is_error: false,
terminal_reason: 'completed',
duration_ms: 999
})
expect([...recorded.bodies.values()].find((body) => readAgentJournalTurn(body))).toMatchObject({
state: 'interrupted',
durationMs: 654
})
expect(controller.signal.aborted).toBe(true)
expect(recorded.tombstones).toHaveLength(0)
})
it('does not synthesize terminal lifecycle for ordinary Stop', async () => {
const events: ClaudeStructuredSessionEvent[] = []
const adapter = await acquired(fakeClaude({ replayUuid: 'turn-1' }), {}, events)
@@ -519,110 +255,6 @@ describe('Claude live prompt ownership', () => {
).toBe(false)
})
it('does not report success or release the claim when prompt lifecycle admission fails', async () => {
const controller = new AbortController()
const claude = fakeClaude({
replayUuid: 'turn-1',
routes: { interrupt: () => controller.abort() }
})
const recorded = lifecycleRecorder(false)
const adapter = adapterFor(claude)
await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9',
events: recorded.sink
})
await startTurn(adapter)
const connection = claude.connections[0]
if (!connection) {
throw new Error('expected Claude connection')
}
invokeCanUseTool(connection, 'Bash', 'permission-1', 'tool-1', {
input: { command: 'git status' },
signal: controller.signal
})
const promptItemId = [...recorded.bodies].find(([, body]) => body.kind === 'approval')?.[0]
if (!promptItemId) {
throw new Error('expected durable Claude prompt')
}
await expect(
adapter.cancelTurn({
sessionId: 'session-1',
turnId: 'turn-1',
fence: 7,
prompt: { itemId: promptItemId }
})
).rejects.toThrow(/lifecycle was not admitted/)
const commit = vi.fn(async () => undefined)
await expect(
adapter.answerPrompt({
sessionId: 'session-1',
itemId: promptItemId,
kind: 'approval',
response: { kind: 'option', optionId: 'allow' },
fence: 7,
commit
})
).rejects.toThrow(/no longer waiting/)
expect(commit).not.toHaveBeenCalled()
})
it('checks the bound item, turn, fence, and current acquisition without callback revival', async () => {
const claude = fakeClaude({ replayUuid: 'turn-1' })
const adapter = await acquired(claude)
await startTurn(adapter)
const connection = claude.connections[0]
if (!connection) {
throw new Error('expected Claude connection')
}
const answered = invokeCanUseTool(connection, 'Bash', 'permission-1', 'tool-1', {
input: { command: 'git status' }
})
adapter.bindPromptItemId('session-1', 'journal-prompt', 'permission-1')
for (const input of [
{ turnId: 'turn-1', fence: 7, itemId: 'other-item' },
{ turnId: 'turn-2', fence: 7, itemId: 'journal-prompt' },
{ turnId: 'turn-1', fence: 6, itemId: 'journal-prompt' }
]) {
await expect(
adapter.cancelTurn({
sessionId: 'session-1',
turnId: input.turnId,
fence: input.fence,
prompt: { itemId: input.itemId }
})
).resolves.toEqual({ cancelled: false })
}
expect(claude.connections[0]?.calls.some((call) => call.subtype === 'interrupt')).toBe(false)
await adapter.acquire({ identity: identityFor(), fence: 8, spawnToken: 'spawn-10' })
await expect(answered.promise).resolves.toBeNull()
await expect(
adapter.cancelTurn({
sessionId: 'session-1',
turnId: 'turn-1',
fence: 8,
prompt: { itemId: 'journal-prompt' }
})
).resolves.toEqual({ cancelled: false })
const commit = vi.fn(async () => undefined)
await expect(
adapter.answerPrompt({
sessionId: 'session-1',
itemId: 'journal-prompt',
kind: 'approval',
response: { kind: 'option', optionId: 'allow' },
fence: 8,
commit
})
).rejects.toThrow(/no longer waiting/)
expect(commit).not.toHaveBeenCalled()
expect(claude.connections[1]?.calls.some((call) => call.subtype === 'interrupt')).toBe(false)
})
it('rejects a grouped prompt batch without partially revising its first row', () => {
const tombstones: string[] = []
const appendTombstone = vi.fn(
@@ -3,45 +3,18 @@ import {
AgentSessionPromptUnavailableError,
type StructuredAgentSessionAdapter
} from '../native-chat/agent-session-wire/structured-agent-session-adapter'
import { CLAUDE_DEFAULT_REQUEST_TIMEOUT_MS } from './claude-agent-sdk-control-requests'
import {
answerClaudePrompt,
cancelClaudeTurn,
supportsClaudeQueuedInterruptCancellation
} from './claude-structured-control-actions'
import type { StructuredAgentSessionAdapterStop } from '../native-chat/agent-session-wire/structured-agent-session-adapter-stop'
import { answerClaudePrompt, cancelClaudeTurn } from './claude-structured-control-actions'
import type { ClaudeLateDispatchSettlement } from './claude-replay-turn-resolution'
import { buildClaudePromptReply } from './claude-structured-prompt-replies'
import { buildClaudePromptReply, claudePromptDismissal } from './claude-structured-prompt-replies'
import type { ClaudeSession } from './claude-structured-session-state'
import type { ClaudePendingPrompt } from './claude-prompt-registry'
import { CLAUDE_STOP_GRACE_MS } from './claude-request-end-wait'
import type { PermissionResult } from '@anthropic-ai/claude-agent-sdk'
type CancelInput = Parameters<StructuredAgentSessionAdapter['cancelTurn']>[0]
type AnswerInput = Parameters<StructuredAgentSessionAdapter['answerPrompt']>[0]
export function admitClaudePromptCancellation(session: ClaudeSession, promptKey: string): boolean {
const admission = session.translator?.journalPrompts.cancel(promptKey)
return admission?.accepted ?? true
}
function waitForClaudePromptCancellation(
observed: Promise<void>,
timeoutMs = CLAUDE_DEFAULT_REQUEST_TIMEOUT_MS
): Promise<void> {
let timer: ReturnType<typeof setTimeout> | null = null
const deadline = new Promise<never>((_resolve, reject) => {
timer = setTimeout(
() => reject(new Error('Claude prompt cancellation abort was not observed')),
timeoutMs
)
timer.unref?.()
})
return Promise.race([observed, deadline]).finally(() => {
if (timer) {
clearTimeout(timer)
}
})
}
function requireSession(sessions: Map<string, ClaudeSession>, sessionId: string): ClaudeSession {
const session = sessions.get(sessionId)
if (!session) {
@@ -87,19 +60,20 @@ function cancelClaudeConversation(
)
}
/** A Stop's interrupt. A card's own Cancel never comes here: `claudePromptCancelRoute` routes it. */
export async function cancelClaudeStructuredTurn(input: {
request: CancelInput
sessions: Map<string, ClaudeSession>
timeoutMs?: number
admitPromptCancellation: (session: ClaudeSession, promptKey: string) => boolean
onDispatchSettledLate?: ClaudeLateDispatchSettlement
}): Promise<{ cancelled: boolean }> {
const { request, sessions, timeoutMs } = input
const { request, sessions } = input
// A Stop ends the child next, so its interrupt shares the grace with Claude's wind-down after it.
const timeoutMs = Math.min(input.timeoutMs ?? CLAUDE_STOP_GRACE_MS, CLAUDE_STOP_GRACE_MS)
const session = requireSession(sessions, request.sessionId)
const acquisitionGeneration = session.acquisitionGeneration
const prompt = request.prompt
// Before startup lands nothing was written, so there is nothing to interrupt.
if (!prompt && session.startup.state === 'pending') {
if (request.prompt || session.startup.state === 'pending') {
return { cancelled: false }
}
const requestedTurnId = request.turnId
@@ -108,11 +82,10 @@ export async function cancelClaudeStructuredTurn(input: {
// pending handover is no reason to hold it; the queue sweep settles that follow-up. With nothing
// live, a written follow-up whose turn has not opened is one the naming client has not seen start.
if (
!prompt &&
(requestedTurnId === undefined ||
(session.translator?.commandTurnId !== requestedTurnId &&
(liveTurnId === requestedTurnId ||
(liveTurnId === null && session.dispatchWaiters.length > 0))))
requestedTurnId === undefined ||
(session.translator?.commandTurnId !== requestedTurnId &&
(liveTurnId === requestedTurnId ||
(liveTurnId === null && session.dispatchWaiters.length > 0)))
) {
return cancelClaudeConversation(
session,
@@ -122,21 +95,6 @@ export async function cancelClaudeStructuredTurn(input: {
input.onDispatchSettledLate
)
}
if (requestedTurnId === undefined) {
return { cancelled: false }
}
if (prompt && session.fence !== request.fence) {
return { cancelled: false }
}
const claim = prompt ? session.prompts.claimBound(prompt.itemId, requestedTurnId) : null
if (prompt && !claim) {
return { cancelled: false }
}
const cancellationObserved = claim ? session.prompts.observeCancellation(claim) : null
if (claim && !cancellationObserved) {
session.prompts.releaseClaim(claim)
return { cancelled: false }
}
// Judge against the published journal, because that is the only turn a client could have been
// shown — but only while it HAS an answer. The journal drains through a serialized async queue,
// so a null read means the row has not landed yet, not that nothing is running; falling back to
@@ -157,55 +115,27 @@ export async function cancelClaudeStructuredTurn(input: {
![...session.dispatchWaiters, ...session.retiredDispatchWaiters].some(
(waiter) => waiter.dispatchSequence === session.dispatchSequence
)
// Prompt cancellation has a separate callback-settlement contract, so only a provider with
// cancelQueued can release its uncertain queued send.
const dispatchAdmissionAllowsCancellation = (): boolean =>
dispatchAdmissionIsCurrent() ||
(Boolean(prompt) && supportsClaudeQueuedInterruptCancellation(session))
const compactionOwnsTurn = (): boolean =>
session.translator !== null && session.translator.commandTurnId === requestedTurnId
const isCurrent = (): boolean =>
sessions.get(request.sessionId) === session &&
session.fence === request.fence &&
session.acquisitionGeneration === acquisitionGeneration &&
(claim && prompt
? ownsRequestedTurn() &&
session.prompts.ownsBoundClaim(claim, prompt.itemId, requestedTurnId) &&
dispatchAdmissionAllowsCancellation()
: compactionOwnsTurn() || (ownsRequestedTurn() && dispatchAdmissionAllowsCancellation()))
let interruptConfirmed = false
try {
// A turn is cancelled only at a client's request, so the stop is the user's.
const result = await cancelClaudeTurn(
session,
timeoutMs,
() => {
const current = isCurrent()
// Read with the result that ends it: a stopped command reports no compaction.
if (current && compactionOwnsTurn()) {
session.translator?.commandInterruptRequested(requestedTurnId)
}
return current
},
input.onDispatchSettledLate,
{ turnId: requestedTurnId, cause: 'user-stop' }
)
if (result.cancelled && claim && cancellationObserved) {
interruptConfirmed = true
await waitForClaudePromptCancellation(cancellationObserved, timeoutMs)
if (!input.admitPromptCancellation(session, claim.found.prompt.promptKey)) {
throw new Error(`Claude prompt cancellation lifecycle was not admitted for ${claim.itemId}`)
// A turn is cancelled only at a client's request, so the stop is the user's.
return cancelClaudeTurn(
session,
timeoutMs,
() => {
const current =
sessions.get(request.sessionId) === session &&
session.fence === request.fence &&
session.acquisitionGeneration === acquisitionGeneration &&
(compactionOwnsTurn() || (ownsRequestedTurn() && dispatchAdmissionIsCurrent()))
// Read with the result that ends it: a stopped command reports no compaction.
if (current && compactionOwnsTurn()) {
session.translator?.commandInterruptRequested(requestedTurnId)
}
} else if (claim) {
session.prompts.releaseClaim(claim)
}
return result
} catch (error) {
if (claim && !interruptConfirmed) {
session.prompts.releaseClaim(claim)
}
throw error
}
return current
},
input.onDispatchSettledLate,
{ turnId: requestedTurnId, cause: 'user-stop' }
)
}
function prepareClaudePromptReply(
@@ -252,3 +182,41 @@ export async function answerClaudeStructuredPrompt(input: {
throw error
}
}
/** The host records the dismissal (`commit`) while the claim is held. A Stop that ends the child
* leaves the request to end with it; only `answer` declines it, so Claude never gets a reply racing
* that Stop's interrupt. */
export async function dismissClaudeStructuredPrompt(input: {
request: Parameters<NonNullable<StructuredAgentSessionAdapterStop['dismissPrompt']>>[0]
sessions: Map<string, ClaudeSession>
}): Promise<void> {
const { request, sessions } = input
const session = sessions.get(request.sessionId)
const claim = session?.fence === request.fence ? session.prompts.claim(request.itemId) : null
if (!session || !claim) {
throw new AgentSessionPromptUnavailableError(request.itemId)
}
const { promptKey } = claim.found.prompt
const journalPrompts = session.translator?.journalPrompts
// Before the commit: Claude's own cancel of the request, landing while the host writes, must
// not write after it.
const handBack = journalPrompts?.handOver(promptKey)
try {
try {
await request.commit()
} catch (error) {
// Nothing recorded the card: it is Claude's again, and a withdrawal Claude made meanwhile,
// which wrote nothing then, closes it now.
handBack?.()
if (!session.prompts.find(request.itemId)) {
journalPrompts?.cancel(promptKey)
}
throw error
}
if (request.answer && session.prompts.ownsClaim(claim)) {
await answerClaudePrompt(session, claim, claudePromptDismissal(claim.found.prompt))
}
} finally {
session.prompts.releaseClaim(claim)
}
}
@@ -3,6 +3,7 @@ import type {
AgentSessionPromptResponse,
AgentSessionQuestionAnswer
} from '../../shared/agent-session-question-answer'
import type { StructuredAgentSessionAdapterStop } from '../native-chat/agent-session-wire/structured-agent-session-adapter-stop'
import {
claudePromptQuestions,
isClaudePromptRecord,
@@ -18,9 +19,32 @@ export {
type ClaudePromptSettle
} from './claude-prompt-registry'
/** No card offers `cancel` any more; an older card's "Stop" option answers it as a dismissal. */
export const CLAUDE_APPROVAL_DECISIONS = ['allow', 'allowForSession', 'deny', 'cancel'] as const
export type ClaudeApprovalDecision = (typeof CLAUDE_APPROVAL_DECISIONS)[number]
/** A card's own Cancel: an approval is dismissed (a tool's with the reply its Deny sends, a plan's
* asking Claude to wait for the user), and a question ends the way the chat's Stop does. Nothing
* on a card interrupts the turn and leaves the child running. */
export const claudePromptCancelRoute: NonNullable<
StructuredAgentSessionAdapterStop['routePromptCancel']
> = ({ prompt }) => (prompt.kind === 'question' ? { kind: 'stop' } : { kind: 'dismiss' })
/** Declines a request the user dismissed. A dismissed plan or question waits on the user: Claude is
* told to end its turn rather than revise or ask again. */
export function claudePromptDismissal(prompt: ClaudePendingPrompt): PermissionResult {
return {
behavior: 'deny',
message:
prompt.kind === 'question'
? 'The user dismissed these questions without answering. End your turn and wait for them.'
: prompt.subject?.kind === 'plan'
? 'The user dismissed this plan without approving it. End your turn and wait for them to say what to change.'
: 'User denied this action.',
toolUseID: prompt.toolUseId
}
}
function isClaudeApprovalDecision(optionId: string): optionId is ClaudeApprovalDecision {
return CLAUDE_APPROVAL_DECISIONS.some((decision) => decision === optionId)
}
@@ -93,15 +117,15 @@ function approvalResponse(prompt: ClaudePendingPrompt, optionId: string): Permis
toolUseID: prompt.toolUseId
}
}
if (decision === 'cancel') {
return claudePromptDismissal(prompt)
}
return {
behavior: 'deny',
message:
decision === 'cancel'
? 'User stopped this turn.'
: prompt.subject?.kind === 'plan'
? 'The user asked you to keep planning. Revise the plan and call ExitPlanMode again.'
: 'User denied this action.',
...(decision === 'cancel' ? { interrupt: true } : {}),
prompt.subject?.kind === 'plan'
? 'The user asked you to keep planning. Revise the plan and call ExitPlanMode again.'
: 'User denied this action.',
toolUseID: prompt.toolUseId
}
}
@@ -2,6 +2,7 @@ import type {
ClaudeSession,
ClaudeStructuredSessionAdapterDeps
} from './claude-structured-session-state'
import { settleClaudeTurnEndWaiters } from './claude-request-end-wait'
/**
* Record a completed turn: its leaf becomes the one close and exit persist, and the durable point
@@ -17,6 +18,8 @@ export function persistClaudeTurnResumePoint(
return
}
session.turnEndLeafUuid = session.leafUuid
// The turn ended: a Stop waiting to end the child re-reads what Claude still has in flight.
settleClaudeTurnEndWaiters(session)
const leafUuid = session.turnEndLeafUuid
const persist = deps.persistResumePoint
if (!persist || leafUuid === null || session.resumePointWrite?.leafUuid === leafUuid) {
@@ -9,6 +9,7 @@ import { openClaudeStreamJsonConnection } from './claude-stream-json-connection'
import { buildClaudePermissionCallbacks } from './claude-structured-inbound-control'
import { resolveClaudeReplayTurn } from './claude-replay-turn-resolution'
import { claudeSessionStateEndsTurn } from './claude-session-state-turn-over'
import { settleClaudeTurnEndWaiters } from './claude-request-end-wait'
import {
readClaudeCapabilities,
readClaudeFrameString,
@@ -120,6 +121,7 @@ export async function acquireClaudeSession({
}
// The CLI's idle releases its doubted sends; a late echo still accepts one it goes on to run.
if (claudeSessionStateEndsTurn(message) && sessions.get(sessionId) === liveSession) {
settleClaudeTurnEndWaiters(liveSession)
deps.onSessionIdle?.({ sessionId })
}
}
@@ -153,7 +155,6 @@ export async function acquireClaudeSession({
const { canUseTool, onUserDialog } = buildClaudePermissionCallbacks({
sessionId,
prompts,
currentTurnId: () => translator?.currentTurnId ?? null,
emit: (event) =>
callbacks.deliver(attempt, sessionId, () => callbacks.emit(liveSession, input.events, event))
})
@@ -28,6 +28,7 @@ import {
type ClaudeStructuredSessionEvent
} from './claude-structured-session-state'
import { closeAllClaudeSessions, closeClaudeSession } from './claude-structured-session-close'
import { claudeStoppedRequestEndWait } from './claude-request-end-wait'
import {
drainClaudeObservedExits,
observeClaudeSessionExit,
@@ -38,10 +39,11 @@ import type { AgentSessionBackgroundTaskState } from '../../shared/agent-session
import { resolveClaudeProviderHistoryWindow } from './claude-structured-history-window'
import { drainClaudeChildWork } from './claude-child-work-evidence'
import {
admitClaudePromptCancellation,
answerClaudeStructuredPrompt,
cancelClaudeStructuredTurn
cancelClaudeStructuredTurn,
dismissClaudeStructuredPrompt
} from './claude-structured-prompt-ownership'
import { claudePromptCancelRoute } from './claude-structured-prompt-replies'
export type { ClaudeStructuredLaunch } from './claude-structured-launch-resolution'
export type {
@@ -174,12 +176,7 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
}
bindPromptItemId(sessionId: string, journalItemId: string, promptKey: string): void {
const session = this.sessions.get(sessionId)
session?.prompts.bindJournalItemId(
journalItemId,
promptKey,
session.translator?.currentTurnId ?? null
)
this.sessions.get(sessionId)?.prompts.bindJournalItemId(journalItemId, promptKey)
}
dispatch: StructuredAgentSessionAdapter['dispatch'] = (input) =>
@@ -192,12 +189,17 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
cancelClaudeStructuredTurn({
request,
sessions: this.sessions,
admitPromptCancellation: (session, promptKey) =>
admitClaudePromptCancellation(session, promptKey),
onDispatchSettledLate: (settlement) =>
this.deps.onDispatchSettledLate?.({ sessionId: request.sessionId, ...settlement }),
...(this.deps.requestTimeoutMs === undefined ? {} : { timeoutMs: this.deps.requestTimeoutMs })
})
// Stop is a session boundary for Claude: an interrupt can answer while background work keeps the
// CLI running, and a refused one leaves the turn running.
stopEndsSession = (): boolean => true
awaitStoppedRequestEnd = claudeStoppedRequestEndWait(this.sessions)
routePromptCancel = claudePromptCancelRoute
dismissPrompt: StructuredAgentSessionAdapter['dismissPrompt'] = (request) =>
dismissClaudeStructuredPrompt({ request, sessions: this.sessions })
stopBackgroundTasks: NonNullable<StructuredAgentSessionAdapter['stopBackgroundTasks']> = async (
input
) => {
@@ -131,8 +131,8 @@ describe('Claude published session close lifecycle', () => {
expect(events.filter((event) => event.type === 'ended')).toHaveLength(1)
expect(events.filter((event) => event.type === 'handle')).toHaveLength(0)
expect(disposeTranslator).toHaveBeenCalledOnce()
// The host's child records hear the session end, which settles the child still running.
expect(childWork).toEqual(['live', 'session-ended'])
// A close Orca asked for stops the child still running, then the host hears the session end.
expect(childWork).toEqual(['live', 'ended', 'session-ended'])
await expect(adapter.closeSession('session-1')).resolves.toBe(true)
expect(persistHandle).toHaveBeenCalledTimes(2)
@@ -18,6 +18,7 @@ import type { ClaudePromptRegistry } from './claude-structured-prompt-replies'
import { closeProcessRegistry } from '../../shared/child-process/close-process-registry'
import { retireClaudeDispatchWaiters } from './claude-structured-dispatch'
import { settledClaudeTurnEndLeaf } from './claude-structured-resume-point'
import { settleClaudeTurnEndWaiters } from './claude-request-end-wait'
/** The root's own exit was seen first-hand. The lease follows the root, so a descendant
* left unverified or seen alive does not hold it. */
@@ -68,6 +69,7 @@ export function settleClaudeExitedSession(session: ClaudeSession): void {
// The child is gone, so no replay can start these turns. Nothing else ends a
// waiter's life now that no deadline does.
retireClaudeDispatchWaiters(session)
settleClaudeTurnEndWaiters(session)
for (const prompt of session.prompts.clear()) {
prompt.settle(null)
}
@@ -93,6 +95,7 @@ async function finalizeClaudePublishedSession(
session: ClaudeSession
): Promise<boolean> {
retireClaudeDispatchWaiters(session)
settleClaudeTurnEndWaiters(session)
// Settle every in-flight permission callback so closing leaves no dangling promise; `null`
// writes no response, and the SDK ignores any post-cleanup answer regardless.
for (const prompt of session.prompts.clear()) {
@@ -115,6 +118,11 @@ async function finalizeClaudePublishedSession(
rootExitVerdict = cleanupError
}
// Queues the session's ending for the host's child records; the adapter delivers it after close.
// A close that proved the whole tree gone stopped what still ran. One that saw a descendant
// survive, like an exit of the session's own, leaves how it ended unknown.
if (connectionClosed === true) {
session.childWork.stopLive()
}
session.childWork.clear()
session.backgroundTasks.clear()
const leafUuid = await settledClaudeTurnEndLeaf(session)
@@ -334,18 +334,6 @@ describe("a user's Stop inside a live turn", () => {
expect(providerRows(state.items)).toBe(1)
})
it('forgets a Stop the CLI refused', () => {
const state = sinkState()
const translator = createClaudeJournalTranslator({ sink: state.sink })
translator.handle(userTurn('user-1'))
translator.recordTurnStop('user-1', 'user-stop')
translator.withdrawTurnStop('user-1')
translator.handle({ type: 'message', sessionId: 'orca-session', message: cutShort })
expect(settledTurn(state.items, 'user-1')).toMatchObject({ outcome: 'failure' })
})
it('does not read a host stop as the user asking', () => {
const state = sinkState()
const translator = createClaudeJournalTranslator({ sink: state.sink })
@@ -93,10 +93,13 @@ function sessionHoldingTurn(turnId: string | null): ReturnType<typeof sessionFor
session.translator = {
handle: vi.fn(),
openTurnInLiveProviderCycle: false,
journalPrompts: { cancel: vi.fn(), resolve: vi.fn() },
journalPrompts: {
resolve: vi.fn(),
handOver: () => () => {},
cancel: () => ({ accepted: true })
},
currentTurnId: turnId,
recordTurnStop: () => true,
withdrawTurnStop: () => {},
commandTurnId: null,
beginCommand: vi.fn(),
forgetCommand: vi.fn(),
@@ -133,8 +136,7 @@ function cancellationOf(
): Promise<{ cancelled: boolean }> {
return cancelClaudeStructuredTurn({
request,
sessions: new Map([['session-1', session]]),
admitPromptCancellation: () => true
sessions: new Map([['session-1', session]])
})
}
@@ -1,5 +1,8 @@
import { describe, expect, it, vi } from 'vitest'
import type { AgentSessionJournalIdentity } from '../../../shared/agent-session-journal-types'
import type {
AgentJournalQuestionItem,
AgentSessionJournalIdentity
} from '../../../shared/agent-session-journal-types'
import type {
AgentSessionAcquisition,
StructuredAgentSessionAdapter
@@ -167,6 +170,52 @@ describe('StructuredAgentSessionAdapterRouter optional lifecycle methods', () =>
)
})
const QUESTION: AgentJournalQuestionItem = {
kind: 'question',
question: 'Which branch?',
options: [{ id: 'main', label: 'main' }],
resolution: { state: 'pending', selectedOptionId: null, resolvedBy: null, resolvedAt: null }
}
describe('StructuredAgentSessionAdapterRouter.stopEndsSession', () => {
it("answers for the session's live owner, and keeps the child with none", async () => {
const claude = adapterOf(vi.fn(async () => true))
claude.stopEndsSession = () => true
const codex = adapterOf(vi.fn(async () => false))
const router = new StructuredAgentSessionAdapterRouter({ claude, codex }, async () => {})
expect(router.stopEndsSession('session-1')).toBe(false)
await router.acquire({ identity: claudeIdentity('session-1'), fence: 1, spawnToken: 'spawn-1' })
expect(router.stopEndsSession('session-1')).toBe(true)
})
it("waits on the session's live owner to wind the Stop down, and on nothing with none", async () => {
const claude = adapterOf(vi.fn(async () => true))
claude.awaitStoppedRequestEnd = vi.fn(async () => undefined)
const codex = adapterOf(vi.fn(async () => false))
const router = new StructuredAgentSessionAdapterRouter({ claude, codex }, async () => {})
await router.awaitStoppedRequestEnd('session-1', 5)
expect(claude.awaitStoppedRequestEnd).not.toHaveBeenCalled()
await router.acquire({ identity: claudeIdentity('session-1'), fence: 1, spawnToken: 'spawn-1' })
await router.awaitStoppedRequestEnd('session-1', 5)
expect(claude.awaitStoppedRequestEnd).toHaveBeenCalledWith('session-1', 5)
})
it("answers a card's Cancel as the session's live owner does, and leaves it to cancelTurn with none", async () => {
const claude = adapterOf(vi.fn(async () => true))
claude.routePromptCancel = () => ({ kind: 'stop' })
const codex = adapterOf(vi.fn(async () => false))
const router = new StructuredAgentSessionAdapterRouter({ claude, codex }, async () => {})
expect(router.routePromptCancel({ sessionId: 'session-1', prompt: QUESTION })).toBeUndefined()
await router.acquire({ identity: claudeIdentity('session-1'), fence: 1, spawnToken: 'spawn-1' })
expect(router.routePromptCancel({ sessionId: 'session-1', prompt: QUESTION })).toEqual({
kind: 'stop'
})
})
})
describe('StructuredAgentSessionAdapterRouter.closeAll', () => {
it('refuses to acquire once the global close proof is published', async () => {
const acquire = vi.fn(async ({ fence, spawnToken }) => acquisition(fence, spawnToken))
@@ -119,6 +119,18 @@ export class StructuredAgentSessionAdapterRouter implements StructuredAgentSessi
holdsDispatch = (sessionId: string): boolean =>
this.liveOwnerOrNull(sessionId)?.holdsDispatch?.(sessionId) ?? false
stopEndsSession = (sessionId: string): boolean =>
this.liveOwnerOrNull(sessionId)?.stopEndsSession?.(sessionId) ?? false
awaitStoppedRequestEnd = async (sessionId: string, stoppedAt: number) =>
this.liveOwnerOrNull(sessionId)?.awaitStoppedRequestEnd?.(sessionId, stoppedAt)
routePromptCancel: NonNullable<StructuredAgentSessionAdapter['routePromptCancel']> = (input) =>
this.liveOwnerOrNull(input.sessionId)?.routePromptCancel?.(input)
dismissPrompt: NonNullable<StructuredAgentSessionAdapter['dismissPrompt']> = async (input) =>
this.owner(input.sessionId).dismissPrompt?.(input)
readCommands: NonNullable<StructuredAgentSessionAdapter['readCommands']> = (sessionId) =>
this.liveOwnerOrNull(sessionId)?.readCommands?.(sessionId)
@@ -0,0 +1,37 @@
// 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.
import type {
AgentJournalApprovalItem,
AgentJournalQuestionItem
} from '../../../shared/agent-session-journal-types'
/** Where a card's Cancel goes: a dismissal (`dismissPrompt`), or the chat's Stop. */
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. */
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
* `stoppedAt` (when the interrupt went out), runs out. */
awaitStoppedRequestEnd?(sessionId: string, stoppedAt: number): Promise<void>
/** Where the pending card's own Cancel goes. Undefined: `cancelTurn` with the prompt. */
routePromptCancel?(input: {
sessionId: string
prompt: AgentJournalApprovalItem | AgentJournalQuestionItem
}): AgentSessionPromptCancelRoute | undefined
/** The user dismissed the pending card; `commit` records it as cancelled, with the claim held.
* `answer` then declines the provider's request; without it the request is left to end with the
* child a Stop ends. Either way the provider records nothing more for the card. */
dismissPrompt?(input: {
sessionId: string
itemId: string
fence: number
answer: boolean
commit: () => Promise<void>
}): Promise<void>
}
@@ -40,6 +40,7 @@ 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'
export type {
StructuredAgentSessionChildEndCause,
StructuredAgentSessionStopCause
@@ -250,7 +251,7 @@ export type AgentSessionCancelOutcome = {
refusal?: { detail?: ProviderDiagnostic }
}
export type StructuredAgentSessionAdapter = {
export type StructuredAgentSessionAdapter = StructuredAgentSessionAdapterStop & {
/** Provider-aware capability check for hosts that route more than one adapter. */
supportsCreate?(location: AgentSessionExecutionLocation, agent: string): boolean
/** Provider/runtime support, kept here so remote enablement changes adapter data, not UI logic. */
@@ -0,0 +1,116 @@
// 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.
import { agentSessionFailureWords } from '../../../shared/agent-session-failure-words'
import { agentSessionFailureFact } from '../../../shared/agent-session-failure'
import type {
AgentSessionCancelResult,
AgentSessionMutationEnvelope,
AgentSessionMutationResult
} from '../../../shared/agent-session-wire'
import {
mutateStructuredAgentSession,
type StructuredAgentSessionMutationContext
} from './structured-agent-session-mutation-context'
import type { StructuredAgentSessionCaller } from './structured-agent-session-host-types'
import type { MutationPlan } from './structured-agent-session-mutation-plans'
import { runStopWithQueuePause } from './structured-agent-session-queued-stop'
import {
openForWrite,
structuredAgentSessionFailureWordsContext
} from './structured-agent-session-send-preparation'
import {
endStoppedStructuredAgentSession,
isMainAgentWorkingOnceFlushed,
performCancel,
type StructuredAgentSessionStopWindDown
} from './structured-agent-session-turns-cancel'
import type { AgentSessionTurnContext, TurnOutcome } from './structured-agent-session-turns'
type ChatStopOutcome = TurnOutcome<AgentSessionCancelResult>
/** What the chat's Stop did, and whether its next step ends the provider's session. */
export type StructuredAgentSessionChatStopRun = { outcome: ChatStopOutcome; endsSession: boolean }
/** Runs `plan` with `run`, which may call `stop` for the chat's Stop. */
export function mutateWithChatStop<TValue>(
context: StructuredAgentSessionMutationContext,
caller: StructuredAgentSessionCaller,
params: { envelope: AgentSessionMutationEnvelope; turnId?: string },
plan: MutationPlan<TValue>,
run: (
ctx: AgentSessionTurnContext,
stop: () => Promise<StructuredAgentSessionChatStopRun>
) => Promise<TurnOutcome<TValue>>
): Promise<AgentSessionMutationResult<TValue>> {
const { envelope, turnId } = params
const { sessionId } = envelope
// Set by the Stop's step only when its provider's session ends; a replay leaves it unset.
let windDown: StructuredAgentSessionStopWindDown | undefined
const named = turnId !== undefined ? { turnId } : {}
// Stop's queue step, the same for every client: once the Stop takes effect the queue is paused.
// The cards stay published; nothing is withdrawn and no text ever rides the answer.
const stop = (ctx: AgentSessionTurnContext): Promise<ChatStopOutcome> =>
runStopWithQueuePause(ctx, async (tookEffect) => {
// Stop withdraws every queued SUBMISSION first, whatever the start or the child is doing.
const withdrawn = await ctx.journal.rejectQueuedSubmissions(
ctx.fence,
agentSessionFailureWords(agentSessionFailureFact('cancelled'), { surface: 'rejection' })
)
const child = context.sessions.get(ctx.sessionId)?.child
if (child?.phase === 'starting') {
// A start that may never land is the one thing here Stop has to end; the chat stays.
await tookEffect()
await context.stopAgent(ctx.sessionId)
return { ok: true, value: { ...named, cancelled: true } }
}
// A Stop naming no turn ends nothing more unless the session reads working, by the rule
// every session list and the chat's own Stop read it.
const inFlight = turnId !== undefined || (await isMainAgentWorkingOnceFlushed(ctx))
const record = context.deps.store.getRecord(ctx.sessionId)
if (!child || !inFlight) {
if (withdrawn.length > 0) {
await tookEffect()
}
return { ok: true, value: { ...named, cancelled: withdrawn.length > 0 } }
}
await tookEffect()
return performCancel(
{ ...ctx, failureTextContext: structuredAgentSessionFailureWordsContext(record) },
{
clientOperationId: envelope.clientOperationId,
...named,
stopChild: () => context.stopAgent(sessionId),
endSession: (owed) => {
windDown = owed
},
withdrewQueued: withdrawn.length > 0
}
)
})
const result = mutateStructuredAgentSession(
context,
caller,
envelope,
{
...plan,
run: (ctx) =>
run(ctx, async () => ({ outcome: await stop(ctx), endsSession: windDown !== undefined }))
},
openForWrite(context, envelope)
)
// Queued in the mutation's own tick, so a send made meanwhile lands behind the child's end.
void context.serialize(sessionId, async () => {
if (windDown) {
await endStoppedStructuredAgentSession(
{ sessionId, adapter: context.deps.adapter },
windDown,
() => context.stopAgent(sessionId),
(error) => context.deps.onEventSinkError?.({ sessionId, error })
)
}
})
return result
}
@@ -208,38 +208,34 @@ it('ends a hung /compact Claude will not interrupt by stopping it, and answers t
})
})
it('ends a stopped /compact on its own interrupted result, and a late copy of that result does not end the next one', async () => {
it('ends a stopped /compact on its own interrupted result, then ends its child; a late copy never reaches the next one', async () => {
claude.routes.interrupt = () => ({})
const first = await compact()
const { connection, uuid: firstUuid } = await sent('/compact')
await expect(stop(first)).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
// The interrupt was taken: the command runs until Claude answers it.
expect((await commandState(first))?.state).toBe('running')
// The interrupt was taken: the child waits for Claude to end the command itself.
expect(connection.closed).toBe(false)
frame(connection, result(firstUuid, INTERRUPTED))
await vi.waitFor(async () =>
expect(await commandState(first)).toMatchObject({
state: 'interrupted',
outcome: 'cancellation'
})
)
await vi.waitFor(() => expect(connection.closed).toBe(true))
expect(await commandState(first)).toMatchObject({
state: 'interrupted',
outcome: 'cancellation'
})
const second = await compact()
await vi.waitFor(() =>
expect(
connection.sent.filter((message) => JSON.stringify(message).includes('/compact'))
).toHaveLength(2)
)
const secondUuid = String(connection.sent.at(-1)?.uuid)
const next = await vi.waitFor(() => {
const started = claude.connections.at(-1)
expect(started).not.toBe(connection)
expect(started?.sent.some((message) => JSON.stringify(message).includes('/compact'))).toBe(true)
return started!
})
const secondUuid = String(next.sent.at(-1)?.uuid)
frame(connection, result(firstUuid, INTERRUPTED))
frame(
connection,
result('an-earlier-send', { subtype: 'error_during_execution', is_error: true })
)
expect((await commandState(second))?.state).toBe('running')
frame(connection, { type: 'system', subtype: 'compact_boundary', uuid: 'boundary' })
frame(connection, result(secondUuid))
frame(next, { type: 'system', subtype: 'compact_boundary', uuid: 'boundary' })
frame(next, result(secondUuid))
await vi.waitFor(async () =>
expect(await commandState(second)).toMatchObject({ state: 'completed', outcome: 'success' })
)
@@ -0,0 +1,870 @@
// A Stop ends Claude's child on the shipping adapter, whatever Claude answered the interrupt. The
// Stop answers on the interrupt; its next serialized step ends the child once Claude has wound down
// what it had in flight, or the grace runs out. The chat rests; the next send resumes it.
import { mkdtemp, rm } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, beforeEach, expect, it, vi } from 'vitest'
import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope'
import { readAgentJournalTurn } from '../../../shared/agent-session-turn-record'
import { DISPATCH_REJECTED_CANCELLED } from '../../../shared/structured-agent-session-dispatch-rejection'
import { activeStructuredAgentSessionTurnId } from '../../../shared/structured-agent-session-live-turn'
import { projectStructuredAgentSessionStatusState } from '../../../shared/structured-agent-session-projection'
import { structuredAgentSessionAgentStatus } from '../../../shared/structured-agent-session-agent-status'
import type { AgentStatusStructuredSessionSubject } from '../../../shared/agent-status-subject'
import { AgentHookServer } from '../../agent-hooks/server'
import {
ClaudeControlRequestError,
runClaudeControl
} from '../../claude/claude-agent-sdk-control-requests'
import { CLAUDE_STOP_GRACE_MS } from '../../claude/claude-request-end-wait'
import { ClaudeStructuredSessionAdapter } from '../../claude/claude-structured-session-adapter'
import type { ClaudeStructuredSessionEvent } from '../../claude/claude-structured-session-state'
import {
fakeClaude,
PROVIDER_SESSION_ID,
type FakeConnection
} from '../../claude/claude-structured-session-test-support'
import { invokeCanUseTool } from '../../claude/claude-can-use-tool-test-support'
import type { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
import { openTestAgentSessionRecordStore } from '../../runtime/agent-session-record-store-test-harness'
import { structuredClaudeLifecycleEvent } from '../../runtime/structured-claude-runtime-adapter'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { StructuredAgentSessionHost } from './structured-agent-session-host'
import {
HOST_TEST_NOW as NOW,
HOST_TEST_SESSION as SESSION,
hostTestAttachParams,
hostTestMessage,
hostTestOperationId,
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
const CALLER = { callerKey: 'client-1' }
// As Claude Code 2.1.280 advertises them on a turn's system/init frame.
const CAPABILITIES = ['interrupt_receipt_v1', 'interrupt_cancel_queued_v1', 'msg_lifecycle_v1']
let root: string
let host: StructuredAgentSessionHost
let adapter: ClaudeStructuredSessionAdapter
let store: AgentSessionRecordStore
let queued: string[]
let claude: ReturnType<typeof fakeClaude>
let events: ClaudeStructuredSessionEvent[]
let sinkErrors: unknown[]
// The host's status row and child records, as the app's hook server holds them.
let server: AgentHookServer
let statusSubject: AgentStatusStructuredSessionSubject | undefined
beforeEach(async () => {
root = await mkdtemp(join(tmpdir(), 'orca-claude-stop-ends-session-'))
resetHostTestOperationIds()
queued = []
events = []
sinkErrors = []
server = new AgentHookServer()
statusSubject = undefined
claude = fakeClaude({
replayUuid: null,
routes: {
interrupt: (params) =>
params?.cancelQueued
? { still_queued: [], cancelled: queued.splice(0) }
: { still_queued: [...queued] }
}
})
const lifecycle: Promise<void>[] = []
adapter = new ClaudeStructuredSessionAdapter({
resolveLaunch: async () => ({
pathToClaudeCodeExecutable: 'claude',
options: {},
cwd: root,
claudeConfigDir: join(root, 'claude-home'),
providerSessionId: PROVIDER_SESSION_ID,
resumeLeafUuid: null,
resumesTranscript: (store.getRecord(SESSION)?.providerHandleChain.length ?? 0) > 0,
continuesChain: (store.getRecord(SESSION)?.providerHandleChain.length ?? 0) > 0
}),
onEvent: (event) => {
events.push(event)
const mapped = structuredClaudeLifecycleEvent(event)
if (mapped) {
lifecycle.push(host.handleAdapterEvent(mapped))
}
},
onDispatchSettledLate: (settlement) => void host.settleLateDispatch(settlement),
onChildWorkEvidence: (sessionId, evidence) =>
host.publishChildWorkEvidence(sessionId, evidence),
openConnection: claude.openConnection,
readProcessStartTime: async () => 1_700_000_000_000,
now: () => NOW
})
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
store,
adapter: Object.assign(adapter, { supportsCreate: () => true }),
journalDatabase: openTestJournalHostDatabase(root),
claimKeyId: 'key-1',
mintSpawnToken: () => 'spawn-a',
onEventSinkError: ({ error }) => sinkErrors.push(error),
statusSink: {
publish: (summary, subject) => {
statusSubject = subject
server.ingestStructuredStatus(summary, subject)
},
forget: (subject) => server.dropStructuredStatus(subject),
publishChildWork: (subject, evidence, provider) =>
server.ingestStructuredChildWork(subject, evidence, provider),
readChildWork: (subject) => server.getStructuredChildWorkViews(subject)
},
now: () => NOW
})
const params = hostTestAttachParams(null, {
provider: 'claude',
agent: 'claude',
accountHome: { variable: 'CLAUDE_CONFIG_DIR', path: join(root, 'claude-home') },
providerHandle: { kind: 'claude', sessionId: PROVIDER_SESSION_ID, leafUuid: null }
})
expect(await host.attach(CALLER, params)).toMatchObject({ ok: true })
await adapter.awaitStarted(SESSION)
await Promise.all(lifecycle)
})
afterEach(async () => {
await adapter.closeAll()
await host.flushAllStreamedEvents()
await rm(root, { recursive: true, force: true })
})
function eventually<T>(assertion: () => T | Promise<T>): Promise<T> {
return vi.waitFor(assertion, { timeout: 10_000 })
}
function envelope(
method: 'agentSession.send' | 'agentSession.cancel',
fields: Record<string, unknown>,
fence = store.getRecord(SESSION)!.lease.runtimeFence
) {
return {
sessionId: SESSION,
clientOperationId: hostTestOperationId(),
expectedRuntimeFence: fence,
payloadFingerprint: computeAgentSessionPayloadFingerprint({
method,
sessionId: SESSION,
fields
})
}
}
async function send(text: string, fence?: number): Promise<string> {
const body = hostTestMessage(text)
const sent = await host.send(CALLER, {
envelope: envelope('agentSession.send', { body }, fence),
body
})
if (!sent.ok) {
throw new Error('send refused')
}
return sent.value.clientMessageId
}
async function dispatch(clientMessageId: string) {
const submission = (await host.journalSnapshot(SESSION)).submissions.find(
(entry) => entry.clientMessageId === clientMessageId
)
return { state: submission?.dispatchState, reason: submission?.reason }
}
function frame(connection: FakeConnection, message: Record<string, unknown>): void {
connection.handlers.onMessage?.({ session_id: PROVIDER_SESSION_ID, ...message })
}
/** Sends a message and lets Claude open its turn and write one reply; returns the turn's id. */
async function openTurn(connection: FakeConnection, text = 'Write a long reply.'): Promise<string> {
const clientMessageId = await send(text)
await eventually(() =>
expect(connection.sent.some((message) => JSON.stringify(message).includes(text))).toBe(true)
)
frame(connection, {
type: 'system',
subtype: 'init',
uuid: `init-${text}`,
model: 'claude-sonnet-5',
capabilities: CAPABILITIES
})
const written = connection.sent.at(-1)!
frame(connection, { ...written, uuid: written.uuid })
frame(connection, {
type: 'assistant',
uuid: 'stopped-turn-leaf',
parent_tool_use_id: null,
message: { id: 'msg-1', role: 'assistant', content: [{ type: 'text', text: 'Working on' }] }
})
await eventually(async () => expect((await dispatch(clientMessageId)).state).toBe('accepted'))
const turnId = activeStructuredAgentSessionTurnId((await host.journalSnapshot(SESSION)).items)
expect(turnId).not.toBeNull()
return turnId!
}
function stop(turnId?: string) {
const fields = turnId === undefined ? {} : { turnId }
return host.cancel(CALLER, { envelope: envelope('agentSession.cancel', fields), ...fields })
}
/** Resolves once everything queued on the session's lane so far has run: a Stop's second step. */
function laneDrained(): Promise<void> {
return host['tasks'].serialize(SESSION, async () => {})
}
function wrote(connection: FakeConnection, text: string): boolean {
return connection.sent.some((message) => JSON.stringify(message).includes(text))
}
const INTERRUPTED_RESULT = {
type: 'result',
subtype: 'error_during_execution',
is_error: true,
terminal_reason: 'aborted_streaming',
uuid: 'interrupted-result'
}
// As the real control surface runs it: no answer ever comes, only the deadline Orca sets.
const NEVER_ANSWERS = (options?: Record<string, unknown>) =>
runClaudeControl(
'interrupt',
() => new Promise(() => {}),
typeof options?.timeoutMs === 'number' ? options.timeoutMs : undefined
)
async function interruptSent(connection: FakeConnection): Promise<void> {
await eventually(() =>
expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(true)
)
}
async function turnOutcome(): Promise<string | undefined> {
await host.flushStreamedEvents(SESSION)
const snapshot = await host.journalSnapshot(SESSION)
return readAgentJournalTurn(snapshot.items.findLast((item) => item.body.kind === 'turn')?.body)
?.outcome
}
async function statusTexts(): Promise<string[]> {
await host.flushStreamedEvents(SESSION)
return (await host.journalSnapshot(SESSION)).items.flatMap((item) =>
item.body.kind === 'status' ? [String(item.body.text)] : []
)
}
it('answers on the interrupt, ends the child once the stopped turn ends, and rests at that turn', async () => {
const connection = claude.connections[0]!
await openTurn(connection)
await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
// The Stop answered on the interrupt Claude took: the child waits for Claude to end the turn.
expect(connection.closed).toBe(false)
expect(await statusTexts()).toEqual(['Cancellation requested.'])
const ended = Date.now()
frame(connection, INTERRUPTED_RESULT)
await laneDrained()
// The turn's own end releases the wait; the grace is not waited out.
expect(Date.now() - ended).toBeLessThan(CLAUDE_STOP_GRACE_MS / 2)
expect(connection.closed).toBe(true)
expect(store.getRecord(SESSION)?.lease.claimStatus).toBe('released')
expect(await turnOutcome()).toBe('cancellation')
// The resume point is the stopped turn's own, so the next send continues after it.
expect(events.findLast((event) => event.type === 'handle')).toMatchObject({
type: 'handle',
providerSessionId: PROVIDER_SESSION_ID,
leafUuid: 'stopped-turn-leaf'
})
})
/** Sends a message Claude takes but has not echoed yet: no turn row, the chat still working. */
async function sendUnechoed(connection: FakeConnection): Promise<string> {
const clientMessageId = await send('Write a long reply.')
await eventually(() => expect(wrote(connection, 'Write a long reply.')).toBe(true))
return clientMessageId
}
function childRecords() {
return statusSubject ? server.getStructuredChildWorkViews(statusSubject) : []
}
async function agentStatus() {
await host.flushStreamedEvents(SESSION)
const { items, submissions } = await host.journalSnapshot(SESSION)
const { summary } = projectStructuredAgentSessionStatusState(
items,
submissions,
store.getRecord(SESSION)!.lease.runtimeFence
)
return summary.status
? structuredAgentSessionAgentStatus({
status: summary.status,
turnOutcome: summary.turnOutcome,
childWork: childRecords()
})
: null
}
it('reads a Stop pressed before Claude echoed the send as interrupted, not as a finished turn', async () => {
const connection = claude.connections[0]!
// As Claude winds down a request interrupted before its echo: echo, marker, aborted result, idle.
claude.routes.interrupt = () => {
setTimeout(() => {
const written = connection.sent.find((message) => message.type === 'user')!
frame(connection, { ...written, uuid: written.uuid })
frame(connection, {
type: 'user',
message: {
role: 'user',
content: [{ type: 'text', text: '[Request interrupted by user]' }]
},
parent_tool_use_id: null,
uuid: 'interrupted-marker'
})
frame(connection, INTERRUPTED_RESULT)
frame(connection, { type: 'system', subtype: 'session_state_changed', state: 'idle' })
}, 5)
return { still_queued: [], cancelled: [] }
}
const clientMessageId = await sendUnechoed(connection)
await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
await laneDrained()
expect(connection.closed).toBe(true)
expect(await dispatch(clientMessageId)).toMatchObject({ state: 'accepted' })
expect(await turnOutcome()).toBe('cancellation')
expect(await agentStatus()).toMatchObject({ mainAgent: { outcome: 'cancellation' } })
})
it('ends the child once the grace runs out when Claude says nothing after a Stop before the echo', async () => {
const connection = claude.connections[0]!
claude.routes.interrupt = () => ({ still_queued: [], cancelled: [] })
const clientMessageId = await sendUnechoed(connection)
const asked = Date.now()
await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
await laneDrained()
// As before the wait: the send Claude never answered is doubt once its child ends.
expect(connection.closed).toBe(true)
expect(Date.now() - asked).toBeLessThan(CLAUDE_STOP_GRACE_MS + 1_500)
expect(await dispatch(clientMessageId)).toMatchObject({ state: 'unknown' })
}, 15_000)
it('withdraws a follow-up Claude queued behind the turn before the child ends, never doubt', async () => {
const connection = claude.connections[0]!
await openTurn(connection)
const followUp = await send('And then this.')
await eventually(() => expect(connection.sent.at(-1)?.message).toBeDefined())
queued.push(String(connection.sent.at(-1)!.uuid))
await eventually(async () => expect((await dispatch(followUp)).state).toBe('pending'))
await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
// Withdrawn by Claude's own receipt while the child still runs, so it is never re-sent and never
// read as delivered.
expect(connection.closed).toBe(false)
expect(await dispatch(followUp)).toEqual({
state: 'rejected',
reason: DISPATCH_REJECTED_CANCELLED
})
frame(connection, INTERRUPTED_RESULT)
await laneDrained()
expect(connection.closed).toBe(true)
})
it('ends the child at once when Claude refuses the interrupt, and says only that the stop was asked', async () => {
claude.routes.interrupt = () => {
throw new ClaudeControlRequestError('interrupt', 'Claude did not answer the interrupt.')
}
const connection = claude.connections[0]!
await openTurn(connection)
const asked = Date.now()
await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
await laneDrained()
// A turn Claude would not interrupt ends only with its child, so there is no grace to wait.
expect(Date.now() - asked).toBeLessThan(CLAUDE_STOP_GRACE_MS / 2)
expect(connection.closed).toBe(true)
expect(await turnOutcome()).toBe('cancellation')
const texts = await statusTexts()
expect(texts).toContain('Cancellation requested.')
expect(texts.some((text) => text.includes("didn't stop"))).toBe(false)
})
it('ends the child when the interrupt fails, with no unconfirmed row', async () => {
claude.routes.interrupt = () => {
throw new Error('control request lost')
}
const connection = claude.connections[0]!
await openTurn(connection)
await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
await laneDrained()
expect(connection.closed).toBe(true)
expect(await turnOutcome()).toBe('cancellation')
expect(await statusTexts()).toEqual(['Cancellation requested.'])
})
it('ends the child within the grace when Claude never answers the interrupt', async () => {
claude.routes.interrupt = NEVER_ANSWERS
const connection = claude.connections[0]!
await openTurn(connection)
const asked = Date.now()
await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
await laneDrained()
expect(connection.closed).toBe(true)
expect(Date.now() - asked).toBeLessThan(CLAUDE_STOP_GRACE_MS + 1_500)
expect(await turnOutcome()).toBe('cancellation')
expect(await statusTexts()).toEqual(['Cancellation requested.'])
}, 15_000)
it('ends background work Claude runs when the Stop ends the child', async () => {
const connection = claude.connections[0]!
await openTurn(connection)
frame(connection, {
type: 'system',
subtype: 'task_started',
uuid: 'task-start',
task_id: 'background-1',
task_type: 'local_agent',
is_backgrounded: true
})
const background = { providerId: 'background-1', kind: 'agent' }
// The host's child record, which the strip, the sidebar and Monitoring all read.
expect(childRecords()).toEqual([
expect.objectContaining({ ...background, state: 'working', membership: 'live' })
])
await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
frame(connection, INTERRUPTED_RESULT)
await laneDrained()
// The work ran inside Claude's process, so it ends with it: its record settles, and the chat
// reads done, Interrupted, with nothing left for Monitoring.
expect(childRecords()).toEqual([
expect.objectContaining({
...background,
state: 'done',
membership: 'settled',
// Stopped with the chat, as a task's own stop reads: Interrupted, not an unknown ending.
outcome: 'cancelled'
})
])
expect(await agentStatus()).toEqual({
state: 'done',
mainAgent: { state: 'done', outcome: 'cancellation' }
})
expect(connection.closed).toBe(true)
})
it('starts a new child for the next send after a Stop, on the same Claude conversation', async () => {
const connection = claude.connections[0]!
await openTurn(connection)
await stop()
frame(connection, INTERRUPTED_RESULT)
await laneDrained()
const next = await send('Carry on.')
const resumed = await eventually(() => {
const started = claude.connections.at(-1)
expect(started).not.toBe(connection)
expect(started && wrote(started, 'Carry on.')).toBe(true)
return started!
})
// The wake resumes the same Claude conversation; the first test pins the leaf it resumes after.
expect(store.getRecord(SESSION)?.providerHandleChain.at(-1)?.handle).toMatchObject({
sessionId: PROVIDER_SESSION_ID
})
expect(resumed.closed).toBe(false)
expect(await dispatch(next)).toMatchObject({ state: 'pending' })
})
it('delivers a send issued with the pre-Stop fence during the Stop to the resumed child only', async () => {
const connection = claude.connections[0]!
await openTurn(connection)
const fence = store.getRecord(SESSION)!.lease.runtimeFence
let childrenWhenClosed: number | undefined
const close = connection.close
connection.close = async () => {
const proven = await close()
childrenWhenClosed = claude.connections.length
return proven
}
let answer!: () => void
claude.routes.interrupt = () =>
new Promise((resolve) => {
answer = () => resolve({ still_queued: [], cancelled: [] })
})
// Issued while the Stop's first step waits on the interrupt: the client still holds the fence
// the rest moves.
const stopped = stop()
await interruptSent(connection)
const sent = send('Typed during the Stop.', fence)
answer()
expect(await stopped).toMatchObject({ ok: true, value: { cancelled: true } })
frame(connection, INTERRUPTED_RESULT)
await expect(sent).resolves.toEqual(expect.any(String))
const resumed = await eventually(() => {
const started = claude.connections.at(-1)!
expect(started).not.toBe(connection)
expect(wrote(started, 'Typed during the Stop.')).toBe(true)
return started
})
// Handed over only after the old child's close resolved, and never to that child.
expect(childrenWhenClosed).toBe(1)
expect(resumed.closed).toBe(false)
expect(wrote(connection, 'Typed during the Stop.')).toBe(false)
expect(store.getRecord(SESSION)!.lease.runtimeFence).not.toBe(fence)
expect((await statusTexts()).filter((text) => text !== 'Cancellation requested.')).toEqual([])
})
it('sends a queue-if-active message issued while the Stop ends the child directly to the resumed child', async () => {
const connection = claude.connections[0]!
await openTurn(connection)
const fence = store.getRecord(SESSION)!.lease.runtimeFence
let childrenWhenClosed: number | undefined
const close = connection.close
connection.close = async () => {
const proven = await close()
childrenWhenClosed = claude.connections.length
return proven
}
await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
// Issued while the Stop's second step waits for the stopped turn, with the fence it moves.
const body = hostTestMessage('Queued during the Stop.')
const sent = host.send(CALLER, {
envelope: envelope('agentSession.send', { body, delivery: 'queue-if-active' }, fence),
body,
delivery: 'queue-if-active',
userSend: true
})
frame(connection, INTERRUPTED_RESULT)
// Nothing runs after the rest, so it goes out as a direct submission, not a queued card.
expect(await sent).toMatchObject({ ok: true, value: { submission: expect.any(Object) } })
await eventually(() => {
const started = claude.connections.at(-1)!
expect(started).not.toBe(connection)
expect(wrote(started, 'Queued during the Stop.')).toBe(true)
})
expect(childrenWhenClosed).toBe(1)
expect(wrote(connection, 'Queued during the Stop.')).toBe(false)
})
it("runs nothing queued during the Stop's first step before the child's end", async () => {
const connection = claude.connections[0]!
await openTurn(connection)
let answer!: () => void
claude.routes.interrupt = () =>
new Promise((resolve) => {
answer = () => resolve({ still_queued: [], cancelled: [] })
})
const stopped = stop()
await interruptSent(connection)
// Any later operation on the chat, such as a prompt answer or an option change, queues here.
let childLiveForNextOperation: boolean | undefined
const next = host['tasks'].serialize(SESSION, async () => {
childLiveForNextOperation = !connection.closed
})
answer()
await stopped
frame(connection, INTERRUPTED_RESULT)
await next
expect(childLiveForNextOperation).toBe(false)
})
it('still answers the Stop, with its row, when the child cannot be proven gone; the failure is reported', async () => {
const connection = claude.connections[0]!
await openTurn(connection)
const close = connection.close
connection.close = async () => false
await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
frame(connection, INTERRUPTED_RESULT)
await laneDrained()
expect(await statusTexts()).toEqual(['Cancellation requested.'])
expect(sinkErrors).toEqual([
expect.objectContaining({
name: 'StructuredAgentSessionEvictionError',
step: 'stop-provider-child'
})
])
connection.close = close
})
it.each([
['takes', undefined],
[
'fails',
() => {
throw new Error('control request lost')
}
],
['never answers', NEVER_ANSWERS]
])(
'ends the child for a Stop naming the turn that just ended when Claude interrupts the follow-up and %s',
async (_answer, interrupt) => {
if (interrupt) {
claude.routes.interrupt = interrupt
}
const connection = claude.connections[0]!
const ended = await openTurn(connection)
frame(connection, { type: 'result', subtype: 'success', is_error: false, uuid: 'result-1' })
await eventually(async () =>
expect(
activeStructuredAgentSessionTurnId((await host.journalSnapshot(SESSION)).items)
).toBeNull()
)
// Handed over but not yet echoed: no turn of its own for a client to name.
await send('Follow-up.')
await eventually(() => expect(wrote(connection, 'Follow-up.')).toBe(true))
// As the phone sends it: the turn it last saw working.
await expect(stop(ended)).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
await laneDrained()
expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(true)
expect(connection.closed).toBe(true)
expect(await statusTexts()).toEqual(['Cancellation requested.'])
},
15_000
)
it('keeps a second Stop pressed while the first ends the child quiet', async () => {
const connection = claude.connections[0]!
await openTurn(connection)
await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } })
const second = stop()
frame(connection, INTERRUPTED_RESULT)
await expect(second).resolves.toMatchObject({ ok: true, value: { cancelled: false } })
expect(connection.closed).toBe(true)
expect(connection.calls.filter((call) => call.subtype === 'interrupt')).toHaveLength(1)
expect(await statusTexts()).toEqual(['Cancellation requested.'])
expect(sinkErrors).toEqual([])
})
const BRANCH_QUESTION = {
questions: [
{
question: 'Which branch?',
header: 'Branch',
multiSelect: false,
options: [{ label: 'main' }, { label: 'dev' }]
}
]
}
/** Claude asks on the open turn; returns its answer and the card the journal shows for it. */
async function ask(
connection: FakeConnection,
toolName: string,
input: Record<string, unknown>,
signal?: AbortSignal
): Promise<{
answered: ReturnType<typeof invokeCanUseTool>
card: { itemId: string; expectedRevision: number }
}> {
const answered = invokeCanUseTool(connection, toolName, 'permission-1', 'tool-1', {
input,
...(signal ? { signal } : {})
})
const card = await eventually(async () => {
await host.flushStreamedEvents(SESSION)
const item = (await host.journalSnapshot(SESSION)).items.find(
(entry) => entry.body.kind === 'approval' || entry.body.kind === 'question'
)
expect(item).toBeDefined()
return { itemId: item!.itemId, expectedRevision: item!.revision }
})
return { answered, card }
}
function cancelCard(turnId: string, prompt: { itemId: string; expectedRevision: number }) {
const fields = { turnId, prompt }
return host.cancel(CALLER, { envelope: envelope('agentSession.cancel', fields), ...fields })
}
async function cardResolution(itemId: string): Promise<unknown> {
await host.flushStreamedEvents(SESSION)
const body = (await host.journalSnapshot(SESSION)).items.find(
(item) => item.itemId === itemId
)?.body
return body?.kind === 'approval' || body?.kind === 'question' ? body.resolution : undefined
}
it('dismisses an approval card on its Cancel with the Deny reply, and the turn and child go on', async () => {
const connection = claude.connections[0]!
const turnId = await openTurn(connection)
const { answered, card } = await ask(connection, 'Bash', { command: 'rm -rf build' })
await expect(cancelCard(turnId, card)).resolves.toMatchObject({
ok: true,
value: { turnId, cancelled: true }
})
await laneDrained()
const reply = await answered.promise
expect(reply).toMatchObject({ behavior: 'deny', message: 'User denied this action.' })
expect(reply).not.toHaveProperty('interrupt')
// Read as cancelled by the user, as every card's Cancel reads, not as a Deny pressed.
expect(await cardResolution(card.itemId)).toMatchObject({
state: 'cancelled',
resolvedBy: CALLER.callerKey
})
expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(false)
expect(connection.closed).toBe(false)
expect(activeStructuredAgentSessionTurnId((await host.journalSnapshot(SESSION)).items)).toBe(
turnId
)
expect(await statusTexts()).toEqual([])
})
it("ends a question card's Cancel the way the chat's Stop does, and the next send resumes", async () => {
const connection = claude.connections[0]!
const turnId = await openTurn(connection)
const request = new AbortController()
const { answered, card } = await ask(
connection,
'AskUserQuestion',
BRANCH_QUESTION,
request.signal
)
await expect(cancelCard(turnId, card)).resolves.toMatchObject({
ok: true,
value: { turnId, cancelled: true }
})
// As Claude does once interrupted: it cancels the request it was holding.
request.abort()
await host.flushStreamedEvents(SESSION)
// Settled with the Stop's first step, before the child ends: not answerable meanwhile.
expect(connection.closed).toBe(false)
const cancelledHere = { state: 'cancelled', resolvedBy: CALLER.callerKey }
expect(await cardResolution(card.itemId)).toMatchObject(cancelledHere)
expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(true)
frame(connection, INTERRUPTED_RESULT)
await laneDrained()
expect(connection.closed).toBe(true)
expect(await turnOutcome()).toBe('cancellation')
// The child's end takes Claude's request with it: no reply raced the interrupt, and nothing
// wrote over the user's cancel.
expect(await answered.promise).not.toMatchObject({ behavior: expect.any(String) })
expect(await cardResolution(card.itemId)).toMatchObject(cancelledHere)
expect(await statusTexts()).toEqual(['Cancellation requested.'])
await send('Carry on.')
await eventually(() => {
const started = claude.connections.at(-1)!
expect(started).not.toBe(connection)
expect(wrote(started, 'Carry on.')).toBe(true)
})
})
it("declines a question card's Cancel itself when the Stop finds nothing to stop", async () => {
const connection = claude.connections[0]!
const ended = await openTurn(connection)
frame(connection, { type: 'result', subtype: 'success', is_error: false, uuid: 'result-1' })
await eventually(async () =>
expect(
activeStructuredAgentSessionTurnId((await host.journalSnapshot(SESSION)).items)
).toBeNull()
)
// Asked after the main turn ended, as a background agent might.
const { answered, card } = await ask(connection, 'AskUserQuestion', BRANCH_QUESTION)
await expect(cancelCard(ended, card)).resolves.toMatchObject({ ok: true })
await laneDrained()
expect(await answered.promise).toMatchObject({
behavior: 'deny',
message: expect.stringMatching(/dismissed/)
})
expect(await cardResolution(card.itemId)).toMatchObject({
state: 'cancelled',
resolvedBy: CALLER.callerKey
})
expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(false)
expect(connection.closed).toBe(false)
expect(await statusTexts()).toEqual([])
})
it('dismisses a question card a finished turn raised, leaving the turn running now alone', async () => {
const connection = claude.connections[0]!
const earlier = await openTurn(connection)
// Asked in the first turn, then outlived it, as a background agent's question might.
const { answered, card } = await ask(connection, 'AskUserQuestion', BRANCH_QUESTION)
frame(connection, { type: 'result', subtype: 'success', is_error: false, uuid: 'result-1' })
await eventually(async () =>
expect(
activeStructuredAgentSessionTurnId((await host.journalSnapshot(SESSION)).items)
).toBeNull()
)
const running = await openTurn(connection, 'Now something else.')
expect(running).not.toBe(earlier)
await expect(cancelCard(earlier, card)).resolves.toMatchObject({ ok: true })
await laneDrained()
expect(await answered.promise).toMatchObject({ behavior: 'deny' })
expect(await cardResolution(card.itemId)).toMatchObject({
state: 'cancelled',
resolvedBy: CALLER.callerKey
})
expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(false)
expect(connection.closed).toBe(false)
expect(activeStructuredAgentSessionTurnId((await host.journalSnapshot(SESSION)).items)).toBe(
running
)
})
it('dismisses a plan card on its Cancel: Claude is told to wait for the user, and keeps running', async () => {
const connection = claude.connections[0]!
const turnId = await openTurn(connection)
const { answered, card } = await ask(connection, 'ExitPlanMode', {
plan: '# Release\n\n- Tag it'
})
await expect(cancelCard(turnId, card)).resolves.toMatchObject({
ok: true,
value: { turnId, cancelled: true }
})
await laneDrained()
const reply = await answered.promise
expect(reply).toMatchObject({ behavior: 'deny' })
expect(reply).not.toHaveProperty('interrupt')
// Not "keep planning": nothing asks Claude to revise and show another plan.
expect(JSON.stringify(reply)).toMatch(/wait for them/)
expect(JSON.stringify(reply)).not.toMatch(/keep planning|ExitPlanMode again/i)
expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(false)
expect(connection.closed).toBe(false)
const cards = (await host.journalSnapshot(SESSION)).items.filter(
(item) => item.body.kind === 'approval'
)
expect(cards).toHaveLength(1)
expect(await cardResolution(card.itemId)).toMatchObject({
state: 'cancelled',
resolvedBy: CALLER.callerKey
})
expect(await statusTexts()).toEqual([])
})
@@ -36,6 +36,9 @@ let host: StructuredAgentSessionHost
let dispatch: Mock<StructuredAgentSessionAdapter['dispatch']>
let cancelTurn: Mock<StructuredAgentSessionAdapter['cancelTurn']>
let awaitStarted: Mock<NonNullable<StructuredAgentSessionAdapter['awaitStarted']>>
let closeSession: Mock<NonNullable<StructuredAgentSessionAdapter['closeSession']>>
/** Codex's answer by default: its Stop keeps the child. */
let stopEndsSession: boolean
let events: StructuredAgentSessionEventSink | undefined
function eventually(assertion: () => void | Promise<void>): Promise<void> {
@@ -49,6 +52,8 @@ beforeEach(async () => {
dispatch = vi.fn(async () => ({ state: 'admitted' as const }))
cancelTurn = vi.fn(async () => ({ cancelled: true }))
awaitStarted = vi.fn(async () => undefined)
closeSession = vi.fn(async () => true)
stopEndsSession = false
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
store,
@@ -74,9 +79,10 @@ beforeEach(async () => {
},
dispatch,
awaitStarted,
closeSession: vi.fn(async () => true),
closeSession,
releaseAcquisition: vi.fn(async () => true),
cancelTurn,
stopEndsSession: () => stopEndsSession,
answerPrompt: vi.fn(async () => undefined),
setOption: vi.fn(async () => undefined)
},
@@ -365,3 +371,78 @@ describe('a Stop that names its turn, as an older client sends it', () => {
expect(await statusRows()).toEqual([])
})
})
describe('a Stop on a provider whose Stop ends its session', () => {
async function handedOver(): Promise<void> {
const { id, result } = send('hello')
await result
await eventually(async () => expect((await submission(id))?.handedOverAt).toBeDefined())
}
/** The Stop's second step, which ends the child, runs next on the session's lane. */
function laneDrained(): Promise<void> {
return host['tasks'].serialize(SESSION, async () => {})
}
it('ends the child after the cancel even when the provider refused it, and says only that it was asked', async () => {
stopEndsSession = true
await handedOver()
cancelTurn.mockResolvedValueOnce({
cancelled: false,
refusal: { detail: { text: 'no active turn to interrupt', audience: 'person' } }
})
expect(await stop()).toMatchObject({ ok: true, value: { cancelled: true } })
await laneDrained()
expect(closeSession).toHaveBeenCalledWith(SESSION, 'user-stop')
expect(await statusRows()).toEqual(['Cancellation requested.'])
})
it('keeps the child of a provider whose Stop is not a session boundary', async () => {
await handedOver()
cancelTurn.mockResolvedValueOnce({ cancelled: true })
expect(await stop()).toMatchObject({ ok: true, value: { cancelled: true } })
await laneDrained()
expect(closeSession).not.toHaveBeenCalled()
})
it('leaves the child alone when the Stop named a turn that is no longer live', async () => {
stopEndsSession = true
cancelTurn.mockResolvedValueOnce({ cancelled: false })
expect(await stop('turn-1')).toMatchObject({ ok: true, value: { cancelled: false } })
await laneDrained()
expect(closeSession).not.toHaveBeenCalled()
expect(await statusRows()).toEqual([ALREADY_FINISHED])
})
it("ends the child when the provider's cancel of a Stop naming a turn no longer live fails", async () => {
stopEndsSession = true
await handedOver()
// The interrupt went out and its answer was lost: what it stopped is unknown.
cancelTurn.mockRejectedValueOnce(new Error('control request lost'))
expect(await stop('turn-1')).toMatchObject({ ok: true, value: { cancelled: true } })
await laneDrained()
expect(closeSession).toHaveBeenCalledWith(SESSION, 'user-stop')
expect(await statusRows()).toEqual(['Cancellation requested.'])
})
it('ends the child when the provider took a Stop naming a turn that is no longer live', async () => {
stopEndsSession = true
await handedOver()
// The interrupt stopped the follow-up in flight, which has no turn a client could name.
cancelTurn.mockResolvedValueOnce({ cancelled: true })
expect(await stop('turn-1')).toMatchObject({ ok: true, value: { cancelled: true } })
await laneDrained()
expect(closeSession).toHaveBeenCalledWith(SESSION, 'user-stop')
expect(await statusRows()).toEqual(['Cancellation requested.'])
})
})
@@ -18,96 +18,34 @@ import type {
AgentSessionThreadGoalChange,
AgentSessionThreadGoalResult
} from '../../../shared/agent-session-wire'
import { agentSessionFailureWords } from '../../../shared/agent-session-failure-words'
import { agentSessionFailureFact } from '../../../shared/agent-session-failure'
import type { AgentSessionPromptRequest } from './structured-agent-session-turns-prompt'
import type { AgentChildWorkView } from '../../../shared/agent-status-child-work-view'
import { threadGoalPlan } from './structured-agent-session-thread-goal'
import {
isMainAgentWorkingOnceFlushed,
performCancel
} from './structured-agent-session-turns-cancel'
import {
admitAndRunAgentSessionMutation,
type AgentSessionMutationRequest,
type AgentSessionMutationSessionPreparation
} from './structured-agent-session-mutation-admission'
mutateStructuredAgentSession,
type StructuredAgentSessionMutationContext
} from './structured-agent-session-mutation-context'
import {
openForWrite,
openWithAgent,
sendPreparation,
structuredAgentSessionFailureWordsContext,
structuredAgentSessionSendBlock
} from './structured-agent-session-send-preparation'
import {
cancelPlan,
promptPlan,
sendPlan,
setOptionPlan,
type MutationPlan
setOptionPlan
} from './structured-agent-session-mutation-plans'
import { runQueueableStructuredAgentSessionSend } from './structured-agent-session-queued-send'
import { runStopWithQueuePause } from './structured-agent-session-queued-stop'
import type {
StructuredAgentSessionCaller,
StructuredAgentSessionHostDeps,
StructuredAgentSessionHostSession
} from './structured-agent-session-host-types'
import { cancelStructuredAgentSessionPrompt } from './structured-agent-session-prompt-cancel'
import { mutateWithChatStop } from './structured-agent-session-chat-stop'
export type { StructuredAgentSessionMutationContext } from './structured-agent-session-mutation-context'
import type { StructuredAgentSessionCaller } from './structured-agent-session-host-types'
import {
readStructuredAgentSessionOptions,
recordStructuredAgentSessionOptionIntent
} from './structured-agent-session-options-read'
export type StructuredAgentSessionMutationContext = {
deps: StructuredAgentSessionHostDeps
sessions: Map<string, StructuredAgentSessionHostSession>
publish: (sessionId: string, journal: StructuredAgentSessionHostSession['journal']) => void
flushStreamedEvents: (sessionId: string) => Promise<void>
/** The host's accessor, for a caller outside the session's serialize. */
conversation: (sessionId: string) => Promise<StructuredAgentSessionHostSession>
/** The session's child records, as the strip reads them; what command admission decides on. */
readChildWork: (sessionId: string) => AgentChildWorkView[] | undefined
serialize: <T>(sessionId: string, task: () => Promise<T>) => Promise<T>
/** The session's conversation, opened when closed; inside the caller's serialize. */
openConversation: (sessionId: string) => Promise<StructuredAgentSessionHostSession | null>
/** Gives the session a provider child; inside the caller's serialize. */
ensureAgent: (sessionId: string) => Promise<AgentSessionMutationSessionPreparation>
/** A message was accepted: the session's delivery loop hands it over. */
wakeDelivery: (sessionId: string) => void
/** Stops the session's provider child, keeping its conversation; inside the caller's serialize. */
stopAgent: (sessionId: string) => Promise<void>
/** Only for gate inputs living in the RECORD store, which can settle with no
* journal commit (a conversation command). Draft-table changes need no call:
* the draft store notifies through the journal's own commit listener. */
wakeQueuedDrain?: (sessionId: string) => void
now: () => number
}
/** Admits the envelope and runs the plan inside the session's serialize. */
export function mutateStructuredAgentSession<TValue>(
context: StructuredAgentSessionMutationContext,
caller: StructuredAgentSessionCaller,
envelope: AgentSessionMutationEnvelope,
plan: MutationPlan<TValue>,
prepareSession?: AgentSessionMutationRequest<TValue>['prepareSession']
): Promise<AgentSessionMutationResult<TValue>> {
return context.serialize(envelope.sessionId, () =>
admitAndRunAgentSessionMutation({
store: context.deps.store,
adapter: context.deps.adapter,
callerKey: caller.callerKey,
envelope,
plan,
journal: () => context.sessions.get(envelope.sessionId)?.journal,
prepareSession,
publish: (journal) => context.publish(envelope.sessionId, journal),
flushStreamedEvents: context.flushStreamedEvents,
providerChildPhase: () => context.sessions.get(envelope.sessionId)?.child?.phase,
now: () => context.now()
})
)
}
export function sendStructuredAgentSessionTurn(
context: StructuredAgentSessionMutationContext,
caller: StructuredAgentSessionCaller,
@@ -157,7 +95,7 @@ export function cancelStructuredAgentSessionTurn(
prompt?: { itemId: string; expectedRevision: number }
}
): Promise<AgentSessionMutationResult<AgentSessionCancelResult>> {
if (params.scope || params.prompt) {
if (params.scope) {
return mutateStructuredAgentSession(
context,
caller,
@@ -166,57 +104,19 @@ export function cancelStructuredAgentSessionTurn(
openForWrite(context, params.envelope)
)
}
const plan = cancelPlan({
...params,
stopChild: () => context.stopAgent(params.envelope.sessionId)
})
return mutateStructuredAgentSession(
context,
caller,
params.envelope,
{
...plan,
// Stop's queue step, the same for every client: once the Stop takes effect
// the queue is paused. The cards stay published; nothing is withdrawn and no
// text ever rides the answer.
run: (ctx) =>
runStopWithQueuePause(ctx, async (tookEffect) => {
// Stop withdraws every queued SUBMISSION first, whatever the start or the child is doing.
const withdrawn = await ctx.journal.rejectQueuedSubmissions(
ctx.fence,
agentSessionFailureWords(agentSessionFailureFact('cancelled'), { surface: 'rejection' })
)
const named = params.turnId !== undefined ? { turnId: params.turnId } : {}
const child = context.sessions.get(ctx.sessionId)?.child
if (child?.phase === 'starting') {
// A start that may never land is the one thing here Stop has to end; the chat stays.
await tookEffect()
await context.stopAgent(ctx.sessionId)
return { ok: true, value: { ...named, cancelled: true } }
}
// A Stop naming no turn ends nothing more unless the session reads working, by the rule
// every session list and the chat's own Stop read it.
const inFlight = params.turnId !== undefined || (await isMainAgentWorkingOnceFlushed(ctx))
const record = context.deps.store.getRecord(ctx.sessionId)
if (!child || !inFlight) {
if (withdrawn.length > 0) {
await tookEffect()
}
return { ok: true, value: { ...named, cancelled: withdrawn.length > 0 } }
}
await tookEffect()
return performCancel(
{ ...ctx, failureTextContext: structuredAgentSessionFailureWordsContext(record) },
{
clientOperationId: params.envelope.clientOperationId,
...named,
stopChild: () => context.stopAgent(params.envelope.sessionId),
withdrewQueued: withdrawn.length > 0
}
)
})
},
openForWrite(context, params.envelope)
const plan = cancelPlan(params)
const { prompt } = params
// A card's Cancel stops whatever the chat has in flight, as the Stop button does; it reaches the
// Stop only for a card the live turn raised (`cancelStructuredAgentSessionPrompt`).
const stopped = prompt ? { envelope: params.envelope } : params
return mutateWithChatStop(context, caller, stopped, plan, (ctx, stop) =>
prompt
? cancelStructuredAgentSessionPrompt(
ctx,
{ ...(params.turnId !== undefined ? { turnId: params.turnId } : {}), prompt },
{ stop, interrupt: () => plan.run(ctx) }
)
: stop().then(({ outcome }) => outcome)
)
}
@@ -0,0 +1,69 @@
// The context every client mutation of a session runs with, and the one path each takes: admit the
// envelope against the lease, then run its plan inside the session's serialize.
import type {
AgentSessionMutationEnvelope,
AgentSessionMutationResult
} from '../../../shared/agent-session-wire'
import type { AgentChildWorkView } from '../../../shared/agent-status-child-work-view'
import {
admitAndRunAgentSessionMutation,
type AgentSessionMutationRequest,
type AgentSessionMutationSessionPreparation
} from './structured-agent-session-mutation-admission'
import type { MutationPlan } from './structured-agent-session-mutation-plans'
import type {
StructuredAgentSessionCaller,
StructuredAgentSessionHostDeps,
StructuredAgentSessionHostSession
} from './structured-agent-session-host-types'
export type StructuredAgentSessionMutationContext = {
deps: StructuredAgentSessionHostDeps
sessions: Map<string, StructuredAgentSessionHostSession>
publish: (sessionId: string, journal: StructuredAgentSessionHostSession['journal']) => void
flushStreamedEvents: (sessionId: string) => Promise<void>
/** The host's accessor, for a caller outside the session's serialize. */
conversation: (sessionId: string) => Promise<StructuredAgentSessionHostSession>
/** The session's child records, as the strip reads them; what command admission decides on. */
readChildWork: (sessionId: string) => AgentChildWorkView[] | undefined
serialize: <T>(sessionId: string, task: () => Promise<T>) => Promise<T>
/** The session's conversation, opened when closed; inside the caller's serialize. */
openConversation: (sessionId: string) => Promise<StructuredAgentSessionHostSession | null>
/** Gives the session a provider child; inside the caller's serialize. */
ensureAgent: (sessionId: string) => Promise<AgentSessionMutationSessionPreparation>
/** A message was accepted: the session's delivery loop hands it over. */
wakeDelivery: (sessionId: string) => void
/** Stops the session's provider child, keeping its conversation; inside the caller's serialize. */
stopAgent: (sessionId: string) => Promise<void>
/** Only for gate inputs living in the RECORD store, which can settle with no
* journal commit (a conversation command). Draft-table changes need no call:
* the draft store notifies through the journal's own commit listener. */
wakeQueuedDrain?: (sessionId: string) => void
now: () => number
}
/** Admits the envelope and runs the plan inside the session's serialize. */
export function mutateStructuredAgentSession<TValue>(
context: StructuredAgentSessionMutationContext,
caller: StructuredAgentSessionCaller,
envelope: AgentSessionMutationEnvelope,
plan: MutationPlan<TValue>,
prepareSession?: AgentSessionMutationRequest<TValue>['prepareSession']
): Promise<AgentSessionMutationResult<TValue>> {
return context.serialize(envelope.sessionId, () =>
admitAndRunAgentSessionMutation({
store: context.deps.store,
adapter: context.deps.adapter,
callerKey: caller.callerKey,
envelope,
plan,
journal: () => context.sessions.get(envelope.sessionId)?.journal,
prepareSession,
publish: (journal) => context.publish(envelope.sessionId, journal),
flushStreamedEvents: context.flushStreamedEvents,
providerChildPhase: () => context.sessions.get(envelope.sessionId)?.child?.phase,
now: () => context.now()
})
)
}
@@ -173,7 +173,6 @@ export function cancelPlan(params: {
scope?: 'background-tasks'
taskId?: string
prompt?: { itemId: string; expectedRevision: number }
stopChild?: () => Promise<void>
/** The session's child records, which name the tasks a background Stop reaches. */
childWork?: () => readonly AgentChildWorkView[] | undefined
}): MutationPlan<AgentSessionCancelResult> {
@@ -196,7 +195,6 @@ export function cancelPlan(params: {
...(params.scope ? { scope: params.scope } : {}),
...(params.taskId ? { taskId: params.taskId } : {}),
...(params.prompt ? { prompt: params.prompt } : {}),
...(params.stopChild ? { stopChild: params.stopChild } : {}),
...(params.childWork ? { childWork: params.childWork } : {})
}),
// Interrupting twice would kill a turn the client never asked to stop, so a
@@ -6,8 +6,13 @@ import { afterEach, describe, expect, it, vi } from 'vitest'
import type { AgentSessionJournalIdentity } from '../../../shared/agent-session-journal-types'
import { createTrackedJournalOpener } from '../agent-session-journal/journal-host-database-test-support'
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter'
import {
AgentSessionPromptUnavailableError,
type StructuredAgentSessionAdapter
} from './structured-agent-session-adapter'
import { performCancel, type AgentSessionTurnContext } from './structured-agent-session-turns'
import { cancelStructuredAgentSessionPrompt } from './structured-agent-session-prompt-cancel'
import type { AgentSessionPromptCancelRoute } from './structured-agent-session-adapter-stop'
const IDENTITY: AgentSessionJournalIdentity = {
sessionId: 'session-1',
@@ -34,16 +39,27 @@ afterEach(async () => {
}
})
async function pendingPrompt(): Promise<{ journal: AgentSessionJournal; itemId: string }> {
async function pendingPrompt(
options = [{ id: 'allow', label: 'Allow' }],
/** Raise the card in a turn that is still running, rather than on the conversation. */
inLiveTurn = false
): Promise<{ journal: AgentSessionJournal; itemId: string }> {
root = await mkdtemp(join(tmpdir(), 'orca-prompt-cancel-'))
const journal = await journals.open({ identity: IDENTITY, stateDirectory: root })
const turn = inLiveTurn
? await journal.appendItem(
{ ...PROMPT_IDENTITY, ordinal: 0 },
{ kind: 'turn', turnId: 'turn-1', state: 'running', startedAt: 1 },
{ fence: 1, turnScope: AGENT_JOURNAL_THREAD_SCOPE }
)
: null
const item = await journal.appendItem(
PROMPT_IDENTITY,
{
kind: 'approval',
title: 'Approve?',
detail: null,
options: [{ id: 'allow', label: 'Allow' }],
options,
resolution: {
state: 'pending',
selectedOptionId: null,
@@ -51,7 +67,10 @@ async function pendingPrompt(): Promise<{ journal: AgentSessionJournal; itemId:
resolvedAt: null
}
},
{ fence: 1, turnScope: AGENT_JOURNAL_THREAD_SCOPE }
{
fence: 1,
turnScope: turn ? { kind: 'turn', turnItemId: turn.itemId } : AGENT_JOURNAL_THREAD_SCOPE
}
)
return { journal, itemId: item.itemId }
}
@@ -212,3 +231,141 @@ describe('performCancel for a pending prompt', () => {
expect(journal.snapshot().items).toHaveLength(1)
})
})
describe("a card's own Cancel, as its provider answers it", () => {
async function cancelCard(
answer: AgentSessionPromptCancelRoute | undefined,
revision = 1,
endsSession = true,
inLiveTurn = true
) {
const { journal, itemId } = await pendingPrompt(
[
{ id: 'allow', label: 'Allow' },
{ id: 'deny', label: 'Deny' }
],
inLiveTurn
)
const ctx = context(
journal,
vi.fn(async () => ({ cancelled: true })),
vi.fn(async () => undefined)
)
const answerPrompt = vi.fn<StructuredAgentSessionAdapter['answerPrompt']>(async (input) => {
await input.commit()
})
const dismissPrompt = vi.fn<NonNullable<StructuredAgentSessionAdapter['dismissPrompt']>>(
async (input) => {
await input.commit()
}
)
Object.assign(ctx.adapter, { answerPrompt, dismissPrompt, routePromptCancel: () => answer })
const routes = {
stop: vi.fn(async () => ({
outcome: { ok: true as const, value: { cancelled: endsSession } },
endsSession
})),
interrupt: vi.fn(async () => ({ ok: true as const, value: { cancelled: true } }))
}
const result = await cancelStructuredAgentSessionPrompt(
ctx,
{ turnId: 'turn-1', prompt: { itemId, expectedRevision: revision } },
routes
)
const card = journal.snapshot().items.find((item) => item.itemId === itemId)?.body
return { result, routes, answerPrompt, dismissPrompt, card }
}
it('interrupts the turn holding the card for a provider that gives no answer', async () => {
const { routes, answerPrompt } = await cancelCard(undefined)
expect(routes.interrupt).toHaveBeenCalledOnce()
expect(routes.stop).not.toHaveBeenCalled()
expect(answerPrompt).not.toHaveBeenCalled()
})
it('records a dismissal as cancelled by the caller and has the provider decline the request', async () => {
const { result, routes, answerPrompt, dismissPrompt, card } = await cancelCard({
kind: 'dismiss'
})
expect(result).toEqual({ ok: true, value: { turnId: 'turn-1', cancelled: true } })
expect(dismissPrompt).toHaveBeenCalledWith(expect.objectContaining({ answer: true }))
expect(card).toMatchObject({
resolution: { state: 'cancelled', selectedOptionId: null, resolvedBy: 'client-1' }
})
expect(answerPrompt).not.toHaveBeenCalled()
expect(routes.stop).not.toHaveBeenCalled()
})
it("runs the chat's Stop and settles the card in its step, leaving the request to the child's end", async () => {
const { result, routes, answerPrompt, dismissPrompt, card } = await cancelCard({
kind: 'stop'
})
expect(result).toEqual({ ok: true, value: { turnId: 'turn-1', cancelled: true } })
expect(routes.stop).toHaveBeenCalledOnce()
expect(routes.interrupt).not.toHaveBeenCalled()
expect(dismissPrompt).toHaveBeenCalledWith(expect.objectContaining({ answer: false }))
expect(answerPrompt).not.toHaveBeenCalled()
expect(card).toMatchObject({ resolution: { state: 'cancelled', resolvedBy: 'client-1' } })
})
it('declines the request itself when the Stop ends nothing', async () => {
const { result, dismissPrompt, card } = await cancelCard({ kind: 'stop' }, 1, false)
expect(result).toEqual({ ok: true, value: { turnId: 'turn-1', cancelled: true } })
expect(dismissPrompt).toHaveBeenCalledWith(expect.objectContaining({ answer: true }))
expect(card).toMatchObject({ resolution: { state: 'cancelled', resolvedBy: 'client-1' } })
})
it('still records the cancel when the provider already let the request go', async () => {
const { journal, itemId } = await pendingPrompt()
const ctx = context(
journal,
vi.fn(async () => ({ cancelled: true })),
vi.fn(async () => undefined)
)
Object.assign(ctx.adapter, {
dismissPrompt: async () => {
throw new AgentSessionPromptUnavailableError(itemId)
},
routePromptCancel: () => ({ kind: 'dismiss' })
})
await expect(
cancelStructuredAgentSessionPrompt(
ctx,
{ turnId: 'turn-1', prompt: { itemId, expectedRevision: 1 } },
{ stop: vi.fn(), interrupt: vi.fn() }
)
).resolves.toMatchObject({ ok: true })
expect(journal.snapshot().items.find((item) => item.itemId === itemId)?.body).toMatchObject({
resolution: { state: 'cancelled', resolvedBy: 'client-1' }
})
})
it('dismisses a card a turn that is no longer live raised, and stops nothing', async () => {
const { result, routes, dismissPrompt, card } = await cancelCard(
{ kind: 'stop' },
1,
true,
false
)
expect(result).toEqual({ ok: true, value: { turnId: 'turn-1', cancelled: true } })
expect(routes.stop).not.toHaveBeenCalled()
expect(dismissPrompt).toHaveBeenCalledWith(expect.objectContaining({ answer: true }))
expect(card).toMatchObject({ resolution: { state: 'cancelled', resolvedBy: 'client-1' } })
})
it('refuses a card that moved on before choosing a route', async () => {
const { result, routes } = await cancelCard({ kind: 'stop' }, 2)
expect(result).toMatchObject({
ok: false,
refusal: { code: 'agent_session_item_revision_stale' }
})
expect(routes.stop).not.toHaveBeenCalled()
})
})
@@ -0,0 +1,146 @@
// A prompt card's own Cancel, routed the way its provider says: a dismissal, or the chat's Stop. A
// provider that says nothing interrupts the turn holding the card. The host decides, so a client of
// any version gets the same Cancel.
import { agentSessionFailureFact } from '../../../shared/agent-session-failure'
import { agentSessionFailureWords } from '../../../shared/agent-session-failure-words'
import { parseAgentJournalItemKey } from '../../../shared/agent-session-journal-item-key'
import { refuse, type AgentSessionCancelResult } from '../../../shared/agent-session-wire'
import { AgentSessionPromptUnavailableError } from './structured-agent-session-adapter'
import type { StructuredAgentSessionChatStopRun } from './structured-agent-session-chat-stop'
import {
validatePendingPrompt,
type PendingPromptValidation
} from './structured-agent-session-prompt-state'
import type { AgentSessionTurnContext, TurnOutcome } from './structured-agent-session-turns'
type CancelOutcome = TurnOutcome<AgentSessionCancelResult>
type PendingPrompt = Extract<PendingPromptValidation, { ok: true }>
export async function cancelStructuredAgentSessionPrompt(
ctx: AgentSessionTurnContext,
input: { turnId?: string; prompt: { itemId: string; expectedRevision: number } },
routes: {
stop: () => Promise<StructuredAgentSessionChatStopRun>
interrupt: () => Promise<CancelOutcome>
}
): Promise<CancelOutcome> {
const validated = validatePendingPrompt(ctx, input.prompt)
if (!validated.ok) {
return validated
}
const route = ctx.adapter.routePromptCancel?.({
sessionId: ctx.sessionId,
prompt: validated.prompt
})
if (!route) {
return routes.interrupt()
}
const cancelled: CancelOutcome = {
ok: true,
value: { ...(input.turnId ? { turnId: input.turnId } : {}), cancelled: true }
}
if (route.kind === 'dismiss') {
const dismissed = await dismissPrompt(ctx, validated, true)
return dismissed.ok ? cancelled : dismissed
}
// Judged once the provider's accepted lifecycle has landed: its own cancel of the card, or the end
// of the turn that raised it, may still be queued. A failed drain judges what has landed.
await ctx.flushStreamedEvents().catch(() => undefined)
const current = validatePendingPrompt(ctx, input.prompt)
if (!current.ok) {
return current
}
if (!raisedByLiveTurn(ctx, current)) {
// A request that outlived its turn, such as a background agent's: the turn running now is not
// the one the user is cancelling, so nothing stops and the request is declined.
const dismissed = await dismissPrompt(ctx, current, true)
return dismissed.ok ? cancelled : dismissed
}
const stopped = await routes.stop()
if (!stopped.outcome.ok) {
return stopped.outcome
}
// Settled in the Stop's own step, so the card is not answerable while the child ends. That end
// takes the provider's request with it; a Stop that ends nothing must answer the request itself.
const dismissed = await dismissPrompt(ctx, current, !stopped.endsSession)
return dismissed.ok ? cancelled : dismissed
}
function raisedByLiveTurn(ctx: AgentSessionTurnContext, pending: PendingPrompt): boolean {
const live = ctx.journal.liveTurnScope()
const raised = pending.item.turnScope
return live.kind === 'turn' && raised?.kind === 'turn' && raised.turnItemId === live.turnItemId
}
/** Records the card as cancelled by the caller, the provider's lifecycle drained first so nothing
* it already sent lands after; `answer` also declines the provider's request. */
async function dismissPrompt(
ctx: AgentSessionTurnContext,
pending: PendingPrompt,
answer: boolean
): Promise<TurnOutcome<null>> {
const { item, prompt } = pending
const identity = parseAgentJournalItemKey(item.itemId)
if (!identity) {
return {
ok: false,
refusal: refuse(
'agent_session_operation_invalid',
{ reason: 'requestMalformed' },
`Item id ${item.itemId} is not a well-formed item key.`
)
}
}
let committed = false
const commit = async (): Promise<void> => {
await ctx.flushStreamedEvents()
await ctx.journal.appendItem(
identity,
{
...prompt,
resolution: {
state: 'cancelled',
selectedOptionId: null,
resolvedBy: ctx.resolvedBy,
resolvedAt: ctx.now()
}
},
// A revision: the prompt keeps the turn it was raised in.
{ fence: ctx.fence, turnScope: ctx.journal.liveTurnScope() }
)
committed = true
}
try {
await ctx.adapter.dismissPrompt?.({
sessionId: ctx.sessionId,
itemId: item.itemId,
fence: ctx.fence,
answer,
commit
})
} catch (error) {
if (!committed && !(error instanceof AgentSessionPromptUnavailableError)) {
throw error
}
if (committed) {
// The adapter's error is Orca's; the row says only what the user needs to know.
await ctx.journal.appendItem(
{ provider: 'orca', clientMessageId: `${item.itemId}#delivery` },
{
kind: 'status',
...agentSessionFailureWords(agentSessionFailureFact('answerUnconfirmed'), {
surface: 'row'
})
},
{ fence: ctx.fence, turnScope: ctx.journal.liveTurnScope() }
)
}
}
// A request the provider already let go of, by its own cancel or the child's end, is still the
// user's to have cancelled.
if (!committed) {
await commit()
}
return { ok: true, value: null }
}
@@ -34,7 +34,7 @@ import { unsettledQueuedMessages } from './structured-agent-session-queued-stop'
import {
mutateStructuredAgentSession,
type StructuredAgentSessionMutationContext
} from './structured-agent-session-host-mutations'
} from './structured-agent-session-mutation-context'
import type { StructuredAgentSessionCaller } from './structured-agent-session-host-types'
import {
openForWrite,
@@ -32,6 +32,33 @@ export async function isMainAgentWorkingOnceFlushed(
)
}
/** 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. */
export type StructuredAgentSessionStopWindDown = { waitsForProvider: boolean; stoppedAt: number }
/**
* A session-ending Stop's second step, queued behind its first in the same tick so nothing sent
* meanwhile reaches the child it ends. The Stop has answered: a failure here is reported. The next
* Stop retries the wind-down it leaves owed, and so does the idle sweep: at its next tick once the
* child is proven gone, else only after the chat idles with no child work left.
*/
export async function endStoppedStructuredAgentSession(
ctx: Pick<AgentSessionTurnContext, 'sessionId' | 'adapter'>,
windDown: StructuredAgentSessionStopWindDown,
stopChild: () => Promise<void>,
onError: (error: unknown) => void
): Promise<void> {
try {
if (windDown.waitsForProvider) {
await ctx.adapter.awaitStoppedRequestEnd?.(ctx.sessionId, windDown.stoppedAt)
}
await stopChild()
} catch (error) {
onError(error)
}
}
export async function performCancel(
ctx: AgentSessionTurnContext,
input: {
@@ -43,6 +70,9 @@ export async function performCancel(
prompt?: { itemId: string; expectedRevision: number }
/** Ends the provider child, for a running command the provider did not take the Stop on. */
stopChild?: () => Promise<void>
/** Hands the child's end to the Stop's next serialized step, for a provider whose Stop ends
* its session. */
endSession?: (windDown: StructuredAgentSessionStopWindDown) => void
/** The host already withdrew queued messages for this Stop. */
withdrewQueued?: boolean
/** The session's child records: a background Stop reaches the tasks they offer a stop. */
@@ -70,6 +100,12 @@ export async function performCancel(
(input.turnId === undefined || input.turnId === liveTurnId)
const stoppedBefore =
runningCommand && structuredAgentSessionCommandWasStopped(ctx.journal, liveTurnId)
// Read while the child is live: a provider whose Stop is a session boundary loses it next.
const endsSession =
input.endSession !== undefined && ctx.adapter.stopEndsSession?.(ctx.sessionId) === true
const stoppedAt = Date.now()
// The provider's own answer; unset when its cancel threw, leaving the effect unknown.
let taken: boolean | undefined
try {
const dispatchStatus = latestJournalDispatchObservation(ctx.journal, ctx.fence)
const outcome: AgentSessionCancelOutcome = stoppedBefore
@@ -94,6 +130,7 @@ export async function performCancel(
...(dispatchStatus ? { dispatchStatus } : {}),
...(input.prompt ? { prompt: { itemId: input.prompt.itemId } } : {})
})
taken = outcome.cancelled
cancelled = outcome.cancelled
if (!cancelled && input.withdrewQueued && !(await isMainAgentWorkingOnceFlushed(ctx))) {
// A Stop that withdrew what was queued and left nothing working ended what it was sent for,
@@ -125,7 +162,21 @@ export async function performCancel(
...agentSessionFailureWords(agentSessionFailureFact('cancelUnconfirmed'), { surface: 'row' })
}
}
if (runningCommand && !cancelled) {
// A Stop naming a turn that has since ended keeps the session only when the provider declined
// it: an interrupt, answered or not, can stop a follow-up whose turn has not opened.
if (
endsSession &&
(input.turnId === undefined || input.turnId === liveTurnId || taken !== false)
) {
// An interrupt the provider took is worth waiting on, turn row or not: a Stop before the echo
// has none, and the echo still opens the turn the Stop interrupted.
input.endSession?.({ waitsForProvider: taken === true, stoppedAt })
cancelled = true
// The child's end confirms the Stop, so a refused or unconfirmed interrupt says nothing more.
if (note !== null) {
note = { kind: 'status', text: 'Cancellation requested.' }
}
} else if (runningCommand && !cancelled) {
await input.stopChild?.()
cancelled = true
note = { kind: 'status', text: 'Cancellation requested.' }
@@ -17,7 +17,7 @@ import type { StructuredAgentSessionHost } from './structured-agent-session-host
import {
mutateStructuredAgentSession,
type StructuredAgentSessionMutationContext
} from './structured-agent-session-host-mutations'
} from './structured-agent-session-mutation-context'
import type { StructuredAgentSessionCaller } from './structured-agent-session-host-types'
import {
conversationCommandPlan,
@@ -780,15 +780,35 @@ describe('a structured Claude session over agentSession.*', () => {
provider: 'claude',
leafUuid: 'assistant-leaf'
})
// A Claude Stop ends its child once Claude ends the stopped turn, so the chat rests; the next
// open resumes the conversation.
const old = claude.live()
const resumed = await ok<{ fence: number }>('agentSession.ensure', ensureParams(created.fence))
expect(resumed.fence).toBe(created.fence + 1)
old.handlers.onMessage?.({
type: 'result',
subtype: 'error_during_execution',
is_error: true,
session_id: PROVIDER_SESSION,
uuid: 'interrupted-result'
})
await vi.waitFor(() => expect(leaseOf(SESSION).claimStatus).toBe('released'))
expect(old.closed).toBe(true)
// Claude owns where the conversation continues; the stored leaf is the last completed turn.
const rested = leaseOf(SESSION)
const resumed = await ok<{ fence: number }>(
'agentSession.ensure',
ensureParams(rested.runtimeFence)
)
expect(resumed.fence).toBe(rested.runtimeFence + 1)
expect(claude.live()).not.toBe(old)
// Claude owns where the conversation continues; the stored leaf is the stopped turn's, which
// Claude ended with its own result.
expect(claude.live().launch.options).toMatchObject({ resume: PROVIDER_SESSION })
expect(claude.live().launch.options).not.toHaveProperty('resumeSessionAt')
const lastCompletedTurn = {
handle: { provider: 'claude', sessionId: PROVIDER_SESSION, leafUuid: 'assistant-leaf' },
handle: {
provider: 'claude',
sessionId: PROVIDER_SESSION,
leafUuid: 'provider-opened-assistant'
},
origin: 'resumed'
}
expect(host.deps.store.getRecord(SESSION).providerHandleChain.at(-1)).toMatchObject(