fix(native-chat): retain Claude task rows through journal pressure

This commit is contained in:
Brennan Benson
2026-09-17 02:14:18 -07:00
parent 11f05cdddb
commit e84c650659
21 changed files with 1403 additions and 228 deletions
@@ -0,0 +1,294 @@
import { describe, expect, it, vi } from 'vitest'
import type {
StructuredAgentSessionEventSink,
StructuredAgentSessionLifecycleJournal
} from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import type { ClaudeBackgroundTaskRow } from './claude-background-task-row-lifecycle'
import {
ClaudeBackgroundTaskIdentityResolver,
writeClaudeBackgroundTaskRow
} from './claude-background-task-row-journal'
import { ClaudeBackgroundTaskRows } from './claude-background-task-rows'
import { START_BASH } from './claude-background-task-row-test-support'
function journal(
epoch: () => string,
visitItems = vi.fn()
): StructuredAgentSessionLifecycleJournal {
return {
get epoch() {
return epoch()
},
visitItems
}
}
function row(): ClaudeBackgroundTaskRow {
return {
block: {
type: 'background-task',
taskId: 'task-1',
kind: 'command',
label: 'Build',
state: 'blocked',
error: 'exit 1'
},
lastSerialized: null,
toolUseId: 'tool-1',
terminalNotificationReceived: true,
generation: 1
}
}
describe('Claude background task row journal', () => {
it('resolves one immutable run once per journal epoch', () => {
let epoch = 'epoch-1'
const visits = vi.fn()
const firstJournal = journal(() => epoch, visits)
const resolver = new ClaudeBackgroundTaskIdentityResolver()
for (let revision = 0; revision < 100; revision += 1) {
resolver.resolve(firstJournal, 'task-1', 'tool-1')
}
expect(visits).toHaveBeenCalledOnce()
resolver.resolve(firstJournal, 'task-1', 'tool-2')
expect(visits).toHaveBeenCalledTimes(2)
epoch = 'epoch-2'
resolver.resolve(firstJournal, 'task-1', 'tool-1')
expect(visits).toHaveBeenCalledTimes(3)
resolver.resolve(
journal(() => 'epoch-2', visits),
'task-1',
'tool-1'
)
expect(visits).toHaveBeenCalledTimes(4)
})
it('bounds resolved run identities with LRU eviction', () => {
const visits = vi.fn()
const boundJournal = journal(() => 'epoch-1', visits)
const resolver = new ClaudeBackgroundTaskIdentityResolver()
for (let index = 0; index < 513; index += 1) {
resolver.resolve(boundJournal, `task-${index}`, `tool-${index}`)
}
resolver.resolve(boundJournal, 'task-0', 'tool-0')
expect(visits).toHaveBeenCalledTimes(514)
})
it('records a revision only after its append and publication are admitted', () => {
const task = row()
const appendAndPublish = vi
.fn()
.mockReturnValueOnce({ accepted: false, reason: 'backpressure' })
.mockReturnValueOnce({ accepted: true })
const sink = {
appendItem: vi.fn(),
appendTombstone: vi.fn(),
publish: vi.fn(),
tryAppendResolvedItemAndPublish: appendAndPublish
}
const resolver = new ClaudeBackgroundTaskIdentityResolver()
expect(writeClaudeBackgroundTaskRow(sink, resolver, 'task-1', task)).toEqual({
accepted: false,
reason: 'backpressure'
})
expect(task.lastSerialized).toBeNull()
expect(writeClaudeBackgroundTaskRow(sink, resolver, 'task-1', task)).toEqual({
accepted: true
})
expect(task.lastSerialized).not.toBeNull()
expect(appendAndPublish).toHaveBeenCalledTimes(2)
})
it('keeps fallback append coalescing separate from its ordered publication', () => {
const calls: { operation: 'append' | 'publish'; coalescingKey?: string }[] = []
const task = row()
const sink: StructuredAgentSessionEventSink = {
appendItem: vi.fn(),
appendTombstone: vi.fn(),
publish: vi.fn(),
tryAppendResolvedItem: vi.fn((_identity, _body, _resolve, options) => {
calls.push({
operation: 'append',
...(options?.coalescingKey ? { coalescingKey: options.coalescingKey } : {})
})
return { accepted: true as const }
}),
tryPublish: vi.fn((options) => {
calls.push({
operation: 'publish',
...(options?.coalescingKey ? { coalescingKey: options.coalescingKey } : {})
})
return { accepted: true as const }
})
}
expect(
writeClaudeBackgroundTaskRow(sink, new ClaudeBackgroundTaskIdentityResolver(), 'task-1', task)
).toEqual({ accepted: true })
expect(calls).toEqual([
{
operation: 'append',
coalescingKey: JSON.stringify(['claude-background-task', 'task-1', 'tool-1'])
},
{ operation: 'publish' }
])
expect(task.lastSerialized).not.toBeNull()
})
it('uses reserved lifecycle capacity when provider exit settles a live row', () => {
const appendAndPublish = vi.fn<
NonNullable<StructuredAgentSessionEventSink['tryAppendResolvedItemAndPublish']>
>(() => ({ accepted: true as const }))
const rows = new ClaudeBackgroundTaskRows({
sink: {
appendItem: vi.fn(),
appendTombstone: vi.fn(),
publish: vi.fn(),
tryAppendResolvedItemAndPublish: appendAndPublish
},
isForwardedParentTool: () => true
})
rows.observe(START_BASH)
rows.settleSession()
expect(appendAndPublish).toHaveBeenCalledTimes(2)
expect(appendAndPublish.mock.calls[1]?.[3]).toMatchObject({ lifecycle: true })
})
it('promotes a backpressured terminal row to lifecycle capacity on provider exit', () => {
const appendAndPublish = vi
.fn<NonNullable<StructuredAgentSessionEventSink['tryAppendResolvedItemAndPublish']>>()
.mockReturnValueOnce({ accepted: false, reason: 'backpressure' })
.mockReturnValueOnce({ accepted: true })
const rows = new ClaudeBackgroundTaskRows({
sink: {
appendItem: vi.fn(),
appendTombstone: vi.fn(),
publish: vi.fn(),
tryAppendResolvedItemAndPublish: appendAndPublish
},
isForwardedParentTool: () => true
})
rows.observe({
type: 'system',
subtype: 'task_notification',
task_id: 'terminal-at-exit',
status: 'failed',
summary: 'failed'
})
rows.settleSession()
expect(appendAndPublish).toHaveBeenCalledTimes(2)
expect(appendAndPublish.mock.calls[1]?.[3]).toMatchObject({ lifecycle: true })
})
it('retries a refused row before later provider messages are read', () => {
const appendAndPublish = vi
.fn()
.mockReturnValueOnce({ accepted: false, reason: 'backpressure' })
.mockReturnValueOnce({ accepted: true })
const rows = new ClaudeBackgroundTaskRows({
sink: {
appendItem: vi.fn(),
appendTombstone: vi.fn(),
publish: vi.fn(),
tryAppendResolvedItemAndPublish: appendAndPublish
},
isForwardedParentTool: () => true
})
expect(rows.observe(START_BASH)).toBe(true)
expect(rows.retryPendingWrites()).toEqual({ accepted: true })
expect(appendAndPublish).toHaveBeenCalledTimes(2)
expect(rows.retryPendingWrites()).toEqual({ accepted: true })
expect(appendAndPublish).toHaveBeenCalledTimes(2)
})
it('abandons a permanently failed retry and reports recovery instead of latching', () => {
const onPersistenceFailure = vi.fn()
const appendAndPublish = vi
.fn()
.mockReturnValueOnce({ accepted: false, reason: 'backpressure' })
.mockReturnValueOnce({ accepted: false, reason: 'failed' })
const rows = new ClaudeBackgroundTaskRows({
sink: {
appendItem: vi.fn(),
appendTombstone: vi.fn(),
publish: vi.fn(),
tryAppendResolvedItemAndPublish: appendAndPublish
},
isForwardedParentTool: () => true,
onPersistenceFailure
})
rows.observe(START_BASH)
expect(rows.retryPendingWrites()).toEqual({ accepted: false, reason: 'failed' })
expect(onPersistenceFailure).toHaveBeenCalledOnce()
expect(rows.retryPendingWrites()).toEqual({ accepted: true })
})
it('hands retry-capacity exhaustion to session recovery without throwing', () => {
const onPersistenceFailure = vi.fn()
const rows = new ClaudeBackgroundTaskRows({
sink: {
appendItem: vi.fn(),
appendTombstone: vi.fn(),
publish: vi.fn(),
tryAppendResolvedItemAndPublish: vi.fn(() => ({
accepted: false as const,
reason: 'backpressure' as const
}))
},
isForwardedParentTool: () => true,
onPersistenceFailure
})
for (let index = 0; index < 512; index += 1) {
rows.observe({
type: 'system',
subtype: 'task_notification',
task_id: `task-${index}`,
status: 'failed',
summary: 'failed'
})
}
expect(() => {
rows.observe({
type: 'system',
subtype: 'task_notification',
task_id: 'task-overflow',
status: 'failed',
summary: 'failed'
})
}).not.toThrow()
expect(onPersistenceFailure).toHaveBeenCalledWith(
expect.objectContaining({
message: 'claude background task journal retry capacity exhausted'
})
)
})
it.each(['closed', 'failed'] as const)('surfaces a %s sink refusal', (reason) => {
const task = row()
const sink = {
appendItem: vi.fn(),
appendTombstone: vi.fn(),
publish: vi.fn(),
tryAppendResolvedItemAndPublish: vi.fn(() => ({ accepted: false as const, reason }))
}
expect(
writeClaudeBackgroundTaskRow(sink, new ClaudeBackgroundTaskIdentityResolver(), 'task-1', task)
).toEqual({ accepted: false, reason })
expect(task.lastSerialized).toBeNull()
})
})
@@ -9,11 +9,15 @@ import {
} from '../../shared/native-chat-types'
import type {
StructuredAgentSessionEventSink,
StructuredAgentSessionLifecycleJournal
StructuredAgentSessionLifecycleJournal,
StructuredAgentSessionSinkAdmission
} from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import type { ClaudeBackgroundTaskRow } from './claude-background-task-row-lifecycle'
import { parseAgentJournalItemKey } from '../../shared/agent-session-journal-item-key'
const MAX_RESOLVED_TASK_IDENTITIES = 512
const ADMITTED: StructuredAgentSessionSinkAdmission = { accepted: true }
/** Durable identity for one RUN of a task.
*
* A provider may reuse a task id for a distinct later invocation, and a row
@@ -78,6 +82,47 @@ export function resolveClaudeBackgroundTaskIdentity(
)
}
/** Resolves each immutable provider run once per bound journal epoch. */
export class ClaudeBackgroundTaskIdentityResolver {
private journal: StructuredAgentSessionLifecycleJournal | null = null
private epoch: string | null = null
private readonly identities = new Map<string, AgentJournalItemIdentity>()
resolve = (
journal: StructuredAgentSessionLifecycleJournal,
id: string,
toolUseId: string | undefined
): AgentJournalItemIdentity => {
if (this.journal !== journal || this.epoch !== journal.epoch) {
this.journal = journal
this.epoch = journal.epoch
this.identities.clear()
}
const key = JSON.stringify([id, toolUseId ?? null])
const cached = this.identities.get(key)
if (cached) {
this.identities.delete(key)
this.identities.set(key, cached)
return cached
}
const identity = resolveClaudeBackgroundTaskIdentity(journal, id, toolUseId)
this.identities.set(key, identity)
if (this.identities.size > MAX_RESOLVED_TASK_IDENTITIES) {
const oldest = this.identities.keys().next()
if (!oldest.done) {
this.identities.delete(oldest.value)
}
}
return identity
}
clear(): void {
this.journal = null
this.epoch = null
this.identities.clear()
}
}
function persistedTaskGeneration(clientMessageId: string, taskId: string): number | null {
const base = `claude-background-task:${taskId}`
if (clientMessageId === base) {
@@ -93,42 +138,58 @@ function persistedTaskGeneration(clientMessageId: string, taskId: string): numbe
export function writeClaudeBackgroundTaskRow(
sink: StructuredAgentSessionEventSink,
identities: ClaudeBackgroundTaskIdentityResolver,
id: string,
row: ClaudeBackgroundTaskRow,
/** Runs only when a row is really appended, so a duplicate delivery that
* changes nothing never opens a turn. */
beforeAppend?: () => void
): void {
/** Runs before admission to preserve turn-before-row ordering; duplicate
* delivery skips it, and a retry reuses the turn the first attempt opened. */
beforeAppend?: () => void,
lifecycle = false
): StructuredAgentSessionSinkAdmission {
const body = claudeBackgroundTaskBody(row.block)
const serialized = JSON.stringify(body)
if (serialized === row.lastSerialized) {
return
return ADMITTED
}
row.lastSerialized = serialized
beforeAppend?.()
const identity = claudeBackgroundTaskIdentity(id, row.generation)
// Generation is translator-local and resets when a provider stream is
// recreated. Keep unresolved writes from distinct provider runs queued side
// by using the provider's parent tool identity as the coalescing discriminator.
const coalescingKey = JSON.stringify(['claude-background-task', id, row.toolUseId ?? null])
const resolveIdentity = sink.tryAppendResolvedItem
if (resolveIdentity) {
const appendOptions = { coalescingKey, ...(lifecycle ? { lifecycle: true } : {}) }
const publishOptions = lifecycle ? { lifecycle: true } : {}
const resolveIdentity = (journal: StructuredAgentSessionLifecycleJournal) =>
identities.resolve(journal, id, row.toolUseId)
const appendAndPublish = sink.tryAppendResolvedItemAndPublish
let admission: StructuredAgentSessionSinkAdmission
if (appendAndPublish) {
// Reserve enough space for any safe generation suffix; the actual identity
// is selected once the deferred sink is bound to the durable journal.
const identitySizeBound = claudeBackgroundTaskIdentity(id, Number.MAX_SAFE_INTEGER)
resolveIdentity(
identitySizeBound,
body,
(journal) => resolveClaudeBackgroundTaskIdentity(journal, id, row.toolUseId),
{
coalescingKey
}
)
sink.publish()
return
admission = appendAndPublish(identitySizeBound, body, resolveIdentity, appendOptions)
} else if (sink.tryAppendResolvedItem) {
const identitySizeBound = claudeBackgroundTaskIdentity(id, Number.MAX_SAFE_INTEGER)
admission = sink.tryAppendResolvedItem(identitySizeBound, body, resolveIdentity, appendOptions)
if (admission.accepted) {
admission = sink.tryPublish
? sink.tryPublish(publishOptions)
: (sink.publish(publishOptions), ADMITTED)
}
} else if (sink.tryAppendItem) {
admission = sink.tryAppendItem(identity, body, appendOptions)
if (admission.accepted) {
admission = sink.tryPublish
? sink.tryPublish(publishOptions)
: (sink.publish(publishOptions), ADMITTED)
}
} else {
sink.appendItem(identity, body, appendOptions)
sink.publish(publishOptions)
admission = ADMITTED
}
sink.appendItem(identity, body, {
coalescingKey
})
sink.publish()
if (admission.accepted) {
row.lastSerialized = serialized
}
return admission
}
@@ -0,0 +1,92 @@
import type {
StructuredAgentSessionEventSink,
StructuredAgentSessionSinkAdmission
} from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import type { ClaudeBackgroundTaskRow } from './claude-background-task-row-lifecycle'
import {
ClaudeBackgroundTaskIdentityResolver,
writeClaudeBackgroundTaskRow
} from './claude-background-task-row-journal'
const MAX_PENDING_TASK_WRITES = 512
export class ClaudeBackgroundTaskRowWriter {
private readonly pending = new Map<
string,
{ id: string; row: ClaudeBackgroundTaskRow; lifecycle: boolean }
>()
private readonly identities = new ClaudeBackgroundTaskIdentityResolver()
constructor(
private readonly sink: StructuredAgentSessionEventSink,
private readonly onPersistenceFailure?: (error: Error) => void
) {}
write(
id: string,
row: ClaudeBackgroundTaskRow,
beforeAppend?: () => void,
lifecycle = false
): void {
const admission = writeClaudeBackgroundTaskRow(
this.sink,
this.identities,
id,
row,
beforeAppend,
lifecycle
)
const key = JSON.stringify([id, row.toolUseId ?? null, row.generation])
if (admission.accepted) {
this.pending.delete(key)
} else if (admission.reason === 'backpressure') {
if (!this.pending.has(key) && this.pending.size >= MAX_PENDING_TASK_WRITES) {
this.pending.clear()
this.onPersistenceFailure?.(
new Error('claude background task journal retry capacity exhausted')
)
return
}
this.pending.set(key, { id, row, lifecycle })
} else if (admission.reason === 'failed') {
this.onPersistenceFailure?.(new Error('claude background task journal sink failed'))
}
}
/** Replays bounded row obligations before provider reading resumes. */
retryPendingWrites(): StructuredAgentSessionSinkAdmission {
for (const [key, pending] of this.pending) {
const admission = writeClaudeBackgroundTaskRow(
this.sink,
this.identities,
pending.id,
pending.row,
undefined,
pending.lifecycle
)
if (!admission.accepted) {
if (admission.reason !== 'backpressure') {
this.pending.clear()
if (admission.reason === 'failed') {
this.onPersistenceFailure?.(new Error('claude background task journal sink failed'))
}
}
return admission
}
this.pending.delete(key)
}
return { accepted: true }
}
settlePendingWrites(): void {
for (const pending of this.pending.values()) {
pending.lifecycle = true
}
this.retryPendingWrites()
}
dispose(): void {
this.pending.clear()
this.identities.clear()
}
}
+40 -113
View File
@@ -3,21 +3,15 @@
import { isSettledBackgroundTaskState } from '../../shared/native-chat-background-task-row'
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import {
classifyClaudeBackgroundTaskKind,
record,
taskAliasId
} from './claude-background-task-frames'
import { record, taskAliasId } from './claude-background-task-frames'
import {
claudeBackgroundTaskPatchChange,
claudeBackgroundTaskToolUseId,
canonicalClaudeBackgroundTaskId,
finalizeClaudeBackgroundTaskRow,
isClaudeBackgroundTranscriptTask,
newClaudeBackgroundTaskRow,
newClaudeBackgroundTaskRowFromNotification,
reviseClaudeBackgroundTaskRow,
shouldRestartClaudeBackgroundTaskRow,
type ClaudeBackgroundTaskChange,
type ClaudeBackgroundTaskRow
} from './claude-background-task-row-lifecycle'
@@ -26,10 +20,10 @@ import {
ensureClaudeBackgroundTaskRowSlot,
type ClaudeBackgroundTaskLedgerSizes
} from './claude-background-task-memory'
import { writeClaudeBackgroundTaskRow } from './claude-background-task-row-journal'
import { ClaudeBackgroundTaskRowWriter } from './claude-background-task-row-writer'
import { observeClaudeBackgroundTaskStart } from './claude-background-task-start'
import { observeClaudeBackgroundTaskRoster } from './claude-background-task-roster'
import { ClaudeSubagentIds } from './claude-subagent-id-aliases'
import { isClaudeSubagentTask } from './claude-subagent-task-frames'
import { ClaudeOverflowTerminalRows } from './claude-overflow-terminal-rows'
const MAX_TASK_ROWS = 64
@@ -52,11 +46,13 @@ export type ClaudeBackgroundTaskRowsDeps = {
* output, so writing one must reopen a turn the provider resumed itself —
* otherwise the session renders the row while reporting idle. */
openOutputTurn?: (frame: Record<string, unknown>, observedAt: number) => void
onPersistenceFailure?: (error: Error) => void
now?: () => number
}
export class ClaudeBackgroundTaskRows {
private readonly rows = new Map<string, ClaudeBackgroundTaskRow>()
private readonly writer: ClaudeBackgroundTaskRowWriter
private readonly overflowTerminalRows: ClaudeOverflowTerminalRows
private readonly ledgers = new ClaudeBackgroundTaskLedgers()
private readonly ids = new ClaudeSubagentIds()
@@ -64,6 +60,7 @@ export class ClaudeBackgroundTaskRows {
constructor(private readonly deps: ClaudeBackgroundTaskRowsDeps) {
this.now = deps.now ?? (() => Date.now())
this.writer = new ClaudeBackgroundTaskRowWriter(deps.sink, deps.onPersistenceFailure)
this.overflowTerminalRows = new ClaudeOverflowTerminalRows(
this.ledgers,
this.now,
@@ -120,7 +117,15 @@ export class ClaudeBackgroundTaskRows {
return false
}
if (message.subtype === 'task_started') {
return this.observeStart(id, message)
return observeClaudeBackgroundTaskStart({
id,
message,
rows: this.rows,
ledgers: this.ledgers,
isForwardedParentTool: this.deps.isForwardedParentTool,
openRow: (taskId, frame) => this.openRow(taskId, frame),
maxRows: MAX_TASK_ROWS
})
}
if (this.ledgers.foreign.has(id)) {
return true
@@ -134,111 +139,23 @@ export class ClaudeBackgroundTaskRows {
settleSession(): void {
for (const [id, row] of this.rows) {
if (!isSettledBackgroundTaskState(row.block.state)) {
this.revise(id, { state: 'unverifiable' })
this.revise(id, { state: 'unverifiable' }, true)
}
}
this.writer.settlePendingWrites()
}
dispose(): void {
this.settleSession()
this.rows.clear()
this.writer.dispose()
this.overflowTerminalRows.clear()
this.ledgers.clear()
this.ids.clear()
}
private observeStart(id: string, message: Record<string, unknown>): boolean {
if (this.ledgers.fallbackTaskIds.has(id)) {
if (this.ledgers.terminalTaskIds.has(id)) {
const previousToolUseId = this.ledgers.terminalToolUseIds.get(id)
const currentToolUseId = claudeBackgroundTaskToolUseId(message)
if (
previousToolUseId !== undefined &&
currentToolUseId !== undefined &&
previousToolUseId !== currentToolUseId
) {
this.ledgers.fallbackTaskIds.delete(id)
} else {
return false
}
} else {
return false
}
}
if (message.ambient === true || message.skip_transcript === true) {
this.ledgers.rememberForeign(id, 'ambient')
return true
}
if (isClaudeSubagentTask(message)) {
this.ledgers.rememberForeign(id, 'roster')
return true
}
const kind = classifyClaudeBackgroundTaskKind(message.task_type)
if (!isClaudeBackgroundTranscriptTask(message, kind)) {
this.ledgers.rememberForeign(id, 'foreground')
return true
}
const existing = this.rows.get(id)
if (existing) {
this.ledgers.foreign.delete(id)
// A task that already exists and has not finished is not re-opened: a
// duplicate announcement is a redelivery, not a second run, and treating
// it as one would restate a row the user is already reading.
if (!isSettledBackgroundTaskState(existing.block.state)) {
return true
}
if (shouldRestartClaudeBackgroundTaskRow(existing, message)) {
this.openRow(id, message)
}
return true
}
let restartedTerminal = false
if (this.ledgers.terminalTaskIds.has(id)) {
const previousToolUseId = this.ledgers.terminalToolUseIds.get(id)
const currentToolUseId = claudeBackgroundTaskToolUseId(message)
// A terminal edge that had no usable tool id cannot prove a later start
// is a new run, so keep the conservative orphan guard. When both runs
// name their parent, a different alias is the provider's restart signal.
if (
previousToolUseId === undefined ||
currentToolUseId === undefined ||
previousToolUseId === currentToolUseId
) {
return true
}
restartedTerminal = true
}
this.ledgers.foreign.delete(id)
if (!this.admitsFirstRun(message)) {
// The refusal is recorded, not forgotten: the task belongs to the
// sidechain that spawned it, so its later frames find an owner here
// instead of looking like a task nothing ever decided about.
this.ledgers.rememberForeign(id, 'sidechain')
return true
}
if (!ensureClaudeBackgroundTaskRowSlot(this.rows, MAX_TASK_ROWS)) {
this.ledgers.rememberFallback(id)
return false
}
if (restartedTerminal) {
this.ledgers.terminalTaskIds.delete(id)
this.ledgers.terminalToolUseIds.delete(id)
}
this.openRow(id, message)
return true
}
/** The gate a task passes ONCE, when its first row is minted. Later frames
* for an admitted task are never re-gated: the decision belongs to the
* announcement, and re-asking it on a patch that carries no `tool_use_id`
* would drop the outcome of a task already on screen. */
private admitsFirstRun(message: Record<string, unknown>): boolean {
const toolUseId = claudeBackgroundTaskToolUseId(message)
// Conditional on the field being PRESENT. An announcement that names a tool
// this session never forwarded is a nested child and is refused; one that
// names no tool at all is admitted, because there is nothing to contradict
// — absence of the field is not evidence of an unforwarded parent.
return toolUseId === undefined || this.deps.isForwardedParentTool(toolUseId)
retryPendingWrites() {
return this.writer.retryPendingWrites()
}
private openRow(id: string, message: Record<string, unknown>): void {
@@ -323,30 +240,40 @@ export class ClaudeBackgroundTaskRows {
return true
}
private revise(id: string, change: ClaudeBackgroundTaskChange): void {
private revise(id: string, change: ClaudeBackgroundTaskChange, lifecycle = false): void {
const row = this.rows.get(id)
if (!row) {
return
}
const wasLive = !isSettledBackgroundTaskState(row.block.state)
reviseClaudeBackgroundTaskRow(row, change, this.now())
this.write(id, wasLive)
this.write(id, wasLive, lifecycle)
}
private write(id: string, openOutputTurn = true): void {
private write(id: string, openOutputTurn = true, lifecycle = false): void {
const row = this.rows.get(id)
if (!row) {
return
}
this.writeRow(id, row, openOutputTurn)
this.writeRow(id, row, openOutputTurn, lifecycle)
}
private writeRow(id: string, row: ClaudeBackgroundTaskRow, openOutputTurn = true): void {
private writeRow(
id: string,
row: ClaudeBackgroundTaskRow,
openOutputTurn = true,
lifecycle = false
): void {
const journaling = this.journaling
writeClaudeBackgroundTaskRow(this.deps.sink, id, row, () => {
if (journaling && openOutputTurn) {
this.deps.openOutputTurn?.(journaling.frame, journaling.observedAt)
}
})
this.writer.write(
id,
row,
() => {
if (journaling && openOutputTurn) {
this.deps.openOutputTurn?.(journaling.frame, journaling.observedAt)
}
},
lifecycle
)
}
}
@@ -0,0 +1,98 @@
import { isSettledBackgroundTaskState } from '../../shared/native-chat-background-task-row'
import { classifyClaudeBackgroundTaskKind } from './claude-background-task-frames'
import {
claudeBackgroundTaskToolUseId,
isClaudeBackgroundTranscriptTask,
shouldRestartClaudeBackgroundTaskRow,
type ClaudeBackgroundTaskRow
} from './claude-background-task-row-lifecycle'
import {
ensureClaudeBackgroundTaskRowSlot,
type ClaudeBackgroundTaskLedgers
} from './claude-background-task-memory'
import { isClaudeSubagentTask } from './claude-subagent-task-frames'
export function observeClaudeBackgroundTaskStart(input: {
id: string
message: Record<string, unknown>
rows: Map<string, ClaudeBackgroundTaskRow>
ledgers: ClaudeBackgroundTaskLedgers
isForwardedParentTool: (toolUseId: string) => boolean
openRow: (id: string, message: Record<string, unknown>) => void
maxRows: number
}): boolean {
const { id, message, rows, ledgers } = input
if (ledgers.fallbackTaskIds.has(id)) {
if (ledgers.terminalTaskIds.has(id)) {
const previousToolUseId = ledgers.terminalToolUseIds.get(id)
const currentToolUseId = claudeBackgroundTaskToolUseId(message)
if (
previousToolUseId !== undefined &&
currentToolUseId !== undefined &&
previousToolUseId !== currentToolUseId
) {
ledgers.fallbackTaskIds.delete(id)
} else {
return false
}
} else {
return false
}
}
if (message.ambient === true || message.skip_transcript === true) {
ledgers.rememberForeign(id, 'ambient')
return true
}
if (isClaudeSubagentTask(message)) {
ledgers.rememberForeign(id, 'roster')
return true
}
const kind = classifyClaudeBackgroundTaskKind(message.task_type)
if (!isClaudeBackgroundTranscriptTask(message, kind)) {
ledgers.rememberForeign(id, 'foreground')
return true
}
const existing = rows.get(id)
if (existing) {
ledgers.foreign.delete(id)
// A live task's duplicate announcement is redelivery, not a new run.
if (!isSettledBackgroundTaskState(existing.block.state)) {
return true
}
if (shouldRestartClaudeBackgroundTaskRow(existing, message)) {
input.openRow(id, message)
}
return true
}
let restartedTerminal = false
if (ledgers.terminalTaskIds.has(id)) {
const previousToolUseId = ledgers.terminalToolUseIds.get(id)
const currentToolUseId = claudeBackgroundTaskToolUseId(message)
// A terminal edge without a usable parent cannot prove a later start is a new run.
if (
previousToolUseId === undefined ||
currentToolUseId === undefined ||
previousToolUseId === currentToolUseId
) {
return true
}
restartedTerminal = true
}
ledgers.foreign.delete(id)
const toolUseId = claudeBackgroundTaskToolUseId(message)
// Absence of a parent is not evidence of an unforwarded parent.
if (toolUseId !== undefined && !input.isForwardedParentTool(toolUseId)) {
ledgers.rememberForeign(id, 'sidechain')
return true
}
if (!ensureClaudeBackgroundTaskRowSlot(rows, input.maxRows)) {
ledgers.rememberFallback(id)
return false
}
if (restartedTerminal) {
ledgers.terminalTaskIds.delete(id)
ledgers.terminalToolUseIds.delete(id)
}
input.openRow(id, message)
return true
}
@@ -37,6 +37,131 @@ function fakeChild(): ChildProcessWithoutNullStreams {
}
describe('Claude stream-json close ordering', () => {
it('stops pulling SDK messages until reading resumes', async () => {
mocks.refresh.mockReset()
mocks.proveClaudeChildExit.mockReset()
mocks.refresh.mockResolvedValue(undefined)
mocks.proveClaudeChildExit.mockResolvedValue(true)
const child = fakeChild()
const first = Promise.withResolvers<Record<string, unknown>>()
const next = vi
.fn<() => Promise<IteratorResult<Record<string, unknown>>>>()
.mockImplementationOnce(async () => ({ value: await first.promise, done: false }))
.mockResolvedValueOnce({ value: { type: 'second' }, done: false })
.mockResolvedValue({ value: undefined, done: true })
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: This injected query exercises only the async iterator used by the connection.
const queryImpl = ((params: Parameters<typeof query>[0]) => {
params.options?.spawnClaudeCodeProcess?.({
command: 'claude',
args: [],
env: {},
signal: new AbortController().signal
})
return {
[Symbol.asyncIterator]: () => ({ next })
}
}) as unknown as typeof query
const seen: string[] = []
let connection: Awaited<ReturnType<typeof openClaudeStreamJsonConnection>>
connection = await openClaudeStreamJsonConnection(
{ pathToClaudeCodeExecutable: 'claude', options: {}, cwd: '/work/repo' },
{
onMessage: (message) => {
seen.push(String(message.type))
if (message.type === 'first') {
connection.pauseReading?.()
}
}
},
() => child,
queryImpl
)
first.resolve({ type: 'first' })
await vi.waitFor(() => expect(seen).toEqual(['first']))
await new Promise((resolve) => setImmediate(resolve))
expect(next).toHaveBeenCalledOnce()
connection.resumeReading?.()
await vi.waitFor(() => expect(seen).toEqual(['first', 'second']))
expect(next).toHaveBeenCalledTimes(3)
await expect(connection.close()).resolves.toBe(true)
})
it('releases a pulled frame when provider exit is reported', async () => {
mocks.refresh.mockReset()
mocks.proveClaudeChildExit.mockReset()
mocks.refresh.mockResolvedValue(undefined)
mocks.proveClaudeChildExit.mockResolvedValue(true)
const child = fakeChild()
const first = Promise.withResolvers<Record<string, unknown>>()
const next = vi
.fn<() => Promise<IteratorResult<Record<string, unknown>>>>()
.mockImplementationOnce(async () => ({ value: await first.promise, done: false }))
.mockResolvedValue({ value: undefined, done: true })
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: This injected query exercises only the async iterator used by the connection.
const queryImpl = ((params: Parameters<typeof query>[0]) => {
params.options?.spawnClaudeCodeProcess?.({
command: 'claude',
args: [],
env: {},
signal: new AbortController().signal
})
return {
[Symbol.asyncIterator]: () => ({ next })
}
}) as unknown as typeof query
const events: string[] = []
const connection = await openClaudeStreamJsonConnection(
{ pathToClaudeCodeExecutable: 'claude', options: {}, cwd: '/work/repo' },
{
onMessage: (message) => events.push(`message:${String(message.type)}`),
onExit: () => events.push('exit')
},
() => child,
queryImpl
)
connection.pauseReading?.()
first.resolve({ type: 'task_notification' })
await vi.waitFor(() => expect(next).toHaveBeenCalledOnce())
expect(events).toEqual([])
child.emit('exit', 1, null)
await vi.waitFor(() => expect(events).toEqual(['exit', 'message:task_notification']))
await expect(connection.close()).resolves.toBe(true)
})
it('returns an unproven close without waiting on a live output reader', async () => {
mocks.refresh.mockReset()
mocks.proveClaudeChildExit.mockReset()
mocks.refresh.mockResolvedValue(undefined)
mocks.proveClaudeChildExit.mockResolvedValue(false)
const child = fakeChild()
const next = vi.fn(() => new Promise<IteratorResult<Record<string, unknown>>>(() => {}))
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: This injected query exercises only the async iterator used by the connection.
const queryImpl = ((params: Parameters<typeof query>[0]) => {
params.options?.spawnClaudeCodeProcess?.({
command: 'claude',
args: [],
env: {},
signal: new AbortController().signal
})
return {
[Symbol.asyncIterator]: () => ({ next })
}
}) as unknown as typeof query
const connection = await openClaudeStreamJsonConnection(
{ pathToClaudeCodeExecutable: 'claude', options: {}, cwd: '/work/repo' },
{},
() => child,
queryImpl
)
await expect(connection.close()).resolves.toBe(false)
expect(next).toHaveBeenCalledOnce()
})
it('waits for the live tree refresh before ending stdin', async () => {
const refreshDone = Promise.withResolvers<void>()
mocks.refresh.mockReturnValueOnce(refreshDone.promise)
@@ -38,6 +38,10 @@ function loadClaudeAgentSdk(): Promise<typeof ClaudeAgentSdk> {
return claudeAgentSdk
}
function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === 'object' && value !== null
}
export type ClaudeStreamJsonLaunch = {
/** Orca's resolved user CLI; the SDK falls back to a bundled binary that is not installed. */
pathToClaudeCodeExecutable: string
@@ -78,6 +82,8 @@ export type ClaudeStreamJsonConnection = ClaudeControlSurface & {
readonly closed: boolean
/** What the ladder has observed so far; read after a `close()` that returned false. */
readonly exitVerdict: ClaudeChildExitVerdict
pauseReading?: () => void
resumeReading?: () => void
send: (message: Record<string, unknown>) => Promise<void>
/** Resolves true after processless settlement, or root exit plus observed tree exit. */
close: () => Promise<boolean>
@@ -141,6 +147,23 @@ export async function openClaudeStreamJsonConnection(
let faultReported = false
let exitReported = false
let closePromise: Promise<boolean> | null = null
let readingBarrier: Promise<void> | null = null
let releaseReadingBarrier: (() => void) | null = null
const pauseReading = (): void => {
if (closing || exited || terminalError || readingBarrier) {
return
}
readingBarrier = new Promise<void>((resolve) => {
releaseReadingBarrier = resolve
})
}
const resumeReading = (): void => {
const release = releaseReadingBarrier
readingBarrier = null
releaseReadingBarrier = null
release?.()
}
const waitUntilReadable = (): Promise<void> => readingBarrier ?? Promise.resolve()
// One reaper per child: every close attempt and error-path reap shares its proof.
const rootSettled = (): boolean => exited || processless
const tree = createClaudeChildTreeReaper(child, { exited: rootSettled })
@@ -180,6 +203,7 @@ export async function openClaudeStreamJsonConnection(
})
const handleUnexpectedEnd = (cause?: Error): void => {
resumeReading()
terminalError ??= exitError(spawner.stderrTail, exitStatus, cause)
inbox.fail(terminalError)
if (!closing && !faultReported) {
@@ -192,18 +216,38 @@ export async function openClaudeStreamJsonConnection(
}
}
void (async () => {
for await (const message of session) {
handlers.onMessage?.(message as unknown as Record<string, unknown>)
const readerDone = (async () => {
try {
const iterator = session[Symbol.asyncIterator]()
let completed = false
try {
for (;;) {
await waitUntilReadable()
const next = await iterator.next()
if (next.done) {
completed = true
break
}
await waitUntilReadable()
if (!isRecord(next.value)) {
throw new Error('claude stream-json yielded a non-object message')
}
handlers.onMessage?.(next.value)
}
} finally {
if (!completed) {
await iterator.return?.()
}
}
} catch (error: unknown) {
// The SDK ends its generator in error when the child dies or the transport
// fails; a transport failure with a live child still has to reap the tree.
if (!closing && !exited) {
void tree.reap()
}
handleUnexpectedEnd(error instanceof Error ? error : new Error(String(error)))
}
})().catch((error: unknown) => {
// The SDK ends its generator in error when the child dies or the transport
// fails; a transport failure with a live child still has to reap the tree.
if (!closing && !exited) {
void tree.reap()
}
handleUnexpectedEnd(error instanceof Error ? error : new Error(String(error)))
})
})()
child.on('error', (error) => {
if (spawner.pid === undefined) {
@@ -251,6 +295,7 @@ export async function openClaudeStreamJsonConnection(
const close = (): Promise<boolean> => {
closePromise ??= (async () => {
closing = true
resumeReading()
// Arm the descendant proof before ending stdin. The SDK may exit the root
// immediately; a post-exit walk cannot recover descendants that reparented.
await (tree.refresh?.() ?? tree.capture())
@@ -264,8 +309,10 @@ export async function openClaudeStreamJsonConnection(
inbox.fail(new Error('claude stream-json connection closed'))
if (!proven) {
closePromise = null
return false
}
return proven
await readerDone
return true
})()
return closePromise
}
@@ -284,6 +331,8 @@ export async function openClaudeStreamJsonConnection(
tree: tree.treeVerdict
} as const
},
pauseReading,
resumeReading,
send,
close
}
@@ -160,6 +160,88 @@ function persistedTarget(
}
describe('claude journal translation — background task rows', () => {
it('persists one terminal row at the hard watermark without an opcode fallback', async () => {
const persisted = new Map<string, AgentJournalItemBody>()
const appendEntered = Promise.withResolvers<void>()
const appendGate = Promise.withResolvers<void>()
const target = persistedTarget(persisted)
const appendItem = target.journal.appendItem.bind(target.journal)
vi.spyOn(target.journal, 'appendItem').mockImplementationOnce(async (...args) => {
appendEntered.resolve()
await appendGate.promise
return appendItem(...args)
})
const deferred = createDeferredStructuredAgentSessionEventSink({
watermarks: {
pauseQueuedOperations: 1,
maxQueuedOperations: 4,
lowQueuedOperations: 0,
maxQueuedBytes: 1_000_000
}
})
const translator = createClaudeJournalTranslator({
sink: deferred.sink,
fallbackIdPrefix: 'hard-watermark'
})
const providerResume = vi.fn()
let sinkPaused = false
deferred.sink.bindReadingControl?.({
pauseReading: () => {
sinkPaused = true
},
resumeReading: () => {
sinkPaused = false
const admission = translator.retryPendingTaskRows?.() ?? { accepted: true }
if (!sinkPaused && (admission.accepted || admission.reason !== 'backpressure')) {
providerResume()
}
}
})
deferred.bind(target)
deferred.sink.appendItem(
{ provider: 'orca', clientMessageId: 'blocked-prefill' },
{ kind: 'message', role: 'system', blocks: [{ type: 'text', text: 'prefill' }] }
)
await appendEntered.promise
const notification = systemFrame({
subtype: 'task_notification',
task_id: 'hard-watermark-task',
tool_use_id: 'toolu-hard-watermark',
status: 'failed',
summary: 'The real provider task failed',
uuid: 'hard-watermark-notification'
})
translator.handle(notification)
translator.handle(notification)
expect(deferred.state().queuedOperations).toBe(4)
expect(translator.retryPendingTaskRows?.()).toEqual({
accepted: false,
reason: 'backpressure'
})
expect(
[...persisted.values()].filter(
(body) => body.kind === 'message' && blockOf(body)?.taskId === 'hard-watermark-task'
)
).toEqual([])
appendGate.resolve()
await vi.waitFor(() => expect(providerResume).toHaveBeenCalledOnce())
await expect(deferred.drained()).resolves.toEqual({ ok: true })
const taskRows = [...persisted.values()].filter(
(body) => body.kind === 'message' && blockOf(body)?.taskId === 'hard-watermark-task'
)
expect(taskRows).toHaveLength(1)
expect(blockOf(taskRows[0])?.error).toBeUndefined()
expect(blockOf(taskRows[0])?.summary).toBe('The real provider task failed')
expect(
[...persisted.values()].some(
(body) =>
body.kind === 'status' && body.providerFrame?.kind.includes('task_notification') === true
)
).toBe(false)
})
it('coalesces an unbound overflow patch and aliased final notification', async () => {
const persisted = new Map<string, AgentJournalItemBody>()
const deferred = createDeferredStructuredAgentSessionEventSink()
@@ -1,5 +1,8 @@
import type { AgentSessionDeltaCoalescerDeps } from '../native-chat/agent-session-wire/agent-session-delta-coalescer'
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import type {
StructuredAgentSessionEventSink,
StructuredAgentSessionSinkAdmission
} from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import type { ClaudeStructuredSessionEvent } from './claude-structured-session-state'
import {
claudeStreamingMessageBody,
@@ -35,6 +38,7 @@ export type ClaudeJournalTranslatorDeps = {
coalesceMs?: number
schedule?: AgentSessionDeltaCoalescerDeps['schedule']
fallbackIdPrefix?: string
onBackgroundTaskJournalFailure?: (error: Error) => void
}
export type ClaudeJournalTranslator = {
@@ -44,6 +48,7 @@ export type ClaudeJournalTranslator = {
* a client's Stop names. Sole owner: no reader keeps a copy to disagree with. */
readonly currentTurnId: string | null
flush: () => void
retryPendingTaskRows?: () => StructuredAgentSessionSinkAdmission
/** Streamed blocks still awaiting a final frame. A settled turn leaves none. */
readonly pendingStreamedBlocks: number
dispose: () => void
@@ -52,12 +57,14 @@ export type ClaudeJournalTranslator = {
export function createClaudeSessionJournalTranslator(
sink: StructuredAgentSessionEventSink | undefined,
prompts: ClaudePromptRegistry,
fallbackIdPrefix: string
fallbackIdPrefix: string,
onBackgroundTaskJournalFailure?: (error: Error) => void
): ClaudeJournalTranslator | null {
return sink
? createClaudeJournalTranslator({
sink,
fallbackIdPrefix,
...(onBackgroundTaskJournalFailure ? { onBackgroundTaskJournalFailure } : {}),
bindPromptItemId: (itemId, promptKey, questionId) =>
prompts.bindJournalItemId(itemId, promptKey, questionId)
})
@@ -89,7 +96,10 @@ export function createClaudeJournalTranslator(
// A typed task row is provider output: journaling one must open a resumed
// turn, or the session shows the row while reading idle.
openOutputTurn: (frame, observedAt) =>
turn.ensureOpen(frame, claudeStreamTurnSource(frame), observedAt)
turn.ensureOpen(frame, claudeStreamTurnSource(frame), observedAt),
...(deps.onBackgroundTaskJournalFailure
? { onPersistenceFailure: deps.onBackgroundTaskJournalFailure }
: {})
})
const streamedText = createClaudeStreamedTextCheckpoints({
...(deps.coalesceMs === undefined ? {} : { coalesceMs: deps.coalesceMs }),
@@ -224,6 +234,7 @@ export function createClaudeJournalTranslator(
return turn.id
},
flush: streamedText.flush,
retryPendingTaskRows: () => backgroundTasks.retryPendingWrites(),
get pendingStreamedBlocks() {
return streamedText.pending
},
@@ -47,6 +47,10 @@ import { resolveClaudeAcquisitionError } from './claude-structured-session-close
import { readClaudeTranscriptEntryUuid } from './claude-tui-exit'
import { withAgentSessionCreatePhase } from '../observability/agent-session-instrumentation'
import { resolveClaudeAcquisitionLaunch } from './claude-structured-acquisition-launch'
import {
bindClaudeJournalReadingControl,
createClaudeJournalFailureHandler
} from './claude-structured-session-journal-control'
export const CLAUDE_STRUCTURED_INIT_TIMEOUT_MS = 10_000
@@ -72,12 +76,8 @@ export async function acquireClaudeSession({
}
const sessionId = input.identity.sessionId
const prompts = new ClaudePromptRegistry()
const translator = createClaudeSessionJournalTranslator(
input.events,
prompts,
String(input.fence)
)
const { previous, attempt } = acquisitions.start(sessionId, prompts)
let unbindReadingControl: (() => void) | undefined
let liveSession: ClaudeSession | null = null
let observedLeafUuid: string | null = null,
expectedProviderSessionId: string | null = null
@@ -85,6 +85,12 @@ export async function acquireClaudeSession({
// this acquisition owns. Keep the check ahead of every stateful consumer.
const initTimeoutMs = deps.initTimeoutMs ?? CLAUDE_STRUCTURED_INIT_TIMEOUT_MS
const initDeadline = createClaudeInitDeadline(sessionId, initTimeoutMs)
const translator = createClaudeSessionJournalTranslator(
input.events,
prompts,
String(input.fence),
createClaudeJournalFailureHandler({ attempt, initDeadline, callbacks, sessionId })
)
const rewind = new ClaudeRewindAttempt(input.rewind, input.rewind?.onProved)
const onMessage = (message: Record<string, unknown>): void => {
@@ -198,6 +204,7 @@ export async function acquireClaudeSession({
)
)
attempt.connection = connection
unbindReadingControl = bindClaudeJournalReadingControl(input.events, connection, translator)
acquisitions.assertCurrent(sessionId, attempt)
initDeadline.start()
const [initialization, init] = await withAgentSessionCreatePhase(
@@ -259,6 +266,7 @@ export async function acquireClaudeSession({
prompts,
translator,
events: input.events,
...(unbindReadingControl ? { unbindReadingControl } : {}),
process,
acquisitionGeneration: mintClaudeAcquisitionGeneration(deps),
options: acquisitionOptions.options,
@@ -284,6 +292,7 @@ export async function acquireClaudeSession({
return acquired
} catch (error) {
initDeadline.clear()
unbindReadingControl?.()
const acquisitionError = await resolveClaudeAcquisitionError({
error,
sessionId,
@@ -26,7 +26,10 @@ import {
closeClaudeSession,
settleClaudeExitedSession
} from './claude-structured-session-close'
import { readClaudeTranscriptLeafWithReproof } from './claude-transcript-branch-proof'
import {
drainClaudeObservedExits,
persistClaudeSessionHandle
} from './claude-structured-session-exit-lifecycle'
import type { AgentSessionBackgroundTaskState } from '../../shared/agent-session-wire'
import { resolveClaudeProviderHistoryWindow } from './claude-structured-history-window'
import {
@@ -80,7 +83,10 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
attempt.buffered.push(event)
return
}
if (this.sessions.get(sessionId)?.connection === attempt.connection) {
if (
this.sessions.get(sessionId)?.connection === attempt.connection ||
this.exits.get(sessionId)?.connection === attempt.connection
) {
event()
}
}
@@ -103,7 +109,12 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
}
this.exits.set(sessionId, exit)
exit.publication = closePromise
.then((proven) => (proven ? this.settleUnexpectedExit(sessionId, exit) : undefined))
.then((proven) => {
if (!proven) {
return undefined
}
return this.settleUnexpectedExit(sessionId, exit)
})
.catch(() => undefined)
}
@@ -112,37 +123,19 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
* retry. Publication trails observation by the close ladder and the
* transcript cursor write, so nothing outside can otherwise tell the two
* apart without guessing at wall-clock. */
drainObservedExits = async (): Promise<void> => {
const awaited = new Set<Promise<void>>()
for (;;) {
const pending = [...this.exits.values()]
.map((exit) => exit.publication)
.filter(
(publication): publication is Promise<void> =>
publication !== undefined && !awaited.has(publication)
)
if (pending.length === 0) {
return
}
for (const publication of pending) {
awaited.add(publication)
}
// A publication can settle an exit that itself observes another; only the
// ones this pass has not already awaited keep the loop going.
await Promise.all(pending)
}
}
drainObservedExits = (): Promise<void> => drainClaudeObservedExits(this.exits)
/** Lifecycle recovery is published only after the child tree proof is true. */
private settleUnexpectedExit(sessionId: string, exit: ClaudeSessionExit): Promise<void> {
exit.settlementPromise ??= (async () => {
exit.session.unbindReadingControl?.()
if (this.exits.get(sessionId) !== exit) {
settleClaudeExitedSession(exit.session)
return
}
// Persist the transcript-derived cursor before publishing the lifecycle
// event that lets the host release and reacquire this exact child.
await this.persistSessionHandle(sessionId, exit.session).catch(() => undefined)
await persistClaudeSessionHandle(sessionId, exit.session, this.deps).catch(() => undefined)
if (this.exits.get(sessionId) !== exit) {
settleClaudeExitedSession(exit.session)
return
@@ -177,30 +170,6 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
this.sessions.has(input.identity.sessionId) || this.exits.has(input.identity.sessionId)
})
private async persistSessionHandle(sessionId: string, session: ClaudeSession): Promise<void> {
try {
const transcriptLeaf = this.deps.readTranscriptLeaf
? await readClaudeTranscriptLeafWithReproof({
readTranscriptLeaf: this.deps.readTranscriptLeaf,
providerSessionId: session.providerSessionId,
previousLeafUuid: session.leafUuid,
claudeConfigDir: session.claudeConfigDir
})
: null
if (transcriptLeaf) {
session.leafUuid = transcriptLeaf
}
} catch {
// A stale or unavailable tail must not overwrite the last observed leaf.
}
await this.deps.persistHandle?.({
sessionId,
providerSessionId: session.providerSessionId,
leafUuid: session.leafUuid,
fence: session.fence
})
}
private emit(session: ClaudeSession | null, event: ClaudeStructuredSessionEvent): void {
const backgroundTasksChanged =
event.type === 'ended'
@@ -97,7 +97,9 @@ async function finalizeClaudePublishedSession(
for (const prompt of session.prompts.clear()) {
prompt.settle(null)
}
if ((await session.connection.close()) !== true) {
const connectionClosed = await session.connection.close()
session.unbindReadingControl?.()
if (connectionClosed !== true) {
const cleanupError = claudeAcquisitionCleanupError(
session.connection,
new Error('provider close unproven')
@@ -0,0 +1,56 @@
import { readClaudeTranscriptLeafWithReproof } from './claude-transcript-branch-proof'
import type {
ClaudeSession,
ClaudeSessionExit,
ClaudeStructuredSessionAdapterDeps
} from './claude-structured-session-state'
/** Wait for each first-hand exit's publication, including exits observed while waiting. */
export async function drainClaudeObservedExits(
exits: Map<string, ClaudeSessionExit>
): Promise<void> {
const awaited = new Set<Promise<void>>()
for (;;) {
const pending = [...exits.values()]
.map((exit) => exit.publication)
.filter(
(publication): publication is Promise<void> =>
publication !== undefined && !awaited.has(publication)
)
if (pending.length === 0) {
return
}
for (const publication of pending) {
awaited.add(publication)
}
await Promise.all(pending)
}
}
export async function persistClaudeSessionHandle(
sessionId: string,
session: ClaudeSession,
deps: Pick<ClaudeStructuredSessionAdapterDeps, 'readTranscriptLeaf' | 'persistHandle'>
): Promise<void> {
try {
const transcriptLeaf = deps.readTranscriptLeaf
? await readClaudeTranscriptLeafWithReproof({
readTranscriptLeaf: deps.readTranscriptLeaf,
providerSessionId: session.providerSessionId,
previousLeafUuid: session.leafUuid,
claudeConfigDir: session.claudeConfigDir
})
: null
if (transcriptLeaf) {
session.leafUuid = transcriptLeaf
}
} catch {
// An unavailable tail must not overwrite the last observed leaf.
}
await deps.persistHandle?.({
sessionId,
providerSessionId: session.providerSessionId,
leafUuid: session.leafUuid,
fence: session.fence
})
}
@@ -0,0 +1,53 @@
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import type { ClaudeStreamJsonConnection } from './claude-stream-json-connection'
import type { ClaudeJournalTranslator } from './claude-structured-journal-translation'
import type { createClaudeInitDeadline } from './claude-structured-init-deadline'
import type {
ClaudeAcquisitionAttempt,
ClaudeAcquireCallbacks
} from './claude-structured-session-state'
export function createClaudeJournalFailureHandler(input: {
attempt: ClaudeAcquisitionAttempt
initDeadline: ReturnType<typeof createClaudeInitDeadline>
callbacks: ClaudeAcquireCallbacks
sessionId: string
}): (error: Error) => void {
return (error) => {
if (!input.attempt.published) {
input.initDeadline.reject(error)
return
}
const connection = input.attempt.connection
if (connection) {
void connection
.close()
.catch(() => false)
.finally(() => input.callbacks.handleExit(input.sessionId, input.attempt, error))
}
}
}
export function bindClaudeJournalReadingControl(
sink: StructuredAgentSessionEventSink | undefined,
connection: ClaudeStreamJsonConnection,
translator: ClaudeJournalTranslator | null
): (() => void) | undefined {
if (!connection.pauseReading || !connection.resumeReading) {
return undefined
}
let sinkPaused = false
return sink?.bindReadingControl?.({
pauseReading: () => {
sinkPaused = true
connection.pauseReading?.()
},
resumeReading: () => {
sinkPaused = false
const retried = translator?.retryPendingTaskRows?.() ?? { accepted: true }
if (!sinkPaused && (retried.accepted || retried.reason !== 'backpressure')) {
connection.resumeReading?.()
}
}
})
}
@@ -19,6 +19,7 @@ export function createClaudeSessionPublication(input: {
prompts: ClaudePromptRegistry
translator: ClaudeJournalTranslator | null
events: ClaudeSession['events']
unbindReadingControl?: () => void
process: AgentSessionAcquisition['process']
linkId?: string
observedAt: number
@@ -83,7 +84,8 @@ export function createClaudeSessionPublication(input: {
]),
restoreSkippedOptions: new Set(),
translator: input.translator,
events: input.events
events: input.events,
...(input.unbindReadingControl ? { unbindReadingControl: input.unbindReadingControl } : {})
}
}
}
@@ -0,0 +1,263 @@
import { describe, expect, it, vi } from 'vitest'
import {
createDeferredStructuredAgentSessionEventSink,
type StructuredAgentSessionEventTarget,
type StructuredAgentSessionEventSink,
type StructuredAgentSessionReadingControl
} from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key'
import type {
AgentJournalItemBody,
AgentJournalItemIdentity
} from '../../shared/agent-session-journal-types'
import type { AgentSessionJournal } from '../native-chat/agent-session-journal/journal-store'
import { blockOf } from './claude-background-task-row-test-support'
import {
adapterFor,
fakeClaude,
identityFor,
PROVIDER_SESSION_ID
} from './claude-structured-session-test-support'
function controlledSink(): {
sink: StructuredAgentSessionEventSink
control: () => StructuredAgentSessionReadingControl | undefined
unbind: ReturnType<typeof vi.fn>
} {
let control: StructuredAgentSessionReadingControl | undefined
const unbind = vi.fn()
return {
sink: {
appendItem: vi.fn(),
appendTombstone: vi.fn(),
publish: vi.fn(),
bindReadingControl: (next) => {
control = next
return unbind
}
},
control: () => control,
unbind
}
}
function persistedTarget(
persisted: Map<string, AgentJournalItemBody>
): StructuredAgentSessionEventTarget {
const journal =
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: This test double implements the journal methods exercised by the deferred sink.
{
appendItem: async (identity: AgentJournalItemIdentity, body: AgentJournalItemBody) => {
persisted.set(agentJournalItemKey(identity), body)
return { cursor: { epoch: 'test', sequence: persisted.size }, itemId: '', revision: 1 }
},
appendTombstone: vi.fn(),
visitItems: (
visit: (itemId: string, sequence: number, body: AgentJournalItemBody) => void
) => {
for (const [itemId, body] of persisted) {
visit(itemId, 0, body)
}
},
epoch: 'test'
} as unknown as AgentSessionJournal
return { journal, fence: 1, publish: vi.fn() }
}
describe('Claude structured reading control', () => {
it('binds sink pressure to SDK reading and unbinds on requested close', async () => {
const claude = fakeClaude()
const adapter = adapterFor(claude)
const events = controlledSink()
await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9',
events: events.sink
})
const pauseReading = vi.spyOn(claude.connections[0], 'pauseReading')
const resumeReading = vi.spyOn(claude.connections[0], 'resumeReading')
events.control()?.pauseReading()
expect(pauseReading).toHaveBeenCalledOnce()
events.control()?.resumeReading()
expect(resumeReading).toHaveBeenCalledOnce()
await expect(adapter.closeSession('session-1')).resolves.toBe(true)
expect(events.unbind).toHaveBeenCalledOnce()
})
it('unbinds when acquisition fails after the connection opens', async () => {
const claude = fakeClaude({ initProof: 'none' })
const adapter = adapterFor(claude, {}, [], [], 1)
const events = controlledSink()
await expect(
adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9',
events: events.sink
})
).rejects.toThrow('did not finish starting')
expect(events.unbind).toHaveBeenCalledOnce()
})
it('unbinds when the published provider exits unexpectedly', async () => {
const claude = fakeClaude()
const adapter = adapterFor(claude)
const events = controlledSink()
await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9',
events: events.sink
})
claude.connections[0].handlers.onExit?.(new Error('provider exited'))
await adapter.drainObservedExits()
expect(events.unbind).toHaveBeenCalledOnce()
})
it('keeps delivery ownership until a reported provider exit finishes closing', async () => {
const claude = fakeClaude()
const adapter = adapterFor(claude)
const events = controlledSink()
await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9',
events: events.sink
})
const connection = claude.connections[0]
connection.handlers.onExit?.(new Error('provider exited'))
connection.handlers.onMessage?.({
type: 'system',
subtype: 'task_notification',
session_id: PROVIDER_SESSION_ID,
task_id: 'held-terminal-frame',
status: 'failed',
summary: 'The final task outcome'
})
expect(events.sink.appendItem).toHaveBeenCalledWith(
expect.anything(),
expect.objectContaining({
kind: 'message',
role: 'system',
blocks: expect.arrayContaining([
expect.objectContaining({
type: 'background-task',
taskId: 'held-terminal-frame',
summary: 'The final task outcome'
})
])
}),
expect.anything()
)
await adapter.drainObservedExits()
expect(events.unbind).toHaveBeenCalledOnce()
})
it('releases SDK reading when a pending row becomes permanently refused', async () => {
const claude = fakeClaude()
const adapter = adapterFor(claude)
const events = controlledSink()
events.sink.tryAppendResolvedItemAndPublish = vi
.fn()
.mockReturnValueOnce({ accepted: false, reason: 'backpressure' })
.mockReturnValueOnce({ accepted: false, reason: 'failed' })
await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9',
events: events.sink
})
const resumeReading = vi.spyOn(claude.connections[0], 'resumeReading')
claude.connections[0].handlers.onMessage?.({
type: 'system',
subtype: 'task_notification',
session_id: PROVIDER_SESSION_ID,
task_id: 'failed-journal-row',
status: 'failed',
summary: 'failed'
})
events.control()?.pauseReading()
events.control()?.resumeReading()
expect(resumeReading).toHaveBeenCalledOnce()
await adapter.drainObservedExits()
})
it('automatically retries a hard-watermark row before resuming SDK reads', async () => {
const persisted = new Map<string, AgentJournalItemBody>()
const target = persistedTarget(persisted)
const deferred = createDeferredStructuredAgentSessionEventSink({
watermarks: {
pauseQueuedOperations: 1,
maxQueuedOperations: 4,
lowQueuedOperations: 0,
maxQueuedBytes: 1_000_000
}
})
deferred.bind(target)
const claude = fakeClaude()
const adapter = adapterFor(claude)
await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9',
events: deferred.sink
})
await deferred.drained()
persisted.clear()
const appendEntered = Promise.withResolvers<void>()
const appendGate = Promise.withResolvers<void>()
const appendItem = target.journal.appendItem.bind(target.journal)
vi.spyOn(target.journal, 'appendItem').mockImplementationOnce(async (...args) => {
appendEntered.resolve()
await appendGate.promise
return appendItem(...args)
})
const resumeReading = vi.spyOn(claude.connections[0], 'resumeReading')
deferred.sink.appendItem(
{ provider: 'orca', clientMessageId: 'blocked-prefill' },
{ kind: 'message', role: 'system', blocks: [{ type: 'text', text: 'prefill' }] }
)
await appendEntered.promise
const notification = {
type: 'system',
subtype: 'task_notification',
session_id: PROVIDER_SESSION_ID,
task_id: 'hard-watermark-task',
tool_use_id: 'toolu-hard-watermark',
status: 'failed',
summary: 'The real provider task failed',
uuid: 'hard-watermark-notification'
}
claude.connections[0].handlers.onMessage?.(notification)
claude.connections[0].handlers.onMessage?.(notification)
expect(deferred.state().queuedOperations).toBe(4)
appendGate.resolve()
await vi.waitFor(() => expect(resumeReading).toHaveBeenCalledOnce())
await expect(deferred.drained()).resolves.toEqual({ ok: true })
const taskRows = [...persisted.values()].filter(
(body) => body.kind === 'message' && blockOf(body)?.taskId === 'hard-watermark-task'
)
expect(taskRows).toHaveLength(1)
expect(blockOf(taskRows[0])?.summary).toBe('The real provider task failed')
expect(
[...persisted.values()].some(
(body) =>
body.kind === 'status' && body.providerFrame?.kind.includes('task_notification') === true
)
).toBe(false)
await expect(adapter.closeSession('session-1')).resolves.toBe(true)
})
})
@@ -170,6 +170,7 @@ export type ClaudeSession = {
closeEnded?: boolean
translator: ClaudeJournalTranslator | null
events: StructuredAgentSessionEventSink | undefined
unbindReadingControl?: () => void
}
export function mintClaudeAcquisitionGeneration(deps: ClaudeStructuredSessionAdapterDeps): string {
@@ -74,6 +74,7 @@ export function fakeClaude(
const route = routes[subtype]
return route ? route(params) : undefined
}
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The fake implements the complete connection contract below.
const openConnection = (async (launch, handlers = {}) => {
const connection: FakeConnection = {
launch,
@@ -83,6 +84,8 @@ export function fakeClaude(
closeCount: 0,
pid: 4321,
closed: false,
pauseReading: () => {},
resumeReading: () => {},
initializationResult: async () => {
connection.calls.push({ subtype: 'initialize' })
if (options.exitBeforeInit) {
@@ -224,6 +224,35 @@ describe('deferred structured agent-session event sink', () => {
expect(log).toHaveLength(2)
})
it('admits a resolved append and publication as one bounded operation', async () => {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink({
watermarks: {
pauseQueuedOperations: 1,
maxQueuedOperations: 2,
lowQueuedOperations: 0,
maxQueuedBytes: 1_000_000
}
})
expect(deferred.sink.tryAppendItem?.(identity(0), BODY)).toEqual({ accepted: true })
expect(
deferred.sink.tryAppendResolvedItemAndPublish?.(identity(1), BODY, () => identity(1))
).toEqual({ accepted: true })
expect(deferred.sink.tryAppendItem?.(identity(2), BODY)).toEqual({
accepted: false,
reason: 'backpressure'
})
deferred.bind(target(5, log))
await deferred.drained()
expect(log).toEqual([
{ call: 'appendItem', fence: 5, ordinal: 0 },
{ call: 'appendItem', fence: 5, ordinal: 1 },
{ call: 'publish', fence: 5 }
])
})
it('pauses provider reading at the soft byte watermark before rejecting writes', async () => {
const log: Recorded[] = []
const changes: boolean[] = []
@@ -8,6 +8,7 @@ import type { AgentSessionJournal } from '../agent-session-journal/journal-store
import type { JournalLifecycleMutationInput } from '../agent-session-journal/journal-row-builders'
import { estimateStructuredAgentSessionItemBytes } from './structured-agent-session-event-sink-estimate'
import { StructuredAgentSessionSinkQueue } from './structured-agent-session-event-sink-queue'
import { createStructuredAgentSessionResolvedAppend } from './structured-agent-session-resolved-append'
export type StructuredAgentSessionSinkAdmission =
| { accepted: true }
@@ -71,6 +72,13 @@ export type StructuredAgentSessionEventSink = {
resolveIdentity: StructuredAgentSessionIdentityResolver,
options?: StructuredAgentSessionAppendOptions
): StructuredAgentSessionSinkAdmission
/** Queues one resolved append and its publication as a single admitted operation. */
tryAppendResolvedItemAndPublish?(
identitySizeBound: AgentJournalItemIdentity,
body: AgentJournalItemBody,
resolveIdentity: StructuredAgentSessionIdentityResolver,
options?: StructuredAgentSessionAppendOptions
): StructuredAgentSessionSinkAdmission
/** Queues one journal-derived lifecycle append; a null resolution is a no-op. */
tryAppendLifecycleTransition?(
identitySizeBound: AgentJournalItemIdentity,
@@ -152,6 +160,7 @@ export function createDeferredStructuredAgentSessionEventSink(
...(deps.readingControl ? { readingControl: deps.readingControl } : {}),
...(deps.onBackpressureChange ? { onBackpressureChange: deps.onBackpressureChange } : {})
})
const resolvedAppend = createStructuredAgentSessionResolvedAppend(queue)
const appendLifecycleBatch = (
settlementId: string,
@@ -213,28 +222,7 @@ export function createDeferredStructuredAgentSessionEventSink(
},
options
),
tryAppendResolvedItem: (identitySizeBound, body, resolveIdentity, options = {}) => {
const bytes = estimateStructuredAgentSessionItemBytes(identitySizeBound, body)
return queue.submit(
{
bytes,
run: async (bound) => {
const identity = resolveIdentity(bound.journal)
if (identity === null) {
return
}
if (estimateStructuredAgentSessionItemBytes(identity, body) > bytes) {
throw new Error('structured agent-session item identity exceeded its reserved size')
}
await bound.journal.appendItem(identity, body, {
fence: bound.fence,
...(options.observedAt === undefined ? {} : { observedAt: options.observedAt })
})
}
},
options
)
},
...resolvedAppend,
tryAppendLifecycleTransition: (identitySizeBound, body, resolveIdentity) => {
const bytes = estimateStructuredAgentSessionItemBytes(identitySizeBound, body)
return queue.submit(
@@ -0,0 +1,61 @@
import { estimateStructuredAgentSessionItemBytes } from './structured-agent-session-event-sink-estimate'
import type { StructuredAgentSessionEventSink } from './structured-agent-session-event-sink'
import type { StructuredAgentSessionSinkQueue } from './structured-agent-session-event-sink-queue'
/** Resolve a queued item's run identity against the journal bound at execution. */
export function createStructuredAgentSessionResolvedAppend(
queue: StructuredAgentSessionSinkQueue
): {
tryAppendResolvedItem: NonNullable<StructuredAgentSessionEventSink['tryAppendResolvedItem']>
tryAppendResolvedItemAndPublish: NonNullable<
StructuredAgentSessionEventSink['tryAppendResolvedItemAndPublish']
>
} {
return {
tryAppendResolvedItem: (identitySizeBound, body, resolveIdentity, options = {}) => {
const bytes = estimateStructuredAgentSessionItemBytes(identitySizeBound, body)
return queue.submit(
{
bytes,
run: async (bound) => {
const identity = resolveIdentity(bound.journal)
if (identity === null) {
return
}
if (estimateStructuredAgentSessionItemBytes(identity, body) > bytes) {
throw new Error('structured agent-session item identity exceeded its reserved size')
}
await bound.journal.appendItem(identity, body, {
fence: bound.fence,
...(options.observedAt === undefined ? {} : { observedAt: options.observedAt })
})
}
},
options
)
},
tryAppendResolvedItemAndPublish: (identitySizeBound, body, resolveIdentity, options = {}) => {
const bytes = estimateStructuredAgentSessionItemBytes(identitySizeBound, body) + 1
return queue.submit(
{
bytes,
run: async (bound) => {
const identity = resolveIdentity(bound.journal)
if (identity === null) {
return
}
if (estimateStructuredAgentSessionItemBytes(identity, body) + 1 > bytes) {
throw new Error('structured agent-session item identity exceeded its reserved size')
}
await bound.journal.appendItem(identity, body, {
fence: bound.fence,
...(options.observedAt === undefined ? {} : { observedAt: options.observedAt })
})
bound.publish()
}
},
options
)
}
}
}