mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 00:02:10 +00:00
fix(claude): resolve background task identity after rebind
This commit is contained in:
@@ -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()
|
||||
|
||||
@@ -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<typeof harness>['tra
|
||||
}
|
||||
|
||||
describe('claude journal translation — background task rows', () => {
|
||||
it('resolves a queued restart identity after the sink rebinds', async () => {
|
||||
const persisted = new Map<string, AgentJournalItemBody>()
|
||||
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<string, AgentJournalItemBody>()
|
||||
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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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(
|
||||
|
||||
Reference in New Issue
Block a user