Luigit
repositories / will

will

owned by admin

src/agent/spool.ts

Raw
import { open as fsOpen, mkdir, readdir, readFile, rename, rm } from 'node:fs/promises'
import { join } from 'node:path'
import type { AuditEvent, ReleaseProvenance } from '../shared/events.ts'

/**
 * Durable audit spool: every agent event is fsynced to disk before anything else
 * happens, then drained idempotently (unique event ids) into the service.
 * On persistent service failure events stay queued forever (spec AC11).
 */

export class AuditSpool {
  readonly #dir: string
  #draining = false

  private constructor(dir: string) {
    this.#dir = dir
  }

  static async open(dir: string): Promise<AuditSpool> {
    await mkdir(dir, { recursive: true })
    return new AuditSpool(dir)
  }

  async append(event: AuditEvent, provenance: ReleaseProvenance): Promise<void> {
    const enriched: AuditEvent = { ...event, provenance: { ...provenance, ...event.provenance } }
    const file = join(this.#dir, `${Date.now()}-${event.id}.json`)
    const tmp = `${file}.tmp`
    const fh = await fsOpen(tmp, 'w')
    try {
      await fh.writeFile(JSON.stringify(enriched), 'utf8')
      await fh.sync()
    } finally {
      await fh.close()
    }
    await rename(tmp, file)
  }

  async list(): Promise<string[]> {
    try {
      return (await readdir(this.#dir)).filter((f: string) => f.endsWith('.json')).sort()
    } catch {
      return []
    }
  }

  async read(name: string): Promise<AuditEvent | undefined> {
    try {
      return JSON.parse(await readFile(join(this.#dir, name), 'utf8')) as AuditEvent
    } catch {
      return undefined
    }
  }

  async remove(name: string): Promise<void> {
    await rm(join(this.#dir, name), { force: true })
  }

  /**
   * Drain the spool into the service ingestion API; idempotent by event id.
   * Returns the number of successfully submitted events.
   */
  async drain(
    submit: (events: AuditEvent[]) => Promise<{ accepted: number; duplicates: number }>,
    batch = 100,
  ): Promise<number> {
    if (this.#draining) return 0
    this.#draining = true
    try {
      let submitted = 0
      const names = await this.list()
      for (let i = 0; i < names.length; i += batch) {
        const slice = names.slice(i, i + batch)
        const events: AuditEvent[] = []
        for (const name of slice) {
          const event = await this.read(name)
          if (event) events.push(event)
        }
        if (events.length === 0) continue
        await submit(events)
        submitted += events.length
        for (const name of slice) await this.remove(name)
      }
      return submitted
    } finally {
      this.#draining = false
    }
  }
}