refactor(storage): rename dsh-domain to dsh-storage-domain
The bare 'domain' name was too generic for a published package. The directory moves to packages/storage/storage-domain, the package becomes @deepseek-ai/dsh-storage-domain, and the plugin/invariant names follow; the ctx surface (ctx.storage.domain), the domain/changed event, and all runtime behavior are unchanged. References, catalogs, graphs, and the bilingual design note move together.
This commit is contained in:
357
packages/storage/storage-domain/src/domain.ts
Normal file
357
packages/storage/storage-domain/src/domain.ts
Normal file
@@ -0,0 +1,357 @@
|
||||
/**
|
||||
* Runtime of one open domain: authoritative in-memory state, the single
|
||||
* per-domain write chain, and change-event emission. Reads are synchronous
|
||||
* from memory; every write queues on the chain, awaits backend durability
|
||||
* FIRST, then mutates memory, then emits `domain/changed` — a rejected
|
||||
* backend write leaves memory untouched (no divergence between reads and the
|
||||
* medium), and events carry values that equal the in-memory state at
|
||||
* emission, in write order.
|
||||
* @module @deepseek-ai/dsh-storage-domain/src/domain
|
||||
*/
|
||||
|
||||
import type { Context } from 'cordis'
|
||||
import type { KvUnit } from '@deepseek-ai/dsh-storage'
|
||||
import { DomainError } from './error.ts'
|
||||
import type { DomainSpec, DomainGlobalSpec, TableKeyOf, TableValueOf } from './spec.ts'
|
||||
import type { DomainChanged } from './events.ts'
|
||||
|
||||
/** Handle on a domain's global singleton. */
|
||||
export interface DomainGlobal<G> {
|
||||
/**
|
||||
* Current value, synchronously from the authoritative in-memory state.
|
||||
* Before the first `set` this is the spec's `initial`.
|
||||
* @returns the current global value.
|
||||
*/
|
||||
get(): G
|
||||
|
||||
/**
|
||||
* Replace the value durably. Queued on the domain's write chain; the first
|
||||
* `set` is what materializes the global on the medium.
|
||||
* @param value - New value; must satisfy the spec's schema (not re-checked
|
||||
* here — validation happens at the durable read boundary).
|
||||
* @returns resolution after durability and event emission.
|
||||
*/
|
||||
set(value: G): Promise<void>
|
||||
}
|
||||
|
||||
/**
|
||||
* Handle on one declared table. Records are plain immutable data: returned
|
||||
* values are the stored objects themselves (no defensive copies) and must not
|
||||
* be mutated in place — replace via `put`/`update`.
|
||||
*/
|
||||
export interface KvTable<K extends string, V> {
|
||||
/**
|
||||
* Read one record, synchronously from memory.
|
||||
* @param key - Record key.
|
||||
* @returns the record, or `undefined` when absent.
|
||||
*/
|
||||
get(key: K): V | undefined
|
||||
|
||||
/**
|
||||
* Snapshot iterator over `[key, record]` pairs. A snapshot, not a live
|
||||
* view: iteration stays stable while queued writes land.
|
||||
* @returns the pair iterator.
|
||||
*/
|
||||
entries(): IterableIterator<[K, V]>
|
||||
|
||||
/**
|
||||
* Snapshot iterator over keys.
|
||||
* @returns the key iterator.
|
||||
*/
|
||||
keys(): IterableIterator<K>
|
||||
|
||||
/** Current record count. */
|
||||
readonly size: number
|
||||
|
||||
/**
|
||||
* Insert or overwrite one record durably.
|
||||
* @param key - Record key.
|
||||
* @param value - The full new record (no partial merge).
|
||||
* @returns resolution after durability and event emission.
|
||||
*/
|
||||
put(key: K, value: V): Promise<void>
|
||||
|
||||
/**
|
||||
* Delete one record durably.
|
||||
* @param key - Record key.
|
||||
* @returns `true` when the record existed, `false` when it was already
|
||||
* absent (no write and no event in that case).
|
||||
*/
|
||||
delete(key: K): Promise<boolean>
|
||||
|
||||
/**
|
||||
* Atomic read-modify-write on the domain's write chain: `fn` sees the
|
||||
* value current at its queue slot, so concurrent updates never interleave.
|
||||
* @param key - Record key; a missing key rejects with `missing-key`.
|
||||
* @param fn - Synchronous pure transform from current to next record.
|
||||
* @returns the stored next record.
|
||||
*/
|
||||
update(key: K, fn: (current: V) => V): Promise<V>
|
||||
}
|
||||
|
||||
/** Global handle of a spec: typed when declared, `never` (inaccessible) when not. */
|
||||
export type DomainGlobalHandleOf<S extends DomainSpec> =
|
||||
S extends { readonly global: DomainGlobalSpec<infer G> } ? DomainGlobal<G> : never
|
||||
|
||||
/** One open domain, typed by its spec. */
|
||||
export interface Domain<S extends DomainSpec> {
|
||||
/** Domain name from the spec. */
|
||||
readonly name: string
|
||||
/** Global singleton handle; a spec without `global` has no usable handle (`never`). */
|
||||
readonly global: DomainGlobalHandleOf<S>
|
||||
/**
|
||||
* Resolve one declared table handle. Handles are stable — repeated calls
|
||||
* return the same instance.
|
||||
* @param name - Declared table name.
|
||||
* @returns the typed table handle.
|
||||
*/
|
||||
table<N extends keyof S['tables'] & string>(name: N): KvTable<TableKeyOf<S, N>, TableValueOf<S, N>>
|
||||
|
||||
/**
|
||||
* Close this domain: reject new writes immediately, drain already-queued
|
||||
* writes (their events still emit), release the backend unit, then free
|
||||
* the domain name for a later open. Idempotent — repeated calls share one
|
||||
* teardown. The consumer owns this call (typically as its own `ctx.effect`
|
||||
* disposer); the facility closes any domain left open when it unmounts.
|
||||
* @returns resolution after the unit is released.
|
||||
*/
|
||||
close(): Promise<void>
|
||||
}
|
||||
|
||||
/** Internal seam handing table handles their domain-owned write machinery. */
|
||||
interface TableHost {
|
||||
readonly domainName: string
|
||||
readonly unit: KvUnit
|
||||
/** Queue one job on the domain's single write chain. */
|
||||
enqueue<T>(job: () => Promise<T>): Promise<T>
|
||||
/** Throw `closed` once the domain has fully closed (reads stay valid while draining). */
|
||||
assertReadable(): void
|
||||
/** Emit `domain/changed` for one durably landed write. */
|
||||
emitChanged(change: DomainChanged): void
|
||||
}
|
||||
|
||||
const noop = () => {}
|
||||
|
||||
/**
|
||||
* The single domain implementation behind the {@link Domain} interface. The
|
||||
* facility constructs it from a validated `loadAll` snapshot and erases it to
|
||||
* `Domain<S>`; nothing outside this package constructs one.
|
||||
*/
|
||||
export class DomainImpl {
|
||||
/** Domain name from the spec. */
|
||||
readonly name: string
|
||||
|
||||
private readonly tables = new Map<string, KvTableImpl<string, unknown>>()
|
||||
private globalValue: unknown
|
||||
private readonly globalHandle?: DomainGlobal<unknown>
|
||||
|
||||
/** Tail of the write chain; every link settles (rejections are observed by the caller's slice). */
|
||||
private chain: Promise<void> = Promise.resolve()
|
||||
/** Set when close begins: new writes reject while already-queued writes drain. */
|
||||
private disposing = false
|
||||
/** Set when close finishes (chain drained, unit closed): reads reject from here on. */
|
||||
private closed = false
|
||||
private disposal?: Promise<void>
|
||||
|
||||
/**
|
||||
* @param ctx - Context that carries `domain/changed` emissions.
|
||||
* @param spec - The domain declaration.
|
||||
* @param unit - The opened backend unit; this instance owns its lifecycle.
|
||||
* @param records - Validated records from the unit's `loadAll`, one entry
|
||||
* per declared table (empty maps included) — the facility builds it from
|
||||
* the spec, so the entry set IS the table set.
|
||||
* @param globalValue - Validated stored global, or the spec's `initial`
|
||||
* when the medium held none; `undefined` when the spec declares no global.
|
||||
* @param onClosed - Facility hook run once after teardown completes; frees
|
||||
* the domain name for a later open.
|
||||
*/
|
||||
constructor(
|
||||
private readonly ctx: Context,
|
||||
spec: DomainSpec,
|
||||
private readonly unit: KvUnit,
|
||||
records: Map<string, Map<string, unknown>>,
|
||||
globalValue: unknown,
|
||||
private readonly onClosed: () => void,
|
||||
) {
|
||||
this.name = spec.name
|
||||
const host: TableHost = {
|
||||
domainName: spec.name,
|
||||
unit,
|
||||
enqueue: job => this.enqueue(job),
|
||||
assertReadable: () => { this.assertReadable() },
|
||||
emitChanged: (change) => { this.emitChanged(change) },
|
||||
}
|
||||
for (const [table, tableRecords] of records) {
|
||||
this.tables.set(table, new KvTableImpl(host, table, tableRecords))
|
||||
}
|
||||
if (spec.global !== undefined) {
|
||||
this.globalValue = globalValue
|
||||
this.globalHandle = {
|
||||
get: () => {
|
||||
this.assertReadable()
|
||||
return this.globalValue
|
||||
},
|
||||
set: value => this.enqueue(async () => {
|
||||
await this.unit.setGlobal(value)
|
||||
this.globalValue = value
|
||||
this.emitChanged({ domain: this.name, table: '', key: '', operation: 'put', value })
|
||||
}),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Global singleton handle; accessing it on a spec that declares no global is a caller bug and throws. */
|
||||
get global(): DomainGlobal<unknown> {
|
||||
if (this.globalHandle === undefined) {
|
||||
throw new Error(`domain '${this.name}' declares no global`)
|
||||
}
|
||||
return this.globalHandle
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve one declared table handle; an undeclared name is a caller bug
|
||||
* and throws.
|
||||
* @param name - Declared table name.
|
||||
* @returns the stable table handle.
|
||||
*/
|
||||
table(name: string): KvTable<string, unknown> {
|
||||
const table = this.tables.get(name)
|
||||
if (table === undefined) {
|
||||
throw new Error(`domain '${this.name}' declares no table '${name}'`)
|
||||
}
|
||||
return table
|
||||
}
|
||||
|
||||
/**
|
||||
* Close this domain: reject new writes immediately, drain already-queued
|
||||
* writes (their events still emit), close the unit, then free the name via
|
||||
* the facility hook. Idempotent — repeated calls share one teardown.
|
||||
* @returns resolution after the unit is released.
|
||||
*/
|
||||
close(): Promise<void> {
|
||||
this.disposal ??= this.runClose()
|
||||
return this.disposal
|
||||
}
|
||||
|
||||
private async runClose(): Promise<void> {
|
||||
this.disposing = true
|
||||
// Chain links never reject (each is settled via then(noop, noop)), so
|
||||
// this await is a pure drain barrier.
|
||||
await this.chain
|
||||
await this.unit.close()
|
||||
this.closed = true
|
||||
this.onClosed()
|
||||
}
|
||||
|
||||
/**
|
||||
* Dispatch one post-durability change notification, containing observer
|
||||
* failures: the write is already committed (medium and memory both hold
|
||||
* the new state), so a throwing listener must not retroactively reject it.
|
||||
*/
|
||||
private emitChanged(change: DomainChanged): void {
|
||||
try {
|
||||
this.ctx.emit('domain/changed', change)
|
||||
} catch (error) {
|
||||
// Swallows synchronous observer exceptions only: emit dispatches
|
||||
// listeners inline and nothing else runs in the try. The event is a
|
||||
// notification, not a transaction participant — the commit point has
|
||||
// passed, so containment (with a log) is the only correct outcome.
|
||||
this.ctx.logger.warn(`domain '${this.name}': domain/changed listener failed: ${String(error)}`)
|
||||
}
|
||||
}
|
||||
|
||||
private enqueue<T>(job: () => Promise<T>): Promise<T> {
|
||||
if (this.disposing) {
|
||||
return Promise.reject(new DomainError('closed', `domain '${this.name}' is closed`))
|
||||
}
|
||||
const result = this.chain.then(job)
|
||||
this.chain = result.then(noop, noop)
|
||||
return result
|
||||
}
|
||||
|
||||
private assertReadable(): void {
|
||||
if (this.closed) {
|
||||
throw new DomainError('closed', `domain '${this.name}' is closed`)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Table handle bound to one in-memory record map and its domain's write chain. */
|
||||
class KvTableImpl<K extends string, V> implements KvTable<K, V> {
|
||||
constructor(
|
||||
private readonly host: TableHost,
|
||||
private readonly tableName: string,
|
||||
private readonly records: Map<string, unknown>,
|
||||
) {}
|
||||
|
||||
get(key: K): V | undefined {
|
||||
this.host.assertReadable()
|
||||
return this.records.get(key) as V | undefined
|
||||
}
|
||||
|
||||
entries(): IterableIterator<[K, V]> {
|
||||
this.host.assertReadable()
|
||||
return ([...this.records.entries()] as [K, V][])[Symbol.iterator]()
|
||||
}
|
||||
|
||||
keys(): IterableIterator<K> {
|
||||
this.host.assertReadable()
|
||||
return ([...this.records.keys()] as K[])[Symbol.iterator]()
|
||||
}
|
||||
|
||||
get size(): number {
|
||||
this.host.assertReadable()
|
||||
return this.records.size
|
||||
}
|
||||
|
||||
put(key: K, value: V): Promise<void> {
|
||||
return this.host.enqueue(async () => {
|
||||
await this.host.unit.putRecord(this.tableName, key, value)
|
||||
this.records.set(key, value)
|
||||
this.emitPut(key, value)
|
||||
})
|
||||
}
|
||||
|
||||
delete(key: K): Promise<boolean> {
|
||||
return this.host.enqueue(async () => {
|
||||
// Existence is decided at this job's chain slot, not at call time: an
|
||||
// earlier queued put of the same key makes this delete observe it.
|
||||
if (!this.records.has(key)) return false
|
||||
await this.host.unit.deleteRecord(this.tableName, key)
|
||||
this.records.delete(key)
|
||||
this.host.emitChanged({
|
||||
domain: this.host.domainName,
|
||||
table: this.tableName,
|
||||
key,
|
||||
operation: 'deleted',
|
||||
})
|
||||
return true
|
||||
})
|
||||
}
|
||||
|
||||
update(key: K, fn: (current: V) => V): Promise<V> {
|
||||
return this.host.enqueue(async () => {
|
||||
if (!this.records.has(key)) {
|
||||
throw new DomainError(
|
||||
'missing-key',
|
||||
`domain '${this.host.domainName}' table '${this.tableName}' has no record '${key}' to update`,
|
||||
)
|
||||
}
|
||||
const next = fn(this.records.get(key) as V)
|
||||
await this.host.unit.putRecord(this.tableName, key, next)
|
||||
this.records.set(key, next)
|
||||
this.emitPut(key, next)
|
||||
return next
|
||||
})
|
||||
}
|
||||
|
||||
private emitPut(key: K, value: V): void {
|
||||
this.host.emitChanged({
|
||||
domain: this.host.domainName,
|
||||
table: this.tableName,
|
||||
key,
|
||||
operation: 'put',
|
||||
value,
|
||||
})
|
||||
}
|
||||
}
|
||||
53
packages/storage/storage-domain/src/error.ts
Normal file
53
packages/storage/storage-domain/src/error.ts
Normal file
@@ -0,0 +1,53 @@
|
||||
/**
|
||||
* Error vocabulary of the domain data form.
|
||||
* @module @deepseek-ai/dsh-storage-domain/src/error
|
||||
*/
|
||||
|
||||
/** Discriminant codes carried by every {@link DomainError}. */
|
||||
export type DomainErrorCode =
|
||||
| 'already-open'
|
||||
| 'facet-unsupported'
|
||||
| 'invalid-record'
|
||||
| 'missing-key'
|
||||
| 'closed'
|
||||
|
||||
/** Location of the record that failed schema validation at the durable boundary. */
|
||||
export interface InvalidRecordDetail {
|
||||
/** Table holding the rejected record; `''` for the global singleton. */
|
||||
readonly table: string
|
||||
/** Key of the rejected record; `''` for the global singleton. */
|
||||
readonly key: string
|
||||
}
|
||||
|
||||
/** Construction options: standard `cause` plus the `invalid-record` location. */
|
||||
export interface DomainErrorOptions extends ErrorOptions {
|
||||
/** Present exactly when `code` is `invalid-record`. */
|
||||
readonly detail?: InvalidRecordDetail
|
||||
}
|
||||
|
||||
/**
|
||||
* Error thrown by the domain layer. The `code` is the stable contract
|
||||
* consumers may switch on; `message` is diagnostic prose. Backend failures
|
||||
* (`backend-not-found`, `version-mismatch`, …) pass through as
|
||||
* `StorageError` — the domain layer does not rewrap them.
|
||||
*/
|
||||
export class DomainError extends Error {
|
||||
override readonly name = 'DomainError'
|
||||
|
||||
/** Present exactly when `code` is `invalid-record`. */
|
||||
readonly detail?: InvalidRecordDetail
|
||||
|
||||
/**
|
||||
* @param code - Stable discriminant for the failure class.
|
||||
* @param message - Human-readable diagnostic detail.
|
||||
* @param options - Standard error options plus the `invalid-record` location.
|
||||
*/
|
||||
constructor(
|
||||
readonly code: DomainErrorCode,
|
||||
message: string,
|
||||
options?: DomainErrorOptions,
|
||||
) {
|
||||
super(message, options)
|
||||
if (options?.detail) this.detail = options.detail
|
||||
}
|
||||
}
|
||||
48
packages/storage/storage-domain/src/events.ts
Normal file
48
packages/storage/storage-domain/src/events.ts
Normal file
@@ -0,0 +1,48 @@
|
||||
/**
|
||||
* Change-event vocabulary of the domain data form. Every durable write emits
|
||||
* one event after the backend resolves durability, carrying the new snapshot
|
||||
* and an operation discriminant — never the old value (a diffing consumer
|
||||
* keeps its own previous snapshot). This is the event source for cross-process
|
||||
* change push (RPC frames) in a later phase.
|
||||
* @module @deepseek-ai/dsh-storage-domain/src/events
|
||||
*/
|
||||
|
||||
/** Shared location fields of one durable domain change. */
|
||||
export interface DomainChangedBase {
|
||||
/** Owning domain name. */
|
||||
readonly domain: string
|
||||
/** Table name; `''` for a global-singleton write. */
|
||||
readonly table: string
|
||||
/** Record key; `''` for a global-singleton write. */
|
||||
readonly key: string
|
||||
}
|
||||
|
||||
/** A record (or the global singleton) was inserted or overwritten. */
|
||||
export interface DomainChangedPut extends DomainChangedBase {
|
||||
readonly operation: 'put'
|
||||
/** The new snapshot. */
|
||||
readonly value: unknown
|
||||
}
|
||||
|
||||
/** A record was deleted; tombstones carry no value. */
|
||||
export interface DomainChangedDeleted extends DomainChangedBase {
|
||||
readonly operation: 'deleted'
|
||||
readonly value?: never
|
||||
}
|
||||
|
||||
/** One durable domain change; a closed union — switch on `operation`. */
|
||||
export type DomainChanged = DomainChangedPut | DomainChangedDeleted
|
||||
|
||||
declare module 'cordis' {
|
||||
interface Events {
|
||||
/**
|
||||
* A domain record or the global singleton changed, emitted once per write
|
||||
* strictly after the backend acknowledged durability. Events of one
|
||||
* domain arrive in its write-chain order.
|
||||
* @param change - domain, table (`''` for global), key (`''` for global),
|
||||
* operation discriminant, and on `put` the new snapshot.
|
||||
* @mode emit
|
||||
*/
|
||||
'domain/changed'(change: DomainChanged): void
|
||||
}
|
||||
}
|
||||
203
packages/storage/storage-domain/src/index.ts
Normal file
203
packages/storage/storage-domain/src/index.ts
Normal file
@@ -0,0 +1,203 @@
|
||||
/**
|
||||
* Domain data form (`ctx.storage.domain`): schema-validated, change-emitting
|
||||
* KV domains over storage backends. The single implementation of the domain
|
||||
* layer — consumers depend on this package and never touch backends directly.
|
||||
* Plugin `Config` is schemastery; record schemas inside domain specs are zod
|
||||
* (see `src/spec.ts` for the split rationale).
|
||||
* @module @deepseek-ai/dsh-storage-domain
|
||||
*/
|
||||
|
||||
import type { Context } from 'cordis'
|
||||
import z from 'schemastery'
|
||||
import { DomainError } from './error.ts'
|
||||
import { descriptorOf } from './spec.ts'
|
||||
import type { DomainSpec } from './spec.ts'
|
||||
import { DomainImpl } from './domain.ts'
|
||||
import type { Domain } from './domain.ts'
|
||||
|
||||
export { DomainError } from './error.ts'
|
||||
export type { DomainErrorCode, DomainErrorOptions, InvalidRecordDetail } from './error.ts'
|
||||
export { defineDomain, domainTable, descriptorOf } from './spec.ts'
|
||||
export type {
|
||||
DomainSpec, DomainGlobalSpec, DomainTableSpec,
|
||||
TableKeyOf, TableValueOf, GlobalValueOf,
|
||||
} from './spec.ts'
|
||||
export type { DomainChanged } from './events.ts'
|
||||
export type { Domain, DomainGlobal, DomainGlobalHandleOf, KvTable } from './domain.ts'
|
||||
|
||||
declare module '@deepseek-ai/dsh-storage' {
|
||||
interface StorageForms {
|
||||
domain: DomainFacility
|
||||
}
|
||||
}
|
||||
|
||||
/** Cordis plugin name. */
|
||||
export const name = 'storage-domain'
|
||||
/** The storage hub must be present before the form can mount. */
|
||||
export const inject = ['storage']
|
||||
|
||||
/**
|
||||
* Plugin config. Which backend serves which domain is decided here, not
|
||||
* globally on the hub: `backend` is the default route and `routes` overrides
|
||||
* it per domain name. A route naming an unregistered backend fails loud at
|
||||
* `open` with `backend-not-found`.
|
||||
*/
|
||||
export interface Config {
|
||||
/** Default backend name for every domain without an explicit route. Required: there is no universally correct medium. */
|
||||
backend: string
|
||||
/** Per-domain overrides: domain name → backend name. */
|
||||
routes?: Record<string, string>
|
||||
}
|
||||
|
||||
export const Config: z<Config> = z.object({
|
||||
backend: z.string().required(),
|
||||
routes: z.dict(z.string()).default({}),
|
||||
})
|
||||
|
||||
/**
|
||||
* The mounted domain facility. Opens declared domains over routed backends;
|
||||
* one facility instance owns the open-domain table and enforces single-open
|
||||
* per domain name.
|
||||
*/
|
||||
export class DomainFacility {
|
||||
private readonly domains = new Map<string, DomainImpl>()
|
||||
/** Names reserved by an in-flight or completed open, so concurrent opens of one name fail loud. */
|
||||
private readonly reserved = new Set<string>()
|
||||
|
||||
/**
|
||||
* @param ctx - Context of the domain plugin; open-domain effects and change
|
||||
* events attach here.
|
||||
* @param config - Validated plugin config.
|
||||
*/
|
||||
constructor(
|
||||
private readonly ctx: Context,
|
||||
private readonly config: Config,
|
||||
) {}
|
||||
|
||||
/**
|
||||
* Open one declared domain. Steps, each failing the whole call: reject a
|
||||
* name that is already open (`already-open`); resolve the backend route
|
||||
* (`backend-not-found` passes through from the hub); require its `kv` facet
|
||||
* (`facet-unsupported`); open the unit projected from the spec (backend
|
||||
* `version-mismatch`/`malformed-medium` pass through); load and validate
|
||||
* every stored record against the spec's zod schemas (`invalid-record`
|
||||
* with the offending table and key); construct the domain.
|
||||
*
|
||||
* Lifecycle: the CALLER owns the returned handle and closes it via
|
||||
* `Domain.close()` (typically as its own `ctx.effect` disposer) — the
|
||||
* facility does not tie the domain to any consumer fiber. Domains still
|
||||
* open when the facility unmounts are closed by the plugin disposer.
|
||||
* @param spec - The domain declaration, typically from `defineDomain`.
|
||||
* @returns the opened domain handle, typed by the spec.
|
||||
*/
|
||||
async open<S extends DomainSpec>(spec: S): Promise<Domain<S>> {
|
||||
if (this.reserved.has(spec.name)) {
|
||||
throw new DomainError('already-open', `domain '${spec.name}' is already open`)
|
||||
}
|
||||
this.reserved.add(spec.name)
|
||||
try {
|
||||
const backendName = this.config.routes?.[spec.name] ?? this.config.backend
|
||||
const backend = this.ctx.storage.backend.get(backendName)
|
||||
if (!backend.kv) {
|
||||
throw new DomainError(
|
||||
'facet-unsupported',
|
||||
`backend '${backendName}' routed for domain '${spec.name}' has no kv facet`,
|
||||
)
|
||||
}
|
||||
const unit = await backend.kv.open(descriptorOf(spec))
|
||||
try {
|
||||
const snapshot = await unit.loadAll()
|
||||
const tables = new Map<string, Map<string, unknown>>()
|
||||
for (const [table, tableSpec] of Object.entries(spec.tables)) {
|
||||
const records = new Map<string, unknown>()
|
||||
for (const [key, raw] of Object.entries(snapshot.tables[table] ?? {})) {
|
||||
records.set(key, parseRecord(spec.name, table, key, () => tableSpec.valueSchema.parse(raw)))
|
||||
}
|
||||
tables.set(table, records)
|
||||
}
|
||||
// A null stored global means "never written": serve `initial` without
|
||||
// materializing it — the first `set` writes.
|
||||
const globalSpec = spec.global
|
||||
const globalValue = globalSpec === undefined
|
||||
? undefined
|
||||
: snapshot.global === null
|
||||
? globalSpec.initial
|
||||
: parseRecord(spec.name, '', '', () => globalSpec.schema.parse(snapshot.global))
|
||||
// The onClosed hook runs strictly after teardown completes: writes
|
||||
// landing during the drain still emit domain/changed, and the domain
|
||||
// stays resolvable (the package invariant cross-checks each event)
|
||||
// until fully closed — only then does the name free up for reopening.
|
||||
const domain: DomainImpl = new DomainImpl(this.ctx, spec, unit, tables, globalValue, () => {
|
||||
this.domains.delete(spec.name)
|
||||
this.reserved.delete(spec.name)
|
||||
})
|
||||
this.domains.set(spec.name, domain)
|
||||
// The single type-erasure point: DomainImpl is the untyped runtime,
|
||||
// Domain<S> the spec-typed view; the unknown hop is required because
|
||||
// S's conditional global-handle type stays unresolved here.
|
||||
return domain as unknown as Domain<S>
|
||||
} catch (error) {
|
||||
await unit.close()
|
||||
throw error
|
||||
}
|
||||
} catch (error) {
|
||||
// Any failure means the domain never registered (nothing can throw
|
||||
// after it), so releasing the name reservation is unconditional.
|
||||
this.reserved.delete(spec.name)
|
||||
throw error
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Look up an open domain by name, untyped. Diagnostic surface (the package
|
||||
* invariant cross-checks change events against live domain state); typed
|
||||
* consumers hold the handle returned by {@link open}.
|
||||
* @param name - Domain name.
|
||||
* @returns the open domain runtime, or `undefined` when not open.
|
||||
*/
|
||||
get(name: string): DomainImpl | undefined {
|
||||
return this.domains.get(name)
|
||||
}
|
||||
|
||||
/**
|
||||
* Close every domain still open on this facility. The unmount path for
|
||||
* consumers that never called `Domain.close()` themselves; closing is
|
||||
* idempotent, so double-closing an already-closed domain is harmless.
|
||||
* @returns resolution after every unit is released.
|
||||
*/
|
||||
async closeAll(): Promise<void> {
|
||||
await Promise.all([...this.domains.values()].map(domain => domain.close()))
|
||||
}
|
||||
}
|
||||
|
||||
/** Run one zod parse, translating failure to `invalid-record` with its location. */
|
||||
function parseRecord<T>(domain: string, table: string, key: string, parse: () => T): T {
|
||||
try {
|
||||
return parse()
|
||||
} catch (error) {
|
||||
const slot = table === '' ? 'global' : `record '${key}' in table '${table}'`
|
||||
throw new DomainError(
|
||||
'invalid-record',
|
||||
`domain '${domain}': stored ${slot} does not match its schema`,
|
||||
{ detail: { table, key }, cause: error },
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Mount the domain data form on the storage hub.
|
||||
* @param ctx - Plugin context.
|
||||
* @param config - Validated plugin config.
|
||||
*/
|
||||
export function apply(ctx: Context, config: Config) {
|
||||
const facility = new DomainFacility(ctx, config)
|
||||
ctx.effect(() => {
|
||||
const unmount = ctx.storage.mount('domain', facility)
|
||||
return async () => {
|
||||
// Close leftovers before unmounting: draining writes still emit
|
||||
// domain/changed, whose invariant resolves the facility through the hub.
|
||||
await facility.closeAll()
|
||||
unmount()
|
||||
}
|
||||
})
|
||||
}
|
||||
67
packages/storage/storage-domain/src/invariant.ts
Normal file
67
packages/storage/storage-domain/src/invariant.ts
Normal file
@@ -0,0 +1,67 @@
|
||||
/**
|
||||
* Package-owned invariant companion for `@deepseek-ai/dsh-storage-domain`: every
|
||||
* `domain/changed` event must agree with the emitting domain's authoritative
|
||||
* in-memory state (the owned event-stream ↔ mutable-data relationship of this
|
||||
* package). Writes emit strictly after mutating memory and the write chain
|
||||
* serializes them, so at emission time the event's snapshot equals the
|
||||
* current read — any divergence means a write path skipped the chain or
|
||||
* emitted a stale value.
|
||||
* @module @deepseek-ai/dsh-storage-domain/invariant
|
||||
*/
|
||||
|
||||
import type { Context } from 'cordis'
|
||||
import type { InvariantFailure, InvariantInstaller } from '@deepseek-ai/dsh-invariants'
|
||||
import type { DomainChanged } from './events.ts'
|
||||
|
||||
const PACKAGE_NAME = '@deepseek-ai/dsh-storage-domain'
|
||||
|
||||
/** Cordis companion plugin name. */
|
||||
export const name = 'storage-domain-invariant'
|
||||
/** Service required before the companion can reserve package ownership. */
|
||||
export const inject = ['invariants']
|
||||
|
||||
/** Install the change-event ↔ memory-state agreement check. */
|
||||
const install: InvariantInstaller = Object.assign((ctx: Context, fail: InvariantFailure) => {
|
||||
ctx.on('domain/changed', (change: DomainChanged) => {
|
||||
const domain = ctx.storage.form('domain').get(change.domain)
|
||||
if (domain === undefined) {
|
||||
return fail(`domain/changed for '${change.domain}' emitted while that domain is not open`)
|
||||
}
|
||||
if (change.table === '') {
|
||||
// Global write: the event snapshot must be the current global value.
|
||||
if (domain.global.get() !== change.value) {
|
||||
return fail(`domain/changed global value for '${change.domain}' differs from the in-memory global`)
|
||||
}
|
||||
return
|
||||
}
|
||||
const current = domain.table(change.table).get(change.key)
|
||||
switch (change.operation) {
|
||||
case 'deleted':
|
||||
if (current !== undefined) {
|
||||
return fail(
|
||||
`domain/changed deletion of '${change.domain}'.'${change.table}'['${change.key}'] `
|
||||
+ 'emitted while the record is still in memory',
|
||||
)
|
||||
}
|
||||
return
|
||||
case 'put':
|
||||
if (current !== change.value) {
|
||||
return fail(
|
||||
`domain/changed value for '${change.domain}'.'${change.table}'['${change.key}'] `
|
||||
+ 'differs from the in-memory record',
|
||||
)
|
||||
}
|
||||
return
|
||||
default:
|
||||
change satisfies never
|
||||
}
|
||||
}, { global: true })
|
||||
}, { inject: ['storage'] })
|
||||
|
||||
/**
|
||||
* Register this package's invariant companion.
|
||||
* @param ctx - Cordis context carrying the invariant service.
|
||||
* @returns the installed registration's disposer after setup succeeds.
|
||||
*/
|
||||
export const apply = (ctx: Context): Promise<() => void> =>
|
||||
Promise.resolve(ctx.invariants.register(PACKAGE_NAME, install))
|
||||
112
packages/storage/storage-domain/src/spec.ts
Normal file
112
packages/storage/storage-domain/src/spec.ts
Normal file
@@ -0,0 +1,112 @@
|
||||
/**
|
||||
* Domain declaration vocabulary. A spec object is the single source of a
|
||||
* domain's identity, layout, and record schemas: the owning package defines
|
||||
* it once with {@link defineDomain} and both the type surface and the runtime
|
||||
* (validation, descriptor projection) derive from it. Record schemas are zod
|
||||
* (`z.infer` keeps types un-duplicated and the same schemas later project to
|
||||
* RPC wire schemas); plugin `Config` stays schemastery.
|
||||
* @module @deepseek-ai/dsh-storage-domain/src/spec
|
||||
*/
|
||||
|
||||
import type { ZodType } from 'zod'
|
||||
import { UNIT_NAME_RE, type KvUnitDescriptor } from '@deepseek-ai/dsh-storage'
|
||||
|
||||
/** Global singleton declaration: schema plus the value used before the first write. */
|
||||
export interface DomainGlobalSpec<G> {
|
||||
/** Validates the stored global at the durable boundary. */
|
||||
readonly schema: ZodType<G>
|
||||
/** Value served when the medium holds no global yet; not written until the first `set`. */
|
||||
readonly initial: G
|
||||
}
|
||||
|
||||
/**
|
||||
* One table declaration. `K` is a phantom key type (typically a branded
|
||||
* string) carried for compile-time projection only; keys are plain strings on
|
||||
* the medium.
|
||||
*/
|
||||
export interface DomainTableSpec<K extends string = string, V = unknown> {
|
||||
/** Validates every stored record at the durable boundary. */
|
||||
readonly valueSchema: ZodType<V>
|
||||
/** Phantom carrier for the key type; never present at runtime. */
|
||||
readonly __key?: K
|
||||
}
|
||||
|
||||
/** Static declaration of one domain: identity, version, and record layout. */
|
||||
export interface DomainSpec {
|
||||
/** Domain name; must match `UNIT_NAME_RE` (doubles as the backend unit name). */
|
||||
readonly name: string
|
||||
/** Domain format version; a medium stamped with a different version rejects at open. */
|
||||
readonly version: number
|
||||
/** Optional global singleton slot. */
|
||||
readonly global?: DomainGlobalSpec<unknown>
|
||||
/** Table declarations keyed by table name; each name must match `UNIT_NAME_RE`. */
|
||||
readonly tables: Record<string, DomainTableSpec>
|
||||
}
|
||||
|
||||
/** Key type of one declared table, recovered from its phantom carrier. */
|
||||
export type TableKeyOf<S extends DomainSpec, N extends keyof S['tables']> =
|
||||
S['tables'][N] extends DomainTableSpec<infer K> ? K : never
|
||||
|
||||
/** Value type of one declared table. */
|
||||
export type TableValueOf<S extends DomainSpec, N extends keyof S['tables']> =
|
||||
S['tables'][N] extends DomainTableSpec<string, infer V> ? V : never
|
||||
|
||||
/** Global value type of a spec; `never` when the spec declares no global. */
|
||||
export type GlobalValueOf<S extends DomainSpec> =
|
||||
S['global'] extends DomainGlobalSpec<infer G> ? G : never
|
||||
|
||||
/**
|
||||
* Declare one table.
|
||||
* @param schema - zod schema validating every stored record of this table.
|
||||
* @returns the table declaration, key-typed by `K`.
|
||||
*/
|
||||
export function domainTable<K extends string, V>(schema: ZodType<V>): DomainTableSpec<K, V> {
|
||||
return { valueSchema: schema }
|
||||
}
|
||||
|
||||
/**
|
||||
* Identity helper that pins a spec's literal types and validates its shape.
|
||||
* Misconfiguration fails loud at the owning package's module load, before any
|
||||
* medium is touched: a domain or table name outside `UNIT_NAME_RE`, a version
|
||||
* that is not a non-negative integer, or a global schema that accepts `null`
|
||||
* all throw. The `null` rejection guards round-tripping: backends store the
|
||||
* global as opaque JSON with `null` as the "never written" sentinel, so a
|
||||
* nullable global would be indistinguishable from an absent one on reopen
|
||||
* (a stored `null` silently reverts to `initial`).
|
||||
* @param spec - The domain declaration.
|
||||
* @returns the same spec, narrowed to its literal type.
|
||||
*/
|
||||
export function defineDomain<S extends DomainSpec>(spec: S): S {
|
||||
if (!UNIT_NAME_RE.test(spec.name)) {
|
||||
throw new Error(`domain name '${spec.name}' must match ${UNIT_NAME_RE}`)
|
||||
}
|
||||
if (!Number.isInteger(spec.version) || spec.version < 0) {
|
||||
throw new Error(`domain '${spec.name}' version must be a non-negative integer, got ${spec.version}`)
|
||||
}
|
||||
for (const table of Object.keys(spec.tables)) {
|
||||
if (!UNIT_NAME_RE.test(table)) {
|
||||
throw new Error(`domain '${spec.name}' table name '${table}' must match ${UNIT_NAME_RE}`)
|
||||
}
|
||||
}
|
||||
if (spec.global !== undefined && spec.global.schema.safeParse(null).success) {
|
||||
throw new Error(
|
||||
`domain '${spec.name}' global schema must not accept null: `
|
||||
+ 'null is the medium\'s "never written" sentinel, so a stored null could not round-trip',
|
||||
)
|
||||
}
|
||||
return spec
|
||||
}
|
||||
|
||||
/**
|
||||
* Project a spec onto the backend-facing unit descriptor.
|
||||
* @param spec - The domain declaration.
|
||||
* @returns the descriptor handed to `KvFacet.open`.
|
||||
*/
|
||||
export function descriptorOf(spec: DomainSpec): KvUnitDescriptor {
|
||||
return {
|
||||
name: spec.name,
|
||||
version: spec.version,
|
||||
tables: Object.keys(spec.tables),
|
||||
hasGlobal: spec.global !== undefined,
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user