| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141 |
- /**
- * One opened JSON unit. The in-memory state is authoritative; every write
- * primitive mutates it and republishes the whole file atomically. Writes are
- * NOT queued here — per the backend contract, write ordering belongs to the
- * caller (the domain layer's write chain); this unit only guarantees that
- * each single call publishes a complete, durable file.
- * @module @deepseek-ai/dsh-storage-json/src/unit
- */
- import { readFile } from 'node:fs/promises'
- import { StorageError } from '@deepseek-ai/dsh-storage'
- import type { KvUnit, KvUnitDescriptor } from '@deepseek-ai/dsh-storage'
- import { writeAtomic } from './atomic.ts'
- import { parse, serialize } from './format.ts'
- import type { UnitState } from './format.ts'
- /**
- * Open (load or lazily create) one unit backed by `path`.
- * @param descriptor - Static identity and shape of the unit.
- * @param path - Absolute unit file path under the backend root.
- * @param onClose - Backend callback releasing the unit's open-slot.
- * @returns the opened unit.
- */
- export async function openJsonUnit(
- descriptor: KvUnitDescriptor,
- path: string,
- onClose: () => void,
- ): Promise<KvUnit> {
- let text: string | undefined
- try {
- text = await readFile(path, 'utf8')
- } catch (error) {
- if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error
- // Missing file = empty unit; materialization defers to the first write.
- }
- const state: UnitState =
- text === undefined
- ? {
- version: descriptor.version,
- global: null,
- tables: new Map(descriptor.tables.map(table => [table, new Map<string, unknown>()])),
- }
- : parse(text, descriptor)
- return new JsonKvUnit(descriptor, path, state, onClose)
- }
- class JsonKvUnit implements KvUnit {
- private closed = false
- /** In-flight publishes; close() drains them before releasing the unit. */
- private readonly inFlight = new Set<Promise<void>>()
- constructor(
- private readonly descriptor: KvUnitDescriptor,
- private readonly path: string,
- private readonly state: UnitState,
- private readonly onClose: () => void,
- ) {}
- // eslint-disable-next-line @typescript-eslint/require-await -- async keeps the closed guard a rejection, not a synchronous throw
- async loadAll(): Promise<{ tables: Record<string, Record<string, unknown>>; global: unknown }> {
- this.assertOpen()
- const tables: Record<string, Record<string, unknown>> = {}
- for (const [table, records] of this.state.tables) {
- tables[table] = Object.fromEntries(records)
- }
- return { tables, global: this.state.global }
- }
- async putRecord(table: string, key: string, value: unknown): Promise<void> {
- this.assertOpen()
- const records = this.records(table)
- const hadKey = records.has(key)
- const previous = records.get(key)
- records.set(key, value)
- // Roll back on a failed publish: memory is authoritative, so a rejected
- // write must not survive in memory (or ride along with the next publish).
- await this.publish().catch((error: unknown) => {
- if (hadKey) records.set(key, previous)
- else records.delete(key)
- throw error
- })
- }
- async deleteRecord(table: string, key: string): Promise<void> {
- this.assertOpen()
- const records = this.records(table)
- if (!records.has(key)) return
- const previous = records.get(key)
- records.delete(key)
- await this.publish().catch((error: unknown) => {
- records.set(key, previous)
- throw error
- })
- }
- async setGlobal(value: unknown): Promise<void> {
- this.assertOpen()
- if (!this.descriptor.hasGlobal) {
- throw new Error(`unit '${this.descriptor.name}' does not declare a global slot`)
- }
- const previous = this.state.global
- this.state.global = value
- await this.publish().catch((error: unknown) => {
- this.state.global = previous
- throw error
- })
- }
- async close(): Promise<void> {
- if (this.closed) {
- await Promise.allSettled(this.inFlight)
- return
- }
- this.closed = true
- await Promise.allSettled(this.inFlight)
- this.onClose()
- }
- private assertOpen(): void {
- if (this.closed) {
- throw new StorageError('closed', `unit '${this.descriptor.name}' is closed`)
- }
- }
- private records(table: string): Map<string, unknown> {
- const records = this.state.tables.get(table)
- if (!records) {
- throw new Error(`unit '${this.descriptor.name}' does not declare table '${table}'`)
- }
- return records
- }
- private publish(): Promise<void> {
- const write = writeAtomic(this.path, serialize(this.descriptor.name, this.state))
- this.inFlight.add(write)
- // Swallow only on the tracking branch: the caller still awaits `write`
- // itself, so rejections stay observed exactly once.
- write.catch(() => {}).finally(() => this.inFlight.delete(write))
- return write
- }
- }
|