fix(native-chat): fence command execution after recovery

This commit is contained in:
Brennan Benson
2026-09-16 17:19:49 -07:00
parent e9690a707f
commit 48dd3fba55
2 changed files with 37 additions and 13 deletions
@@ -127,6 +127,32 @@ describe('host conversation command concurrency', () => {
)
})
it('does not start provider work after recovery advances during event flushing', async () => {
await host.hold(HOST_TEST_SESSION, 'conversation-surface')
const flush = Promise.withResolvers<void>()
const flushing = Promise.withResolvers<void>()
vi.spyOn(host, 'flushStreamedEvents').mockImplementationOnce(() => {
flushing.resolve()
return flush.promise
})
const running = host.conversationCommand(CALLER, commandParams('compact'))
await flushing.promise
const fence = store.getRecord(HOST_TEST_SESSION)!.lease.runtimeFence
await host.handleAdapterEvent({
type: 'ended',
sessionId: HOST_TEST_SESSION,
reason: 'provider exited',
cause: 'unexpected-exit',
fence,
acquisitionGeneration: 'generation-1'
})
flush.resolve()
await expect(running).resolves.toMatchObject({ ok: true, value: { state: 'unknown' } })
expect(compact).not.toHaveBeenCalled()
})
it('retires a provider completion that arrives after recovery advances the generation', async () => {
await host.hold(HOST_TEST_SESSION, 'conversation-surface')
const completion = Promise.withResolvers<Record<string, never>>()
@@ -16,7 +16,7 @@ import {
publishConversationCommandLifecycle
} from './structured-conversation-command-lifecycle'
import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations'
import type { StructuredAgentSessionHost } from './structured-agent-session-host'
import type { StructuredAgentSessionHost as Host } from './structured-agent-session-host'
type ExecutionOwner = {
isCurrent: (entry: PendingConversationCommand) => boolean
@@ -24,16 +24,12 @@ type ExecutionOwner = {
settleWaiter: (entry: PendingConversationCommand, result: ConversationCommandResult) => void
report: (entry: PendingConversationCommand, error: unknown) => void
}
const STALE_COMMAND_COMPLETION = new Error('Conversation operation became stale after recovery.')
type ExecutionHost = Pick<Host, 'attach' | 'close' | 'flushStreamedEvents' | 'hasSession'>
export class StructuredConversationCommandExecution {
constructor(
private readonly context: () => StructuredAgentSessionMutationContext,
private readonly host: Pick<
StructuredAgentSessionHost,
'attach' | 'close' | 'flushStreamedEvents' | 'hasSession'
>,
private readonly host: ExecutionHost,
private readonly owner: ExecutionOwner
) {}
@@ -46,7 +42,8 @@ export class StructuredConversationCommandExecution {
try {
await this.publishLifecycle(entry, execution.prepared, 'running')
await this.host.flushStreamedEvents(execution.turn.sessionId)
if (!this.owner.isCurrent(entry)) {
if (!this.canSettle(entry, execution)) {
await this.markUnknown(entry, new Error(CONVERSATION_COMMAND_ABANDONED), true)
return
}
providerMaySettleLate = entry.command === 'compact'
@@ -132,7 +129,8 @@ export class StructuredConversationCommandExecution {
return
}
}
if (!this.owner.isCurrent(entry)) {
if (!this.canSettle(entry, execution)) {
await this.markUnknown(entry, new Error(CONVERSATION_COMMAND_ABANDONED), true)
return
}
if (effectiveOptions) {
@@ -142,7 +140,8 @@ export class StructuredConversationCommandExecution {
}
})
}
if (!this.owner.isCurrent(entry)) {
if (!this.canSettle(entry, execution)) {
await this.markUnknown(entry, new Error(CONVERSATION_COMMAND_ABANDONED), true)
return
}
const replacementSessionId = execution.prepared.replacementSessionId!
@@ -288,9 +287,8 @@ export class StructuredConversationCommandExecution {
entry: PendingConversationCommand,
execution: PreparedConversationCommand
): boolean {
const command = this.context().deps.store.getRecord(
execution.turn.sessionId
)?.conversationCommand
const store = this.context().deps.store
const command = store.getRecord(execution.turn.sessionId)?.conversationCommand
return this.ownsExecution(entry, execution) && command?.phase === 'prepared'
}