diff --git a/src/main/claude/claude-open-turn.ts b/src/main/claude/claude-open-turn.ts new file mode 100644 index 00000000000..6e9d4235db6 --- /dev/null +++ b/src/main/claude/claude-open-turn.ts @@ -0,0 +1,107 @@ +// The session's open turn, and the lifecycle row that publishes it. +// +// Sole owner of turn identity: the row this writes carries the same id it holds, +// and that row's id is what a client's Stop names. Readers ask here rather than +// keeping a copy, so there is nothing to disagree with. + +import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' +import { + claudeTurnLifecycleItem, + type ClaudeCurrentTurn, + type ClaudeTurnEnd +} from './claude-turn-lifecycle-item' +import { createClaudeTurnOpener, type ClaudeTurnSource } from './claude-turn-opening' + +export type ClaudeOpenTurnDeps = { + sink: StructuredAgentSessionEventSink + /** Settles the superseded turn's children; they get no later event of their own. */ + settleChildren: (groupKey: string | null) => void +} + +export class ClaudeOpenTurn { + private current: ClaudeCurrentTurn | null = null + /** Provider output may not reopen a turn after the session ended or a turn + * failed: nothing would ever close the turn it opened, and the row would read + * working for the life of the session. Only an accepted send lifts it. */ + private reopenSuppressed = false + private readonly opener: ( + frame: Record, + source: ClaudeTurnSource | null, + observedAt: number + ) => void + + constructor(private readonly deps: ClaudeOpenTurnDeps) { + this.opener = createClaudeTurnOpener({ + isTurnOpen: () => this.isOpen, + isSuppressed: () => this.reopenSuppressed, + open: (turn, observedAt) => this.open(turn, observedAt) + }) + } + + get id(): string | null { + return this.current?.turnId ?? null + } + + get groupKey(): string | null { + return this.current ? `${this.current.sessionId}:${this.current.turnId}` : null + } + + get isOpen(): boolean { + return this.current !== null + } + + /** Open a turn, ending whichever one was still open. A new turn starting is the + * only end the previous one gets when its result never arrives; settling it + * later would sweep THIS turn. */ + open(turn: ClaudeCurrentTurn, observedAt: number): void { + if (this.current) { + this.deps.settleChildren(this.groupKey) + this.publish(this.current, { state: 'interrupted', completedAt: observedAt }) + } + this.current = turn + this.publish(turn) + this.deps.sink.setActivity?.(null) + } + + /** The provider produced, so a turn is running. Idempotent: every frame of one + * reply stays inside the turn its first frame opened. A subagent's output is + * its parent turn's work and never a turn of its own. */ + ensureOpen( + frame: Record, + source: ClaudeTurnSource | null, + observedAt: number + ): void { + this.opener(frame, source, observedAt) + } + + /** End the open turn, if one is open, and clear the live activity line. */ + settle(end: ClaudeTurnEnd): void { + if (this.current) { + this.publish(this.current, end) + this.current = null + } + this.deps.sink.setActivity?.(null) + } + + /** An accepted send is the only thing that lifts the latch. */ + allowReopen(): void { + this.reopenSuppressed = false + } + + suppressReopen(): void { + this.reopenSuppressed = true + } + + /** A turn that failed is not resumed by whatever the provider says next; the + * next send is what resumes it. The latch only ever sets here. */ + suppressReopenOnFailure(failed: boolean): void { + this.reopenSuppressed ||= failed + } + + private publish(turn: ClaudeCurrentTurn, end?: ClaudeTurnEnd): void { + const item = claudeTurnLifecycleItem(turn, end) + this.deps.sink.appendItem(item.identity, item.body, item.options) + // Preserve first-work evidence when completion arrives before the journal drains. + this.deps.sink.publish({ coalescingKey: item.publishCoalescingKey }) + } +} diff --git a/src/main/claude/claude-structured-control-actions.test.ts b/src/main/claude/claude-structured-control-actions.test.ts index e9afa9bac25..9d9b8310206 100644 --- a/src/main/claude/claude-structured-control-actions.test.ts +++ b/src/main/claude/claude-structured-control-actions.test.ts @@ -197,6 +197,7 @@ describe('answerClaudePrompt', () => { cancel: vi.fn(() => ({ accepted: true as const })), resolve: resolvePrompt }, + currentTurnId: null, flush: vi.fn(), pendingStreamedBlocks: 0, dispose: vi.fn() diff --git a/src/main/claude/claude-structured-dispatch-admission.test.ts b/src/main/claude/claude-structured-dispatch-admission.test.ts index 6ed8f91f072..57e7a56e59f 100644 --- a/src/main/claude/claude-structured-dispatch-admission.test.ts +++ b/src/main/claude/claude-structured-dispatch-admission.test.ts @@ -48,7 +48,7 @@ describe('Claude structured dispatch admission', () => { clientMessageId: 'client-2', providerIdentity: { provider: 'claude', sessionId: 'provider-session', uuid: queuedUuid } }) - expect(session.activeTurnId).toBe(queuedUuid) + expect(session.dispatchWaiters).toHaveLength(0) } finally { vi.useRealTimers() } diff --git a/src/main/claude/claude-structured-dispatch.test.ts b/src/main/claude/claude-structured-dispatch.test.ts index 05a3ce8092f..d5b8bc2451f 100644 --- a/src/main/claude/claude-structured-dispatch.test.ts +++ b/src/main/claude/claude-structured-dispatch.test.ts @@ -38,7 +38,7 @@ describe('Claude structured dispatch image limits', () => { } ) - it('takes the active turn identity from a replay that lands after dispatch returned', async () => { + it('settles the waiter from a replay that lands after dispatch returned', async () => { const session = sessionFor() const dispatched = dispatchClaudeTurn(session, { clientMessageId: 'client-1', @@ -47,11 +47,9 @@ describe('Claude structured dispatch image limits', () => { await vi.waitFor(() => expect(session.dispatchWaiters).toHaveLength(1)) const sentUuid = (session.dispatchWaiters[0] as { sentUuid?: string }).sentUuid await expect(dispatched).resolves.toEqual({ state: 'admitted' }) - expect(session.activeTurnId).toBeUndefined() expect(resolveClaudeReplayWaiter(session, userReplayFrame(sentUuid!, 'one'))).toBe(true) - expect(session.activeTurnId).toBe(sentUuid) - expect(session.activeTurnSequence).toBe(session.dispatchSequence) + expect(session.dispatchWaiters).toHaveLength(0) }) it('recovers the active identity when a replay lands after the child died', async () => { @@ -68,8 +66,7 @@ describe('Claude structured dispatch image limits', () => { expect(session.retiredDispatchWaiters).toHaveLength(1) expect(resolveClaudeReplayWaiter(session, userReplayFrame(sentUuid!, 'one'))).toBe(true) - expect(session.activeTurnId).toBe(sentUuid) - expect(session.activeTurnSequence).toBe(session.dispatchSequence) + expect(session.retiredDispatchWaiters).toHaveLength(0) }) it('settles the send the replay proves was delivered, whenever it arrives', async () => { @@ -410,7 +407,7 @@ describe('Claude structured dispatch image limits', () => { expect(resolveClaudeReplayWaiter(session, userReplayFrame('fresh-replay', 'retry me'))).toBe( true ) - expect(session.activeTurnId).toBe('fresh-replay') + expect(session.dispatchWaiters).toHaveLength(0) }) it('does not claim an SDK-pulled frame was unwritten when its write outcome is ambiguous', async () => { diff --git a/src/main/claude/claude-structured-dispatch.ts b/src/main/claude/claude-structured-dispatch.ts index 84253b765f6..6a41d42c81c 100644 --- a/src/main/claude/claude-structured-dispatch.ts +++ b/src/main/claude/claude-structured-dispatch.ts @@ -30,6 +30,14 @@ const MAX_ACTIVE_DISPATCH_WAITERS = 64 /** Settles a provider-proven late outcome; replay rows independently reconcile acceptance. */ export type ClaudeLateDispatchSettlement = (input: ClaudeLateDispatchOutcome) => void +/** A send is still awaiting its echo, so an interrupt would let it loose as an + * unexpected turn unless the CLI cancels the queue in the same round trip. + * Derived from the live waiters: a retired one is no longer awaited, and gating + * Stop on it would strand the user for the life of the session. */ +export function claudeHasUnsettledDispatch(session: ClaudeSession): boolean { + return session.dispatchWaiters.length > 0 +} + export function resolveClaudeReplayWaiter( session: ClaudeSession, message: Record, @@ -144,18 +152,13 @@ function settleWaiter( } waiter.settledUuid = uuid waiter.resolve(uuid) - // Dispatch returned on admission. Settle delivery unfenced while the sequence - // still fences which turn owns the identity; see `recoverLateIdentity`. + // Dispatch returned on admission, so the replay is what settles delivery. if (waiter.clientMessageId) { onSettledLate?.({ clientMessageId: waiter.clientMessageId, providerIdentity: { provider: 'claude', sessionId: session.providerSessionId, uuid } }) } - if (waiter.dispatchSequence === session.dispatchSequence) { - session.activeTurnId = uuid - session.activeTurnSequence = waiter.dispatchSequence - } } function forgetRetiredWaiter(session: ClaudeSession, waiter: ClaudeDispatchWaiter): void { @@ -176,18 +179,14 @@ function recoverLateIdentity( return false } // The provider acted on this dispatch, so the send it came from is delivered. - // Unfenced on purpose: the dispatch-sequence check below only decides which - // turn owns the identity, while delivery is settled for good either way. + // Unfenced on purpose: the dispatch-sequence check below only decides whether + // this replay still opens a turn, while delivery is settled for good either way. if (waiter.clientMessageId) { onSettledLate?.({ clientMessageId: waiter.clientMessageId, providerIdentity: { provider: 'claude', sessionId: session.providerSessionId, uuid } }) } - if (waiter.dispatchSequence === session.dispatchSequence) { - session.activeTurnId = uuid - session.activeTurnSequence = waiter.dispatchSequence - } return isUserReplay && waiter.dispatchSequence === session.dispatchSequence } @@ -293,7 +292,7 @@ export async function dispatchClaudeTurn( if (session.dispatchWaiters.length >= MAX_ACTIVE_DISPATCH_WAITERS) { return { state: 'rejected', reason: DISPATCH_REJECTED_QUEUE_FULL } } - const dispatchSequence = ++session.dispatchSequence + ++session.dispatchSequence // Read the sent content, not the journal blocks: only the mapped trailing prompt decides // whether Claude runs a command, so the two cannot disagree about which frame settles this. const acceptsResult = claudeDispatchInvokesSlashCommand(content) @@ -319,8 +318,6 @@ export async function dispatchClaudeTurn( if (waiter.settledUuid) { const uuid = await replayed if (uuid) { - session.activeTurnId = uuid - session.activeTurnSequence = dispatchSequence return { state: 'accepted', providerIdentity: { provider: 'claude', sessionId: session.providerSessionId, uuid } diff --git a/src/main/claude/claude-structured-journal-translation.ts b/src/main/claude/claude-structured-journal-translation.ts index 216e760315c..e34b4aed111 100644 --- a/src/main/claude/claude-structured-journal-translation.ts +++ b/src/main/claude/claude-structured-journal-translation.ts @@ -36,16 +36,11 @@ import { claudeStreamTurnStartSource, claudeStreamTurnSource, claudeTurnOpenedBySendEcho, - createClaudeTurnOpener, isRootClaudeFrame, type ClaudeTurnSource } from './claude-turn-opening' -import { - claudeTurnEndForResult, - claudeTurnLifecycleItem, - type ClaudeCurrentTurn, - type ClaudeTurnEnd -} from './claude-turn-lifecycle-item' +import { claudeTurnEndForResult } from './claude-turn-lifecycle-item' +import { ClaudeOpenTurn } from './claude-open-turn' import { ClaudeJournalPrompts } from './claude-structured-journal-prompts' export type ClaudeJournalTranslatorDeps = { @@ -59,6 +54,9 @@ export type ClaudeJournalTranslatorDeps = { export type ClaudeJournalTranslator = { handle: (event: ClaudeStructuredSessionEvent) => void 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 flush: () => void /** Streamed blocks still awaiting a final frame. A settled turn leaves none. */ readonly pendingStreamedBlocks: number @@ -86,20 +84,17 @@ export function createClaudeJournalTranslator( const tools = new Map() const prompts = new ClaudeJournalPrompts(deps) const streamedBlocks = createClaudeStreamedBlockRegistry() - let currentTurn: ClaudeCurrentTurn | null = null - /** Provider output may not reopen a turn after the session ended or a turn - * failed: nothing would ever close the turn it opened, and the row would read - * working for the life of the session. Only an accepted send lifts it. */ - let reopenSuppressed = false - const groupKeyOf = (turn: ClaudeCurrentTurn | null): string | null => - turn ? `${turn.sessionId}:${turn.turnId}` : null + const turn = new ClaudeOpenTurn({ + sink: deps.sink, + settleChildren: (groupKey) => subagents.settleTurn(groupKey) + }) const providerFallback = createClaudeProviderFrameFallback( deps.sink, deps.fallbackIdPrefix ?? 'acquisition' ) const subagents = new ClaudeSubagentRoster({ sink: deps.sink, - currentGroupKey: () => groupKeyOf(currentTurn) + currentGroupKey: () => turn.groupKey }) const streamedText = createClaudeStreamedTextCheckpoints({ ...(deps.coalesceMs === undefined ? {} : { coalesceMs: deps.coalesceMs }), @@ -110,42 +105,14 @@ export function createClaudeJournalTranslator( } }) - const publishLifecycle = (turn: ClaudeCurrentTurn, end?: ClaudeTurnEnd): void => { - const item = claudeTurnLifecycleItem(turn, end) - deps.sink.appendItem(item.identity, item.body, item.options) - // Preserve first-work evidence when completion arrives before the journal drains. - deps.sink.publish({ coalescingKey: item.publishCoalescingKey }) - } - - /** Open a turn, ending whichever one was still open. A new turn starting is the - * only end the previous one gets when its result never arrives; settling it - * later would sweep THIS turn. */ - const openTurn = (turn: ClaudeCurrentTurn, observedAt: number): void => { - if (currentTurn) { - subagents.settleTurn(groupKeyOf(currentTurn)) - publishLifecycle(currentTurn, { state: 'interrupted', completedAt: observedAt }) - } - currentTurn = turn - publishLifecycle(turn) - deps.sink.setActivity?.(null) - } - - /** The provider produced, so a turn is running. Idempotent: every frame of one - * reply stays inside the turn its first frame opened. A subagent's output is - * its parent turn's work and never a turn of its own. */ - const ensureTurnOpen = createClaudeTurnOpener({ - isTurnOpen: () => currentTurn !== null, - isSuppressed: () => reopenSuppressed, - open: openTurn - }) - const publishActivity = (kind: string, payload: unknown): void => { - if (!currentTurn) { + const turnId = turn.id + if (turnId === null) { return } const text = claudeProviderFrameActivity(kind, payload) if (text !== undefined) { - deps.sink.setActivity?.(text ? { turnId: currentTurn.turnId, text } : null) + deps.sink.setActivity?.(text ? { turnId, text } : null) } } @@ -154,7 +121,7 @@ export function createClaudeJournalTranslator( // `message_start` is the provider's turn boundary. Keep the first text // delta as a compatibility fallback for streams that omit it. const source = delta ? claudeStreamTurnSource(message) : claudeStreamTurnStartSource(message) - ensureTurnOpen(message, source, observedAt) + turn.ensureOpen(message, source, observedAt) if (!delta) { return false } @@ -188,17 +155,17 @@ export function createClaudeJournalTranslator( uuid: envelope.uuid, assistant: envelope.role === 'assistant' } - const openOutputTurn = (): void => ensureTurnOpen(message, source, observedAt) + const openOutputTurn = (): void => turn.ensureOpen(message, source, observedAt) if (body) { // Opening before the append is what brackets a turn around its own first // output; a reader that scans back to the turn record and stops would // otherwise look straight past the row that opened it. - ensureTurnOpen(message, source, observedAt) + turn.ensureOpen(message, source, observedAt) deps.sink.appendItem(identity, body) changed = true } for (const tool of claudeToolUses(outputEnvelope)) { - ensureTurnOpen(message, source, observedAt) + turn.ensureOpen(message, source, observedAt) tools.set(tool.id, tool) deps.sink.appendItem( claudeToolIdentity(envelope.sessionId, tool.id), @@ -223,7 +190,7 @@ export function createClaudeJournalTranslator( changed = true } if (thinking) { - ensureTurnOpen(message, source, observedAt) + turn.ensureOpen(message, source, observedAt) deps.sink.appendItem(claudeThinkingIdentity(envelope.sessionId, envelope.uuid), { kind: 'message', role: 'reasoning', @@ -244,8 +211,8 @@ export function createClaudeJournalTranslator( userItemId: agentJournalItemKey(identity) }) if (sendEchoTurn) { - reopenSuppressed = false - openTurn(sendEchoTurn, observedAt) + turn.allowReopen() + turn.open(sendEchoTurn, observedAt) } if (changed) { deps.sink.publish() @@ -260,18 +227,11 @@ export function createClaudeJournalTranslator( streamedText.flush() // No event will ever settle a child once the provider is gone. subagents.settleSession() - if (currentTurn) { - // The host saw the child end, so the turn's end is observed, not lost. - publishLifecycle(currentTurn, { - state: 'interrupted', - completedAt: event.observedAt ?? Date.now() - }) - currentTurn = null - } + // The host saw the child end, so the turn's end is observed, not lost. + turn.settle({ state: 'interrupted', completedAt: event.observedAt ?? Date.now() }) // A frame that arrives after the child is gone must not open a turn no // event can close. - reopenSuppressed = true - deps.sink.setActivity?.(null) + turn.suppressReopen() return } if (event.type === 'message' && handleStream(event.message, event.observedAt ?? Date.now())) { @@ -291,21 +251,11 @@ export function createClaudeJournalTranslator( const settlesTurn = isRootClaudeFrame(event.message) if (settlesTurn) { prompts.retryPendingCancellations() + turn.suppressReopenOnFailure(event.message.is_error === true) // The turn is over however it ended, so a foreground child still // reported as working will never be settled by an event. - // A turn that failed, or that the user stopped, is not resumed by - // whatever the provider says next; the next send is what resumes it. - // The latch only ever sets here; an accepted send is what lifts it. - reopenSuppressed ||= event.message.is_error === true - subagents.settleTurn(groupKeyOf(currentTurn)) - if (currentTurn) { - publishLifecycle( - currentTurn, - claudeTurnEndForResult(event.message, event.observedAt ?? Date.now()) - ) - currentTurn = null - } - deps.sink.setActivity?.(null) + subagents.settleTurn(turn.groupKey) + turn.settle(claudeTurnEndForResult(event.message, event.observedAt ?? Date.now())) // The turn is over. A block still awaiting its final keeps the text the // flush above journaled, but its live state goes: an interrupted turn // would otherwise retain that text for the life of the session. @@ -335,6 +285,9 @@ export function createClaudeJournalTranslator( } }, journalPrompts: prompts, + get currentTurnId() { + return turn.id + }, flush: streamedText.flush, get pendingStreamedBlocks() { return streamedText.pending diff --git a/src/main/claude/claude-structured-prompt-ownership.ts b/src/main/claude/claude-structured-prompt-ownership.ts index 1dd73f23552..7694fe3d8eb 100644 --- a/src/main/claude/claude-structured-prompt-ownership.ts +++ b/src/main/claude/claude-structured-prompt-ownership.ts @@ -9,7 +9,10 @@ import { cancelClaudeTurn, supportsClaudeQueuedInterruptCancellation } from './claude-structured-control-actions' -import type { ClaudeLateDispatchSettlement } from './claude-structured-dispatch' +import { + claudeHasUnsettledDispatch, + type ClaudeLateDispatchSettlement +} from './claude-structured-dispatch' import type { ClaudeSession } from './claude-structured-session-state' type CancelInput = Parameters[0] @@ -71,20 +74,22 @@ export async function cancelClaudeStructuredTurn(input: { session.prompts.releaseClaim(claim) return { cancelled: false } } + // The translator owns turn identity. A session with no journal has published no + // turn row for a client to name, so it holds no identity this request can contradict. + const ownsRequestedTurn = (): boolean => { + const currentTurnId = session.translator?.currentTurnId ?? null + return currentTurnId === null || currentTurnId === request.turnId + } const isCurrent = (): boolean => sessions.get(request.sessionId) === session && session.fence === request.fence && session.acquisitionGeneration === acquisitionGeneration && (claim && prompt - ? session.activeTurnId === request.turnId && + ? ownsRequestedTurn() && session.prompts.ownsBoundClaim(claim, prompt.itemId, request.turnId) && - (session.activeTurnSequence === session.dispatchSequence || - supportsClaudeQueuedInterruptCancellation(session)) + (!claudeHasUnsettledDispatch(session) || supportsClaudeQueuedInterruptCancellation(session)) : compactions.ownsTurn(request.sessionId, request.turnId) || - (session.activeTurnId === undefined - ? session.dispatchSequence === 0 - : session.activeTurnId === request.turnId && - session.activeTurnSequence === session.dispatchSequence)) + (ownsRequestedTurn() && !claudeHasUnsettledDispatch(session))) let interruptConfirmed = false try { const result = await cancelClaudeTurn( diff --git a/src/main/claude/claude-structured-session-acquisition.ts b/src/main/claude/claude-structured-session-acquisition.ts index 8870a0e1daa..423205c4817 100644 --- a/src/main/claude/claude-structured-session-acquisition.ts +++ b/src/main/claude/claude-structured-session-acquisition.ts @@ -142,7 +142,7 @@ export async function acquireClaudeSession({ const { canUseTool, onUserDialog } = buildClaudePermissionCallbacks({ sessionId, prompts, - currentTurnId: () => liveSession?.activeTurnId ?? null, + 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 f62541b50e9..9357e9635f6 100644 --- a/src/main/claude/claude-structured-session-adapter.ts +++ b/src/main/claude/claude-structured-session-adapter.ts @@ -232,7 +232,7 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda journalItemId, promptKey, questionId, - session.activeTurnId ?? null + session.translator?.currentTurnId ?? null ) } diff --git a/src/main/claude/claude-structured-session-state.ts b/src/main/claude/claude-structured-session-state.ts index 30e830af64e..4c4808b64d8 100644 --- a/src/main/claude/claude-structured-session-state.ts +++ b/src/main/claude/claude-structured-session-state.ts @@ -147,16 +147,12 @@ export type ClaudeSession = { restoreSkippedOptions: Set /** CLI-advertised protocol capabilities from init; gates interrupt-receipt handling. */ capabilities: readonly string[] - /** Provider uuid of the most recently admitted turn, if one is active. */ - activeTurnId?: string backgroundTasks: ClaudeBackgroundTaskTracker /** The `/` surface the CLI reports for itself; seeded from init, kept current * by later init and `commands_changed` frames. */ commands: ClaudeSlashCommandCatalog /** Monotonic fence advanced when a dispatch starts, including unresolved dispatches. */ dispatchSequence: number - /** Dispatch sequence that admitted activeTurnId. */ - activeTurnSequence?: number /** Fences overlapping option writes so a late completion cannot restore stale state. */ optionMutationSequence: number /** Shared durable-close write; a failed write clears this for a retry. */ diff --git a/src/main/claude/claude-structured-session-test-support.ts b/src/main/claude/claude-structured-session-test-support.ts index e728263d058..352a7ed18ec 100644 --- a/src/main/claude/claude-structured-session-test-support.ts +++ b/src/main/claude/claude-structured-session-test-support.ts @@ -14,6 +14,7 @@ import { type ClaudeStructuredLaunch, type ClaudeStructuredSessionEvent } from './claude-structured-session-adapter' +import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' export const PROVIDER_SESSION_ID = '819cf9f8-e43c-4ad7-b50f-54aa158a726a' @@ -245,10 +246,21 @@ export async function acquired( undefined, onDispatchSettledLate ) - await adapter.acquire({ identity: identityFor(), fence: 7, spawnToken: 'spawn-9' }) + await adapter.acquire({ + identity: identityFor(), + fence: 7, + spawnToken: 'spawn-9', + // Production acquires with a journal sink, and turn identity lives on the + // translator it builds; without one this fixture models no session that ships. + events: recordingJournalSink() + }) return adapter } +export function recordingJournalSink(): StructuredAgentSessionEventSink { + return { appendItem: () => {}, appendTombstone: () => {}, publish: () => {} } +} + export function tick(): Promise { return new Promise((resolve) => setImmediate(resolve)) } diff --git a/src/main/claude/claude-turn-ownership.test.ts b/src/main/claude/claude-turn-ownership.test.ts new file mode 100644 index 00000000000..05bd15bc630 --- /dev/null +++ b/src/main/claude/claude-turn-ownership.test.ts @@ -0,0 +1,152 @@ +// Which turn a Stop is allowed to interrupt, for turns the provider opened on its +// own as well as turns Orca's own send echo opened. + +import { describe, expect, it, vi } from 'vitest' +import type { AgentJournalItemBody } from '../../shared/agent-session-journal-types' +import { readAgentJournalTurn } from '../../shared/agent-session-turn-record' +import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' +import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key' +import { + PROVIDER_SESSION_ID, + USER_MESSAGE, + adapterFor, + fakeClaude, + identityFor, + type FakeConnection +} from './claude-structured-session-test-support' + +function journalSink(): { + sink: StructuredAgentSessionEventSink + bodies: Map +} { + const bodies = new Map() + return { + bodies, + sink: { + appendItem: (identity, body) => bodies.set(agentJournalItemKey(identity), body), + appendTombstone: (identity) => bodies.delete(agentJournalItemKey(identity)), + publish: vi.fn() + } + } +} + +/** The turn row a client would read, which is the id its Stop carries. */ +function runningTurnId(bodies: Map): string | null { + for (const body of bodies.values()) { + const turn = readAgentJournalTurn(body) + if (turn?.state === 'running') { + return turn.turnId + } + } + return null +} + +async function acquiredWithJournal(claude: ReturnType): Promise<{ + adapter: ReturnType + bodies: Map + connection: FakeConnection +}> { + const { sink, bodies } = journalSink() + const adapter = adapterFor(claude) + await adapter.acquire({ + identity: identityFor(), + fence: 7, + spawnToken: 'spawn-9', + events: sink + }) + const connection = claude.connections[0] + if (!connection) { + throw new Error('expected Claude connection') + } + return { adapter, bodies, connection } +} + +function completeTurn(connection: FakeConnection, uuid: string): void { + connection.handlers.onMessage?.({ + type: 'result', + subtype: 'success', + uuid, + session_id: PROVIDER_SESSION_ID, + is_error: false, + terminal_reason: 'completed', + duration_ms: 12 + }) +} + +/** The provider resuming on its own — a background task reporting in wakes the agent. */ +function providerOutput(connection: FakeConnection, uuid: string): void { + connection.handlers.onMessage?.({ + type: 'assistant', + uuid, + session_id: PROVIDER_SESSION_ID, + parent_tool_use_id: null, + message: { role: 'assistant', content: [{ type: 'text', text: 'picking this back up' }] } + }) +} + +describe('Claude turn ownership', () => { + it('stops a turn the provider opened after the session already dispatched once', async () => { + const claude = fakeClaude({ replayUuid: 'echo-turn' }) + const { adapter, bodies, connection } = await acquiredWithJournal(claude) + + await adapter.dispatch({ + sessionId: 'session-1', + clientMessageId: 'client-1', + body: USER_MESSAGE, + fence: 7 + }) + expect(runningTurnId(bodies)).toBe('echo-turn') + completeTurn(connection, 'result-1') + expect(runningTurnId(bodies)).toBeNull() + + providerOutput(connection, 'provider-turn') + // The client cancels with the journal row's id, which is the provider frame's. + expect(runningTurnId(bodies)).toBe('provider-turn') + + await expect( + adapter.cancelTurn({ sessionId: 'session-1', turnId: 'provider-turn', fence: 7 }) + ).resolves.toEqual({ cancelled: true }) + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(true) + }) + + it('still stops an echo-opened turn', async () => { + const claude = fakeClaude({ replayUuid: 'echo-turn' }) + const { adapter, bodies, connection } = await acquiredWithJournal(claude) + + await adapter.dispatch({ + sessionId: 'session-1', + clientMessageId: 'client-1', + body: USER_MESSAGE, + fence: 7 + }) + expect(runningTurnId(bodies)).toBe('echo-turn') + + await expect( + adapter.cancelTurn({ sessionId: 'session-1', turnId: 'echo-turn', fence: 7 }) + ).resolves.toEqual({ cancelled: true }) + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(true) + }) + + it('refuses a stale turn id once the provider opened a newer turn', async () => { + const claude = fakeClaude({ replayUuid: 'echo-turn' }) + const { adapter, bodies, connection } = await acquiredWithJournal(claude) + + await adapter.dispatch({ + sessionId: 'session-1', + clientMessageId: 'client-1', + body: USER_MESSAGE, + fence: 7 + }) + completeTurn(connection, 'result-1') + providerOutput(connection, 'provider-turn') + expect(runningTurnId(bodies)).toBe('provider-turn') + + await expect( + adapter.cancelTurn({ sessionId: 'session-1', turnId: 'echo-turn', fence: 7 }) + ).resolves.toEqual({ cancelled: false }) + await expect( + adapter.cancelTurn({ sessionId: 'session-1', turnId: 'not-a-turn', fence: 7 }) + ).resolves.toEqual({ cancelled: false }) + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(false) + }) +})