index.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272
  1. /**
  2. * SQLite durable session-persistence backend (`@deepseek-ai/dsh-session-persistence-sqlite`).
  3. *
  4. * A SECOND {@link SessionPersistence} implementation, built to validate that the
  5. * abstract seam + the shared `runPersistenceContract` suite are genuinely
  6. * backend-agnostic: the same append-only / contiguous-seq / lazy-materialization
  7. * / interrupted-turn-close-on-load semantics the JSONL backend expresses over
  8. * file bytes, expressed here over `node:sqlite` rows. Each `SessionEvent` maps
  9. * 1:1 onto a row `(session_id, seq, type, time, data, source_event_seqs, surface_op)`.
  10. *
  11. * Like the JSONL backend it supplies ONLY the storage primitives (the
  12. * {@link PersistenceBackend} hooks below — INSERT/DELETE/SELECT inside
  13. * transactions); all the write-path orchestration lives in the backend-agnostic
  14. * {@link PersistenceCoordinator} this class composes. The four public
  15. * {@link SessionPersistence} methods delegate to the coordinator.
  16. *
  17. * @module @deepseek-ai/dsh-session-persistence-sqlite
  18. */
  19. import { Context } from 'cordis'
  20. import z from 'schemastery'
  21. import { DatabaseSync } from 'node:sqlite'
  22. import { mkdir } from 'node:fs/promises'
  23. import { dirname, resolve } from 'node:path'
  24. import {
  25. SessionPersistence, PersistenceCoordinator,
  26. type PersistenceBackend, type StoredPrefix,
  27. } from '@deepseek-ai/dsh-session-persistence'
  28. import type { Session, SessionEvent, SurfaceEventType, SessionId, SessionHeader } from '@deepseek-ai/dsh-session'
  29. import {
  30. type JournalMode, openDatabase, rowToMeta, scanRows, type EventRow, type SessionRow,
  31. } from './schema.ts'
  32. export { SCHEMA_VERSION } from './schema.ts'
  33. /**
  34. * Serialize an event's surface-metadata fields for SQL binding. Both fields are
  35. * nullable TEXT columns — null when the event has no surface metadata (non-surface
  36. * events, events written before surface support).
  37. */
  38. function surfaceBindings(event: SessionEvent): [string | null, string | null] {
  39. const se = event as SessionEvent<SurfaceEventType>
  40. return [
  41. se.sourceEventSeqs ? JSON.stringify(se.sourceEventSeqs) : null,
  42. se.surfaceOp !== undefined ? JSON.stringify(se.surfaceOp) : null,
  43. ]
  44. }
  45. /** Plugin configuration. */
  46. export interface Config {
  47. /**
  48. * Filesystem path to the SQLite database file. The special value `:memory:`
  49. * opens an in-process database (tests); a file path is created (with parent
  50. * dirs) on construction.
  51. */
  52. path: string
  53. /**
  54. * SQLite `journal_mode` pragma. `wal` (the default) is the recorded
  55. * durability model; pick a rollback-journal mode (`delete`/`truncate`/
  56. * `persist`) on filesystems where WAL's shared-memory files do not work
  57. * (network mounts). See {@link JournalMode}.
  58. */
  59. journalMode?: JournalMode
  60. }
  61. /**
  62. * The SQLite persistence backend. Load as a plugin; it registers as
  63. * `ctx.sessionPersistence` and (via the coordinator) installs the write-path
  64. * listeners. Its torn-tail marker is the seq to delete from.
  65. */
  66. export class SessionPersistenceSqlite extends SessionPersistence implements PersistenceBackend<number> {
  67. static inject = ['sessions']
  68. static Config: z<Config> = z.object({
  69. path: z.string().required(),
  70. journalMode: z.union(['wal', 'delete', 'truncate', 'persist'] as const).default('wal'),
  71. })
  72. /**
  73. * Backend label for the coordinator's dispose diagnostics. Intentionally
  74. * shadows cordis `Service.name` (set to `'sessionPersistence'` by the base);
  75. * see the JSONL backend for why this does not affect service resolution.
  76. */
  77. override readonly name = 'session-persistence-sqlite'
  78. private db!: DatabaseSync
  79. private ready: Promise<void>
  80. private coordinator: PersistenceCoordinator<number>
  81. constructor(ctx: Context, public config: Config) {
  82. super(ctx)
  83. // Open the database asynchronously (the parent directory may need creating);
  84. // every hook awaits `ready` first. Opening synchronously would force a sync
  85. // mkdir and block plugin apply. schemastery (static Config) has already
  86. // filled `journalMode`; the cast records that runtime fact.
  87. this.ready = this.openDb(config.path, (config as Required<Config>).journalMode)
  88. this.coordinator = new PersistenceCoordinator<number>(this.ctx, this)
  89. }
  90. private async openDb(path: string, journalMode: JournalMode): Promise<void> {
  91. if (path !== ':memory:') {
  92. const abs = resolve(path)
  93. await mkdir(dirname(abs), { recursive: true, mode: 0o700 })
  94. this.db = openDatabase(abs, journalMode)
  95. } else {
  96. this.db = openDatabase(path, journalMode)
  97. }
  98. }
  99. // --- SessionPersistence service surface (delegated to the coordinator) ---
  100. create(meta: SessionHeader): Promise<void> {
  101. return this.coordinator.create(meta)
  102. }
  103. append(id: SessionId, events: readonly SessionEvent[]): Promise<void> {
  104. return this.coordinator.append(id, events)
  105. }
  106. load(id: SessionId): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  107. return this.coordinator.load(id)
  108. }
  109. // `list` is BOTH the public service method and the PersistenceBackend hook —
  110. // one method (the SELECT below). The coordinator adds no orchestration for
  111. // listing, so routing it through the coordinator would just recurse. Defined
  112. // once, in the "PersistenceBackend hooks" section.
  113. /**
  114. * The per-session init promises, exposed for white-box tests that await a
  115. * specific session's onCreated (there is no public API to await one init).
  116. */
  117. get inits(): Map<Session, Promise<void>> {
  118. return this.coordinator.inits
  119. }
  120. // --- PersistenceBackend hooks (the SQLite storage primitives) ---
  121. /** Read a stored prefix by id (ids are globally unique — no scope to scan). */
  122. loadStored(id: SessionId): Promise<StoredPrefix<number> | undefined> {
  123. return this.readPrefix(id)
  124. }
  125. /** Read a stored prefix; `cwd` is ignored (the id is globally unique in SQLite). */
  126. loadLive(id: SessionId, _cwd: string | undefined): Promise<StoredPrefix<number> | undefined> {
  127. return this.readPrefix(id)
  128. }
  129. /**
  130. * Read a session's row + ordered events into a {@link StoredPrefix}. The
  131. * torn-tail marker is the seq from which a never-committed tail must be deleted
  132. * (`scanRows` already returns it as `number | undefined`).
  133. */
  134. private async readPrefix(id: SessionId): Promise<StoredPrefix<number> | undefined> {
  135. await this.ready
  136. const row = this.rowFor(id)
  137. if (row === undefined) return undefined
  138. const meta = rowToMeta(row)
  139. const eventRows = this.db
  140. .prepare('SELECT seq, type, time, data, source_event_seqs, surface_op FROM events WHERE session_id = ? ORDER BY seq')
  141. .all(id) as unknown as EventRow[]
  142. const { preserved, tornFrom } = scanRows(eventRows)
  143. return { meta, events: preserved, ...tornFrom !== undefined ? { tornMarker: tornFrom } : {} }
  144. }
  145. /**
  146. * Durably append a batch in ONE transaction: materialize the sessions row (if
  147. * lazy) and INSERT every event, or roll back entirely. The transaction is the
  148. * atomicity + durability boundary, so a mid-batch failure (a UNIQUE violation
  149. * on a duplicated seq) leaves the stored log untouched.
  150. */
  151. async appendBatch(meta: SessionHeader, events: readonly SessionEvent[], isMaterialized: boolean): Promise<void> {
  152. await this.ready
  153. const insertEvent = this.db.prepare(
  154. 'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op) VALUES (?, ?, ?, ?, ?, ?, ?)',
  155. )
  156. this.db.exec('BEGIN')
  157. try {
  158. if (!isMaterialized) this.writeRow(meta)
  159. for (const event of events) {
  160. const [surfaceSeqs, surfaceOp] = surfaceBindings(event)
  161. insertEvent.run(meta.id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp)
  162. }
  163. this.db.exec('COMMIT')
  164. } catch (error) {
  165. this.db.exec('ROLLBACK')
  166. throw error
  167. }
  168. }
  169. /**
  170. * Make a crash repair durable in ONE transaction: DELETE the torn tail (from
  171. * `tornMarker`) and INSERT the synthetic `closers`. After COMMIT the stored rows
  172. * == the balanced log.
  173. */
  174. async commitRepair(meta: SessionHeader, tornMarker: number | undefined, closers: readonly SessionEvent[]): Promise<void> {
  175. await this.ready
  176. this.db.exec('BEGIN')
  177. try {
  178. if (tornMarker !== undefined) {
  179. this.db.prepare('DELETE FROM events WHERE session_id = ? AND seq >= ?').run(meta.id, tornMarker)
  180. }
  181. if (closers.length > 0) {
  182. const insertEvent = this.db.prepare(
  183. 'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op) VALUES (?, ?, ?, ?, ?, ?, ?)',
  184. )
  185. for (const event of closers) {
  186. const [surfaceSeqs, surfaceOp] = surfaceBindings(event)
  187. insertEvent.run(meta.id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp)
  188. }
  189. }
  190. this.db.exec('COMMIT')
  191. } catch (error) {
  192. // The DELETE+INSERT cannot collide (a row at a closer's seq is preserved or
  193. // deleted as torn first); this rolls back a DB-level failure (disk full,
  194. // etc.), unreachable in test.
  195. /* v8 ignore start */
  196. this.db.exec('ROLLBACK')
  197. throw error
  198. /* v8 ignore stop */
  199. }
  200. }
  201. /** List all materialized sessions' metadata (every row is a materialized session). */
  202. async list(): Promise<SessionHeader[]> {
  203. await this.ready
  204. const rows = this.db
  205. .prepare('SELECT * FROM sessions')
  206. .all() as unknown as SessionRow[]
  207. return rows.map(rowToMeta)
  208. }
  209. /** Close the database handle (awaited by the coordinator's dispose, post-drain). */
  210. async close(): Promise<void> {
  211. await this.ready
  212. this.db.close()
  213. }
  214. // --- row helpers ---
  215. /** Fetch a session's row, or undefined if absent. */
  216. private rowFor(id: SessionId): SessionRow | undefined {
  217. return this.db.prepare('SELECT * FROM sessions WHERE id = ?').get(id) as unknown as SessionRow | undefined
  218. }
  219. /**
  220. * Insert-or-replace a session's metadata row. The only caller is the first
  221. * materializing `appendBatch`, so writing the row IS the materialization (its
  222. * existence is the signal `list` reads).
  223. */
  224. private writeRow(meta: SessionHeader): void {
  225. this.db.prepare(`
  226. INSERT INTO sessions (id, version, created_at, cwd, parent_session, seed_length)
  227. VALUES (?, ?, ?, ?, ?, ?)
  228. ON CONFLICT(id) DO UPDATE SET
  229. version = excluded.version,
  230. created_at = excluded.created_at,
  231. cwd = excluded.cwd,
  232. parent_session = excluded.parent_session,
  233. seed_length = excluded.seed_length
  234. `).run(
  235. meta.id,
  236. meta.version,
  237. meta.createdAt,
  238. meta.cwd ?? null,
  239. meta.parentSession ?? null,
  240. meta.seedLength ?? null,
  241. )
  242. }
  243. }
  244. export default SessionPersistenceSqlite