sqlite.spec.ts 79 KB

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