diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-event-recovery.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-event-recovery.ts index 5f9894c3280..bf4b2b88d56 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-event-recovery.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-event-recovery.ts @@ -14,6 +14,7 @@ import { export class StructuredAgentSessionEventRecovery { private readonly sinkFailures = new Set() + private readonly pendingSinkRecoveries = new Set>() constructor( private readonly context: { @@ -35,7 +36,7 @@ export class StructuredAgentSessionEventRecovery { return } this.sinkFailures.add(sessionId) - void this.context + const recovery = this.context .serialize(sessionId, async () => { const session = this.context.sessions.get(sessionId) const stop = @@ -61,6 +62,18 @@ export class StructuredAgentSessionEventRecovery { .then((event) => (event ? this.handle(event) : undefined)) .catch((recoveryError) => this.context.onBarrierError(sessionId, recoveryError)) .finally(() => this.sinkFailures.delete(sessionId)) + .then(() => undefined) + this.pendingSinkRecoveries.add(recovery) + void recovery.then( + () => this.pendingSinkRecoveries.delete(recovery), + () => this.pendingSinkRecoveries.delete(recovery) + ) + } + + async drainSinkRecoveries(): Promise { + while (this.pendingSinkRecoveries.size > 0) { + await Promise.allSettled(this.pendingSinkRecoveries) + } } async handle(event: StructuredAgentSessionLifecycleEvent): Promise { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts index aef76c16cdb..f2c005fd7b4 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts @@ -267,7 +267,8 @@ export class StructuredAgentSessionHost { run: () => withTimeout(this.handoffs.drain(), HANDOFF_DRAIN_TIMEOUT_MS, undefined) }, { name: 'drain-attaches', run: () => this.tasks.drainAttaches() }, - { name: 'flush-event-sinks', run: () => this.runtimeState.flushAllEventSinks() } + { name: 'flush-event-sinks', run: () => this.runtimeState.flushAllEventSinks() }, + { name: 'drain-sink-recoveries', run: () => this.eventRecovery.drainSinkRecoveries() } ], sessions: this.sessions }) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-surface-lifetime.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-surface-lifetime.test.ts index 230e39cdaef..d6eaa92176e 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-surface-lifetime.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-surface-lifetime.test.ts @@ -273,6 +273,51 @@ describe('a session evicted and opened again', () => { }) describe('an unexpected provider exit', () => { + it('drains sink-failure lease persistence before closing the host', async () => { + await attach() + const session = ( + host as unknown as { + sessions: Map Promise } }> + } + ).sessions.get(SESSION)! + const settlementStarted = Promise.withResolvers() + const settlementGate = Promise.withResolvers() + const originalTransition = store.transitionHandoff.bind(store) + vi.spyOn(store, 'transitionHandoff').mockImplementationOnce(async (...args) => { + settlementStarted.resolve() + await settlementGate.promise + return originalTransition(...args) + }) + vi.spyOn(session.journal, 'appendItem').mockRejectedValueOnce(new Error('disk unavailable')) + sink?.appendItem( + { provider: 'codex', threadId: THREAD, turnId: 'turn-1', ordinal: 1 }, + { kind: 'message', role: 'assistant', blocks: [{ type: 'text', text: 'lost write' }] } + ) + await settlementStarted.promise + const runtimeState = ( + host as unknown as { + runtimeState: { eventSinkFor: (id: string) => unknown } + } + ).runtimeState + runtimeState.eventSinkFor(SESSION) + let drained = false + const teardown = host.flushAllStreamedEvents().then(() => { + drained = true + }) + try { + // Probe quiescence while persistence is held, matching the handoff teardown contract. + for (let tick = 0; tick < 20; tick += 1) { + await new Promise((resolve) => setTimeout(resolve, 0)) + } + expect(drained).toBe(false) + } finally { + settlementGate.resolve() + await teardown + } + expect(store.getRecord(SESSION)?.lease.claimStatus).toBe('released') + expect(host.hasSession(SESSION)).toBe(false) + }) + it('turns a journal sink failure into observed-exit settlement and lease release', async () => { await attach() const session = (