diff --git a/src/main/claude/claude-child-work-decoder.test.ts b/src/main/claude/claude-child-work-decoder.test.ts index 9891f6c3f84..0d11ba7a6a2 100644 --- a/src/main/claude/claude-child-work-decoder.test.ts +++ b/src/main/claude/claude-child-work-decoder.test.ts @@ -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() diff --git a/src/main/claude/claude-child-work-decoder.ts b/src/main/claude/claude-child-work-decoder.ts index c80c87710a6..e4cee787469 100644 --- a/src/main/claude/claude-child-work-decoder.ts +++ b/src/main/claude/claude-child-work-decoder.ts @@ -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() diff --git a/src/main/claude/claude-journal-translator-contract.ts b/src/main/claude/claude-journal-translator-contract.ts index ed508b30846..24884205ac2 100644 --- a/src/main/claude/claude-journal-translator-contract.ts +++ b/src/main/claude/claude-journal-translator-contract.ts @@ -11,15 +11,13 @@ import type { ClaudeCommandStart } from './claude-command-turn' export type ClaudeJournalTranslator = { handle: (event: ClaudeStructuredSessionEvent) => void - journalPrompts: Pick + journalPrompts: Pick /** 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. */ diff --git a/src/main/claude/claude-open-turn.ts b/src/main/claude/claude-open-turn.ts index 5aca644143b..dc5a9c8dcf7 100644 --- a/src/main/claude/claude-open-turn.ts +++ b/src/main/claude/claude-open-turn.ts @@ -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. */ diff --git a/src/main/claude/claude-prompt-registry.ts b/src/main/claude/claude-prompt-registry.ts index 5cd2f6b1a7b..bb8febab21b 100644 --- a/src/main/claude/claude-prompt-registry.ts +++ b/src/main/claude/claude-prompt-registry.ts @@ -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 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 - resolve: () => void -} - export function isClaudePromptRecord(value: unknown): value is Record { return typeof value === 'object' && value !== null && !Array.isArray(value) } @@ -78,12 +66,9 @@ function questionId(question: Record, index: number): string { /** Session-local callback ownership; none of this state is reconstructed from the transcript. */ export class ClaudePromptRegistry { private readonly prompts = new Map() - private readonly journalBindings = new Map() + /** Journal item id to the prompt key it shows. */ + private readonly journalBindings = new Map() private readonly claims = new Map() - 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 | null { - if (!this.ownsClaim(claim)) { - return null - } - let observation = this.cancellationObservations.get(claim.found.prompt) - if (!observation) { - let resolve = (): void => {} - const promise = new Promise((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 } } diff --git a/src/main/claude/claude-request-end-wait.test.ts b/src/main/claude/claude-request-end-wait.test.ts new file mode 100644 index 00000000000..7cca7a190dc --- /dev/null +++ b/src/main/claude/claude-request-end-wait.test.ts @@ -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, ms: number): Promise { + 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) + }) +}) diff --git a/src/main/claude/claude-request-end-wait.ts b/src/main/claude/claude-request-end-wait.ts new file mode 100644 index 00000000000..a176c2f96f1 --- /dev/null +++ b/src/main/claude/claude-request-end-wait.ts @@ -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 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; 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((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 { + let elapsed = false + let timer: ReturnType | undefined + const budget = new Promise((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 +): (sessionId: string, stoppedAt: number) => Promise { + 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() + } +} diff --git a/src/main/claude/claude-structured-child-work-captures.test.ts b/src/main/claude/claude-structured-child-work-captures.test.ts index d69ba120846..4ab9dd61b62 100644 --- a/src/main/claude/claude-structured-child-work-captures.test.ts +++ b/src/main/claude/claude-structured-child-work-captures.test.ts @@ -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' } ]) }) diff --git a/src/main/claude/claude-structured-child-work-producer.test.ts b/src/main/claude/claude-structured-child-work-producer.test.ts index 7c493f9e8c8..86d1d9a7ac4 100644 --- a/src/main/claude/claude-structured-child-work-producer.test.ts +++ b/src/main/claude/claude-structured-child-work-producer.test.ts @@ -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() }) diff --git a/src/main/claude/claude-structured-control-actions.test.ts b/src/main/claude/claude-structured-control-actions.test.ts index d2e29801020..b8629af1f3f 100644 --- a/src/main/claude/claude-structured-control-actions.test.ts +++ b/src/main/claude/claude-structured-control-actions.test.ts @@ -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(), diff --git a/src/main/claude/claude-structured-control-actions.ts b/src/main/claude/claude-structured-control-actions.ts index 37d2158f765..adf4d866fc3 100644 --- a/src/main/claude/claude-structured-control-actions.ts +++ b/src/main/claude/claude-structured-control-actions.ts @@ -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 diff --git a/src/main/claude/claude-structured-in-turn-stop.test.ts b/src/main/claude/claude-structured-in-turn-stop.test.ts index c5bbc51d8bc..04afb55bad6 100644 --- a/src/main/claude/claude-structured-in-turn-stop.test.ts +++ b/src/main/claude/claude-structured-in-turn-stop.test.ts @@ -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) } ) }) diff --git a/src/main/claude/claude-structured-inbound-control.ts b/src/main/claude/claude-structured-inbound-control.ts index fc2810e8a75..93e91bf6d18 100644 --- a/src/main/claude/claude-structured-inbound-control.ts +++ b/src/main/claude/claude-structured-inbound-control.ts @@ -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)) diff --git a/src/main/claude/claude-structured-journal-prompts.ts b/src/main/claude/claude-structured-journal-prompts.ts index 92301af106c..caf6de4dea5 100644 --- a/src/main/claude/claude-structured-journal-prompts.ts +++ b/src/main/claude/claude-structured-journal-prompts.ts @@ -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 diff --git a/src/main/claude/claude-structured-journal-translation.ts b/src/main/claude/claude-structured-journal-translation.ts index dcf3a3a6720..0d09485a987 100644 --- a/src/main/claude/claude-structured-journal-translation.ts +++ b/src/main/claude/claude-structured-journal-translation.ts @@ -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 }, diff --git a/src/main/claude/claude-structured-prompt-items.test.ts b/src/main/claude/claude-structured-prompt-items.test.ts index f032d4bcf8d..56d55a7e236 100644 --- a/src/main/claude/claude-structured-prompt-items.test.ts +++ b/src/main/claude/claude-structured-prompt-items.test.ts @@ -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' } }) diff --git a/src/main/claude/claude-structured-prompt-items.ts b/src/main/claude/claude-structured-prompt-items.ts index d50735a76bf..c74e4741fc1 100644 --- a/src/main/claude/claude-structured-prompt-items.ts +++ b/src/main/claude/claude-structured-prompt-items.ts @@ -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 = { - 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 } }) } diff --git a/src/main/claude/claude-structured-prompt-ownership.test.ts b/src/main/claude/claude-structured-prompt-ownership.test.ts index 31fd703088f..93f38f221ee 100644 --- a/src/main/claude/claude-structured-prompt-ownership.test.ts +++ b/src/main/claude/claude-structured-prompt-ownership.test.ts @@ -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( diff --git a/src/main/claude/claude-structured-prompt-ownership.ts b/src/main/claude/claude-structured-prompt-ownership.ts index 1014455324e..48b462ed174 100644 --- a/src/main/claude/claude-structured-prompt-ownership.ts +++ b/src/main/claude/claude-structured-prompt-ownership.ts @@ -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[0] type AnswerInput = Parameters[0] -export function admitClaudePromptCancellation(session: ClaudeSession, promptKey: string): boolean { - const admission = session.translator?.journalPrompts.cancel(promptKey) - return admission?.accepted ?? true -} - -function waitForClaudePromptCancellation( - observed: Promise, - timeoutMs = CLAUDE_DEFAULT_REQUEST_TIMEOUT_MS -): Promise { - let timer: ReturnType | null = null - const deadline = new Promise((_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, 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 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>[0] + sessions: Map +}): Promise { + 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) + } +} diff --git a/src/main/claude/claude-structured-prompt-replies.ts b/src/main/claude/claude-structured-prompt-replies.ts index 3bb99e3638a..a6007e4afdf 100644 --- a/src/main/claude/claude-structured-prompt-replies.ts +++ b/src/main/claude/claude-structured-prompt-replies.ts @@ -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 } } diff --git a/src/main/claude/claude-structured-resume-point.ts b/src/main/claude/claude-structured-resume-point.ts index e8124f52d27..b15d2e3507c 100644 --- a/src/main/claude/claude-structured-resume-point.ts +++ b/src/main/claude/claude-structured-resume-point.ts @@ -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) { diff --git a/src/main/claude/claude-structured-session-acquisition.ts b/src/main/claude/claude-structured-session-acquisition.ts index 0a160d99b93..92f84f6c489 100644 --- a/src/main/claude/claude-structured-session-acquisition.ts +++ b/src/main/claude/claude-structured-session-acquisition.ts @@ -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)) }) diff --git a/src/main/claude/claude-structured-session-adapter.ts b/src/main/claude/claude-structured-session-adapter.ts index 83c49e5a77a..46e1d993635 100644 --- a/src/main/claude/claude-structured-session-adapter.ts +++ b/src/main/claude/claude-structured-session-adapter.ts @@ -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 = async ( input ) => { diff --git a/src/main/claude/claude-structured-session-close.test.ts b/src/main/claude/claude-structured-session-close.test.ts index 34e4b6320f8..8881c6ae74b 100644 --- a/src/main/claude/claude-structured-session-close.test.ts +++ b/src/main/claude/claude-structured-session-close.test.ts @@ -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) diff --git a/src/main/claude/claude-structured-session-close.ts b/src/main/claude/claude-structured-session-close.ts index f1735125197..1e19bf5600e 100644 --- a/src/main/claude/claude-structured-session-close.ts +++ b/src/main/claude/claude-structured-session-close.ts @@ -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 { 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) diff --git a/src/main/claude/claude-turn-outcome.test.ts b/src/main/claude/claude-turn-outcome.test.ts index c8de34cee08..0dfb3080710 100644 --- a/src/main/claude/claude-turn-outcome.test.ts +++ b/src/main/claude/claude-turn-outcome.test.ts @@ -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 }) diff --git a/src/main/claude/claude-turn-ownership.test.ts b/src/main/claude/claude-turn-ownership.test.ts index a7255aab558..f150dd362b6 100644 --- a/src/main/claude/claude-turn-ownership.test.ts +++ b/src/main/claude/claude-turn-ownership.test.ts @@ -93,10 +93,13 @@ function sessionHoldingTurn(turnId: string | null): ReturnType () => {}, + 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]]) }) } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.test.ts index fb6c001b320..7de5dc5033b 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.test.ts @@ -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)) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.ts index ddf695630eb..8e43cfb914f 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.ts @@ -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 = (input) => + this.liveOwnerOrNull(input.sessionId)?.routePromptCancel?.(input) + + dismissPrompt: NonNullable = async (input) => + this.owner(input.sessionId).dismissPrompt?.(input) + readCommands: NonNullable = (sessionId) => this.liveOwnerOrNull(sessionId)?.readCommands?.(sessionId) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-stop.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-stop.ts new file mode 100644 index 00000000000..0d73ec723ce --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-stop.ts @@ -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 + /** 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 + }): Promise +} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts index 1d26ea427ae..4525cda08c3 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts @@ -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. */ diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-chat-stop.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-chat-stop.ts new file mode 100644 index 00000000000..af618040fa3 --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-chat-stop.ts @@ -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 + +/** 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( + context: StructuredAgentSessionMutationContext, + caller: StructuredAgentSessionCaller, + params: { envelope: AgentSessionMutationEnvelope; turnId?: string }, + plan: MutationPlan, + run: ( + ctx: AgentSessionTurnContext, + stop: () => Promise + ) => Promise> +): Promise> { + 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 => + 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 +} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-claude-compact-stop.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-claude-compact-stop.test.ts index b456bacac26..041b6580ea2 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-claude-compact-stop.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-claude-compact-stop.test.ts @@ -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' }) ) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-claude-stop-ends-session.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-claude-stop-ends-session.test.ts new file mode 100644 index 00000000000..ad314376b1a --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-claude-stop-ends-session.test.ts @@ -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 +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[] = [] + 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(assertion: () => T | Promise): Promise { + return vi.waitFor(assertion, { timeout: 10_000 }) +} + +function envelope( + method: 'agentSession.send' | 'agentSession.cancel', + fields: Record, + 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 { + 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): 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 { + 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 { + 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) => + runClaudeControl( + 'interrupt', + () => new Promise(() => {}), + typeof options?.timeoutMs === 'number' ? options.timeoutMs : undefined + ) + +async function interruptSent(connection: FakeConnection): Promise { + await eventually(() => + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(true) + ) +} + +async function turnOutcome(): Promise { + 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 { + 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 { + 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, + signal?: AbortSignal +): Promise<{ + answered: ReturnType + 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 { + 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([]) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-stop.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-stop.test.ts index 2eeae4787ce..488dcdfd2d0 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-stop.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-stop.test.ts @@ -36,6 +36,9 @@ let host: StructuredAgentSessionHost let dispatch: Mock let cancelTurn: Mock let awaitStarted: Mock> +let closeSession: Mock> +/** Codex's answer by default: its Stop keeps the child. */ +let stopEndsSession: boolean let events: StructuredAgentSessionEventSink | undefined function eventually(assertion: () => void | Promise): Promise { @@ -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 { + 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 { + 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.']) + }) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts index c88ae3a2f52..5067bcb3fb2 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts @@ -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 - publish: (sessionId: string, journal: StructuredAgentSessionHostSession['journal']) => void - flushStreamedEvents: (sessionId: string) => Promise - /** The host's accessor, for a caller outside the session's serialize. */ - conversation: (sessionId: string) => Promise - /** The session's child records, as the strip reads them; what command admission decides on. */ - readChildWork: (sessionId: string) => AgentChildWorkView[] | undefined - serialize: (sessionId: string, task: () => Promise) => Promise - /** The session's conversation, opened when closed; inside the caller's serialize. */ - openConversation: (sessionId: string) => Promise - /** Gives the session a provider child; inside the caller's serialize. */ - ensureAgent: (sessionId: string) => Promise - /** 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 - /** 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( - context: StructuredAgentSessionMutationContext, - caller: StructuredAgentSessionCaller, - envelope: AgentSessionMutationEnvelope, - plan: MutationPlan, - prepareSession?: AgentSessionMutationRequest['prepareSession'] -): Promise> { - 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> { - 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) ) } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-context.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-context.ts new file mode 100644 index 00000000000..73b47f1fe6e --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-context.ts @@ -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 + publish: (sessionId: string, journal: StructuredAgentSessionHostSession['journal']) => void + flushStreamedEvents: (sessionId: string) => Promise + /** The host's accessor, for a caller outside the session's serialize. */ + conversation: (sessionId: string) => Promise + /** The session's child records, as the strip reads them; what command admission decides on. */ + readChildWork: (sessionId: string) => AgentChildWorkView[] | undefined + serialize: (sessionId: string, task: () => Promise) => Promise + /** The session's conversation, opened when closed; inside the caller's serialize. */ + openConversation: (sessionId: string) => Promise + /** Gives the session a provider child; inside the caller's serialize. */ + ensureAgent: (sessionId: string) => Promise + /** 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 + /** 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( + context: StructuredAgentSessionMutationContext, + caller: StructuredAgentSessionCaller, + envelope: AgentSessionMutationEnvelope, + plan: MutationPlan, + prepareSession?: AgentSessionMutationRequest['prepareSession'] +): Promise> { + 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() + }) + ) +} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts index e18e65e6ab6..a81ea5270c8 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts @@ -173,7 +173,6 @@ export function cancelPlan(params: { scope?: 'background-tasks' taskId?: string prompt?: { itemId: string; expectedRevision: number } - stopChild?: () => Promise /** The session's child records, which name the tasks a background Stop reaches. */ childWork?: () => readonly AgentChildWorkView[] | undefined }): MutationPlan { @@ -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 diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-prompt-cancel.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-prompt-cancel.test.ts index b9a7bc7c5c4..6e0a44e60ca 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-prompt-cancel.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-prompt-cancel.test.ts @@ -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(async (input) => { + await input.commit() + }) + const dismissPrompt = vi.fn>( + 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() + }) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-prompt-cancel.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-prompt-cancel.ts new file mode 100644 index 00000000000..5e3272b4252 --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-prompt-cancel.ts @@ -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 +type PendingPrompt = Extract + +export async function cancelStructuredAgentSessionPrompt( + ctx: AgentSessionTurnContext, + input: { turnId?: string; prompt: { itemId: string; expectedRevision: number } }, + routes: { + stop: () => Promise + interrupt: () => Promise + } +): Promise { + 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> { + 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 => { + 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 } +} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-mutations.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-mutations.ts index 140573c5ff3..5b7f4bdf16c 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-mutations.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-mutations.ts @@ -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, diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts index ee6e7bf1179..cdb70ff41d4 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts @@ -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, + windDown: StructuredAgentSessionStopWindDown, + stopChild: () => Promise, + onError: (error: unknown) => void +): Promise { + 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 + /** 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.' } diff --git a/src/main/native-chat/agent-session-wire/structured-conversation-compaction.ts b/src/main/native-chat/agent-session-wire/structured-conversation-compaction.ts index 066afdc08db..98c04a2e840 100644 --- a/src/main/native-chat/agent-session-wire/structured-conversation-compaction.ts +++ b/src/main/native-chat/agent-session-wire/structured-conversation-compaction.ts @@ -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, diff --git a/src/main/runtime/claude-structured-session-integration.test.ts b/src/main/runtime/claude-structured-session-integration.test.ts index dc88bf858af..47199f6199f 100644 --- a/src/main/runtime/claude-structured-session-integration.test.ts +++ b/src/main/runtime/claude-structured-session-integration.test.ts @@ -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(