Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
19 kB · 467 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468import { randomUUID } from "node:crypto";import fs from "node:fs/promises";import os from "node:os";import path from "node:path";import YAML from "yaml";import { FilesystemConnector } from "../connectors/filesystem.js";import { sha256, type JsonObject } from "../core/json.js";import { stableKey } from "../core/ids.js";import type { ThoughtEvent } from "../events/types.js";import type { JazzThoughtStore } from "../jazz/store.js";import { MEMORY_MATERIALIZATION_FAILED_EVENT_TYPE, MEMORY_MATERIALIZED_EVENT_TYPE, MEMORY_PROPOSAL_EVENT_TYPE, memoryProposalPayloadSchema, proposalDecisionPayloadSchema, type MemoryMaterializationFailureCode,} from "./contracts.js";import { decisionsForProposal, requireProposal } from "./review.js";
export const STREAM_MEMORY_SOURCE = "filesystem:telegram-agent-context";export const STREAM_MEMORY_PATH = "memory.md";const MATERIALIZER = "trusted-local-memory-materializer@1";const DEFAULT_MAX_BYTES = 1_000_000;
type FailureCode = MemoryMaterializationFailureCode;
export interface MemoryMaterializerOptions { contextRoot: string; source?: typeof STREAM_MEMORY_SOURCE | undefined; maxFileBytes?: number | undefined; beforeRename?: (() => void | Promise<void>) | undefined;}
export interface MemoryMaterializationResult { status: "materialized" | "failed"; event: ThoughtEvent;}
class MaterializationFailure extends Error { constructor(readonly code: FailureCode) { super(code); this.name = "MaterializationFailure"; }}
export async function materializeMemoryDecision( store: JazzThoughtStore, decisionEventId: string, options: MemoryMaterializerOptions,): Promise<MemoryMaterializationResult> { const decision = await store.getEvent(decisionEventId); if (!decision || decision.type !== "stream.thought.agent.proposal.decision") { throw new Error("Memory materialization requires a proposal decision event"); } const decisionPayload = proposalDecisionPayloadSchema.parse(decision.payload); if (decisionPayload.proposalType !== "memory-change") throw new Error("Decision does not target a memory proposal"); const proposal = await requireProposal(store, decisionPayload.proposalEventId); if (proposal.type !== MEMORY_PROPOSAL_EVENT_TYPE) throw new Error("Memory decision targets the wrong proposal kind"); const proposalPayload = memoryProposalPayloadSchema.parse(proposal.payload); const decisions = await decisionsForProposal(store, proposal.id); if (decisions.length !== 1 || decisions[0]!.id !== decision.id) { const failed = await appendFailure(store, proposal, decision, proposalPayload.target.versionId, proposalPayload.target.sha256, "already-decided"); return { status: "failed", event: failed }; } if (decision.parentEventId !== proposal.id || decision.rootEventId !== proposal.rootEventId) { throw new Error("Memory decision lineage is inconsistent"); } const existing = await materializationEvents(store, decision.id); const completed = existing.find((event) => event.type === MEMORY_MATERIALIZED_EVENT_TYPE); if (completed) return { status: "materialized", event: completed }; if (decisionPayload.disposition === "reject") { const failed = await appendFailure(store, proposal, decision, proposalPayload.target.versionId, proposalPayload.target.sha256, "decision-rejected"); return { status: "failed", event: failed }; }
const source = options.source ?? STREAM_MEMORY_SOURCE; if (source !== STREAM_MEMORY_SOURCE || proposalPayload.target.source !== STREAM_MEMORY_SOURCE || proposalPayload.target.path !== STREAM_MEMORY_PATH) { const failed = await appendFailure(store, proposal, decision, proposalPayload.target.versionId, proposalPayload.target.sha256, "root-invalid"); return { status: "failed", event: failed }; }
try { const newContent = await materializeContent(store, proposalPayload, decisionPayload, options); const connector = new FilesystemConnector({ id: STREAM_MEMORY_SOURCE, root: options.contextRoot, privacy: "sensitive", storeContent: true, maxFileBytes: options.maxFileBytes ?? DEFAULT_MAX_BYTES, }); let scan; try { scan = await connector.scan(store); } catch { throw new MaterializationFailure("filesystem-scan-failed"); } const expectedSha256 = sha256(newContent); const current = (await store.listCurrentDocuments(STREAM_MEMORY_SOURCE)).find((item) => ( !item.deleted && item.path === STREAM_MEMORY_PATH )); if (!current || current.documentId !== proposalPayload.target.documentId || current.sha256 !== expectedSha256 || current.versionId !== stableKey("version", STREAM_MEMORY_SOURCE, current.documentId, expectedSha256)) { throw new MaterializationFailure("receipt-mismatch"); } const version = await store.getDocumentVersion(current.versionId); if (!version || version.source !== STREAM_MEMORY_SOURCE || version.documentId !== current.documentId || version.path !== STREAM_MEMORY_PATH || version.sha256 !== expectedSha256 || version.content !== newContent || version.sizeBytes !== Buffer.byteLength(newContent)) { throw new MaterializationFailure("receipt-mismatch"); } const filesystemEvent = scan.events.find((event) => event.payload.versionId === current.versionId) ?? (await store.listEvents({ source: STREAM_MEMORY_SOURCE })).find((event) => event.payload.versionId === current.versionId); if (!filesystemEvent) throw new MaterializationFailure("receipt-mismatch"); const materialized = (await store.appendEvent({ type: MEMORY_MATERIALIZED_EVENT_TYPE, schemaVersion: 1, source: "materializer:agent-context", sourceKind: "system", externalId: decision.id, idempotencyKey: stableKey("memory-proposal-materialized", decision.id, current.versionId), occurredAt: new Date().toISOString(), actor: "operator:local", rootEventId: proposal.rootEventId, parentEventId: decision.id, correlationId: proposal.id, privacy: "sensitive", payload: { proposalEventId: proposal.id, decisionEventId: decision.id, operation: proposalPayload.operation, base: proposalPayload.target, result: { documentId: current.documentId, versionId: current.versionId, sha256: current.sha256, sizeBytes: current.sizeBytes, filesystemEventId: filesystemEvent.id, }, materializedBy: MATERIALIZER, }, createdByRuntime: "thoughtstream-memory-materializer-v1", })).event; return { status: "materialized", event: materialized }; } catch (error) { const code = error instanceof MaterializationFailure ? error.code : "write-failed"; const failed = await appendFailure(store, proposal, decision, proposalPayload.target.versionId, proposalPayload.target.sha256, code); return { status: "failed", event: failed }; }}
async function materializeContent( store: JazzThoughtStore, proposal: ReturnType<typeof memoryProposalPayloadSchema.parse>, decision: ReturnType<typeof proposalDecisionPayloadSchema.parse>, options: MemoryMaterializerOptions,): Promise<string> { const current = (await store.listCurrentDocuments(STREAM_MEMORY_SOURCE)).find((item) => ( !item.deleted && item.path === STREAM_MEMORY_PATH )); if (!current || current.documentId !== proposal.target.documentId || current.versionId !== proposal.target.versionId || current.sha256 !== proposal.target.sha256 || current.contentType !== "text/markdown") { throw new MaterializationFailure("stale-base"); } const base = await store.getDocumentVersion(proposal.target.versionId); if (!base || base.source !== STREAM_MEMORY_SOURCE || base.documentId !== proposal.target.documentId || base.path !== STREAM_MEMORY_PATH || base.contentType !== "text/markdown" || base.sha256 !== proposal.target.sha256 || sha256(base.content) !== base.sha256 || Buffer.byteLength(base.content) !== base.sizeBytes) { throw new MaterializationFailure("base-evidence-invalid"); } const root = path.resolve(options.contextRoot); const rootStat = await fs.lstat(root).catch(() => undefined); if (!rootStat || !rootStat.isDirectory() || rootStat.isSymbolicLink()) throw new MaterializationFailure("root-invalid"); const realRoot = await fs.realpath(root).catch(() => undefined); if (!realRoot || realRoot !== root) throw new MaterializationFailure("root-invalid"); const target = path.join(root, STREAM_MEMORY_PATH); if (path.dirname(target) !== root) throw new MaterializationFailure("target-invalid"); const targetStat = await fs.lstat(target).catch(() => undefined); if (!targetStat) throw new MaterializationFailure("target-invalid"); if (targetStat.isSymbolicLink()) throw new MaterializationFailure("symlink-refused"); if (!targetStat.isFile()) throw new MaterializationFailure("target-invalid"); const realTarget = await fs.realpath(target).catch(() => undefined); if (realTarget !== target) throw new MaterializationFailure("symlink-refused"); const diskContent = normalize(await fs.readFile(target, "utf8")); const baseId = frontmatterId(base.content); if (!baseId) throw new MaterializationFailure("frontmatter-invalid"); const approvedText = decision.disposition === "edit" ? decision.replacementText! : proposal.proposedText; const newContent = proposal.operation === "append" ? appendText(base.content, approvedText) : normalize(approvedText); const newId = frontmatterId(newContent); if (!newId) throw new MaterializationFailure("frontmatter-invalid"); if (newId !== baseId) throw new MaterializationFailure("document-identity-changed"); if (Buffer.byteLength(newContent) > (options.maxFileBytes ?? DEFAULT_MAX_BYTES)) { throw new MaterializationFailure("content-too-large"); } const diskSha256 = sha256(diskContent); const newSha256 = sha256(newContent); if (diskSha256 !== base.sha256 && diskSha256 !== newSha256) throw new MaterializationFailure("stale-base"); if (diskSha256 === newSha256) { if ((targetStat.mode & 0o777) !== 0o600) throw new MaterializationFailure("write-failed"); return newContent; } await atomicReplaceMemory(root, target, base.sha256, newContent, { rootDev: rootStat.dev, rootIno: rootStat.ino, targetDev: targetStat.dev, targetIno: targetStat.ino, }, options.beforeRename); const written = await fs.lstat(target).catch(() => undefined); if (!written || !written.isFile() || written.isSymbolicLink() || (written.mode & 0o777) !== 0o600) { throw new MaterializationFailure("write-failed"); } if (sha256(normalize(await fs.readFile(target, "utf8"))) !== newSha256) { throw new MaterializationFailure("write-failed"); } return newContent;}
async function atomicReplaceMemory( root: string, target: string, expectedBaseSha256: string, content: string, expectedIdentity: { rootDev: number | bigint; rootIno: number | bigint; targetDev: number | bigint; targetIno: number | bigint }, beforeRename: MemoryMaterializerOptions["beforeRename"],): Promise<void> { let lock: MaterializerLock | undefined; let temporary: string | undefined; try { lock = await acquireMaterializerLock(root); temporary = path.join(root, `.memory.md.thoughtstream.${process.pid}.${randomUUID()}.tmp`); const temporaryHandle = await fs.open(temporary, "wx", 0o600); try { await temporaryHandle.writeFile(content, "utf8"); await temporaryHandle.sync(); await temporaryHandle.chmod(0o600); } finally { await temporaryHandle.close(); } await beforeRename?.(); const rootStat = await fs.lstat(root); const targetStat = await fs.lstat(target); if (!rootStat.isDirectory() || rootStat.isSymbolicLink() || targetStat.isSymbolicLink()) { throw new MaterializationFailure("symlink-refused"); } if (!targetStat.isFile()) throw new MaterializationFailure("target-invalid"); if ( rootStat.dev !== expectedIdentity.rootDev || rootStat.ino !== expectedIdentity.rootIno || targetStat.dev !== expectedIdentity.targetDev || targetStat.ino !== expectedIdentity.targetIno ) { throw new MaterializationFailure("stale-base"); } if (sha256(normalize(await fs.readFile(target, "utf8"))) !== expectedBaseSha256) { throw new MaterializationFailure("stale-base"); } await fs.rename(temporary, target); temporary = undefined; const rootHandle = await fs.open(root, "r"); try { await rootHandle.sync(); } finally { await rootHandle.close(); } } finally { if (temporary) await fs.rm(temporary, { force: true }).catch(() => undefined); await lock?.release().catch(() => undefined); }}
interface MaterializerLockRecord { bootId: string; pid: number; processStart: string | null; token: string; createdAt: string;}
interface MaterializerLock { release(): Promise<void>;}
async function acquireMaterializerLock(root: string): Promise<MaterializerLock> { const lockPath = path.join(root, ".thoughtstream-memory-materializer.lock"); const record: MaterializerLockRecord = { bootId: await readBootId(), pid: process.pid, processStart: await readProcessStart(process.pid), token: randomUUID(), createdAt: new Date().toISOString(), }; for (let attempt = 0; attempt < 3; attempt += 1) { try { const handle = await fs.open(lockPath, "wx", 0o600); try { await handle.writeFile(`${JSON.stringify(record)}\n`, "utf8"); await handle.sync(); } finally { await handle.close(); } return { release: async () => { const existing = await readMaterializerLock(lockPath).catch(() => undefined); if (existing?.record.token === record.token && existing.record.pid === record.pid) { await unlinkSameLock(lockPath, existing.dev, existing.ino); } }, }; } catch (error) { if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw new MaterializationFailure("write-failed"); const existing = await readMaterializerLock(lockPath).catch(() => undefined); if (!existing) throw new MaterializationFailure("write-failed"); const currentBootId = await readBootId(); const currentStart = existing.record.bootId === currentBootId ? await readProcessStart(existing.record.pid) : null; const live = existing.record.bootId === currentBootId && ( existing.record.processStart !== null ? currentStart === existing.record.processStart : processIsAlive(existing.record.pid) ); if (live) throw new MaterializationFailure("write-failed"); await unlinkSameLock(lockPath, existing.dev, existing.ino); } } throw new MaterializationFailure("write-failed");}
async function readMaterializerLock(lockPath: string): Promise<{ record: MaterializerLockRecord; dev: number; ino: number;}> { const stat = await fs.lstat(lockPath); if (!stat.isFile() || stat.isSymbolicLink() || stat.size < 2 || stat.size > 4_096) { throw new Error("Materializer lock is invalid"); } const value = JSON.parse(await fs.readFile(lockPath, "utf8")) as Partial<MaterializerLockRecord>; if (typeof value.bootId !== "string" || !Number.isSafeInteger(value.pid) || Number(value.pid) < 1 || (value.processStart !== null && typeof value.processStart !== "string") || typeof value.token !== "string" || value.token.length < 16 || typeof value.createdAt !== "string") { throw new Error("Materializer lock is invalid"); } return { record: value as MaterializerLockRecord, dev: stat.dev, ino: stat.ino, };}
async function unlinkSameLock(lockPath: string, dev: number, ino: number): Promise<void> { const current = await fs.lstat(lockPath).catch(() => undefined); if (!current) return; if (!current.isFile() || current.isSymbolicLink() || current.dev !== dev || current.ino !== ino) { throw new MaterializationFailure("write-failed"); } await fs.unlink(lockPath);}
async function readBootId(): Promise<string> { try { return (await fs.readFile("/proc/sys/kernel/random/boot_id", "utf8")).trim(); } catch { return `${os.hostname()}:${Math.floor(Date.now() - os.uptime() * 1_000)}`; }}
async function readProcessStart(pid: number): Promise<string | null> { try { const stat = await fs.readFile(`/proc/${pid}/stat`, "utf8"); const close = stat.lastIndexOf(")"); if (close < 0) return null; return stat.slice(close + 2).split(" ")[19] ?? null; } catch { return null; }}
function processIsAlive(pid: number): boolean { try { process.kill(pid, 0); return true; } catch (error) { return (error as NodeJS.ErrnoException).code !== "ESRCH"; }}
async function appendFailure( store: JazzThoughtStore, proposal: ThoughtEvent, decision: ThoughtEvent, baseVersionId: string, baseSha256: string, reasonCode: FailureCode,): Promise<ThoughtEvent> { return (await store.appendEvent({ type: MEMORY_MATERIALIZATION_FAILED_EVENT_TYPE, schemaVersion: 1, source: "materializer:agent-context", sourceKind: "system", externalId: decision.id, idempotencyKey: stableKey("memory-proposal-materialization-failed", decision.id, reasonCode), occurredAt: new Date().toISOString(), actor: "operator:local", rootEventId: proposal.rootEventId, parentEventId: decision.id, correlationId: proposal.id, privacy: "sensitive", payload: { proposalEventId: proposal.id, decisionEventId: decision.id, baseVersionId, baseSha256, reasonCode, contentRedacted: true, materializedBy: MATERIALIZER, }, createdByRuntime: "thoughtstream-memory-materializer-v1", })).event;}
async function materializationEvents(store: JazzThoughtStore, decisionEventId: string): Promise<ThoughtEvent[]> { return (await store.listEvents({ types: [MEMORY_MATERIALIZED_EVENT_TYPE, MEMORY_MATERIALIZATION_FAILED_EVENT_TYPE], })).filter((event) => event.payload.decisionEventId === decisionEventId);}
function normalize(value: string): string { return value.replace(/\r\n/g, "\n");}
function appendText(base: string, addition: string): string { return `${normalize(base).trimEnd()}\n\n${normalize(addition).trim()}\n`;}
function frontmatterId(content: string): string | undefined { const normalized = normalize(content); if (!normalized.startsWith("---\n")) return undefined; const end = normalized.indexOf("\n---\n", 4); if (end < 0) return undefined; try { const parsed = YAML.parse(normalized.slice(4, end)); if (!parsed || typeof parsed !== "object" || Array.isArray(parsed) || !("id" in parsed)) return undefined; const id = String((parsed as { id: unknown }).id).trim(); return id || undefined; } catch { return undefined; }}