sqlite.spec.ts 81 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819
  1. import { createAssistantMessage, createUserMessage } from '@deepseek-ai/dsh-llm'
  2. import { afterEach, describe, expect, it, vi } from 'vitest'
  3. import { Context, type Fiber } from '@deepseek-ai/cordis'
  4. import { DatabaseSync } from 'node:sqlite'
  5. import { chmod, mkdtemp, rm, stat, writeFile } from 'node:fs/promises'
  6. import { tmpdir } from 'node:os'
  7. import { dirname, join } from 'node:path'
  8. import SessionStore, { SESSION_FORMAT_VERSION, SessionId } from '@deepseek-ai/dsh-session'
  9. import type { SessionEvent, SessionHeader, SessionId as SessionIdType } from '@deepseek-ai/dsh-session'
  10. import SessionPersistence, { SessionPersistenceRevision } from '@deepseek-ai/dsh-session-persistence'
  11. import type { SessionPersistenceSnapshot } from '@deepseek-ai/dsh-session-persistence'
  12. import SqliteSessionPersistence from '@deepseek-ai/dsh-session-persistence-sqlite'
  13. import SqliteSessionQueryEngine, {
  14. SESSION_QUERY_SQLITE_SCHEMA_VERSION,
  15. } from '@deepseek-ai/dsh-session-query-sqlite'
  16. import {
  17. SESSION_QUERY_DEFAULT_PERSISTED_INSPECT_CONCURRENCY,
  18. SessionQueryError,
  19. SessionSearchCursor,
  20. type SessionAvailability,
  21. type SessionQueryErrorCode,
  22. type SessionSearchRequest,
  23. } from '@deepseek-ai/dsh-session-query'
  24. const temporaryDirectories: string[] = []
  25. afterEach(async () => {
  26. for (const directory of temporaryDirectories.splice(0)) {
  27. await rm(directory, { recursive: true, force: true })
  28. }
  29. })
  30. async function temporaryPath(name = 'search.db'): Promise<string> {
  31. const directory = await mkdtemp(join(tmpdir(), 'dsh-session-search-'))
  32. temporaryDirectories.push(directory)
  33. return join(directory, name)
  34. }
  35. function header(id: string, createdAt = 1, extra: Partial<SessionHeader> = {}): SessionHeader {
  36. return { version: SESSION_FORMAT_VERSION, id: SessionId(id), createdAt, ...extra }
  37. }
  38. function messageEvents(text: string, time = 1): SessionEvent[] {
  39. return [{
  40. type: 'user/message',
  41. seq: 0,
  42. time,
  43. data: createUserMessage({
  44. content: [{ type: 'text', text }], source: { kind: 'user' },
  45. }),
  46. surfaceOp: 'append',
  47. }]
  48. }
  49. function expectCode(code: SessionQueryErrorCode): Error {
  50. return expect.objectContaining({ code }) as Error
  51. }
  52. function replaceCursorOffset(
  53. cursor: ReturnType<typeof SessionSearchCursor>,
  54. offset: number,
  55. ): ReturnType<typeof SessionSearchCursor> {
  56. const payload = JSON.parse(
  57. Buffer.from(cursor, 'base64url').toString('utf8'),
  58. ) as Record<string, unknown>
  59. return SessionSearchCursor(Buffer.from(JSON.stringify({ ...payload, offset }), 'utf8').toString('base64url'))
  60. }
  61. class TestPersistence extends SessionPersistence {
  62. override readonly supportsRawArtifacts = false
  63. static entries = new Map<SessionIdType, { meta: SessionHeader; events: SessionEvent[] }>()
  64. static revisions = new Map<SessionIdType, number>()
  65. static nextRevision = 0
  66. static loads = new Map<SessionIdType, number>()
  67. static inspections = new Map<SessionIdType, number>()
  68. static inspectSignals: Array<AbortSignal | undefined> = []
  69. static snapshotSignals: Array<AbortSignal | undefined> = []
  70. static loadEffect: ((entry: { meta: SessionHeader; events: SessionEvent[] }) => void) | undefined
  71. static inspectEffect: ((
  72. entry: { meta: SessionHeader; events: SessionEvent[] },
  73. signal?: AbortSignal,
  74. ) => void | Promise<void>) | undefined
  75. static listGate: Promise<void> | undefined
  76. static listStarted: (() => void) | undefined
  77. static snapshotEffect: ((signal?: AbortSignal) => void | Promise<void>) | undefined
  78. static snapshotOverride: (() => SessionPersistenceSnapshot[]) | undefined
  79. static failure: unknown
  80. locate(_meta: SessionHeader): undefined {
  81. return undefined
  82. }
  83. static reset(entries: readonly { meta: SessionHeader; events: SessionEvent[] }[] = []): void {
  84. this.entries = new Map()
  85. this.revisions = new Map()
  86. this.loads = new Map()
  87. this.inspections = new Map()
  88. this.inspectSignals = []
  89. this.snapshotSignals = []
  90. this.loadEffect = undefined
  91. this.inspectEffect = undefined
  92. for (const entry of entries) this.set(entry)
  93. this.listGate = undefined
  94. this.listStarted = undefined
  95. this.snapshotEffect = undefined
  96. this.snapshotOverride = undefined
  97. this.failure = undefined
  98. }
  99. static set(entry: { meta: SessionHeader; events: SessionEvent[] }): void {
  100. this.entries.set(entry.meta.id, structuredClone(entry))
  101. this.revisions.set(entry.meta.id, ++this.nextRevision)
  102. }
  103. create(meta: SessionHeader): Promise<void> {
  104. TestPersistence.set({ meta, events: [] })
  105. return Promise.resolve()
  106. }
  107. append(id: SessionIdType, events: readonly SessionEvent[]): Promise<void> {
  108. const entry = TestPersistence.entries.get(id)
  109. if (entry === undefined) return Promise.reject(new Error('missing test session'))
  110. entry.events.push(...structuredClone(events))
  111. TestPersistence.revisions.set(id, ++TestPersistence.nextRevision)
  112. return Promise.resolve()
  113. }
  114. async load(id: SessionIdType): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  115. TestPersistence.loads.set(id, (TestPersistence.loads.get(id) ?? 0) + 1)
  116. if (TestPersistence.failure !== undefined) throw TestPersistence.failure
  117. const entry = TestPersistence.entries.get(id)
  118. if (entry === undefined) throw new Error('missing test session')
  119. if (TestPersistence.loadEffect !== undefined) {
  120. const effect = TestPersistence.loadEffect
  121. TestPersistence.loadEffect = undefined
  122. effect(entry)
  123. TestPersistence.revisions.set(id, ++TestPersistence.nextRevision)
  124. }
  125. return structuredClone(entry)
  126. }
  127. async inspect(id: SessionIdType, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  128. TestPersistence.inspections.set(id, (TestPersistence.inspections.get(id) ?? 0) + 1)
  129. TestPersistence.inspectSignals.push(signal)
  130. if (TestPersistence.failure !== undefined) throw TestPersistence.failure
  131. const entry = TestPersistence.entries.get(id)
  132. if (entry === undefined) throw new Error('missing test session')
  133. await TestPersistence.inspectEffect?.(entry, signal)
  134. TestPersistence.inspectEffect = undefined
  135. return structuredClone(entry)
  136. }
  137. async readFrom(id: SessionIdType, fromSeq: number, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  138. const whole = await this.inspect(id, signal)
  139. return { meta: whole.meta, events: whole.events.filter(event => event.seq >= fromSeq) }
  140. }
  141. async list(): Promise<SessionHeader[]> {
  142. TestPersistence.listStarted?.()
  143. await TestPersistence.listGate
  144. if (TestPersistence.failure !== undefined) throw TestPersistence.failure
  145. return [...TestPersistence.entries.values()].map(entry => structuredClone(entry.meta))
  146. }
  147. async listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]> {
  148. TestPersistence.snapshotSignals.push(signal)
  149. TestPersistence.listStarted?.()
  150. await TestPersistence.listGate
  151. if (TestPersistence.failure !== undefined) throw TestPersistence.failure
  152. const snapshots = TestPersistence.snapshotOverride?.()
  153. ?? [...TestPersistence.entries.values()].map(entry => ({
  154. header: structuredClone(entry.meta),
  155. revision: SessionPersistenceRevision(`test:${TestPersistence.revisions.get(entry.meta.id)}`),
  156. }))
  157. await TestPersistence.snapshotEffect?.(signal)
  158. return snapshots
  159. }
  160. }
  161. async function liveContext(config: ConstructorParameters<typeof SqliteSessionQueryEngine>[1] = { path: ':memory:' }): Promise<Context> {
  162. const ctx = new Context()
  163. await ctx.plugin(SessionStore)
  164. await ctx.plugin(SqliteSessionQueryEngine, config)
  165. return ctx
  166. }
  167. describe('SQLite session search', () => {
  168. it('defaults and validates opening policy and persisted inspection concurrency through its Cordis config', async () => {
  169. const defaultCtx = await liveContext()
  170. expect((defaultCtx.sessionQuery as SqliteSessionQueryEngine).config.openAt).toBe('startup')
  171. expect((defaultCtx.sessionQuery as SqliteSessionQueryEngine).config.persistedInspectConcurrency)
  172. .toBe(SESSION_QUERY_DEFAULT_PERSISTED_INSPECT_CONCURRENCY)
  173. const configuredValue = 2
  174. const configured = new SqliteSessionQueryEngine.Config({
  175. path: ':memory:',
  176. openAt: 'first-search',
  177. persistedInspectConcurrency: configuredValue,
  178. })
  179. expect(configured.openAt).toBe('first-search')
  180. expect(configured.persistedInspectConcurrency).toBe(configuredValue)
  181. const configuredCtx = await liveContext(configured)
  182. expect((configuredCtx.sessionQuery as SqliteSessionQueryEngine).config.persistedInspectConcurrency)
  183. .toBe(configuredValue)
  184. for (const persistedInspectConcurrency of [0, Number.MAX_SAFE_INTEGER + 1]) {
  185. expect(() => new SqliteSessionQueryEngine.Config({
  186. path: ':memory:',
  187. persistedInspectConcurrency,
  188. })).toThrow()
  189. }
  190. expect(() => new SqliteSessionQueryEngine.Config({
  191. path: ':memory:',
  192. openAt: 'later' as never,
  193. })).toThrow()
  194. })
  195. it('mounts and disposes first-search mode without opening its database', async () => {
  196. const path = await temporaryPath('unopened.db')
  197. const ctx = new Context()
  198. await ctx.plugin(SessionStore)
  199. const search = await ctx.plugin(SqliteSessionQueryEngine, {
  200. path,
  201. openAt: 'first-search',
  202. })
  203. await expect(stat(path)).rejects.toMatchObject({ code: 'ENOENT' })
  204. await search.dispose()
  205. await expect(stat(path)).rejects.toMatchObject({ code: 'ENOENT' })
  206. })
  207. it('refuses search in never mode while inherited reads and traces keep working', async () => {
  208. const path = await temporaryPath('never-mode.db')
  209. const ctx = new Context()
  210. await ctx.plugin(SessionStore)
  211. const search = await ctx.plugin(SqliteSessionQueryEngine, { path, openAt: 'never' })
  212. const service = ctx.sessionQuery as SqliteSessionQueryEngine
  213. expect(service.config.openAt).toBe('never')
  214. const parent = SessionId('never-parent')
  215. const child = SessionId('never-child')
  216. ctx.sessions.create(parent, { seed: messageEvents('never opened needle'), meta: { createdAt: 10 } })
  217. ctx.sessions.create(child, { meta: { parentSession: parent, createdAt: 20 } })
  218. await expect(service.searchSessions({ query: 'needle' }))
  219. .rejects.toThrow(expectCode('SESSION_QUERY_SEARCH_DISABLED'))
  220. await expect(service.searchEvents({ sessionId: parent, query: 'needle' }))
  221. .rejects.toThrow(expectCode('SESSION_QUERY_SEARCH_DISABLED'))
  222. expect((await service.listSessions()).map(record => record.header.id).sort())
  223. .toEqual([child, parent])
  224. const lineage = await service.traceSession(parent)
  225. expect(lineage.complete).toBe(true)
  226. expect(lineage.descendants.map(node => node.session.header.id)).toEqual([child])
  227. // The disabled index never touches the filesystem, in mount, use, or disposal.
  228. await expect(stat(path)).rejects.toMatchObject({ code: 'ENOENT' })
  229. await search.dispose()
  230. await expect(stat(path)).rejects.toMatchObject({ code: 'ENOENT' })
  231. })
  232. it('opens once on the first search and reuses readiness for later searches', async () => {
  233. const ctx = new Context()
  234. await ctx.plugin(SessionStore)
  235. await ctx.plugin(SqliteSessionQueryEngine, {
  236. path: ':memory:',
  237. openAt: 'first-search',
  238. })
  239. const service = ctx.sessionQuery as SqliteSessionQueryEngine
  240. const internals = service as unknown as { _open(): Promise<void> }
  241. const open = vi.spyOn(internals, '_open')
  242. await expect(service.searchSessions({ query: 'first' })).resolves.toEqual({ items: [] })
  243. await expect(service.searchSessions({ query: 'second' })).resolves.toEqual({ items: [] })
  244. expect(open).toHaveBeenCalledOnce()
  245. })
  246. it('shares one readiness promise across concurrent first searches', async () => {
  247. const ctx = new Context()
  248. await ctx.plugin(SessionStore)
  249. await ctx.plugin(SqliteSessionQueryEngine, {
  250. path: ':memory:',
  251. openAt: 'first-search',
  252. })
  253. const service = ctx.sessionQuery as SqliteSessionQueryEngine
  254. const internals = service as unknown as { _open(): Promise<void> }
  255. const originalOpen = internals._open.bind(internals)
  256. const release = Promise.withResolvers<undefined>()
  257. const started = Promise.withResolvers<undefined>()
  258. const open = vi.spyOn(internals, '_open').mockImplementation(async () => {
  259. started.resolve(undefined)
  260. await release.promise
  261. await originalOpen()
  262. })
  263. const first = service.searchSessions({ query: 'first' })
  264. const second = service.searchSessions({ query: 'second' })
  265. await started.promise
  266. expect(open).toHaveBeenCalledOnce()
  267. release.resolve(undefined)
  268. await expect(Promise.all([first, second])).resolves.toEqual([
  269. { items: [] },
  270. { items: [] },
  271. ])
  272. expect(open).toHaveBeenCalledOnce()
  273. })
  274. it('searches two-character Unicode61 tokens in live-only sessions', async () => {
  275. const ctx = await liveContext({ path: ':memory:', snippetChars: 20 })
  276. const session = ctx.sessions.create(SessionId('live'), {
  277. // agentPreset rides along: the index rebuilds the header a caller reads,
  278. // and a session listed under the wrong composition is a lie about what it
  279. // ran. The full-header comparison below is what pins every column.
  280. meta: { cwd: '/work', createdAt: 10, seedLength: 1, delegationDepth: 2, agentPreset: 'minimal' },
  281. })
  282. session.append(
  283. 'user/message',
  284. createUserMessage({
  285. content: [{ type: 'text', text: 'An AI helper' }], source: { kind: 'user' },
  286. }),
  287. { surfaceOp: 'append' },
  288. )
  289. await expect(ctx.sessionQuery.searchEvents({ sessionId: session.id, query: 'AI' }))
  290. .resolves.toMatchObject({
  291. session: { ...session.header, seedLength: 1 },
  292. items: [{ sessionId: session.id, seq: 0, snippet: 'An AI helper' }],
  293. })
  294. await expect(ctx.sessionQuery.searchSessions({ query: 'AI' }))
  295. .resolves.toMatchObject({ items: [{ header: { ...session.header, seedLength: 1 }, live: true, persisted: false }] })
  296. })
  297. it('excludes assistant reasoning while indexing visible answer text', async () => {
  298. const ctx = await liveContext()
  299. const session = ctx.sessions.create(SessionId('reasoning'))
  300. session.append(
  301. 'assistant/message',
  302. {
  303. turn: 1,
  304. step: 1,
  305. message: createAssistantMessage({
  306. content: [
  307. { type: 'reasoning', text: 'private-chain-marker' },
  308. { type: 'text', text: 'visible-answer-marker' },
  309. ],
  310. source: { provider: 'mock', model: 'mock' },
  311. }),
  312. },
  313. { surfaceOp: 'append' },
  314. )
  315. await expect(ctx.sessionQuery.searchSessions({ query: 'private-chain-marker' }))
  316. .resolves.toEqual({ items: [] })
  317. await expect(ctx.sessionQuery.searchSessions({ query: 'visible-answer-marker' }))
  318. .resolves.toMatchObject({
  319. items: [{
  320. header: { id: session.id },
  321. bestMatch: { snippet: 'visible-answer-marker' },
  322. }],
  323. })
  324. })
  325. it('searches all surfaces by default and applies metadata before ranking', async () => {
  326. const ctx = await liveContext({ path: ':memory:', defaultLimit: 10, maxLimit: 20 })
  327. const parent = SessionId('parent')
  328. const events: SessionEvent[] = [
  329. { type: 'user/message', seq: 0, time: 10, data: createUserMessage({
  330. content: [{ type: 'text', text: 'needle original' }], source: { kind: 'user' },
  331. }), surfaceOp: 'append' },
  332. { type: 'assistant/chunk', seq: 1, time: 11, data: { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'needle raw' } } },
  333. { type: 'user/message', seq: 2, time: 12, data: createUserMessage({
  334. content: [{ type: 'text', text: 'needle summary' }], source: { kind: 'plugin', plugin: 'test' },
  335. }), surfaceOp: { op: 'replace', start: 0, end: 0 }, sourceEventSeqs: [0] },
  336. { type: 'turn/end', seq: 3, time: 13, data: { turn: 1, reason: { kind: 'error', error: { message: 'needle failure', code: 'UNKNOWN' } } } },
  337. ]
  338. ctx.sessions.create(SessionId('a'), { seed: events, meta: { cwd: '/a', parentSession: parent, createdAt: 20 } })
  339. ctx.sessions.create(SessionId('b'), { seed: messageEvents('needle peer', 12), meta: { createdAt: 20 } })
  340. const all = await ctx.sessionQuery.searchEvents({ sessionId: SessionId('a'), query: 'needle' })
  341. expect(new Set(all.items.map(item => item.surface))).toEqual(new Set(['current', 'shadowed', 'log-only']))
  342. await expect(ctx.sessionQuery.searchEvents({
  343. sessionId: SessionId('a'),
  344. query: 'needle',
  345. filters: [
  346. { kind: 'seq', from: 2, to: 2 },
  347. { kind: 'time', from: 12, to: 12 },
  348. { kind: 'type', values: ['user/message'] },
  349. { kind: 'surface', values: ['current'] },
  350. ],
  351. })).resolves.toMatchObject({ items: [{ seq: 2, surface: 'current' }] })
  352. const grouped = await ctx.sessionQuery.searchSessions({
  353. query: 'needle',
  354. sessionFilters: [
  355. { kind: 'id', values: [SessionId('a')] },
  356. { kind: 'cwd', values: ['/a'] },
  357. { kind: 'created-at', from: 20, to: 20 },
  358. { kind: 'parent', values: [parent] },
  359. { kind: 'availability', values: ['live'] },
  360. ],
  361. eventFilters: [{ kind: 'surface', values: ['shadowed'] }],
  362. })
  363. expect(grouped.items).toHaveLength(1)
  364. expect(grouped.items[0]).toMatchObject({
  365. header: { id: SessionId('a'), cwd: '/a', parentSession: parent },
  366. live: true,
  367. persisted: false,
  368. bestMatch: { seq: 0, surface: 'shadowed' },
  369. })
  370. })
  371. it('searches at the supported FTS5 outer-predicate boundary in both scopes', async () => {
  372. const ctx = await liveContext()
  373. const session = ctx.sessions.create(SessionId('predicate-boundary'), {
  374. seed: messageEvents('needle'),
  375. meta: { cwd: '/work' },
  376. })
  377. const sessionFilters = Array.from(
  378. { length: 14 },
  379. () => ({ kind: 'cwd' as const, values: ['/work', null] }),
  380. )
  381. const eventFilters = Array.from(
  382. { length: 13 },
  383. () => ({ kind: 'type' as const, values: ['user/message' as const] }),
  384. )
  385. await expect(ctx.sessionQuery.searchSessions({ query: 'needle', sessionFilters }))
  386. .resolves.toMatchObject({ items: [{ header: { id: session.id } }] })
  387. await expect(ctx.sessionQuery.searchEvents({
  388. sessionId: session.id,
  389. query: 'needle',
  390. filters: eventFilters,
  391. })).resolves.toMatchObject({ items: [{ sessionId: session.id, seq: 0 }] })
  392. })
  393. it('rejects unsupported FTS5 outer-predicate counts with typed errors', async () => {
  394. const ctx = await liveContext()
  395. const session = ctx.sessions.create(SessionId('predicate-limit'), { seed: messageEvents('needle') })
  396. const sessionFilters = Array.from(
  397. { length: 1_100 },
  398. () => ({ kind: 'id' as const, values: [session.id] }),
  399. )
  400. const eventFilters = Array.from(
  401. { length: 1_100 },
  402. () => ({ kind: 'type' as const, values: ['user/message' as const] }),
  403. )
  404. await expect(ctx.sessionQuery.searchSessions({ query: 'needle', sessionFilters }))
  405. .rejects.toThrow(expectCode('SESSION_QUERY_INVALID_FILTER'))
  406. await expect(ctx.sessionQuery.searchEvents({
  407. sessionId: session.id,
  408. query: 'needle',
  409. filters: eventFilters,
  410. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_FILTER'))
  411. await expect(ctx.sessionQuery.searchSessions({
  412. query: 'needle',
  413. sessionFilters: sessionFilters.slice(0, 7),
  414. eventFilters: eventFilters.slice(0, 8),
  415. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_FILTER'))
  416. await expect(ctx.sessionQuery.searchEvents({
  417. sessionId: session.id,
  418. query: 'needle',
  419. filters: eventFilters.slice(0, 14),
  420. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_FILTER'))
  421. })
  422. it('uses literal phrase tokens, stable ties, and bounded Unicode snippets', async () => {
  423. const ctx = await liveContext({ path: ':memory:', defaultLimit: 10, maxLimit: 10, snippetChars: 5 })
  424. ctx.sessions.create(SessionId('a'), { seed: messageEvents('😀😀 alpha beta BRAID 😀😀', 10), meta: { createdAt: 1 } })
  425. ctx.sessions.create(SessionId('b'), { seed: messageEvents('alpha beta', 10), meta: { createdAt: 1 } })
  426. ctx.sessions.create(SessionId('c'), { seed: messageEvents('alpha middle beta', 10), meta: { createdAt: 1 } })
  427. ctx.sessions.create(SessionId('d'), { seed: messageEvents('alpha beta', 10), meta: { createdAt: 1 } })
  428. ctx.sessions.create(SessionId('operator'), { seed: messageEvents('needle OR absent', 10), meta: { createdAt: 1 } })
  429. ctx.sessions.create(SessionId('only'), { seed: messageEvents('needle only', 10), meta: { createdAt: 1 } })
  430. ctx.sessions.create(SessionId('quote'), { seed: messageEvents('say "needle" exactly', 10), meta: { createdAt: 1 } })
  431. const phrase = await ctx.sessionQuery.searchSessions({ query: 'alpha beta' })
  432. expect(phrase.items.map(item => item.header.id)).toEqual([SessionId('b'), SessionId('d'), SessionId('a')])
  433. expect(phrase.items.every(item => Array.from(item.bestMatch.snippet).length <= 5)).toBe(true)
  434. await expect(ctx.sessionQuery.searchSessions({ query: 'AI' })).resolves.toEqual({ items: [] })
  435. await expect(ctx.sessionQuery.searchSessions({ query: 'needle OR absent' }))
  436. .resolves.toMatchObject({ items: [{ header: { id: SessionId('operator') } }] })
  437. await expect(ctx.sessionQuery.searchSessions({ query: 'say "needle"' }))
  438. .resolves.toMatchObject({ items: [{ header: { id: SessionId('quote') } }] })
  439. await expect(ctx.sessionQuery.searchSessions({ query: '*' })).resolves.toEqual({ items: [] })
  440. })
  441. it('ranks live and persisted matches on one source-comparable contract', async () => {
  442. const persisted = header('z-persisted')
  443. TestPersistence.reset([
  444. { meta: persisted, events: messageEvents('needle needle', 10) },
  445. ...Array.from({ length: 12 }, (_, index) => ({
  446. meta: header(`filler-${index}`),
  447. events: messageEvents('needle', 10),
  448. })),
  449. ])
  450. const ctx = await liveContext()
  451. const persistence = await ctx.plugin(TestPersistence)
  452. ctx.sessions.create(SessionId('a-live'), {
  453. seed: messageEvents('needle needle', 10),
  454. meta: { createdAt: persisted.createdAt },
  455. })
  456. const result = await ctx.sessionQuery.searchSessions({
  457. query: 'needle',
  458. sessionFilters: [{ kind: 'id', values: [SessionId('a-live'), persisted.id] }],
  459. })
  460. expect(result.items.map(item => item.header.id)).toEqual([SessionId('a-live'), persisted.id])
  461. await persistence.dispose()
  462. })
  463. it('positions snippets from FTS5 matches across diacritics and punctuation', async () => {
  464. const ctx = await liveContext({ path: ':memory:', snippetChars: 14 })
  465. const session = ctx.sessions.create(SessionId('snippet'), {
  466. seed: messageEvents('long long long—café,\nnext value', 10),
  467. })
  468. const page = await ctx.sessionQuery.searchEvents({ sessionId: session.id, query: 'CAFE' })
  469. expect(page.items).toHaveLength(1)
  470. expect(page.items[0]!.snippet).toContain('café')
  471. expect(page.items[0]!.snippet).toContain('—')
  472. expect(page.items[0]!.snippet).not.toContain('\n')
  473. expect(Array.from(page.items[0]!.snippet).length).toBeLessThanOrEqual(14)
  474. })
  475. it('binds cursors to requests and only invalidates within-session pages for target changes', async () => {
  476. const ctx = await liveContext({ path: ':memory:', defaultLimit: 1, maxLimit: 5 })
  477. const target = ctx.sessions.create(SessionId('target'), {
  478. seed: [
  479. ...messageEvents('needle one', 10),
  480. { ...messageEvents('needle two', 11)[0]!, seq: 1 },
  481. { ...messageEvents('needle three', 12)[0]!, seq: 2 },
  482. ],
  483. })
  484. ctx.sessions.create(SessionId('other'), { seed: messageEvents('needle other', 10) })
  485. const eventPage = await ctx.sessionQuery.searchEvents({ sessionId: target.id, query: 'needle', limit: 1 })
  486. const sessionPage = await ctx.sessionQuery.searchSessions({ query: 'needle', limit: 1 })
  487. expect(eventPage.nextCursor).toEqual(expect.any(String))
  488. expect(sessionPage.nextCursor).toEqual(expect.any(String))
  489. if (eventPage.nextCursor === undefined || sessionPage.nextCursor === undefined) throw new Error('expected cursors')
  490. const unsafeOffsetCursor = replaceCursorOffset(eventPage.nextCursor, 1e100)
  491. await expect(ctx.sessionQuery.searchEvents({
  492. sessionId: target.id,
  493. query: 'needle',
  494. limit: 1,
  495. cursor: unsafeOffsetCursor,
  496. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_CURSOR'))
  497. const eventKeys = eventPage.items.map(item => `${item.sessionId}:${item.seq}`)
  498. let eventCursor: ReturnType<typeof SessionSearchCursor> | undefined = eventPage.nextCursor
  499. while (eventCursor !== undefined) {
  500. const next = await ctx.sessionQuery.searchEvents({
  501. sessionId: target.id,
  502. query: 'needle',
  503. limit: 1,
  504. cursor: eventCursor,
  505. })
  506. eventKeys.push(...next.items.map(item => `${item.sessionId}:${item.seq}`))
  507. eventCursor = next.nextCursor
  508. }
  509. expect(eventKeys).toHaveLength(3)
  510. expect(new Set(eventKeys).size).toBe(eventKeys.length)
  511. const sessionIds = sessionPage.items.map(item => item.header.id)
  512. let sessionCursor: ReturnType<typeof SessionSearchCursor> | undefined = sessionPage.nextCursor
  513. while (sessionCursor !== undefined) {
  514. const next = await ctx.sessionQuery.searchSessions({ query: 'needle', limit: 1, cursor: sessionCursor })
  515. sessionIds.push(...next.items.map(item => item.header.id))
  516. sessionCursor = next.nextCursor
  517. }
  518. expect(sessionIds).toHaveLength(2)
  519. expect(new Set(sessionIds).size).toBe(sessionIds.length)
  520. ctx.sessions.create(SessionId('unrelated'), { seed: messageEvents('needle unrelated', 20) })
  521. await expect(ctx.sessionQuery.searchEvents({
  522. sessionId: target.id,
  523. query: 'needle',
  524. limit: 1,
  525. cursor: eventPage.nextCursor,
  526. })).resolves.toMatchObject({ items: [{ sessionId: target.id }] })
  527. await expect(ctx.sessionQuery.searchSessions({ query: 'needle', limit: 1, cursor: sessionPage.nextCursor }))
  528. .rejects.toThrow(expectCode('SESSION_QUERY_STALE_CURSOR'))
  529. await expect(ctx.sessionQuery.searchEvents({
  530. sessionId: target.id,
  531. query: 'different',
  532. limit: 1,
  533. cursor: eventPage.nextCursor,
  534. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_CURSOR'))
  535. target.append('user/message', createUserMessage({
  536. content: [{ type: 'text', text: 'needle four' }], source: { kind: 'user' },
  537. }), { surfaceOp: 'append' })
  538. await expect(ctx.sessionQuery.searchEvents({
  539. sessionId: target.id,
  540. query: 'needle',
  541. limit: 1,
  542. cursor: eventPage.nextCursor,
  543. })).rejects.toThrow(expectCode('SESSION_QUERY_STALE_CURSOR'))
  544. })
  545. it('invalidates session cursors after transient persistence topology changes', async () => {
  546. TestPersistence.reset()
  547. const ctx = await liveContext({ path: ':memory:', defaultLimit: 1, maxLimit: 5 })
  548. ctx.sessions.create(SessionId('first'), { seed: messageEvents('needle first') })
  549. ctx.sessions.create(SessionId('second'), { seed: messageEvents('needle second') })
  550. const page = await ctx.sessionQuery.searchSessions({ query: 'needle', limit: 1 })
  551. if (page.nextCursor === undefined) throw new Error('expected cursor')
  552. const persistence = await ctx.plugin(TestPersistence)
  553. await persistence.dispose()
  554. await expect(ctx.sessionQuery.searchSessions({
  555. query: 'needle',
  556. limit: 1,
  557. cursor: page.nextCursor,
  558. })).rejects.toThrow(expectCode('SESSION_QUERY_STALE_CURSOR'))
  559. })
  560. it('rejects invalid requests, filters, cursors, and direct config', async () => {
  561. const ctx = await liveContext({ path: ':memory:', defaultLimit: 2, maxLimit: 3 })
  562. const session = ctx.sessions.create(SessionId('valid'), { seed: messageEvents('needle') })
  563. for (const request of [
  564. { sessionId: session.id, query: '' },
  565. { sessionId: session.id, query: 'needle', limit: 0 },
  566. { sessionId: session.id, query: 'needle', limit: 4 },
  567. { sessionId: session.id, query: 'needle', filters: [{ kind: 'seq', from: 2, to: 1 }] },
  568. { sessionId: session.id, query: 'needle', filters: [{ kind: 'surface', values: ['future'] }] },
  569. { sessionId: session.id, query: 'bad\0query' },
  570. ] as const) {
  571. await expect(ctx.sessionQuery.searchEvents(request as never)).rejects.toBeInstanceOf(Error)
  572. }
  573. await expect(ctx.sessionQuery.searchSessions({
  574. query: 'needle',
  575. sessionFilters: [{ kind: 'availability', values: ['remote' as never] }],
  576. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_FILTER'))
  577. await expect(ctx.sessionQuery.searchSessions({
  578. query: 'needle',
  579. sessionFilters: [{ kind: 'future' } as never],
  580. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_FILTER'))
  581. await expect(ctx.sessionQuery.searchSessions({
  582. query: 'needle',
  583. eventFilters: [{ kind: 'future' } as never],
  584. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_FILTER'))
  585. await expect(ctx.sessionQuery.searchEvents({
  586. sessionId: session.id,
  587. query: 'needle',
  588. filters: [{ kind: 'future' } as never],
  589. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_FILTER'))
  590. await expect(ctx.sessionQuery.searchEvents({
  591. sessionId: session.id,
  592. query: 'needle',
  593. cursor: SessionSearchCursor('not-json'),
  594. }))
  595. .rejects.toThrow(expectCode('SESSION_QUERY_INVALID_CURSOR'))
  596. await expect(ctx.sessionQuery.searchEvents({ sessionId: SessionId('absent'), query: 'needle' }))
  597. .rejects.toThrow(expectCode('SESSION_QUERY_SESSION_NOT_FOUND'))
  598. for (const config of [
  599. { path: '' },
  600. { path: ':memory:', defaultLimit: 0 },
  601. { path: ':memory:', maxLimit: 0 },
  602. { path: ':memory:', defaultLimit: 1e100 },
  603. { path: ':memory:', maxLimit: 1e100 },
  604. { path: ':memory:', snippetChars: 0 },
  605. { path: ':memory:', readWindowMax: -1 },
  606. { path: ':memory:', persistedInspectConcurrency: 0 },
  607. { path: ':memory:', persistedInspectConcurrency: Number.MAX_SAFE_INTEGER + 1 },
  608. { path: ':memory:', defaultLimit: 3, maxLimit: 2 },
  609. { path: ':memory:', openAt: 'later' },
  610. { path: ':memory:', journalMode: 'memory' },
  611. ]) {
  612. const direct = new Context()
  613. await direct.plugin(SessionStore)
  614. expect(() => new SqliteSessionQueryEngine(direct, config as never))
  615. .toThrow(expectCode('SESSION_QUERY_INVALID_CONFIG'))
  616. expect(direct.sessionQuery).toBeUndefined()
  617. }
  618. })
  619. it('rejects aggregate filter bindings above SQLite\'s portable variable limit', async () => {
  620. const ctx = await liveContext()
  621. const session = ctx.sessions.create(SessionId('binding-limit'), { seed: messageEvents('needle') })
  622. // Each clause is below the ceiling; combined with its sibling and fixed
  623. // query bindings, the complete statement is not portable.
  624. const halfPortableLimit = 16_383
  625. const ids = Array.from(
  626. { length: halfPortableLimit },
  627. (_, index) => SessionId(`binding-${index}`),
  628. )
  629. const types = Array.from({ length: halfPortableLimit }, () => 'user/message' as const)
  630. const surfaces = Array.from({ length: halfPortableLimit }, () => 'current' as const)
  631. await expect(ctx.sessionQuery.searchSessions({
  632. query: 'needle',
  633. sessionFilters: [{ kind: 'id', values: ids }],
  634. eventFilters: [{ kind: 'type', values: types }],
  635. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_FILTER'))
  636. await expect(ctx.sessionQuery.searchEvents({
  637. sessionId: session.id,
  638. query: 'needle',
  639. filters: [
  640. { kind: 'type', values: types },
  641. { kind: 'surface', values: surfaces },
  642. ],
  643. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_FILTER'))
  644. })
  645. it('rejects one 125,000-value filter list with a typed error', async () => {
  646. const ctx = await liveContext()
  647. const ids = Array.from(
  648. { length: 125_000 },
  649. (_, index) => SessionId(`oversized-binding-${index}`),
  650. )
  651. await expect(ctx.sessionQuery.searchSessions({
  652. query: 'needle',
  653. sessionFilters: [{ kind: 'id', values: ids }],
  654. })).rejects.toThrow(expectCode('SESSION_QUERY_INVALID_FILTER'))
  655. })
  656. })
  657. describe('SQLite reconciliation and source lifecycle', () => {
  658. it('owns queued request and filter values before waiting for the serializer', async () => {
  659. const durable = header('owned')
  660. TestPersistence.reset([{ meta: durable, events: messageEvents('durable needle') }])
  661. const ctx = await liveContext()
  662. const persistence = await ctx.plugin(TestPersistence)
  663. let release!: () => void
  664. TestPersistence.listGate = new Promise<void>((resolve) => { release = resolve })
  665. let markStarted!: () => void
  666. const started = new Promise<void>((resolve) => { markStarted = resolve })
  667. TestPersistence.listStarted = () => {
  668. TestPersistence.listStarted = undefined
  669. markStarted()
  670. }
  671. const blocking = ctx.sessionQuery.searchSessions({ query: 'needle' })
  672. await started
  673. const availability: SessionAvailability[] = ['persisted']
  674. const request: SessionSearchRequest = {
  675. query: 'needle',
  676. sessionFilters: [{ kind: 'availability', values: availability }],
  677. }
  678. const queued = ctx.sessionQuery.searchSessions(request)
  679. request.query = 'absent'
  680. availability[0] = 'live'
  681. release()
  682. await expect(blocking).resolves.toMatchObject({ items: [{ header: durable }] })
  683. await expect(queued).resolves.toMatchObject({ items: [{ header: durable }] })
  684. await persistence.dispose()
  685. })
  686. it('mounts persistence dynamically, shadows with TEMP live rows, reveals, and hides on unmount', async () => {
  687. const shared = header('shared', 10, { cwd: '/work' })
  688. const durable = header('durable', 5)
  689. TestPersistence.reset([
  690. { meta: shared, events: messageEvents('persisted needle') },
  691. { meta: durable, events: messageEvents('durable needle') },
  692. ])
  693. const ctx = await liveContext()
  694. await expect(ctx.sessionQuery.searchSessions({ query: 'durable' })).resolves.toEqual({ items: [] })
  695. const persistenceFiber = await ctx.plugin(TestPersistence)
  696. await expect(ctx.sessionQuery.searchSessions({ query: 'durable' }))
  697. .resolves.toMatchObject({ items: [{ header: durable, live: false, persisted: true }] })
  698. const live = ctx.sessions.prepare(shared.id, { meta: { createdAt: 10, cwd: '/work' } })
  699. live.append('user/message', createUserMessage({
  700. content: [{ type: 'text', text: 'live needle' }], source: { kind: 'user' },
  701. }), { surfaceOp: 'append' })
  702. const detach = ctx.sessions.enter(live)
  703. ctx.sessions.announce(live)
  704. await expect(ctx.sessionQuery.searchSessions({ query: 'persisted' })).resolves.toEqual({ items: [] })
  705. await expect(ctx.sessionQuery.searchSessions({ query: 'live' }))
  706. .resolves.toMatchObject({ items: [{ header: shared, live: true, persisted: true }] })
  707. detach()
  708. await expect(ctx.sessionQuery.searchSessions({ query: 'persisted' }))
  709. .resolves.toMatchObject({ items: [{ header: shared, live: false, persisted: true }] })
  710. await persistenceFiber.dispose()
  711. await expect(ctx.sessionQuery.searchSessions({ query: 'durable' })).resolves.toEqual({ items: [] })
  712. await expect(ctx.sessionQuery.searchEvents({ sessionId: durable.id, query: 'needle' }))
  713. .rejects.toThrow(expectCode('SESSION_QUERY_SESSION_NOT_FOUND'))
  714. })
  715. it('does not load a persisted log while the same session is live', async () => {
  716. const shared = header('checkpointed-live', 10)
  717. TestPersistence.reset([{ meta: shared, events: messageEvents('persisted needle') }])
  718. const ctx = await liveContext()
  719. const live = ctx.sessions.prepare(shared.id, {
  720. seed: messageEvents('live needle'),
  721. meta: { createdAt: shared.createdAt },
  722. })
  723. const detach = ctx.sessions.enter(live)
  724. ctx.sessions.announce(live)
  725. const persistence = await ctx.plugin(TestPersistence)
  726. await expect(ctx.sessionQuery.searchSessions({
  727. query: 'live',
  728. sessionFilters: [{ kind: 'availability', values: ['persisted'] }],
  729. })).resolves.toMatchObject({
  730. items: [{ header: shared, live: true, persisted: true }],
  731. })
  732. expect(TestPersistence.loads.get(shared.id)).toBeUndefined()
  733. expect(TestPersistence.inspections.get(shared.id)).toBeUndefined()
  734. detach()
  735. await expect(ctx.sessionQuery.searchSessions({ query: 'persisted' }))
  736. .resolves.toMatchObject({ items: [{ header: shared, live: false, persisted: true }] })
  737. expect(TestPersistence.loads.get(shared.id)).toBeUndefined()
  738. expect(TestPersistence.inspections.get(shared.id)).toBe(1)
  739. await persistence.dispose()
  740. })
  741. it('retries when a live owner attaches during persistence observation', async () => {
  742. TestPersistence.reset()
  743. const ctx = await liveContext()
  744. await ctx.plugin(TestPersistence)
  745. TestPersistence.snapshotEffect = () => {
  746. TestPersistence.snapshotEffect = undefined
  747. ctx.sessions.create(SessionId('attached'), { seed: messageEvents('attached needle') })
  748. }
  749. await expect(ctx.sessionQuery.searchSessions({ query: 'attached' }))
  750. .resolves.toMatchObject({ items: [{ header: { id: SessionId('attached') } }] })
  751. })
  752. it('cannot crash-repair a log when live ownership begins during persisted inspection', async () => {
  753. const shared = header('attach-during-inspect', 10)
  754. const persistedEvents = messageEvents('persisted needle')
  755. TestPersistence.reset([{ meta: shared, events: persistedEvents }])
  756. const ctx = await liveContext()
  757. await ctx.plugin(TestPersistence)
  758. TestPersistence.loadEffect = (entry) => {
  759. entry.events = messageEvents('incorrect repair')
  760. }
  761. TestPersistence.inspectEffect = () => {
  762. ctx.sessions.create(shared.id, {
  763. seed: messageEvents('live needle'),
  764. meta: { createdAt: shared.createdAt },
  765. })
  766. }
  767. await expect(ctx.sessionQuery.searchSessions({ query: 'live' }))
  768. .resolves.toMatchObject({ items: [{ header: shared, live: true, persisted: true }] })
  769. expect(TestPersistence.loads.get(shared.id)).toBeUndefined()
  770. expect(TestPersistence.entries.get(shared.id)?.events).toEqual(persistedEvents)
  771. })
  772. it('retries when one live owner replaces another during persistence observation', async () => {
  773. TestPersistence.reset()
  774. const ctx = await liveContext()
  775. const first = ctx.sessions.prepare(SessionId('first'), { seed: messageEvents('first needle') })
  776. const detachFirst = ctx.sessions.enter(first)
  777. ctx.sessions.announce(first)
  778. await ctx.plugin(TestPersistence)
  779. TestPersistence.snapshotEffect = () => {
  780. TestPersistence.snapshotEffect = undefined
  781. detachFirst()
  782. ctx.sessions.create(SessionId('second'), { seed: messageEvents('second needle') })
  783. }
  784. await expect(ctx.sessionQuery.searchSessions({ query: 'second' }))
  785. .resolves.toMatchObject({ items: [{ header: { id: SessionId('second') } }] })
  786. })
  787. it('uses the reconciled persistence binding through the query boundary', async () => {
  788. const durable = header('post-reconcile-unmount')
  789. TestPersistence.reset([{ meta: durable, events: [
  790. ...messageEvents('durable needle', 1),
  791. { ...messageEvents('durable needle again', 2)[0]!, seq: 1 },
  792. ] }])
  793. const ctx = await liveContext({ path: ':memory:', defaultLimit: 1, maxLimit: 2 })
  794. const persistence = await ctx.plugin(TestPersistence)
  795. const internals = ctx.sessionQuery as unknown as {
  796. _reconcile(signal: AbortSignal | undefined): Promise<{
  797. identity: symbol
  798. service?: SessionPersistence
  799. }>
  800. }
  801. const reconcile = internals._reconcile.bind(internals)
  802. const boundary = vi.spyOn(internals, '_reconcile').mockImplementation(async (signal) => {
  803. const binding = await reconcile(signal)
  804. await persistence.dispose()
  805. return binding
  806. })
  807. const page = await ctx.sessionQuery.searchEvents({
  808. sessionId: durable.id,
  809. query: 'needle',
  810. limit: 1,
  811. })
  812. expect(page.items).toMatchObject([{ sessionId: durable.id }])
  813. expect(page.nextCursor).toEqual(expect.any(String))
  814. boundary.mockRestore()
  815. await expect(ctx.sessionQuery.searchEvents({ sessionId: durable.id, query: 'needle' }))
  816. .rejects.toThrow(expectCode('SESSION_QUERY_SESSION_NOT_FOUND'))
  817. })
  818. it('discards a stale list rejection when persistence unmounts during observation', async () => {
  819. const durable = header('racing')
  820. TestPersistence.reset([{ meta: durable, events: messageEvents('durable needle') }])
  821. const ctx = await liveContext()
  822. const persistenceFiber = await ctx.plugin(TestPersistence)
  823. let release!: () => void
  824. TestPersistence.listGate = new Promise<void>((resolve) => { release = resolve })
  825. let markStarted!: () => void
  826. const started = new Promise<void>((resolve) => { markStarted = resolve })
  827. TestPersistence.listStarted = () => {
  828. TestPersistence.listStarted = undefined
  829. markStarted()
  830. }
  831. const search = ctx.sessionQuery.searchSessions({ query: 'needle' })
  832. await started
  833. await persistenceFiber.dispose()
  834. TestPersistence.failure = new Error('stale backend rejection')
  835. release()
  836. await expect(search).resolves.toEqual({ items: [] })
  837. })
  838. it('retries against a replacement after the prior binding rejects', async () => {
  839. const durable = header('replacement')
  840. TestPersistence.reset([{ meta: durable, events: messageEvents('durable needle') }])
  841. const ctx = await liveContext()
  842. const prior = await ctx.plugin(TestPersistence)
  843. let rejectPrior!: (reason: unknown) => void
  844. TestPersistence.listGate = new Promise<void>((_resolve, reject) => { rejectPrior = reject })
  845. let markStarted!: () => void
  846. const started = new Promise<void>((resolve) => { markStarted = resolve })
  847. TestPersistence.listStarted = () => {
  848. TestPersistence.listStarted = undefined
  849. markStarted()
  850. }
  851. const search = ctx.sessionQuery.searchSessions({ query: 'needle' })
  852. await started
  853. await prior.dispose()
  854. TestPersistence.listGate = undefined
  855. const replacement = await ctx.plugin(TestPersistence)
  856. rejectPrior(new Error('stale prior binding'))
  857. await expect(search).resolves.toMatchObject({ items: [{ header: durable }] })
  858. await replacement.dispose()
  859. })
  860. it('reloads a replacement source even when its opaque revisions collide', async () => {
  861. const durable = header('colliding-replacement')
  862. TestPersistence.reset([{ meta: durable, events: messageEvents('old content') }])
  863. const revision = TestPersistence.revisions.get(durable.id)!
  864. const ctx = await liveContext()
  865. const prior = await ctx.plugin(TestPersistence)
  866. await expect(ctx.sessionQuery.searchSessions({ query: 'old' }))
  867. .resolves.toMatchObject({ items: [{ header: durable }] })
  868. await prior.dispose()
  869. TestPersistence.set({ meta: durable, events: messageEvents('new needle') })
  870. TestPersistence.revisions.set(durable.id, revision)
  871. const replacement = await ctx.plugin(TestPersistence)
  872. const page = await ctx.sessionQuery.searchSessions({ query: 'new needle' })
  873. expect(TestPersistence.inspections.get(durable.id)).toBe(2)
  874. expect(page).toMatchObject({ items: [{ header: durable }] })
  875. await expect(ctx.sessionQuery.searchSessions({ query: 'old' })).resolves.toEqual({ items: [] })
  876. expect(TestPersistence.inspections.get(durable.id)).toBe(2)
  877. await replacement.dispose()
  878. })
  879. it('retries when a successful observation belongs to a source unmounted during listing', async () => {
  880. const durable = header('successful-unmount')
  881. TestPersistence.reset([{ meta: durable, events: messageEvents('durable needle') }])
  882. const ctx = await liveContext()
  883. const persistence = await ctx.plugin(TestPersistence)
  884. let lists = 0
  885. TestPersistence.snapshotEffect = async () => {
  886. lists += 1
  887. if (lists === 2) await persistence.dispose()
  888. }
  889. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' })).resolves.toEqual({ items: [] })
  890. expect(lists).toBe(2)
  891. })
  892. it('retries when the snapshot population changes during observation', async () => {
  893. const first = header('first')
  894. const added = header('added-during-list')
  895. TestPersistence.reset([{ meta: first, events: messageEvents('first needle') }])
  896. const ctx = await liveContext()
  897. await ctx.plugin(TestPersistence)
  898. TestPersistence.snapshotEffect = () => {
  899. TestPersistence.snapshotEffect = undefined
  900. TestPersistence.set({ meta: added, events: messageEvents('added needle') })
  901. }
  902. const page = await ctx.sessionQuery.searchSessions({ query: 'needle' })
  903. expect(page.items.map(item => item.header.id).sort()).toEqual([added.id, first.id].sort())
  904. expect(TestPersistence.inspections.get(first.id)).toBe(2)
  905. expect(TestPersistence.inspections.get(added.id)).toBe(1)
  906. })
  907. it('fails after one retry when persistence snapshots keep changing', async () => {
  908. const durable = header('continuous-mutation')
  909. TestPersistence.reset([{ meta: durable, events: messageEvents('durable needle') }])
  910. const ctx = await liveContext()
  911. await ctx.plugin(TestPersistence)
  912. let lists = 0
  913. TestPersistence.snapshotEffect = () => {
  914. lists += 1
  915. TestPersistence.set({ meta: durable, events: messageEvents(`durable needle ${lists}`) })
  916. }
  917. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' }))
  918. .rejects.toThrow(expectCode('SESSION_QUERY_PERSISTENCE_FAILED'))
  919. expect(lists).toBe(4)
  920. })
  921. it('retries if the persistence binding changes while live sessions are observed', async () => {
  922. const durable = header('live-boundary-retry')
  923. TestPersistence.reset([{ meta: durable, events: messageEvents('durable needle') }])
  924. const ctx = await liveContext()
  925. await ctx.plugin(TestPersistence)
  926. const internals = ctx.sessionQuery as unknown as {
  927. _persistenceBinding: { identity: symbol; service?: SessionPersistence }
  928. }
  929. const originalList = ctx.sessions.list.bind(ctx.sessions)
  930. let bumped = false
  931. const list = vi.spyOn(ctx.sessions, 'list').mockImplementation(() => {
  932. if (!bumped) {
  933. bumped = true
  934. internals._persistenceBinding = {
  935. ...internals._persistenceBinding,
  936. identity: Symbol(),
  937. }
  938. }
  939. return originalList()
  940. })
  941. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' }))
  942. .resolves.toMatchObject({ items: [{ header: durable }] })
  943. expect(TestPersistence.inspections.get(durable.id)).toBe(2)
  944. list.mockRestore()
  945. })
  946. it('rejects malformed snapshots and preserves typed persistence failures', async () => {
  947. const durable = header('invalid-snapshot')
  948. TestPersistence.reset([{ meta: durable, events: messageEvents('durable needle') }])
  949. const ctx = await liveContext()
  950. await ctx.plugin(TestPersistence)
  951. TestPersistence.snapshotOverride = () => 'not-an-array' as never
  952. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' }))
  953. .rejects.toThrow(expectCode('SESSION_QUERY_PERSISTENCE_FAILED'))
  954. TestPersistence.snapshotOverride = () => [{ header: durable, revision: 1 as never }]
  955. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' }))
  956. .rejects.toThrow(expectCode('SESSION_QUERY_PERSISTENCE_FAILED'))
  957. TestPersistence.snapshotOverride = () => [
  958. { header: durable, revision: SessionPersistenceRevision('duplicate:1') },
  959. { header: durable, revision: SessionPersistenceRevision('duplicate:2') },
  960. ]
  961. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' }))
  962. .rejects.toThrow(expectCode('SESSION_QUERY_PERSISTENCE_FAILED'))
  963. TestPersistence.snapshotOverride = undefined
  964. const typed = new SessionQueryError('typed persistence failure', 'SESSION_QUERY_PERSISTENCE_FAILED')
  965. TestPersistence.failure = typed
  966. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' })).rejects.toBe(typed)
  967. })
  968. it('rejects immutable header conflicts between live and persisted sources', async () => {
  969. const shared = header('conflict', 10, { delegationDepth: 1 })
  970. TestPersistence.reset([{ meta: shared, events: messageEvents('persisted needle') }])
  971. const ctx = await liveContext()
  972. await ctx.plugin(TestPersistence)
  973. ctx.sessions.create(shared.id, {
  974. seed: messageEvents('live needle'),
  975. meta: { createdAt: 10, delegationDepth: 2 },
  976. })
  977. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' }))
  978. .rejects.toThrow(expectCode('SESSION_QUERY_SOURCE_CONFLICT'))
  979. })
  980. it('preserves unchanged persisted generations while reconciling new, changed, and deleted rows', { timeout: 20_000 }, async () => {
  981. const path = await temporaryPath()
  982. const unchanged = header('unchanged')
  983. const changed = header('changed')
  984. const deleted = header('deleted')
  985. TestPersistence.reset([
  986. { meta: unchanged, events: messageEvents('unchanged needle') },
  987. { meta: changed, events: messageEvents('old needle') },
  988. { meta: deleted, events: messageEvents('deleted needle') },
  989. ])
  990. const first = new Context()
  991. await first.plugin(SessionStore)
  992. const firstPersistence = await first.plugin(TestPersistence)
  993. const firstSearch = await first.plugin(SqliteSessionQueryEngine, { path })
  994. await first.sessionQuery.searchSessions({ query: 'needle' })
  995. expect(Object.fromEntries(TestPersistence.inspections)).toEqual({ unchanged: 1, changed: 1, deleted: 1 })
  996. await first.sessionQuery.searchSessions({ query: 'needle' })
  997. expect(Object.fromEntries(TestPersistence.inspections)).toEqual({ unchanged: 1, changed: 1, deleted: 1 })
  998. await firstSearch.dispose()
  999. await firstPersistence.dispose()
  1000. const beforeDb = new DatabaseSync(path)
  1001. const beforeRows = beforeDb.prepare('SELECT id, generation FROM persisted_sessions ORDER BY id').all() as Array<{ id: string; generation: number }>
  1002. beforeDb.close()
  1003. const before = new Map(beforeRows.map(row => [row.id, row.generation]))
  1004. const added = header('added')
  1005. TestPersistence.entries.delete(deleted.id)
  1006. TestPersistence.set({ meta: changed, events: messageEvents('changed needle') })
  1007. TestPersistence.set({ meta: added, events: messageEvents('added needle') })
  1008. const second = new Context()
  1009. await second.plugin(SessionStore)
  1010. const secondPersistence = await second.plugin(TestPersistence)
  1011. const secondSearch = await second.plugin(SqliteSessionQueryEngine, { path })
  1012. const result = await second.sessionQuery.searchSessions({ query: 'needle' })
  1013. expect(result.items.map(item => item.header.id).sort()).toEqual([added.id, changed.id, unchanged.id].sort())
  1014. expect(Object.fromEntries(TestPersistence.inspections)).toEqual({
  1015. unchanged: 1,
  1016. changed: 2,
  1017. deleted: 1,
  1018. added: 1,
  1019. })
  1020. await secondSearch.dispose()
  1021. await secondPersistence.dispose()
  1022. const afterDb = new DatabaseSync(path)
  1023. const afterRows = afterDb.prepare('SELECT id, generation FROM persisted_sessions ORDER BY id').all() as Array<{ id: string; generation: number }>
  1024. afterDb.close()
  1025. const after = new Map(afterRows.map(row => [row.id, row.generation]))
  1026. expect(after.get(unchanged.id)).toBe(before.get(unchanged.id))
  1027. expect(after.get(changed.id)).toBeGreaterThan(before.get(changed.id)!)
  1028. expect(after.has(deleted.id)).toBe(false)
  1029. expect(after.has(added.id)).toBe(true)
  1030. })
  1031. it('drops connection-local live overlays on reopen and retains persistent bases', async () => {
  1032. const path = await temporaryPath()
  1033. const shared = header('shared', 10)
  1034. TestPersistence.reset([{ meta: shared, events: messageEvents('persisted needle') }])
  1035. const first = new Context()
  1036. await first.plugin(SessionStore)
  1037. const persistence = await first.plugin(TestPersistence)
  1038. const live = first.sessions.create(shared.id, { seed: messageEvents('live needle'), meta: { createdAt: 10 } })
  1039. const search = await first.plugin(SqliteSessionQueryEngine, { path })
  1040. await expect(first.sessionQuery.searchEvents({ sessionId: live.id, query: 'live' })).resolves.toMatchObject({ items: [{}] })
  1041. await search.dispose()
  1042. await persistence.dispose()
  1043. const second = new Context()
  1044. await second.plugin(SessionStore)
  1045. const persistenceAgain = await second.plugin(TestPersistence)
  1046. const searchAgain = await second.plugin(SqliteSessionQueryEngine, { path })
  1047. await expect(second.sessionQuery.searchSessions({ query: 'live' })).resolves.toEqual({ items: [] })
  1048. await expect(second.sessionQuery.searchSessions({ query: 'persisted' }))
  1049. .resolves.toMatchObject({ items: [{ header: shared, live: false, persisted: true }] })
  1050. expect(TestPersistence.inspections.get(shared.id)).toBe(1)
  1051. await searchAgain.dispose()
  1052. await persistenceAgain.dispose()
  1053. })
  1054. it('refreshes after an external mutating load repair without loading from the query path', async () => {
  1055. const durable = header('repair')
  1056. TestPersistence.reset([{ meta: durable, events: messageEvents('before repair') }])
  1057. const ctx = await liveContext()
  1058. const persistence = await ctx.plugin(TestPersistence)
  1059. await expect(ctx.sessionQuery.searchSessions({ query: 'before' }))
  1060. .resolves.toMatchObject({ items: [{ header: durable }] })
  1061. TestPersistence.loadEffect = (entry) => {
  1062. entry.events = messageEvents('repaired needle')
  1063. }
  1064. await ctx.sessionPersistence.load(durable.id)
  1065. await expect(ctx.sessionQuery.searchSessions({ query: 'repaired' }))
  1066. .resolves.toMatchObject({ items: [{ header: durable }] })
  1067. expect(TestPersistence.inspections.get(durable.id)).toBe(2)
  1068. await ctx.sessionQuery.searchSessions({ query: 'repaired' })
  1069. expect(TestPersistence.inspections.get(durable.id)).toBe(2)
  1070. expect(TestPersistence.loads.get(durable.id)).toBe(1)
  1071. await persistence.dispose()
  1072. })
  1073. it('recovers on the next search after source and SQLite transaction failures', async () => {
  1074. TestPersistence.reset([{ meta: header('durable'), events: messageEvents('durable needle') }])
  1075. const ctx = await liveContext()
  1076. await ctx.plugin(TestPersistence)
  1077. TestPersistence.failure = 'offline'
  1078. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' }))
  1079. .rejects.toThrow(expectCode('SESSION_QUERY_PERSISTENCE_FAILED'))
  1080. const signal = new AbortController().signal
  1081. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal }))
  1082. .rejects.toThrow(expectCode('SESSION_QUERY_PERSISTENCE_FAILED'))
  1083. TestPersistence.failure = new Error('still offline')
  1084. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal }))
  1085. .rejects.toThrow(expectCode('SESSION_QUERY_PERSISTENCE_FAILED'))
  1086. TestPersistence.failure = undefined
  1087. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' })).resolves.toMatchObject({ items: [{}] })
  1088. const live = ctx.sessions.create(SessionId('live'), { seed: messageEvents('base') })
  1089. await ctx.sessionQuery.searchEvents({ sessionId: live.id, query: 'base' })
  1090. const db = (ctx.sessionQuery as unknown as { _db: DatabaseSync })._db
  1091. db.exec('PRAGMA query_only = ON')
  1092. live.append('user/message', createUserMessage({
  1093. content: [{ type: 'text', text: 'retry needle' }], source: { kind: 'user' },
  1094. }), { surfaceOp: 'append' })
  1095. await expect(ctx.sessionQuery.searchEvents({ sessionId: live.id, query: 'needle' }))
  1096. .rejects.toThrow(expectCode('SESSION_QUERY_INDEX_FAILED'))
  1097. db.exec('PRAGMA query_only = OFF')
  1098. // seq 2: one-event seed, end-seed, then the live message.
  1099. await expect(ctx.sessionQuery.searchEvents({ sessionId: live.id, query: 'needle' }))
  1100. .resolves.toMatchObject({ items: [{ seq: 2 }] })
  1101. })
  1102. })
  1103. describe('SQLite schema, cancellation, and real persistence integration', () => {
  1104. it('creates a new database and WAL sidecars owner-only without changing its parent mode', async () => {
  1105. if (process.platform === 'win32') return
  1106. const path = await temporaryPath()
  1107. const directory = dirname(path)
  1108. await chmod(directory, 0o755)
  1109. const ctx = await liveContext({ path })
  1110. await ctx.sessionQuery.searchSessions({ query: 'needle' })
  1111. expect((await stat(directory)).mode & 0o777).toBe(0o755)
  1112. expect((await stat(path)).mode & 0o777).toBe(0o600)
  1113. expect((await stat(`${path}-wal`)).mode & 0o777).toBe(0o600)
  1114. expect((await stat(`${path}-shm`)).mode & 0o777).toBe(0o600)
  1115. await (ctx.sessionQuery as SqliteSessionQueryEngine).close()
  1116. })
  1117. it('creates a persistent rollback journal owner-only', async () => {
  1118. if (process.platform === 'win32') return
  1119. const path = await temporaryPath()
  1120. const ctx = await liveContext({ path, journalMode: 'persist' })
  1121. await ctx.sessionQuery.searchSessions({ query: 'needle' })
  1122. expect((await stat(path)).mode & 0o777).toBe(0o600)
  1123. expect((await stat(`${path}-journal`)).mode & 0o777).toBe(0o600)
  1124. await (ctx.sessionQuery as SqliteSessionQueryEngine).close()
  1125. })
  1126. it('preserves the mode of an existing database file', async () => {
  1127. if (process.platform === 'win32') return
  1128. const path = await temporaryPath()
  1129. await writeFile(path, '', { mode: 0o644 })
  1130. await chmod(path, 0o644)
  1131. const ctx = await liveContext({ path, journalMode: 'delete' })
  1132. await ctx.sessionQuery.searchSessions({ query: 'needle' })
  1133. expect((await stat(path)).mode & 0o777).toBe(0o644)
  1134. await (ctx.sessionQuery as SqliteSessionQueryEngine).close()
  1135. })
  1136. it('surfaces filesystem failures while pre-creating the database', async () => {
  1137. const path = `${await temporaryPath()}\0`
  1138. const ctx = new Context()
  1139. await ctx.plugin(SessionStore)
  1140. await expect(ctx.plugin(SqliteSessionQueryEngine, { path })).rejects.toMatchObject({
  1141. code: 'SESSION_QUERY_INDEX_FAILED',
  1142. cause: { code: 'ERR_INVALID_ARG_VALUE' },
  1143. })
  1144. expect(ctx.sessionQuery).toBeUndefined()
  1145. })
  1146. it('resets a recognized incompatible schema but refuses unknown or foreign tables', { timeout: 20_000 }, async () => {
  1147. const stalePath = await temporaryPath('stale.db')
  1148. const staleOwner = await liveContext({ path: stalePath })
  1149. await (staleOwner.sessionQuery as SqliteSessionQueryEngine).close()
  1150. const stale = new DatabaseSync(stalePath)
  1151. stale.exec(`PRAGMA user_version = ${SESSION_QUERY_SQLITE_SCHEMA_VERSION - 1}`)
  1152. stale.close()
  1153. const staleCtx = await liveContext({ path: stalePath })
  1154. staleCtx.sessions.create(SessionId('live'), { seed: messageEvents('needle') })
  1155. await staleCtx.sessionQuery.searchSessions({ query: 'needle' })
  1156. await (staleCtx.sessionQuery as SqliteSessionQueryEngine).close()
  1157. const rebuilt = new DatabaseSync(stalePath)
  1158. expect((rebuilt.prepare('PRAGMA user_version').get() as { user_version: number }).user_version)
  1159. .toBe(SESSION_QUERY_SQLITE_SCHEMA_VERSION)
  1160. rebuilt.close()
  1161. const augmentedPath = await temporaryPath('augmented.db')
  1162. const augmentedOwner = await liveContext({ path: augmentedPath })
  1163. await (augmentedOwner.sessionQuery as SqliteSessionQueryEngine).close()
  1164. const augmented = new DatabaseSync(augmentedPath)
  1165. augmented.exec('CREATE TABLE unrelated(value TEXT)')
  1166. augmented.exec("INSERT INTO unrelated VALUES ('safe')")
  1167. augmented.exec('PRAGMA user_version = 999')
  1168. augmented.close()
  1169. const augmentedCtx = new Context()
  1170. await augmentedCtx.plugin(SessionStore)
  1171. await expect(augmentedCtx.plugin(SqliteSessionQueryEngine, { path: augmentedPath }))
  1172. .rejects.toThrow(expectCode('SESSION_QUERY_INDEX_FAILED'))
  1173. expect(augmentedCtx.sessionQuery).toBeUndefined()
  1174. const stillAugmented = new DatabaseSync(augmentedPath)
  1175. expect(stillAugmented.prepare('SELECT value FROM unrelated').get()).toEqual({ value: 'safe' })
  1176. expect(stillAugmented.prepare('PRAGMA user_version').get()).toEqual({ user_version: 999 })
  1177. stillAugmented.close()
  1178. const currentAugmentedPath = await temporaryPath('current-augmented.db')
  1179. const currentAugmentedOwner = await liveContext({ path: currentAugmentedPath })
  1180. await (currentAugmentedOwner.sessionQuery as SqliteSessionQueryEngine).close()
  1181. const currentAugmented = new DatabaseSync(currentAugmentedPath)
  1182. currentAugmented.exec('CREATE TABLE unrelated(value TEXT)')
  1183. currentAugmented.exec("INSERT INTO unrelated VALUES ('safe')")
  1184. currentAugmented.close()
  1185. const currentAugmentedCtx = new Context()
  1186. await currentAugmentedCtx.plugin(SessionStore)
  1187. await expect(currentAugmentedCtx.plugin(SqliteSessionQueryEngine, {
  1188. path: currentAugmentedPath,
  1189. journalMode: 'delete',
  1190. })).rejects.toThrow(expectCode('SESSION_QUERY_INDEX_FAILED'))
  1191. expect(currentAugmentedCtx.sessionQuery).toBeUndefined()
  1192. const stillCurrentAugmented = new DatabaseSync(currentAugmentedPath)
  1193. expect(stillCurrentAugmented.prepare('SELECT value FROM unrelated').get()).toEqual({ value: 'safe' })
  1194. expect(stillCurrentAugmented.prepare('PRAGMA user_version').get())
  1195. .toEqual({ user_version: SESSION_QUERY_SQLITE_SCHEMA_VERSION })
  1196. expect(stillCurrentAugmented.prepare('PRAGMA journal_mode').get()).toEqual({ journal_mode: 'wal' })
  1197. stillCurrentAugmented.close()
  1198. const foreignPath = await temporaryPath('foreign.db')
  1199. const foreign = new DatabaseSync(foreignPath)
  1200. foreign.exec('PRAGMA journal_mode = WAL')
  1201. foreign.exec('CREATE TABLE canonical(value TEXT)')
  1202. foreign.exec("INSERT INTO canonical VALUES ('safe')")
  1203. foreign.close()
  1204. const foreignCtx = new Context()
  1205. await foreignCtx.plugin(SessionStore)
  1206. await expect(foreignCtx.plugin(SqliteSessionQueryEngine, { path: foreignPath, journalMode: 'delete' }))
  1207. .rejects.toThrow(expectCode('SESSION_QUERY_INDEX_FAILED'))
  1208. expect(foreignCtx.sessionQuery).toBeUndefined()
  1209. const stillForeign = new DatabaseSync(foreignPath)
  1210. expect(stillForeign.prepare('SELECT value FROM canonical').get()).toEqual({ value: 'safe' })
  1211. expect(stillForeign.prepare('PRAGMA journal_mode').get()).toEqual({ journal_mode: 'wal' })
  1212. stillForeign.close()
  1213. const wildcardPath = await temporaryPath('sqlite-wildcard.db')
  1214. const wildcard = new DatabaseSync(wildcardPath)
  1215. wildcard.exec('PRAGMA journal_mode = WAL')
  1216. wildcard.exec('CREATE TABLE sqliteX(value TEXT)')
  1217. wildcard.exec("INSERT INTO sqliteX VALUES ('safe')")
  1218. wildcard.close()
  1219. const wildcardCtx = new Context()
  1220. await wildcardCtx.plugin(SessionStore)
  1221. await expect(wildcardCtx.plugin(SqliteSessionQueryEngine, {
  1222. path: wildcardPath,
  1223. journalMode: 'delete',
  1224. })).rejects.toThrow(expectCode('SESSION_QUERY_INDEX_FAILED'))
  1225. expect(wildcardCtx.sessionQuery).toBeUndefined()
  1226. const stillWildcard = new DatabaseSync(wildcardPath)
  1227. expect(stillWildcard.prepare('SELECT value FROM sqliteX').get()).toEqual({ value: 'safe' })
  1228. expect(stillWildcard.prepare('PRAGMA application_id').get()).toEqual({ application_id: 0 })
  1229. expect(stillWildcard.prepare('PRAGMA user_version').get()).toEqual({ user_version: 0 })
  1230. expect(stillWildcard.prepare('PRAGMA journal_mode').get()).toEqual({ journal_mode: 'wal' })
  1231. stillWildcard.close()
  1232. const otherAppPath = await temporaryPath('other-app.db')
  1233. const otherApp = new DatabaseSync(otherAppPath)
  1234. otherApp.exec('PRAGMA application_id = 123')
  1235. otherApp.close()
  1236. const otherAppCtx = new Context()
  1237. await otherAppCtx.plugin(SessionStore)
  1238. await expect(otherAppCtx.plugin(SqliteSessionQueryEngine, { path: otherAppPath }))
  1239. .rejects.toThrow(expectCode('SESSION_QUERY_INDEX_FAILED'))
  1240. expect(otherAppCtx.sessionQuery).toBeUndefined()
  1241. })
  1242. it('fails plugin initialization without an unhandled rejection or partial service', async () => {
  1243. const path = await temporaryPath('never-queried.db')
  1244. const foreign = new DatabaseSync(path)
  1245. foreign.exec('CREATE TABLE canonical(value TEXT)')
  1246. foreign.close()
  1247. const unhandled: unknown[] = []
  1248. const onUnhandled = (reason: unknown) => { unhandled.push(reason) }
  1249. process.on('unhandledRejection', onUnhandled)
  1250. try {
  1251. const ctx = new Context()
  1252. await ctx.plugin(SessionStore)
  1253. await expect(ctx.plugin(SqliteSessionQueryEngine, { path }))
  1254. .rejects.toThrow(expectCode('SESSION_QUERY_INDEX_FAILED'))
  1255. await new Promise<void>((resolve) => { setImmediate(resolve) })
  1256. expect(unhandled).toEqual([])
  1257. expect(ctx.sessionQuery).toBeUndefined()
  1258. } finally {
  1259. process.off('unhandledRejection', onUnhandled)
  1260. }
  1261. })
  1262. it('defers an invalid database failure only in first-search mode', async () => {
  1263. const path = await temporaryPath('lazy-invalid.db')
  1264. const foreign = new DatabaseSync(path)
  1265. foreign.exec('CREATE TABLE canonical(value TEXT)')
  1266. foreign.close()
  1267. const lazyCtx = new Context()
  1268. await lazyCtx.plugin(SessionStore)
  1269. const lazy = await lazyCtx.plugin(SqliteSessionQueryEngine, {
  1270. path,
  1271. openAt: 'first-search',
  1272. })
  1273. expect(lazyCtx.sessionQuery).toBeInstanceOf(SqliteSessionQueryEngine)
  1274. await expect(lazyCtx.sessionQuery.searchSessions({ query: 'needle' }))
  1275. .rejects.toThrow(expectCode('SESSION_QUERY_INDEX_FAILED'))
  1276. await lazy.dispose()
  1277. const eagerCtx = new Context()
  1278. await eagerCtx.plugin(SessionStore)
  1279. await expect(eagerCtx.plugin(SqliteSessionQueryEngine, { path }))
  1280. .rejects.toThrow(expectCode('SESSION_QUERY_INDEX_FAILED'))
  1281. expect(eagerCtx.sessionQuery).toBeUndefined()
  1282. })
  1283. it.each(['sessions', 'events'] as const)(
  1284. 'forwards one exact reconciliation signal through both snapshot lists and persisted inspection for %s search',
  1285. async (scope) => {
  1286. const durable = header(`signal-${scope}`)
  1287. TestPersistence.reset([{ meta: durable, events: messageEvents('signal needle') }])
  1288. const ctx = await liveContext()
  1289. await ctx.plugin(TestPersistence)
  1290. const controller = new AbortController()
  1291. const result = scope === 'sessions'
  1292. ? await ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
  1293. : await ctx.sessionQuery.searchEvents(
  1294. { sessionId: durable.id, query: 'needle' },
  1295. { signal: controller.signal },
  1296. )
  1297. expect(result.items).toHaveLength(1)
  1298. expect(TestPersistence.snapshotSignals).toEqual([controller.signal, controller.signal])
  1299. expect(TestPersistence.inspectSignals).toEqual([controller.signal])
  1300. },
  1301. )
  1302. it.each(['sessions', 'events'] as const)(
  1303. 'starts no persistence observation for a pre-aborted %s search',
  1304. async (scope) => {
  1305. const durable = header(`pre-aborted-${scope}`)
  1306. TestPersistence.reset([{ meta: durable, events: messageEvents('needle') }])
  1307. const ctx = await liveContext()
  1308. await ctx.plugin(TestPersistence)
  1309. const controller = new AbortController()
  1310. controller.abort(new Error(`pre-aborted ${scope}`))
  1311. const pending = scope === 'sessions'
  1312. ? ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
  1313. : ctx.sessionQuery.searchEvents(
  1314. { sessionId: durable.id, query: 'needle' },
  1315. { signal: controller.signal },
  1316. )
  1317. await expect(pending).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
  1318. expect(TestPersistence.snapshotSignals).toEqual([])
  1319. expect(TestPersistence.inspectSignals).toEqual([])
  1320. },
  1321. )
  1322. it('awaits cooperative snapshot-list cancellation cleanup without starting another observation step', async () => {
  1323. const durable = header('cooperative-list-abort')
  1324. TestPersistence.reset([{ meta: durable, events: messageEvents('needle') }])
  1325. const ctx = await liveContext()
  1326. await ctx.plugin(TestPersistence)
  1327. const started = Promise.withResolvers<AbortSignal>()
  1328. const abortObserved = Promise.withResolvers<undefined>()
  1329. const cleanup = Promise.withResolvers<undefined>()
  1330. TestPersistence.snapshotEffect = async (signal) => {
  1331. TestPersistence.snapshotEffect = undefined
  1332. if (signal === undefined) throw new Error('expected reconciliation signal')
  1333. started.resolve(signal)
  1334. await new Promise<void>((resolve) => {
  1335. signal.addEventListener('abort', () => { resolve() }, { once: true })
  1336. })
  1337. abortObserved.resolve(undefined)
  1338. await cleanup.promise
  1339. signal.throwIfAborted()
  1340. }
  1341. const controller = new AbortController()
  1342. const pending = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
  1343. expect(await started.promise).toBe(controller.signal)
  1344. let settled = false
  1345. void pending.then(
  1346. () => { settled = true },
  1347. () => { settled = true },
  1348. )
  1349. controller.abort(new Error('cooperative list cancellation'))
  1350. await abortObserved.promise
  1351. expect(settled).toBe(false)
  1352. expect(TestPersistence.snapshotSignals).toEqual([controller.signal])
  1353. expect(TestPersistence.inspectSignals).toEqual([])
  1354. cleanup.resolve(undefined)
  1355. await expect(pending).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
  1356. })
  1357. it('keeps a second search serialized while an abort-ignoring snapshot list finishes', async () => {
  1358. const durable = header('serialized-list-abort')
  1359. TestPersistence.reset([{ meta: durable, events: messageEvents('needle') }])
  1360. const ctx = await liveContext()
  1361. await ctx.plugin(TestPersistence)
  1362. const cleanup = Promise.withResolvers<undefined>()
  1363. const started = Promise.withResolvers<undefined>()
  1364. TestPersistence.listGate = cleanup.promise
  1365. TestPersistence.listStarted = () => {
  1366. TestPersistence.listStarted = undefined
  1367. started.resolve(undefined)
  1368. }
  1369. const controller = new AbortController()
  1370. const first = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
  1371. await started.promise
  1372. let firstSettled = false
  1373. let secondSettled = false
  1374. void first.then(
  1375. () => { firstSettled = true },
  1376. () => { firstSettled = true },
  1377. )
  1378. controller.abort(new Error('ignored list cancellation'))
  1379. const second = ctx.sessionQuery.searchEvents({ sessionId: durable.id, query: 'needle' })
  1380. void second.then(
  1381. () => { secondSettled = true },
  1382. () => { secondSettled = true },
  1383. )
  1384. await Promise.resolve()
  1385. expect(firstSettled).toBe(false)
  1386. expect(secondSettled).toBe(false)
  1387. expect(TestPersistence.snapshotSignals).toEqual([controller.signal])
  1388. expect(TestPersistence.inspectSignals).toEqual([])
  1389. cleanup.resolve(undefined)
  1390. await expect(first).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
  1391. await expect(second).resolves.toMatchObject({ items: [{ sessionId: durable.id }] })
  1392. })
  1393. it('awaits an abort-ignoring inspection and starts neither another inspection nor the after-list', async () => {
  1394. const first = header('ignored-inspect-first')
  1395. const second = header('ignored-inspect-second')
  1396. TestPersistence.reset([
  1397. { meta: first, events: messageEvents('first needle') },
  1398. { meta: second, events: messageEvents('second needle') },
  1399. ])
  1400. const ctx = await liveContext()
  1401. await ctx.plugin(TestPersistence)
  1402. const started = Promise.withResolvers<AbortSignal>()
  1403. const cleanup = Promise.withResolvers<undefined>()
  1404. TestPersistence.inspectEffect = async (_entry, signal) => {
  1405. TestPersistence.inspectEffect = undefined
  1406. if (signal === undefined) throw new Error('expected reconciliation signal')
  1407. started.resolve(signal)
  1408. await cleanup.promise
  1409. }
  1410. const controller = new AbortController()
  1411. const pending = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
  1412. expect(await started.promise).toBe(controller.signal)
  1413. let settled = false
  1414. void pending.then(
  1415. () => { settled = true },
  1416. () => { settled = true },
  1417. )
  1418. controller.abort(new Error('ignored inspect cancellation'))
  1419. await Promise.resolve()
  1420. expect(settled).toBe(false)
  1421. expect(TestPersistence.snapshotSignals).toEqual([controller.signal])
  1422. expect(TestPersistence.inspections.get(first.id)).toBe(1)
  1423. expect(TestPersistence.inspections.get(second.id)).toBeUndefined()
  1424. cleanup.resolve(undefined)
  1425. await expect(pending).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
  1426. expect(TestPersistence.snapshotSignals).toEqual([controller.signal])
  1427. expect(TestPersistence.inspections.get(second.id)).toBeUndefined()
  1428. })
  1429. it('cancels both queued and in-flight source waits without committing them', async () => {
  1430. TestPersistence.reset()
  1431. const ctx = await liveContext()
  1432. await ctx.plugin(TestPersistence)
  1433. const boundaryController = new AbortController()
  1434. const boundary = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: boundaryController.signal })
  1435. queueMicrotask(() => { boundaryController.abort() })
  1436. await expect(boundary).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
  1437. const readyController = new AbortController()
  1438. readyController.abort()
  1439. const internals = ctx.sessionQuery as unknown as {
  1440. _ensureReady(signal: AbortSignal): Promise<void>
  1441. }
  1442. await expect(internals._ensureReady(readyController.signal))
  1443. .rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
  1444. let releaseBlocking!: () => void
  1445. TestPersistence.listGate = new Promise<void>((resolve) => { releaseBlocking = resolve })
  1446. let markBlockingStarted!: () => void
  1447. const blockingStarted = new Promise<void>((resolve) => { markBlockingStarted = resolve })
  1448. TestPersistence.listStarted = () => {
  1449. TestPersistence.listStarted = undefined
  1450. markBlockingStarted()
  1451. }
  1452. const blocking = ctx.sessionQuery.searchSessions({ query: 'needle' })
  1453. await blockingStarted
  1454. const queuedController = new AbortController()
  1455. const queued = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: queuedController.signal })
  1456. queuedController.abort()
  1457. await expect(queued).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
  1458. releaseBlocking()
  1459. await expect(blocking).resolves.toEqual({ items: [] })
  1460. TestPersistence.set({
  1461. meta: header('uncommitted'),
  1462. events: messageEvents('durable needle'),
  1463. })
  1464. let releaseActive!: () => void
  1465. TestPersistence.listGate = new Promise<void>((resolve) => { releaseActive = resolve })
  1466. let markActiveStarted!: () => void
  1467. const activeStarted = new Promise<void>((resolve) => { markActiveStarted = resolve })
  1468. TestPersistence.listStarted = () => {
  1469. TestPersistence.listStarted = undefined
  1470. markActiveStarted()
  1471. }
  1472. const activeController = new AbortController()
  1473. const active = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: activeController.signal })
  1474. await activeStarted
  1475. activeController.abort()
  1476. let activeSettled = false
  1477. void active.then(
  1478. () => { activeSettled = true },
  1479. () => { activeSettled = true },
  1480. )
  1481. await Promise.resolve()
  1482. expect(activeSettled).toBe(false)
  1483. releaseActive()
  1484. await expect(active).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
  1485. const db = (ctx.sessionQuery as unknown as { _db: DatabaseSync })._db
  1486. expect(db.prepare('SELECT COUNT(*) AS count FROM persisted_sessions').get()).toEqual({ count: 0 })
  1487. await expect(ctx.sessionQuery.searchSessions({ query: 'needle' }))
  1488. .resolves.toMatchObject({ items: [{ header: { id: SessionId('uncommitted') } }] })
  1489. })
  1490. it.each([
  1491. [new Error('ready error'), 'ready error'],
  1492. ['non-error ready failure', 'session-search dependency rejected with a non-Error value'],
  1493. ])('normalizes a rejected readiness wait before mapping it to an index error', async (failure, detail) => {
  1494. TestPersistence.reset()
  1495. const ctx = await liveContext()
  1496. const internals = ctx.sessionQuery as unknown as {
  1497. _ready: Promise<void>
  1498. _ensureReady(signal: AbortSignal): Promise<void>
  1499. }
  1500. internals._ready = Promise.resolve().then(() => {
  1501. throw failure
  1502. })
  1503. await expect(internals._ensureReady(new AbortController().signal))
  1504. .rejects.toThrow(`session-search SQLite index failed to open: ${detail}`)
  1505. })
  1506. it('checks cancellation after readiness before reconciliation accesses SQLite', async () => {
  1507. TestPersistence.reset()
  1508. const ctx = await liveContext()
  1509. const internals = ctx.sessionQuery as unknown as {
  1510. _db: DatabaseSync
  1511. _ready: Promise<void>
  1512. _ensureReady(signal: AbortSignal | undefined): Promise<void>
  1513. }
  1514. const readiness = Promise.withResolvers<undefined>()
  1515. internals._ready = readiness.promise
  1516. const readyWaitStarted = Promise.withResolvers<undefined>()
  1517. const ensureReady = internals._ensureReady.bind(internals)
  1518. vi.spyOn(internals, '_ensureReady').mockImplementation(async (signal) => {
  1519. const pending = ensureReady(signal)
  1520. readyWaitStarted.resolve(undefined)
  1521. return pending
  1522. })
  1523. const prepare = vi.spyOn(internals._db, 'prepare')
  1524. const reason = new Error('cancelled after readiness')
  1525. const controller = new AbortController()
  1526. const pending = ctx.sessionQuery.searchSessions({ query: 'needle' }, { signal: controller.signal })
  1527. await readyWaitStarted.promise
  1528. const queueBoundaryAbort = readiness.promise.then(() => {
  1529. queueMicrotask(() => { controller.abort(reason) })
  1530. })
  1531. readiness.resolve(undefined)
  1532. await queueBoundaryAbort
  1533. await expect(pending).rejects.toThrow(expectCode('SESSION_QUERY_ABORTED'))
  1534. expect(prepare).not.toHaveBeenCalled()
  1535. })
  1536. it('rejects queued and future work when close waits for an accepted operation', async () => {
  1537. TestPersistence.reset()
  1538. let release!: () => void
  1539. TestPersistence.listGate = new Promise<void>((resolve) => { release = resolve })
  1540. let markStarted!: () => void
  1541. const started = new Promise<void>((resolve) => { markStarted = resolve })
  1542. TestPersistence.listStarted = () => {
  1543. TestPersistence.listStarted = undefined
  1544. markStarted()
  1545. }
  1546. const ctx = await liveContext()
  1547. await ctx.plugin(TestPersistence)
  1548. const search = ctx.sessionQuery as SqliteSessionQueryEngine
  1549. const accepted = search.searchSessions({ query: 'needle' })
  1550. await started
  1551. const queued = search.searchSessions({ query: 'needle' })
  1552. const closing = search.close()
  1553. const repeatedClose = search.close()
  1554. expect(repeatedClose).toBe(closing)
  1555. release()
  1556. await expect(accepted).resolves.toEqual({ items: [] })
  1557. await expect(queued).rejects.toThrow(expectCode('SESSION_QUERY_INDEX_FAILED'))
  1558. await Promise.all([closing, repeatedClose])
  1559. await expect(search.searchSessions({ query: 'needle' }))
  1560. .rejects.toThrow(expectCode('SESSION_QUERY_INDEX_FAILED'))
  1561. expect(search.close()).toBe(closing)
  1562. })
  1563. it('awaits optional-persistence child-fiber quiescence on disposal', async () => {
  1564. TestPersistence.reset()
  1565. const ctx = new Context()
  1566. await ctx.plugin(SessionStore)
  1567. const search = await ctx.plugin(SqliteSessionQueryEngine, { path: ':memory:' })
  1568. const persistence = await ctx.plugin(TestPersistence)
  1569. const optional = (ctx.sessionQuery as unknown as {
  1570. _optionalPersistenceFiber: Fiber
  1571. })._optionalPersistenceFiber
  1572. let release!: () => void
  1573. const cleanup = new Promise<void>((resolve) => { release = resolve })
  1574. optional.ctx.effect(() => () => cleanup)
  1575. let settled = false
  1576. const disposing = search.dispose().then(() => { settled = true })
  1577. await Promise.resolve()
  1578. expect(settled).toBe(false)
  1579. release()
  1580. await disposing
  1581. await persistence.dispose()
  1582. })
  1583. it('combines the real SQLite persistence backend with the real search service keylessly', async () => {
  1584. const persistencePath = await temporaryPath('canonical.db')
  1585. const searchPath = await temporaryPath('derived.db')
  1586. const ctx = new Context()
  1587. await ctx.plugin(SessionStore)
  1588. const persistence = await ctx.plugin(SqliteSessionPersistence, { path: persistencePath })
  1589. const search = await ctx.plugin(SqliteSessionQueryEngine, { path: searchPath })
  1590. const meta = header('real', 10, { cwd: '/work' })
  1591. await ctx.sessionPersistence.create(meta)
  1592. await ctx.sessionPersistence.append(meta.id, messageEvents('real SQLite needle'))
  1593. await expect(ctx.sessionQuery.searchSessions({ query: 'SQLite needle' }))
  1594. .resolves.toMatchObject({ items: [{ header: meta, persisted: true, live: false }] })
  1595. await expect(ctx.sessionQuery.searchEvents({ sessionId: meta.id, query: 'SQLite needle' }))
  1596. .resolves.toMatchObject({ session: meta, items: [{ sessionId: meta.id, seq: 0 }] })
  1597. await expect(ctx.sessionQuery.searchEvents({ sessionId: SessionId('absent'), query: 'needle' }))
  1598. .rejects.toThrow(expectCode('SESSION_QUERY_SESSION_NOT_FOUND'))
  1599. await search.dispose()
  1600. await expect(ctx.sessionPersistence.load(meta.id)).resolves.toMatchObject({ meta, events: [{ seq: 0 }] })
  1601. await persistence.dispose()
  1602. })
  1603. it('reconciles colliding local revisions when a derived index reopens against another SQLite store', async () => {
  1604. const persistencePathA = await temporaryPath('canonical-a.db')
  1605. const persistencePathB = await temporaryPath('canonical-b.db')
  1606. const searchPath = await temporaryPath('derived-collision.db')
  1607. const shared = header('same-id', 10)
  1608. const first = new Context()
  1609. await first.plugin(SessionStore)
  1610. const persistenceA = await first.plugin(SqliteSessionPersistence, { path: persistencePathA })
  1611. await first.sessionPersistence.create(shared)
  1612. await first.sessionPersistence.append(shared.id, messageEvents('alpha source'))
  1613. const inspectA = vi.spyOn(first.sessionPersistence, 'inspect')
  1614. const searchA = await first.plugin(SqliteSessionQueryEngine, { path: searchPath })
  1615. await expect(first.sessionQuery.searchSessions({ query: 'alpha' }))
  1616. .resolves.toMatchObject({ items: [{ header: shared }] })
  1617. expect(inspectA).toHaveBeenCalledTimes(1)
  1618. await searchA.dispose()
  1619. await persistenceA.dispose()
  1620. const reopened = new Context()
  1621. await reopened.plugin(SessionStore)
  1622. const persistenceAAgain = await reopened.plugin(SqliteSessionPersistence, { path: persistencePathA })
  1623. const reopenedInspect = vi.spyOn(reopened.sessionPersistence, 'inspect')
  1624. const searchAAgain = await reopened.plugin(SqliteSessionQueryEngine, { path: searchPath })
  1625. await expect(reopened.sessionQuery.searchSessions({ query: 'alpha' }))
  1626. .resolves.toMatchObject({ items: [{ header: shared }] })
  1627. expect(reopenedInspect).not.toHaveBeenCalled()
  1628. await searchAAgain.dispose()
  1629. await persistenceAAgain.dispose()
  1630. const second = new Context()
  1631. await second.plugin(SessionStore)
  1632. const persistenceB = await second.plugin(SqliteSessionPersistence, { path: persistencePathB })
  1633. await second.sessionPersistence.create(shared)
  1634. await second.sessionPersistence.append(shared.id, messageEvents('bravo source'))
  1635. const inspectB = vi.spyOn(second.sessionPersistence, 'inspect')
  1636. const searchB = await second.plugin(SqliteSessionQueryEngine, { path: searchPath })
  1637. await expect(second.sessionQuery.searchSessions({ query: 'bravo' }))
  1638. .resolves.toMatchObject({ items: [{ header: shared }] })
  1639. await expect(second.sessionQuery.searchSessions({ query: 'alpha' })).resolves.toEqual({ items: [] })
  1640. expect(inspectB).toHaveBeenCalledTimes(1)
  1641. await searchB.dispose()
  1642. await persistenceB.dispose()
  1643. })
  1644. })