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.
This commit is contained in:
@@ -282,8 +282,13 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
|||||||
* @returns stored header and events before any synthetic recovery closers.
|
* @returns stored header and events before any synthetic recovery closers.
|
||||||
*/
|
*/
|
||||||
inspect(id: SessionId, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
|
inspect(id: SessionId, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
|
||||||
return Promise.resolve(this.retirements.get(id))
|
// Waiting for an in-flight retirement drain must honor cancellation too: a
|
||||||
.then(() => this.serialize(id, () => this.inspectCore(id, signal), signal))
|
// 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(
|
private async inspectCore(
|
||||||
|
|||||||
@@ -50,6 +50,7 @@ interface CoordinatorInternals {
|
|||||||
states: Map<unknown, unknown>
|
states: Map<unknown, unknown>
|
||||||
live: Map<unknown, { pending: unknown[]; flush: Promise<void> | undefined }>
|
live: Map<unknown, { pending: unknown[]; flush: Promise<void> | undefined }>
|
||||||
chains: Map<unknown, unknown>
|
chains: Map<unknown, unknown>
|
||||||
|
retirements: Map<unknown, Promise<void>>
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -451,6 +452,53 @@ describe('PersistenceCoordinator observation cancellation', () => {
|
|||||||
await ctx.fiber.dispose()
|
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<never>
|
||||||
|
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<boolean>()
|
||||||
|
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', () => {
|
describe('PersistenceCoordinator retirement', () => {
|
||||||
|
|||||||
Reference in New Issue
Block a user