Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
11 kB · 303 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304import { createCipheriv, createDecipheriv, createHash, randomBytes } from "node:crypto";import fs from "node:fs/promises";import path from "node:path";
interface StoredEntry<T> { value: T; expiresAt?: number;}
interface StoreDocument<T> { version: 1; entries: Record<string, StoredEntry<T>>;}
interface EncryptedEnvelope { version: 1; algorithm: "aes-256-gcm"; iv: string; ciphertext: string; tag: string;}
export interface SecureJsonStoreOptions<T> { directory: string; name: string; key: Buffer; maxEntries: number; maxSerializedBytes: number; ttlMs?: number; /** Flush file and directory before acknowledging an admission write. */ durable?: boolean; now?: () => number; validate?: (value: unknown) => T;}
export interface StoreReadResult<T> { value: T; expired: boolean;}
export class SecureJsonStore<T> { private readonly filePath: string; private readonly now: () => number; private readonly validate: (value: unknown) => T; private document: StoreDocument<T> = { version: 1, entries: {} }; private operation = Promise.resolve(); private initialization?: Promise<void>; private initialized = false;
constructor(private readonly options: SecureJsonStoreOptions<T>) { if (!path.isAbsolute(options.directory)) throw new Error("Secure store directory must be absolute"); if (!/^[a-z][a-z0-9-]*$/.test(options.name)) throw new Error("Secure store name is invalid"); if (options.key.length !== 32) throw new Error("Secure store key must contain exactly 32 bytes"); if (!Number.isSafeInteger(options.maxEntries) || options.maxEntries < 1 || options.maxEntries > 100_000) { throw new Error("Secure store maxEntries must be between 1 and 100000"); } if (!Number.isSafeInteger(options.maxSerializedBytes) || options.maxSerializedBytes < 1_024 || options.maxSerializedBytes > 64 * 1024 * 1024) { throw new Error("Secure store maxSerializedBytes must be between 1024 and 67108864"); } if (options.ttlMs !== undefined && (!Number.isSafeInteger(options.ttlMs) || options.ttlMs < 1_000)) { throw new Error("Secure store TTL must be at least one second"); } this.filePath = path.join(options.directory, `${options.name}.enc.json`); this.now = options.now ?? Date.now; this.validate = options.validate ?? ((value) => value as T); }
async initialize(): Promise<void> { if (this.initialized) return; this.initialization ??= this.initializeOnce(); return this.initialization; }
async get(key: string): Promise<T | undefined> { const result = await this.getWithExpiration(key); if (!result) return undefined; if (result.expired) { await this.del(key); return undefined; } return result.value; }
async getWithExpiration(key: string): Promise<StoreReadResult<T> | undefined> { await this.initialize(); await this.operation; const entry = this.document.entries[key]; if (!entry) return undefined; return { value: structuredClone(entry.value), expired: entry.expiresAt !== undefined && entry.expiresAt <= this.now(), }; }
async set(key: string, value: T, guard?: () => boolean): Promise<void> { await this.mutate(() => { this.document.entries[key] = { value: structuredClone(value), ...(this.options.ttlMs ? { expiresAt: this.now() + this.options.ttlMs } : {}), }; }, guard); }
async replaceAll(key: string, value: T, guard?: () => boolean): Promise<void> { await this.mutate(() => { this.document.entries = { [key]: { value: structuredClone(value), ...(this.options.ttlMs ? { expiresAt: this.now() + this.options.ttlMs } : {}), }, }; }, guard); }
async del(key: string): Promise<void> { await this.mutate(() => { delete this.document.entries[key]; }); }
async replaceIf(key: string, predicate: (value: T) => boolean, value: T): Promise<boolean> { let replaced = false; await this.mutate(() => { const entry = this.document.entries[key]; if (!entry || !predicate(entry.value)) return; this.document.entries[key] = { value: structuredClone(value), ...(this.options.ttlMs ? { expiresAt: this.now() + this.options.ttlMs } : {}), }; replaced = true; }); return replaced; }
async deleteIf(key: string, predicate: (value: T) => boolean): Promise<boolean> { let deleted = false; await this.mutate(() => { const entry = this.document.entries[key]; if (!entry || !predicate(entry.value)) return; delete this.document.entries[key]; deleted = true; }); return deleted; }
async take(key: string, guard?: () => boolean): Promise<T | undefined> { let value: T | undefined; await this.mutate(() => { const entry = this.document.entries[key]; if (!entry || (entry.expiresAt !== undefined && entry.expiresAt <= this.now())) { delete this.document.entries[key]; return; } value = structuredClone(entry.value); delete this.document.entries[key]; }, guard); return value; }
async size(): Promise<number> { await this.initialize(); await this.operation; return Object.keys(this.document.entries).length; }
private async initializeOnce(): Promise<void> { await fs.mkdir(this.options.directory, { recursive: true, mode: 0o700 }); await fs.chmod(this.options.directory, 0o700); const stat = await fs.stat(this.filePath).catch((error: NodeJS.ErrnoException) => { if (error.code === "ENOENT") return undefined; throw error; }); if (stat) { const maxEnvelopeBytes = Math.ceil(this.options.maxSerializedBytes * 1.5) + 4_096; if (!stat.isFile() || stat.size > maxEnvelopeBytes) throw new Error("Secure store envelope exceeds its configured bound"); this.document = this.decryptDocument(await fs.readFile(this.filePath, "utf8")); } this.removeExpired(); this.assertBounds(); this.initialized = true; }
private async mutate(change: () => void, guard?: () => boolean): Promise<void> { await this.initialize(); const next = this.operation.then(async () => { const previous = structuredClone(this.document); try { if (guard && !guard()) throw new Error("Secure store write lost authority"); change(); this.removeExpired(); this.assertBounds(); if (guard && !guard()) throw new Error("Secure store write lost authority"); await this.persist(guard); } catch (error) { this.document = previous; throw error; } }); this.operation = next.catch(() => undefined); return next; }
private removeExpired(): void { const now = this.now(); for (const [key, entry] of Object.entries(this.document.entries)) { if (entry.expiresAt !== undefined && entry.expiresAt <= now) delete this.document.entries[key]; } }
private assertBounds(): void { const count = Object.keys(this.document.entries).length; if (count > this.options.maxEntries) throw new Error("Secure store entry limit exceeded"); const bytes = Buffer.byteLength(JSON.stringify(this.document), "utf8"); if (bytes > this.options.maxSerializedBytes) throw new Error("Secure store serialized byte limit exceeded"); }
private async persist(guard?: () => boolean): Promise<void> { const iv = randomBytes(12); const cipher = createCipheriv("aes-256-gcm", this.options.key, iv); const plaintext = Buffer.from(JSON.stringify(this.document), "utf8"); if (plaintext.length > this.options.maxSerializedBytes) throw new Error("Secure store serialized byte limit exceeded"); const ciphertext = Buffer.concat([cipher.update(plaintext), cipher.final()]); const envelope: EncryptedEnvelope = { version: 1, algorithm: "aes-256-gcm", iv: iv.toString("base64"), ciphertext: ciphertext.toString("base64"), tag: cipher.getAuthTag().toString("base64"), }; const tempPath = `${this.filePath}.tmp-${process.pid}-${randomBytes(6).toString("hex")}`; await fs.writeFile(tempPath, `${JSON.stringify(envelope)}\n`, { mode: 0o600, flag: "wx" }); if (guard && !guard()) { await fs.unlink(tempPath); throw new Error("Secure store write lost authority"); } if (this.options.durable) { const file = await fs.open(tempPath, "r"); try { await file.sync(); } finally { await file.close(); } } if (guard && !guard()) { await fs.unlink(tempPath); throw new Error("Secure store write lost authority"); } await fs.rename(tempPath, this.filePath); if (this.options.durable) { const directory = await fs.open(this.options.directory, "r"); try { await directory.sync(); } finally { await directory.close(); } } }
private decryptDocument(raw: string): StoreDocument<T> { let envelope: EncryptedEnvelope; try { envelope = JSON.parse(raw) as EncryptedEnvelope; } catch { throw new Error("Secure store envelope is invalid"); } if (envelope.version !== 1 || envelope.algorithm !== "aes-256-gcm") { throw new Error("Secure store envelope version is unsupported"); } try { const decipher = createDecipheriv("aes-256-gcm", this.options.key, Buffer.from(envelope.iv, "base64")); decipher.setAuthTag(Buffer.from(envelope.tag, "base64")); const plaintext = Buffer.concat([ decipher.update(Buffer.from(envelope.ciphertext, "base64")), decipher.final(), ]); if (plaintext.length > this.options.maxSerializedBytes) throw new Error("document bytes"); const parsed = JSON.parse(plaintext.toString("utf8")) as StoreDocument<unknown>; if (parsed.version !== 1 || !parsed.entries || typeof parsed.entries !== "object" || Array.isArray(parsed.entries)) { throw new Error("document shape"); } if (Object.keys(parsed.entries).length > this.options.maxEntries) throw new Error("entry count"); const entries: Record<string, StoredEntry<T>> = {}; for (const [key, entry] of Object.entries(parsed.entries)) { if (!entry || typeof entry !== "object" || Array.isArray(entry) || !("value" in entry)) throw new Error("entry shape"); const expiresAt = "expiresAt" in entry ? Number(entry.expiresAt) : undefined; if (expiresAt !== undefined && !Number.isSafeInteger(expiresAt)) throw new Error("entry expiry"); entries[key] = { value: this.validate(entry.value), ...(expiresAt !== undefined ? { expiresAt } : {}), }; } return { version: 1, entries }; } catch { throw new Error("Secure store could not be authenticated or decoded"); } }}
export function decodeKey32(encoded: string | undefined, label: string): Buffer { if (!encoded) throw new Error(`${label} is required`); const compact = encoded.trim(); const bytes = Buffer.from(compact, "base64"); if (bytes.length !== 32 || bytes.toString("base64") !== compact) { throw new Error(`${label} must be canonical base64 for exactly 32 bytes`); } return bytes;}
export function opaqueKey(value: string): string { return createHash("sha256").update(value, "utf8").digest("hex");}