fix(session-query): contain sync cancellation failures
This commit is contained in:
@@ -225,7 +225,12 @@ export class SessionProviderCoordinator {
|
|||||||
private _syncLive(state: ProviderState, session: Session): Promise<void> {
|
private _syncLive(state: ProviderState, session: Session): Promise<void> {
|
||||||
const existing = state.liveSync.get(session.id)
|
const existing = state.liveSync.get(session.id)
|
||||||
if (existing !== undefined) return existing
|
if (existing !== undefined) return existing
|
||||||
const snapshot = this._snapshotLive(session)
|
let snapshot: ReturnType<SessionTextExtractors['buildSnapshot']>
|
||||||
|
try {
|
||||||
|
snapshot = this._snapshotLive(session)
|
||||||
|
} catch (error: unknown) {
|
||||||
|
return Promise.reject(this._synchronizationError(state, error))
|
||||||
|
}
|
||||||
const promise = this._enqueue(state, async () => {
|
const promise = this._enqueue(state, async () => {
|
||||||
/* v8 ignore next -- a provider can be disposed while queued behind an in-flight update */
|
/* v8 ignore next -- a provider can be disposed while queued behind an in-flight update */
|
||||||
if (!state.active) return
|
if (!state.active) return
|
||||||
@@ -324,7 +329,12 @@ export class SessionProviderCoordinator {
|
|||||||
|
|
||||||
function waitFor<T>(work: Promise<T>, signal: AbortSignal | undefined): Promise<T> {
|
function waitFor<T>(work: Promise<T>, signal: AbortSignal | undefined): Promise<T> {
|
||||||
if (signal === undefined) return work
|
if (signal === undefined) return work
|
||||||
if (signal.aborted) return Promise.reject(aborted())
|
if (signal.aborted) {
|
||||||
|
// Cancellation supersedes the caller's result, but shared work must still
|
||||||
|
// have a rejection observer when it has already failed synchronously.
|
||||||
|
void work.catch((_supersededError: unknown) => undefined)
|
||||||
|
return Promise.reject(aborted())
|
||||||
|
}
|
||||||
return new Promise<T>((resolve, reject) => {
|
return new Promise<T>((resolve, reject) => {
|
||||||
const onAbort = () => { reject(aborted()) }
|
const onAbort = () => { reject(aborted()) }
|
||||||
signal.addEventListener('abort', onAbort, { once: true })
|
signal.addEventListener('abort', onAbort, { once: true })
|
||||||
|
|||||||
@@ -697,6 +697,60 @@ describe('provider selection and synchronization', () => {
|
|||||||
expect(asError(thrown).message).toContain(`provider "${provider.id}"`)
|
expect(asError(thrown).message).toContain(`provider "${provider.id}"`)
|
||||||
expect(provider.sessionRequests).toEqual([])
|
expect(provider.sessionRequests).toEqual([])
|
||||||
})
|
})
|
||||||
|
|
||||||
|
it('observes synchronous synchronization failure when the caller is already aborted', async () => {
|
||||||
|
const ctx = await liveContext()
|
||||||
|
const session = ctx.sessions.create(SessionId('aborted-throwing-extractor'))
|
||||||
|
session.append('test/note', { note: 'unreachable' })
|
||||||
|
const provider = new FakeProvider()
|
||||||
|
ctx.sessionQuery.registerSearchProvider(provider)
|
||||||
|
ctx.sessionQuery.registerEventTextExtractor('test/note', {
|
||||||
|
version: 'aborted-throwing-v1',
|
||||||
|
extract: () => { throw new Error('superseded extraction failure') },
|
||||||
|
})
|
||||||
|
const controller = new AbortController()
|
||||||
|
controller.abort()
|
||||||
|
const unhandled: unknown[] = []
|
||||||
|
const onUnhandled = (reason: unknown) => { unhandled.push(reason) }
|
||||||
|
process.on('unhandledRejection', onUnhandled)
|
||||||
|
try {
|
||||||
|
await expect(ctx.sessionQuery.searchSessions({ query: 'x' }, { signal: controller.signal }))
|
||||||
|
.rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
|
||||||
|
await new Promise<void>((resolve) => { setImmediate(resolve) })
|
||||||
|
expect(unhandled).toEqual([])
|
||||||
|
} finally {
|
||||||
|
process.off('unhandledRejection', onUnhandled)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
it('types synchronous live-target extraction failures and leaves retries clean', async () => {
|
||||||
|
const ctx = await liveContext()
|
||||||
|
const session = ctx.sessions.create(SessionId('throwing-live-extractor'))
|
||||||
|
session.append('test/note', { note: 'unreachable' })
|
||||||
|
const provider = new FakeProvider()
|
||||||
|
ctx.sessionQuery.registerSearchProvider(provider)
|
||||||
|
const cause = new Error('live extractor failed')
|
||||||
|
const disposeExtractor = ctx.sessionQuery.registerEventTextExtractor('test/note', {
|
||||||
|
version: 'live-throwing-v1',
|
||||||
|
extract: () => { throw cause },
|
||||||
|
})
|
||||||
|
|
||||||
|
let thrown: unknown
|
||||||
|
try {
|
||||||
|
await ctx.sessionQuery.searchEvents({ sessionId: session.id, query: 'x' })
|
||||||
|
} catch (error: unknown) {
|
||||||
|
thrown = error
|
||||||
|
}
|
||||||
|
expect(thrown).toBeInstanceOf(SessionQueryError)
|
||||||
|
expect(thrown).toMatchObject({ code: 'SESSION_QUERY_INDEX_FAILED', cause })
|
||||||
|
expect(asError(thrown).message).toContain(`provider "${provider.id}"`)
|
||||||
|
expect(provider.eventRequests).toEqual([])
|
||||||
|
|
||||||
|
disposeExtractor()
|
||||||
|
await expect(ctx.sessionQuery.searchEvents({ sessionId: session.id, query: 'x' }))
|
||||||
|
.resolves.toMatchObject({ providerId: provider.id })
|
||||||
|
expect(provider.eventRequests).toHaveLength(1)
|
||||||
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
describe('semantic text extractors', () => {
|
describe('semantic text extractors', () => {
|
||||||
|
|||||||
Reference in New Issue
Block a user