sqlite.spec.ts 93 KB

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