返回源码地图

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

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

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

1/**
2 * SQLite storage backend for the storage hub: one database file hosts every
3 * routed unit, document-per-row (`key TEXT` / `value TEXT` JSON). Registers
4 * as backend `sqlite`; the disposer unregisters first, then closes the medium.
5 * @module @deepseek-ai/dsh-storage-sqlite
6 */
7
8import type { Context } from '@deepseek-ai/cordis'
9import z from '@deepseek-ai/schemastery'
10import type { DatabaseSync } from 'node:sqlite'
11import { StorageError, UNIT_NAME_RE, storageBackendServiceKey } from '@deepseek-ai/dsh-storage'
12import type { KvFacet, KvUnit, KvUnitDescriptor, StorageBackend } from '@deepseek-ai/dsh-storage'
13import { openDatabase, recordTableName, type JournalMode } from './schema.ts'
14import { SqliteKvUnit } from './unit.ts'
15
16export { STORAGE_SQLITE_SCHEMA_VERSION, type JournalMode } from './schema.ts'
17
18/** Cordis plugin name. */
19export const name = 'storage-sqlite'
20/** The backend registers on the storage hub. */
21export const inject = ['storage']
22
23/** Plugin configuration. */
24export 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 path
29 * modes are preserved. Filesystem setup errors other than an existing
30 * database fail the open. The backend does not protect confidentiality or
31 * integrity when another principal can replace the database entry in its
32 * parent directory.
33 */
34 path: string
35 /**
36 * SQLite `journal_mode` pragma. `wal` (the default) suits local disks; pick
37 * a rollback-journal mode (`delete`/`truncate`/`persist`) on filesystems
38 * where WAL's shared-memory files do not work (network mounts). See
39 * {@link JournalMode}.
40 */
41 journalMode?: JournalMode
42}
43
44/** Schemastery validator for {@link Config}. */
45export const Config: z<Config> = z.object({
46 path: z.string().required(),
47 journalMode: z.union(['wal', 'delete', 'truncate', 'persist'] as const).default('wal'),
48})
49
50/**
51 * The SQLite {@link StorageBackend}. Owns one `DatabaseSync` connection and
52 * the open-unit table; `kv.open` validates names, enforces the per-unit
53 * version stamp in `units`, and ensures the unit's record tables.
54 */
55export class SqliteStorageBackend implements StorageBackend {
56 /** The key-value facet; the only shape this backend serves. */
57 readonly kv: KvFacet = { open: descriptor => this.openUnit(descriptor) }
58
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> | undefined
63
64 /**
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 an
70 // open failure still surfaces to each caller; this guard only prevents an
71 // unhandled-rejection crash when the failure precedes the first use.
72 this.ready.catch(() => {})
73 }
74
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 same
91 // 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 pending
96 }
97
98 private async materializeUnit(descriptor: KvUnitDescriptor): Promise<SqliteKvUnit> {
99 const db = await this.ready
100 const row = db.prepare('SELECT version FROM units WHERE name = ?').get(descriptor.name) as
101 | { version: number }
102 | undefined
103 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 NULL
117 ) STRICT
118 `)
119 }
120 return new SqliteKvUnit(db, descriptor, () => {
121 this.units.delete(descriptor.name)
122 })
123 }
124
125 /**
126 * Close every open unit and release the database. Idempotent; concurrent
127 * 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.closing
133 }
134
135 private async doClose(): Promise<void> {
136 let db: DatabaseSync
137 try {
138 db = await this.ready
139 } catch {
140 // The medium never opened; that failure already rejected the opener and
141 // every unit call, so there is nothing left to release here.
142 return
143 }
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}
151
152/**
153 * Register the SQLite backend as `sqlite` on the storage hub. The disposer
154 * unregisters the name first, then closes the backend.
155 * @param ctx - Plugin context (must inject `storage`).
156 * @param config - Validated plugin configuration.
157 */
158export 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}