sqlite.spec.ts 79 KB

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