index.ts 10 KB

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