sqlite.spec.ts 36 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888
  1. import { afterEach, describe, expect, it, vi } from 'vitest'
  2. import { spawn } from 'node:child_process'
  3. import { Context } from '@deepseek-ai/cordis'
  4. import { once } from 'node:events'
  5. import { chmod, mkdir, mkdtemp, rm, stat, symlink, writeFile } from 'node:fs/promises'
  6. import { tmpdir } from 'node:os'
  7. import { join } from 'node:path'
  8. import { performance } from 'node:perf_hooks'
  9. import { pathToFileURL } from 'node:url'
  10. import { DatabaseSync } from 'node:sqlite'
  11. import Loader from '@deepseek-ai/cordis-plugin-loader'
  12. import Include from '@deepseek-ai/cordis-plugin-include'
  13. import SessionStore, { SessionId, type SessionEvent } from '@deepseek-ai/dsh-session'
  14. import SessionPersistenceSqlite, {
  15. DEFAULT_BUSY_TIMEOUT_MS,
  16. SCHEMA_VERSION,
  17. } from '@deepseek-ai/dsh-session-persistence-sqlite'
  18. import {
  19. runCoordinatorContract,
  20. type CoordinatorFixture,
  21. } from '../../session-persistence/tests/coordinator-contract.ts'
  22. import {
  23. meta,
  24. runPersistenceContract,
  25. } from '../../session-persistence/tests/contract.ts'
  26. import { MAX_PACKED_DATA_BYTES } from '../src/codec.ts'
  27. import {
  28. decodeEventRow,
  29. decodeSessionRow,
  30. decodeStoreIdentity,
  31. openDatabase,
  32. validateSchemaForMutation,
  33. rowToMeta,
  34. SESSION_PERSISTENCE_SQLITE_APPLICATION_ID,
  35. type SessionRow,
  36. } from '../src/schema.ts'
  37. import { SqliteStore } from '../src/store.ts'
  38. import { sql } from '../src/sql.ts'
  39. import { testSql } from './test-sql.ts'
  40. const dirs: string[] = []
  41. afterEach(async () => {
  42. for (const directory of dirs.splice(0)) await rm(directory, { recursive: true, force: true })
  43. })
  44. async function freshDbPath(prefix = 'dsh-sqlite-'): Promise<string> {
  45. const directory = await mkdtemp(join(tmpdir(), prefix))
  46. dirs.push(directory)
  47. return join(directory, 'sessions.db')
  48. }
  49. async function backendFailure(path: string): Promise<unknown> {
  50. const ctx = new Context()
  51. await ctx.plugin(SessionStore)
  52. try {
  53. await ctx.plugin(SessionPersistenceSqlite, { path })
  54. await ctx.sessionPersistence.list()
  55. return undefined
  56. } catch (error: unknown) {
  57. return error
  58. } finally {
  59. await ctx.fiber.dispose()
  60. }
  61. }
  62. function errorMessage(error: unknown): string {
  63. return error instanceof Error ? error.message : String(error)
  64. }
  65. function databaseWithJournalFailure(
  66. nextFailure: () => Error | undefined,
  67. ): typeof DatabaseSync {
  68. return class JournalFailureDatabase extends DatabaseSync {
  69. override prepare(source: string) {
  70. if (source !== sql('journal-mode-wal')) return super.prepare(source)
  71. const statement = super.prepare(sql('journal-mode-wal'))
  72. const get = statement.get.bind(statement)
  73. Object.defineProperty(statement, 'get', {
  74. value: () => {
  75. const failure = nextFailure()
  76. if (failure !== undefined) throw failure
  77. return get()
  78. },
  79. })
  80. return statement
  81. }
  82. }
  83. }
  84. function chunk(seq: number, text = `token-${seq}`): SessionEvent {
  85. return {
  86. type: 'assistant/chunk',
  87. seq,
  88. time: 1_000 + seq,
  89. data: {
  90. turn: 1,
  91. step: 1,
  92. chunk: { type: 'text-delta', index: 0, text },
  93. },
  94. }
  95. }
  96. function chunkLog(count: number): SessionEvent[] {
  97. return [
  98. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
  99. { type: 'step/start', seq: 1, time: 2, data: { turn: 1, step: 1 } },
  100. ...Array.from({ length: count }, (_, index) => chunk(index + 2)),
  101. { type: 'step/end', seq: count + 2, time: count + 3, data: { turn: 1, step: 1 } },
  102. {
  103. type: 'turn/end',
  104. seq: count + 3,
  105. time: count + 4,
  106. data: { turn: 1, reason: { kind: 'completed' } },
  107. },
  108. ]
  109. }
  110. async function measureWriteTraffic(
  111. path: string,
  112. events: readonly SessionEvent[],
  113. ): Promise<{
  114. readonly walBytes: number
  115. readonly idleWalBytes: number
  116. readonly rows: number
  117. readonly largest: number
  118. readonly inserted: number
  119. readonly changed: number
  120. readonly removed: number
  121. }> {
  122. interface PhysicalRow {
  123. readonly rowid: number
  124. readonly seq: number
  125. readonly type: string
  126. readonly time: number
  127. readonly data: string | Uint8Array
  128. readonly source_event_seqs: Uint8Array | null
  129. readonly surface_op: string | null
  130. readonly is_packed: number
  131. }
  132. const sameValue = (left: string | Uint8Array | null, right: string | Uint8Array | null): boolean => (
  133. typeof left === 'string' || left === null
  134. ? left === right
  135. : right instanceof Uint8Array && Buffer.from(left).equals(Buffer.from(right))
  136. )
  137. const sameRow = (left: PhysicalRow, right: PhysicalRow): boolean => (
  138. left.rowid === right.rowid
  139. && left.seq === right.seq
  140. && left.type === right.type
  141. && left.time === right.time
  142. && sameValue(left.data, right.data)
  143. && sameValue(left.source_event_seqs, right.source_event_seqs)
  144. && left.surface_op === right.surface_op
  145. && left.is_packed === right.is_packed
  146. )
  147. const ctx = new Context()
  148. await ctx.plugin(SessionStore)
  149. await ctx.plugin(SessionPersistenceSqlite, { path, writeBatchMaxDelayMs: 200 })
  150. try {
  151. const header = meta('traffic')
  152. await ctx.sessionPersistence.create(header)
  153. let previous = new Map<number, PhysicalRow>()
  154. let inserted = 0
  155. let changed = 0
  156. let removed = 0
  157. const probe = new DatabaseSync(path, { readOnly: true })
  158. try {
  159. const selectRows = probe.prepare(testSql('select-event-rows'))
  160. for (let offset = 0; offset < events.length; offset += 40) {
  161. await ctx.sessionPersistence.append(header.id, events.slice(offset, offset + 40))
  162. const current = new Map((selectRows.all(header.id) as unknown as PhysicalRow[])
  163. .map(row => [row.seq, row]))
  164. for (const [seq, row] of current) {
  165. const old = previous.get(seq)
  166. if (old === undefined) inserted += 1
  167. else if (!sameRow(old, row)) changed += 1
  168. }
  169. for (const seq of previous.keys()) if (!current.has(seq)) removed += 1
  170. previous = current
  171. }
  172. } finally {
  173. probe.close()
  174. }
  175. const db = new DatabaseSync(path, { readOnly: true })
  176. const measured = db.prepare(testSql('measure-write-traffic')).get() as { rows: number; largest: number }
  177. db.close()
  178. const walBytes = (await stat(`${path}-wal`)).size
  179. await new Promise(resolve => setTimeout(resolve, 250))
  180. return {
  181. walBytes,
  182. idleWalBytes: (await stat(`${path}-wal`)).size,
  183. rows: measured.rows,
  184. largest: measured.largest,
  185. inserted,
  186. changed,
  187. removed,
  188. }
  189. } finally {
  190. await ctx.fiber.dispose()
  191. }
  192. }
  193. runPersistenceContract('sqlite', async () => {
  194. const ctx = new Context()
  195. await ctx.plugin(SessionStore)
  196. const fiber = await ctx.plugin(SessionPersistenceSqlite, { path: ':memory:' })
  197. return {
  198. persistence: ctx.sessionPersistence,
  199. dispose: async () => { await fiber.dispose() },
  200. }
  201. })
  202. runCoordinatorContract('sqlite', async (): Promise<CoordinatorFixture> => {
  203. const directory = await mkdtemp(join(tmpdir(), 'dsh-sqlite-coord-'))
  204. const path = join(directory, 'sessions.db')
  205. return {
  206. mount: async ctx => ctx.plugin(SessionPersistenceSqlite, { path }),
  207. corruptTail: async (id) => {
  208. const db = new DatabaseSync(path)
  209. const last = db.prepare(testSql('select-last-event'))
  210. .get(id) as { seq: number; type: string; data: string }
  211. const logicalLength = last.type === 'text-chunks'
  212. ? (JSON.parse(last.data) as { texts: string[] }).texts.length
  213. : 1
  214. const next = last.seq + logicalLength
  215. db.prepare(testSql('insert-corrupt-event'))
  216. .run(id, next, 'assistant/chunk', 99, '{not valid json', 0)
  217. db.close()
  218. },
  219. cleanup: async () => { await rm(directory, { recursive: true, force: true }) },
  220. }
  221. })
  222. describe('SessionPersistenceSqlite physical packing', () => {
  223. it('loads from cordis.yml and packs through the assembled service', async () => {
  224. const path = await freshDbPath('dsh-sqlite-loader-')
  225. const configPath = join(path, '..', 'cordis.yml')
  226. await writeFile(configPath, [
  227. "- name: '@deepseek-ai/dsh-session'",
  228. "- name: '@deepseek-ai/dsh-session-persistence-sqlite'",
  229. ' config:',
  230. ` path: ${JSON.stringify(path)}`,
  231. '',
  232. ].join('\n'))
  233. const ctx = new Context()
  234. ctx.baseUrl = pathToFileURL(join(path, '..')).href + '/'
  235. await ctx.plugin(Loader)
  236. ctx.loader.builtins.include = Include
  237. ctx.loader.internal = {
  238. version: 'sqlite',
  239. async import(specifier: string) {
  240. if (specifier === '@deepseek-ai/dsh-session') return SessionStore
  241. if (specifier === '@deepseek-ai/dsh-session-persistence-sqlite') {
  242. return SessionPersistenceSqlite
  243. }
  244. throw new Error(`unexpected Loader import: ${specifier}`)
  245. },
  246. } as unknown as NonNullable<typeof ctx.loader.internal>
  247. await ctx.loader.create({
  248. name: 'cordis:include',
  249. config: { path: pathToFileURL(configPath).href },
  250. })
  251. await ctx.loader.await()
  252. const header = meta('loader')
  253. const events = chunkLog(4)
  254. await ctx.sessionPersistence.create(header)
  255. await ctx.sessionPersistence.append(header.id, events)
  256. expect((await ctx.sessionPersistence.inspect(header.id)).events).toEqual(events)
  257. await ctx.fiber.dispose()
  258. const db = new DatabaseSync(path)
  259. expect(db.prepare(testSql('count-packed-events')).get())
  260. .toEqual({ count: 1 })
  261. db.close()
  262. })
  263. it('packs each append once without rewriting earlier rows and seeks inside packed rows', async () => {
  264. const path = await freshDbPath()
  265. const ctx = new Context()
  266. await ctx.plugin(SessionStore)
  267. const fiber = await ctx.plugin(SessionPersistenceSqlite, { path })
  268. const header = meta('packed')
  269. const events = chunkLog(100)
  270. await ctx.sessionPersistence.create(header)
  271. await ctx.sessionPersistence.append(header.id, events.slice(0, 3))
  272. await ctx.sessionPersistence.append(header.id, events.slice(3, 4))
  273. const before = new DatabaseSync(path, { readOnly: true })
  274. const originalRows = before.prepare(testSql('select-event-rowids')).all()
  275. before.close()
  276. await ctx.sessionPersistence.append(header.id, events.slice(4))
  277. const inspected = await ctx.sessionPersistence.inspect(header.id)
  278. expect(inspected.events).toEqual(events)
  279. for (const fromSeq of [0, 2, 25, 101, 104, 105]) {
  280. expect((await ctx.sessionPersistence.readFrom(header.id, fromSeq)).events)
  281. .toEqual(events.filter(event => event.seq >= fromSeq))
  282. }
  283. await fiber.dispose()
  284. const db = new DatabaseSync(path)
  285. expect(db.prepare(testSql('select-user-version')).get()).toEqual({ user_version: SCHEMA_VERSION })
  286. expect(db.prepare(testSql('select-page-size')).get()).toEqual({ page_size: 65_536 })
  287. expect(db.prepare(testSql('count-events')).get()).toEqual({ count: 7 })
  288. expect(db.prepare(testSql('count-packed-events')).get())
  289. .toEqual({ count: 1 })
  290. expect(db.prepare(testSql('select-event-rowids')).all().slice(0, originalRows.length))
  291. .toEqual(originalRows)
  292. db.close()
  293. })
  294. it.runIf(process.platform !== 'win32')('bounds paced-stream WAL extent without rewriting committed rows', async () => {
  295. const events = chunkLog(1_000)
  296. const measured = await measureWriteTraffic(await freshDbPath('dsh-sqlite-traffic-'), events)
  297. expect(measured).toMatchObject({ rows: 31, inserted: 31, changed: 0, removed: 0 })
  298. expect(measured.inserted).toBe(measured.rows)
  299. expect(measured.largest).toBeLessThanOrEqual(MAX_PACKED_DATA_BYTES)
  300. expect(measured.idleWalBytes).toBe(measured.walBytes)
  301. })
  302. it('includes a packed predecessor when an overlapping scalar tail hides it', async () => {
  303. const path = await freshDbPath('dsh-sqlite-overlap-')
  304. const store = new SqliteStore({ path, journalMode: 'wal', busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS })
  305. const header = meta('overlap')
  306. await store.appendBatch(header, [chunk(0), chunk(1), chunk(2)], false)
  307. const db = new DatabaseSync(path)
  308. db.prepare(testSql('insert-corrupt-event'))
  309. .run(header.id, 1, 'assistant/chunk', 2, JSON.stringify(chunk(1).data), 0)
  310. db.close()
  311. expect((await store.loadStoredFrom(header.id, 2))?.events).toEqual([chunk(2)])
  312. const malformed = new DatabaseSync(path)
  313. malformed.prepare(testSql('delete-session-events')).run(header.id)
  314. malformed.prepare(testSql('insert-corrupt-event'))
  315. .run(header.id, 0, 'text-chunks', 1, '{not json', 1)
  316. malformed.close()
  317. expect((await store.loadStoredFrom(header.id, 2))?.events).toEqual([])
  318. await store.close()
  319. })
  320. it('waits for a competing process within the configured busy timeout', async () => {
  321. const path = await freshDbPath('dsh-sqlite-busy-')
  322. const store = new SqliteStore({ path, journalMode: 'wal', busyTimeoutMs: 1_000 })
  323. const header = meta('busy')
  324. await store.appendBatch(header, [chunk(0)], false)
  325. const holder = spawn(process.execPath, ['--input-type=module', '-e', String.raw`
  326. import { DatabaseSync } from 'node:sqlite';
  327. const db = new DatabaseSync(process.argv[1]);
  328. db.exec('BEGIN IMMEDIATE');
  329. process.stdout.write('locked\n');
  330. setTimeout(() => { db.exec('COMMIT'); db.close(); }, 100);
  331. `, path], { stdio: ['ignore', 'pipe', 'pipe'] })
  332. const exited = new Promise<number | null>((resolve, reject) => {
  333. holder.once('error', reject)
  334. holder.once('exit', resolve)
  335. })
  336. try {
  337. await once(holder.stdout, 'data')
  338. await expect(store.appendBatch(header, [chunk(1)], true)).resolves.toBeUndefined()
  339. const code = await exited
  340. expect(code).toBe(0)
  341. expect((await store.loadStored(header.id))?.events).toEqual([chunk(0), chunk(1)])
  342. } finally {
  343. if (holder.exitCode === null) holder.kill()
  344. await store.close()
  345. }
  346. })
  347. it('rejects an older SQLite physical schema', async () => {
  348. const path = await freshDbPath('dsh-sqlite-old-schema-')
  349. const seed = await openDatabase(DatabaseSync, path, 'wal', DEFAULT_BUSY_TIMEOUT_MS)
  350. seed.exec(testSql('set-user-version-17'))
  351. seed.close()
  352. await chmod(path, 0o600)
  353. await expect(openDatabase(DatabaseSync, path, 'wal', DEFAULT_BUSY_TIMEOUT_MS))
  354. .rejects.toThrow(/schema version 17.*incompatible/)
  355. })
  356. it('keeps the page size of an established schema 19 database', async () => {
  357. const path = await freshDbPath('dsh-sqlite-page-size-')
  358. const seed = await openDatabase(DatabaseSync, path, 'delete', DEFAULT_BUSY_TIMEOUT_MS)
  359. seed.close()
  360. const resize = new DatabaseSync(path)
  361. resize.exec(testSql('set-page-size-4096'))
  362. resize.exec(testSql('vacuum'))
  363. expect(resize.prepare(testSql('select-page-size')).get()).toEqual({ page_size: 4_096 })
  364. resize.close()
  365. const reopened = await openDatabase(DatabaseSync, path, 'wal', DEFAULT_BUSY_TIMEOUT_MS)
  366. expect(reopened.prepare(testSql('select-page-size')).get()).toEqual({ page_size: 4_096 })
  367. reopened.close()
  368. })
  369. it('rejects a stale physical append without replacing the winning tail', async () => {
  370. const path = await freshDbPath('dsh-sqlite-stale-')
  371. const first = new SqliteStore({ path, journalMode: 'wal', busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS })
  372. const second = new SqliteStore({ path, journalMode: 'wal', busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS })
  373. const header = meta(SessionId('stale'))
  374. await first.appendBatch(header, [chunk(0)], false)
  375. await second.appendBatch(header, [chunk(1)], true)
  376. await expect(first.appendBatch(header, [chunk(1)], true)).rejects.toThrow(/stored next seq is 2/)
  377. expect((await first.loadStored(header.id))?.events).toEqual([chunk(0), chunk(1)])
  378. await first.close()
  379. await second.close()
  380. })
  381. it('rolls back lazy integer-key materialization after a rejected append', async () => {
  382. const store = new SqliteStore({
  383. path: await freshDbPath('dsh-sqlite-key-rollback-'),
  384. journalMode: 'wal',
  385. busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS,
  386. })
  387. const header = meta(SessionId('key-rollback'))
  388. await expect(store.appendBatch(header, [chunk(1)], false)).rejects.toThrow(/stored next seq is 0/)
  389. await expect(store.appendBatch(header, [chunk(0)], true)).rejects.toThrow(/metadata row is missing/)
  390. await expect(store.appendBatch(header, [chunk(0)], false)).resolves.toBeUndefined()
  391. expect((await store.loadStored(header.id))?.events).toEqual([chunk(0)])
  392. await store.close()
  393. })
  394. it('rejects a stale repair without deleting a newer winning tail', async () => {
  395. const path = await freshDbPath('dsh-sqlite-stale-repair-')
  396. const stale = new SqliteStore({ path, journalMode: 'wal', busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS })
  397. const winner = new SqliteStore({ path, journalMode: 'wal', busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS })
  398. const header = meta(SessionId('stale-repair'))
  399. await stale.appendBatch(header, [chunk(0)], false)
  400. const db = new DatabaseSync(path)
  401. db.prepare(testSql('insert-corrupt-event')).run(header.id, 1, 'assistant/chunk', 2, '{not json', 0)
  402. db.close()
  403. expect((await stale.loadStored(header.id))?.tornMarker).toBe(1)
  404. await winner.commitRepair(header, 1, [])
  405. await winner.appendBatch(header, [chunk(1), chunk(2)], true)
  406. await expect(stale.commitRepair(header, 1, [])).rejects.toThrow(/repair is stale/)
  407. expect((await stale.loadStored(header.id))?.events).toEqual([chunk(0), chunk(1), chunk(2)])
  408. await stale.close()
  409. await winner.close()
  410. })
  411. })
  412. describe('SessionPersistenceSqlite schema ownership', () => {
  413. it('accepts every configured journal mode and SQLite memory mode result', async () => {
  414. const resources = {
  415. wal: 'journal-mode-wal',
  416. delete: 'journal-mode-delete',
  417. truncate: 'journal-mode-truncate',
  418. persist: 'journal-mode-persist',
  419. } as const
  420. for (const mode of ['wal', 'delete', 'truncate', 'persist'] as const) {
  421. ;(await openDatabase(DatabaseSync, ':memory:', mode, DEFAULT_BUSY_TIMEOUT_MS)).close()
  422. const path = await freshDbPath(`dsh-sqlite-journal-${mode}-`)
  423. const db = await openDatabase(DatabaseSync, path, mode, DEFAULT_BUSY_TIMEOUT_MS)
  424. expect(db.prepare(sql(resources[mode])).get()).toEqual({ journal_mode: mode })
  425. expect(db.prepare(sql('select-trusted-schema')).get()).toEqual({ trusted_schema: 0 })
  426. expect(db.prepare(sql('select-mmap-size')).get()).toEqual({ mmap_size: 0 })
  427. expect(db.prepare(sql('select-synchronous')).get()).toEqual({ synchronous: 2 })
  428. db.close()
  429. }
  430. })
  431. it('retries a busy journal-mode transition within its retry budget', async () => {
  432. const path = await freshDbPath('dsh-sqlite-journal-busy-')
  433. let attempts = 0
  434. const BusyOnceDatabase = databaseWithJournalFailure(() => {
  435. attempts += 1
  436. return attempts === 1
  437. ? Object.assign(new Error('database is locked'), {
  438. code: 'ERR_SQLITE_ERROR',
  439. errcode: 5,
  440. errstr: 'database is locked',
  441. })
  442. : undefined
  443. })
  444. const db = await openDatabase(BusyOnceDatabase, path, 'wal', 100)
  445. expect(attempts).toBe(2)
  446. expect(db.prepare(sql('journal-mode-wal')).get()).toEqual({ journal_mode: 'wal' })
  447. expect(db.prepare(sql('select-trusted-schema')).get()).toEqual({ trusted_schema: 0 })
  448. expect(db.prepare(sql('select-mmap-size')).get()).toEqual({ mmap_size: 0 })
  449. expect(db.prepare(sql('select-synchronous')).get()).toEqual({ synchronous: 2 })
  450. db.close()
  451. })
  452. it('does not retry journal failures outside the available busy budget', async () => {
  453. for (const { errcode, timeout } of [
  454. { errcode: 5, timeout: 0 },
  455. { errcode: 6, timeout: 100 },
  456. ]) {
  457. let attempts = 0
  458. const FailingDatabase = databaseWithJournalFailure(() => {
  459. attempts += 1
  460. return Object.assign(new Error(`SQLite error ${errcode}`), { errcode })
  461. })
  462. await expect(openDatabase(
  463. FailingDatabase,
  464. await freshDbPath(`dsh-sqlite-journal-failure-${errcode}-`),
  465. 'wal',
  466. timeout,
  467. )).rejects.toThrow(`SQLite error ${errcode}`)
  468. expect(attempts).toBe(1)
  469. }
  470. })
  471. it('starts no journal retry after its open-relative cutoff', async () => {
  472. let attempts = 0
  473. const BusyDatabase = databaseWithJournalFailure(() => {
  474. attempts += 1
  475. return Object.assign(new Error('database is locked'), { errcode: 5 })
  476. })
  477. const clock = vi.spyOn(performance, 'now')
  478. .mockReturnValueOnce(0)
  479. .mockReturnValueOnce(50)
  480. .mockReturnValueOnce(100)
  481. try {
  482. await expect(openDatabase(
  483. BusyDatabase,
  484. await freshDbPath('dsh-sqlite-journal-cutoff-'),
  485. 'wal',
  486. 100,
  487. )).rejects.toThrow('database is locked')
  488. } finally {
  489. clock.mockRestore()
  490. }
  491. expect(attempts).toBe(1)
  492. })
  493. it('paces repeated busy journal-mode attempts', async () => {
  494. const attemptedAt: number[] = []
  495. const BusyTwiceDatabase = databaseWithJournalFailure(() => {
  496. attemptedAt.push(performance.now())
  497. return attemptedAt.length <= 2
  498. ? Object.assign(new Error('database is locked'), { errcode: 5 })
  499. : undefined
  500. })
  501. const db = await openDatabase(
  502. BusyTwiceDatabase,
  503. await freshDbPath('dsh-sqlite-journal-paced-'),
  504. 'wal',
  505. DEFAULT_BUSY_TIMEOUT_MS,
  506. )
  507. db.close()
  508. expect(attemptedAt).toHaveLength(3)
  509. for (let index = 1; index < attemptedAt.length; index += 1) {
  510. const previous = attemptedAt[index - 1]
  511. const current = attemptedAt[index]
  512. if (previous === undefined || current === undefined) throw new Error('missing journal attempt timestamp')
  513. expect(current - previous).toBeGreaterThanOrEqual(5)
  514. }
  515. })
  516. it('rejects unversioned, incompatible, and foreign-application databases', async () => {
  517. const unversionedPath = await freshDbPath('dsh-sqlite-unversioned-')
  518. const unversioned = new DatabaseSync(unversionedPath)
  519. unversioned.exec(testSql('create-unrelated-table'))
  520. unversioned.close()
  521. await expect(openDatabase(DatabaseSync, unversionedPath, 'wal', DEFAULT_BUSY_TIMEOUT_MS)).rejects.toThrow(/unversioned schema/)
  522. const incompatiblePath = await freshDbPath('dsh-sqlite-incompatible-')
  523. const incompatible = new DatabaseSync(incompatiblePath)
  524. incompatible.exec(testSql('set-user-version-17'))
  525. incompatible.close()
  526. await expect(openDatabase(DatabaseSync, incompatiblePath, 'wal', DEFAULT_BUSY_TIMEOUT_MS)).rejects.toThrow(/incompatible with this build/)
  527. const foreignPath = await freshDbPath('dsh-sqlite-foreign-')
  528. const foreign = new DatabaseSync(foreignPath)
  529. foreign.exec(testSql('set-user-version-19'))
  530. foreign.exec(testSql('set-application-id-12345'))
  531. foreign.close()
  532. await expect(openDatabase(DatabaseSync, foreignPath, 'wal', DEFAULT_BUSY_TIMEOUT_MS)).rejects.toThrow(/has application id 12345/)
  533. })
  534. it('rejects changed columns and non-strict owned tables', async () => {
  535. const changedPath = await freshDbPath('dsh-sqlite-columns-')
  536. ;(await openDatabase(DatabaseSync, changedPath, 'wal', DEFAULT_BUSY_TIMEOUT_MS)).close()
  537. const changed = new DatabaseSync(changedPath)
  538. changed.exec(testSql('add-unexpected-column'))
  539. changed.close()
  540. await expect(openDatabase(DatabaseSync, changedPath, 'wal', DEFAULT_BUSY_TIMEOUT_MS)).rejects.toThrow(/required schema objects/)
  541. const nonStrictPath = await freshDbPath('dsh-sqlite-nonstrict-')
  542. ;(await openDatabase(DatabaseSync, nonStrictPath, 'wal', DEFAULT_BUSY_TIMEOUT_MS)).close()
  543. const nonStrict = new DatabaseSync(nonStrictPath)
  544. nonStrict.exec(testSql('replace-events-with-nonstrict-table'))
  545. nonStrict.close()
  546. await expect(openDatabase(DatabaseSync, nonStrictPath, 'wal', DEFAULT_BUSY_TIMEOUT_MS)).rejects.toThrow(/required schema objects/)
  547. const loosePath = await freshDbPath('dsh-sqlite-loose-')
  548. const loose = new DatabaseSync(loosePath)
  549. loose.exec(testSql('create-loose-schema'))
  550. loose.close()
  551. await expect(openDatabase(DatabaseSync, loosePath, 'wal', DEFAULT_BUSY_TIMEOUT_MS)).rejects.toThrow(/required schema objects/)
  552. })
  553. it('rejects schema ownership changes observed at mutation time', async () => {
  554. const changedVersion = await openDatabase(DatabaseSync, ':memory:', 'wal', DEFAULT_BUSY_TIMEOUT_MS)
  555. changedVersion.exec(testSql('set-user-version-17'))
  556. expect(() => { validateSchemaForMutation(DatabaseSync, changedVersion, ':memory:') })
  557. .toThrow(/schema changed before mutation/)
  558. changedVersion.close()
  559. const changedApplication = await openDatabase(DatabaseSync, ':memory:', 'wal', DEFAULT_BUSY_TIMEOUT_MS)
  560. changedApplication.exec(testSql('set-application-id-12345'))
  561. expect(() => { validateSchemaForMutation(DatabaseSync, changedApplication, ':memory:') })
  562. .toThrow(/application id changed before mutation/)
  563. changedApplication.close()
  564. })
  565. it('validates creation time and restores every optional header field', () => {
  566. const base: SessionRow = {
  567. id: 'stored-header',
  568. version: 0,
  569. created_at: 1,
  570. cwd: '/project',
  571. parent_session: 'parent',
  572. seed_length: 4,
  573. origin: 'subagent',
  574. incarnation: '00000000-0000-4000-8000-000000000000',
  575. revision: 1,
  576. delegation_depth: 2,
  577. agent_preset: 'minimal',
  578. }
  579. expect(rowToMeta(decodeSessionRow(base))).toMatchObject({
  580. cwd: '/project',
  581. parentSession: 'parent',
  582. seedLength: 4,
  583. origin: 'subagent',
  584. delegationDepth: 2,
  585. agentPreset: 'minimal',
  586. })
  587. expect(() => decodeSessionRow({ ...base, created_at: -1 })).toThrow(/created_at/)
  588. expect(() => decodeSessionRow({ ...base, origin: 'external' })).toThrow(/origin/)
  589. expect(() => decodeSessionRow({ ...base, delegation_depth: -1 })).toThrow(/delegation_depth/)
  590. })
  591. it('rejects malformed SQLite row primitives generically', () => {
  592. const base: SessionRow = {
  593. id: 'stored-header',
  594. version: 0,
  595. created_at: 1,
  596. cwd: '/project',
  597. parent_session: null,
  598. seed_length: null,
  599. origin: null,
  600. incarnation: '00000000-0000-4000-8000-000000000000',
  601. revision: 1,
  602. delegation_depth: null,
  603. agent_preset: null,
  604. }
  605. for (const [value, message] of [
  606. [null, /object/],
  607. [{ ...base, id: 1 }, /id.*string/],
  608. [{ ...base, id: '' }, /id.*empty/],
  609. [{ ...base, version: '0' }, /version.*safe integer/],
  610. [{ ...base, cwd: 'relative' }, /cwd.*absolute/],
  611. [{ ...base, cwd: 1 }, /cwd.*string or null/],
  612. [{ ...base, incarnation: 'invalid' }, /incarnation.*UUID/],
  613. [{ ...base, seed_length: '1' }, /seed_length.*safe integer or null/],
  614. [{ ...base, agent_preset: 1 }, /agent_preset.*string or null/],
  615. ] as const) {
  616. expect(() => decodeSessionRow(value)).toThrow(message)
  617. }
  618. const eventRow = {
  619. seq: 0, type: 'turn/start', time: 1, data: '{}',
  620. source_event_seqs: null, surface_op: null, is_packed: 0,
  621. }
  622. for (const [value, message] of [
  623. [null, /object/],
  624. [{ ...eventRow, seq: '0' }, /seq.*safe integer/],
  625. [{ ...eventRow, type: '' }, /type.*empty/],
  626. [{ ...eventRow, time: '1' }, /time.*safe integer/],
  627. [{ ...eventRow, data: 1 }, /data.*string or blob/],
  628. [{ ...eventRow, source_event_seqs: 1 }, /source_event_seqs.*blob or null/],
  629. [{ ...eventRow, is_packed: 2 }, /is_packed.*0 or 1/],
  630. ] as const) {
  631. expect(() => decodeEventRow(value)).toThrow(message)
  632. }
  633. expect(() => decodeStoreIdentity({ store_id: 'invalid' })).toThrow(/store_id.*UUID/)
  634. })
  635. it('rejects invalid durable metadata before exposing a session header', async () => {
  636. const path = await freshDbPath('dsh-sqlite-metadata-')
  637. const store = new SqliteStore({ path, journalMode: 'wal', busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS })
  638. const header = meta('invalid-metadata')
  639. await store.appendBatch(header, [chunk(0)], false)
  640. const db = new DatabaseSync(path)
  641. db.prepare(testSql('update-invalid-session-metadata')).run(header.id)
  642. db.close()
  643. await expect(store.list()).rejects.toThrow(/seed_length|origin|delegation_depth/)
  644. await expect(store.loadStored(header.id)).rejects.toThrow(/seed_length|origin|delegation_depth/)
  645. await store.close()
  646. })
  647. it('uses the shared persistence application identity', () => {
  648. expect(SESSION_PERSISTENCE_SQLITE_APPLICATION_ID).toBe(0x44534850)
  649. })
  650. })
  651. describe('SessionPersistenceSqlite edge behavior', () => {
  652. it('materializes an explicitly durable empty live session', async () => {
  653. const path = await freshDbPath('dsh-sqlite-empty-')
  654. const ctx = new Context()
  655. await ctx.plugin(SessionStore)
  656. await ctx.plugin(SessionPersistenceSqlite, { path })
  657. const session = ctx.sessions.create(SessionId('empty'), { meta: { cwd: '/workspace' } })
  658. await ctx.sessionPersistence.ensureMaterialized(session)
  659. await expect(ctx.sessionPersistence.list()).resolves.toEqual([session.header])
  660. await expect(ctx.sessionPersistence.load(session.id)).resolves.toEqual({ meta: session.header, events: [] })
  661. await ctx.fiber.dispose()
  662. })
  663. it('keeps a fresh database unopened until the first persistence operation', async () => {
  664. const path = await freshDbPath('dsh-sqlite-lazy-')
  665. const ctx = new Context()
  666. await ctx.plugin(SessionStore)
  667. await ctx.plugin(SessionPersistenceSqlite, { path })
  668. await expect(stat(path)).rejects.toMatchObject({ code: 'ENOENT' })
  669. const emitWarning = Reflect.get(process, 'emitWarning')
  670. expect(await ctx.sessionPersistence.list()).toEqual([])
  671. expect(Reflect.get(process, 'emitWarning')).toBe(emitWarning)
  672. expect(typeof (await stat(path)).size).toBe('number')
  673. await ctx.fiber.dispose()
  674. })
  675. it('disposes after path validation without opening the database', async () => {
  676. const path = await freshDbPath('dsh-sqlite-unused-')
  677. const ctx = new Context()
  678. await ctx.plugin(SessionStore)
  679. await ctx.plugin(SessionPersistenceSqlite, { path })
  680. await ctx.fiber.dispose()
  681. await expect(stat(path)).rejects.toMatchObject({ code: 'ENOENT' })
  682. const untouchedPath = await freshDbPath('dsh-sqlite-never-validated-')
  683. const untouched = new SqliteStore({
  684. path: untouchedPath,
  685. journalMode: 'wal',
  686. busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS,
  687. })
  688. await untouched.close()
  689. await expect(stat(untouchedPath)).rejects.toMatchObject({ code: 'ENOENT' })
  690. })
  691. it('uses constructor defaults and exposes locate and prepare directly', 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, { path: ':memory:' })
  697. }, { inject: ['sessions'] }))
  698. const header = meta('direct-provider')
  699. const events = chunkLog(3)
  700. expect(persistence.locate(header)).toBeUndefined()
  701. await persistence.create(header)
  702. await persistence.append(header.id, events)
  703. const preparation = await persistence.prepare(header.id)
  704. expect(preparation.session.header).toEqual(header)
  705. preparation[Symbol.dispose]()
  706. await ctx.fiber.dispose()
  707. })
  708. it('keeps empty mutations inert and rolls back a repair without metadata', async () => {
  709. const store = new SqliteStore({
  710. path: ':memory:',
  711. journalMode: 'wal',
  712. busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS,
  713. })
  714. const header = meta('empty-store')
  715. await store.appendBatch(header, [], false)
  716. await store.commitRepair(header, undefined, [])
  717. expect(await store.readStoredRevision(header.id)).toBeUndefined()
  718. await expect(store.commitRepair(header, 0, [])).rejects.toThrow(/metadata row is missing/)
  719. await store.close()
  720. })
  721. it('rejects omitted torn markers and stale closer positions', async () => {
  722. const path = await freshDbPath('dsh-sqlite-repair-validation-')
  723. const store = new SqliteStore({ path, journalMode: 'wal', busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS })
  724. const header = meta('repair-validation')
  725. await store.appendBatch(header, [chunk(0)], false)
  726. const db = new DatabaseSync(path)
  727. db.prepare(testSql('insert-corrupt-event')).run(header.id, 1, 'assistant/chunk', 2, '{not json', 0)
  728. db.close()
  729. await expect(store.commitRepair(header, undefined, [chunk(1)])).rejects.toThrow(/omitted current torn tail/)
  730. await store.commitRepair(header, 1, [])
  731. await expect(store.commitRepair(header, undefined, [chunk(2)])).rejects.toThrow(/closer starts at seq 2/)
  732. const cleared = new DatabaseSync(path)
  733. cleared.prepare(testSql('delete-session-events')).run(header.id)
  734. cleared.close()
  735. await store.commitRepair(header, undefined, [chunk(0)])
  736. expect((await store.loadStored(header.id))?.events).toEqual([chunk(0)])
  737. await store.close()
  738. })
  739. it('rejects malformed physical tail rows before appending', async () => {
  740. const path = await freshDbPath('dsh-sqlite-tail-')
  741. const store = new SqliteStore({ path, journalMode: 'wal', busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS })
  742. const header = meta('invalid-tail')
  743. await store.appendBatch(header, [chunk(0)], false)
  744. const db = new DatabaseSync(path)
  745. db.prepare(testSql('insert-corrupt-event'))
  746. .run(header.id, 1, 'assistant/chunk', 2, '{not json', 0)
  747. db.close()
  748. await expect(store.appendBatch(header, [chunk(2)], true)).rejects.toThrow(/invalid physical tail/)
  749. await store.close()
  750. })
  751. it('rejects missing and empty store identities', async () => {
  752. for (const mode of ['missing', 'empty'] as const) {
  753. const path = await freshDbPath(`dsh-sqlite-identity-${mode}-`)
  754. const db = await openDatabase(DatabaseSync, path, 'wal', DEFAULT_BUSY_TIMEOUT_MS)
  755. if (mode === 'missing') db.exec(testSql('delete-persistence-state'))
  756. else db.exec(testSql('empty-store-id'))
  757. db.close()
  758. await chmod(path, 0o600)
  759. expect(errorMessage(await backendFailure(path))).toMatch(/no valid store identity/)
  760. }
  761. })
  762. it('rejects invalid paths during service initialization', async () => {
  763. const path = await freshDbPath('dsh-sqlite-invalid-path-')
  764. const ctx = new Context()
  765. await ctx.plugin(SessionStore)
  766. await expect(ctx.plugin(SessionPersistenceSqlite, { path: `${path}\0` })).rejects.toMatchObject({
  767. code: 'ERR_INVALID_ARG_VALUE',
  768. })
  769. await ctx.fiber.dispose()
  770. })
  771. it('rejects non-files and symbolic links', async () => {
  772. const directoryPath = await freshDbPath('dsh-sqlite-directory-')
  773. await mkdir(directoryPath)
  774. expect(errorMessage(await backendFailure(directoryPath)))
  775. .toMatch(/must be a regular file/)
  776. const linkPath = await freshDbPath('dsh-sqlite-link-')
  777. const target = join(linkPath, '..', 'target.db')
  778. await writeFile(target, '')
  779. await symlink(target, linkPath)
  780. expect(errorMessage(await backendFailure(linkPath)))
  781. .toMatch(/not a symbolic link/)
  782. const parentLinkPath = await freshDbPath('dsh-sqlite-parent-link-')
  783. const realParent = join(parentLinkPath, '..', 'real-parent')
  784. const linkedParent = join(parentLinkPath, '..', 'linked-parent')
  785. await mkdir(realParent, { mode: 0o700 })
  786. await symlink(realParent, linkedParent)
  787. expect(errorMessage(await backendFailure(join(linkedParent, 'sessions.db'))))
  788. .toMatch(/must be a real directory/)
  789. })
  790. it.runIf(
  791. process.getuid !== undefined && process.getuid() !== 0,
  792. )('rejects permissive files and writable parents', async () => {
  793. const permissivePath = await freshDbPath('dsh-sqlite-permissive-')
  794. await writeFile(permissivePath, '')
  795. await chmod(permissivePath, 0o644)
  796. expect(errorMessage(await backendFailure(permissivePath)))
  797. .toMatch(/accessible only by that user/)
  798. const writableParentPath = await freshDbPath('dsh-sqlite-parent-')
  799. await chmod(join(writableParentPath, '..'), 0o770)
  800. expect(errorMessage(await backendFailure(writableParentPath)))
  801. .toMatch(/not group\/world-writable/)
  802. })
  803. it('surfaces database creation failures after path validation', async () => {
  804. const path = await freshDbPath('dsh-sqlite-create-failure-')
  805. const store = new SqliteStore({ path, journalMode: 'wal', busyTimeoutMs: DEFAULT_BUSY_TIMEOUT_MS })
  806. await store.validatePath()
  807. const parent = join(path, '..')
  808. await rm(parent, { recursive: true })
  809. await writeFile(parent, 'not a directory')
  810. await expect(store.open()).rejects.toThrow(/ENOENT|ENOTDIR/)
  811. await store.close()
  812. })
  813. })