From f4b0dba3809f43944ebae8abbb0ebb7b1e879433 Mon Sep 17 00:00:00 2001 From: Hypatia May Date: Sat, 11 Jul 2026 10:59:54 +0800 Subject: [PATCH] fix(session-query): contain sync cancellation failures --- .../session-query/src/provider.ts | 14 ++++- .../session-query/tests/session-query.spec.ts | 54 +++++++++++++++++++ 2 files changed, 66 insertions(+), 2 deletions(-) diff --git a/packages/session-query/session-query/src/provider.ts b/packages/session-query/session-query/src/provider.ts index a8b6e3038c..7152e9f5e5 100644 --- a/packages/session-query/session-query/src/provider.ts +++ b/packages/session-query/session-query/src/provider.ts @@ -225,7 +225,12 @@ export class SessionProviderCoordinator { private _syncLive(state: ProviderState, session: Session): Promise { const existing = state.liveSync.get(session.id) if (existing !== undefined) return existing - const snapshot = this._snapshotLive(session) + let snapshot: ReturnType + try { + snapshot = this._snapshotLive(session) + } catch (error: unknown) { + return Promise.reject(this._synchronizationError(state, error)) + } const promise = this._enqueue(state, async () => { /* v8 ignore next -- a provider can be disposed while queued behind an in-flight update */ if (!state.active) return @@ -324,7 +329,12 @@ export class SessionProviderCoordinator { function waitFor(work: Promise, signal: AbortSignal | undefined): Promise { 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((resolve, reject) => { const onAbort = () => { reject(aborted()) } signal.addEventListener('abort', onAbort, { once: true }) diff --git a/packages/session-query/session-query/tests/session-query.spec.ts b/packages/session-query/session-query/tests/session-query.spec.ts index 6d7e9cab33..5249cd22ae 100644 --- a/packages/session-query/session-query/tests/session-query.spec.ts +++ b/packages/session-query/session-query/tests/session-query.spec.ts @@ -697,6 +697,60 @@ describe('provider selection and synchronization', () => { expect(asError(thrown).message).toContain(`provider "${provider.id}"`) 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((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', () => {