1
/**2
* SQLite storage backend for the storage hub: one database file hosts every3
* routed unit, document-per-row (`key TEXT` / `value TEXT` JSON). Registers4
* as backend `sqlite`; the disposer unregisters first, then closes the medium.5
* @module @deepseek-ai/dsh-storage-sqlite6
*/8
import type { Context } from '@deepseek-ai/cordis'9
import z from '@deepseek-ai/schemastery'10
import type { DatabaseSync } from 'node:sqlite'11
import { StorageError, UNIT_NAME_RE, storageBackendServiceKey } from '@deepseek-ai/dsh-storage'12
import type { KvFacet, KvUnit, KvUnitDescriptor, StorageBackend } from '@deepseek-ai/dsh-storage'13
import { openDatabase, recordTableName, type JournalMode } from './schema.ts'14
import { SqliteKvUnit } from './unit.ts'16
export { STORAGE_SQLITE_SCHEMA_VERSION, type JournalMode } from './schema.ts'18
/** Cordis plugin name. */19
export const name = 'storage-sqlite'20
/** The backend registers on the storage hub. */21
export const inject = ['storage']23
/** Plugin configuration. */24
export interface Config {25
/**26
* Filesystem path to the SQLite database file. The special value `:memory:`27
* opens an in-process database (tests). On filesystems with POSIX modes,28
* missing directories and databases are created owner-only; existing path29
* modes are preserved. Filesystem setup errors other than an existing30
* database fail the open. The backend does not protect confidentiality or31
* integrity when another principal can replace the database entry in its32
* parent directory.33
*/34
path: string35
/**36
* SQLite `journal_mode` pragma. `wal` (the default) suits local disks; pick37
* a rollback-journal mode (`delete`/`truncate`/`persist`) on filesystems38
* where WAL's shared-memory files do not work (network mounts). See39
* {@link JournalMode}.40
*/41
journalMode?: JournalMode42
}44
/** Schemastery validator for {@link Config}. */45
export const Config: z<Config> = z.object({46
path: z.string().required(),47
journalMode: z.union(['wal', 'delete', 'truncate', 'persist'] as const).default('wal'),48
})50
/**51
* The SQLite {@link StorageBackend}. Owns one `DatabaseSync` connection and52
* the open-unit table; `kv.open` validates names, enforces the per-unit53
* version stamp in `units`, and ensures the unit's record tables.54
*/55
export class SqliteStorageBackend implements StorageBackend {56
/** The key-value facet; the only shape this backend serves. */57
readonly kv: KvFacet = { open: descriptor => this.openUnit(descriptor) }59
private readonly ready: Promise<DatabaseSync>60
/** Open (or still-opening) units by name; presence is the double-open guard. */61
private readonly units = new Map<string, Promise<SqliteKvUnit>>()62
private closing: Promise<void> | undefined64
/**65
* @param config - Validated plugin configuration.66
*/67
constructor(config: Config) {68
this.ready = openDatabase(config.path, (config as Required<Config>).journalMode)69
// Mark the rejection handled: every primitive re-awaits `ready`, so an70
// open failure still surfaces to each caller; this guard only prevents an71
// unhandled-rejection crash when the failure precedes the first use.72
this.ready.catch(() => {})73
}75
private openUnit(descriptor: KvUnitDescriptor): Promise<KvUnit> {76
if (this.closing !== undefined) {77
return Promise.reject(new StorageError('closed', 'sqlite storage backend is closed'))78
}79
if (!UNIT_NAME_RE.test(descriptor.name)) {80
return Promise.reject(new Error(`kv unit name '${descriptor.name}' violates ${UNIT_NAME_RE}`))81
}82
for (const table of descriptor.tables) {83
if (!UNIT_NAME_RE.test(table)) {84
return Promise.reject(new Error(`kv table name '${table}' in unit '${descriptor.name}' violates ${UNIT_NAME_RE}`))85
}86
}87
if (this.units.has(descriptor.name)) {88
return Promise.reject(new Error(`kv unit '${descriptor.name}' is already open (double-open is a caller bug)`))89
}90
// Reserve the name synchronously so a concurrent second open of the same91
// name rejects instead of racing past the guard during the awaits below.92
const pending = this.materializeUnit(descriptor)93
this.units.set(descriptor.name, pending)94
pending.catch(() => this.units.delete(descriptor.name))95
return pending96
}98
private async materializeUnit(descriptor: KvUnitDescriptor): Promise<SqliteKvUnit> {99
const db = await this.ready100
const row = db.prepare('SELECT version FROM units WHERE name = ?').get(descriptor.name) as101
| { version: number }102
| undefined103
if (row === undefined) {104
db.prepare('INSERT INTO units (name, version) VALUES (?, ?)').run(descriptor.name, descriptor.version)105
} else if (row.version !== descriptor.version) {106
throw new StorageError(107
'version-mismatch',108
`kv unit '${descriptor.name}' is stamped version ${row.version} on the medium, incompatible with descriptor version ${descriptor.version}`,109
)110
}111
for (const table of descriptor.tables) {112
// Both segments passed UNIT_NAME_RE, so the identifier is safe in DDL.113
db.exec(`114
CREATE TABLE IF NOT EXISTS "${recordTableName(descriptor.name, table)}" (115
key TEXT PRIMARY KEY,116
value TEXT NOT NULL117
) STRICT118
`)119
}120
return new SqliteKvUnit(db, descriptor, () => {121
this.units.delete(descriptor.name)122
})123
}125
/**126
* Close every open unit and release the database. Idempotent; concurrent127
* and repeated calls resolve once teardown finishes.128
* @returns resolution after the medium is released.129
*/130
close(): Promise<void> {131
this.closing ??= this.doClose()132
return this.closing133
}135
private async doClose(): Promise<void> {136
let db: DatabaseSync137
try {138
db = await this.ready139
} catch {140
// The medium never opened; that failure already rejected the opener and141
// every unit call, so there is nothing left to release here.142
return143
}144
for (const pending of [...this.units.values()]) {145
const unit = await pending.catch(() => undefined)146
await unit?.close()147
}148
db.close()149
}150
}152
/**153
* Register the SQLite backend as `sqlite` on the storage hub. The disposer154
* unregisters the name first, then closes the backend.155
* @param ctx - Plugin context (must inject `storage`).156
* @param config - Validated plugin configuration.157
*/158
export function apply(ctx: Context, config: Config) {159
const backend = new SqliteStorageBackend(config)160
ctx.effect(() => {161
const dispose = ctx.storage.backend.register('sqlite', backend)162
return async () => {163
dispose()164
await backend.close()165
}166
}, 'storage-sqlite.registerBackend')167
ctx.provide(storageBackendServiceKey('sqlite'), backend)168
}