sqlite.spec.ts 46 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043
  1. import { createUserMessage, createMessage } from '@deepseek-ai/dsh-llm'
  2. import { afterEach, describe, expect, it } from 'vitest'
  3. import { Context } from 'cordis'
  4. import { existsSync } from 'node:fs'
  5. import { chmod, mkdtemp, rm, stat, symlink, writeFile } from 'node:fs/promises'
  6. import { tmpdir } from 'node:os'
  7. import { dirname, join } from 'node:path'
  8. import { DatabaseSync } from 'node:sqlite'
  9. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  10. import type { Session, SessionEvent, SurfaceEvent, SurfaceEventType } from '@deepseek-ai/dsh-session'
  11. import SessionPersistenceSqlite, { SCHEMA_VERSION } from '@deepseek-ai/dsh-session-persistence-sqlite'
  12. import {
  13. openDatabase,
  14. rowToEvent,
  15. rowToMeta,
  16. scanRows,
  17. SESSION_PERSISTENCE_SQLITE_APPLICATION_ID,
  18. type EventRow,
  19. } from '../src/schema.ts'
  20. import { runPersistenceContract, meta, oneTurnLog, appendLog } from '../../session-persistence/tests/contract.ts'
  21. import { runCoordinatorContract, type CoordinatorFixture } from '../../session-persistence/tests/coordinator-contract.ts'
  22. const dirs: string[] = []
  23. afterEach(async () => { for (const d of dirs.splice(0)) await rm(d, { recursive: true, force: true }) })
  24. async function expectFlushError(promise: Promise<unknown>, message: RegExp): Promise<void> {
  25. try {
  26. await promise
  27. } catch (error) {
  28. expect(error).toBeInstanceOf(Error)
  29. expect((error as Error).message).toMatch(message)
  30. return
  31. }
  32. throw new Error('expected flush to reject')
  33. }
  34. async function freshDbPath(): Promise<string> {
  35. const dir = await mkdtemp(join(tmpdir(), 'dsh-sqlite-'))
  36. dirs.push(dir)
  37. return join(dir, 'sessions.db')
  38. }
  39. /** Create the exact owned v13 layout without passing through the v14 opener. */
  40. function createV13Database(path: string): DatabaseSync {
  41. const db = new DatabaseSync(path)
  42. db.exec(`
  43. PRAGMA foreign_keys = ON;
  44. CREATE TABLE persistence_state (
  45. singleton INTEGER PRIMARY KEY CHECK (singleton = 1),
  46. store_id TEXT NOT NULL
  47. ) STRICT;
  48. CREATE TABLE sessions (
  49. id TEXT PRIMARY KEY,
  50. version INTEGER NOT NULL,
  51. created_at INTEGER NOT NULL,
  52. cwd TEXT,
  53. parent_session TEXT,
  54. seed_length INTEGER,
  55. origin TEXT,
  56. delegation_depth INTEGER,
  57. incarnation TEXT NOT NULL,
  58. revision INTEGER NOT NULL
  59. ) STRICT;
  60. CREATE TABLE events (
  61. session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
  62. seq INTEGER NOT NULL,
  63. type TEXT NOT NULL,
  64. time INTEGER NOT NULL,
  65. data TEXT NOT NULL,
  66. source_event_seqs TEXT,
  67. surface_op TEXT,
  68. PRIMARY KEY (session_id, seq)
  69. ) STRICT;
  70. PRAGMA application_id = ${SESSION_PERSISTENCE_SQLITE_APPLICATION_ID};
  71. PRAGMA user_version = 13;
  72. `)
  73. db.prepare('INSERT INTO persistence_state (singleton, store_id) VALUES (1, ?)').run('v13-fixture-store')
  74. return db
  75. }
  76. /** A context with the session store + SQLite backend, plus a teardown. */
  77. async function backend(path = ':memory:'): Promise<{ ctx: Context; dispose: () => Promise<void> }> {
  78. const ctx = new Context()
  79. await ctx.plugin(SessionStore)
  80. const fiber = await ctx.plugin(SessionPersistenceSqlite, { path })
  81. return { ctx, dispose: () => fiber.dispose() }
  82. }
  83. // Run the same backend-agnostic contract as JSONL to pin identical semantics.
  84. runPersistenceContract('sqlite', async () => {
  85. const ctx = new Context()
  86. await ctx.plugin(SessionStore)
  87. const fiber = await ctx.plugin(SessionPersistenceSqlite, { path: ':memory:' })
  88. return {
  89. persistence: ctx.sessionPersistence,
  90. dispose: async () => { await fiber.dispose() },
  91. }
  92. })
  93. // A file-backed database lets two mounts share rows across reload. `corruptTail` inserts invalid
  94. // JSON past the committed seq, exercising coordinator repair against real database rows.
  95. runCoordinatorContract('sqlite', async (): Promise<CoordinatorFixture> => {
  96. const dir = await mkdtemp(join(tmpdir(), 'dsh-sqlite-coord-'))
  97. const path = join(dir, 'sessions.db')
  98. return {
  99. mount: async ctx => ctx.plugin(SessionPersistenceSqlite, { path }),
  100. corruptTail: async (id) => {
  101. // A row past the committed region whose `data` does not parse: scanRows
  102. // bounds the preserved prefix at it and returns its seq as tornFrom, which
  103. // the backend surfaces to the coordinator as the tornMarker to delete from.
  104. const db = openDatabase(path, 'wal')
  105. const next = (db.prepare('SELECT COALESCE(MAX(seq), -1) + 1 AS n FROM events WHERE session_id = ?')
  106. .get(id) as { n: number }).n
  107. db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, ?, ?, ?, ?)')
  108. .run(id, next, 'assistant/chunk', 99, '{not valid json')
  109. db.close()
  110. },
  111. cleanup: async () => { await rm(dir, { recursive: true, force: true }) },
  112. }
  113. })
  114. describe('scanRows', () => {
  115. // scanRows works off EventRows (data is a JSON string column); build them from SessionEvents
  116. // so the unit tests read in terms of the event vocabulary. Surface metadata is serialized to
  117. // its nullable columns so the conversion remains faithful.
  118. const rows = (events: SessionEvent[]): EventRow[] =>
  119. events.map((e) => {
  120. const se = e as SessionEvent<SurfaceEventType>
  121. return {
  122. seq: e.seq, type: e.type, time: e.time, data: JSON.stringify(e.data),
  123. source_event_seqs: se.sourceEventSeqs !== undefined ? JSON.stringify(se.sourceEventSeqs) : null,
  124. surface_op: se.surfaceOp !== undefined ? JSON.stringify(se.surfaceOp) : null,
  125. }
  126. })
  127. it('preserves the full log when it ends exactly on a turn/end (no torn tail)', () => {
  128. const { preserved, tornFrom } = scanRows(rows(oneTurnLog()))
  129. expect(preserved).toEqual(oneTurnLog())
  130. expect(tornFrom).toBeUndefined()
  131. })
  132. it('PRESERVES the real events of an interrupted turn after the last turn/end', () => {
  133. // turn 1 committed (0..5) + a crashed turn 2 (turn/start 6, step/start 7, no
  134. // close): all 8 rows are intact, so the whole prefix is preserved and there
  135. // is no torn fragment to delete. (load() then synthesizes the closers.)
  136. const withOpenTurn: SessionEvent[] = [
  137. ...oneTurnLog(),
  138. { type: 'turn/start', seq: 6, time: 7, data: { turn: 2 } },
  139. { type: 'step/start', seq: 7, time: 8, data: { turn: 2, step: 1 } },
  140. ]
  141. const { preserved, tornFrom } = scanRows(rows(withOpenTurn))
  142. expect(preserved.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  143. expect(tornFrom).toBeUndefined()
  144. })
  145. it('preserves the contiguous prefix and flags a torn tail at a seq gap', () => {
  146. // A gap after seq 0 (no committed turn/end): seq 0 is the preserved
  147. // interrupted-turn event; the gap bounds it and marks the torn fragment.
  148. const gapped: SessionEvent[] = [
  149. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
  150. { type: 'step/start', seq: 2, time: 2, data: { turn: 1, step: 1 } }, // seq 1 missing
  151. ]
  152. const { preserved, tornFrom } = scanRows(rows(gapped))
  153. expect(preserved.map(e => e.seq)).toEqual([0])
  154. expect(tornFrom).toBe(1)
  155. })
  156. it('an empty log preserves nothing and has no torn tail', () => {
  157. expect(scanRows([])).toEqual({ preserved: [] })
  158. })
  159. it('throws on a seq gap inside the committed region (before the last turn/end)', () => {
  160. const gapped: SessionEvent[] = [
  161. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
  162. { type: 'step/start', seq: 2, time: 2, data: { turn: 1, step: 1 } }, // seq 1 missing
  163. { type: 'turn/end', seq: 3, time: 3, data: { turn: 1, reason: { kind: 'completed' } } },
  164. ]
  165. expect(() => scanRows(rows(gapped))).toThrow(/seq gap in committed region/)
  166. })
  167. it('throws on an unparsable row inside the committed region', () => {
  168. const withCorruptCommitted: EventRow[] = [
  169. { seq: 0, type: 'turn/start', time: 1, data: '{not json', source_event_seqs: null, surface_op: null }, // corrupt, sits before a turn/end
  170. { seq: 1, type: 'turn/end', time: 2, data: JSON.stringify({ turn: 1, reason: { kind: 'completed' } }), source_event_seqs: null, surface_op: null },
  171. ]
  172. expect(() => scanRows(withCorruptCommitted)).toThrow(/unparsable committed event/)
  173. })
  174. it('tolerates an unparsable torn-tail row after the last turn/end', () => {
  175. const withCorruptTail: EventRow[] = [
  176. ...rows(oneTurnLog()),
  177. { seq: 6, type: 'turn/start', time: 7, data: '{not json', source_event_seqs: null, surface_op: null }, // torn fragment, no committed turn/end after
  178. ]
  179. const { preserved, tornFrom } = scanRows(withCorruptTail)
  180. expect(preserved).toEqual(oneTurnLog())
  181. expect(tornFrom).toBe(6)
  182. })
  183. })
  184. describe('rowToMeta', () => {
  185. it('restores optional origin metadata', () => {
  186. expect(rowToMeta({
  187. id: 'with-origin',
  188. version: 0,
  189. created_at: 1,
  190. cwd: null,
  191. time_zone: 'Asia/Shanghai',
  192. parent_session: null,
  193. seed_length: null,
  194. origin: 'subagent',
  195. incarnation: 'with-origin',
  196. revision: 1,
  197. delegation_depth: null,
  198. })).toMatchObject({ id: 'with-origin', origin: 'subagent', timeZone: 'Asia/Shanghai' })
  199. })
  200. it('rejects fractional stored creation metadata', () => {
  201. expect(() => rowToMeta({
  202. id: 'fractional',
  203. version: 0,
  204. created_at: 1.5,
  205. cwd: null,
  206. time_zone: null,
  207. parent_session: null,
  208. seed_length: null,
  209. origin: null,
  210. incarnation: 'fractional',
  211. revision: 1,
  212. delegation_depth: null,
  213. })).toThrow('stored session createdAt must be a non-negative safe integer')
  214. })
  215. })
  216. describe('SessionPersistenceSqlite: durability and crash semantics', () => {
  217. it('rejects a stored v0 log containing a legacy request/header-delta event', async () => {
  218. const path = await freshDbPath()
  219. const m = meta('legacy-header-delta', '/legacy')
  220. const db = openDatabase(path, 'wal')
  221. db.prepare('INSERT INTO sessions (id, version, created_at, cwd, parent_session, seed_length, delegation_depth, incarnation, revision) VALUES (?, ?, ?, ?, NULL, NULL, NULL, ?, 1)')
  222. .run(m.id, m.version, m.createdAt, m.cwd ?? null, 'legacy-header-delta')
  223. const insert = db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, ?, ?, ?, ?)')
  224. insert.run(m.id, 0, 'turn/start', 1, JSON.stringify({ turn: 1 }))
  225. insert.run(m.id, 1, 'request/header-delta', 2, JSON.stringify({ config: { model: 'legacy' } }))
  226. insert.run(m.id, 2, 'turn/end', 3, JSON.stringify({ turn: 1, reason: { kind: 'completed' } }))
  227. db.close()
  228. const mounted = await backend(path)
  229. await expect(mounted.ctx.sessionPersistence.load(m.id)).rejects.toThrow(/unsupported legacy request\/header-delta event at seq 1/)
  230. await mounted.dispose()
  231. })
  232. it('rejects a stored v0 full header carrying the legacy fallback reason', async () => {
  233. const path = await freshDbPath()
  234. const m = meta('legacy-header-fallback', '/legacy')
  235. const db = openDatabase(path, 'wal')
  236. db.prepare('INSERT INTO sessions (id, version, created_at, cwd, parent_session, seed_length, delegation_depth, incarnation, revision) VALUES (?, ?, ?, ?, NULL, NULL, NULL, ?, 1)')
  237. .run(m.id, m.version, m.createdAt, m.cwd ?? null, 'legacy-header-fallback')
  238. db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, ?, ?, ?, ?)')
  239. .run(m.id, 0, 'request/header', 1, JSON.stringify({
  240. header: { config: { model: 'legacy' } },
  241. reason: 'fallback',
  242. }))
  243. db.close()
  244. const mounted = await backend(path)
  245. await expect(mounted.ctx.sessionPersistence.load(m.id))
  246. .rejects.toThrow(/unsupported legacy request\/header reason "fallback" at seq 0/)
  247. await mounted.dispose()
  248. })
  249. it('has no independent per-session log location', async () => {
  250. const { ctx, dispose } = await backend()
  251. expect(ctx.sessionPersistence.locate(meta('sqlite-location'))).toBeUndefined()
  252. await dispose()
  253. })
  254. it('an interrupted turn (rows after the last turn/end) is PRESERVED and closed during load', async () => {
  255. const path = await freshDbPath()
  256. const m = meta('crash')
  257. // Run 1: persist a complete turn, then a half-written second turn (no turn/end).
  258. const ctx1 = new Context()
  259. await ctx1.plugin(SessionStore)
  260. const fiber1 = await ctx1.plugin(SessionPersistenceSqlite, { path })
  261. await ctx1.sessionPersistence.create(m)
  262. await ctx1.sessionPersistence.append(m.id, oneTurnLog())
  263. await ctx1.sessionPersistence.append(m.id, [
  264. { type: 'turn/start', seq: 6, time: 7, data: { turn: 2 } },
  265. { type: 'step/start', seq: 7, time: 8, data: { turn: 2, step: 1 } },
  266. ])
  267. await fiber1.dispose()
  268. // Run 2: load PRESERVES the interrupted turn's real events (a turn can be huge
  269. // — never truncated) and closes the orphaned turn with synthetic boundary
  270. // events: step/end (the step was open) then turn/end {interrupted}.
  271. const ctx2 = new Context()
  272. await ctx2.plugin(SessionStore)
  273. const fiber2 = await ctx2.plugin(SessionPersistenceSqlite, { path })
  274. const loaded = await ctx2.sessionPersistence.load(m.id)
  275. expect(loaded.events.map(e => e.type)).toEqual([
  276. 'turn/start', 'user/message', 'step/start', 'assistant/message', 'step/end', 'turn/end', // turn 1
  277. 'turn/start', 'step/start', 'step/end', 'turn/end', // turn 2: real events + synthetic closers
  278. ])
  279. expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7, 8, 9])
  280. const last = loaded.events.at(-1)!
  281. expect(last.type === 'turn/end' && last.data.reason).toEqual({ kind: 'interrupted' })
  282. // load durably closed the turn, so the next append continues at the balanced
  283. // length (seq 10) and a reload round-trips identically.
  284. await ctx2.sessionPersistence.append(m.id, [
  285. { type: 'turn/start', seq: 10, time: 9, data: { turn: 3 } },
  286. { type: 'turn/end', seq: 11, time: 10, data: { turn: 3, reason: { kind: 'completed' } } },
  287. ])
  288. const reloaded = await ctx2.sessionPersistence.load(m.id)
  289. expect(reloaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11])
  290. await fiber2.dispose()
  291. })
  292. it('load() durably closes the interrupted turn: the synthetic closers are on disk after load', async () => {
  293. const path = await freshDbPath()
  294. const m = meta('load-closes')
  295. const b1 = await backend(path)
  296. await b1.ctx.sessionPersistence.create(m)
  297. await b1.ctx.sessionPersistence.append(m.id, oneTurnLog()) // seqs 0..5
  298. await b1.dispose()
  299. // Hand-write an interrupted turn (turn/start seq 6, no turn/end).
  300. const db = openDatabase(path, 'wal')
  301. db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, 6, ?, 7, ?)')
  302. .run(m.id, 'turn/start', JSON.stringify({ turn: 2 }))
  303. db.close()
  304. const b2 = await backend(path)
  305. const loaded = await b2.ctx.sessionPersistence.load(m.id)
  306. // turn 2's real turn/start (seq 6) is preserved + a synthetic turn/end (seq 7).
  307. expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  308. expect(loaded.events.at(-1)!.type).toBe('turn/end')
  309. // load() is mutating: the synthetic turn/end MUST be on disk so the stored log
  310. // is balanced and the cursor is truthful (contract: load closes, not defers).
  311. const probe = openDatabase(path, 'wal')
  312. const stored = probe.prepare('SELECT seq, type FROM events WHERE session_id = ? ORDER BY seq').all(m.id) as { seq: number; type: string }[]
  313. probe.close()
  314. expect(stored.map(r => r.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  315. expect(stored.at(-1)!.type).toBe('turn/end')
  316. await b2.dispose()
  317. })
  318. it('all-tail load: a session whose only turn never closed is preserved and closed on load', async () => {
  319. const path = await freshDbPath()
  320. const m = meta('all-tail')
  321. const b1 = await backend(path)
  322. await b1.ctx.sessionPersistence.create(m)
  323. // A first turn that NEVER completed: turn/start + user/message, no turn/end.
  324. await b1.ctx.sessionPersistence.append(m.id, [
  325. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
  326. { type: 'user/message', seq: 1, time: 2, data: createUserMessage({
  327. content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' },
  328. }), surfaceOp: 'append' },
  329. ])
  330. await b1.dispose()
  331. // A fresh backend loads it: the interrupted (only) turn's real events are
  332. // preserved and closed with a synthetic turn/end {interrupted} — NOT
  333. // truncated. The session was materialized, so list() reports it present.
  334. const b2 = await backend(path)
  335. const loaded = await b2.ctx.sessionPersistence.load(m.id)
  336. expect(loaded.events.map(e => e.type)).toEqual(['turn/start', 'user/message', 'turn/end'])
  337. expect(loaded.events.at(-1)!.type === 'turn/end' && loaded.events.at(-1)!.data).toMatchObject({ reason: { kind: 'interrupted' } })
  338. expect((await b2.ctx.sessionPersistence.list()).map(x => x.id)).toContain(m.id)
  339. await b2.dispose()
  340. })
  341. it('rejects opening a database whose schema version is neither v13 nor the current build', async () => {
  342. const path = await freshDbPath()
  343. openDatabase(path, 'wal').close() // stamp user_version = SCHEMA_VERSION
  344. // Bump user_version past what this build supports.
  345. const dbNewer = openDatabase(path, 'wal')
  346. dbNewer.exec(`PRAGMA user_version = ${SCHEMA_VERSION + 1}`)
  347. dbNewer.close()
  348. expect(() => openDatabase(path, 'wal')).toThrow(/incompatible with this build/)
  349. // Versions older than the one explicit migration remain unsupported.
  350. const olderPath = await freshDbPath()
  351. openDatabase(olderPath, 'wal').close()
  352. const dbOlder = openDatabase(olderPath, 'wal')
  353. dbOlder.exec(`PRAGMA user_version = ${SCHEMA_VERSION - 2}`)
  354. dbOlder.close()
  355. expect(() => openDatabase(olderPath, 'wal')).toThrow(/incompatible with this build/)
  356. })
  357. it('atomically migrates an owned v13 fixture and leaves old rows headerless', async () => {
  358. const path = await freshDbPath()
  359. const old = meta('v13-headerless', '/work')
  360. const legacy = createV13Database(path)
  361. legacy.prepare(`
  362. INSERT INTO sessions
  363. (id, version, created_at, cwd, parent_session, seed_length, origin, delegation_depth, incarnation, revision)
  364. VALUES (?, ?, ?, ?, NULL, NULL, NULL, NULL, ?, 1)
  365. `).run(old.id, old.version, old.createdAt, old.cwd ?? null, 'v13-headerless-incarnation')
  366. const insertEvent = legacy.prepare(
  367. 'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op) VALUES (?, ?, ?, ?, ?, ?, ?)',
  368. )
  369. for (const event of oneTurnLog()) {
  370. const surface = event as SessionEvent<SurfaceEventType>
  371. insertEvent.run(
  372. old.id,
  373. event.seq,
  374. event.type,
  375. event.time,
  376. JSON.stringify(event.data),
  377. surface.sourceEventSeqs !== undefined ? JSON.stringify(surface.sourceEventSeqs) : null,
  378. surface.surfaceOp !== undefined ? JSON.stringify(surface.surfaceOp) : null,
  379. )
  380. }
  381. legacy.close()
  382. const migrated = openDatabase(path, 'wal')
  383. expect(migrated.prepare('PRAGMA user_version').get()).toEqual({ user_version: 14 })
  384. expect(migrated.prepare('SELECT time_zone FROM sessions WHERE id = ?').get(old.id))
  385. .toEqual({ time_zone: null })
  386. migrated.close()
  387. const mounted = await backend(path)
  388. try {
  389. const loaded = await mounted.ctx.sessionPersistence.load(old.id)
  390. expect(loaded.meta.timeZone).toBeUndefined()
  391. expect(loaded.events).toEqual(oneTurnLog())
  392. const zoned = meta('v14-zoned', '/work', 'Asia/Shanghai')
  393. await mounted.ctx.sessionPersistence.create(zoned)
  394. await mounted.ctx.sessionPersistence.append(zoned.id, oneTurnLog())
  395. expect((await mounted.ctx.sessionPersistence.load(zoned.id)).meta.timeZone).toBe('Asia/Shanghai')
  396. } finally {
  397. await mounted.dispose()
  398. }
  399. })
  400. it('rejects a spoofed v13 layout without changing its schema or version', async () => {
  401. const path = await freshDbPath()
  402. const malformed = new DatabaseSync(path)
  403. malformed.exec(`
  404. CREATE TABLE sessions (id TEXT);
  405. PRAGMA application_id = ${SESSION_PERSISTENCE_SQLITE_APPLICATION_ID};
  406. PRAGMA user_version = 13;
  407. `)
  408. malformed.close()
  409. expect(() => openDatabase(path, 'wal')).toThrow(/does not match the owned v13 schema/)
  410. const unchanged = new DatabaseSync(path)
  411. const columns = unchanged.prepare('PRAGMA table_info(sessions)').all() as Array<{ name: string }>
  412. expect(columns.map(column => column.name)).toEqual(['id'])
  413. expect(unchanged.prepare('PRAGMA user_version').get()).toEqual({ user_version: 13 })
  414. expect(unchanged.prepare(
  415. "SELECT name FROM sqlite_schema WHERE name IN ('persistence_state', 'events')",
  416. ).all()).toEqual([])
  417. unchanged.close()
  418. })
  419. it('rejects a table-backed unversioned database before stamping or changing journal mode', async () => {
  420. const path = await freshDbPath()
  421. const legacy = new DatabaseSync(path)
  422. legacy.exec('CREATE TABLE sessions (id TEXT PRIMARY KEY)')
  423. legacy.close()
  424. expect(() => openDatabase(path, 'wal')).toThrow(/unversioned schema or application identity/)
  425. const unchanged = new DatabaseSync(path)
  426. expect(unchanged.prepare('PRAGMA user_version').get()).toEqual({ user_version: 0 })
  427. expect(unchanged.prepare('PRAGMA journal_mode').get()).toEqual({ journal_mode: 'delete' })
  428. expect(unchanged.prepare(
  429. "SELECT name FROM sqlite_schema WHERE type = 'table' AND name = 'sessions'",
  430. ).get()).toEqual({ name: 'sessions' })
  431. unchanged.close()
  432. })
  433. it('counts a sqliteX table as user-owned instead of mistaking it for SQLite metadata', async () => {
  434. const path = await freshDbPath()
  435. const unrelated = new DatabaseSync(path)
  436. unrelated.exec('CREATE TABLE sqliteX (value TEXT)')
  437. unrelated.exec("INSERT INTO sqliteX VALUES ('safe')")
  438. unrelated.close()
  439. expect(() => openDatabase(path, 'wal')).toThrow(/unversioned schema or application identity/)
  440. const unchanged = new DatabaseSync(path)
  441. expect(unchanged.prepare('SELECT value FROM sqliteX').get()).toEqual({ value: 'safe' })
  442. expect(unchanged.prepare('PRAGMA application_id').get()).toEqual({ application_id: 0 })
  443. expect(unchanged.prepare('PRAGMA user_version').get()).toEqual({ user_version: 0 })
  444. expect(unchanged.prepare('PRAGMA journal_mode').get()).toEqual({ journal_mode: 'delete' })
  445. unchanged.close()
  446. })
  447. it('rejects view-only and foreign-application unversioned databases without mutation', async () => {
  448. const viewPath = await freshDbPath()
  449. const viewOnly = new DatabaseSync(viewPath)
  450. viewOnly.exec('CREATE VIEW foreign_view AS SELECT 1 AS value')
  451. viewOnly.close()
  452. expect(() => openDatabase(viewPath, 'wal')).toThrow(/unversioned schema or application identity/)
  453. const unchangedView = new DatabaseSync(viewPath)
  454. expect(unchangedView.prepare('PRAGMA journal_mode').get()).toEqual({ journal_mode: 'delete' })
  455. expect(unchangedView.prepare(
  456. "SELECT type FROM sqlite_schema WHERE name = 'foreign_view'",
  457. ).get()).toEqual({ type: 'view' })
  458. unchangedView.close()
  459. const applicationPath = await freshDbPath()
  460. const foreignApplication = new DatabaseSync(applicationPath)
  461. foreignApplication.exec('PRAGMA application_id = 12345')
  462. foreignApplication.close()
  463. expect(() => openDatabase(applicationPath, 'wal')).toThrow(/unversioned schema or application identity/)
  464. const unchangedApplication = new DatabaseSync(applicationPath)
  465. expect(unchangedApplication.prepare('PRAGMA application_id').get()).toEqual({ application_id: 12345 })
  466. expect(unchangedApplication.prepare('PRAGMA user_version').get()).toEqual({ user_version: 0 })
  467. expect(unchangedApplication.prepare('PRAGMA journal_mode').get()).toEqual({ journal_mode: 'delete' })
  468. unchangedApplication.close()
  469. })
  470. it.each([13, SCHEMA_VERSION])('rejects a schema-v%i database with a foreign application identity', async (version) => {
  471. const path = await freshDbPath()
  472. const foreign = new DatabaseSync(path)
  473. foreign.exec('PRAGMA application_id = 12345')
  474. foreign.exec(`PRAGMA user_version = ${version}`)
  475. foreign.close()
  476. expect(() => openDatabase(path, 'wal')).toThrow(/has application id 12345/)
  477. const unchanged = new DatabaseSync(path)
  478. expect(unchanged.prepare('PRAGMA application_id').get()).toEqual({ application_id: 12345 })
  479. expect(unchanged.prepare('PRAGMA user_version').get()).toEqual({ user_version: version })
  480. expect(unchanged.prepare('PRAGMA journal_mode').get()).toEqual({ journal_mode: 'delete' })
  481. unchanged.close()
  482. })
  483. it('rolls back tables created before persistence-state initialization fails', async () => {
  484. const path = await freshDbPath()
  485. const conflicting = new DatabaseSync(path)
  486. conflicting.exec(`PRAGMA application_id = ${SESSION_PERSISTENCE_SQLITE_APPLICATION_ID}`)
  487. conflicting.exec(`PRAGMA user_version = ${SCHEMA_VERSION}`)
  488. conflicting.exec("CREATE VIEW persistence_state AS SELECT 1 AS singleton, 'foreign' AS store_id")
  489. conflicting.close()
  490. expect(() => openDatabase(path, 'wal')).toThrow()
  491. const unchanged = new DatabaseSync(path)
  492. expect(unchanged.prepare(
  493. "SELECT type FROM sqlite_schema WHERE name = 'persistence_state'",
  494. ).get()).toEqual({ type: 'view' })
  495. expect(unchanged.prepare(
  496. "SELECT type FROM sqlite_schema WHERE name = 'sessions'",
  497. ).get()).toBeUndefined()
  498. expect(unchanged.prepare(
  499. "SELECT type FROM sqlite_schema WHERE name = 'events'",
  500. ).get()).toBeUndefined()
  501. expect(unchanged.prepare('PRAGMA application_id').get())
  502. .toEqual({ application_id: SESSION_PERSISTENCE_SQLITE_APPLICATION_ID })
  503. expect(unchanged.prepare('PRAGMA user_version').get()).toEqual({ user_version: SCHEMA_VERSION })
  504. expect(unchanged.prepare('PRAGMA journal_mode').get()).toEqual({ journal_mode: 'delete' })
  505. unchanged.close()
  506. })
  507. it('stamps the persistence application identity with the schema version', async () => {
  508. const path = await freshDbPath()
  509. openDatabase(path, 'wal').close()
  510. const db = new DatabaseSync(path)
  511. expect(db.prepare('PRAGMA application_id').get())
  512. .toEqual({ application_id: SESSION_PERSISTENCE_SQLITE_APPLICATION_ID })
  513. expect(db.prepare('PRAGMA user_version').get()).toEqual({ user_version: SCHEMA_VERSION })
  514. expect(db.prepare('PRAGMA table_info(sessions)').all()).toContainEqual(expect.objectContaining({
  515. name: 'time_zone',
  516. type: 'TEXT',
  517. notnull: 0,
  518. }))
  519. db.close()
  520. })
  521. it('rejects a sibling v3 database (the merge-collided version) rather than opening it against missing columns', async () => {
  522. // Version 3 identified two incompatible sibling layouts, so it is always rejected.
  523. const path = await freshDbPath()
  524. openDatabase(path, 'wal').close() // creates + stamps user_version = SCHEMA_VERSION
  525. const db = openDatabase(path, 'wal')
  526. db.exec('PRAGMA user_version = 3')
  527. db.close()
  528. expect(() => openDatabase(path, 'wal')).toThrow(/schema version 3, incompatible with this build/)
  529. })
  530. it('a corrupt-JSON row in the uncommitted tail is discarded on load, not unloadable', async () => {
  531. const path = await freshDbPath()
  532. const m = meta('corrupt-tail')
  533. const b1 = await backend(path)
  534. await b1.ctx.sessionPersistence.create(m)
  535. await b1.ctx.sessionPersistence.append(m.id, oneTurnLog()) // committed: seqs 0..5
  536. await b1.dispose()
  537. // A torn row after the last committed turn has invalid JSON. `scanRows` locates the boundary
  538. // from seq/type columns without parsing the tail, preserves the committed prefix, and load
  539. // deletes the row; invalid JSON inside the committed region would remain fatal.
  540. const db = openDatabase(path, 'wal')
  541. db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, 6, ?, 7, ?)')
  542. .run(m.id, 'turn/start', '{not valid json')
  543. db.close()
  544. const b2 = await backend(path)
  545. const loaded = await b2.ctx.sessionPersistence.load(m.id)
  546. expect(loaded.events).toEqual(oneTurnLog()) // torn tail discarded, committed intact (turn 1 already balanced → no closers)
  547. // load physically deleted the corrupt tail row, so a fresh append continues.
  548. await b2.ctx.sessionPersistence.append(m.id, [
  549. { type: 'turn/start', seq: 6, time: 8, data: { turn: 2 } },
  550. { type: 'turn/end', seq: 7, time: 9, data: { turn: 2, reason: { kind: 'completed' } } },
  551. ])
  552. const reloaded = await b2.ctx.sessionPersistence.load(m.id)
  553. expect(reloaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  554. await b2.dispose()
  555. })
  556. it('append rolls back the whole batch on a mid-batch seq collision (transaction)', async () => {
  557. const ctx = new Context()
  558. await ctx.plugin(SessionStore)
  559. const fiber = await ctx.plugin(SessionPersistenceSqlite, { path: ':memory:' })
  560. const m = meta('rollback')
  561. await ctx.sessionPersistence.create(m)
  562. await ctx.sessionPersistence.append(m.id, oneTurnLog()) // seqs 0..5
  563. // A batch that re-states an already-stored seq must be rejected and leave
  564. // the stored log unchanged (the UNIQUE (session_id, seq) constraint fires
  565. // inside the transaction → ROLLBACK).
  566. await expect(ctx.sessionPersistence.append(m.id, oneTurnLog())).rejects.toThrow()
  567. const loaded = await ctx.sessionPersistence.load(m.id)
  568. expect(loaded.events).toEqual(oneTurnLog()) // unchanged
  569. await fiber.dispose()
  570. })
  571. it('persists across separate backend instances over the same file', async () => {
  572. const path = await freshDbPath()
  573. const m = meta('persist', '/proj')
  574. const ctx1 = new Context()
  575. await ctx1.plugin(SessionStore)
  576. const fiber1 = await ctx1.plugin(SessionPersistenceSqlite, { path })
  577. await ctx1.sessionPersistence.create(m)
  578. await ctx1.sessionPersistence.append(m.id, oneTurnLog())
  579. await fiber1.dispose()
  580. const ctx2 = new Context()
  581. await ctx2.plugin(SessionStore)
  582. const fiber2 = await ctx2.plugin(SessionPersistenceSqlite, { path })
  583. expect((await ctx2.sessionPersistence.list()).map(x => x.id)).toContain(m.id)
  584. const loaded = await ctx2.sessionPersistence.load(m.id)
  585. expect(loaded.meta).toMatchObject({ id: m.id, cwd: '/proj' })
  586. expect(loaded.events).toEqual(oneTurnLog())
  587. await fiber2.dispose()
  588. })
  589. it('source-qualifies revisions across stores while preserving same-file reopen identity', async () => {
  590. const pathA = await freshDbPath()
  591. const pathB = await freshDbPath()
  592. const m = meta('revision-source')
  593. const a = await backend(pathA)
  594. await a.ctx.sessionPersistence.create(m)
  595. await a.ctx.sessionPersistence.append(m.id, oneTurnLog())
  596. const revisionA = (await a.ctx.sessionPersistence.listSnapshots())[0]?.revision
  597. await a.dispose()
  598. const probeA = openDatabase(pathA, 'wal')
  599. const storeIdA = (probeA.prepare(
  600. 'SELECT store_id FROM persistence_state WHERE singleton = 1',
  601. ).get() as { store_id: string }).store_id
  602. probeA.close()
  603. const aliasA = `${pathA}.alias`
  604. await symlink(pathA, aliasA)
  605. const reopenedA = await backend(aliasA)
  606. expect((await reopenedA.ctx.sessionPersistence.listSnapshots())[0]?.revision).toBe(revisionA)
  607. await reopenedA.dispose()
  608. const b = await backend(pathB)
  609. await b.ctx.sessionPersistence.create(m)
  610. await b.ctx.sessionPersistence.append(m.id, oneTurnLog())
  611. const revisionB = (await b.ctx.sessionPersistence.listSnapshots())[0]?.revision
  612. const probeB = openDatabase(pathB, 'wal')
  613. const storeIdB = (probeB.prepare(
  614. 'SELECT store_id FROM persistence_state WHERE singleton = 1',
  615. ).get() as { store_id: string }).store_id
  616. probeB.close()
  617. expect(storeIdB).not.toBe(storeIdA)
  618. expect(revisionB).not.toBe(revisionA)
  619. expect(String(revisionA)).toMatch(/:revision:1$/)
  620. expect(String(revisionB)).toMatch(/:revision:1$/)
  621. await b.dispose()
  622. })
  623. it('binds a full stored prefix to the same revision as a lightweight read', async () => {
  624. const b = await backend()
  625. const m = meta('stored-prefix-revision')
  626. await b.ctx.sessionPersistence.create(m)
  627. await b.ctx.sessionPersistence.append(m.id, oneTurnLog())
  628. const persistence = b.ctx.sessionPersistence as SessionPersistenceSqlite
  629. const stored = await persistence.loadStored(m.id)
  630. expect(stored?.revision).toBe(await persistence.readStoredRevision(m.id))
  631. expect(await persistence.readStoredRevision(SessionId('missing-revision'))).toBeUndefined()
  632. await b.dispose()
  633. })
  634. it('changes revisions when a deleted session id is materialized again in the same database', async () => {
  635. const path = await freshDbPath()
  636. const m = meta('recreated-revision')
  637. const first = await backend(path)
  638. await first.ctx.sessionPersistence.create(m)
  639. await first.ctx.sessionPersistence.append(m.id, oneTurnLog())
  640. const before = (await first.ctx.sessionPersistence.listSnapshots())[0]?.revision
  641. await first.dispose()
  642. const cleanup = openDatabase(path, 'wal')
  643. cleanup.prepare('DELETE FROM sessions WHERE id = ?').run(m.id)
  644. cleanup.close()
  645. const second = await backend(path)
  646. await second.ctx.sessionPersistence.create(m)
  647. await second.ctx.sessionPersistence.append(m.id, oneTurnLog())
  648. const after = (await second.ctx.sessionPersistence.listSnapshots())[0]?.revision
  649. expect(after).not.toBe(before)
  650. expect(String(before)).toMatch(/:revision:1$/)
  651. expect(String(after)).toMatch(/:revision:1$/)
  652. await second.dispose()
  653. })
  654. it('awaits in-flight readiness before surfacing snapshot-list cancellation', async () => {
  655. const b = await backend()
  656. const internals = b.ctx.sessionPersistence as unknown as { ready: Promise<void> }
  657. const originalReady = internals.ready
  658. const readiness = Promise.withResolvers<undefined>()
  659. internals.ready = readiness.promise
  660. const reason = new Error('SQLite snapshot readiness cancelled')
  661. const controller = new AbortController()
  662. const pending = b.ctx.sessionPersistence.listSnapshots(controller.signal)
  663. let settled = false
  664. void pending.then(
  665. () => { settled = true },
  666. () => { settled = true },
  667. )
  668. controller.abort(reason)
  669. await Promise.resolve()
  670. expect(settled).toBe(false)
  671. readiness.resolve(undefined)
  672. await expect(pending).rejects.toBe(reason)
  673. internals.ready = originalReady
  674. await b.dispose()
  675. })
  676. it('exposes the schema version constant', () => {
  677. expect(SCHEMA_VERSION).toBe(14)
  678. })
  679. it('keeps the revision stable for an empty repair hook', async () => {
  680. const b = await backend()
  681. const m = meta('empty-repair')
  682. await b.ctx.sessionPersistence.create(m)
  683. await b.ctx.sessionPersistence.append(m.id, oneTurnLog())
  684. const before = await b.ctx.sessionPersistence.listSnapshots()
  685. await (b.ctx.sessionPersistence as SessionPersistenceSqlite).commitRepair(m, undefined, [])
  686. expect(await b.ctx.sessionPersistence.listSnapshots()).toEqual(before)
  687. await b.dispose()
  688. })
  689. })
  690. describe('SessionPersistenceSqlite: edge cases', () => {
  691. it('resolves the preparation-cache default without schema normalization', async () => {
  692. const ctx = new Context()
  693. await ctx.plugin(SessionStore)
  694. let persistence!: SessionPersistenceSqlite
  695. await ctx.plugin(Object.assign((inner: Context) => {
  696. persistence = new SessionPersistenceSqlite(inner, {
  697. path: ':memory:',
  698. journalMode: 'wal',
  699. })
  700. }, { inject: ['sessions'] }))
  701. expect(await persistence.list()).toEqual([])
  702. await ctx.fiber.dispose()
  703. })
  704. it('uses the configured preparation cache through the public service', async () => {
  705. const ctx = new Context()
  706. await ctx.plugin(SessionStore)
  707. const fiber = await ctx.plugin(SessionPersistenceSqlite, {
  708. path: ':memory:',
  709. preparedSessionCacheSize: 1,
  710. writeBatchMaxDelayMs: 1,
  711. })
  712. const m = meta('sqlite-preparation-cache')
  713. await ctx.sessionPersistence.create(m)
  714. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  715. const preparation = await ctx.sessionPersistence.prepare(m.id)
  716. expect(preparation.session.header).toEqual(m)
  717. preparation[Symbol.dispose]()
  718. await fiber.dispose()
  719. })
  720. it('rejects and closes a current-schema database with an invalid store identity', async () => {
  721. const path = await freshDbPath()
  722. const db = openDatabase(path, 'wal')
  723. db.exec("UPDATE persistence_state SET store_id = '' WHERE singleton = 1")
  724. db.close()
  725. const b = await backend(path)
  726. await expect(b.ctx.sessionPersistence.listSnapshots()).rejects.toThrow(/no valid store identity/)
  727. await expect(b.dispose()).resolves.toBeUndefined()
  728. })
  729. it('creates a new database and WAL sidecars with owner-only modes without changing its parent mode', async () => {
  730. if (process.platform === 'win32') return
  731. const path = await freshDbPath()
  732. const dir = dirname(path)
  733. await chmod(dir, 0o755)
  734. const b = await backend(path)
  735. await b.ctx.sessionPersistence.list()
  736. expect((await stat(dir)).mode & 0o777).toBe(0o755)
  737. expect((await stat(path)).mode & 0o777).toBe(0o600)
  738. expect((await stat(`${path}-wal`)).mode & 0o777).toBe(0o600)
  739. expect((await stat(`${path}-shm`)).mode & 0o777).toBe(0o600)
  740. await b.dispose()
  741. })
  742. it('creates a persistent rollback journal with owner-only mode', async () => {
  743. if (process.platform === 'win32') return
  744. const path = await freshDbPath()
  745. const ctx = new Context()
  746. await ctx.plugin(SessionStore)
  747. const fiber = await ctx.plugin(SessionPersistenceSqlite, { path, journalMode: 'persist' })
  748. const m = meta('persist-permissions')
  749. await ctx.sessionPersistence.create(m)
  750. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  751. expect((await stat(path)).mode & 0o777).toBe(0o600)
  752. expect((await stat(`${path}-journal`)).mode & 0o777).toBe(0o600)
  753. await fiber.dispose()
  754. })
  755. it('preserves the mode of an existing database file', async () => {
  756. if (process.platform === 'win32') return
  757. const path = await freshDbPath()
  758. await writeFile(path, '', { mode: 0o644 })
  759. await chmod(path, 0o644)
  760. const ctx = new Context()
  761. await ctx.plugin(SessionStore)
  762. const fiber = await ctx.plugin(SessionPersistenceSqlite, { path, journalMode: 'delete' })
  763. await ctx.sessionPersistence.list()
  764. expect((await stat(path)).mode & 0o777).toBe(0o644)
  765. await fiber.dispose()
  766. })
  767. it('surfaces an invalid database path during pre-creation', async () => {
  768. const path = await freshDbPath()
  769. const b = await backend(`${path}\0`)
  770. await expect(b.ctx.sessionPersistence.list()).rejects.toMatchObject({ code: 'ERR_INVALID_ARG_VALUE' })
  771. await b.dispose()
  772. })
  773. it('append rolls back and rethrows when an event INSERT fails inside the transaction', async () => {
  774. const path = await freshDbPath()
  775. const m = meta('rollback-insert')
  776. const b1 = await backend(path)
  777. await b1.ctx.sessionPersistence.create(m)
  778. await b1.ctx.sessionPersistence.append(m.id, oneTurnLog())
  779. // A SECOND backend over the same file loads the session first, so it adopts
  780. // cursor 6 (the committed length) into its OWN in-memory state.
  781. const b2 = await backend(path)
  782. await b2.ctx.sessionPersistence.load(m.id) // cursor 6 in b2
  783. const turn2: SessionEvent[] = [
  784. { type: 'turn/start', seq: 6, time: 7, data: { turn: 2 } },
  785. { type: 'turn/end', seq: 7, time: 8, data: { turn: 2, reason: { kind: 'completed' } } },
  786. ]
  787. // b1 commits seq 6..7 first.
  788. await b1.ctx.sessionPersistence.append(m.id, turn2)
  789. // b2 still thinks its cursor is 6, so this batch passes the contiguity check
  790. // but its INSERT of seq 6 hits the UNIQUE (session_id, seq) constraint
  791. // mid-transaction → ROLLBACK + rethrow.
  792. await expect(b2.ctx.sessionPersistence.append(m.id, turn2)).rejects.toThrow(/UNIQUE/)
  793. // b1's turn is intact; b2's rolled-back attempt left nothing extra.
  794. const loaded = await b1.ctx.sessionPersistence.load(m.id)
  795. expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  796. await b1.dispose()
  797. await b2.dispose()
  798. })
  799. it('journalMode config reaches the database (default wal, rollback modes selectable)', async () => {
  800. // :memory: databases always report journal_mode=memory, so probe file DBs.
  801. const walPath = await freshDbPath()
  802. const bWal = await backend(walPath)
  803. await bWal.ctx.sessionPersistence.create(meta('jm-wal'))
  804. const probe = openDatabase(walPath, 'wal')
  805. expect((probe.prepare('PRAGMA journal_mode').get() as { journal_mode: string }).journal_mode).toBe('wal')
  806. probe.close()
  807. await bWal.dispose()
  808. const deletePath = await freshDbPath()
  809. const ctx = new Context()
  810. await ctx.plugin(SessionStore)
  811. const fiber = await ctx.plugin(SessionPersistenceSqlite, { path: deletePath, journalMode: 'delete' })
  812. await ctx.sessionPersistence.create(meta('jm-delete'))
  813. // Probe through a second connection: journal_mode=delete is a per-database
  814. // property only insofar as no WAL files exist — assert the world, not the
  815. // backend's self-report (no -wal sidecar after writes in delete mode).
  816. const db = openDatabase(deletePath, 'delete')
  817. expect((db.prepare('PRAGMA journal_mode').get() as { journal_mode: string }).journal_mode).toBe('delete')
  818. db.close()
  819. expect(existsSync(`${deletePath}-wal`)).toBe(false)
  820. await fiber.dispose()
  821. })
  822. it('HMR: a DIFFERENT session colliding with a materialized on-disk id is rejected', async () => {
  823. const path = await freshDbPath()
  824. // Instance 1 materializes a session and disposes.
  825. const b1 = await backend(path)
  826. const s1 = b1.ctx.sessions.create(SessionId('hmr-collide'))
  827. appendLog(s1, oneTurnLog())
  828. await b1.ctx.sessions.flush(s1)
  829. await b1.dispose()
  830. // A fresh context with an UNRELATED live session reusing the id meets a
  831. // materialized row that is NOT a prefix of its events → reject.
  832. const ctx = new Context()
  833. await ctx.plugin(SessionStore)
  834. let session!: Session
  835. await ctx.plugin(Object.assign((inner: Context) => {
  836. session = inner.sessions.create(SessionId('hmr-collide'))
  837. }, { inject: ['sessions'] }))
  838. session.append('turn/start', { turn: 1 })
  839. await ctx.plugin(SessionPersistenceSqlite, { path })
  840. await expectFlushError(ctx.sessions.flush(session), /id collision/)
  841. await ctx.fiber.dispose()
  842. })
  843. })
  844. describe('surface field round-trip', () => {
  845. it('rowToEvent parses surface fields from EventRow columns', () => {
  846. const row: EventRow = {
  847. seq: 0, type: 'assistant/message', time: 1,
  848. data: JSON.stringify({ turn: 1, step: 1, content: [] }),
  849. source_event_seqs: JSON.stringify([3, 5]),
  850. surface_op: JSON.stringify('append'),
  851. }
  852. const event = rowToEvent(row)
  853. expect((event as SurfaceEvent).sourceEventSeqs).toEqual([3, 5])
  854. expect((event as SurfaceEvent).surfaceOp).toBe('append')
  855. })
  856. it('rowToEvent handles replace surfaceOp object', () => {
  857. const row: EventRow = {
  858. seq: 0, type: 'assistant/message', time: 1,
  859. data: JSON.stringify({ turn: 1, step: 1, content: [] }),
  860. source_event_seqs: JSON.stringify([0, 1]),
  861. surface_op: JSON.stringify({ op: 'replace', start: 0, end: 1 }),
  862. }
  863. const event = rowToEvent(row)
  864. expect((event as SurfaceEvent).sourceEventSeqs).toEqual([0, 1])
  865. expect((event as SurfaceEvent).surfaceOp).toEqual({ op: 'replace', start: 0, end: 1 })
  866. })
  867. it('scanRows with surface columns reconstructs events with surface fields', () => {
  868. const rows: EventRow[] = [
  869. { seq: 0, type: 'user/message', time: 1,
  870. data: JSON.stringify({ content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' } }),
  871. source_event_seqs: null, surface_op: '{"op":"replace","start":0,"end":0}' },
  872. { seq: 1, type: 'turn/end', time: 2,
  873. data: JSON.stringify({ turn: 1, reason: { kind: 'completed' } }),
  874. source_event_seqs: null, surface_op: null },
  875. ]
  876. const { preserved } = scanRows(rows)
  877. expect(preserved).toHaveLength(2)
  878. expect((preserved[0]! as SurfaceEvent).surfaceOp).toEqual({ op: 'replace', start: 0, end: 0 })
  879. expect((preserved[0]! as SurfaceEvent).sourceEventSeqs).toBeUndefined()
  880. expect((preserved[1] as SessionEvent<SurfaceEventType>).surfaceOp).toBeUndefined()
  881. })
  882. it('append and load round-trips surface fields through SQLite', async () => {
  883. const ctx = new Context()
  884. await ctx.plugin(SessionStore)
  885. const fiber = await ctx.plugin(SessionPersistenceSqlite, { path: ':memory:' })
  886. const session = ctx.sessions.create(SessionId('roundtrip-surface'))
  887. session.append('turn/start', { turn: 1 })
  888. session.append('step/start', { turn: 1, step: 1 })
  889. session.append('user/message', createUserMessage({
  890. content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' },
  891. }), { surfaceOp: 'append' })
  892. session.append('assistant/message', {
  893. turn: 1, step: 1,
  894. message: createMessage({
  895. role: 'assistant',
  896. content: [],
  897. source: {
  898. kind: 'model',
  899. ...{ provider: 'mock', model: 'mock' },
  900. },
  901. }),
  902. }, { surfaceOp: 'append', sourceEventSeqs: [2] })
  903. session.append('step/end', { turn: 1, step: 1 })
  904. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  905. await ctx.sessions.flush(session)
  906. const loaded = await ctx.sessionPersistence.load(SessionId('roundtrip-surface'))
  907. expect(loaded.events).toHaveLength(6)
  908. const um = loaded.events[2]!
  909. expect((um as SurfaceEvent).surfaceOp).toBe('append')
  910. expect((um as SurfaceEvent).sourceEventSeqs).toBeUndefined()
  911. const am = loaded.events[3]!
  912. expect((am as SurfaceEvent).surfaceOp).toBe('append')
  913. expect((am as SurfaceEvent).sourceEventSeqs).toEqual([2])
  914. await fiber.dispose()
  915. })
  916. it('persists events with surfaceOp but no sourceEventSeqs (covers null branch in surfaceBindings)', async () => {
  917. const ctx = new Context()
  918. await ctx.plugin(SessionStore)
  919. const fiber = await ctx.plugin(SessionPersistenceSqlite, { path: ':memory:' })
  920. const session = ctx.sessions.create(SessionId('surface-noseq'))
  921. session.append('turn/start', { turn: 1 })
  922. session.append('user/message', createUserMessage({
  923. content: [],
  924. source: { kind: 'user' },
  925. }), { surfaceOp: 'append' })
  926. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  927. await ctx.sessions.flush(session)
  928. const loaded = await ctx.sessionPersistence.load(SessionId('surface-noseq'))
  929. expect((loaded.events[1]! as SurfaceEvent).surfaceOp).toBe('append')
  930. expect((loaded.events[1]! as SurfaceEvent).sourceEventSeqs).toBeUndefined()
  931. await fiber.dispose()
  932. })
  933. })