返回源码地图

packages/storage/storage-domain/src/domain.ts

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

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

1/**
2 * Runtime of one open domain: authoritative in-memory state, the single
3 * per-domain write chain, and change-event emission. Reads are synchronous
4 * from memory; every write queues on the chain, awaits backend durability
5 * FIRST, then mutates memory, then emits `domain/changed` — a rejected
6 * backend write leaves memory untouched (no divergence between reads and the
7 * medium), and events carry values that equal the in-memory state at
8 * emission, in write order.
9 * @module @deepseek-ai/dsh-storage-domain/src/domain
10 */
11
12import type { Context } from '@deepseek-ai/cordis'
13import type { KvUnit } from '@deepseek-ai/dsh-storage'
14import { DomainError } from './error.ts'
15import type { DomainSpec, DomainGlobalSpec, TableKeyOf, TableValueOf } from './spec.ts'
16import type { DomainChanged } from './events.ts'
17
18/** Handle on a domain's global singleton. */
19export interface DomainGlobal<G> {
20 /**
21 * Current value, synchronously from the authoritative in-memory state.
22 * Before the first `set` this is the spec's `initial`.
23 * @returns the current global value.
24 */
25 get(): G
26
27 /**
28 * Replace the value durably. Queued on the domain's write chain; the first
29 * `set` is what materializes the global on the medium.
30 * @param value - New value; must satisfy the spec's schema (not re-checked
31 * here — validation happens at the durable read boundary).
32 * @returns resolution after durability and event emission.
33 */
34 set(value: G): Promise<void>
35}
36
37/**
38 * Handle on one declared table. Records are plain immutable data: returned
39 * values are the stored objects themselves (no defensive copies) and must not
40 * be mutated in place — replace via `put`/`update`.
41 */
42export interface KvTable<K extends string, V> {
43 /**
44 * Read one record, synchronously from memory.
45 * @param key - Record key.
46 * @returns the record, or `undefined` when absent.
47 */
48 get(key: K): V | undefined
49
50 /**
51 * Snapshot iterator over `[key, record]` pairs. A snapshot, not a live
52 * view: iteration stays stable while queued writes land.
53 * @returns the pair iterator.
54 */
55 entries(): IterableIterator<[K, V]>
56
57 /**
58 * Snapshot iterator over keys.
59 * @returns the key iterator.
60 */
61 keys(): IterableIterator<K>
62
63 /** Current record count. */
64 readonly size: number
65
66 /**
67 * Insert or overwrite one record durably.
68 * @param key - Record key.
69 * @param value - The full new record (no partial merge).
70 * @returns resolution after durability and event emission.
71 */
72 put(key: K, value: V): Promise<void>
73
74 /**
75 * Delete one record durably.
76 * @param key - Record key.
77 * @returns `true` when the record existed, `false` when it was already
78 * absent (no write and no event in that case).
79 */
80 delete(key: K): Promise<boolean>
81
82 /**
83 * Atomic read-modify-write on the domain's write chain: `fn` sees the
84 * value current at its queue slot, so concurrent updates never interleave.
85 * @param key - Record key; a missing key rejects with `missing-key`.
86 * @param fn - Synchronous pure transform from current to next record.
87 * @returns the stored next record.
88 */
89 update(key: K, fn: (current: V) => V): Promise<V>
90}
91
92/** Global handle of a spec: typed when declared, `never` (inaccessible) when not. */
93export type DomainGlobalHandleOf<S extends DomainSpec> =
94 S extends { readonly global: DomainGlobalSpec<infer G> } ? DomainGlobal<G> : never
95
96/** One open domain, typed by its spec. */
97export interface Domain<S extends DomainSpec> {
98 /** Domain name from the spec. */
99 readonly name: string
100 /** Global singleton handle; a spec without `global` has no usable handle (`never`). */
101 readonly global: DomainGlobalHandleOf<S>
102 /**
103 * Resolve one declared table handle. Handles are stable — repeated calls
104 * return the same instance.
105 * @param name - Declared table name.
106 * @returns the typed table handle.
107 */
108 table<N extends keyof S['tables'] & string>(name: N): KvTable<TableKeyOf<S, N>, TableValueOf<S, N>>
109
110 /**
111 * Close this domain: reject new writes immediately, drain already-queued
112 * writes (their events still emit), release the backend unit, then free
113 * the domain name for a later open. Idempotent — repeated calls share one
114 * teardown. The consumer owns this call (typically as its own `ctx.effect`
115 * disposer); the facility closes any domain left open when it unmounts.
116 * @returns resolution after the unit is released.
117 */
118 close(): Promise<void>
119}
120
121/** Internal boundary handing table handles their domain-owned write machinery. */
122interface TableHost {
123 readonly domainName: string
124 readonly unit: KvUnit
125 /** Queue one job on the domain's single write chain. */
126 enqueue<T>(job: () => Promise<T>): Promise<T>
127 /** Throw `closed` once the domain has fully closed (reads stay valid while draining). */
128 assertReadable(): void
129 /** Emit `domain/changed` for one durably landed write. */
130 emitChanged(change: DomainChanged): void
131}
132
133const noop = () => {}
134
135/**
136 * The single domain implementation behind the {@link Domain} interface. The
137 * facility constructs it from a validated `loadAll` snapshot and erases it to
138 * `Domain<S>`; nothing outside this package constructs one.
139 */
140export class DomainImpl {
141 /** Domain name from the spec. */
142 readonly name: string
143
144 private readonly tables = new Map<string, KvTableImpl<string, unknown>>()
145 private globalValue: unknown
146 private readonly globalHandle?: DomainGlobal<unknown>
147
148 /** Tail of the write chain; every link settles (rejections are observed by the caller's slice). */
149 private chain: Promise<void> = Promise.resolve()
150 /** Set when close begins: new writes reject while already-queued writes drain. */
151 private disposing = false
152 /** Set when close finishes (chain drained, unit closed): reads reject from here on. */
153 private closed = false
154 private disposal?: Promise<void>
155
156 /**
157 * @param ctx - Context that carries `domain/changed` emissions.
158 * @param spec - The domain declaration.
159 * @param unit - The opened backend unit; this instance owns its lifecycle.
160 * @param records - Validated records from the unit's `loadAll`, one entry
161 * per declared table (empty maps included) — the facility builds it from
162 * the spec, so the entry set IS the table set.
163 * @param globalValue - Validated stored global, or the spec's `initial`
164 * when the medium held none; `undefined` when the spec declares no global.
165 * @param onClosed - Facility hook run once after teardown completes; frees
166 * the domain name for a later open.
167 */
168 constructor(
169 private readonly ctx: Context,
170 spec: DomainSpec,
171 private readonly unit: KvUnit,
172 records: Map<string, Map<string, unknown>>,
173 globalValue: unknown,
174 private readonly onClosed: () => void,
175 ) {
176 this.name = spec.name
177 const host: TableHost = {
178 domainName: spec.name,
179 unit,
180 enqueue: job => this.enqueue(job),
181 assertReadable: () => { this.assertReadable() },
182 emitChanged: (change) => { this.emitChanged(change) },
183 }
184 for (const [table, tableRecords] of records) {
185 this.tables.set(table, new KvTableImpl(host, table, tableRecords))
186 }
187 if (spec.global !== undefined) {
188 this.globalValue = globalValue
189 this.globalHandle = {
190 get: () => {
191 this.assertReadable()
192 return this.globalValue
193 },
194 set: value => this.enqueue(async () => {
195 await this.unit.setGlobal(value)
196 this.globalValue = value
197 this.emitChanged({ domain: this.name, table: '', key: '', operation: 'put', value })
198 }),
199 }
200 }
201 }
202
203 /** Global singleton handle; accessing it on a spec that declares no global is a caller bug and throws. */
204 get global(): DomainGlobal<unknown> {
205 if (this.globalHandle === undefined) {
206 throw new Error(`domain '${this.name}' declares no global`)
207 }
208 return this.globalHandle
209 }
210
211 /**
212 * Resolve one declared table handle; an undeclared name is a caller bug
213 * and throws.
214 * @param name - Declared table name.
215 * @returns the stable table handle.
216 */
217 table(name: string): KvTable<string, unknown> {
218 const table = this.tables.get(name)
219 if (table === undefined) {
220 throw new Error(`domain '${this.name}' declares no table '${name}'`)
221 }
222 return table
223 }
224
225 /**
226 * Close this domain: reject new writes immediately, drain already-queued
227 * writes (their events still emit), close the unit, then free the name via
228 * the facility hook. Idempotent — repeated calls share one teardown.
229 * @returns resolution after the unit is released.
230 */
231 close(): Promise<void> {
232 this.disposal ??= this.runClose()
233 return this.disposal
234 }
235
236 private async runClose(): Promise<void> {
237 this.disposing = true
238 // Chain links never reject (each is settled via then(noop, noop)), so
239 // this await is a pure drain barrier.
240 await this.chain
241 await this.unit.close()
242 this.closed = true
243 this.onClosed()
244 }
245
246 /**
247 * Dispatch one post-durability change notification, containing observer
248 * failures: the write is already committed (medium and memory both hold
249 * the new state), so a throwing listener must not retroactively reject it.
250 */
251 private emitChanged(change: DomainChanged): void {
252 try {
253 this.ctx.emit('domain/changed', change)
254 } catch (error) {
255 // Swallows synchronous observer exceptions only: emit dispatches
256 // listeners inline and nothing else runs in the try. The event is a
257 // notification, not a transaction participant — the commit point has
258 // passed, so containment (with a log) is the only correct outcome.
259 this.ctx.logger.warn(`domain '${this.name}': domain/changed listener failed: ${String(error)}`)
260 }
261 }
262
263 private enqueue<T>(job: () => Promise<T>): Promise<T> {
264 if (this.disposing) {
265 return Promise.reject(new DomainError('closed', `domain '${this.name}' is closed`))
266 }
267 const result = this.chain.then(job)
268 this.chain = result.then(noop, noop)
269 return result
270 }
271
272 private assertReadable(): void {
273 if (this.closed) {
274 throw new DomainError('closed', `domain '${this.name}' is closed`)
275 }
276 }
277}
278
279/** Table handle bound to one in-memory record map and its domain's write chain. */
280class KvTableImpl<K extends string, V> implements KvTable<K, V> {
281 constructor(
282 private readonly host: TableHost,
283 private readonly tableName: string,
284 private readonly records: Map<string, unknown>,
285 ) {}
286
287 get(key: K): V | undefined {
288 this.host.assertReadable()
289 return this.records.get(key) as V | undefined
290 }
291
292 entries(): IterableIterator<[K, V]> {
293 this.host.assertReadable()
294 return ([...this.records.entries()] as [K, V][])[Symbol.iterator]()
295 }
296
297 keys(): IterableIterator<K> {
298 this.host.assertReadable()
299 return ([...this.records.keys()] as K[])[Symbol.iterator]()
300 }
301
302 get size(): number {
303 this.host.assertReadable()
304 return this.records.size
305 }
306
307 put(key: K, value: V): Promise<void> {
308 return this.host.enqueue(async () => {
309 await this.host.unit.putRecord(this.tableName, key, value)
310 this.records.set(key, value)
311 this.emitPut(key, value)
312 })
313 }
314
315 delete(key: K): Promise<boolean> {
316 return this.host.enqueue(async () => {
317 // Existence is decided at this job's chain slot, not at call time: an
318 // earlier queued put of the same key makes this delete observe it.
319 if (!this.records.has(key)) return false
320 await this.host.unit.deleteRecord(this.tableName, key)
321 this.records.delete(key)
322 this.host.emitChanged({
323 domain: this.host.domainName,
324 table: this.tableName,
325 key,
326 operation: 'deleted',
327 })
328 return true
329 })
330 }
331
332 update(key: K, fn: (current: V) => V): Promise<V> {
333 return this.host.enqueue(async () => {
334 if (!this.records.has(key)) {
335 throw new DomainError(
336 'missing-key',
337 `domain '${this.host.domainName}' table '${this.tableName}' has no record '${key}' to update`,
338 )
339 }
340 const next = fn(this.records.get(key) as V)
341 await this.host.unit.putRecord(this.tableName, key, next)
342 this.records.set(key, next)
343 this.emitPut(key, next)
344 return next
345 })
346 }
347
348 private emitPut(key: K, value: V): void {
349 this.host.emitChanged({
350 domain: this.host.domainName,
351 table: this.tableName,
352 key,
353 operation: 'put',
354 value,
355 })
356 }
357}