Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/klaus/src/cache.ts

Raw
import { createHash, randomUUID } from "node:crypto";
import {
	mkdir,
	open,
	readdir,
	readFile,
	rename,
	writeFile,
} from "node:fs/promises";
import { join } from "node:path";
import { dbg } from "./debug.js";
import { SDK_VERSION } from "./protocol.js";

export interface StoreKey {
	projectKey: string;
	sessionId: string;
	subpath?: string;
}

export interface OpaqueEntry {
	type: string;
	uuid?: string;
	[key: string]: unknown;
}

export interface KlausSessionStore {
	append(key: StoreKey, entries: OpaqueEntry[]): Promise<void>;
	load(key: StoreKey): Promise<OpaqueEntry[] | null>;
	listSubkeys(key: {
		projectKey: string;
		sessionId: string;
	}): Promise<string[]>;
}

const keyName = (key: StoreKey): string =>
	createHash("sha256")
		.update(`${key.projectKey}\0${key.sessionId}\0${key.subpath ?? ""}`)
		.digest("hex");

export class KlausCacheCorruptionError extends Error {
	constructor(message: string, options?: ErrorOptions) {
		super(`Klaus cache corruption: ${message}`, options);
		this.name = "KlausCacheCorruptionError";
	}
}

function validateEntry(value: unknown): OpaqueEntry {
	dbg?.("cache.validateEntry");
	if (
		typeof value !== "object" ||
		value === null ||
		Array.isArray(value) ||
		typeof (value as { type?: unknown }).type !== "string"
	) {
		throw new KlausCacheCorruptionError("invalid opaque session entry.");
	}
	return value as OpaqueEntry;
}

export class MemorySessionStore implements KlausSessionStore {
	private readonly entries = new Map<string, OpaqueEntry[]>();

	async append(key: StoreKey, entries: OpaqueEntry[]): Promise<void> {
		dbg?.("cache.memory.append.start", { entryCount: entries.length });
		const name = keyName(key);
		const current = this.entries.get(name) ?? [];
		const uuids = new Set(
			current.flatMap((entry) => (entry.uuid ? [entry.uuid] : [])),
		);
		for (const entry of entries) {
			const valid = validateEntry(structuredClone(entry));
			if (valid.uuid && uuids.has(valid.uuid)) continue;
			if (valid.uuid) uuids.add(valid.uuid);
			current.push(valid);
		}
		this.entries.set(name, current);
		dbg?.("cache.memory.append.end", { storedCount: current.length });
	}

	async load(key: StoreKey): Promise<OpaqueEntry[] | null> {
		dbg?.("cache.memory.load");
		const entries = this.entries.get(keyName(key));
		return entries ? structuredClone(entries) : null;
	}

	async listSubkeys(): Promise<string[]> {
		dbg?.("cache.memory.listSubkeys");
		return [];
	}
}

export interface KlausCheckpoint {
	sdkSessionId: string;
	fingerprint: string;
	position?: string;
	piLeafId?: string;
	messageCount: number;
	protocol: 1;
}

function validateCheckpoint(value: unknown): KlausCheckpoint {
	const checkpoint = value as Partial<KlausCheckpoint>;
	if (
		typeof value !== "object" ||
		value === null ||
		Array.isArray(value) ||
		typeof checkpoint.sdkSessionId !== "string" ||
		typeof checkpoint.fingerprint !== "string" ||
		!Number.isInteger(checkpoint.messageCount) ||
		(checkpoint.messageCount ?? -1) < 0 ||
		(checkpoint.position !== undefined &&
			typeof checkpoint.position !== "string") ||
		(checkpoint.piLeafId !== undefined &&
			typeof checkpoint.piLeafId !== "string") ||
		checkpoint.protocol !== 1
	) {
		throw new Error("Klaus checkpoint is invalid.");
	}
	return checkpoint as KlausCheckpoint;
}

function checkpointPath(root: string, piSessionId: string): string {
	return join(
		root,
		`${keyName({ projectKey: "pi", sessionId: piSessionId })}.json`,
	);
}

export async function loadCheckpoints(
	root: string,
	piSessionId: string,
): Promise<KlausCheckpoint[]> {
	dbg?.("cache.checkpoints.load.start");
	try {
		const value: unknown = JSON.parse(
			await readFile(checkpointPath(root, piSessionId), "utf8"),
		);
		const checkpoints = Array.isArray(value)
			? value.map(validateCheckpoint)
			: [validateCheckpoint(value)];
		dbg?.("cache.checkpoints.load.end", { count: checkpoints.length });
		return checkpoints;
	} catch (error) {
		if ((error as NodeJS.ErrnoException).code === "ENOENT") return [];
		throw error;
	}
}

export async function loadCheckpoint(
	root: string,
	piSessionId: string,
): Promise<KlausCheckpoint | undefined> {
	dbg?.("cache.checkpoint.load");
	return (await loadCheckpoints(root, piSessionId)).at(-1);
}

const checkpointWrites = new Map<string, Promise<void>>();

export async function saveCheckpoint(
	root: string,
	piSessionId: string,
	checkpoint: KlausCheckpoint,
): Promise<void> {
	dbg?.("cache.checkpoint.save.start");
	const path = checkpointPath(root, piSessionId);
	const previous = checkpointWrites.get(path) ?? Promise.resolve();
	const write = previous
		.catch(() => undefined)
		.then(async () => {
			await mkdir(root, { recursive: true, mode: 0o700 });
			const checkpoints = await loadCheckpoints(root, piSessionId);
			const existing = checkpoints.findIndex(
				(item) =>
					item.piLeafId === checkpoint.piLeafId &&
					item.fingerprint === checkpoint.fingerprint,
			);
			if (existing >= 0) checkpoints[existing] = checkpoint;
			else checkpoints.push(checkpoint);
			const temporary = `${path}.${process.pid}.${randomUUID()}.tmp`;
			await writeFile(temporary, JSON.stringify(checkpoints), { mode: 0o600 });
			await rename(temporary, path);
		});
	checkpointWrites.set(path, write);
	try {
		await write;
	} finally {
		if (checkpointWrites.get(path) === write) checkpointWrites.delete(path);
	}
	dbg?.("cache.checkpoint.save.end");
}

export class StagingSessionStore implements KlausSessionStore {
	private readonly batches: Array<{ key: StoreKey; entries: OpaqueEntry[] }> =
		[];

	constructor(private readonly target: KlausSessionStore) {
		dbg?.("cache.staging.create");
	}

	async append(key: StoreKey, entries: OpaqueEntry[]): Promise<void> {
		dbg?.("cache.staging.append", { entryCount: entries.length });
		this.batches.push({
			key: structuredClone(key),
			entries: structuredClone(entries),
		});
	}

	async load(key: StoreKey): Promise<OpaqueEntry[] | null> {
		dbg?.("cache.staging.load");
		const stored = (await this.target.load(key)) ?? [];
		const staged = this.batches
			.filter((batch) => keyName(batch.key) === keyName(key))
			.flatMap((batch) => batch.entries);
		return stored.length || staged.length
			? structuredClone([...stored, ...staged])
			: null;
	}

	async listSubkeys(key: {
		projectKey: string;
		sessionId: string;
	}): Promise<string[]> {
		dbg?.("cache.staging.listSubkeys");
		return this.target.listSubkeys(key);
	}

	async commit(): Promise<void> {
		dbg?.("cache.staging.commit.start", { batchCount: this.batches.length });
		for (const batch of this.batches) {
			await this.target.append(batch.key, batch.entries);
			dbg?.("cache.staging.commit.batch", { entryCount: batch.entries.length });
		}
		this.batches.length = 0;
		dbg?.("cache.staging.commit.end");
	}
}

export class FileSessionStore implements KlausSessionStore {
	private readonly writes = new Map<string, Promise<void>>();
	private ready?: Promise<void>;

	constructor(private readonly root: string) {
		dbg?.("cache.file.create");
	}

	private ensureManifest(): Promise<void> {
		dbg?.("cache.file.ensureManifest", { existing: Boolean(this.ready) });
		this.ready ??= (async () => {
			await mkdir(this.root, { recursive: true, mode: 0o700 });
			const path = join(this.root, "manifest.json");
			try {
				const value: unknown = JSON.parse(await readFile(path, "utf8"));
				if (
					typeof value !== "object" ||
					value === null ||
					(value as { schema?: unknown }).schema !== 1 ||
					(value as { sdk?: unknown }).sdk !== SDK_VERSION
				) {
					throw new Error("Klaus cache requires a tested migration.");
				}
			} catch (error) {
				if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error;
				const existing = await readdir(this.root);
				if (existing.length > 0) {
					throw new Error("Klaus cache requires a tested migration.");
				}
				await writeFile(path, JSON.stringify({ schema: 1, sdk: SDK_VERSION }), {
					mode: 0o600,
				});
			}
		})();
		return this.ready;
	}

	async prepare(): Promise<void> {
		dbg?.("cache.file.prepare.start");
		await this.ensureManifest();
		dbg?.("cache.file.prepare.end");
	}

	async append(key: StoreKey, entries: OpaqueEntry[]): Promise<void> {
		dbg?.("cache.file.append.start", { entryCount: entries.length });
		await this.ensureManifest();
		const name = keyName(key);
		const previous = this.writes.get(name) ?? Promise.resolve();
		const write = previous.then(async () => {
			await mkdir(this.root, { recursive: true, mode: 0o700 });
			const existing = (await this.load(key)) ?? [];
			const uuids = new Set(
				existing.flatMap((entry) => (entry.uuid ? [entry.uuid] : [])),
			);
			const fresh: OpaqueEntry[] = [];
			for (const entry of entries) {
				const valid = validateEntry(structuredClone(entry));
				if (valid.uuid && uuids.has(valid.uuid)) continue;
				if (valid.uuid) uuids.add(valid.uuid);
				fresh.push(valid);
			}
			if (fresh.length === 0) return;
			const file = await open(join(this.root, `${name}.ndjson`), "a", 0o600);
			try {
				await file.write(
					`${fresh.map((entry) => JSON.stringify(entry)).join("\n")}\n`,
				);
			} finally {
				await file.close();
			}
		});
		this.writes.set(name, write);
		try {
			await write;
		} finally {
			if (this.writes.get(name) === write) this.writes.delete(name);
		}
		dbg?.("cache.file.append.end");
	}

	async load(key: StoreKey): Promise<OpaqueEntry[] | null> {
		dbg?.("cache.file.load.start");
		await this.ensureManifest();
		try {
			const text = await readFile(
				join(this.root, `${keyName(key)}.ndjson`),
				"utf8",
			);
			return text
				.split("\n")
				.filter(Boolean)
				.map((line) => validateEntry(JSON.parse(line) as unknown));
		} catch (error) {
			if ((error as NodeJS.ErrnoException).code === "ENOENT") return null;
			if (error instanceof KlausCacheCorruptionError) throw error;
			if (error instanceof SyntaxError) {
				throw new KlausCacheCorruptionError("invalid transcript JSON.", {
					cause: error,
				});
			}
			throw error;
		}
	}

	async listSubkeys(): Promise<string[]> {
		dbg?.("cache.file.listSubkeys");
		return [];
	}
}