返回源码地图

packages/storage/storage-sqlite/src/unit.ts

main snapshot · da00f7f5358f · 正文引用章节 16;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

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