fix(session-query): address review round 1
This commit is contained in:
@@ -374,9 +374,12 @@ describe('SessionStore', () => {
|
||||
expect(observations).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('contains rejected session/removed listeners during teardown', async () => {
|
||||
it('contains failing session/removed listeners without starving later observers', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
const observed: SessionId[] = []
|
||||
ctx.on('session/removed', () => { throw new Error('synchronous observer failed') })
|
||||
ctx.on('session/removed', header => void observed.push(header.id))
|
||||
ctx.on('session/removed', () => Promise.reject(new Error('observer failed')))
|
||||
const session = ctx.sessions.prepare(SessionId('contained'))
|
||||
const detach = ctx.sessions.enter(session)
|
||||
@@ -385,6 +388,7 @@ describe('SessionStore', () => {
|
||||
await Promise.resolve()
|
||||
await Promise.resolve()
|
||||
expect(ctx.sessions.get(session.id)).toBeUndefined()
|
||||
expect(observed).toEqual([session.id])
|
||||
})
|
||||
|
||||
it('rolls back the session (and onAppend) when a session/created listener throws (P1-1)', async () => {
|
||||
|
||||
@@ -128,11 +128,12 @@ export function runCoordinatorContract(name: string, makeFixture: () => Promise<
|
||||
const fix = await makeFixture()
|
||||
const { ctx, fiber } = await freshCtx(fix)
|
||||
const observed: Array<{ headerId: SessionId; change: SessionPersistedChange }> = []
|
||||
ctx.on('session/persisted', () => { throw new Error('synchronous derived read model failed') })
|
||||
ctx.on('session/persisted', (header, change) => {
|
||||
observed.push({ headerId: header.id, change: structuredClone(change) })
|
||||
header.createdAt = -1
|
||||
return Promise.reject(new Error('derived read model failed'))
|
||||
})
|
||||
ctx.on('session/persisted', () => Promise.reject(new Error('asynchronous derived read model failed')))
|
||||
try {
|
||||
const m = meta('notifications', WORK)
|
||||
await ctx.sessionPersistence.create(m)
|
||||
|
||||
@@ -22,7 +22,7 @@ Session filters cover id, exact cwd, inclusive creation time, parent id/root, an
|
||||
|
||||
## Full-text providers
|
||||
|
||||
`registerSearchProvider(provider)` is effect-scoped and ids are unique. Without `searchProvider`, exactly one locally available provider must be registered; explicit selection fails loudly when the named provider is missing or unavailable. Search pages default to 20 hits and reject limits above 100. Provider scores never cross the public API: event hits carry a plain snippet, while each session hit carries exactly one best matching event.
|
||||
`registerSearchProvider(provider)` is effect-scoped and ids are unique. Without `searchProvider`, exactly one locally available provider must be registered; explicit selection fails loudly when the named provider is missing or unavailable. Search pages default to 20 hits and reject limits above 100; a provider returning more hits than the normalized request limit fails with a typed provider error rather than silently dropping cursor-addressable results. Provider scores never cross the public API: event hits carry a plain snippet, while each session hit carries exactly one best matching event.
|
||||
|
||||
The service feeds providers two independent layers: a durable persisted base (`persistedInventory`, `replacePersisted`, `removePersisted`, `setPersistedActive`) and an ephemeral live override (`replaceLive`, `removeLive`). A search waits for the relevant source state observed before its call: the whole corpus for session search, only the target for a live event search. Failed derived updates do not fail session writes; affected searches receive `SESSION_QUERY_INDEX_FAILED`, and a later search retries the dirty state. `AbortSignal` lets a caller stop waiting and is also passed to provider search.
|
||||
|
||||
|
||||
@@ -12,10 +12,18 @@ interface PersistenceBinding {
|
||||
token: symbol
|
||||
service: SessionPersistence
|
||||
headers: Map<SessionId, SessionHeader>
|
||||
/** Notifications retained until a list that began after them completes. */
|
||||
observations: Map<SessionId, PersistedObservation>
|
||||
observationGeneration: number
|
||||
error?: unknown
|
||||
refreshing: Promise<void> | undefined
|
||||
}
|
||||
|
||||
interface PersistedObservation {
|
||||
generation: number
|
||||
header: SessionHeader
|
||||
}
|
||||
|
||||
/** Active persistence view used by provider reconciliation. */
|
||||
export interface PersistenceView {
|
||||
/** Canonical headers in deterministic creation order. */
|
||||
@@ -132,6 +140,8 @@ export class SessionCorpus {
|
||||
token: Symbol('session-query-persistence'),
|
||||
service,
|
||||
headers: new Map(),
|
||||
observations: new Map(),
|
||||
observationGeneration: 0,
|
||||
refreshing: undefined,
|
||||
}
|
||||
this._persistence = binding
|
||||
@@ -140,7 +150,10 @@ export class SessionCorpus {
|
||||
ctx.on('session/persisted', (header) => {
|
||||
/* v8 ignore next -- a stale notification can race optional-service disposal */
|
||||
if (this._persistence?.token !== binding.token) return
|
||||
binding.headers.set(header.id, structuredClone(header))
|
||||
const snapshot = structuredClone(header)
|
||||
const observation = { generation: ++binding.observationGeneration, header: snapshot }
|
||||
binding.headers.set(header.id, snapshot)
|
||||
binding.observations.set(header.id, observation)
|
||||
this._onPersistenceChange(true)
|
||||
})
|
||||
ctx.effect(() => () => { this._detachPersistence(binding) }, 'sessionQuery.persistenceBinding')
|
||||
@@ -155,10 +168,21 @@ export class SessionCorpus {
|
||||
|
||||
private _refreshPersistence(binding: PersistenceBinding): Promise<void> {
|
||||
if (binding.refreshing !== undefined) return binding.refreshing
|
||||
const startGeneration = binding.observationGeneration
|
||||
const refresh = binding.service.list().then((headers) => {
|
||||
/* v8 ignore next -- a list completion can race optional-service disposal */
|
||||
if (this._persistence?.token !== binding.token) return
|
||||
binding.headers = new Map(headers.map(header => [header.id, structuredClone(header)]))
|
||||
const nextHeaders = new Map(headers.map(header => [header.id, structuredClone(header)]))
|
||||
for (const [id, observation] of binding.observations) {
|
||||
// A notification newer than this list's snapshot is the authoritative
|
||||
// read-your-writes layer; older ones must already be present in list().
|
||||
if (observation.generation > startGeneration) {
|
||||
nextHeaders.set(id, structuredClone(observation.header))
|
||||
} else {
|
||||
binding.observations.delete(id)
|
||||
}
|
||||
}
|
||||
binding.headers = nextHeaders
|
||||
binding.error = undefined
|
||||
this._onPersistenceChange(true)
|
||||
}).catch((error: unknown) => {
|
||||
|
||||
@@ -292,7 +292,10 @@ export class SessionProviderCoordinator {
|
||||
if (page.providerId !== state.provider.id) {
|
||||
throw new SessionQueryError(`session-query provider "${state.provider.id}" returned providerId "${page.providerId}"`, 'SESSION_QUERY_PROVIDER_ERROR')
|
||||
}
|
||||
return page.items.length <= limit ? page : { ...page, items: page.items.slice(0, limit) }
|
||||
if (page.items.length > limit) {
|
||||
throw new SessionQueryError(`session-query provider "${state.provider.id}" returned ${page.items.length} items for limit ${limit}`, 'SESSION_QUERY_PROVIDER_ERROR')
|
||||
}
|
||||
return page
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -118,7 +118,7 @@ export interface SessionSearchHit extends SessionRecord {
|
||||
export interface SessionSearchPage<T> {
|
||||
/** Stable id of the provider that produced this page. */
|
||||
providerId: string
|
||||
/** Ranked hits in deterministic provider order. */
|
||||
/** Ranked hits in deterministic provider order, no longer than the requested limit. */
|
||||
items: readonly T[]
|
||||
/** Opaque next-page cursor, absent when the result is exhausted. */
|
||||
nextCursor?: string
|
||||
|
||||
@@ -52,11 +52,15 @@ class TestPersistence extends SessionPersistence {
|
||||
static entries = new Map<SessionIdType, { meta: SessionHeader; events: SessionEvent[] }>()
|
||||
static listFailure: unknown
|
||||
static loadFailure: unknown
|
||||
static listBarrier: Promise<void> | undefined
|
||||
static onList: (() => void) | undefined
|
||||
|
||||
static reset(entries: readonly { meta: SessionHeader; events: SessionEvent[] }[] = []): void {
|
||||
this.entries = new Map(entries.map(entry => [entry.meta.id, structuredClone(entry)]))
|
||||
this.listFailure = undefined
|
||||
this.loadFailure = undefined
|
||||
this.listBarrier = undefined
|
||||
this.onList = undefined
|
||||
}
|
||||
|
||||
create(meta: SessionHeader): Promise<void> {
|
||||
@@ -80,7 +84,9 @@ class TestPersistence extends SessionPersistence {
|
||||
|
||||
list(): Promise<SessionHeader[]> {
|
||||
if (TestPersistence.listFailure !== undefined) return Promise.reject(asError(TestPersistence.listFailure))
|
||||
return Promise.resolve([...TestPersistence.entries.values()].map(entry => structuredClone(entry.meta)))
|
||||
const snapshot = [...TestPersistence.entries.values()].map(entry => structuredClone(entry.meta))
|
||||
TestPersistence.onList?.()
|
||||
return (TestPersistence.listBarrier ?? Promise.resolve()).then(() => snapshot)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -182,6 +188,12 @@ function asError(value: unknown): Error {
|
||||
return value instanceof Error ? value : new Error(String(value))
|
||||
}
|
||||
|
||||
function deferred(): { promise: Promise<void>; resolve: () => void } {
|
||||
let resolve!: () => void
|
||||
const promise = new Promise<void>((done) => { resolve = done })
|
||||
return { promise, resolve }
|
||||
}
|
||||
|
||||
describe('pure result filters', () => {
|
||||
it('chains session filters as AND while values within one filter are OR', () => {
|
||||
const root: SessionRecord = { header: header('root', 1, { cwd: '/a' }), live: true, persisted: false }
|
||||
@@ -379,8 +391,8 @@ describe('provider selection and synchronization', () => {
|
||||
{ ...record, bestMatch }, { ...record, bestMatch }, { ...record, bestMatch },
|
||||
], nextCursor: 'next' }
|
||||
|
||||
const page = await ctx.sessionQuery.searchSessions({ query: ' hello ', sessionFilters: [{ kind: 'availability', values: ['live'] }] })
|
||||
expect(page.items).toHaveLength(2)
|
||||
await expect(ctx.sessionQuery.searchSessions({ query: ' hello ', sessionFilters: [{ kind: 'availability', values: ['live'] }] }))
|
||||
.rejects.toThrow(expectCode('SESSION_QUERY_PROVIDER_ERROR'))
|
||||
expect(provider.sessionRequests[0]).toMatchObject({ query: 'hello', limit: 2 })
|
||||
expect(provider.live.get(session.id)?.documents[0]?.text).toBe('hello')
|
||||
await expect(ctx.sessionQuery.searchSessions({ query: ' ' })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_QUERY'))
|
||||
@@ -534,6 +546,30 @@ describe('provider selection and synchronization', () => {
|
||||
expect(provider.persisted.has(persisted.id)).toBe(true)
|
||||
})
|
||||
|
||||
it('preserves persisted observations that race an older inventory listing', async () => {
|
||||
TestPersistence.reset()
|
||||
const listStarted = deferred()
|
||||
const releaseList = deferred()
|
||||
TestPersistence.onList = listStarted.resolve
|
||||
TestPersistence.listBarrier = releaseList.promise
|
||||
const ctx = await liveContext()
|
||||
const provider = new FakeProvider()
|
||||
ctx.sessionQuery.registerSearchProvider(provider)
|
||||
await ctx.plugin(TestPersistence)
|
||||
await listStarted.promise
|
||||
|
||||
const announced = header('racing-announcement', 3)
|
||||
TestPersistence.entries.set(announced.id, { meta: announced, events: eventLog('after durable notification') })
|
||||
await ctx.parallel('session/persisted', announced, { kind: 'append', fromSeq: 0, toSeq: 0 })
|
||||
const search = ctx.sessionQuery.searchSessions({ query: 'notification' })
|
||||
releaseList.resolve()
|
||||
|
||||
await expect(search).resolves.toMatchObject({ providerId: provider.id })
|
||||
expect(provider.persisted.get(announced.id)?.documents[0]?.text).toBe('after durable notification')
|
||||
TestPersistence.listBarrier = undefined
|
||||
TestPersistence.onList = undefined
|
||||
})
|
||||
|
||||
it('synchronizes only a live target for event search and retries dirty failures', async () => {
|
||||
const ctx = await liveContext()
|
||||
const session = ctx.sessions.create(SessionId('target'))
|
||||
@@ -646,11 +682,9 @@ describe('semantic text extractors', () => {
|
||||
session.append('user/message', { content: [{ type: 'test/text', value: 'block note' }], source: { kind: 'user' } }, { surfaceOp: 'append' })
|
||||
const provider = new FakeProvider()
|
||||
ctx.sessionQuery.registerSearchProvider(provider)
|
||||
let disposeEvent!: () => void
|
||||
let disposeContent!: () => void
|
||||
const extractorFiber = await ctx.plugin(Object.assign((inner: Context) => {
|
||||
disposeEvent = inner.sessionQuery.registerEventTextExtractor('test/note', { version: 'event-v1', extract: event => [event.data.note] })
|
||||
disposeContent = inner.sessionQuery.registerContentTextExtractor('test/text', { version: 'block-v1', extract: block => [block.value] })
|
||||
inner.sessionQuery.registerEventTextExtractor('test/note', { version: 'event-v1', extract: event => [event.data.note] })
|
||||
inner.sessionQuery.registerContentTextExtractor('test/text', { version: 'block-v1', extract: block => [block.value] })
|
||||
}, { inject: ['sessionQuery'] }))
|
||||
|
||||
await ctx.sessionQuery.searchEvents({ sessionId: session.id, query: 'x' })
|
||||
@@ -663,13 +697,21 @@ describe('semantic text extractors', () => {
|
||||
expect(() => ctx.sessionQuery.registerContentTextExtractor('test/text', { version: ' ', extract: () => [] }))
|
||||
.toThrow(expectCode('SESSION_QUERY_INVALID_EXTRACTOR'))
|
||||
|
||||
disposeEvent()
|
||||
disposeContent()
|
||||
await extractorFiber.dispose()
|
||||
await ctx.sessionQuery.searchEvents({ sessionId: session.id, query: 'x' })
|
||||
const second = provider.live.get(session.id)
|
||||
expect(second?.documents).toEqual([])
|
||||
expect(second?.fingerprint).not.toBe(first?.fingerprint)
|
||||
await extractorFiber.dispose()
|
||||
|
||||
const replacementFiber = await ctx.plugin(Object.assign((inner: Context) => {
|
||||
inner.sessionQuery.registerEventTextExtractor('test/note', { version: 'event-v2', extract: event => [`replacement ${event.data.note}`] })
|
||||
inner.sessionQuery.registerContentTextExtractor('test/text', { version: 'block-v2', extract: block => [`replacement ${block.value}`] })
|
||||
}, { inject: ['sessionQuery'] }))
|
||||
await ctx.sessionQuery.searchEvents({ sessionId: session.id, query: 'x' })
|
||||
const third = provider.live.get(session.id)
|
||||
expect(third?.documents.map(document => document.text)).toEqual(['replacement event note', 'replacement block note'])
|
||||
expect(third?.fingerprint).not.toBe(second?.fingerprint)
|
||||
await replacementFiber.dispose()
|
||||
})
|
||||
})
|
||||
|
||||
|
||||
Reference in New Issue
Block a user