返回源码地图

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

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

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

1/**
2 * Domain data form (`ctx.storage.domain`): schema-validated, change-emitting
3 * KV domains over storage backends. The single implementation of the domain
4 * layer — consumers depend on this package and never touch backends directly.
5 * Plugin `Config` is schemastery; record schemas inside domain specs are zod
6 * (see `src/spec.ts` for the split rationale).
7 * @module @deepseek-ai/dsh-storage-domain
8 */
9
10import type { Context } from '@deepseek-ai/cordis'
11import z from '@deepseek-ai/schemastery'
12import { storageBackendServiceKey } from '@deepseek-ai/dsh-storage'
13import { DomainError } from './error.ts'
14import { descriptorOf } from './spec.ts'
15import type { DomainSpec } from './spec.ts'
16import { DomainImpl } from './domain.ts'
17import type { Domain } from './domain.ts'
18
19export { DomainError } from './error.ts'
20export type { DomainErrorCode, DomainErrorOptions, InvalidRecordDetail } from './error.ts'
21export { defineDomain, domainTable, descriptorOf } from './spec.ts'
22export type {
23 DomainSpec, DomainGlobalSpec, DomainTableSpec,
24 TableKeyOf, TableValueOf, GlobalValueOf,
25} from './spec.ts'
26export type { DomainChanged } from './events.ts'
27export type { Domain, DomainGlobal, DomainGlobalHandleOf, KvTable } from './domain.ts'
28
29declare module '@deepseek-ai/dsh-storage' {
30 interface StorageForms {
31 domain: DomainFacility
32 }
33}
34
35declare module '@deepseek-ai/cordis' {
36 interface Context {
37 storageDomain: DomainFacility
38 }
39}
40
41/** Cordis plugin name. */
42export const name = 'storage-domain'
43/** The storage hub must be present before the form can mount. */
44export const inject = ['storage']
45
46/**
47 * Plugin config. Which backend serves which domain is decided here, not
48 * globally on the hub: `backend` is the default route and `routes` overrides
49 * it per domain name. A route naming an unregistered backend fails loud at
50 * `open` with `backend-not-found`.
51 */
52export interface Config {
53 /** Default backend name for every domain without an explicit route. Required: there is no universally correct medium. */
54 backend: string
55 /** Per-domain overrides: domain name → backend name. */
56 routes?: Record<string, string>
57}
58
59export const Config: z<Config> = z.object({
60 backend: z.string().required(),
61 routes: z.dict(z.string()).default({}),
62})
63
64/**
65 * The mounted domain facility. Opens declared domains over routed backends;
66 * one facility instance owns the open-domain table and enforces single-open
67 * per domain name.
68 */
69export class DomainFacility {
70 private readonly domains = new Map<string, DomainImpl>()
71 /** Names reserved by an in-flight or completed open, so concurrent opens of one name fail loud. */
72 private readonly reserved = new Set<string>()
73
74 /**
75 * @param ctx - Context of the domain plugin; open-domain effects and change
76 * events attach here.
77 * @param config - Validated plugin config.
78 */
79 constructor(
80 private readonly ctx: Context,
81 private readonly config: Config,
82 ) {}
83
84 /**
85 * Open one declared domain. Steps, each failing the whole call: reject a
86 * name that is already open (`already-open`); resolve the backend route
87 * (`backend-not-found` passes through from the hub); require its `kv` facet
88 * (`facet-unsupported`); open the unit projected from the spec (backend
89 * `version-mismatch`/`malformed-medium` pass through); load and validate
90 * every stored record against the spec's zod schemas (`invalid-record`
91 * with the offending table and key — unless the spec declares
92 * `invalidRecords: 'backup-and-skip'` and the unit can move documents aside, in
93 * which case the failing record is backed up, logged, and skipped);
94 * construct the domain.
95 *
96 * Lifecycle: the CALLER owns the returned handle and closes it via
97 * `Domain.close()` (typically as its own `ctx.effect` disposer) — the
98 * facility does not tie the domain to any consumer fiber. Domains still
99 * open when the facility unmounts are closed by the plugin disposer.
100 * @param spec - The domain declaration, typically from `defineDomain`.
101 * @returns the opened domain handle, typed by the spec.
102 */
103 async open<S extends DomainSpec>(spec: S): Promise<Domain<S>> {
104 if (this.reserved.has(spec.name)) {
105 throw new DomainError('already-open', `domain '${spec.name}' is already open`)
106 }
107 this.reserved.add(spec.name)
108 try {
109 const backendName = this.config.routes?.[spec.name] ?? this.config.backend
110 const backend = this.ctx.storage.backend.get(backendName)
111 if (!backend.kv) {
112 throw new DomainError(
113 'facet-unsupported',
114 `backend '${backendName}' routed for domain '${spec.name}' has no kv facet`,
115 )
116 }
117 const unit = await backend.kv.open(descriptorOf(spec))
118 try {
119 const snapshot = await unit.loadAll()
120 const tables = new Map<string, Map<string, unknown>>()
121 for (const [table, tableSpec] of Object.entries(spec.tables)) {
122 const records = new Map<string, unknown>()
123 for (const [key, raw] of Object.entries(snapshot.tables[table] ?? {})) {
124 let parsed: unknown
125 try {
126 parsed = parseRecord(spec.name, table, key, () => tableSpec.valueSchema.parse(raw))
127 } catch (error) {
128 // Backup-and-skip policy (disposable derived data): move the record's
129 // document aside, log the concrete failure, and open without the
130 // record. Backends that cannot move a document keep the loud path.
131 if (spec.invalidRecords !== 'backup-and-skip' || unit.backupRecord === undefined) throw error
132 const moved = await unit.backupRecord(table, key)
133 // parseRecord always wraps the zod failure as the cause.
134 this.ctx.logger.error(
135 `domain '${spec.name}': stored record '${key}' in table '${table}' failed schema validation; `
136 + `moved to '${moved}' and treated as absent. Cause: ${String((error as DomainError).cause)}`,
137 )
138 continue
139 }
140 records.set(key, parsed)
141 }
142 tables.set(table, records)
143 }
144 // A null stored global means "never written": serve `initial` without
145 // materializing it — the first `set` writes.
146 const globalSpec = spec.global
147 const globalValue = globalSpec === undefined
148 ? undefined
149 : snapshot.global === null
150 ? globalSpec.initial
151 : parseRecord(spec.name, '', '', () => globalSpec.schema.parse(snapshot.global))
152 // The onClosed hook runs strictly after teardown completes: writes
153 // landing during the drain still emit domain/changed, and the name
154 // frees up for reopening only once the domain is fully closed.
155 const domain: DomainImpl = new DomainImpl(this.ctx, spec, unit, tables, globalValue, () => {
156 this.domains.delete(spec.name)
157 this.reserved.delete(spec.name)
158 })
159 this.domains.set(spec.name, domain)
160 // The single type-erasure point: DomainImpl is the untyped runtime,
161 // Domain<S> the spec-typed view; the unknown hop is required because
162 // S's conditional global-handle type stays unresolved here.
163 return domain as unknown as Domain<S>
164 } catch (error) {
165 await unit.close()
166 throw error
167 }
168 } catch (error) {
169 // Any failure means the domain never registered (nothing can throw
170 // after it), so releasing the name reservation is unconditional.
171 this.reserved.delete(spec.name)
172 throw error
173 }
174 }
175
176 /**
177 * Look up an open domain by name, untyped. Diagnostic surface; typed
178 * consumers hold the handle returned by {@link open}.
179 * @param name - Domain name.
180 * @returns the open domain runtime, or `undefined` when not open.
181 */
182 get(name: string): DomainImpl | undefined {
183 return this.domains.get(name)
184 }
185
186 /**
187 * Close every domain still open on this facility. The unmount path for
188 * consumers that never called `Domain.close()` themselves; closing is
189 * idempotent, so double-closing an already-closed domain is harmless.
190 * @returns resolution after every unit is released.
191 */
192 async closeAll(): Promise<void> {
193 await Promise.all([...this.domains.values()].map(domain => domain.close()))
194 }
195}
196
197/** Run one zod parse, translating failure to `invalid-record` with its location. */
198function parseRecord<T>(domain: string, table: string, key: string, parse: () => T): T {
199 try {
200 return parse()
201 } catch (error) {
202 const slot = table === '' ? 'global' : `record '${key}' in table '${table}'`
203 throw new DomainError(
204 'invalid-record',
205 `domain '${domain}': stored ${slot} does not match its schema`,
206 { detail: { table, key }, cause: error },
207 )
208 }
209}
210
211/**
212 * Mount the domain data form on the storage hub.
213 * @param ctx - Plugin context.
214 * @param config - Validated plugin config.
215 * @returns resolution after an already-available backend set activates the form.
216 */
217export function apply(ctx: Context, config: Config): Promise<void> {
218 const backendServices = [...new Set([
219 config.backend,
220 ...Object.values(config.routes ?? {}),
221 ])].map(storageBackendServiceKey)
222
223 const fiber = ctx.inject(backendServices, (domainCtx) => {
224 const facility = new DomainFacility(domainCtx, config)
225 domainCtx.effect(() => {
226 const unmount = domainCtx.storage.mount('domain', facility)
227 return async () => {
228 // Close leftovers before unmounting: draining writes still emit
229 // domain/changed.
230 await facility.closeAll()
231 unmount()
232 }
233 })
234 domainCtx.provide('storageDomain', facility)
235 })
236 return Promise.resolve(fiber).then(() => {})
237}