Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/decay/record-log.ts

Raw
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<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));
}