repositories / will
will
owned by admin
src/agent/spool.ts
Rawimport { 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
}
}
}