From 63c8494d69b55e32bd5b20f65128b361b63737d4 Mon Sep 17 00:00:00 2001 From: _Kerman Date: Mon, 27 Jul 2026 02:00:26 +0800 Subject: [PATCH] fix(session-persistence): honor inspect cancellation during retirement drain inspect() awaited an in-flight retirement drain unconditionally before entering serialize(), so a slow drain pinned a cancelled inspect until it finished, past the documented cancellation boundary. Race the retirement wait against the signal with observeQueuedAbort, matching serialize()'s queued-read behaviour. Adds a regression that cancels an inspect while a gated retirement is pending and asserts prompt rejection without a backend read. --- .../session-persistence/src/coordinator.ts | 9 +++- .../tests/persistence.spec.ts | 48 +++++++++++++++++++ 2 files changed, 55 insertions(+), 2 deletions(-) diff --git a/packages/session-persistence/session-persistence/src/coordinator.ts b/packages/session-persistence/session-persistence/src/coordinator.ts index cc23b71f37..c55f360f0a 100644 --- a/packages/session-persistence/session-persistence/src/coordinator.ts +++ b/packages/session-persistence/session-persistence/src/coordinator.ts @@ -282,8 +282,13 @@ export class PersistenceCoordinator { * @returns stored header and events before any synthetic recovery closers. */ inspect(id: SessionId, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> { - return Promise.resolve(this.retirements.get(id)) - .then(() => this.serialize(id, () => this.inspectCore(id, signal), signal)) + // Waiting for an in-flight retirement drain must honor cancellation too: a + // slow drain would otherwise pin a cancelled inspect until it finishes, + // past the documented boundary. serialize() already races the signal for + // the queued read; do the same for the retirement wait. + const retired = Promise.resolve(this.retirements.get(id)) + const waited = signal === undefined ? retired : observeQueuedAbort(retired, signal, () => false) + return waited.then(() => this.serialize(id, () => this.inspectCore(id, signal), signal)) } private async inspectCore( diff --git a/packages/session-persistence/session-persistence/tests/persistence.spec.ts b/packages/session-persistence/session-persistence/tests/persistence.spec.ts index eedd1b0156..2bb49948fb 100644 --- a/packages/session-persistence/session-persistence/tests/persistence.spec.ts +++ b/packages/session-persistence/session-persistence/tests/persistence.spec.ts @@ -50,6 +50,7 @@ interface CoordinatorInternals { states: Map live: Map | undefined }> chains: Map + retirements: Map> } /** @@ -451,6 +452,53 @@ describe('PersistenceCoordinator observation cancellation', () => { await ctx.fiber.dispose() } }) + + it('rejects a cancelled inspect while an in-flight retirement drain is still pending', async () => { + const ctx = new Context() + await ctx.plugin(SessionStore) + const backend = new ControlledBackend() + let coordinator!: PersistenceCoordinator + const backendFiber = await ctx.plugin(Object.assign((inner: Context) => { + coordinator = new PersistenceCoordinator(inner, backend) + }, { inject: ['sessions'] })) + const internals = coordinator as unknown as CoordinatorInternals + const appendGate = Promise.withResolvers() + backend.beforeAppend = async () => { await appendGate.promise } + + try { + const id = SessionId('retiring-inspect') + let session!: Session + const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => { + session = inner.sessions.create(id) + }, { inject: ['sessions'] })) + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + session.append('turn/end', { turn: 1, reason: { kind: 'completed' } }) + // Dispose the session so retirement starts; its append is gated, so the + // retirement promise stays pending in the coordinator. + await sessionFiber.dispose() + await vi.waitFor(() => { expect(internals.retirements.has(id)).toBe(true) }) + + const controller = new AbortController() + const reason = new Error('inspect cancelled during retirement') + const pending = coordinator.inspect(id, controller.signal) + let observedReason: unknown + const observed = pending.catch((error: unknown) => { observedReason = error }) + + // Cancel before the gated retirement can settle: the inspect must reject + // promptly instead of waiting for the drain, and must never reach the + // backend read. + controller.abort(reason) + await vi.waitFor(() => { expect(observedReason).toBe(reason) }) + expect(backend.loadAttempts).toBe(0) + + appendGate.resolve(true) + await observed + } finally { + appendGate.resolve(true) + await backendFiber.dispose() + await ctx.fiber.dispose() + } + }) }) describe('PersistenceCoordinator retirement', () => {