unit.ts 4.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141
  1. /**
  2. * One opened JSON unit. The in-memory state is authoritative; every write
  3. * primitive mutates it and republishes the whole file atomically. Writes are
  4. * NOT queued here — per the backend contract, write ordering belongs to the
  5. * caller (the domain layer's write chain); this unit only guarantees that
  6. * each single call publishes a complete, durable file.
  7. * @module @deepseek-ai/dsh-storage-json/src/unit
  8. */
  9. import { readFile } from 'node:fs/promises'
  10. import { StorageError } from '@deepseek-ai/dsh-storage'
  11. import type { KvUnit, KvUnitDescriptor } from '@deepseek-ai/dsh-storage'
  12. import { writeAtomic } from './atomic.ts'
  13. import { parse, serialize } from './format.ts'
  14. import type { UnitState } from './format.ts'
  15. /**
  16. * Open (load or lazily create) one unit backed by `path`.
  17. * @param descriptor - Static identity and shape of the unit.
  18. * @param path - Absolute unit file path under the backend root.
  19. * @param onClose - Backend callback releasing the unit's open-slot.
  20. * @returns the opened unit.
  21. */
  22. export async function openJsonUnit(
  23. descriptor: KvUnitDescriptor,
  24. path: string,
  25. onClose: () => void,
  26. ): Promise<KvUnit> {
  27. let text: string | undefined
  28. try {
  29. text = await readFile(path, 'utf8')
  30. } catch (error) {
  31. if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error
  32. // Missing file = empty unit; materialization defers to the first write.
  33. }
  34. const state: UnitState =
  35. text === undefined
  36. ? {
  37. version: descriptor.version,
  38. global: null,
  39. tables: new Map(descriptor.tables.map(table => [table, new Map<string, unknown>()])),
  40. }
  41. : parse(text, descriptor)
  42. return new JsonKvUnit(descriptor, path, state, onClose)
  43. }
  44. class JsonKvUnit implements KvUnit {
  45. private closed = false
  46. /** In-flight publishes; close() drains them before releasing the unit. */
  47. private readonly inFlight = new Set<Promise<void>>()
  48. constructor(
  49. private readonly descriptor: KvUnitDescriptor,
  50. private readonly path: string,
  51. private readonly state: UnitState,
  52. private readonly onClose: () => void,
  53. ) {}
  54. // oxlint-disable-next-line typescript/require-await -- async keeps the closed guard a rejection, not a synchronous throw
  55. async loadAll(): Promise<{ tables: Record<string, Record<string, unknown>>; global: unknown }> {
  56. this.assertOpen()
  57. const tables: Record<string, Record<string, unknown>> = {}
  58. for (const [table, records] of this.state.tables) {
  59. tables[table] = Object.fromEntries(records)
  60. }
  61. return { tables, global: this.state.global }
  62. }
  63. async putRecord(table: string, key: string, value: unknown): Promise<void> {
  64. this.assertOpen()
  65. const records = this.records(table)
  66. const hadKey = records.has(key)
  67. const previous = records.get(key)
  68. records.set(key, value)
  69. // Roll back on a failed publish: memory is authoritative, so a rejected
  70. // write must not survive in memory (or ride along with the next publish).
  71. await this.publish().catch((error: unknown) => {
  72. if (hadKey) records.set(key, previous)
  73. else records.delete(key)
  74. throw error
  75. })
  76. }
  77. async deleteRecord(table: string, key: string): Promise<void> {
  78. this.assertOpen()
  79. const records = this.records(table)
  80. if (!records.has(key)) return
  81. const previous = records.get(key)
  82. records.delete(key)
  83. await this.publish().catch((error: unknown) => {
  84. records.set(key, previous)
  85. throw error
  86. })
  87. }
  88. async setGlobal(value: unknown): Promise<void> {
  89. this.assertOpen()
  90. if (!this.descriptor.hasGlobal) {
  91. throw new Error(`unit '${this.descriptor.name}' does not declare a global slot`)
  92. }
  93. const previous = this.state.global
  94. this.state.global = value
  95. await this.publish().catch((error: unknown) => {
  96. this.state.global = previous
  97. throw error
  98. })
  99. }
  100. async close(): Promise<void> {
  101. if (this.closed) {
  102. await Promise.allSettled(this.inFlight)
  103. return
  104. }
  105. this.closed = true
  106. await Promise.allSettled(this.inFlight)
  107. this.onClose()
  108. }
  109. private assertOpen(): void {
  110. if (this.closed) {
  111. throw new StorageError('closed', `unit '${this.descriptor.name}' is closed`)
  112. }
  113. }
  114. private records(table: string): Map<string, unknown> {
  115. const records = this.state.tables.get(table)
  116. if (!records) {
  117. throw new Error(`unit '${this.descriptor.name}' does not declare table '${table}'`)
  118. }
  119. return records
  120. }
  121. private publish(): Promise<void> {
  122. const write = writeAtomic(this.path, serialize(this.descriptor.name, this.state))
  123. this.inFlight.add(write)
  124. // Swallow only on the tracking branch: the caller still awaits `write`
  125. // itself, so rejections stay observed exactly once.
  126. write.catch(() => {}).finally(() => this.inFlight.delete(write))
  127. return write
  128. }
  129. }