diff --git a/src/main/claude/claude-background-task-row-journal.test.ts b/src/main/claude/claude-background-task-row-journal.test.ts new file mode 100644 index 00000000000..ecd23480605 --- /dev/null +++ b/src/main/claude/claude-background-task-row-journal.test.ts @@ -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 + >(() => ({ 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>() + .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() + }) +}) diff --git a/src/main/claude/claude-background-task-row-journal.ts b/src/main/claude/claude-background-task-row-journal.ts index 78b842200f3..a4779898204 100644 --- a/src/main/claude/claude-background-task-row-journal.ts +++ b/src/main/claude/claude-background-task-row-journal.ts @@ -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() + + 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 } diff --git a/src/main/claude/claude-background-task-row-writer.ts b/src/main/claude/claude-background-task-row-writer.ts new file mode 100644 index 00000000000..bb973a1ae70 --- /dev/null +++ b/src/main/claude/claude-background-task-row-writer.ts @@ -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() + } +} diff --git a/src/main/claude/claude-background-task-rows.ts b/src/main/claude/claude-background-task-rows.ts index be0a03469b1..64829d0493a 100644 --- a/src/main/claude/claude-background-task-rows.ts +++ b/src/main/claude/claude-background-task-rows.ts @@ -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, observedAt: number) => void + onPersistenceFailure?: (error: Error) => void now?: () => number } export class ClaudeBackgroundTaskRows { private readonly rows = new Map() + 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): 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): 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): 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 + ) } } diff --git a/src/main/claude/claude-background-task-start.ts b/src/main/claude/claude-background-task-start.ts new file mode 100644 index 00000000000..33bcc3a4a60 --- /dev/null +++ b/src/main/claude/claude-background-task-start.ts @@ -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 + rows: Map + ledgers: ClaudeBackgroundTaskLedgers + isForwardedParentTool: (toolUseId: string) => boolean + openRow: (id: string, message: Record) => 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 +} diff --git a/src/main/claude/claude-stream-json-connection-close.test.ts b/src/main/claude/claude-stream-json-connection-close.test.ts index e824139a776..22076f3600e 100644 --- a/src/main/claude/claude-stream-json-connection-close.test.ts +++ b/src/main/claude/claude-stream-json-connection-close.test.ts @@ -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>() + const next = vi + .fn<() => Promise>>>() + .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[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> + 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>() + const next = vi + .fn<() => Promise>>>() + .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[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>>(() => {})) + // 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[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() mocks.refresh.mockReturnValueOnce(refreshDone.promise) diff --git a/src/main/claude/claude-stream-json-connection.ts b/src/main/claude/claude-stream-json-connection.ts index 1a99bb92064..89255094e59 100644 --- a/src/main/claude/claude-stream-json-connection.ts +++ b/src/main/claude/claude-stream-json-connection.ts @@ -38,6 +38,10 @@ function loadClaudeAgentSdk(): Promise { return claudeAgentSdk } +function isRecord(value: unknown): value is Record { + 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) => Promise /** Resolves true after processless settlement, or root exit plus observed tree exit. */ close: () => Promise @@ -141,6 +147,23 @@ export async function openClaudeStreamJsonConnection( let faultReported = false let exitReported = false let closePromise: Promise | null = null + let readingBarrier: Promise | null = null + let releaseReadingBarrier: (() => void) | null = null + const pauseReading = (): void => { + if (closing || exited || terminalError || readingBarrier) { + return + } + readingBarrier = new Promise((resolve) => { + releaseReadingBarrier = resolve + }) + } + const resumeReading = (): void => { + const release = releaseReadingBarrier + readingBarrier = null + releaseReadingBarrier = null + release?.() + } + const waitUntilReadable = (): Promise => 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) + 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 => { 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 } diff --git a/src/main/claude/claude-structured-journal-translation-background-tasks.test.ts b/src/main/claude/claude-structured-journal-translation-background-tasks.test.ts index a4d542bf0f9..aef4c3ece0a 100644 --- a/src/main/claude/claude-structured-journal-translation-background-tasks.test.ts +++ b/src/main/claude/claude-structured-journal-translation-background-tasks.test.ts @@ -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() + const appendEntered = Promise.withResolvers() + const appendGate = Promise.withResolvers() + 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() const deferred = createDeferredStructuredAgentSessionEventSink() diff --git a/src/main/claude/claude-structured-journal-translation.ts b/src/main/claude/claude-structured-journal-translation.ts index d267a4eeffe..112f4319b17 100644 --- a/src/main/claude/claude-structured-journal-translation.ts +++ b/src/main/claude/claude-structured-journal-translation.ts @@ -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 }, diff --git a/src/main/claude/claude-structured-session-acquisition.ts b/src/main/claude/claude-structured-session-acquisition.ts index 622b7618926..a10a2dca671 100644 --- a/src/main/claude/claude-structured-session-acquisition.ts +++ b/src/main/claude/claude-structured-session-acquisition.ts @@ -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): 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, diff --git a/src/main/claude/claude-structured-session-adapter.ts b/src/main/claude/claude-structured-session-adapter.ts index 9357e9635f6..a719745aa50 100644 --- a/src/main/claude/claude-structured-session-adapter.ts +++ b/src/main/claude/claude-structured-session-adapter.ts @@ -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 => { - const awaited = new Set>() - for (;;) { - const pending = [...this.exits.values()] - .map((exit) => exit.publication) - .filter( - (publication): publication is Promise => - 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 => drainClaudeObservedExits(this.exits) /** Lifecycle recovery is published only after the child tree proof is true. */ private settleUnexpectedExit(sessionId: string, exit: ClaudeSessionExit): Promise { 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 { - 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' diff --git a/src/main/claude/claude-structured-session-close.ts b/src/main/claude/claude-structured-session-close.ts index 431d6e38ab8..f6a58c36357 100644 --- a/src/main/claude/claude-structured-session-close.ts +++ b/src/main/claude/claude-structured-session-close.ts @@ -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') diff --git a/src/main/claude/claude-structured-session-exit-lifecycle.ts b/src/main/claude/claude-structured-session-exit-lifecycle.ts new file mode 100644 index 00000000000..20cbc1f3125 --- /dev/null +++ b/src/main/claude/claude-structured-session-exit-lifecycle.ts @@ -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 +): Promise { + const awaited = new Set>() + for (;;) { + const pending = [...exits.values()] + .map((exit) => exit.publication) + .filter( + (publication): publication is Promise => + 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 +): Promise { + 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 + }) +} diff --git a/src/main/claude/claude-structured-session-journal-control.ts b/src/main/claude/claude-structured-session-journal-control.ts new file mode 100644 index 00000000000..dc8a5c9d262 --- /dev/null +++ b/src/main/claude/claude-structured-session-journal-control.ts @@ -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 + 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?.() + } + } + }) +} diff --git a/src/main/claude/claude-structured-session-publication.ts b/src/main/claude/claude-structured-session-publication.ts index 395335332e7..2434f522661 100644 --- a/src/main/claude/claude-structured-session-publication.ts +++ b/src/main/claude/claude-structured-session-publication.ts @@ -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 } : {}) } } } diff --git a/src/main/claude/claude-structured-session-reading-control.test.ts b/src/main/claude/claude-structured-session-reading-control.test.ts new file mode 100644 index 00000000000..7113abd5bf5 --- /dev/null +++ b/src/main/claude/claude-structured-session-reading-control.test.ts @@ -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 +} { + 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 +): 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() + 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() + const appendGate = Promise.withResolvers() + 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) + }) +}) diff --git a/src/main/claude/claude-structured-session-state.ts b/src/main/claude/claude-structured-session-state.ts index 055a477c74c..de697192074 100644 --- a/src/main/claude/claude-structured-session-state.ts +++ b/src/main/claude/claude-structured-session-state.ts @@ -170,6 +170,7 @@ export type ClaudeSession = { closeEnded?: boolean translator: ClaudeJournalTranslator | null events: StructuredAgentSessionEventSink | undefined + unbindReadingControl?: () => void } export function mintClaudeAcquisitionGeneration(deps: ClaudeStructuredSessionAdapterDeps): string { diff --git a/src/main/claude/claude-structured-session-test-support.ts b/src/main/claude/claude-structured-session-test-support.ts index 352a7ed18ec..50b8e766ce4 100644 --- a/src/main/claude/claude-structured-session-test-support.ts +++ b/src/main/claude/claude-structured-session-test-support.ts @@ -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) { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.test.ts index 3d0a5e8a269..8a2485473f8 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.test.ts @@ -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[] = [] diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts index 8215fc0b916..b27e466d94b 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts @@ -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( diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-resolved-append.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-resolved-append.ts new file mode 100644 index 00000000000..b8fb2f1986c --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-resolved-append.ts @@ -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 + 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 + ) + } + } +}