import { 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) | 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 { 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, decision: ReturnType, options: MemoryMaterializerOptions, ): Promise { 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 { 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; } async function acquireMaterializerLock(root: string): Promise { 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; 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 { 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 { 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 { 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 { 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 { 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; } }