schema.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426
  1. /**
  2. * SQLite schema ownership and durable-row validation.
  3. * @module @deepseek-ai/dsh-session-persistence-sqlite/schema
  4. */
  5. import { randomUUID } from 'node:crypto'
  6. import { isAbsolute } from 'node:path'
  7. import { performance } from 'node:perf_hooks'
  8. import type { DatabaseSync } from 'node:sqlite'
  9. import { setTimeout as delay } from 'node:timers/promises'
  10. import { brandString } from '@deepseek-ai/dsh-brand'
  11. import {
  12. type SessionHeader,
  13. type SessionId,
  14. } from '@deepseek-ai/dsh-session'
  15. import { sql } from './sql.ts'
  16. /** Current physical-record schema with packed and compressed event rows. */
  17. export const SCHEMA_VERSION = 19
  18. /** Application id reserved for DeepSeek Harness SQLite session databases. */
  19. export const SESSION_PERSISTENCE_SQLITE_APPLICATION_ID = 0x44534850
  20. /** A materialized session's metadata and monotonic revision. */
  21. export interface SessionRow {
  22. readonly id: string
  23. readonly version: number
  24. readonly created_at: number
  25. readonly cwd: string | null
  26. readonly parent_session: string | null
  27. readonly seed_length: number | null
  28. readonly origin: 'subagent' | null
  29. readonly incarnation: string
  30. readonly revision: number
  31. readonly delegation_depth: number | null
  32. readonly agent_preset: string | null
  33. }
  34. /** One physical event row; packed rows may represent multiple logical events. */
  35. export interface EventRow {
  36. readonly seq: number
  37. readonly type: string
  38. readonly time: number
  39. readonly data: string | Uint8Array
  40. readonly source_event_seqs: Uint8Array | null
  41. readonly surface_op: string | null
  42. readonly is_packed: 0 | 1
  43. }
  44. /** Durable journal modes accepted by the backend. */
  45. export type JournalMode = 'wal' | 'delete' | 'truncate' | 'persist'
  46. interface SchemaObjectRow {
  47. readonly type: string
  48. readonly name: string
  49. readonly tbl_name: string
  50. readonly sql: string
  51. }
  52. const UUID = /^[0-9a-f]{8}-[0-9a-f]{4}-[1-8][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/iu
  53. const JOURNAL_BUSY_RETRY_INTERVAL_MS = 10
  54. type DatabaseSyncConstructor = typeof import('node:sqlite')['DatabaseSync']
  55. /**
  56. * Open and validate a SQLite session database.
  57. * @param Database - lazily imported Node SQLite constructor.
  58. * @param path - SQLite path, including `:memory:`.
  59. * @param journalMode - validated journal pragma.
  60. * @param busyTimeoutMs - validated maximum wait for a competing SQLite lock.
  61. * @returns the configured database handle.
  62. * @throws when connection settings, schema ownership, or SQLite setup cannot be validated.
  63. */
  64. export async function openDatabase(
  65. Database: DatabaseSyncConstructor,
  66. path: string,
  67. journalMode: JournalMode,
  68. busyTimeoutMs: number,
  69. ): Promise<DatabaseSync> {
  70. const deadline = performance.now() + busyTimeoutMs
  71. const db = new Database(path, { timeout: busyTimeoutMs })
  72. try {
  73. configureConnectionSecurity(db, path)
  74. configureDatabase(Database, db, path)
  75. await selectJournalMode(db, path, journalMode, deadline)
  76. configureDurability(db, path)
  77. return db
  78. } catch (error: unknown) {
  79. db.close()
  80. throw error
  81. }
  82. }
  83. function configureConnectionSecurity(db: DatabaseSync, path: string): void {
  84. db.exec(sql('trusted-schema-off'))
  85. const trustedSchema = integerField(db.prepare(sql('select-trusted-schema')).get(), 'trusted_schema')
  86. /* v8 ignore next 3 -- supported SQLite versions return the fixed setting. */
  87. if (trustedSchema !== 0) {
  88. throw new Error(`session database at "${path}" retained trusted_schema=${trustedSchema}, expected 0`)
  89. }
  90. db.exec(sql('mmap-off'))
  91. if (path === ':memory:') return
  92. const mmapSize = integerField(db.prepare(sql('select-mmap-size')).get(), 'mmap_size')
  93. /* v8 ignore next 3 -- supported file-backed SQLite connections return the fixed setting. */
  94. if (mmapSize !== 0) {
  95. throw new Error(`session database at "${path}" retained mmap_size=${mmapSize}, expected 0`)
  96. }
  97. }
  98. function configureDatabase(
  99. Database: DatabaseSyncConstructor,
  100. db: DatabaseSync,
  101. path: string,
  102. ): void {
  103. db.exec(sql('page-size'))
  104. db.exec(sql('foreign-keys-on'))
  105. let began = false
  106. try {
  107. db.exec(sql('begin-immediate'))
  108. began = true
  109. const onDisk = integerField(db.prepare(sql('select-user-version')).get(), 'user_version')
  110. const applicationId = integerField(db.prepare(sql('select-application-id')).get(), 'application_id')
  111. const userObjectCount = integerField(db.prepare(sql('select-user-object-count')).get(), 'count')
  112. if (onDisk === 0 && (applicationId !== 0 || userObjectCount > 0)) {
  113. throw new Error(`session database at "${path}" has an unversioned schema or application identity`)
  114. }
  115. if (onDisk !== 0 && onDisk !== SCHEMA_VERSION) {
  116. throw new Error(
  117. `session database at "${path}" has schema version ${onDisk}, incompatible with this build (${SCHEMA_VERSION})`,
  118. )
  119. }
  120. if (onDisk !== 0 && applicationId !== SESSION_PERSISTENCE_SQLITE_APPLICATION_ID) {
  121. throw new Error(
  122. `session database at "${path}" has application id ${applicationId}, expected ${SESSION_PERSISTENCE_SQLITE_APPLICATION_ID}`,
  123. )
  124. }
  125. if (onDisk === 0) initializeDatabase(db)
  126. validateRequiredSchema(Database, db, path)
  127. db.exec(sql('commit'))
  128. began = false
  129. } catch (error: unknown) {
  130. /* v8 ignore else -- a failed begin leaves no transaction to roll back. */
  131. if (began) {
  132. /* v8 ignore next 5 -- retain the original ownership failure if rollback fails too. */
  133. try {
  134. db.exec(sql('rollback'))
  135. } catch {
  136. // The original database-ownership failure remains actionable.
  137. }
  138. }
  139. throw error
  140. }
  141. }
  142. async function selectJournalMode(
  143. db: DatabaseSync,
  144. path: string,
  145. journalMode: JournalMode,
  146. deadline: number,
  147. ): Promise<void> {
  148. let result: unknown
  149. while (true) {
  150. try {
  151. result = db.prepare(sql(journalResource(journalMode))).get()
  152. break
  153. } catch (error: unknown) {
  154. const remainingMs = Math.max(0, Math.ceil(deadline - performance.now()))
  155. if (!isSqliteBusy(error) || remainingMs === 0) throw error
  156. await delay(Math.min(JOURNAL_BUSY_RETRY_INTERVAL_MS, remainingMs))
  157. if (performance.now() >= deadline) throw error
  158. }
  159. }
  160. const selected = stringField(result, 'journal_mode').toLowerCase()
  161. const expected = path === ':memory:' ? 'memory' : journalMode
  162. /* v8 ignore next 3 -- SQLite returns the selected mode from these fixed, valid pragmas. */
  163. if (selected !== expected) {
  164. throw new Error(`session database at "${path}" selected journal mode ${selected}, expected ${expected}`)
  165. }
  166. }
  167. function configureDurability(db: DatabaseSync, path: string): void {
  168. db.exec(sql('synchronous-full'))
  169. const synchronous = integerField(db.prepare(sql('select-synchronous')).get(), 'synchronous')
  170. /* v8 ignore next 3 -- supported SQLite versions return the fixed setting. */
  171. if (synchronous !== 2) {
  172. throw new Error(`session database at "${path}" retained synchronous=${synchronous}, expected FULL (2)`)
  173. }
  174. }
  175. function isSqliteBusy(error: unknown): boolean {
  176. return typeof error === 'object'
  177. && error !== null
  178. && Reflect.get(error, 'errcode') === 5
  179. }
  180. function journalResource(mode: JournalMode):
  181. | 'journal-mode-wal'
  182. | 'journal-mode-delete'
  183. | 'journal-mode-truncate'
  184. | 'journal-mode-persist' {
  185. switch (mode) {
  186. case 'wal': return 'journal-mode-wal'
  187. case 'delete': return 'journal-mode-delete'
  188. case 'truncate': return 'journal-mode-truncate'
  189. case 'persist': return 'journal-mode-persist'
  190. }
  191. }
  192. function initializeDatabase(db: DatabaseSync): void {
  193. db.exec(sql('schema'))
  194. db.prepare(sql('insert-persistence-state')).run(randomUUID())
  195. db.exec(sql('set-application-id'))
  196. db.exec(sql('set-user-version-19'))
  197. }
  198. let canonicalSchema: readonly SchemaObjectRow[] | undefined
  199. function expectedSchema(Database: DatabaseSyncConstructor): readonly SchemaObjectRow[] {
  200. if (canonicalSchema !== undefined) return canonicalSchema
  201. const reference = new Database(':memory:')
  202. try {
  203. reference.exec(sql('foreign-keys-on'))
  204. reference.exec(sql('schema'))
  205. canonicalSchema = schemaObjects(reference)
  206. return canonicalSchema
  207. } finally {
  208. reference.close()
  209. }
  210. }
  211. function schemaObjects(db: DatabaseSync): SchemaObjectRow[] {
  212. return db.prepare(sql('select-schema-objects')).all().map((value) => {
  213. const row = record(value, 'schema object')
  214. return {
  215. type: stringField(row, 'type'),
  216. name: stringField(row, 'name'),
  217. tbl_name: stringField(row, 'tbl_name'),
  218. sql: normalizeSql(stringField(row, 'sql')),
  219. }
  220. })
  221. }
  222. function normalizeSql(value: string): string {
  223. return value.replaceAll(/\s+/gu, ' ').trim()
  224. }
  225. function validateRequiredSchema(
  226. Database: DatabaseSyncConstructor,
  227. db: DatabaseSync,
  228. path: string,
  229. ): void {
  230. if (JSON.stringify(schemaObjects(db)) !== JSON.stringify(expectedSchema(Database))) {
  231. throw new Error(`session database at "${path}" does not contain the required schema objects`)
  232. }
  233. }
  234. /**
  235. * Recheck schema ownership inside the caller's mutation transaction.
  236. * @param Database - constructor used to validate the canonical schema.
  237. * @param db - open owned database with an active immediate transaction.
  238. * @param path - database location used in ownership diagnostics.
  239. * @throws when another writer changed the application identity, schema, or version.
  240. */
  241. export function validateSchemaForMutation(
  242. Database: DatabaseSyncConstructor,
  243. db: DatabaseSync,
  244. path: string,
  245. ): void {
  246. const version = integerField(db.prepare(sql('select-user-version')).get(), 'user_version')
  247. const applicationId = integerField(db.prepare(sql('select-application-id')).get(), 'application_id')
  248. if (applicationId !== SESSION_PERSISTENCE_SQLITE_APPLICATION_ID) {
  249. throw new Error(
  250. `session database application id changed before mutation (expected ${SESSION_PERSISTENCE_SQLITE_APPLICATION_ID}, got ${applicationId})`,
  251. )
  252. }
  253. validateRequiredSchema(Database, db, path)
  254. if (version !== SCHEMA_VERSION) {
  255. throw new Error(`session database schema changed before mutation (expected ${SCHEMA_VERSION}, got ${version})`)
  256. }
  257. }
  258. /**
  259. * Decode and validate one durable session row.
  260. * @param value - value returned by SQLite.
  261. * @returns a validated session row.
  262. */
  263. export function decodeSessionRow(value: unknown): SessionRow {
  264. const row = record(value, 'stored session metadata')
  265. const id = nonemptyStringField(row, 'id')
  266. const version = safeIntegerField(row, 'version')
  267. const cwd = nullableStringField(row, 'cwd')
  268. if (cwd !== null && !isAbsolute(cwd)) throw new Error('stored session cwd must be absolute')
  269. const parent = nullableStringField(row, 'parent_session')
  270. const origin = nullableStringField(row, 'origin')
  271. if (origin !== null && origin !== 'subagent') throw new Error('stored session origin must be subagent or null')
  272. const incarnation = nonemptyStringField(row, 'incarnation')
  273. if (!UUID.test(incarnation)) throw new Error('stored session incarnation must be a UUID')
  274. return {
  275. id,
  276. version,
  277. created_at: nonnegativeSafeIntegerField(row, 'created_at'),
  278. cwd,
  279. parent_session: parent,
  280. seed_length: nullableNonnegativeSafeIntegerField(row, 'seed_length'),
  281. origin,
  282. delegation_depth: nullableNonnegativeSafeIntegerField(row, 'delegation_depth'),
  283. agent_preset: nullableStringField(row, 'agent_preset'),
  284. incarnation,
  285. revision: nonnegativeSafeIntegerField(row, 'revision'),
  286. }
  287. }
  288. /**
  289. * Decode and validate one durable event row before JSON interpretation.
  290. * @param value - value returned by SQLite.
  291. * @returns a validated physical event row.
  292. */
  293. export function decodeEventRow(value: unknown): EventRow {
  294. const row = record(value, 'stored event')
  295. const isPacked = safeIntegerField(row, 'is_packed')
  296. if (isPacked !== 0 && isPacked !== 1) {
  297. throw new Error('stored event is_packed must be 0 or 1')
  298. }
  299. return {
  300. seq: nonnegativeSafeIntegerField(row, 'seq'),
  301. type: nonemptyStringField(row, 'type'),
  302. time: safeIntegerField(row, 'time'),
  303. data: stringOrBlobField(row, 'data'),
  304. source_event_seqs: nullableBlobField(row, 'source_event_seqs'),
  305. surface_op: nullableStringField(row, 'surface_op'),
  306. is_packed: isPacked,
  307. }
  308. }
  309. /**
  310. * Validate the singleton identity read from durable storage.
  311. * @param value - value returned by SQLite.
  312. * @returns the UUID store identity.
  313. */
  314. export function decodeStoreIdentity(value: unknown): string {
  315. const identity = nonemptyStringField(value, 'store_id')
  316. if (!UUID.test(identity)) throw new Error('stored store_id must be a UUID')
  317. return identity
  318. }
  319. /**
  320. * Reconstruct an immutable session header from a validated metadata row.
  321. * @param row - validated stored metadata row.
  322. * @returns the session header.
  323. */
  324. export function rowToMeta(row: SessionRow): SessionHeader {
  325. return {
  326. version: row.version,
  327. id: brandString<SessionId>(row.id),
  328. createdAt: row.created_at,
  329. ...row.cwd === null ? {} : { cwd: row.cwd },
  330. ...row.parent_session === null ? {} : { parentSession: brandString<SessionId>(row.parent_session) },
  331. ...row.seed_length === null ? {} : { seedLength: row.seed_length },
  332. ...row.origin === null ? {} : { origin: row.origin },
  333. ...row.delegation_depth === null ? {} : { delegationDepth: row.delegation_depth },
  334. ...row.agent_preset === null ? {} : { agentPreset: row.agent_preset },
  335. }
  336. }
  337. function record(value: unknown, label: string): Record<string, unknown> {
  338. if (typeof value !== 'object' || value === null) throw new Error(`${label} must be an object`)
  339. return value as Record<string, unknown>
  340. }
  341. function stringField(value: unknown, key: string): string {
  342. const field = record(value, 'SQLite row')[key]
  343. if (typeof field !== 'string') throw new Error(`stored ${key} must be a string`)
  344. return field
  345. }
  346. function nonemptyStringField(value: unknown, key: string): string {
  347. const field = stringField(value, key)
  348. if (field.length === 0) throw new Error(`stored ${key} must not be empty`)
  349. return field
  350. }
  351. function nullableStringField(value: unknown, key: string): string | null {
  352. const field = record(value, 'SQLite row')[key]
  353. if (field === null) return null
  354. if (typeof field !== 'string') throw new Error(`stored ${key} must be a string or null`)
  355. return field
  356. }
  357. function stringOrBlobField(value: unknown, key: string): string | Uint8Array {
  358. const field = record(value, 'SQLite row')[key]
  359. if (typeof field === 'string' || field instanceof Uint8Array) return field
  360. throw new Error(`stored ${key} must be a string or blob`)
  361. }
  362. function nullableBlobField(value: unknown, key: string): Uint8Array | null {
  363. const field = record(value, 'SQLite row')[key]
  364. if (field === null || field instanceof Uint8Array) return field
  365. throw new Error(`stored ${key} must be a blob or null`)
  366. }
  367. function integerField(value: unknown, key: string): number {
  368. const field = record(value, 'SQLite row')[key]
  369. if (!Number.isSafeInteger(field)) throw new Error(`stored ${key} must be a safe integer`)
  370. return field as number
  371. }
  372. function safeIntegerField(value: unknown, key: string): number {
  373. return integerField(value, key)
  374. }
  375. function nonnegativeSafeIntegerField(value: unknown, key: string): number {
  376. const field = integerField(value, key)
  377. if (field < 0) throw new Error(`stored ${key} must be non-negative`)
  378. return field
  379. }
  380. function nullableSafeIntegerField(value: unknown, key: string): number | null {
  381. const field = record(value, 'SQLite row')[key]
  382. if (field === null) return null
  383. if (!Number.isSafeInteger(field)) throw new Error(`stored ${key} must be a safe integer or null`)
  384. return field as number
  385. }
  386. function nullableNonnegativeSafeIntegerField(value: unknown, key: string): number | null {
  387. const field = nullableSafeIntegerField(value, key)
  388. if (field !== null && field < 0) throw new Error(`stored ${key} must be non-negative or null`)
  389. return field
  390. }