From a9232e8db6d16a2c9f3e8e694ff804f92f0c1541 Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Tue, 15 Sep 2026 16:54:43 -0700 Subject: [PATCH 1/3] fix(claude): single-own turn identity so Stop reaches a provider-opened turn (#20794) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(claude): single-own turn identity so Stop reaches a provider-opened turn Stop silently failed on any Claude turn the provider opened on its own — a background task reporting in wakes the agent — once the session had dispatched at least once. The transcript read "The provider had already finished this turn." while the model kept working. Turn identity was minted twice from the same stream by two components that never talked. The journal translator writes turnId into the durable turn row, which is the id every client's Stop carries. settleWaiter separately wrote session.activeTurnId, only ever on the dispatch-echo path, and nothing cleared it. Cancel read the adapter's copy; prompt binding, status and both clients read the journal's. They agreed only when a send echo opened the turn. Turn identity is now single-owned. The open turn moves out of the translator's closure into ClaudeOpenTurn, which holds the turn and publishes its lifecycle row, so the id readers ask for is the id the row carries. activeTurnId and activeTurnSequence are deleted rather than widened, so the second writer goes with them instead of a second guard being added beside the first. activeTurnSequence was never turn identity: it asked whether a send was still awaiting its echo, which an interrupt would release as an unexpected turn. That is now derived from the live dispatch waiters. Deriving it also retires a latch — a retired waiter left the stored sequence permanently behind the dispatch sequence, refusing every later Stop for the life of the session. Also fixes the mirror defect the same hazard caused: a stale turn id was accepted against a newer provider-opened turn, because activeTurnId was never cleared when a turn ended. The Claude adapter fixture now acquires with a journal sink, as production does; without one it modelled a session that never ships. * fix(claude): reject stale stop after turn settles * fix(claude): preserve dispatch cancellation fence * test(claude): cover provider-opened stop integration * fix(claude): derive dispatch cancellation fence from journal * fix(claude): honor journal dispatch status before local sends * fix(native-chat): omit absent dispatch observation * fix(claude): release unresolved stop fence after deadline * fix(claude): bound and poll dispatch admission wait * test(claude): cover dispatch admission fast path --- src/main/claude/claude-open-turn.ts | 107 +++++ .../claude-structured-control-actions.test.ts | 1 + ...aude-structured-dispatch-admission.test.ts | 2 +- .../claude/claude-structured-dispatch.test.ts | 11 +- src/main/claude/claude-structured-dispatch.ts | 19 +- .../claude-structured-journal-translation.ts | 105 ++--- .../claude-structured-prompt-ownership.ts | 96 ++++- .../claude-structured-session-acquisition.ts | 2 +- .../claude-structured-session-adapter.ts | 2 +- .../claude/claude-structured-session-state.ts | 4 - .../claude-structured-session-test-support.ts | 14 +- src/main/claude/claude-turn-ownership.test.ts | 374 ++++++++++++++++++ .../journal-dispatch-observation.test.ts | 65 +++ .../journal-dispatch-observation.ts | 29 ++ .../structured-agent-session-adapter.ts | 3 + .../structured-agent-session-host-handoff.ts | 15 +- .../structured-agent-session-turns.ts | 3 + ...ude-structured-session-integration.test.ts | 26 +- 18 files changed, 757 insertions(+), 121 deletions(-) create mode 100644 src/main/claude/claude-open-turn.ts create mode 100644 src/main/claude/claude-turn-ownership.test.ts create mode 100644 src/main/native-chat/agent-session-journal/journal-dispatch-observation.test.ts create mode 100644 src/main/native-chat/agent-session-journal/journal-dispatch-observation.ts 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..a69508ce071 100644 --- a/src/main/claude/claude-structured-dispatch.ts +++ b/src/main/claude/claude-structured-dispatch.ts @@ -144,18 +144,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 +171,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 +284,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 +310,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..a76825887e1 100644 --- a/src/main/claude/claude-structured-prompt-ownership.ts +++ b/src/main/claude/claude-structured-prompt-ownership.ts @@ -12,6 +12,11 @@ import { import type { ClaudeLateDispatchSettlement } from './claude-structured-dispatch' import type { ClaudeSession } from './claude-structured-session-state' +/** Conservative user-facing window: below the 10s init and 30s control deadlines, trading + * residual slow-pump risk for ensuring delivery bookkeeping cannot block Stop indefinitely. */ +export const CLAUDE_DISPATCH_ADMISSION_TIMEOUT_MS = 3_000 +const CLAUDE_DISPATCH_ADMISSION_POLL_MS = 50 + type CancelInput = Parameters[0] type AnswerInput = Parameters[0] @@ -47,6 +52,40 @@ function requireSession(sessions: Map, sessionId: string) return session } +function waitForClaudeDispatchAdmission( + admitted: () => boolean, + timeoutMs = CLAUDE_DISPATCH_ADMISSION_TIMEOUT_MS +): Promise { + return new Promise((resolve) => { + let settled = false + let deadline: ReturnType | null = null + let poll: ReturnType | null = null + const finish = (value: boolean): void => { + if (settled) { + return + } + settled = true + if (deadline) { + clearTimeout(deadline) + } + if (poll) { + clearInterval(poll) + } + resolve(value) + } + const check = (): void => { + if (admitted()) { + finish(true) + } + } + deadline = setTimeout(() => finish(false), timeoutMs) + poll = setInterval(check, CLAUDE_DISPATCH_ADMISSION_POLL_MS) + check() + deadline.unref?.() + poll.unref?.() + }) +} + export async function cancelClaudeStructuredTurn(input: { request: CancelInput sessions: Map @@ -71,20 +110,61 @@ 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 translator = session.translator + if (!translator) { + return session.dispatchSequence === 0 + } + const currentTurnId = translator.currentTurnId + return currentTurnId === null + ? session.dispatchSequence === 0 + : currentTurnId === request.turnId + } + // The host supplies the durable latest submission; direct adapter callers fall back to + // the current in-memory waiter so an unknown dispatch remains fenced without a latch. + const dispatchAdmissionIsCurrent = (): boolean => + request.dispatchStatus + ? request.dispatchStatus.state === 'accepted' || + request.dispatchStatus.state === 'rejected' || + (request.dispatchStatus.state === 'unknown' && request.dispatchStatus.recovered) + : session.dispatchSequence === 0 || + ![...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. Ordinary Stop gets a bounded escape below. + const dispatchAdmissionAllowsCancellation = (): boolean => + dispatchAdmissionIsCurrent() || + (Boolean(prompt) && supportsClaudeQueuedInterruptCancellation(session)) + const compactionOwnsTurn = (): boolean => compactions.ownsTurn(request.sessionId, request.turnId) + const currentDispatchHasRetiredWaiter = (): boolean => + session.retiredDispatchWaiters.some( + (waiter) => waiter.dispatchSequence === session.dispatchSequence + ) + let dispatchAdmissionExpired = false + if ( + !prompt && + !compactionOwnsTurn() && + !dispatchAdmissionAllowsCancellation() && + (request.dispatchStatus !== undefined || currentDispatchHasRetiredWaiter()) + ) { + dispatchAdmissionExpired = !(await waitForClaudeDispatchAdmission( + dispatchAdmissionAllowsCancellation + )) + } 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)) - : compactions.ownsTurn(request.sessionId, request.turnId) || - (session.activeTurnId === undefined - ? session.dispatchSequence === 0 - : session.activeTurnId === request.turnId && - session.activeTurnSequence === session.dispatchSequence)) + (dispatchAdmissionAllowsCancellation() || dispatchAdmissionExpired) + : compactionOwnsTurn() || + (ownsRequestedTurn() && + (dispatchAdmissionAllowsCancellation() || dispatchAdmissionExpired))) 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 eb7bc9b7251..2895bcf557c 100644 --- a/src/main/claude/claude-structured-session-acquisition.ts +++ b/src/main/claude/claude-structured-session-acquisition.ts @@ -138,7 +138,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..661138daf96 --- /dev/null +++ b/src/main/claude/claude-turn-ownership.test.ts @@ -0,0 +1,374 @@ +// 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 { StructuredSessionCompaction } from '../native-chat/agent-session-wire/structured-session-compaction' +import { + CLAUDE_DISPATCH_ADMISSION_TIMEOUT_MS, + cancelClaudeStructuredTurn +} from './claude-structured-prompt-ownership' +import { sessionFor } from './claude-structured-dispatch-test-support' +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('refuses a stale id after the owned turn settles', 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() + + await expect( + adapter.cancelTurn({ sessionId: 'session-1', turnId: 'echo-turn', fence: 7 }) + ).resolves.toEqual({ cancelled: false }) + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(false) + }) + + it('keeps the prior dispatch fence after an unknown later send', async () => { + vi.useFakeTimers() + try { + 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') + const sendFirst = connection.send + connection.send = async (message) => { + if (connection.sent.length > 0) { + throw new Error('input pump stopped') + } + await sendFirst(message) + } + + await expect( + adapter.dispatch({ + sessionId: 'session-1', + clientMessageId: 'client-2', + body: USER_MESSAGE, + fence: 7 + }) + ).resolves.toMatchObject({ state: 'unknown' }) + + const cancellation = adapter.cancelTurn({ + sessionId: 'session-1', + turnId: 'echo-turn', + fence: 7 + }) + await vi.advanceTimersByTimeAsync(CLAUDE_DISPATCH_ADMISSION_TIMEOUT_MS - 1) + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(false) + await vi.advanceTimersByTimeAsync(1) + await expect(cancellation).resolves.toEqual({ cancelled: true }) + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(true) + } finally { + vi.useRealTimers() + } + }) + + it('lets a queued-cancel provider release an unresolved ordinary Stop', async () => { + const claude = fakeClaude({ + replayUuid: 'echo-turn', + capabilities: ['interrupt_cancel_queued_v1'] + }) + 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, + dispatchStatus: { state: 'unknown', recovered: false } + }) + ).resolves.toEqual({ cancelled: true }) + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(true) + }) + + it('lets ordinary Stop proceed after the unresolved delivery fence expires', async () => { + vi.useFakeTimers() + try { + 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') + + const cancellation = adapter.cancelTurn({ + sessionId: 'session-1', + turnId: 'echo-turn', + fence: 7, + dispatchStatus: { state: 'unknown', recovered: false } + }) + await vi.advanceTimersByTimeAsync(CLAUDE_DISPATCH_ADMISSION_TIMEOUT_MS - 1) + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(false) + + await vi.advanceTimersByTimeAsync(1) + await expect(cancellation).resolves.toEqual({ cancelled: true }) + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(true) + } finally { + vi.useRealTimers() + } + }) + + it('releases ordinary Stop as soon as a retired delivery fence settles', async () => { + vi.useFakeTimers() + try { + const session = sessionFor() + session.dispatchSequence = 1 + session.translator = { + handle: vi.fn(), + journalPrompts: { cancel: vi.fn(), resolve: vi.fn() }, + currentTurnId: 'turn-1', + flush: vi.fn(), + pendingStreamedBlocks: 0, + dispose: vi.fn() + } + session.retiredDispatchWaiters = [ + { + acceptsResult: false, + clientMessageId: 'client-2', + sentUuid: 'uncertain', + dispatchSequence: 1, + replayContentKey: 'ship-it', + resolve: vi.fn(), + retired: true + } + ] + const interrupt = vi.fn().mockResolvedValue(undefined) + session.connection.interrupt = interrupt + const cancellation = cancelClaudeStructuredTurn({ + request: { sessionId: 'session-1', turnId: 'turn-1', fence: 1 }, + sessions: new Map([['session-1', session]]), + compactions: new StructuredSessionCompaction(), + admitPromptCancellation: () => true + }) + await vi.advanceTimersByTimeAsync(100) + session.retiredDispatchWaiters = [] + await vi.advanceTimersByTimeAsync(100) + const settledBeforeDeadline = interrupt.mock.calls.length > 0 + if (!settledBeforeDeadline) { + await vi.advanceTimersByTimeAsync(CLAUDE_DISPATCH_ADMISSION_TIMEOUT_MS) + await cancellation + } + expect(settledBeforeDeadline).toBe(true) + await expect(cancellation).resolves.toEqual({ cancelled: true }) + } finally { + vi.useRealTimers() + } + }) + + it('does not wait when the dispatch admission is already current', async () => { + vi.useFakeTimers() + try { + 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) + } finally { + vi.useRealTimers() + } + }) + + it('honors an unresolved journal submission before the first in-memory dispatch', async () => { + vi.useFakeTimers() + try { + const claude = fakeClaude({ replayUuid: null }) + const { adapter, bodies, connection } = await acquiredWithJournal(claude) + providerOutput(connection, 'provider-turn') + expect(runningTurnId(bodies)).toBe('provider-turn') + + const cancellation = adapter.cancelTurn({ + sessionId: 'session-1', + turnId: 'provider-turn', + fence: 7, + dispatchStatus: { state: 'pending', recovered: false } + }) + await vi.advanceTimersByTimeAsync(CLAUDE_DISPATCH_ADMISSION_TIMEOUT_MS - 1) + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(false) + await vi.advanceTimersByTimeAsync(1) + await expect(cancellation).resolves.toEqual({ cancelled: true }) + expect(connection.calls.some((call) => call.subtype === 'interrupt')).toBe(true) + } finally { + vi.useRealTimers() + } + }) + + 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) + }) +}) diff --git a/src/main/native-chat/agent-session-journal/journal-dispatch-observation.test.ts b/src/main/native-chat/agent-session-journal/journal-dispatch-observation.test.ts new file mode 100644 index 00000000000..f13f2a423e7 --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-dispatch-observation.test.ts @@ -0,0 +1,65 @@ +import { describe, expect, it } from 'vitest' +import type { AgentJournalSubmission } from '../../../shared/agent-session-journal-types' +import { latestJournalDispatchObservation } from './journal-dispatch-observation' + +describe('latestJournalDispatchObservation', () => { + it('uses the newest submission in the requested fence', () => { + const submissions = [ + { + clientMessageId: 'unknown-7', + fence: 7, + payloadFingerprint: 'unknown-7', + dispatchState: 'unknown' as const, + recovered: true as const, + providerItemId: null, + reason: null, + resolvedAt: null, + submittedAt: 1 + }, + { + clientMessageId: 'pending-8', + fence: 8, + payloadFingerprint: 'pending-8', + dispatchState: 'pending' as const, + providerItemId: null, + reason: null, + resolvedAt: null, + submittedAt: 2 + }, + { + clientMessageId: 'pending-7', + fence: 7, + payloadFingerprint: 'pending-7', + dispatchState: 'pending' as const, + providerItemId: null, + reason: null, + resolvedAt: null, + submittedAt: 2 + }, + { + clientMessageId: 'accepted-7', + fence: 7, + payloadFingerprint: 'accepted-7', + dispatchState: 'accepted' as const, + providerItemId: 'item-7', + reason: null, + resolvedAt: 3, + submittedAt: 4 + } + ] satisfies AgentJournalSubmission[] + const journal = { submissions: () => submissions } + + expect(latestJournalDispatchObservation(journal, 7)).toEqual({ + state: 'accepted', + recovered: false + }) + expect(latestJournalDispatchObservation(journal, 8)).toEqual({ + state: 'pending', + recovered: false + }) + }) + + it('returns no observation when the fence has no submission', () => { + expect(latestJournalDispatchObservation({ submissions: () => [] }, 7)).toBeNull() + }) +}) diff --git a/src/main/native-chat/agent-session-journal/journal-dispatch-observation.ts b/src/main/native-chat/agent-session-journal/journal-dispatch-observation.ts new file mode 100644 index 00000000000..d5248b4163a --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-dispatch-observation.ts @@ -0,0 +1,29 @@ +import type { + AgentJournalDispatchState, + AgentJournalSubmission +} from '../../../shared/agent-session-journal-types' + +export type AgentJournalDispatchObservation = { + state: AgentJournalDispatchState + recovered: boolean +} + +/** Returns the latest write-ahead submission for the execution fence. */ +export function latestJournalDispatchObservation( + journal: { + submissions: () => readonly AgentJournalSubmission[] + }, + fence: number +): AgentJournalDispatchObservation | null { + const latest = journal + .submissions() + .reduce( + (current, submission) => + submission.fence === fence && + (current === null || submission.submittedAt >= current.submittedAt) + ? submission + : current, + null + ) + return latest ? { state: latest.dispatchState, recovered: latest.recovered === true } : null +} 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 e5daa9991d9..729b0e212f1 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 @@ -14,6 +14,7 @@ import type { AgentJournalItemIdentity, AgentJournalItemBody, AgentJournalMessageItem, + AgentJournalDispatchState, AgentSessionJournalIdentity } from '../../../shared/agent-session-journal-types' import type { AgentSessionProviderHandleLink } from '../../../shared/agent-session-provider-handle' @@ -204,6 +205,8 @@ export type StructuredAgentSessionAdapter = { turnId: string fence: number prompt?: { itemId: string } + /** Latest journal submission for this fence, when the host has one. */ + dispatchStatus?: { state: AgentJournalDispatchState; recovered: boolean } | null }): Promise<{ cancelled: boolean }> stopBackgroundTasks?(input: { sessionId: string diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-handoff.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-handoff.ts index 9a3266541be..4fad01d32a4 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-handoff.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-handoff.ts @@ -16,6 +16,7 @@ import type { AgentSessionSubscribers } from './structured-agent-session-subscri import { StructuredTuiTranscriptCatchup } from './structured-tui-transcript-catchup' import { adapterSupportsCreateIfDeclared } from './structured-agent-session-provider-support' import { retryLoadedStructuredAgentSessionSettlement } from './structured-agent-session-settlement-retry' +import { latestJournalDispatchObservation } from '../agent-session-journal/journal-dispatch-observation' type HostHandoffAccess = { session: (sessionId: string) => StructuredAgentSessionHostSession @@ -101,8 +102,18 @@ export function createStructuredAgentSessionHostHandoff( }, acknowledgeNativeRelease: (sessionId) => deps.adapter.acknowledgeSessionRelease?.(sessionId), acquireNative: (input) => acquireNativeHandoffOwner(deps, host, input), - acquireNativeStop: async (sessionId, turnId, fence) => - (await deps.adapter.cancelTurn({ sessionId, turnId, fence })).cancelled, + acquireNativeStop: async (sessionId, turnId, fence) => { + const session = host.session(sessionId) + const dispatchStatus = latestJournalDispatchObservation(session.journal, fence) + return ( + await deps.adapter.cancelTurn({ + sessionId, + turnId, + fence, + ...(dispatchStatus ? { dispatchStatus } : {}) + }) + ).cancelled + }, importTuiHistory: (input) => importTuiHistory(deps, host, input), retryPendingSettlement: (sessionId) => retryLoadedStructuredAgentSessionSettlement({ diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts index c666047e1e7..fb254bbf1e5 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts @@ -17,6 +17,7 @@ import type { } from '../../../shared/agent-session-wire' import { DISPATCH_DOUBT_PERSISTENCE_FAILED } from '../agent-session-journal/journal-dispatch-doubt-reasons' import type { AgentSessionJournal } from '../agent-session-journal/journal-store' +import { latestJournalDispatchObservation } from '../agent-session-journal/journal-dispatch-observation' import type { AgentSessionDispatchOutcome, StructuredAgentSessionAdapter @@ -201,6 +202,7 @@ export async function performCancel( let cancelled = false let note = 'Cancellation requested.' try { + const dispatchStatus = latestJournalDispatchObservation(ctx.journal, ctx.fence) cancelled = input.scope ? ( await ctx.adapter.stopBackgroundTasks?.({ @@ -214,6 +216,7 @@ export async function performCancel( sessionId: ctx.sessionId, turnId: input.turnId, fence: ctx.fence, + ...(dispatchStatus ? { dispatchStatus } : {}), ...(input.prompt ? { prompt: { itemId: input.prompt.itemId } } : {}) }) ).cancelled diff --git a/src/main/runtime/claude-structured-session-integration.test.ts b/src/main/runtime/claude-structured-session-integration.test.ts index d464d87e8f7..f6b540bc5af 100644 --- a/src/main/runtime/claude-structured-session-integration.test.ts +++ b/src/main/runtime/claude-structured-session-integration.test.ts @@ -600,6 +600,18 @@ describe('a structured Claude session over agentSession.*', () => { `claude:${PROVIDER_SESSION}:assistant-leaf` ) + // A background task can wake Claude after the preceding dispatch settled. + // This assistant frame opens the provider-owned turn without an Orca send + // echo; Stop must target that frame's id rather than the settled user row. + claude.live().handlers.onMessage?.({ + type: 'assistant', + session_id: PROVIDER_SESSION, + uuid: 'provider-opened-assistant', + parent_tool_use_id: null, + message: { role: 'assistant', content: [{ type: 'text', text: 'Background task update.' }] } + }) + await getStructuredAgentSessionHost()?.flushStreamedEvents(SESSION) + claude.live().handlers.onMessage?.({ type: 'system', subtype: 'background_tasks_changed', @@ -672,10 +684,14 @@ describe('a structured Claude session over agentSession.*', () => { await expect( ok('agentSession.cancel', { - envelope: envelope('agentSession.cancel', { turnId: 'user-1' }, created.fence), - turnId: 'user-1' + envelope: envelope( + 'agentSession.cancel', + { turnId: 'provider-opened-assistant' }, + created.fence + ), + turnId: 'provider-opened-assistant' }) - ).resolves.toMatchObject({ turnId: 'user-1', cancelled: true }) + ).resolves.toMatchObject({ turnId: 'provider-opened-assistant', cancelled: true }) expect(claude.live().calls.at(-1)).toMatchObject({ subtype: 'interrupt' }) const host = getStructuredAgentSessionHost() as unknown as { @@ -700,13 +716,13 @@ describe('a structured Claude session over agentSession.*', () => { }) expect(claude.live().launch.options).toMatchObject({ resume: PROVIDER_SESSION, - resumeSessionAt: 'assistant-leaf' + resumeSessionAt: 'provider-opened-assistant' }) expect(host.deps.store.getRecord(SESSION).providerHandleChain.at(-1)).toMatchObject({ handle: { provider: 'claude', sessionId: PROVIDER_SESSION, - leafUuid: 'assistant-leaf' + leafUuid: 'provider-opened-assistant' }, origin: 'resumed' }) From 3520e8eb41e9206e312bf66dee3a353a56a4a437 Mon Sep 17 00:00:00 2001 From: BAEK'space <112856532+100space@users.noreply.github.com> Date: Wed, 16 Sep 2026 09:32:54 +0900 Subject: [PATCH 2/3] fix: highlight bash fences in Markdown source mode (#20592) * fix: highlight bash fences in Markdown source mode * refactor: trim shell fence alias registration Drop the speculative exports and document the alias-resolution rationale in one WHY comment; the idempotency guard stays. --------- Co-authored-by: Jinjing <6427696+AmethystLiang@users.noreply.github.com> --- .../register-shell-markdown-aliases.test.ts | 32 ++++++++ .../register-shell-markdown-aliases.ts | 24 ++++++ src/renderer/src/lib/monaco-setup.ts | 2 + .../markdown-source-bash-highlighting.spec.ts | 78 +++++++++++++++++++ 4 files changed, 136 insertions(+) create mode 100644 src/renderer/src/lib/monaco-languages/register-shell-markdown-aliases.test.ts create mode 100644 src/renderer/src/lib/monaco-languages/register-shell-markdown-aliases.ts create mode 100644 tests/e2e/markdown-source-bash-highlighting.spec.ts diff --git a/src/renderer/src/lib/monaco-languages/register-shell-markdown-aliases.test.ts b/src/renderer/src/lib/monaco-languages/register-shell-markdown-aliases.test.ts new file mode 100644 index 00000000000..d977b4322e5 --- /dev/null +++ b/src/renderer/src/lib/monaco-languages/register-shell-markdown-aliases.test.ts @@ -0,0 +1,32 @@ +import { describe, expect, it, vi } from 'vitest' +import { registerShellMarkdownAliases } from './register-shell-markdown-aliases' + +function createMonacoMock(aliases: string[] = ['Shell', 'sh']) { + return { + languages: { + getLanguages: vi.fn(() => [{ id: 'shell', aliases }]), + register: vi.fn() + } + } +} + +describe('registerShellMarkdownAliases', () => { + it('registers bash alongside the built-in shell aliases', () => { + const monaco = createMonacoMock() + + registerShellMarkdownAliases(monaco) + + expect(monaco.languages.register).toHaveBeenCalledWith({ + id: 'shell', + aliases: ['Shell', 'sh', 'bash'] + }) + }) + + it('does not register the alias again when Monaco already exposes it', () => { + const monaco = createMonacoMock(['Shell', 'sh', 'Bash']) + + registerShellMarkdownAliases(monaco) + + expect(monaco.languages.register).not.toHaveBeenCalled() + }) +}) diff --git a/src/renderer/src/lib/monaco-languages/register-shell-markdown-aliases.ts b/src/renderer/src/lib/monaco-languages/register-shell-markdown-aliases.ts new file mode 100644 index 00000000000..1cad43f1dfd --- /dev/null +++ b/src/renderer/src/lib/monaco-languages/register-shell-markdown-aliases.ts @@ -0,0 +1,24 @@ +import type * as Monaco from 'monaco-editor' + +type MonacoModule = typeof Monaco + +// Why: Monaco resolves Markdown fences by alias (never extension) and its shell +// language declares `bash` only as an extension, so ```bash rendered plain while +// ```sh highlighted. Re-registering id 'shell' merges the alias and keeps the +// built-in tokenizer; `Shell` stays first because Monaco uses the first alias as +// the language's display name. +export function registerShellMarkdownAliases(monaco: { + languages: Pick +}): void { + const bashAlreadyRegistered = monaco.languages + .getLanguages() + .some( + ({ id, aliases }) => + id === 'shell' && aliases?.some((alias) => alias.toLowerCase() === 'bash') + ) + if (bashAlreadyRegistered) { + return + } + + monaco.languages.register({ id: 'shell', aliases: ['Shell', 'sh', 'bash'] }) +} diff --git a/src/renderer/src/lib/monaco-setup.ts b/src/renderer/src/lib/monaco-setup.ts index 523ecf8e77b..57a42fc230f 100644 --- a/src/renderer/src/lib/monaco-setup.ts +++ b/src/renderer/src/lib/monaco-setup.ts @@ -10,6 +10,7 @@ import tsWorker from 'monaco-editor/esm/vs/language/typescript/ts.worker?worker' import { registerAstroLanguage } from './monaco-languages/register-astro' import { registerJsonlLanguage } from './monaco-languages/register-jsonl' import { registerNimLanguage } from './monaco-languages/register-nim' +import { registerShellMarkdownAliases } from './monaco-languages/register-shell-markdown-aliases' import { registerSvelteLanguage } from './monaco-languages/register-svelte' import { registerVueLanguage } from './monaco-languages/register-vue' import { installMonacoDelayerCancellationGuard } from './monaco-delayer-cancellation-guard' @@ -79,6 +80,7 @@ registerSvelteLanguage(monaco) registerAstroLanguage(monaco) registerNimLanguage(monaco) registerJsonlLanguage(monaco) +registerShellMarkdownAliases(monaco) installMonacoDelayerCancellationGuard() installMonacoDiffEditorDisposalGuard(monaco) installMonacoPeekReferencesPreviewOptions() diff --git a/tests/e2e/markdown-source-bash-highlighting.spec.ts b/tests/e2e/markdown-source-bash-highlighting.spec.ts new file mode 100644 index 00000000000..33b0e0ed82e --- /dev/null +++ b/tests/e2e/markdown-source-bash-highlighting.spec.ts @@ -0,0 +1,78 @@ +import type { Locator, Page } from '@stablyai/playwright-test' +import { test, expect } from './helpers/orca-app' +import { waitForActiveWorktree, waitForSessionReady } from './helpers/store' +import { + cleanupMarkdownFixture, + createMarkdownFixture, + getActiveWorktreeContext, + openMarkdownFixture, + waitForRichMarkdownEditor +} from './helpers/markdown-editor-fixture' + +const MARKDOWN = `\`\`\`bash +printf '%s\\n' "build complete" # bash-highlight-marker +\`\`\` + +\`\`\`sh +printf '%s\\n' "build complete" # shell-control-marker +\`\`\` +` + +async function switchToSourceMode(page: Page): Promise { + await page.evaluate(() => { + const store = window.__store + if (!store) { + throw new Error('window.__store is not available') + } + const state = store.getState() + if (!state.activeFileId) { + throw new Error('No active editor file') + } + state.setMarkdownViewMode(state.activeFileId, 'source') + }) +} + +async function distinctLeafTokenColors(line: Locator): Promise { + return line.locator('span').evaluateAll((spans) => { + const colors = spans + .filter((span) => span.childElementCount === 0 && span.textContent?.trim()) + .map((span) => window.getComputedStyle(span).color) + return new Set(colors).size + }) +} + +test('highlights bash and sh fences in Markdown Source mode', async ({ orcaPage }, testInfo) => { + await waitForSessionReady(orcaPage) + await waitForActiveWorktree(orcaPage) + + const context = await getActiveWorktreeContext(orcaPage) + let filePath: string | null = null + + try { + filePath = await createMarkdownFixture( + context, + '.orca-e2e-markdown-source-highlighting', + 'bash-and-sh', + testInfo.workerIndex, + MARKDOWN + ) + await openMarkdownFixture(orcaPage, context, filePath) + await waitForRichMarkdownEditor(orcaPage) + await switchToSourceMode(orcaPage) + + const monaco = orcaPage.locator('.monaco-editor').first() + await expect(monaco).toBeVisible({ timeout: 25_000 }) + + for (const marker of ['bash-highlight-marker', 'shell-control-marker']) { + const line = monaco.locator('.view-line').filter({ hasText: marker }) + await expect(line).toHaveCount(1) + await expect + .poll(() => distinctLeafTokenColors(line), { + message: `${marker} should render with distinct shell token colors` + }) + .toBeGreaterThan(1) + } + } finally { + await cleanupMarkdownFixture(filePath) + } +}) From 7ec2986fd11e2de3964546657c0069a00f0676d6 Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Tue, 15 Sep 2026 17:36:08 -0700 Subject: [PATCH 3/3] fix(lint): merge the duplicate agent-status contract type imports (#20907) main's tip fails audit:code-quality:native on import(no-duplicates), which reds the static analysis and verify jobs of every open PR via the merge ref. --- src/shared/agent-status-store-snapshot-budget.ts | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/src/shared/agent-status-store-snapshot-budget.ts b/src/shared/agent-status-store-snapshot-budget.ts index 6ae4dc8f040..dfb4665a01a 100644 --- a/src/shared/agent-status-store-snapshot-budget.ts +++ b/src/shared/agent-status-store-snapshot-budget.ts @@ -4,10 +4,12 @@ import { } from './agent-status-store-contract' import type { AgentChildWorkAliasRecord } from './agent-status-child-work-alias' import type { AgentChildWorkRecord } from './agent-status-child-work' -import type { AgentStatusFactRecord } from './agent-status-store-contract' +import type { + AgentStatusFactRecord, + AgentStatusTombstoneRecord +} from './agent-status-store-contract' import type { AgentStatusParentRecord } from './agent-status-store-parent' import type { AgentStatusStoreState } from './agent-status-store-state' -import type { AgentStatusTombstoneRecord } from './agent-status-store-contract' import { getUtf8ByteLength } from './utf8-byte-limits' type AgentStatusSnapshotRecord =