1
/**2
* Domain data form (`ctx.storage.domain`): schema-validated, change-emitting3
* KV domains over storage backends. The single implementation of the domain4
* layer — consumers depend on this package and never touch backends directly.5
* Plugin `Config` is schemastery; record schemas inside domain specs are zod6
* (see `src/spec.ts` for the split rationale).7
* @module @deepseek-ai/dsh-storage-domain8
*/10
import type { Context } from '@deepseek-ai/cordis'11
import z from '@deepseek-ai/schemastery'12
import { storageBackendServiceKey } from '@deepseek-ai/dsh-storage'13
import { DomainError } from './error.ts'14
import { descriptorOf } from './spec.ts'15
import type { DomainSpec } from './spec.ts'16
import { DomainImpl } from './domain.ts'17
import type { Domain } from './domain.ts'19
export { DomainError } from './error.ts'20
export type { DomainErrorCode, DomainErrorOptions, InvalidRecordDetail } from './error.ts'21
export { defineDomain, domainTable, descriptorOf } from './spec.ts'22
export type {23
DomainSpec, DomainGlobalSpec, DomainTableSpec,24
TableKeyOf, TableValueOf, GlobalValueOf,25
} from './spec.ts'26
export type { DomainChanged } from './events.ts'27
export type { Domain, DomainGlobal, DomainGlobalHandleOf, KvTable } from './domain.ts'29
declare module '@deepseek-ai/dsh-storage' {30
interface StorageForms {31
domain: DomainFacility32
}33
}35
declare module '@deepseek-ai/cordis' {36
interface Context {37
storageDomain: DomainFacility38
}39
}41
/** Cordis plugin name. */42
export const name = 'storage-domain'43
/** The storage hub must be present before the form can mount. */44
export const inject = ['storage']46
/**47
* Plugin config. Which backend serves which domain is decided here, not48
* globally on the hub: `backend` is the default route and `routes` overrides49
* it per domain name. A route naming an unregistered backend fails loud at50
* `open` with `backend-not-found`.51
*/52
export interface Config {53
/** Default backend name for every domain without an explicit route. Required: there is no universally correct medium. */54
backend: string55
/** Per-domain overrides: domain name → backend name. */56
routes?: Record<string, string>57
}59
export const Config: z<Config> = z.object({60
backend: z.string().required(),61
routes: z.dict(z.string()).default({}),62
})64
/**65
* The mounted domain facility. Opens declared domains over routed backends;66
* one facility instance owns the open-domain table and enforces single-open67
* per domain name.68
*/69
export 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>()74
/**75
* @param ctx - Context of the domain plugin; open-domain effects and change76
* events attach here.77
* @param config - Validated plugin config.78
*/79
constructor(80
private readonly ctx: Context,81
private readonly config: Config,82
) {}84
/**85
* Open one declared domain. Steps, each failing the whole call: reject a86
* name that is already open (`already-open`); resolve the backend route87
* (`backend-not-found` passes through from the hub); require its `kv` facet88
* (`facet-unsupported`); open the unit projected from the spec (backend89
* `version-mismatch`/`malformed-medium` pass through); load and validate90
* every stored record against the spec's zod schemas (`invalid-record`91
* with the offending table and key — unless the spec declares92
* `invalidRecords: 'backup-and-skip'` and the unit can move documents aside, in93
* 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 via97
* `Domain.close()` (typically as its own `ctx.effect` disposer) — the98
* facility does not tie the domain to any consumer fiber. Domains still99
* 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.backend110
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: unknown125
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's129
// document aside, log the concrete failure, and open without the130
// record. Backends that cannot move a document keep the loud path.131
if (spec.invalidRecords !== 'backup-and-skip' || unit.backupRecord === undefined) throw error132
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
continue139
}140
records.set(key, parsed)141
}142
tables.set(table, records)143
}144
// A null stored global means "never written": serve `initial` without145
// materializing it — the first `set` writes.146
const globalSpec = spec.global147
const globalValue = globalSpec === undefined148
? undefined149
: snapshot.global === null150
? globalSpec.initial151
: parseRecord(spec.name, '', '', () => globalSpec.schema.parse(snapshot.global))152
// The onClosed hook runs strictly after teardown completes: writes153
// landing during the drain still emit domain/changed, and the name154
// 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 because162
// 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 error167
}168
} catch (error) {169
// Any failure means the domain never registered (nothing can throw170
// after it), so releasing the name reservation is unconditional.171
this.reserved.delete(spec.name)172
throw error173
}174
}176
/**177
* Look up an open domain by name, untyped. Diagnostic surface; typed178
* 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
}186
/**187
* Close every domain still open on this facility. The unmount path for188
* consumers that never called `Domain.close()` themselves; closing is189
* 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
}197
/** Run one zod parse, translating failure to `invalid-record` with its location. */198
function 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
}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
*/217
export function apply(ctx: Context, config: Config): Promise<void> {218
const backendServices = [...new Set([219
config.backend,220
...Object.values(config.routes ?? {}),221
])].map(storageBackendServiceKey)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 emit229
// domain/changed.230
await facility.closeAll()231
unmount()232
}233
})234
domainCtx.provide('storageDomain', facility)235
})236
return Promise.resolve(fiber).then(() => {})237
}