1
/**2
* One opened SQLite KV unit: prepared per-table statements over the3
* `u_<unit>_<table>` record tables plus this unit's row in the shared4
* `unit_globals` table. Each primitive is a single statement, so atomicity5
* comes from SQLite itself — no explicit transactions, and no write queue6
* (write ordering is the caller's responsibility per the KV contract).7
* @module @deepseek-ai/dsh-storage-sqlite/unit8
*/10
import type { DatabaseSync, StatementSync } from 'node:sqlite'11
import { StorageError } from '@deepseek-ai/dsh-storage'12
import type { KvUnit, KvUnitDescriptor } from '@deepseek-ai/dsh-storage'13
import { recordTableName } from './schema.ts'15
/** Prepared statements for one declared table. */16
interface TableStatements {17
upsert: StatementSync18
remove: StatementSync19
selectAll: StatementSync20
}22
/**23
* The SQLite {@link KvUnit}. Constructed by the backend AFTER the unit's24
* record tables exist; statements are prepared once here and reused for every25
* primitive. Values are stored as JSON text in the `value` column.26
*/27
export class SqliteKvUnit implements KvUnit {28
private readonly tables = new Map<string, TableStatements>()29
private readonly globalUpsert: StatementSync | undefined30
private readonly globalSelect: StatementSync | undefined31
private closed = false33
/**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.hasGlobal56
? db.prepare(57
'INSERT INTO unit_globals (unit, value) VALUES (?, ?) ON CONFLICT(unit) DO UPDATE SET value = excluded.value',58
)59
: undefined60
this.globalSelect = descriptor.hasGlobal61
? db.prepare('SELECT value FROM unit_globals WHERE unit = ?')62
: undefined63
}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] = records76
}77
let global: unknown = null78
if (this.globalSelect !== undefined) {79
const row = this.globalSelect.get(this.descriptor.name) as { value: string } | undefined80
if (row !== undefined) global = this.parseValue(row.value, 'global slot')81
}82
return { tables, global }83
})84
}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
}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
}105
deleteRecord(table: string, key: string): Promise<void> {106
return this.settle(() => {107
this.statementsFor(table).remove.run(key)108
})109
}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
}120
close(): Promise<void> {121
if (!this.closed) {122
this.closed = true123
this.onClose()124
}125
return Promise.resolve()126
}128
/**129
* Run one synchronous primitive behind the closed guard, mapping a throw to130
* 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 a138
// 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
}143
private ensureOpen(): void {144
if (this.closed) {145
throw new StorageError('closed', `kv unit '${this.descriptor.name}' is closed`)146
}147
}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 statements155
}156
}