refactor(session-query): simplify persistence binding (round 1)
This commit is contained in:
@@ -97,9 +97,12 @@ interface ObservedPersistedSession {
|
|||||||
loaded?: ObservedSession
|
loaded?: ObservedSession
|
||||||
}
|
}
|
||||||
|
|
||||||
|
interface PersistenceBinding {
|
||||||
|
readonly service?: SessionPersistence
|
||||||
|
}
|
||||||
|
|
||||||
interface Observation {
|
interface Observation {
|
||||||
persistence: SessionPersistence | undefined
|
persistenceBinding: PersistenceBinding
|
||||||
persistenceRevision: number
|
|
||||||
persisted: Map<SessionId, ObservedPersistedSession>
|
persisted: Map<SessionId, ObservedPersistedSession>
|
||||||
live: Map<SessionId, ObservedSession>
|
live: Map<SessionId, ObservedSession>
|
||||||
}
|
}
|
||||||
@@ -161,10 +164,8 @@ export class SessionSearchSqlite extends SessionSearchService {
|
|||||||
private readonly _instance = randomUUID()
|
private readonly _instance = randomUUID()
|
||||||
private readonly _ready: Promise<void>
|
private readonly _ready: Promise<void>
|
||||||
private _db: DatabaseSync | undefined
|
private _db: DatabaseSync | undefined
|
||||||
private _persistence: SessionPersistence | undefined
|
private _persistenceBinding: PersistenceBinding = {}
|
||||||
private _persistenceBinding: object | undefined
|
private _lastPersistenceBinding: PersistenceBinding | undefined
|
||||||
private _persistenceRevision = 0
|
|
||||||
private _lastPersistenceRevision: number | undefined
|
|
||||||
private _persistenceEpoch = 0
|
private _persistenceEpoch = 0
|
||||||
private _globalGeneration = 0
|
private _globalGeneration = 0
|
||||||
private _localGeneration = 0
|
private _localGeneration = 0
|
||||||
@@ -182,16 +183,12 @@ export class SessionSearchSqlite extends SessionSearchService {
|
|||||||
void this._ready.catch(() => undefined)
|
void this._ready.catch(() => undefined)
|
||||||
this._optionalPersistenceFiber = ctx.inject(['sessionPersistence'], (childCtx: Context) => {
|
this._optionalPersistenceFiber = ctx.inject(['sessionPersistence'], (childCtx: Context) => {
|
||||||
const service = childCtx.sessionPersistence
|
const service = childCtx.sessionPersistence
|
||||||
const binding = {}
|
const binding = { service }
|
||||||
this._persistenceBinding = binding
|
this._persistenceBinding = binding
|
||||||
this._persistence = service
|
|
||||||
this._persistenceRevision += 1
|
|
||||||
childCtx.effect(() => () => {
|
childCtx.effect(() => () => {
|
||||||
/* v8 ignore next -- a stale optional-service disposer cannot clear a replacement */
|
/* v8 ignore next -- a stale optional-service disposer cannot clear a replacement */
|
||||||
if (this._persistenceBinding !== binding) return
|
if (this._persistenceBinding !== binding) return
|
||||||
this._persistenceBinding = undefined
|
this._persistenceBinding = {}
|
||||||
this._persistence = undefined
|
|
||||||
this._persistenceRevision += 1
|
|
||||||
}, 'sessionSearchSqlite.persistenceBinding')
|
}, 'sessionSearchSqlite.persistenceBinding')
|
||||||
})
|
})
|
||||||
ctx.effect(() => {
|
ctx.effect(() => {
|
||||||
@@ -330,16 +327,16 @@ export class SessionSearchSqlite extends SessionSearchService {
|
|||||||
const liveById = new Map(liveRows.map(row => [row.id as SessionId, row]))
|
const liveById = new Map(liveRows.map(row => [row.id as SessionId, row]))
|
||||||
const observation = await this._observeStable(persistedById, signal)
|
const observation = await this._observeStable(persistedById, signal)
|
||||||
assertNotAborted(signal)
|
assertNotAborted(signal)
|
||||||
const persistentChanges = observation.persistence === undefined
|
const persistentChanges = observation.persistenceBinding.service === undefined
|
||||||
? []
|
? []
|
||||||
: [...observation.persisted.values()].filter(entry => entry.loaded !== undefined)
|
: [...observation.persisted.values()].filter(entry => entry.loaded !== undefined)
|
||||||
const persistentDeletes = observation.persistence === undefined
|
const persistentDeletes = observation.persistenceBinding.service === undefined
|
||||||
? []
|
? []
|
||||||
: persistedRows.filter(row => !observation.persisted.has(row.id as SessionId))
|
: persistedRows.filter(row => !observation.persisted.has(row.id as SessionId))
|
||||||
const liveChanges = [...observation.live.values()].filter(entry => liveById.get(entry.header.id)?.fingerprint !== entry.fingerprint)
|
const liveChanges = [...observation.live.values()].filter(entry => liveById.get(entry.header.id)?.fingerprint !== entry.fingerprint)
|
||||||
const liveDeletes = liveRows.filter(row => !observation.live.has(row.id as SessionId))
|
const liveDeletes = liveRows.filter(row => !observation.live.has(row.id as SessionId))
|
||||||
const pointerChanged = this._lastPersistenceRevision !== undefined
|
const pointerChanged = this._lastPersistenceBinding !== undefined
|
||||||
&& this._lastPersistenceRevision !== observation.persistenceRevision
|
&& this._lastPersistenceBinding !== observation.persistenceBinding
|
||||||
const hasWrites = persistentChanges.length > 0
|
const hasWrites = persistentChanges.length > 0
|
||||||
|| persistentDeletes.length > 0
|
|| persistentDeletes.length > 0
|
||||||
|| liveChanges.length > 0
|
|| liveChanges.length > 0
|
||||||
@@ -393,7 +390,7 @@ export class SessionSearchSqlite extends SessionSearchService {
|
|||||||
if (hasWrites || pointerChanged) this._globalGeneration += 1
|
if (hasWrites || pointerChanged) this._globalGeneration += 1
|
||||||
if (pointerChanged) this._persistenceEpoch += 1
|
if (pointerChanged) this._persistenceEpoch += 1
|
||||||
this._localGeneration = nextLocalGeneration
|
this._localGeneration = nextLocalGeneration
|
||||||
this._lastPersistenceRevision = observation.persistenceRevision
|
this._lastPersistenceBinding = observation.persistenceBinding
|
||||||
}
|
}
|
||||||
|
|
||||||
private async _observeStable(
|
private async _observeStable(
|
||||||
@@ -402,13 +399,13 @@ export class SessionSearchSqlite extends SessionSearchService {
|
|||||||
): Promise<Observation> {
|
): Promise<Observation> {
|
||||||
for (;;) {
|
for (;;) {
|
||||||
assertNotAborted(signal)
|
assertNotAborted(signal)
|
||||||
const persistence = this._persistence
|
const persistenceBinding = this._persistenceBinding
|
||||||
const persistenceRevision = this._persistenceRevision
|
const persistence = persistenceBinding.service
|
||||||
let persisted = new Map<SessionId, ObservedPersistedSession>()
|
let persisted = new Map<SessionId, ObservedPersistedSession>()
|
||||||
if (persistence !== undefined) {
|
if (persistence !== undefined) {
|
||||||
try {
|
try {
|
||||||
const canReuseIndexed = this._lastPersistenceRevision === undefined
|
const canReuseIndexed = this._lastPersistenceBinding === undefined
|
||||||
|| this._lastPersistenceRevision === persistenceRevision
|
|| this._lastPersistenceBinding === persistenceBinding
|
||||||
const before = await waitWithAbort(persistence.listSnapshots(), signal)
|
const before = await waitWithAbort(persistence.listSnapshots(), signal)
|
||||||
persisted = materializePersistenceSnapshots(before)
|
persisted = materializePersistenceSnapshots(before)
|
||||||
for (const entry of persisted.values()) {
|
for (const entry of persisted.values()) {
|
||||||
@@ -421,14 +418,14 @@ export class SessionSearchSqlite extends SessionSearchService {
|
|||||||
await waitWithAbort(persistence.listSnapshots(), signal),
|
await waitWithAbort(persistence.listSnapshots(), signal),
|
||||||
)
|
)
|
||||||
if (!samePersistenceSnapshots(persisted, after)) continue
|
if (!samePersistenceSnapshots(persisted, after)) continue
|
||||||
if (this._persistenceRevision !== persistenceRevision) continue
|
if (this._persistenceBinding !== persistenceBinding) continue
|
||||||
} catch (error: unknown) {
|
} catch (error: unknown) {
|
||||||
if (isAbort(error) || signal?.aborted) {
|
if (isAbort(error) || signal?.aborted) {
|
||||||
throw new SessionQueryError('session-search aborted', 'SESSION_QUERY_ABORTED', {
|
throw new SessionQueryError('session-search aborted', 'SESSION_QUERY_ABORTED', {
|
||||||
cause: error,
|
cause: error,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
if (this._persistenceRevision !== persistenceRevision) continue
|
if (this._persistenceBinding !== persistenceBinding) continue
|
||||||
if (error instanceof SessionQueryError) throw error
|
if (error instanceof SessionQueryError) throw error
|
||||||
throw new SessionQueryError(
|
throw new SessionQueryError(
|
||||||
`session-search persistence observation failed: ${errorMessage(error)}`,
|
`session-search persistence observation failed: ${errorMessage(error)}`,
|
||||||
@@ -444,8 +441,8 @@ export class SessionSearchSqlite extends SessionSearchService {
|
|||||||
if (durable !== undefined) assertSessionHeadersCompatible(observed.header, durable.header)
|
if (durable !== undefined) assertSessionHeadersCompatible(observed.header, durable.header)
|
||||||
live.set(session.id, observed)
|
live.set(session.id, observed)
|
||||||
}
|
}
|
||||||
if (this._persistenceRevision === persistenceRevision) {
|
if (this._persistenceBinding === persistenceBinding) {
|
||||||
return { persistence, persistenceRevision, persisted, live }
|
return { persistenceBinding, persisted, live }
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -564,7 +561,7 @@ export class SessionSearchSqlite extends SessionSearchService {
|
|||||||
ORDER BY match_count DESC, document_length ASC, time DESC, session_id ASC, seq DESC
|
ORDER BY match_count DESC, document_length ASC, time DESC, session_id ASC, seq DESC
|
||||||
LIMIT ? OFFSET ?
|
LIMIT ? OFFSET ?
|
||||||
`).all(
|
`).all(
|
||||||
...selectedDocumentsParams(request.query, this._persistence !== undefined),
|
...selectedDocumentsParams(request.query, this._persistenceBinding.service !== undefined),
|
||||||
...sessionWhere.params,
|
...sessionWhere.params,
|
||||||
...eventWhere.params,
|
...eventWhere.params,
|
||||||
request.limit + 1,
|
request.limit + 1,
|
||||||
@@ -583,7 +580,7 @@ export class SessionSearchSqlite extends SessionSearchService {
|
|||||||
ORDER BY match_count DESC, document_length ASC, time DESC, seq DESC
|
ORDER BY match_count DESC, document_length ASC, time DESC, seq DESC
|
||||||
LIMIT ? OFFSET ?
|
LIMIT ? OFFSET ?
|
||||||
`).all(
|
`).all(
|
||||||
...selectedDocumentsParams(request.query, this._persistence !== undefined),
|
...selectedDocumentsParams(request.query, this._persistenceBinding.service !== undefined),
|
||||||
request.sessionId,
|
request.sessionId,
|
||||||
...eventWhere.params,
|
...eventWhere.params,
|
||||||
request.limit + 1,
|
request.limit + 1,
|
||||||
@@ -597,7 +594,7 @@ export class SessionSearchSqlite extends SessionSearchService {
|
|||||||
'SELECT generation FROM temp.live_sessions WHERE id = ?',
|
'SELECT generation FROM temp.live_sessions WHERE id = ?',
|
||||||
).get(sessionId) as { generation: number } | undefined
|
).get(sessionId) as { generation: number } | undefined
|
||||||
if (live !== undefined) return `live:${live.generation}`
|
if (live !== undefined) return `live:${live.generation}`
|
||||||
if (this._persistence !== undefined) {
|
if (this._persistenceBinding.service !== undefined) {
|
||||||
const persisted = db.prepare(
|
const persisted = db.prepare(
|
||||||
'SELECT generation FROM persisted_sessions WHERE id = ?',
|
'SELECT generation FROM persisted_sessions WHERE id = ?',
|
||||||
).get(sessionId) as { generation: number } | undefined
|
).get(sessionId) as { generation: number } | undefined
|
||||||
@@ -713,10 +710,7 @@ function selectedDocumentsParams(query: string, persistenceVisible: boolean): Ar
|
|||||||
}
|
}
|
||||||
|
|
||||||
function observeLive(session: Session): ObservedSession {
|
function observeLive(session: Session): ObservedSession {
|
||||||
return observeSession(
|
return observeSession(session.header, session.events)
|
||||||
structuredClone(session.header),
|
|
||||||
session.events.map(event => structuredClone(event)),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
function observeSession(header: SessionHeader, events: readonly SessionEvent[]): ObservedSession {
|
function observeSession(header: SessionHeader, events: readonly SessionEvent[]): ObservedSession {
|
||||||
|
|||||||
@@ -326,6 +326,24 @@ describe('SQLite session search', () => {
|
|||||||
})).rejects.toThrow(expectCode('SESSION_QUERY_STALE_CURSOR'))
|
})).rejects.toThrow(expectCode('SESSION_QUERY_STALE_CURSOR'))
|
||||||
})
|
})
|
||||||
|
|
||||||
|
it('invalidates session cursors after transient persistence topology changes', async () => {
|
||||||
|
TestPersistence.reset()
|
||||||
|
const ctx = await liveContext({ path: ':memory:', defaultLimit: 1, maxLimit: 5 })
|
||||||
|
ctx.sessions.create(SessionId('first'), { seed: messageEvents('needle first') })
|
||||||
|
ctx.sessions.create(SessionId('second'), { seed: messageEvents('needle second') })
|
||||||
|
const page = await ctx.sessionSearch.searchSessions({ query: 'needle', limit: 1 })
|
||||||
|
if (page.nextCursor === undefined) throw new Error('expected cursor')
|
||||||
|
|
||||||
|
const persistence = await ctx.plugin(TestPersistence)
|
||||||
|
await persistence.dispose()
|
||||||
|
|
||||||
|
await expect(ctx.sessionSearch.searchSessions({
|
||||||
|
query: 'needle',
|
||||||
|
limit: 1,
|
||||||
|
cursor: page.nextCursor,
|
||||||
|
})).rejects.toThrow(expectCode('SESSION_QUERY_STALE_CURSOR'))
|
||||||
|
})
|
||||||
|
|
||||||
it('rejects invalid requests, filters, cursors, and direct config', async () => {
|
it('rejects invalid requests, filters, cursors, and direct config', async () => {
|
||||||
const ctx = await liveContext({ path: ':memory:', defaultLimit: 2, maxLimit: 3 })
|
const ctx = await liveContext({ path: ':memory:', defaultLimit: 2, maxLimit: 3 })
|
||||||
const session = ctx.sessions.create(SessionId('valid'), { seed: messageEvents('needle') })
|
const session = ctx.sessions.create(SessionId('valid'), { seed: messageEvents('needle') })
|
||||||
@@ -503,11 +521,6 @@ describe('SQLite reconciliation and source lifecycle', () => {
|
|||||||
TestPersistence.set({ meta: durable, events: messageEvents('new needle') })
|
TestPersistence.set({ meta: durable, events: messageEvents('new needle') })
|
||||||
TestPersistence.revisions.set(durable.id, revision)
|
TestPersistence.revisions.set(durable.id, revision)
|
||||||
const replacement = await ctx.plugin(TestPersistence)
|
const replacement = await ctx.plugin(TestPersistence)
|
||||||
const internals = ctx.sessionSearch as unknown as {
|
|
||||||
_lastPersistenceRevision: number
|
|
||||||
_persistenceRevision: number
|
|
||||||
}
|
|
||||||
expect(internals._persistenceRevision).not.toBe(internals._lastPersistenceRevision)
|
|
||||||
const page = await ctx.sessionSearch.searchSessions({ query: 'new needle' })
|
const page = await ctx.sessionSearch.searchSessions({ query: 'new needle' })
|
||||||
expect(TestPersistence.loads.get(durable.id)).toBe(2)
|
expect(TestPersistence.loads.get(durable.id)).toBe(2)
|
||||||
expect(page).toMatchObject({ items: [{ header: durable }] })
|
expect(page).toMatchObject({ items: [{ header: durable }] })
|
||||||
@@ -553,13 +566,15 @@ describe('SQLite reconciliation and source lifecycle', () => {
|
|||||||
TestPersistence.reset([{ meta: durable, events: messageEvents('durable needle') }])
|
TestPersistence.reset([{ meta: durable, events: messageEvents('durable needle') }])
|
||||||
const ctx = await liveContext()
|
const ctx = await liveContext()
|
||||||
await ctx.plugin(TestPersistence)
|
await ctx.plugin(TestPersistence)
|
||||||
const internals = ctx.sessionSearch as unknown as { _persistenceRevision: number }
|
const internals = ctx.sessionSearch as unknown as {
|
||||||
|
_persistenceBinding: { service?: SessionPersistence }
|
||||||
|
}
|
||||||
const originalList = ctx.sessions.list.bind(ctx.sessions)
|
const originalList = ctx.sessions.list.bind(ctx.sessions)
|
||||||
let bumped = false
|
let bumped = false
|
||||||
const list = vi.spyOn(ctx.sessions, 'list').mockImplementation(() => {
|
const list = vi.spyOn(ctx.sessions, 'list').mockImplementation(() => {
|
||||||
if (!bumped) {
|
if (!bumped) {
|
||||||
bumped = true
|
bumped = true
|
||||||
internals._persistenceRevision += 1
|
internals._persistenceBinding = { ...internals._persistenceBinding }
|
||||||
}
|
}
|
||||||
return originalList()
|
return originalList()
|
||||||
})
|
})
|
||||||
|
|||||||
Reference in New Issue
Block a user