unit.ts 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156
  1. /**
  2. * One opened SQLite KV unit: prepared per-table statements over the
  3. * `u_<unit>_<table>` record tables plus this unit's row in the shared
  4. * `unit_globals` table. Each primitive is a single statement, so atomicity
  5. * comes from SQLite itself — no explicit transactions, and no write queue
  6. * (write ordering is the caller's responsibility per the KV contract).
  7. * @module @deepseek-ai/dsh-storage-sqlite/unit
  8. */
  9. import type { DatabaseSync, StatementSync } from 'node:sqlite'
  10. import { StorageError } from '@deepseek-ai/dsh-storage'
  11. import type { KvUnit, KvUnitDescriptor } from '@deepseek-ai/dsh-storage'
  12. import { recordTableName } from './schema.ts'
  13. /** Prepared statements for one declared table. */
  14. interface TableStatements {
  15. upsert: StatementSync
  16. remove: StatementSync
  17. selectAll: StatementSync
  18. }
  19. /**
  20. * The SQLite {@link KvUnit}. Constructed by the backend AFTER the unit's
  21. * record tables exist; statements are prepared once here and reused for every
  22. * primitive. Values are stored as JSON text in the `value` column.
  23. */
  24. export class SqliteKvUnit implements KvUnit {
  25. private readonly tables = new Map<string, TableStatements>()
  26. private readonly globalUpsert: StatementSync | undefined
  27. private readonly globalSelect: StatementSync | undefined
  28. private closed = false
  29. /**
  30. * @param db - Open database handle owned by the backend (never closed here).
  31. * @param descriptor - Validated descriptor whose record tables already exist.
  32. * @param onClose - Backend callback releasing this unit's open-name slot.
  33. */
  34. constructor(
  35. db: DatabaseSync,
  36. private readonly descriptor: KvUnitDescriptor,
  37. private readonly onClose: () => void,
  38. ) {
  39. for (const table of descriptor.tables) {
  40. // Both name segments are validated against UNIT_NAME_RE by the backend,
  41. // so the physical identifier is safe to interpolate into statement text.
  42. const physical = recordTableName(descriptor.name, table)
  43. this.tables.set(table, {
  44. upsert: db.prepare(
  45. `INSERT INTO "${physical}" (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value`,
  46. ),
  47. remove: db.prepare(`DELETE FROM "${physical}" WHERE key = ?`),
  48. selectAll: db.prepare(`SELECT key, value FROM "${physical}"`),
  49. })
  50. }
  51. this.globalUpsert = descriptor.hasGlobal
  52. ? db.prepare(
  53. 'INSERT INTO unit_globals (unit, value) VALUES (?, ?) ON CONFLICT(unit) DO UPDATE SET value = excluded.value',
  54. )
  55. : undefined
  56. this.globalSelect = descriptor.hasGlobal
  57. ? db.prepare('SELECT value FROM unit_globals WHERE unit = ?')
  58. : undefined
  59. }
  60. loadAll(): Promise<{ tables: Record<string, Record<string, unknown>>; global: unknown }> {
  61. return this.settle(() => {
  62. const tables: Record<string, Record<string, unknown>> = {}
  63. for (const [name, statements] of this.tables) {
  64. // Null prototype: record keys are arbitrary strings, so '__proto__'
  65. // must land as an own property instead of mutating the prototype.
  66. const records: Record<string, unknown> = Object.create(null) as Record<string, unknown>
  67. for (const row of statements.selectAll.all() as unknown as Array<{ key: string; value: string }>) {
  68. records[row.key] = this.parseValue(row.value, `table '${name}' key '${row.key}'`)
  69. }
  70. tables[name] = records
  71. }
  72. let global: unknown = null
  73. if (this.globalSelect !== undefined) {
  74. const row = this.globalSelect.get(this.descriptor.name) as { value: string } | undefined
  75. if (row !== undefined) global = this.parseValue(row.value, 'global slot')
  76. }
  77. return { tables, global }
  78. })
  79. }
  80. /** Parse one stored value column, mapping bad JSON to `malformed-medium`. */
  81. private parseValue(text: string, slot: string): unknown {
  82. try {
  83. return JSON.parse(text)
  84. } catch (error) {
  85. throw new StorageError(
  86. 'malformed-medium',
  87. `kv unit '${this.descriptor.name}' holds unparsable JSON at ${slot}`,
  88. { cause: error },
  89. )
  90. }
  91. }
  92. putRecord(table: string, key: string, value: unknown): Promise<void> {
  93. return this.settle(() => {
  94. this.statementsFor(table).upsert.run(key, JSON.stringify(value))
  95. })
  96. }
  97. deleteRecord(table: string, key: string): Promise<void> {
  98. return this.settle(() => {
  99. this.statementsFor(table).remove.run(key)
  100. })
  101. }
  102. setGlobal(value: unknown): Promise<void> {
  103. return this.settle(() => {
  104. if (this.globalUpsert === undefined) {
  105. throw new Error(`kv unit '${this.descriptor.name}' declared no global slot`)
  106. }
  107. this.globalUpsert.run(this.descriptor.name, JSON.stringify(value))
  108. })
  109. }
  110. close(): Promise<void> {
  111. if (!this.closed) {
  112. this.closed = true
  113. this.onClose()
  114. }
  115. return Promise.resolve()
  116. }
  117. /**
  118. * Run one synchronous primitive behind the closed guard, mapping a throw to
  119. * a rejection so the Promise-returning contract never throws synchronously.
  120. */
  121. private settle<T>(operation: () => T): Promise<T> {
  122. try {
  123. this.ensureOpen()
  124. return Promise.resolve(operation())
  125. } catch (error) {
  126. // Non-Error throws can only enter through JSON.stringify propagating a
  127. // value's own toJSON throw; wrap those, preserve every real Error.
  128. return Promise.reject(error instanceof Error ? error : new Error(String(error)))
  129. }
  130. }
  131. private ensureOpen(): void {
  132. if (this.closed) {
  133. throw new StorageError('closed', `kv unit '${this.descriptor.name}' is closed`)
  134. }
  135. }
  136. private statementsFor(table: string): TableStatements {
  137. const statements = this.tables.get(table)
  138. if (statements === undefined) {
  139. throw new Error(`kv unit '${this.descriptor.name}' declared no table '${table}'`)
  140. }
  141. return statements
  142. }
  143. }