1
0

sqlite.spec.ts 85 KB

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