Combine pending PR #18974 for deflake CI validation

This commit is contained in:
Neil
2026-09-05 19:03:33 -07:00
3 changed files with 61 additions and 2 deletions
@@ -14,6 +14,7 @@ import {
export class StructuredAgentSessionEventRecovery {
private readonly sinkFailures = new Set<string>()
private readonly pendingSinkRecoveries = new Set<Promise<void>>()
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<void> {
while (this.pendingSinkRecoveries.size > 0) {
await Promise.allSettled(this.pendingSinkRecoveries)
}
}
async handle(event: StructuredAgentSessionLifecycleEvent): Promise<void> {
@@ -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
})
@@ -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<string, { journal: { appendItem: (...args: never[]) => Promise<unknown> } }>
}
).sessions.get(SESSION)!
const settlementStarted = Promise.withResolvers<void>()
const settlementGate = Promise.withResolvers<void>()
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 = (