repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/decay/record-log.ts
Rawimport { 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<void> = 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<CanonicalRecordLog> {
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<void> {
if (this.closed) throw new Error("Decay record log is closed");
await this.file.appendFile(`${JSON.stringify(record)}\n`, "utf8");
}
async drain(): Promise<void> {
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<void> {
this.drainQueue = this.drainQueue.then(() => this.readAvailable());
return this.drainQueue;
}
private async readAvailable(): Promise<void> {
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<number> {
return (await stat(this.path)).mode & 0o777;
}
async close(): Promise<void> {
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));
}