index.ts 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412
  1. /**
  2. * SQLite durable session-persistence backend. It maps each session header and
  3. * event to rows, and delegates write-path orchestration to
  4. * {@link PersistenceCoordinator}. It has no independent per-session artifact,
  5. * so its locator returns `undefined`.
  6. * @module @deepseek-ai/dsh-session-persistence-sqlite
  7. */
  8. import { Context } from '@deepseek-ai/cordis'
  9. import z from '@deepseek-ai/schemastery'
  10. import { randomUUID } from 'node:crypto'
  11. import { statSync } from 'node:fs'
  12. import { DatabaseSync } from 'node:sqlite'
  13. import { mkdir, open } from 'node:fs/promises'
  14. import { dirname, resolve } from 'node:path'
  15. import {
  16. DEFAULT_PREPARED_SESSION_CACHE_SIZE, DEFAULT_WRITE_BATCH_MAX_DELAY_MS, MAX_WRITE_BATCH_DELAY_MS,
  17. SessionPersistence, SessionPersistenceRevision, PersistenceCoordinator,
  18. type PersistenceBackend, type SessionLocation, type SessionPersistenceSnapshot,
  19. type SessionInspection, type SessionPersistenceRevision as PersistenceRevision,
  20. type StoredPrefix, type StoredSuffix,
  21. } from '@deepseek-ai/dsh-session-persistence'
  22. import type { SessionEvent, SurfaceEventType, SessionId, SessionHeader, SessionPreparation } from '@deepseek-ai/dsh-session'
  23. import {
  24. type JournalMode, openDatabase, rowToMeta, scanRows, type EventRow, type SessionRow,
  25. } from './schema.ts'
  26. export { SCHEMA_VERSION } from './schema.ts'
  27. /**
  28. * Serialize an event's optional envelope fields for SQL binding. The surface
  29. * fields are nullable TEXT columns — null when the event has no surface
  30. * metadata (non-surface events, events written before surface support); the
  31. * ignorable marker is a nullable INTEGER column — `1` iff the envelope carries
  32. * `ignorable: true`.
  33. */
  34. function envelopeBindings(event: SessionEvent): [string | null, string | null, number | null] {
  35. const se = event as SessionEvent<SurfaceEventType>
  36. return [
  37. se.sourceEventSeqs ? JSON.stringify(se.sourceEventSeqs) : null,
  38. se.surfaceOp !== undefined ? JSON.stringify(se.surfaceOp) : null,
  39. event.ignorable === true ? 1 : null,
  40. ]
  41. }
  42. /** Build the source-qualified revision shared by full and lightweight reads. */
  43. function sqliteRevision(storeIdentity: string, row: SessionRow): PersistenceRevision {
  44. return SessionPersistenceRevision(
  45. `${storeIdentity}:incarnation:${row.incarnation}:revision:${row.revision}`,
  46. )
  47. }
  48. /**
  49. * Exclusively create a missing database file with owner-only permissions.
  50. * Existing files retain their modes, and errors other than `EEXIST` propagate.
  51. * `DatabaseSync` reopens by path, so this does not protect confidentiality or
  52. * integrity when another principal can replace the database entry in its parent
  53. * directory.
  54. */
  55. async function createDatabaseFile(path: string): Promise<void> {
  56. try {
  57. const handle = await open(path, 'wx', 0o600)
  58. await handle.close()
  59. } catch (error) {
  60. if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error
  61. }
  62. }
  63. /** Plugin configuration. */
  64. export interface Config {
  65. /**
  66. * Filesystem path to the SQLite database file. The special value `:memory:`
  67. * opens an in-process database (tests). On filesystems with POSIX modes,
  68. * missing directories and databases are created owner-only; existing path
  69. * modes are preserved. Filesystem setup errors other than an existing database
  70. * fail initialization. The backend does not protect confidentiality or
  71. * integrity when another principal can replace the database entry in its
  72. * parent directory.
  73. */
  74. path: string
  75. /**
  76. * SQLite `journal_mode` pragma. `wal` (the default) is the recorded
  77. * durability model; pick a rollback-journal mode (`delete`/`truncate`/
  78. * `persist`) on filesystems where WAL's shared-memory files do not work
  79. * (network mounts). See {@link JournalMode}.
  80. */
  81. journalMode?: JournalMode
  82. /** Maximum cold Session preparations retained for history-to-resume reuse. */
  83. preparedSessionCacheSize?: number
  84. /** Fixed live-event coalescing window; not a backend completion deadline. */
  85. writeBatchMaxDelayMs?: number
  86. }
  87. /**
  88. * The SQLite persistence backend. Load as a plugin; it registers as
  89. * `ctx.sessionPersistence` and (via the coordinator) installs the write-path
  90. * listeners. Its torn-tail marker is the seq to delete from.
  91. */
  92. export class SessionPersistenceSqlite extends SessionPersistence implements PersistenceBackend<number> {
  93. static inject = ['sessions']
  94. static Config: z<Config> = z.object({
  95. path: z.string().required(),
  96. journalMode: z.union(['wal', 'delete', 'truncate', 'persist'] as const).default('wal'),
  97. preparedSessionCacheSize: z.number().step(1).min(1).default(DEFAULT_PREPARED_SESSION_CACHE_SIZE),
  98. writeBatchMaxDelayMs: z.number().step(1).min(1).max(MAX_WRITE_BATCH_DELAY_MS)
  99. .default(DEFAULT_WRITE_BATCH_MAX_DELAY_MS),
  100. })
  101. /**
  102. * Backend label for the coordinator's dispose diagnostics. Intentionally
  103. * shadows cordis `Service.name` (set to `'sessionPersistence'` by the base);
  104. * see the JSONL backend for why this does not affect service resolution.
  105. */
  106. override readonly name = 'session-persistence-sqlite'
  107. private db!: DatabaseSync
  108. private storeIdentity!: string
  109. private ready: Promise<void>
  110. private coordinator: PersistenceCoordinator<number>
  111. constructor(ctx: Context, public config: Config) {
  112. super(ctx)
  113. // Programmatic wrappers may construct the backend without Schemastery normalization.
  114. const preparedSessionCacheSize = config.preparedSessionCacheSize
  115. ?? DEFAULT_PREPARED_SESSION_CACHE_SIZE
  116. const writeBatchMaxDelayMs = config.writeBatchMaxDelayMs
  117. ?? DEFAULT_WRITE_BATCH_MAX_DELAY_MS
  118. // Open asynchronously so directory creation does not block plugin apply;
  119. // every storage hook awaits the same readiness promise.
  120. this.ready = this.openDb(config.path, (config as Required<Config>).journalMode)
  121. this.coordinator = new PersistenceCoordinator<number>(this.ctx, this, {
  122. preparedSessionCacheSize,
  123. writeBatchMaxDelayMs,
  124. })
  125. }
  126. private async openDb(path: string, journalMode: JournalMode): Promise<void> {
  127. const actual = path === ':memory:' ? path : resolve(path)
  128. if (actual !== ':memory:') {
  129. await mkdir(dirname(actual), { recursive: true, mode: 0o700 })
  130. await createDatabaseFile(actual)
  131. }
  132. this.db = openDatabase(actual, journalMode)
  133. try {
  134. const row = this.db.prepare(
  135. 'SELECT store_id FROM persistence_state WHERE singleton = 1',
  136. ).get() as { store_id: string } | undefined
  137. /* v8 ignore next -- openDatabase inserts the singleton before returning. */
  138. if (row === undefined) {
  139. throw new Error(`session database at "${actual}" has no store identity`)
  140. }
  141. if (row.store_id.length === 0) {
  142. throw new Error(`session database at "${actual}" has no valid store identity`)
  143. }
  144. if (actual !== ':memory:') {
  145. const identity = statSync(actual, { bigint: true })
  146. this.storeIdentity = `file:${identity.dev}:${identity.ino}:${identity.birthtimeNs}:store:${row.store_id}`
  147. } else {
  148. this.storeIdentity = `memory:store:${row.store_id}`
  149. }
  150. } catch (error: unknown) {
  151. this.db.close()
  152. throw error
  153. }
  154. }
  155. // --- SessionPersistence service surface (delegated to the coordinator) ---
  156. /** SQLite has one database, not an independent local artifact per session. */
  157. locate(_meta: SessionHeader): SessionLocation | undefined {
  158. return undefined
  159. }
  160. create(meta: SessionHeader): Promise<void> {
  161. return this.coordinator.create(meta)
  162. }
  163. append(id: SessionId, events: readonly SessionEvent[]): Promise<void> {
  164. return this.coordinator.append(id, events)
  165. }
  166. override prepare(id: SessionId, signal?: AbortSignal): Promise<SessionPreparation> {
  167. return this.coordinator.prepare(id, signal)
  168. }
  169. load(id: SessionId): Promise<SessionInspection> {
  170. return this.coordinator.load(id)
  171. }
  172. inspect(id: SessionId, signal?: AbortSignal): Promise<SessionInspection> {
  173. return this.coordinator.inspect(id, signal)
  174. }
  175. readFrom(id: SessionId, fromSeq: number, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  176. return this.coordinator.readFrom(id, fromSeq, signal)
  177. }
  178. // One method serves both public `list` and the backend hook; delegating it to
  179. // the coordinator would call this hook recursively.
  180. // --- PersistenceBackend hooks (the SQLite storage primitives) ---
  181. /** Read a stored prefix by id (ids are globally unique — no scope to scan). */
  182. loadStored(id: SessionId, signal?: AbortSignal): Promise<StoredPrefix<number> | undefined> {
  183. return this.readPrefix(id, signal)
  184. }
  185. /** Read one row's revision without loading its events. */
  186. async readStoredRevision(id: SessionId, signal?: AbortSignal): Promise<PersistenceRevision | undefined> {
  187. signal?.throwIfAborted()
  188. await this.ready
  189. signal?.throwIfAborted()
  190. const row = this.rowFor(id)
  191. return row === undefined ? undefined : sqliteRevision(this.storeIdentity, row)
  192. }
  193. /**
  194. * Seek-capable suffix read: SQL selects `seq >= fromSeq` directly, so the
  195. * read scales with the suffix, not the log. Torn rows past the preserved
  196. * region are dropped, never repaired (non-mutating read).
  197. */
  198. async loadStoredFrom(id: SessionId, fromSeq: number, signal?: AbortSignal): Promise<StoredSuffix | undefined> {
  199. signal?.throwIfAborted()
  200. await this.ready
  201. signal?.throwIfAborted()
  202. const row = this.rowFor(id)
  203. if (row === undefined) return undefined
  204. const meta = rowToMeta(row)
  205. const eventRows = this.db
  206. .prepare('SELECT seq, type, time, data, source_event_seqs, surface_op, ignorable FROM events WHERE session_id = ? AND seq >= ? ORDER BY seq')
  207. .all(id, fromSeq) as unknown as EventRow[]
  208. signal?.throwIfAborted()
  209. const { preserved } = scanRows(eventRows, fromSeq)
  210. return { meta, events: preserved }
  211. }
  212. /**
  213. * Read a session's row + ordered events into a {@link StoredPrefix}. The
  214. * torn-tail marker is the seq from which a never-committed tail must be deleted
  215. * (`scanRows` already returns it as `number | undefined`).
  216. */
  217. private async readPrefix(id: SessionId, signal?: AbortSignal): Promise<StoredPrefix<number> | undefined> {
  218. signal?.throwIfAborted()
  219. await this.ready
  220. signal?.throwIfAborted()
  221. this.db.exec('BEGIN')
  222. let snapshot: { row: SessionRow; eventRows: EventRow[] } | undefined
  223. try {
  224. const row = this.rowFor(id)
  225. if (row !== undefined) {
  226. const eventRows = this.db
  227. .prepare('SELECT seq, type, time, data, source_event_seqs, surface_op, ignorable FROM events WHERE session_id = ? ORDER BY seq')
  228. .all(id) as unknown as EventRow[]
  229. snapshot = { row, eventRows }
  230. }
  231. this.db.exec('COMMIT')
  232. } catch (error: unknown) {
  233. /* v8 ignore start -- synchronous read failures only need transaction cleanup before propagation. */
  234. this.db.exec('ROLLBACK')
  235. throw error
  236. /* v8 ignore stop */
  237. }
  238. signal?.throwIfAborted()
  239. if (snapshot === undefined) return undefined
  240. const { row, eventRows } = snapshot
  241. const { preserved, tornFrom } = scanRows(eventRows)
  242. return {
  243. meta: rowToMeta(row),
  244. events: preserved,
  245. revision: sqliteRevision(this.storeIdentity, row),
  246. ...tornFrom !== undefined ? { tornMarker: tornFrom } : {},
  247. }
  248. }
  249. /**
  250. * Durably append a batch in ONE transaction: materialize the sessions row (if
  251. * lazy) and INSERT every event, or roll back entirely. The transaction is the
  252. * atomicity + durability boundary, so a mid-batch failure (a UNIQUE violation
  253. * on a duplicated seq) leaves the stored log untouched.
  254. */
  255. async appendBatch(meta: SessionHeader, events: readonly SessionEvent[], isMaterialized: boolean): Promise<void> {
  256. await this.ready
  257. const insertEvent = this.db.prepare(
  258. 'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op, ignorable) VALUES (?, ?, ?, ?, ?, ?, ?, ?)',
  259. )
  260. this.db.exec('BEGIN')
  261. try {
  262. if (!isMaterialized) this.writeRow(meta)
  263. for (const event of events) {
  264. const [surfaceSeqs, surfaceOp, ignorable] = envelopeBindings(event)
  265. insertEvent.run(meta.id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp, ignorable)
  266. }
  267. this.db.prepare('UPDATE sessions SET revision = revision + 1 WHERE id = ?').run(meta.id)
  268. this.db.exec('COMMIT')
  269. } catch (error) {
  270. this.db.exec('ROLLBACK')
  271. throw error
  272. }
  273. }
  274. /**
  275. * Make a crash repair durable in ONE transaction: DELETE the torn tail (from
  276. * `tornMarker`) and INSERT the synthetic `closers`. After COMMIT the stored rows
  277. * == the balanced log.
  278. */
  279. async commitRepair(meta: SessionHeader, tornMarker: number | undefined, closers: readonly SessionEvent[]): Promise<void> {
  280. await this.ready
  281. this.db.exec('BEGIN')
  282. try {
  283. if (tornMarker !== undefined) {
  284. this.db.prepare('DELETE FROM events WHERE session_id = ? AND seq >= ?').run(meta.id, tornMarker)
  285. }
  286. if (closers.length > 0) {
  287. const insertEvent = this.db.prepare(
  288. 'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op, ignorable) VALUES (?, ?, ?, ?, ?, ?, ?, ?)',
  289. )
  290. for (const event of closers) {
  291. const [surfaceSeqs, surfaceOp, ignorable] = envelopeBindings(event)
  292. insertEvent.run(meta.id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp, ignorable)
  293. }
  294. }
  295. if (tornMarker !== undefined || closers.length > 0) {
  296. this.db.prepare('UPDATE sessions SET revision = revision + 1 WHERE id = ?').run(meta.id)
  297. }
  298. this.db.exec('COMMIT')
  299. } catch (error) {
  300. // The DELETE+INSERT cannot collide (a row at a closer's seq is preserved or
  301. // deleted as torn first); this rolls back a DB-level failure (disk full,
  302. // etc.), unreachable in test.
  303. /* v8 ignore start */
  304. this.db.exec('ROLLBACK')
  305. throw error
  306. /* v8 ignore stop */
  307. }
  308. }
  309. /** List all materialized sessions' metadata (every row is a materialized session). */
  310. async list(signal?: AbortSignal): Promise<SessionHeader[]> {
  311. signal?.throwIfAborted()
  312. await this.ready
  313. signal?.throwIfAborted()
  314. const rows = this.db
  315. .prepare('SELECT * FROM sessions')
  316. .all() as unknown as SessionRow[]
  317. signal?.throwIfAborted()
  318. return rows.map(rowToMeta)
  319. }
  320. /** List metadata with a source-qualified monotonic revision per session. */
  321. async listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]> {
  322. signal?.throwIfAborted()
  323. await this.ready
  324. signal?.throwIfAborted()
  325. const rows = this.db.prepare('SELECT * FROM sessions').all() as unknown as SessionRow[]
  326. signal?.throwIfAborted()
  327. return rows.map(row => ({
  328. header: rowToMeta(row),
  329. revision: SessionPersistenceRevision(
  330. `${this.storeIdentity}:incarnation:${row.incarnation}:revision:${row.revision}`,
  331. ),
  332. }))
  333. }
  334. /** Close the database handle (awaited by the coordinator's dispose, post-drain). */
  335. async close(): Promise<void> {
  336. await this.ready
  337. this.db.close()
  338. }
  339. // --- row helpers ---
  340. /** Fetch a session's row, or undefined if absent. */
  341. private rowFor(id: SessionId): SessionRow | undefined {
  342. return this.db.prepare('SELECT * FROM sessions WHERE id = ?').get(id) as unknown as SessionRow | undefined
  343. }
  344. /**
  345. * Insert-or-replace a session's metadata row. The only caller is the first
  346. * materializing `appendBatch`, so writing the row IS the materialization (its
  347. * existence is the signal `list` reads).
  348. */
  349. private writeRow(meta: SessionHeader): void {
  350. this.db.prepare(`
  351. INSERT INTO sessions
  352. (id, version, created_at, cwd, parent_session, seed_length, origin, delegation_depth, agent_preset, incarnation, revision)
  353. VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 0)
  354. ON CONFLICT(id) DO UPDATE SET
  355. version = excluded.version,
  356. created_at = excluded.created_at,
  357. cwd = excluded.cwd,
  358. parent_session = excluded.parent_session,
  359. seed_length = excluded.seed_length,
  360. origin = excluded.origin,
  361. delegation_depth = excluded.delegation_depth,
  362. agent_preset = excluded.agent_preset
  363. `).run(
  364. meta.id,
  365. meta.version,
  366. meta.createdAt,
  367. meta.cwd ?? null,
  368. meta.parentSession ?? null,
  369. meta.seedLength ?? null,
  370. meta.origin ?? null,
  371. meta.delegationDepth ?? null,
  372. meta.agentPreset ?? null,
  373. randomUUID(),
  374. )
  375. }
  376. }
  377. export default SessionPersistenceSqlite