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; load(key: StoreKey): Promise; listSubkeys(key: { projectKey: string; sessionId: string; }): Promise; } 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(); async append(key: StoreKey, entries: OpaqueEntry[]): Promise { 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 { dbg?.("cache.memory.load"); const entries = this.entries.get(keyName(key)); return entries ? structuredClone(entries) : null; } async listSubkeys(): Promise { 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; 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 { 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 { dbg?.("cache.checkpoint.load"); return (await loadCheckpoints(root, piSessionId)).at(-1); } const checkpointWrites = new Map>(); export async function saveCheckpoint( root: string, piSessionId: string, checkpoint: KlausCheckpoint, ): Promise { 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 { dbg?.("cache.staging.append", { entryCount: entries.length }); this.batches.push({ key: structuredClone(key), entries: structuredClone(entries), }); } async load(key: StoreKey): Promise { 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 { dbg?.("cache.staging.listSubkeys"); return this.target.listSubkeys(key); } async commit(): Promise { 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>(); private ready?: Promise; constructor(private readonly root: string) { dbg?.("cache.file.create"); } private ensureManifest(): Promise { 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 { dbg?.("cache.file.prepare.start"); await this.ensureManifest(); dbg?.("cache.file.prepare.end"); } async append(key: StoreKey, entries: OpaqueEntry[]): Promise { 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 { 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 { dbg?.("cache.file.listSubkeys"); return []; } }