index.ts 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414
  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. override readonly supportsRawArtifacts = false
  94. static inject = ['sessions']
  95. static Config: z<Config> = z.object({
  96. path: z.string().required(),
  97. journalMode: z.union(['wal', 'delete', 'truncate', 'persist'] as const).default('wal'),
  98. preparedSessionCacheSize: z.number().step(1).min(1).default(DEFAULT_PREPARED_SESSION_CACHE_SIZE),
  99. writeBatchMaxDelayMs: z.number().step(1).min(1).max(MAX_WRITE_BATCH_DELAY_MS)
  100. .default(DEFAULT_WRITE_BATCH_MAX_DELAY_MS),
  101. })
  102. /**
  103. * Backend label for the coordinator's dispose diagnostics. Intentionally
  104. * shadows cordis `Service.name` (set to `'sessionPersistence'` by the base);
  105. * see the JSONL backend for why this does not affect service resolution.
  106. */
  107. override readonly name = 'session-persistence-sqlite'
  108. private db!: DatabaseSync
  109. private storeIdentity!: string
  110. private ready: Promise<void>
  111. private coordinator: PersistenceCoordinator<number>
  112. constructor(ctx: Context, public config: Config) {
  113. super(ctx)
  114. // Programmatic wrappers may construct the backend without Schemastery normalization.
  115. const preparedSessionCacheSize = config.preparedSessionCacheSize
  116. ?? DEFAULT_PREPARED_SESSION_CACHE_SIZE
  117. const writeBatchMaxDelayMs = config.writeBatchMaxDelayMs
  118. ?? DEFAULT_WRITE_BATCH_MAX_DELAY_MS
  119. // Open asynchronously so directory creation does not block plugin apply;
  120. // every storage hook awaits the same readiness promise.
  121. this.ready = this.openDb(config.path, (config as Required<Config>).journalMode)
  122. this.coordinator = new PersistenceCoordinator<number>(this.ctx, this, {
  123. preparedSessionCacheSize,
  124. writeBatchMaxDelayMs,
  125. })
  126. }
  127. private async openDb(path: string, journalMode: JournalMode): Promise<void> {
  128. const actual = path === ':memory:' ? path : resolve(path)
  129. if (actual !== ':memory:') {
  130. await mkdir(dirname(actual), { recursive: true, mode: 0o700 })
  131. await createDatabaseFile(actual)
  132. }
  133. this.db = openDatabase(actual, journalMode)
  134. try {
  135. const row = this.db.prepare(
  136. 'SELECT store_id FROM persistence_state WHERE singleton = 1',
  137. ).get() as { store_id: string } | undefined
  138. /* v8 ignore next -- openDatabase inserts the singleton before returning. */
  139. if (row === undefined) {
  140. throw new Error(`session database at "${actual}" has no store identity`)
  141. }
  142. if (row.store_id.length === 0) {
  143. throw new Error(`session database at "${actual}" has no valid store identity`)
  144. }
  145. if (actual !== ':memory:') {
  146. const identity = statSync(actual, { bigint: true })
  147. this.storeIdentity = `file:${identity.dev}:${identity.ino}:${identity.birthtimeNs}:store:${row.store_id}`
  148. } else {
  149. this.storeIdentity = `memory:store:${row.store_id}`
  150. }
  151. } catch (error: unknown) {
  152. this.db.close()
  153. throw error
  154. }
  155. }
  156. // --- SessionPersistence service API (delegated to the coordinator) ---
  157. /** SQLite has one database, not an independent local artifact per session. */
  158. locate(_meta: SessionHeader): SessionLocation | undefined {
  159. return undefined
  160. }
  161. create(meta: SessionHeader): Promise<void> {
  162. return this.coordinator.create(meta)
  163. }
  164. append(id: SessionId, events: readonly SessionEvent[]): Promise<void> {
  165. return this.coordinator.append(id, events)
  166. }
  167. override prepare(id: SessionId, signal?: AbortSignal): Promise<SessionPreparation> {
  168. return this.coordinator.prepare(id, signal)
  169. }
  170. load(id: SessionId): Promise<SessionInspection> {
  171. return this.coordinator.load(id)
  172. }
  173. inspect(id: SessionId, signal?: AbortSignal): Promise<SessionInspection> {
  174. return this.coordinator.inspect(id, signal)
  175. }
  176. readFrom(id: SessionId, fromSeq: number, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  177. return this.coordinator.readFrom(id, fromSeq, signal)
  178. }
  179. // One method serves both public `list` and the backend hook; delegating it to
  180. // the coordinator would call this hook recursively.
  181. // --- PersistenceBackend hooks (the SQLite storage primitives) ---
  182. /** Read a stored prefix by id (ids are globally unique — no scope to scan). */
  183. loadStored(id: SessionId, signal?: AbortSignal): Promise<StoredPrefix<number> | undefined> {
  184. return this.readPrefix(id, signal)
  185. }
  186. /** Read one row's revision without loading its events. */
  187. async readStoredRevision(id: SessionId, signal?: AbortSignal): Promise<PersistenceRevision | undefined> {
  188. signal?.throwIfAborted()
  189. await this.ready
  190. signal?.throwIfAborted()
  191. const row = this.rowFor(id)
  192. return row === undefined ? undefined : sqliteRevision(this.storeIdentity, row)
  193. }
  194. /**
  195. * Seek-capable suffix read: SQL selects `seq >= fromSeq` directly, so the
  196. * read scales with the suffix, not the log. Torn rows past the preserved
  197. * region are dropped, never repaired (non-mutating read).
  198. */
  199. async loadStoredFrom(id: SessionId, fromSeq: number, signal?: AbortSignal): Promise<StoredSuffix | undefined> {
  200. signal?.throwIfAborted()
  201. await this.ready
  202. signal?.throwIfAborted()
  203. const row = this.rowFor(id)
  204. if (row === undefined) return undefined
  205. const meta = rowToMeta(row)
  206. const eventRows = this.db
  207. .prepare('SELECT seq, type, time, data, source_event_seqs, surface_op, ignorable FROM events WHERE session_id = ? AND seq >= ? ORDER BY seq')
  208. .all(id, fromSeq) as unknown as EventRow[]
  209. signal?.throwIfAborted()
  210. const { preserved } = scanRows(eventRows, fromSeq)
  211. return { meta, events: preserved }
  212. }
  213. /**
  214. * Read a session's row + ordered events into a {@link StoredPrefix}. The
  215. * torn-tail marker is the seq from which a never-committed tail must be deleted
  216. * (`scanRows` already returns it as `number | undefined`).
  217. */
  218. private async readPrefix(id: SessionId, signal?: AbortSignal): Promise<StoredPrefix<number> | undefined> {
  219. signal?.throwIfAborted()
  220. await this.ready
  221. signal?.throwIfAborted()
  222. this.db.exec('BEGIN')
  223. let snapshot: { row: SessionRow; eventRows: EventRow[] } | undefined
  224. try {
  225. const row = this.rowFor(id)
  226. if (row !== undefined) {
  227. const eventRows = this.db
  228. .prepare('SELECT seq, type, time, data, source_event_seqs, surface_op, ignorable FROM events WHERE session_id = ? ORDER BY seq')
  229. .all(id) as unknown as EventRow[]
  230. snapshot = { row, eventRows }
  231. }
  232. this.db.exec('COMMIT')
  233. } catch (error: unknown) {
  234. /* v8 ignore start -- synchronous read failures only need transaction cleanup before propagation. */
  235. this.db.exec('ROLLBACK')
  236. throw error
  237. /* v8 ignore stop */
  238. }
  239. signal?.throwIfAborted()
  240. if (snapshot === undefined) return undefined
  241. const { row, eventRows } = snapshot
  242. const { preserved, tornFrom } = scanRows(eventRows)
  243. return {
  244. meta: rowToMeta(row),
  245. events: preserved,
  246. revision: sqliteRevision(this.storeIdentity, row),
  247. ...tornFrom !== undefined ? { tornMarker: tornFrom } : {},
  248. }
  249. }
  250. /**
  251. * Durably append a batch in ONE transaction: materialize the sessions row (if
  252. * lazy) and INSERT every event, or roll back entirely. The transaction is the
  253. * atomicity + durability boundary, so a mid-batch failure (a UNIQUE violation
  254. * on a duplicated seq) leaves the stored log untouched.
  255. */
  256. async appendBatch(meta: SessionHeader, events: readonly SessionEvent[], isMaterialized: boolean): Promise<void> {
  257. await this.ready
  258. const insertEvent = this.db.prepare(
  259. 'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op, ignorable) VALUES (?, ?, ?, ?, ?, ?, ?, ?)',
  260. )
  261. this.db.exec('BEGIN')
  262. try {
  263. if (!isMaterialized) this.writeRow(meta)
  264. for (const event of events) {
  265. const [surfaceSeqs, surfaceOp, ignorable] = envelopeBindings(event)
  266. insertEvent.run(meta.id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp, ignorable)
  267. }
  268. this.db.prepare('UPDATE sessions SET revision = revision + 1 WHERE id = ?').run(meta.id)
  269. this.db.exec('COMMIT')
  270. } catch (error) {
  271. this.db.exec('ROLLBACK')
  272. throw error
  273. }
  274. }
  275. /**
  276. * Make a crash repair durable in ONE transaction: DELETE the torn tail (from
  277. * `tornMarker`) and INSERT the synthetic `closers`. After COMMIT the stored rows
  278. * == the balanced log.
  279. */
  280. async commitRepair(meta: SessionHeader, tornMarker: number | undefined, closers: readonly SessionEvent[]): Promise<void> {
  281. await this.ready
  282. this.db.exec('BEGIN')
  283. try {
  284. if (tornMarker !== undefined) {
  285. this.db.prepare('DELETE FROM events WHERE session_id = ? AND seq >= ?').run(meta.id, tornMarker)
  286. }
  287. if (closers.length > 0) {
  288. const insertEvent = this.db.prepare(
  289. 'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op, ignorable) VALUES (?, ?, ?, ?, ?, ?, ?, ?)',
  290. )
  291. for (const event of closers) {
  292. const [surfaceSeqs, surfaceOp, ignorable] = envelopeBindings(event)
  293. insertEvent.run(meta.id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp, ignorable)
  294. }
  295. }
  296. if (tornMarker !== undefined || closers.length > 0) {
  297. this.db.prepare('UPDATE sessions SET revision = revision + 1 WHERE id = ?').run(meta.id)
  298. }
  299. this.db.exec('COMMIT')
  300. } catch (error) {
  301. // The DELETE+INSERT cannot collide (a row at a closer's seq is preserved or
  302. // deleted as torn first); this rolls back a DB-level failure (disk full,
  303. // etc.), unreachable in test.
  304. /* v8 ignore start */
  305. this.db.exec('ROLLBACK')
  306. throw error
  307. /* v8 ignore stop */
  308. }
  309. }
  310. /** List all materialized sessions' metadata (every row is a materialized session). */
  311. async list(signal?: AbortSignal): Promise<SessionHeader[]> {
  312. signal?.throwIfAborted()
  313. await this.ready
  314. signal?.throwIfAborted()
  315. const rows = this.db
  316. .prepare('SELECT * FROM sessions')
  317. .all() as unknown as SessionRow[]
  318. signal?.throwIfAborted()
  319. return rows.map(rowToMeta)
  320. }
  321. /** List metadata with a source-qualified monotonic revision per session. */
  322. async listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]> {
  323. signal?.throwIfAborted()
  324. await this.ready
  325. signal?.throwIfAborted()
  326. const rows = this.db.prepare('SELECT * FROM sessions').all() as unknown as SessionRow[]
  327. signal?.throwIfAborted()
  328. return rows.map(row => ({
  329. header: rowToMeta(row),
  330. revision: SessionPersistenceRevision(
  331. `${this.storeIdentity}:incarnation:${row.incarnation}:revision:${row.revision}`,
  332. ),
  333. }))
  334. }
  335. /** Close the database handle (awaited by the coordinator's dispose, post-drain). */
  336. async close(): Promise<void> {
  337. await this.ready
  338. this.db.close()
  339. }
  340. // --- row helpers ---
  341. /** Fetch a session's row, or undefined if absent. */
  342. private rowFor(id: SessionId): SessionRow | undefined {
  343. return this.db.prepare('SELECT * FROM sessions WHERE id = ?').get(id) as unknown as SessionRow | undefined
  344. }
  345. /**
  346. * Insert-or-replace a session's metadata row. The only caller is the first
  347. * materializing `appendBatch`, so writing the row IS the materialization (its
  348. * existence is the signal `list` reads).
  349. */
  350. private writeRow(meta: SessionHeader): void {
  351. this.db.prepare(`
  352. INSERT INTO sessions
  353. (id, version, created_at, cwd, parent_session, seed_length, origin, delegation_depth, agent_preset, incarnation, revision)
  354. VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 0)
  355. ON CONFLICT(id) DO UPDATE SET
  356. version = excluded.version,
  357. created_at = excluded.created_at,
  358. cwd = excluded.cwd,
  359. parent_session = excluded.parent_session,
  360. seed_length = excluded.seed_length,
  361. origin = excluded.origin,
  362. delegation_depth = excluded.delegation_depth,
  363. agent_preset = excluded.agent_preset
  364. `).run(
  365. meta.id,
  366. meta.version,
  367. meta.createdAt,
  368. meta.cwd ?? null,
  369. meta.parentSession ?? null,
  370. meta.seedLength ?? null,
  371. meta.origin ?? null,
  372. meta.delegationDepth ?? null,
  373. meta.agentPreset ?? null,
  374. randomUUID(),
  375. )
  376. }
  377. }
  378. export default SessionPersistenceSqlite