1
/**2
* Runtime of one open domain: authoritative in-memory state, the single3
* per-domain write chain, and change-event emission. Reads are synchronous4
* from memory; every write queues on the chain, awaits backend durability5
* FIRST, then mutates memory, then emits `domain/changed` — a rejected6
* backend write leaves memory untouched (no divergence between reads and the7
* medium), and events carry values that equal the in-memory state at8
* emission, in write order.9
* @module @deepseek-ai/dsh-storage-domain/src/domain10
*/12
import type { Context } from '@deepseek-ai/cordis'13
import type { KvUnit } from '@deepseek-ai/dsh-storage'14
import { DomainError } from './error.ts'15
import type { DomainSpec, DomainGlobalSpec, TableKeyOf, TableValueOf } from './spec.ts'16
import type { DomainChanged } from './events.ts'18
/** Handle on a domain's global singleton. */19
export 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(): G27
/**28
* Replace the value durably. Queued on the domain's write chain; the first29
* `set` is what materializes the global on the medium.30
* @param value - New value; must satisfy the spec's schema (not re-checked31
* here — validation happens at the durable read boundary).32
* @returns resolution after durability and event emission.33
*/34
set(value: G): Promise<void>35
}37
/**38
* Handle on one declared table. Records are plain immutable data: returned39
* values are the stored objects themselves (no defensive copies) and must not40
* be mutated in place — replace via `put`/`update`.41
*/42
export 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 | undefined50
/**51
* Snapshot iterator over `[key, record]` pairs. A snapshot, not a live52
* view: iteration stays stable while queued writes land.53
* @returns the pair iterator.54
*/55
entries(): IterableIterator<[K, V]>57
/**58
* Snapshot iterator over keys.59
* @returns the key iterator.60
*/61
keys(): IterableIterator<K>63
/** Current record count. */64
readonly size: number66
/**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>74
/**75
* Delete one record durably.76
* @param key - Record key.77
* @returns `true` when the record existed, `false` when it was already78
* absent (no write and no event in that case).79
*/80
delete(key: K): Promise<boolean>82
/**83
* Atomic read-modify-write on the domain's write chain: `fn` sees the84
* 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
}92
/** Global handle of a spec: typed when declared, `never` (inaccessible) when not. */93
export type DomainGlobalHandleOf<S extends DomainSpec> =94
S extends { readonly global: DomainGlobalSpec<infer G> } ? DomainGlobal<G> : never96
/** One open domain, typed by its spec. */97
export interface Domain<S extends DomainSpec> {98
/** Domain name from the spec. */99
readonly name: string100
/** 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 calls104
* 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>>110
/**111
* Close this domain: reject new writes immediately, drain already-queued112
* writes (their events still emit), release the backend unit, then free113
* the domain name for a later open. Idempotent — repeated calls share one114
* 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
}121
/** Internal boundary handing table handles their domain-owned write machinery. */122
interface TableHost {123
readonly domainName: string124
readonly unit: KvUnit125
/** 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(): void129
/** Emit `domain/changed` for one durably landed write. */130
emitChanged(change: DomainChanged): void131
}133
const noop = () => {}135
/**136
* The single domain implementation behind the {@link Domain} interface. The137
* facility constructs it from a validated `loadAll` snapshot and erases it to138
* `Domain<S>`; nothing outside this package constructs one.139
*/140
export class DomainImpl {141
/** Domain name from the spec. */142
readonly name: string144
private readonly tables = new Map<string, KvTableImpl<string, unknown>>()145
private globalValue: unknown146
private readonly globalHandle?: DomainGlobal<unknown>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 = false152
/** Set when close finishes (chain drained, unit closed): reads reject from here on. */153
private closed = false154
private disposal?: Promise<void>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 entry161
* per declared table (empty maps included) — the facility builds it from162
* 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; frees166
* 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.name177
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 = globalValue189
this.globalHandle = {190
get: () => {191
this.assertReadable()192
return this.globalValue193
},194
set: value => this.enqueue(async () => {195
await this.unit.setGlobal(value)196
this.globalValue = value197
this.emitChanged({ domain: this.name, table: '', key: '', operation: 'put', value })198
}),199
}200
}201
}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.globalHandle209
}211
/**212
* Resolve one declared table handle; an undeclared name is a caller bug213
* 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 table223
}225
/**226
* Close this domain: reject new writes immediately, drain already-queued227
* writes (their events still emit), close the unit, then free the name via228
* 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.disposal234
}236
private async runClose(): Promise<void> {237
this.disposing = true238
// Chain links never reject (each is settled via then(noop, noop)), so239
// this await is a pure drain barrier.240
await this.chain241
await this.unit.close()242
this.closed = true243
this.onClosed()244
}246
/**247
* Dispatch one post-durability change notification, containing observer248
* failures: the write is already committed (medium and memory both hold249
* 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 dispatches256
// listeners inline and nothing else runs in the try. The event is a257
// notification, not a transaction participant — the commit point has258
// 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
}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 result270
}272
private assertReadable(): void {273
if (this.closed) {274
throw new DomainError('closed', `domain '${this.name}' is closed`)275
}276
}277
}279
/** Table handle bound to one in-memory record map and its domain's write chain. */280
class 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
) {}287
get(key: K): V | undefined {288
this.host.assertReadable()289
return this.records.get(key) as V | undefined290
}292
entries(): IterableIterator<[K, V]> {293
this.host.assertReadable()294
return ([...this.records.entries()] as [K, V][])[Symbol.iterator]()295
}297
keys(): IterableIterator<K> {298
this.host.assertReadable()299
return ([...this.records.keys()] as K[])[Symbol.iterator]()300
}302
get size(): number {303
this.host.assertReadable()304
return this.records.size305
}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
}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: an318
// earlier queued put of the same key makes this delete observe it.319
if (!this.records.has(key)) return false320
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 true329
})330
}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 next345
})346
}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
}