index.ts 13 KB

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