import { type FSWatcher, watch } from "node:fs"; import { type FileHandle, mkdtemp, open, readFile, rm, stat, } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; import type { RecordChunk } from "./core.js"; import { validateRecordChunk } from "./core.js"; export class CanonicalRecordLog { readonly records: RecordChunk[] = []; readonly directory: string; readonly path: string; private readonly file: FileHandle; private readonly onRecord: (record: RecordChunk) => void; private watcher: FSWatcher | undefined; private byteOffset = 0; private remainder = Buffer.alloc(0); private drainQueue: Promise = Promise.resolve(); private watcherError: Error | undefined; private closed = false; private constructor( directory: string, path: string, file: FileHandle, onRecord: (record: RecordChunk) => void, ) { this.directory = directory; this.path = path; this.file = file; this.onRecord = onRecord; } static async create( onRecord: (record: RecordChunk) => void, ): Promise { const directory = await mkdtemp(join(tmpdir(), "pi-decay-")); const path = join(directory, "classification.jsonl"); try { const file = await open(path, "wx+", 0o600); const log = new CanonicalRecordLog(directory, path, file, onRecord); log.startWatching(); return log; } catch (error) { await rm(directory, { recursive: true, force: true }); throw error; } } private startWatching(): void { this.watcher = watch(this.path, { persistent: false }, () => { void this.scheduleDrain().catch((error: unknown) => { this.watcherError = asError(error); }); }); this.watcher.on("error", (error) => { this.watcherError = error; }); } async append(record: RecordChunk): Promise { if (this.closed) throw new Error("Decay record log is closed"); await this.file.appendFile(`${JSON.stringify(record)}\n`, "utf8"); } async drain(): Promise { await this.scheduleDrain(); if (this.watcherError) throw this.watcherError; if (this.remainder.length > 0) { throw new Error("Decay record log ended with an incomplete JSONL record"); } } private scheduleDrain(): Promise { this.drainQueue = this.drainQueue.then(() => this.readAvailable()); return this.drainQueue; } private async readAvailable(): Promise { if (this.closed) return; const content = await readFile(this.path); if (content.length < this.byteOffset) throw new Error("Decay record log was truncated"); if (content.length === this.byteOffset) return; const unread = content.subarray(this.byteOffset); this.byteOffset = content.length; const buffered = Buffer.concat([this.remainder, unread]); let lineStart = 0; for (;;) { const newline = buffered.indexOf(0x0a, lineStart); if (newline < 0) break; const line = buffered.subarray(lineStart, newline); lineStart = newline + 1; if (line.length === 0) continue; const parsed: unknown = JSON.parse(line.toString("utf8")); const record = validateRecordChunk(parsed); if (!record) throw new Error("Decay record log contains an invalid record"); this.records.push(record); this.onRecord(record); } this.remainder = Buffer.from(buffered.subarray(lineStart)); } async mode(): Promise { return (await stat(this.path)).mode & 0o777; } async close(): Promise { if (this.closed) return; this.watcher?.close(); this.closed = true; try { await this.drainQueue.catch(() => {}); await this.file.close(); } finally { await rm(this.directory, { recursive: true, force: true }); } } } function asError(value: unknown): Error { return value instanceof Error ? value : new Error(String(value)); }