From 1e0843011bbf127a5f7cbccc59d29d67db2f4fd1 Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Mon, 14 Sep 2026 18:35:12 -0700 Subject: [PATCH] fix(claude): resolve background task identity after rebind --- .../claude-background-task-row-journal.ts | 76 ++++++++- ...urnal-translation-background-tasks.test.ts | 154 +++++++++++++++++- .../agent-session-journal/journal-store.ts | 6 +- .../structured-agent-session-event-sink.ts | 36 +++- 4 files changed, 263 insertions(+), 9 deletions(-) diff --git a/src/main/claude/claude-background-task-row-journal.ts b/src/main/claude/claude-background-task-row-journal.ts index ee0c170ce1a..9026abfb023 100644 --- a/src/main/claude/claude-background-task-row-journal.ts +++ b/src/main/claude/claude-background-task-row-journal.ts @@ -3,9 +3,16 @@ import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types' import { backgroundTaskFallbackText } from '../../shared/native-chat-background-task-row' -import type { NativeChatBackgroundTaskBlock } from '../../shared/native-chat-types' -import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' +import { + isBackgroundTaskBlock, + type NativeChatBackgroundTaskBlock +} from '../../shared/native-chat-types' +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 { parseAgentJournalItemKey } from '../../shared/agent-session-journal-item-key' /** Durable identity for one RUN of a task. * @@ -35,6 +42,51 @@ export function claudeBackgroundTaskBody( } } +/** Reconcile one queued row against the durable run identity after a rebind. */ +export function resolveClaudeBackgroundTaskIdentity( + journal: StructuredAgentSessionLifecycleJournal, + id: string, + toolUseId: string | undefined +): AgentJournalItemIdentity { + let maxGeneration = 0 + let matchingGeneration: number | undefined + journal.visitItems((itemId, _sequence, body) => { + const identity = parseAgentJournalItemKey(itemId) + if (!identity || identity.provider !== 'orca') { + return + } + const taskBlock = body.kind === 'message' ? body.blocks.find(isBackgroundTaskBlock) : undefined + if (!taskBlock || taskBlock.taskId !== id) { + return + } + const generation = persistedTaskGeneration(identity.clientMessageId, id) + if (generation === null) { + return + } + maxGeneration = Math.max(maxGeneration, generation) + if (taskBlock.parentToolUseId === toolUseId) { + matchingGeneration = Math.max(matchingGeneration ?? 0, generation) + } + }) + return claudeBackgroundTaskIdentity( + id, + matchingGeneration ?? (maxGeneration === 0 ? 1 : maxGeneration + 1) + ) +} + +function persistedTaskGeneration(clientMessageId: string, taskId: string): number | null { + const base = `claude-background-task:${taskId}` + if (clientMessageId === base) { + return 1 + } + const prefix = `${base}#` + if (!clientMessageId.startsWith(prefix)) { + return null + } + const generation = Number(clientMessageId.slice(prefix.length)) + return Number.isSafeInteger(generation) && generation > 1 ? generation : null +} + export function writeClaudeBackgroundTaskRow( sink: StructuredAgentSessionEventSink, id: string, @@ -51,9 +103,25 @@ export function writeClaudeBackgroundTaskRow( row.lastSerialized = serialized beforeAppend?.() const identity = claudeBackgroundTaskIdentity(id, row.generation) + const resolveIdentity = sink.tryAppendResolvedItem + if (resolveIdentity) { + // 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), + { + // Keyed per RUN, not per task: sharing one key across generations would + // let a restart's queued append evict the finished run's final revision. + coalescingKey: identity.provider === 'orca' ? identity.clientMessageId : `task:${id}` + } + ) + sink.publish() + return + } sink.appendItem(identity, body, { - // Keyed per RUN, not per task: sharing one key across generations would let - // a restart's queued append evict the finished run's final revision. coalescingKey: identity.provider === 'orca' ? identity.clientMessageId : `task:${id}` }) sink.publish() 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 6485bab2e71..5db6240e06a 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 @@ -1,9 +1,16 @@ import { describe, expect, it, vi } from 'vitest' +import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key' import type { AgentJournalItemBody, AgentJournalItemIdentity } from '../../shared/agent-session-journal-types' -import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' +import { + createDeferredStructuredAgentSessionEventSink, + type StructuredAgentSessionEventSink, + type StructuredAgentSessionEventTarget, + type StructuredAgentSessionLifecycleJournal +} from '../native-chat/agent-session-wire/structured-agent-session-event-sink' +import type { AgentSessionJournal } from '../native-chat/agent-session-journal/journal-store' import { createClaudeJournalTranslator } from './claude-structured-journal-translation' // The frames below are the ones the reported session actually carried: two real @@ -112,6 +119,151 @@ function playFailedBackgroundCommand(translator: ReturnType['tra } describe('claude journal translation — background task rows', () => { + it('resolves a queued restart identity after the sink rebinds', async () => { + const persisted = new Map() + const journal = (): AgentSessionJournal => + // 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 + const target = (): StructuredAgentSessionEventTarget => ({ + journal: journal(), + fence: 1, + publish: vi.fn() + }) + const deferred = createDeferredStructuredAgentSessionEventSink() + deferred.bind(target()) + + const first = createClaudeJournalTranslator({ sink: deferred.sink, fallbackIdPrefix: 'first' }) + spawnToolCall(first, 'toolu-first') + first.handle( + systemFrame({ + subtype: 'task_started', + task_id: 'queued-restart', + tool_use_id: 'toolu-first', + task_type: 'local_bash', + is_backgrounded: true + }) + ) + await deferred.drained() + first.dispose() + await deferred.drained() + + const restarted = createDeferredStructuredAgentSessionEventSink() + const second = createClaudeJournalTranslator({ + sink: restarted.sink, + fallbackIdPrefix: 'second' + }) + spawnToolCall(second, 'toolu-second') + second.handle( + systemFrame({ + subtype: 'task_started', + task_id: 'queued-restart', + tool_use_id: 'toolu-second', + task_type: 'local_bash', + is_backgrounded: true + }) + ) + restarted.bind(target()) + await restarted.drained() + + expect([...persisted.keys()].filter((key) => key.includes('queued-restart'))).toEqual([ + 'orca:claude-background-task%3Aqueued-restart', + 'orca:claude-background-task%3Aqueued-restart%232' + ]) + }) + + it('does not overwrite a prior run when a new translator sees a reused task id', () => { + const persisted = new Map() + const journal: StructuredAgentSessionLifecycleJournal = { + epoch: '', + visitItems: (visit) => { + for (const [itemId, body] of persisted) { + visit(itemId, 0, body) + } + } + } + const sink: StructuredAgentSessionEventSink = { + appendItem: (identity, body) => persisted.set(agentJournalItemKey(identity), body), + appendTombstone: vi.fn(), + publish: vi.fn(), + tryAppendResolvedItem: (_identitySizeBound, body, resolveIdentity) => { + const identity = resolveIdentity(journal) + if (identity) { + persisted.set(agentJournalItemKey(identity), body) + } + return { accepted: true } + } + } + const first = createClaudeJournalTranslator({ sink, fallbackIdPrefix: 'first' }) + spawnToolCall(first, 'toolu-first') + first.handle( + systemFrame({ + subtype: 'task_started', + task_id: 'reused-after-reconnect', + tool_use_id: 'toolu-first', + task_type: 'local_bash', + is_backgrounded: true + }) + ) + first.handle( + systemFrame({ + subtype: 'task_notification', + task_id: 'reused-after-reconnect', + tool_use_id: 'toolu-first', + status: 'failed', + summary: 'first run failed' + }) + ) + first.dispose() + + const resumed = createClaudeJournalTranslator({ sink, fallbackIdPrefix: 'resumed' }) + spawnToolCall(resumed, 'toolu-first') + resumed.handle( + systemFrame({ + subtype: 'task_started', + task_id: 'reused-after-reconnect', + tool_use_id: 'toolu-first', + task_type: 'local_bash', + is_backgrounded: true + }) + ) + expect([...persisted.keys()].filter((key) => key.includes('reused-after-reconnect'))).toEqual([ + 'orca:claude-background-task%3Areused-after-reconnect' + ]) + resumed.dispose() + + const second = createClaudeJournalTranslator({ sink, fallbackIdPrefix: 'second' }) + spawnToolCall(second, 'toolu-second') + second.handle( + systemFrame({ + subtype: 'task_started', + task_id: 'reused-after-reconnect', + tool_use_id: 'toolu-second', + task_type: 'local_bash', + is_backgrounded: true + }) + ) + + const rows = [...persisted.entries()].filter(([key]) => key.includes('reused-after-reconnect')) + expect(rows.map(([key]) => key)).toEqual([ + 'orca:claude-background-task%3Areused-after-reconnect', + 'orca:claude-background-task%3Areused-after-reconnect%232' + ]) + }) + it('prints the provider sentence once instead of the opcode twice', () => { const { translator, fallbackRows, taskRowIds, taskRowTexts } = harness() playFailedBackgroundCommand(translator) diff --git a/src/main/native-chat/agent-session-journal/journal-store.ts b/src/main/native-chat/agent-session-journal/journal-store.ts index 2e372f9abae..558af8d5c96 100644 --- a/src/main/native-chat/agent-session-journal/journal-store.ts +++ b/src/main/native-chat/agent-session-journal/journal-store.ts @@ -163,9 +163,11 @@ export class AgentSessionJournal { snapshot = (): AgentJournalSnapshot => renderJournalState(this.state) /** Visits reduced items without allocating and sorting a full snapshot. */ - visitItems = (visit: (itemId: string, sequence: number) => void): void => { + visitItems = ( + visit: (itemId: string, sequence: number, body: AgentJournalItemBody) => void + ): void => { for (const item of this.state.items.values()) { - visit(item.itemId, item.sequence) + visit(item.itemId, item.sequence, item.body) } } 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 857aa118fe8..8215fc0b916 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 @@ -36,10 +36,13 @@ export type StructuredAgentSessionLifecycleJournal = Pick< 'epoch' | 'visitItems' > -export type StructuredAgentSessionLifecycleIdentityResolver = ( +export type StructuredAgentSessionIdentityResolver = ( journal: StructuredAgentSessionLifecycleJournal ) => AgentJournalItemIdentity | null +/** Compatibility alias for lifecycle callers that already use this resolver. */ +export type StructuredAgentSessionLifecycleIdentityResolver = StructuredAgentSessionIdentityResolver + export type StructuredAgentSessionEventSink = { appendItem( identity: AgentJournalItemIdentity, @@ -61,11 +64,18 @@ export type StructuredAgentSessionEventSink = { body: AgentJournalItemBody, options?: StructuredAgentSessionAppendOptions ): StructuredAgentSessionSinkAdmission + /** Queues an ordinary append whose identity is resolved after journal bind. */ + tryAppendResolvedItem?( + 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, body: AgentJournalItemBody, - resolveIdentity: StructuredAgentSessionLifecycleIdentityResolver + resolveIdentity: StructuredAgentSessionIdentityResolver ): StructuredAgentSessionSinkAdmission /** Current durable epoch, when this deferred sink is bound to its journal. */ journalEpoch?(): string | null @@ -203,6 +213,28 @@ 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 + ) + }, tryAppendLifecycleTransition: (identitySizeBound, body, resolveIdentity) => { const bytes = estimateStructuredAgentSessionItemBytes(identitySizeBound, body) return queue.submit(