import { chmodSync, closeSync, existsSync, lstatSync, mkdirSync, openSync, readFileSync, renameSync, rmSync, writeFileSync, } from "node:fs"; import { createHash, randomUUID } from "node:crypto"; import { homedir } from "node:os"; import { dirname, join } from "node:path"; import type { AttentionView, NodeChange, WorkflowState } from "../domain/types.ts"; import { canonicalJson, snapshotHash } from "../domain/snapshot.ts"; import type { ReconstructedCheckpoint } from "../application/branch-reconstruction.ts"; import { deriveNodeChanges } from "../projection/history.ts"; import { createMermaidProjection } from "../projection/mermaid.ts"; export interface ProjectionRevisionRecord { revision: number; checkpointEntryId: string; parentEntryId?: string; path: string; } export interface CurrentProjectionFile { schemaVersion: 2; generationId: string; sessionId: string; projectCwd: string; branchHeadEntryId?: string; revision: number; generatedAt: string; state: WorkflowState; attention: AttentionView; mermaid: ReturnType; changes: NodeChange[]; revisions: ProjectionRevisionRecord[]; renderError?: string; } export function defaultProjectionRoot(): string { return join(process.env.XDG_STATE_HOME ?? join(homedir(), ".local", "state"), "intent-petri", "sessions"); } function safeSegment(value: string): string { return value.replace(/[^A-Za-z0-9_.-]/g, "_").slice(0, 120) || "unknown"; } export function sessionProjectionDir(sessionId: string, root = defaultProjectionRoot()): string { return join(root, safeSegment(sessionId), "active"); } export interface ActivityProjectionFile { schemaVersion: 1; status: "idle" | "running" | "completed" | "failed" | "closed"; summary: string; activeCount: number; observedAt: string; toolName?: string; elapsedMs?: number; idleMs?: number; } function secureDirectory(path: string): void { if (existsSync(path)) { const stat = lstatSync(path); if (stat.isSymbolicLink() || !stat.isDirectory()) throw new Error(`Unsafe projection directory: ${path}`); } else { mkdirSync(path, { recursive: true, mode: 0o700 }); } chmodSync(path, 0o700); } function atomicWrite(path: string, content: string): void { const directory = dirname(path); secureDirectory(directory); const temporary = `${path}.tmp-${randomUUID()}`; const file = openSync(temporary, "wx", 0o600); try { writeFileSync(file, content, { encoding: "utf8" }); } finally { closeSync(file); } chmodSync(temporary, 0o600); renameSync(temporary, path); chmodSync(path, 0o600); } export function writeActivityProjection( sessionId: string, activity: Omit, root = defaultProjectionRoot(), observedAt = new Date().toISOString(), ): ActivityProjectionFile { const payload: ActivityProjectionFile = { schemaVersion: 1, ...activity, observedAt, }; atomicWrite(join(sessionProjectionDir(sessionId, root), "activity.json"), `${JSON.stringify(payload, null, 2)}\n`); return payload; } function revisionFileName(checkpoint: ReconstructedCheckpoint): string { const hash = snapshotHash(checkpoint.state).replace("sha256:", "").slice(0, 16); return `${String(checkpoint.state.revision).padStart(8, "0")}-${safeSegment(checkpoint.entryId)}-${hash}.json`; } export interface WriteProjectionInput { sessionId: string; projectCwd?: string; branchHeadEntryId?: string; state: WorkflowState; attention: AttentionView; history: ReconstructedCheckpoint[]; root?: string; now?: string; } export function writeProjectionFiles(input: WriteProjectionInput): CurrentProjectionFile { const activeDir = sessionProjectionDir(input.sessionId, input.root); secureDirectory(activeDir); const generationsDir = join(activeDir, "generations"); const revisionCacheDir = join(activeDir, "revision-cache"); secureDirectory(generationsDir); secureDirectory(revisionCacheDir); const changes = deriveNodeChanges(input.history); const revisions: ProjectionRevisionRecord[] = []; for (const checkpoint of input.history) { const fileName = revisionFileName(checkpoint); const relativePath = `revision-cache/${fileName}`; const absolutePath = join(activeDir, relativePath); const record: ProjectionRevisionRecord = { revision: checkpoint.state.revision, checkpointEntryId: checkpoint.entryId, ...(checkpoint.parentEntryId ? { parentEntryId: checkpoint.parentEntryId } : {}), path: relativePath, }; revisions.push(record); if (existsSync(absolutePath)) continue; const projection = createMermaidProjection(checkpoint.state); atomicWrite( absolutePath, `${JSON.stringify( { schemaVersion: 2, sessionId: input.sessionId, ...record, state: checkpoint.state, mermaid: projection, changes: changes.filter((change) => change.revision <= checkpoint.state.revision), }, null, 2, )}\n`, ); } const mermaid = createMermaidProjection(input.state); const stateHash = snapshotHash(input.state).replace("sha256:", ""); const lineageHash = createHash("sha256") .update( canonicalJson({ branchHeadEntryId: input.branchHeadEntryId ?? null, projectCwd: input.projectCwd ?? process.cwd(), stateHash, attention: input.attention, changes, revisions, }), ) .digest("hex"); const generationId = `${String(input.state.revision).padStart(8, "0")}-${stateHash.slice(0, 12)}-${lineageHash.slice(0, 12)}`; const generationDir = join(generationsDir, generationId); secureDirectory(generationDir); const generationCurrentPath = join(generationDir, "current.json"); const generatedAt = existsSync(generationCurrentPath) ? (JSON.parse(readFileSync(generationCurrentPath, "utf8")) as CurrentProjectionFile).generatedAt : input.now ?? new Date().toISOString(); const current: CurrentProjectionFile = { schemaVersion: 2, generationId, sessionId: input.sessionId, projectCwd: input.projectCwd ?? process.cwd(), ...(input.branchHeadEntryId ? { branchHeadEntryId: input.branchHeadEntryId } : {}), revision: input.state.revision, generatedAt, state: input.state, attention: input.attention, mermaid, changes, revisions, }; const events = changes.map((change) => JSON.stringify(change)).join("\n") + (changes.length ? "\n" : ""); const currentJson = `${JSON.stringify(current, null, 2)}\n`; // The generation is immutable. Readers follow active/current.json only after every artifact exists. if (!existsSync(generationCurrentPath)) { atomicWrite(join(generationDir, "current.mmd"), mermaid.full); atomicWrite(join(generationDir, "causal.mmd"), mermaid.causal); atomicWrite(join(generationDir, "local.mmd"), mermaid.local); atomicWrite(join(generationDir, "events.jsonl"), events); atomicWrite(generationCurrentPath, currentJson); } // Portable convenience mirrors; the viewer never uses them as a multi-file transaction. atomicWrite(join(activeDir, "current.mmd"), mermaid.full); atomicWrite(join(activeDir, "causal.mmd"), mermaid.causal); atomicWrite(join(activeDir, "local.mmd"), mermaid.local); atomicWrite(join(activeDir, "events.jsonl"), events); rmSync(join(activeDir, "current.svg"), { force: true }); rmSync(join(activeDir, "render-error.txt"), { force: true }); // This single pointer/payload is the cross-file commit marker and is written last. atomicWrite(join(activeDir, "current.json"), currentJson); return current; }