repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/klaus/src/cache.ts
Rawimport { 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 [];
}
}