mirror of
https://github.com/stablyai/orca.git
synced 2026-10-06 16:02:25 +00:00
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.
This commit is contained in:
@@ -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<string, unknown>,
|
||||
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<string, unknown>,
|
||||
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 })
|
||||
}
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
|
||||
@@ -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 () => {
|
||||
|
||||
@@ -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<string, unknown>,
|
||||
@@ -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 }
|
||||
|
||||
@@ -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<ClaudeJournalPrompts, 'cancel' | 'resolve'>
|
||||
/** 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<string, ClaudeToolUse>()
|
||||
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
|
||||
|
||||
@@ -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<StructuredAgentSessionAdapter['cancelTurn']>[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(
|
||||
|
||||
@@ -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))
|
||||
})
|
||||
|
||||
@@ -232,7 +232,7 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
|
||||
journalItemId,
|
||||
promptKey,
|
||||
questionId,
|
||||
session.activeTurnId ?? null
|
||||
session.translator?.currentTurnId ?? null
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -147,16 +147,12 @@ export type ClaudeSession = {
|
||||
restoreSkippedOptions: Set<string>
|
||||
/** 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. */
|
||||
|
||||
@@ -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<void> {
|
||||
return new Promise((resolve) => setImmediate(resolve))
|
||||
}
|
||||
|
||||
@@ -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<string, AgentJournalItemBody>
|
||||
} {
|
||||
const bodies = new Map<string, AgentJournalItemBody>()
|
||||
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, AgentJournalItemBody>): 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<typeof fakeClaude>): Promise<{
|
||||
adapter: ReturnType<typeof adapterFor>
|
||||
bodies: Map<string, AgentJournalItemBody>
|
||||
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)
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user