返回源码地图

packages/session-query/session-query-sqlite/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 16;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * Concrete session-query service with SQLite FTS5 over the live-preferred corpus.
3 *
4 * @module @deepseek-ai/dsh-session-query-sqlite
5 */
6
7import { createHash, randomUUID } from 'node:crypto'
8import { SESSION_FORMAT_VERSION, SessionSeq } from '@deepseek-ai/dsh-session'
9import type { DatabaseSync } from 'node:sqlite'
10import { Context, Service, type Fiber } from '@deepseek-ai/cordis'
11import z from '@deepseek-ai/schemastery'
12import type { Session, SessionEvent, SessionHeader, SessionId, SessionLogOffset } from '@deepseek-ai/dsh-session'
13import type SessionPersistence from '@deepseek-ai/dsh-session-persistence'
14import type {
15 SessionPersistenceRevision,
16 SessionPersistenceSnapshot,
17} from '@deepseek-ai/dsh-session-persistence'
18import SessionQueryEngine, {
19 SESSION_QUERY_DEFAULT_PERSISTED_INSPECT_CONCURRENCY,
20 SESSION_QUERY_DEFAULT_PREPARED_SESSION_CACHE_SIZE,
21 SESSION_QUERY_READ_WINDOW_MAX,
22 SessionQueryError,
23 SessionSearchCursor,
24 assertSessionHeadersCompatible,
25 buildSessionEventSearchDocuments,
26 readColdSessionLog,
27} from '@deepseek-ai/dsh-session-query'
28import type {
29 Config as SessionQueryConfig,
30 SessionEventSearchDocument,
31 SessionEventSearchHit,
32 SessionEventSearchPage,
33 SessionEventSearchRequest,
34 SessionSearchExecContext,
35 SessionSearchHit,
36 SessionSearchCursor as SessionSearchCursorValue,
37 SessionSearchPage,
38 SessionSearchRequest,
39} from '@deepseek-ai/dsh-session-query'
40import {
41 type JournalMode,
42 openSearchDatabase,
43} from './schema.ts'
44import {
45 type NormalizedEventRequest,
46 type NormalizedSessionRequest,
47 FTS_HIGHLIGHT_END,
48 FTS_HIGHLIGHT_START,
49 assertFts5OuterPredicateCount,
50 assertPortableBindingCount,
51 buildEventWhere,
52 buildSessionWhere,
53 makeSnippet,
54 normalizeEventRequest,
55 normalizeSessionRequest,
56 quoteFtsData,
57 requestFingerprint,
58 sanitizeFtsText,
59 SQLITE_MAX_PAGE_LIMIT,
60} from './query.ts'
61
62export {
63 SESSION_QUERY_SQLITE_APPLICATION_ID,
64 SESSION_QUERY_SQLITE_SCHEMA_VERSION,
65 type JournalMode,
66} from './schema.ts'
67
68/** Boot-context slot for a launcher-owned absolute path to this process's derived query index. */
69export const SESSION_QUERY_SQLITE_PATH_KEY = 'launcherSessionQueryPath'
70
71declare module '@deepseek-ai/cordis' {
72 interface Context {
73 /** Launcher-owned absolute path to this process's disposable derived query index. */
74 launcherSessionQueryPath?: string
75 }
76}
77
78/** Default result page size. */
79export const SESSION_QUERY_SQLITE_DEFAULT_LIMIT = 20
80/** Maximum accepted result page size. */
81export const SESSION_QUERY_SQLITE_MAX_LIMIT = 100
82/** Default maximum snippet length in Unicode code points. */
83export const SESSION_QUERY_SQLITE_SNIPPET_CHARS = 240
84
85// One transient source change gets a retry; repeated churn fails rather than monopolizing the queue.
86const STABLE_OBSERVATION_ATTEMPTS = 2
87
88/** SQLite module/handle opening phase; `never` disables full-text search entirely. */
89export type OpenAt = 'startup' | 'first-search' | 'never'
90
91/** Combined session-query configuration backed by SQLite full-text search. */
92export interface Config extends SessionQueryConfig {
93 /**
94 * Dedicated derived-index path; `:memory:` is supported for ephemeral
95 * indexes. Missing directories and database files are created owner-only on
96 * POSIX filesystems; existing modes are preserved.
97 */
98 path: string
99 /**
100 * Open the SQLite module and handle at service activation or the first
101 * search, or `never` to disable full-text search: the inherited exact
102 * reads, filters, and traces stay available, while `searchSessions` and
103 * `searchEvents` fail with `SESSION_QUERY_SEARCH_DISABLED` and SQLite is
104 * never imported or opened. Defaults to `startup`.
105 */
106 openAt?: OpenAt
107 /** SQLite journal mode. Defaults to `wal`. */
108 journalMode?: JournalMode
109 /** Page size when a request omits `limit`. At most `Number.MAX_SAFE_INTEGER - 1`; defaults to 20. */
110 defaultLimit?: number
111 /** Largest accepted page size. At most `Number.MAX_SAFE_INTEGER - 1`; defaults to 100. */
112 maxLimit?: number
113 /** Maximum snippet length in Unicode code points. Defaults to 240. */
114 snippetChars?: number
115 /** Maximum concurrent persisted-log reads in one inherited batch read. Defaults to 4. */
116 persistedReadConcurrency?: number
117 /** Maximum cold prepared-Session observations the inherited reader retains for reuse. Defaults to 5. */
118 preparedSessionCacheSize?: number
119}
120
121interface ResolvedConfig {
122 path: string
123 openAt: OpenAt
124 journalMode: JournalMode
125 defaultLimit: number
126 maxLimit: number
127 snippetChars: number
128 readWindowMax: number
129 persistedReadConcurrency: number
130 preparedSessionCacheSize: number
131}
132
133interface ObservedSession {
134 header: SessionHeader
135 inheritedEventCount: SessionLogOffset
136 documents: SessionEventSearchDocument[]
137 fingerprint: string
138}
139
140interface ObservedPersistedSession {
141 header: SessionHeader
142 revision: SessionPersistenceRevision
143 loaded?: ObservedSession
144}
145
146interface PersistenceBinding {
147 readonly identity: symbol
148 readonly service?: SessionPersistence
149}
150
151interface Observation {
152 persistenceBinding: PersistenceBinding
153 persisted: Map<SessionId, ObservedPersistedSession>
154 live: Map<SessionId, ObservedSession>
155}
156
157interface IndexedPersistedRow {
158 id: string
159 revision: string
160 generation: number
161}
162
163interface IndexedLiveRow {
164 id: string
165 fingerprint: string
166 persisted: number
167 generation: number
168}
169
170interface SessionHeaderRow {
171 session_id: string
172 version: number
173 created_at: number
174 cwd: string | null
175 parent_session: string | null
176 seed_length: number | null
177 delegation_depth: number | null
178 agent_preset: string | null
179}
180
181interface SearchRow extends SessionHeaderRow {
182 live: number
183 persisted: number
184 seq: number
185 type: string
186 time: number
187 surface: string
188 marked_text: string
189 match_count: number
190 document_length: number
191}
192
193interface CursorPayload {
194 version: 1
195 instance: string
196 scope: 'sessions' | 'events'
197 fingerprint: string
198 generation: string
199 offset: number
200}
201
202/** Concrete SQLite owner of the combined `ctx.sessionQuery` service. */
203export class SqliteSessionQueryEngine extends SessionQueryEngine {
204 static override inject = ['sessions']
205
206 static Config: z<Config> = z.object({
207 path: z.string().required(),
208 openAt: z.union(['startup', 'first-search', 'never'] as const).default('startup'),
209 journalMode: z.union(['wal', 'delete', 'truncate', 'persist'] as const).default('wal'),
210 defaultLimit: z.number().step(1).min(1).max(SQLITE_MAX_PAGE_LIMIT).default(SESSION_QUERY_SQLITE_DEFAULT_LIMIT),
211 maxLimit: z.number().step(1).min(1).max(SQLITE_MAX_PAGE_LIMIT).default(SESSION_QUERY_SQLITE_MAX_LIMIT),
212 snippetChars: z.number().step(1).min(1).default(SESSION_QUERY_SQLITE_SNIPPET_CHARS),
213 readWindowMax: z.number().step(1).min(0).default(SESSION_QUERY_READ_WINDOW_MAX),
214 persistedReadConcurrency: z.number()
215 .step(1)
216 .min(1)
217 .max(Number.MAX_SAFE_INTEGER)
218 .default(SESSION_QUERY_DEFAULT_PERSISTED_INSPECT_CONCURRENCY),
219 preparedSessionCacheSize: z.number()
220 .step(1)
221 .min(1)
222 .max(Number.MAX_SAFE_INTEGER)
223 .default(SESSION_QUERY_DEFAULT_PREPARED_SESSION_CACHE_SIZE),
224 })
225
226 /** Validated and defaulted backend configuration. */
227 readonly config: ResolvedConfig
228
229 private readonly _instance = randomUUID()
230 private _ready: Promise<void> | undefined
231 private _db: DatabaseSync | undefined
232 private _persistenceBinding: PersistenceBinding = { identity: Symbol() }
233 private _lastPersistenceIdentity: symbol | undefined
234 private _persistenceEpoch = 0
235 private _globalGeneration = 0
236 private _localGeneration = 0
237 private _tail: Promise<void> = Promise.resolve()
238 private _closed = false
239 private _closePromise: Promise<void> | undefined
240 private readonly _optionalPersistenceFiber: Fiber
241
242 constructor(ctx: Context, config: Config) {
243 // The assignment expression resolves before the base constructor can
244 // register `ctx.sessionQuery`; keep that same validated value afterward.
245 super(ctx, config = resolveConfig(config))
246 this.config = config as ResolvedConfig
247 this._optionalPersistenceFiber = ctx.inject(['sessionPersistence'], (childCtx: Context) => {
248 const service = childCtx.sessionPersistence
249 const binding = { identity: Symbol(), service }
250 this._persistenceBinding = binding
251 childCtx.effect(() => () => {
252 /* v8 ignore next -- a stale optional-service disposer cannot clear a replacement */
253 if (this._persistenceBinding !== binding) return
254 this._persistenceBinding = { identity: Symbol() }
255 }, 'sessionQuerySqlite.persistenceBinding')
256 })
257 ctx.effect(() => {
258 return () => this._optionalPersistenceFiber.dispose()
259 }, 'sessionQuerySqlite.optionalPersistence')
260 ctx.effect(() => async () => this.close(), 'sessionQuerySqlite.close')
261 }
262
263 /** Open eagerly only when activation owns the configured readiness boundary. */
264 protected async [Service.init](): Promise<void> {
265 if (this.config.openAt === 'startup') await this._ensureReady(undefined)
266 }
267
268 override async searchSessions(
269 request: SessionSearchRequest,
270 exec?: SessionSearchExecContext,
271 ): Promise<SessionSearchPage<SessionSearchHit>> {
272 this._assertSearchEnabled()
273 const normalized = normalizeSessionRequest(request, this.config)
274 const signal = exec?.signal
275 return this._serialized(signal, async () => {
276 await this._ensureReady(signal)
277 const persistenceBinding = await this._reconcile(signal)
278 assertNotAborted(signal)
279 const generation = String(this._globalGeneration)
280 const fingerprint = requestFingerprint(normalized)
281 const offset = normalized.cursor === undefined
282 ? 0
283 : decodeCursor(normalized.cursor, this._instance, 'sessions', fingerprint, generation)
284 const rows = this._querySessions(normalized, offset, persistenceBinding)
285 return page(rows, normalized.limit, row => this._sessionHit(row), cursorOffset => encodeCursor({
286 version: 1,
287 instance: this._instance,
288 scope: 'sessions',
289 fingerprint,
290 generation,
291 offset: cursorOffset,
292 }), offset)
293 })
294 }
295
296 override async searchEvents(
297 request: SessionEventSearchRequest,
298 exec?: SessionSearchExecContext,
299 ): Promise<SessionEventSearchPage> {
300 this._assertSearchEnabled()
301 const normalized = normalizeEventRequest(request, this.config)
302 const signal = exec?.signal
303 return this._serialized(signal, async () => {
304 await this._ensureReady(signal)
305 const persistenceBinding = await this._reconcile(signal)
306 assertNotAborted(signal)
307 const target = this._targetObservation(normalized.sessionId, persistenceBinding)
308 const fingerprint = requestFingerprint(normalized)
309 const offset = normalized.cursor === undefined
310 ? 0
311 : decodeCursor(normalized.cursor, this._instance, 'events', fingerprint, target.generation)
312 const rows = this._queryEvents(normalized, offset, persistenceBinding)
313 return {
314 session: target.header,
315 ...page(rows, normalized.limit, row => this._eventHit(row), cursorOffset => encodeCursor({
316 version: 1,
317 instance: this._instance,
318 scope: 'events',
319 fingerprint,
320 generation: target.generation,
321 offset: cursorOffset,
322 }), offset),
323 }
324 })
325 }
326
327 /** Close the database after every accepted operation reaches quiescence. */
328 close(): Promise<void> {
329 this._closePromise ??= this._close()
330 return this._closePromise
331 }
332
333 /**
334 * Refuse full-text calls under `openAt: 'never'` before any request
335 * normalization or SQLite work, so a disabled deployment never imports
336 * node:sqlite, opens the index, or observes sources.
337 */
338 private _assertSearchEnabled(): void {
339 if (this.config.openAt !== 'never') return
340 throw new SessionQueryError(
341 'session search is disabled: this deployment configures the session-query index with openAt "never"',
342 'SESSION_QUERY_SEARCH_DISABLED',
343 )
344 }
345
346 private async _close(): Promise<void> {
347 this._closed = true
348 await this._tail
349 if (this._ready !== undefined) {
350 try {
351 await this._ready
352 } catch {
353 // Opening already closed a partially-created handle; disposal only waits.
354 }
355 }
356 this._db?.close()
357 this._db = undefined
358 }
359
360 private async _open(): Promise<void> {
361 this._db = await openSearchDatabase(this.config.path, this.config.journalMode)
362 const state = this._db.prepare(
363 'SELECT global_generation FROM search_state WHERE singleton = 1',
364 ).get() as { global_generation: number }
365 this._globalGeneration = state.global_generation
366 this._localGeneration = state.global_generation
367 }
368
369 private async _ensureReady(signal: AbortSignal | undefined): Promise<void> {
370 this._ready ??= this._open()
371 try {
372 await waitWithAbort(this._ready, signal)
373 } catch (error: unknown) {
374 if (isAbort(error)) throw error
375 throw new SessionQueryError(
376 `session-search SQLite index failed to open: ${errorMessage(error)}`,
377 'SESSION_QUERY_INDEX_FAILED',
378 { cause: error },
379 )
380 }
381 }
382
383 private async _serialized<T>(signal: AbortSignal | undefined, operation: () => Promise<T>): Promise<T> {
384 if (this._isClosed()) throw indexClosed()
385 let release!: () => void
386 const gate = new Promise<void>((resolve) => { release = resolve })
387 const prior = this._tail
388 this._tail = prior.then(() => gate)
389 try {
390 await waitWithAbort(prior, signal)
391 } catch (error: unknown) {
392 release()
393 throw error
394 }
395 if (this._isClosed()) {
396 release()
397 throw indexClosed()
398 }
399 try {
400 assertNotAborted(signal)
401 return await operation()
402 } finally {
403 release()
404 }
405 }
406
407 private async _reconcile(signal: AbortSignal | undefined): Promise<PersistenceBinding> {
408 assertNotAborted(signal)
409 const db = this._requireDb()
410 const persistedRows = db.prepare(
411 'SELECT id, revision, generation FROM persisted_sessions',
412 ).all() as unknown as IndexedPersistedRow[]
413 const liveRows = db.prepare(
414 'SELECT id, fingerprint, persisted, generation FROM temp.live_sessions',
415 ).all() as unknown as IndexedLiveRow[]
416 const persistedById = new Map(persistedRows.map(row => [row.id as SessionId, row]))
417 const liveById = new Map(liveRows.map(row => [row.id as SessionId, row]))
418 const observation = await this._observeStable(persistedById, signal)
419 assertNotAborted(signal)
420 const persistentChanges = observation.persistenceBinding.service === undefined
421 ? []
422 : [...observation.persisted.values()].filter(entry => entry.loaded !== undefined)
423 const persistentDeletes = observation.persistenceBinding.service === undefined
424 ? []
425 : persistedRows.filter(row => !observation.persisted.has(row.id as SessionId))
426 const liveChanges = [...observation.live.values()].filter((entry) => {
427 const indexed = liveById.get(entry.header.id)
428 const persisted = observation.persisted.has(entry.header.id) ? 1 : 0
429 return indexed?.fingerprint !== entry.fingerprint || indexed.persisted !== persisted
430 })
431 const liveDeletes = liveRows.filter(row => !observation.live.has(row.id as SessionId))
432 const pointerChanged = this._lastPersistenceIdentity !== undefined
433 && this._lastPersistenceIdentity !== observation.persistenceBinding.identity
434 const hasWrites = persistentChanges.length > 0
435 || persistentDeletes.length > 0
436 || liveChanges.length > 0
437 || liveDeletes.length > 0
438
439 let nextMainGeneration = this._mainGeneration()
440 let nextLocalGeneration = this._localGeneration
441 if (persistentChanges.length > 0 || persistentDeletes.length > 0) nextMainGeneration += 1
442 const liveReplacements = liveChanges.map((entry) => {
443 nextLocalGeneration = Math.max(nextLocalGeneration, nextMainGeneration) + 1
444 return {
445 entry,
446 generation: nextLocalGeneration,
447 persisted: observation.persisted.has(entry.header.id),
448 }
449 })
450
451 if (hasWrites) {
452 let began = false
453 try {
454 db.exec('BEGIN IMMEDIATE')
455 began = true
456 for (const row of persistentDeletes) this._deleteSession('persisted', row.id as SessionId)
457 for (const entry of persistentChanges) {
458 /* v8 ignore next -- observation loads every entry whose revision differs */
459 if (entry.loaded === undefined) throw new Error(`missing loaded revision for session "${entry.header.id}"`)
460 this._replacePersistedSession(entry.loaded, entry.revision, nextMainGeneration)
461 }
462 if (persistentChanges.length > 0 || persistentDeletes.length > 0) {
463 db.prepare('UPDATE search_state SET global_generation = ? WHERE singleton = 1').run(nextMainGeneration)
464 }
465 for (const row of liveDeletes) this._deleteSession('live', row.id as SessionId)
466 for (const { entry, generation, persisted } of liveReplacements) {
467 this._replaceLiveSession(entry, generation, persisted)
468 }
469 db.exec('COMMIT')
470 } catch (error: unknown) {
471 /* v8 ignore next -- a BEGIN failure has no transaction to roll back; the common wrapper still reports it. */
472 if (began) {
473 /* v8 ignore next 5 -- ROLLBACK failure requires a SQLite double fault; the original failure remains actionable. */
474 try {
475 db.exec('ROLLBACK')
476 } catch {
477 // The original SQLite failure remains the actionable cause.
478 }
479 }
480 throw new SessionQueryError(
481 `session-search reconciliation failed: ${errorMessage(error)}`,
482 'SESSION_QUERY_INDEX_FAILED',
483 { cause: error },
484 )
485 }
486 }
487
488 if (hasWrites || pointerChanged) this._globalGeneration += 1
489 if (pointerChanged) this._persistenceEpoch += 1
490 this._localGeneration = nextLocalGeneration
491 this._lastPersistenceIdentity = observation.persistenceBinding.identity
492 return observation.persistenceBinding
493 }
494
495 private async _observeStable(
496 indexed: ReadonlyMap<SessionId, IndexedPersistedRow>,
497 signal: AbortSignal | undefined,
498 ): Promise<Observation> {
499 for (let attempt = 0; attempt < STABLE_OBSERVATION_ATTEMPTS; attempt += 1) {
500 assertNotAborted(signal)
501 const persistenceBinding = this._persistenceBinding
502 const persistence = persistenceBinding.service
503 const initiallyLive = new Set(this.ctx.sessions.list().map(session => session.id))
504 let persisted = new Map<SessionId, ObservedPersistedSession>()
505 if (persistence !== undefined) {
506 try {
507 const canReuseIndexed = this._lastPersistenceIdentity === undefined
508 || this._lastPersistenceIdentity === persistenceBinding.identity
509 const listOptions = signal === undefined ? undefined : { signal }
510 const before = await persistence.list(listOptions)
511 assertNotAborted(signal)
512 persisted = materializePersistenceSnapshots(before)
513 for (const entry of persisted.values()) {
514 if (canReuseIndexed && indexed.get(entry.header.id)?.revision === entry.revision) continue
515 // Skip work already shadowed by a live owner. The cold read is
516 // non-mutating (interrupted turns are balanced in memory only), so
517 // an owner attaching after this check cannot cause side effects;
518 // the live-membership retry below makes the returned observation
519 // live-preferred.
520 if (initiallyLive.has(entry.header.id) || this.ctx.sessions.get(entry.header.id) !== undefined) continue
521 assertNotAborted(signal)
522 const loaded = await readColdSessionLog(persistence, entry.header.id, signal)
523 assertNotAborted(signal)
524 assertSessionHeadersCompatible(entry.header, loaded.header)
525 entry.loaded = observeSession(loaded.header, loaded.inheritedEventCount, loaded.events)
526 }
527 assertNotAborted(signal)
528 const afterSnapshots = await persistence.list(listOptions)
529 assertNotAborted(signal)
530 const after = materializePersistenceSnapshots(afterSnapshots)
531 if (!samePersistenceSnapshots(persisted, after)) continue
532 if (this._persistenceBinding !== persistenceBinding) continue
533 } catch (error: unknown) {
534 if (isAbort(error) || signal?.aborted) {
535 throw new SessionQueryError('session-search aborted', 'SESSION_QUERY_ABORTED', {
536 cause: error,
537 })
538 }
539 if (this._persistenceBinding !== persistenceBinding) continue
540 if (error instanceof SessionQueryError) throw error
541 throw new SessionQueryError(
542 `session-search persistence observation failed: ${errorMessage(error)}`,
543 'SESSION_QUERY_PERSISTENCE_FAILED',
544 { cause: error },
545 )
546 }
547 }
548 const live = new Map<SessionId, ObservedSession>()
549 for (const session of this.ctx.sessions.list()) {
550 const observed = observeLive(session)
551 const durable = persisted.get(session.id)
552 if (durable !== undefined) assertSessionHeadersCompatible(observed.header, durable.header)
553 live.set(session.id, observed)
554 }
555 if (!sameSessionIds(initiallyLive, live)) continue
556 return { persistenceBinding, persisted, live }
557 }
558 throw new SessionQueryError(
559 'session-search persistence observation did not stabilize after one retry',
560 'SESSION_QUERY_PERSISTENCE_FAILED',
561 )
562 }
563
564 private _mainGeneration(): number {
565 const row = this._requireDb().prepare(
566 'SELECT global_generation FROM search_state WHERE singleton = 1',
567 ).get() as { global_generation: number }
568 return row.global_generation
569 }
570
571 private _deleteSession(source: 'persisted' | 'live', id: SessionId): void {
572 const db = this._requireDb()
573 if (source === 'persisted') {
574 db.prepare('DELETE FROM persisted_docs WHERE session_id = ?').run(id)
575 db.prepare('DELETE FROM persisted_sessions WHERE id = ?').run(id)
576 } else {
577 db.prepare('DELETE FROM temp.live_docs WHERE session_id = ?').run(id)
578 db.prepare('DELETE FROM temp.live_sessions WHERE id = ?').run(id)
579 }
580 }
581
582 private _replacePersistedSession(
583 entry: ObservedSession,
584 revision: SessionPersistenceRevision,
585 generation: number,
586 ): void {
587 this._deleteSession('persisted', entry.header.id)
588 const db = this._requireDb()
589 db.prepare(`
590 INSERT INTO persisted_sessions
591 (id, version, created_at, cwd, parent_session, seed_length, delegation_depth, agent_preset, revision, generation)
592 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
593 `).run(
594 ...headerBindings(entry.header, entry.inheritedEventCount),
595 revision,
596 generation,
597 )
598 const insert = db.prepare(`
599 INSERT INTO persisted_docs (text, session_id, seq, type, time, surface, codepoint_length)
600 VALUES (?, ?, ?, ?, ?, ?, ?)
601 `)
602 for (const document of entry.documents) {
603 const text = sanitizeFtsText(document.text)
604 insert.run(
605 text,
606 document.sessionId,
607 document.seq,
608 document.type,
609 document.time,
610 document.surface,
611 Array.from(text).length,
612 )
613 }
614 }
615
616 private _replaceLiveSession(entry: ObservedSession, generation: number, persisted: boolean): void {
617 this._deleteSession('live', entry.header.id)
618 const db = this._requireDb()
619 db.prepare(`
620 INSERT INTO temp.live_sessions
621 (id, version, created_at, cwd, parent_session, seed_length, delegation_depth, agent_preset, fingerprint, persisted, generation)
622 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
623 `).run(
624 ...headerBindings(entry.header, entry.inheritedEventCount),
625 entry.fingerprint,
626 persisted ? 1 : 0,
627 generation,
628 )
629 const insert = db.prepare(`
630 INSERT INTO temp.live_docs (text, session_id, seq, type, time, surface, codepoint_length)
631 VALUES (?, ?, ?, ?, ?, ?, ?)
632 `)
633 for (const document of entry.documents) {
634 const text = sanitizeFtsText(document.text)
635 insert.run(
636 text,
637 document.sessionId,
638 document.seq,
639 document.type,
640 document.time,
641 document.surface,
642 Array.from(text).length,
643 )
644 }
645 }
646
647 private _querySessions(
648 request: NormalizedSessionRequest,
649 offset: number,
650 persistenceBinding: PersistenceBinding,
651 ): SearchRow[] {
652 const selected = selectedDocumentsSql()
653 const sessionWhere = buildSessionWhere(request.sessionFilters)
654 const eventWhere = buildEventWhere(request.eventFilters)
655 assertFts5OuterPredicateCount(sessionWhere.predicateCount + eventWhere.predicateCount)
656 const where = [sessionWhere.sql, eventWhere.sql].filter(Boolean).join(' AND ')
657 const bindings = [
658 ...selectedDocumentsParams(request.query, persistenceBinding.service !== undefined),
659 ...sessionWhere.params,
660 ...eventWhere.params,
661 request.limit + 1,
662 offset,
663 ]
664 assertPortableBindingCount(bindings.length)
665 return this._requireDb().prepare(`
666 ${selected.sql},
667 filtered AS (
668 SELECT * FROM matched ${where.length === 0 ? '' : `WHERE ${where}`}
669 ),
670 ranked AS (
671 SELECT *, ROW_NUMBER() OVER (
672 PARTITION BY session_id
673 ORDER BY match_count DESC, document_length ASC, time DESC, seq DESC
674 ) AS event_rank
675 FROM filtered
676 )
677 SELECT * FROM ranked
678 WHERE event_rank = 1
679 ORDER BY match_count DESC, document_length ASC, time DESC, session_id ASC, seq DESC
680 LIMIT ? OFFSET ?
681 `).all(...bindings) as unknown as SearchRow[]
682 }
683
684 private _queryEvents(
685 request: NormalizedEventRequest,
686 offset: number,
687 persistenceBinding: PersistenceBinding,
688 ): SearchRow[] {
689 const selected = selectedDocumentsSql()
690 const eventWhere = buildEventWhere(request.filters)
691 assertFts5OuterPredicateCount(1 + eventWhere.predicateCount)
692 const where = ['session_id = ?', eventWhere.sql].filter(Boolean).join(' AND ')
693 const bindings = [
694 ...selectedDocumentsParams(request.query, persistenceBinding.service !== undefined),
695 request.sessionId,
696 ...eventWhere.params,
697 request.limit + 1,
698 offset,
699 ]
700 assertPortableBindingCount(bindings.length)
701 return this._requireDb().prepare(`
702 ${selected.sql}
703 SELECT * FROM matched
704 WHERE ${where}
705 ORDER BY match_count DESC, document_length ASC, time DESC, seq DESC
706 LIMIT ? OFFSET ?
707 `).all(...bindings) as unknown as SearchRow[]
708 }
709
710 private _targetObservation(
711 sessionId: SessionId,
712 persistenceBinding: PersistenceBinding,
713 ): { header: SessionHeader; generation: string } {
714 const db = this._requireDb()
715 const live = db.prepare(
716 `SELECT
717 id AS session_id, version, created_at, cwd, parent_session, seed_length, delegation_depth, agent_preset, generation
718 FROM temp.live_sessions
719 WHERE id = ?`,
720 ).get(sessionId) as (SessionHeaderRow & { generation: number }) | undefined
721 if (live !== undefined) {
722 return { header: rowHeader(live), generation: `live:${live.generation}` }
723 }
724 if (persistenceBinding.service !== undefined) {
725 const persisted = db.prepare(
726 `SELECT
727 id AS session_id, version, created_at, cwd, parent_session, seed_length, delegation_depth, agent_preset, generation
728 FROM persisted_sessions
729 WHERE id = ?`,
730 ).get(sessionId) as (SessionHeaderRow & { generation: number }) | undefined
731 if (persisted !== undefined) {
732 return {
733 header: rowHeader(persisted),
734 generation: `persisted:${this._persistenceEpoch}:${persisted.generation}`,
735 }
736 }
737 }
738 throw new SessionQueryError(
739 `session "${sessionId}" not found`,
740 'SESSION_QUERY_SESSION_NOT_FOUND',
741 )
742 }
743
744 private _sessionHit(row: SearchRow): SessionSearchHit {
745 return {
746 header: rowHeader(row),
747 live: row.live === 1,
748 persisted: row.persisted === 1,
749 bestMatch: this._eventHit(row),
750 }
751 }
752
753 private _eventHit(row: SearchRow): SessionEventSearchHit {
754 return {
755 sessionId: row.session_id as SessionId,
756 seq: SessionSeq(row.seq),
757 type: row.type as SessionEventSearchHit['type'],
758 time: row.time,
759 surface: row.surface as SessionEventSearchHit['surface'],
760 snippet: makeSnippet(row.marked_text, this.config.snippetChars),
761 }
762 }
763
764 private _requireDb(): DatabaseSync {
765 /* v8 ignore next -- callers await `_ready`; this guards lifecycle misuse */
766 if (this._db === undefined) throw indexClosed()
767 return this._db
768 }
769
770 private _isClosed(): boolean {
771 return this._closed
772 }
773}
774
775/**
776 * The header columns both session upserts bind, in the order their INSERT
777 * lists them. The two statements differ only in what they append after these.
778 * @param header - the session header being written.
779 * @returns one bound value per header column.
780 */
781function headerBindings(
782 header: SessionHeader,
783 inheritedEventCount: SessionLogOffset,
784): (string | number | null)[] {
785 return [
786 header.id,
787 header.version,
788 header.createdAt,
789 header.cwd ?? null,
790 header.parentSession ?? null,
791 header.isSeeded ? inheritedEventCount : null,
792 header.delegationDepth ?? null,
793 header.agentPreset ?? null,
794 ]
795}
796
797function selectedDocumentsSql(): { sql: string } {
798 return {
799 sql: `WITH candidates AS (
800 SELECT
801 pd.session_id AS session_id,
802 ps.version AS version,
803 ps.created_at AS created_at,
804 ps.cwd AS cwd,
805 ps.parent_session AS parent_session,
806 ps.seed_length AS seed_length,
807 ps.delegation_depth AS delegation_depth,
808 ps.agent_preset AS agent_preset,
809 0 AS live,
810 1 AS persisted,
811 CAST(pd.seq AS INTEGER) AS seq,
812 pd.type AS type,
813 CAST(pd.time AS INTEGER) AS time,
814 pd.surface AS surface,
815 highlight(persisted_docs, 0, ?, ?) AS marked_text,
816 CAST(pd.codepoint_length AS INTEGER) AS document_length
817 FROM persisted_docs AS pd
818 JOIN persisted_sessions AS ps ON ps.id = pd.session_id
819 WHERE persisted_docs MATCH ?
820 AND ? = 1
821 AND NOT EXISTS (SELECT 1 FROM temp.live_sessions AS ls WHERE ls.id = pd.session_id)
822 UNION ALL
823 SELECT
824 ld.session_id AS session_id,
825 ls.version AS version,
826 ls.created_at AS created_at,
827 ls.cwd AS cwd,
828 ls.parent_session AS parent_session,
829 ls.seed_length AS seed_length,
830 ls.delegation_depth AS delegation_depth,
831 ls.agent_preset AS agent_preset,
832 1 AS live,
833 CASE WHEN ? = 1 THEN ls.persisted ELSE 0 END AS persisted,
834 CAST(ld.seq AS INTEGER) AS seq,
835 ld.type AS type,
836 CAST(ld.time AS INTEGER) AS time,
837 ld.surface AS surface,
838 highlight(live_docs, 0, ?, ?) AS marked_text,
839 CAST(ld.codepoint_length AS INTEGER) AS document_length
840 FROM temp.live_docs AS ld
841 JOIN temp.live_sessions AS ls ON ls.id = ld.session_id
842 WHERE live_docs MATCH ?
843 ), matched AS (
844 SELECT *,
845 (
846 length(CAST(marked_text AS BLOB))
847 - length(CAST(replace(marked_text, ?, '') AS BLOB))
848 ) / ? AS match_count
849 FROM candidates
850 )`,
851 }
852}
853
854function selectedDocumentsParams(query: string, persistenceVisible: boolean): Array<string | number> {
855 const expression = quoteFtsData(query)
856 const visible = persistenceVisible ? 1 : 0
857 return [
858 FTS_HIGHLIGHT_START,
859 FTS_HIGHLIGHT_END,
860 expression,
861 visible,
862 visible,
863 FTS_HIGHLIGHT_START,
864 FTS_HIGHLIGHT_END,
865 expression,
866 FTS_HIGHLIGHT_START,
867 Buffer.byteLength(FTS_HIGHLIGHT_START, 'utf8'),
868 ]
869}
870
871function observeLive(session: Session): ObservedSession {
872 // oxlint-disable-next-line typescript/no-deprecated -- Existing Session history read; migration deferred.
873 return observeSession(session.header, session.inheritedEventCount, session.snapshotEvents())
874}
875
876function observeSession(
877 header: SessionHeader,
878 inheritedEventCount: SessionLogOffset,
879 events: readonly SessionEvent[],
880): ObservedSession {
881 const detachedHeader = structuredClone(header)
882 const detachedEvents = events.map(event => structuredClone(event))
883 return {
884 header: detachedHeader,
885 inheritedEventCount,
886 documents: buildSessionEventSearchDocuments(detachedHeader.id, detachedEvents),
887 fingerprint: createHash('sha256')
888 .update(JSON.stringify({ header: detachedHeader, inheritedEventCount, events: detachedEvents }))
889 .digest('base64url'),
890 }
891}
892
893function materializePersistenceSnapshots(
894 snapshots: readonly SessionPersistenceSnapshot[],
895): Map<SessionId, ObservedPersistedSession> {
896 if (!isRuntimeArray(snapshots)) throw new Error('persistence snapshots must be an array')
897 const result = new Map<SessionId, ObservedPersistedSession>()
898 for (const snapshot of snapshots) {
899 if (typeof snapshot.revision !== 'string') {
900 throw new Error('persistence snapshot revision must be a string')
901 }
902 const header = structuredClone(snapshot.header)
903 if (result.has(header.id)) {
904 throw new Error(`persistence listed duplicate session "${header.id}"`)
905 }
906 result.set(header.id, { header, revision: snapshot.revision })
907 }
908 return result
909}
910
911function samePersistenceSnapshots(
912 before: ReadonlyMap<SessionId, ObservedPersistedSession>,
913 after: ReadonlyMap<SessionId, ObservedPersistedSession>,
914): boolean {
915 if (before.size !== after.size) return false
916 for (const [id, first] of before) {
917 const second = after.get(id)
918 if (
919 second === undefined
920 || first.revision !== second.revision
921 || !sameHeader(first.header, second.header)
922 ) return false
923 }
924 return true
925}
926
927function sameSessionIds(
928 before: ReadonlySet<SessionId>,
929 after: ReadonlyMap<SessionId, ObservedSession>,
930): boolean {
931 if (before.size !== after.size) return false
932 for (const id of before) {
933 if (!after.has(id)) return false
934 }
935 return true
936}
937
938function sameHeader(a: SessionHeader, b: SessionHeader): boolean {
939 return a.id === b.id
940 && a.createdAt === b.createdAt
941 && a.cwd === b.cwd
942 && a.parentSession === b.parentSession
943 && a.isSeeded === b.isSeeded
944 && (a.delegationDepth ?? 0) === (b.delegationDepth ?? 0)
945 && a.agentPreset === b.agentPreset
946}
947
948function rowHeader(row: SessionHeaderRow): SessionHeader {
949 return {
950 version: SESSION_FORMAT_VERSION,
951 id: row.session_id as SessionId,
952 createdAt: row.created_at,
953 ...row.cwd === null ? {} : { cwd: row.cwd },
954 ...row.parent_session === null ? {} : { parentSession: row.parent_session as SessionId },
955 isSeeded: row.seed_length !== null,
956 ...row.delegation_depth === null ? {} : { delegationDepth: row.delegation_depth },
957 ...row.agent_preset === null ? {} : { agentPreset: row.agent_preset },
958 }
959}
960
961function page<Row, Item>(
962 rows: readonly Row[],
963 limit: number,
964 convert: (row: Row) => Item,
965 nextCursor: (offset: number) => SessionSearchCursorValue,
966 offset: number,
967): SessionSearchPage<Item> {
968 const hasMore = rows.length > limit
969 return {
970 items: rows.slice(0, limit).map(convert),
971 ...hasMore ? { nextCursor: nextCursor(offset + limit) } : {},
972 }
973}
974
975function encodeCursor(payload: CursorPayload): SessionSearchCursorValue {
976 return SessionSearchCursor(Buffer.from(JSON.stringify(payload), 'utf8').toString('base64url'))
977}
978
979function decodeCursor(
980 cursor: SessionSearchCursorValue,
981 instance: string,
982 scope: CursorPayload['scope'],
983 fingerprint: string,
984 generation: string,
985): number {
986 let decoded: Partial<CursorPayload>
987 try {
988 decoded = JSON.parse(Buffer.from(cursor, 'base64url').toString('utf8')) as Partial<CursorPayload>
989 } catch (error: unknown) {
990 throw invalidCursor(error)
991 }
992 if (
993 decoded.version !== 1
994 || decoded.instance !== instance
995 || decoded.scope !== scope
996 || decoded.fingerprint !== fingerprint
997 || !Number.isSafeInteger(decoded.offset)
998 || decoded.offset === undefined
999 || decoded.offset < 0
1000 ) {
1001 throw invalidCursor(new Error('cursor does not belong to this normalized request'))
1002 }
1003 if (decoded.generation !== generation) {
1004 throw new SessionQueryError(
1005 'session-search cursor is stale because its relevant corpus changed',
1006 'SESSION_QUERY_STALE_CURSOR',
1007 )
1008 }
1009 return decoded.offset
1010}
1011
1012function invalidCursor(cause: unknown): SessionQueryError {
1013 return new SessionQueryError(
1014 'session-search cursor is invalid',
1015 'SESSION_QUERY_INVALID_CURSOR',
1016 { cause },
1017 )
1018}
1019
1020function resolveConfig(config: Config): ResolvedConfig {
1021 const resolved: ResolvedConfig = {
1022 path: config.path,
1023 openAt: config.openAt ?? 'startup',
1024 journalMode: config.journalMode ?? 'wal',
1025 defaultLimit: config.defaultLimit ?? SESSION_QUERY_SQLITE_DEFAULT_LIMIT,
1026 maxLimit: config.maxLimit ?? SESSION_QUERY_SQLITE_MAX_LIMIT,
1027 snippetChars: config.snippetChars ?? SESSION_QUERY_SQLITE_SNIPPET_CHARS,
1028 readWindowMax: config.readWindowMax ?? SESSION_QUERY_READ_WINDOW_MAX,
1029 persistedReadConcurrency: config.persistedReadConcurrency
1030 ?? SESSION_QUERY_DEFAULT_PERSISTED_INSPECT_CONCURRENCY,
1031 preparedSessionCacheSize: config.preparedSessionCacheSize
1032 ?? SESSION_QUERY_DEFAULT_PREPARED_SESSION_CACHE_SIZE,
1033 }
1034 if (typeof resolved.path !== 'string' || resolved.path.trim().length === 0) {
1035 throw invalidConfig('path must not be blank')
1036 }
1037 const openPhases: readonly string[] = ['startup', 'first-search', 'never']
1038 if (!openPhases.includes(resolved.openAt)) throw invalidConfig('openAt is not supported')
1039 assertPageLimit('defaultLimit', resolved.defaultLimit)
1040 assertPageLimit('maxLimit', resolved.maxLimit)
1041 assertPositiveInteger('snippetChars', resolved.snippetChars)
1042 if (!Number.isInteger(resolved.readWindowMax) || resolved.readWindowMax < 0) {
1043 throw invalidConfig('readWindowMax must be a non-negative integer')
1044 }
1045 if (
1046 !Number.isSafeInteger(resolved.persistedReadConcurrency)
1047 || resolved.persistedReadConcurrency < 1
1048 ) {
1049 throw invalidConfig('persistedReadConcurrency must be a positive safe integer')
1050 }
1051 if (
1052 !Number.isSafeInteger(resolved.preparedSessionCacheSize)
1053 || resolved.preparedSessionCacheSize < 1
1054 ) {
1055 throw invalidConfig('preparedSessionCacheSize must be a positive safe integer')
1056 }
1057 if (resolved.defaultLimit > resolved.maxLimit) {
1058 throw invalidConfig('defaultLimit must be less than or equal to maxLimit')
1059 }
1060 const journalModes: readonly string[] = ['wal', 'delete', 'truncate', 'persist']
1061 if (!journalModes.includes(resolved.journalMode)) throw invalidConfig('journalMode is not supported')
1062 return resolved
1063}
1064
1065function assertPositiveInteger(name: string, value: number): void {
1066 if (!Number.isInteger(value) || value < 1) throw invalidConfig(`${name} must be a positive integer`)
1067}
1068
1069function assertPageLimit(name: string, value: number): void {
1070 if (!Number.isSafeInteger(value) || value < 1 || value > SQLITE_MAX_PAGE_LIMIT) {
1071 throw invalidConfig(`${name} must be an integer between 1 and ${SQLITE_MAX_PAGE_LIMIT}`)
1072 }
1073}
1074
1075function invalidConfig(detail: string): SessionQueryError {
1076 return new SessionQueryError(
1077 `session-search SQLite config: ${detail}`,
1078 'SESSION_QUERY_INVALID_CONFIG',
1079 )
1080}
1081
1082function indexClosed(): SessionQueryError {
1083 return new SessionQueryError('session-search SQLite index is closed', 'SESSION_QUERY_INDEX_FAILED')
1084}
1085
1086function assertNotAborted(signal: AbortSignal | undefined): void {
1087 if (signal?.aborted) {
1088 throw new SessionQueryError('session-search aborted', 'SESSION_QUERY_ABORTED')
1089 }
1090}
1091
1092function waitWithAbort<T>(promise: Promise<T>, signal: AbortSignal | undefined): Promise<T> {
1093 if (signal === undefined) return promise
1094 if (signal.aborted) return Promise.reject(new SessionQueryError('session-search aborted', 'SESSION_QUERY_ABORTED'))
1095 return new Promise<T>((resolve, reject) => {
1096 const onAbort = () => {
1097 reject(new SessionQueryError('session-search aborted', 'SESSION_QUERY_ABORTED'))
1098 }
1099 signal.addEventListener('abort', onAbort, { once: true })
1100 promise.then(
1101 (value) => {
1102 signal.removeEventListener('abort', onAbort)
1103 resolve(value)
1104 },
1105 (error: unknown) => {
1106 signal.removeEventListener('abort', onAbort)
1107 reject(asError(error))
1108 },
1109 )
1110 })
1111}
1112
1113function isAbort(error: unknown): boolean {
1114 return error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED'
1115}
1116
1117function asError(error: unknown): Error {
1118 return error instanceof Error
1119 ? error
1120 : new Error('session-search dependency rejected with a non-Error value', { cause: error })
1121}
1122
1123function errorMessage(error: unknown): string {
1124 return error instanceof Error ? error.message : 'unknown error'
1125}
1126
1127function isRuntimeArray(value: unknown): boolean {
1128 return Array.isArray(value)
1129}
1130
1131export default SqliteSessionQueryEngine