import * as fs from "node:fs"; import * as path from "node:path"; import { deleteRunDirectoryAt } from "./run-storage.ts"; import { isTerminalStatus } from "./types.ts"; import type { ApiRetryRecord, ClusterEvent, ClusterRunResult, ClusterSnapshot, ClusterState, ReviewResult, TaskAttempt, TaskRuntime, UsageStats, WorkerHistoryMessage, WorkerResult, WorkerToolCall, } from "./types.ts"; export interface RunPersistence { runDir: string; snapshotPath: string; writeQueue?: Promise; } function safeName(value: string): string { return value.replace(/[^a-zA-Z0-9._-]+/g, "_").slice(0, 80) || "task"; } function projectPathName(cwd: string): string { const resolvedCwd = path.resolve(cwd); return `--${resolvedCwd.replace(/^[/\\]/, "").replace(/[/\\:]/g, "-")}--`; } async function runsRoot(cwd: string): Promise { const { getAgentDir } = await import("@earendil-works/pi-coding-agent"); return path.join(getAgentDir(), "subagent-cluster", "runs", projectPathName(cwd)); } export async function createRunPersistence(cwd: string, runId: string): Promise { const runDir = path.join(await runsRoot(cwd), runId); await fs.promises.mkdir(runDir, { recursive: true }); return { runDir, snapshotPath: path.join(runDir, "snapshot.json") }; } export async function deleteHistoricalRun(cwd: string, runId: string): Promise { return deleteRunDirectoryAt(await runsRoot(cwd), runId); } function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } function isFiniteNumber(value: unknown): value is number { return typeof value === "number" && Number.isFinite(value); } function isClusterLevel(value: unknown): value is TaskRuntime["level"] { return value === "low" || value === "medium" || value === "high"; } function isClusterStatus(value: unknown): value is ClusterState["status"] { return value === "running" || value === "paused" || value === "completed" || value === "failed" || value === "cancelled"; } function isTaskStatus(value: unknown): value is TaskRuntime["status"] { return ( value === "queued" || value === "running" || value === "reviewing" || value === "paused_for_user" || value === "completed" || value === "timed_out" || value === "failed" || value === "blocked" || value === "cancelled" ); } function isNonEmptyString(value: unknown): value is string { return typeof value === "string" && value.trim() !== ""; } function isOptionalString(record: Record, key: string): boolean { return record[key] === undefined || typeof record[key] === "string"; } function isOptionalBoolean(record: Record, key: string): boolean { return record[key] === undefined || typeof record[key] === "boolean"; } function isOptionalFiniteNumber(record: Record, key: string): boolean { return record[key] === undefined || isFiniteNumber(record[key]); } function isUsageStats(value: unknown, requireComplete = false): value is UsageStats { if (!isRecord(value)) return false; return ( isFiniteNumber(value.input) && isFiniteNumber(value.output) && isFiniteNumber(value.cacheRead) && isFiniteNumber(value.cacheWrite) && (requireComplete ? isFiniteNumber(value.totalTokens) && value.totalTokens === value.input + value.output + value.cacheRead + value.cacheWrite : isOptionalFiniteNumber(value, "totalTokens")) && isFiniteNumber(value.cost) && isFiniteNumber(value.contextTokens) && isFiniteNumber(value.turns) ); } function isWorkerToolCall(value: unknown): value is WorkerToolCall { return isRecord(value) && typeof value.name === "string" && isRecord(value.args); } function isWorkerHistoryUsage(value: unknown): boolean { if (!isRecord(value)) return false; if ( !isOptionalFiniteNumber(value, "input") || !isOptionalFiniteNumber(value, "output") || !isOptionalFiniteNumber(value, "cacheRead") || !isOptionalFiniteNumber(value, "cacheWrite") || !isOptionalFiniteNumber(value, "reasoning") || !isOptionalFiniteNumber(value, "totalTokens") ) return false; if (value.cost === undefined) return true; return isRecord(value.cost) && isOptionalFiniteNumber(value.cost, "total"); } function isWorkerHistoryMessage(value: unknown): value is WorkerHistoryMessage { if (!isRecord(value)) return false; return ( isOptionalString(value, "role") && isOptionalString(value, "toolName") && isOptionalString(value, "toolCallId") && isOptionalBoolean(value, "isError") && (value.usage === undefined || isWorkerHistoryUsage(value.usage)) && isOptionalString(value, "model") && isOptionalString(value, "stopReason") && isOptionalString(value, "errorMessage") ); } function isApiRetryRecord(value: unknown): value is ApiRetryRecord { return ( isRecord(value) && isFiniteNumber(value.timestamp) && isFiniteNumber(value.attempt) && Number.isInteger(value.attempt) && value.attempt > 0 && isFiniteNumber(value.delayMs) && Number.isInteger(value.delayMs) && value.delayMs >= 0 && typeof value.errorSummary === "string" ); } function isWorkerResult(value: unknown, requireComplete = false): value is WorkerResult { if (!isRecord(value)) return false; return ( isFiniteNumber(value.exitCode) && typeof value.output === "string" && typeof value.stderr === "string" && Array.isArray(value.toolCalls) && value.toolCalls.every(isWorkerToolCall) && isOptionalString(value, "model") && isOptionalString(value, "stopReason") && isOptionalString(value, "errorMessage") && isUsageStats(value.usage, requireComplete) && isOptionalString(value, "outputPath") && isOptionalBoolean(value, "outputTruncated") && Array.isArray(value.history) && value.history.every(isWorkerHistoryMessage) && (requireComplete ? Array.isArray(value.apiRetries) && value.apiRetries.every(isApiRetryRecord) : value.apiRetries === undefined || (Array.isArray(value.apiRetries) && value.apiRetries.every(isApiRetryRecord))) ); } function isReviewResult(value: unknown): value is ReviewResult { if (!isRecord(value)) return false; return ( (value.decision === "pass" || value.decision === "retry" || value.decision === "escalate" || value.decision === "ask_user" || value.decision === "timeout" || value.decision === "error") && typeof value.reason === "string" && Array.isArray(value.missingCriteria) && value.missingCriteria.every((criterion) => typeof criterion === "string") && typeof value.nextInstruction === "string" && isOptionalString(value, "rawOutput") ); } function isTaskAttempt(value: unknown, requireComplete = false): value is TaskAttempt { if (!isRecord(value)) return false; return ( isClusterLevel(value.level) && isFiniteNumber(value.attempt) && isOptionalString(value, "model") && isFiniteNumber(value.startedAt) && isOptionalFiniteNumber(value, "finishedAt") && (value.worker === undefined || isWorkerResult(value.worker, requireComplete)) && (value.reviewer === undefined || isWorkerResult(value.reviewer, requireComplete)) && (value.review === undefined || isReviewResult(value.review)) ); } function isTaskRuntime(value: unknown, requireComplete = false): value is TaskRuntime { if (!isRecord(value)) return false; return ( typeof value.id === "string" && typeof value.title === "string" && isNonEmptyString(value.taskType) && typeof value.task === "string" && Array.isArray(value.acceptanceCriteria) && value.acceptanceCriteria.every((criterion) => typeof criterion === "string") && Array.isArray(value.dependsOn) && value.dependsOn.every((dependency) => typeof dependency === "string") && isOptionalString(value, "cwd") && isClusterLevel(value.requestedLevel) && isClusterLevel(value.initialLevel) && isNonEmptyString(value.levelSelectionReason) && isClusterLevel(value.level) && isTaskStatus(value.status) && Array.isArray(value.attempts) && value.attempts.every((attempt) => isTaskAttempt(attempt, requireComplete)) && isOptionalString(value, "output") && isOptionalString(value, "outputPath") && (value.review === undefined || isReviewResult(value.review)) && isOptionalString(value, "error") && isOptionalFiniteNumber(value, "startedAt") && isOptionalFiniteNumber(value, "finishedAt") ); } function isClusterEvent(value: unknown): value is ClusterEvent { if (!isRecord(value)) return false; return ( isFiniteNumber(value.timestamp) && isOptionalString(value, "taskId") && (value.kind === "run" || value.kind === "worker" || value.kind === "review" || value.kind === "state" || value.kind === "control" || value.kind === "error") && typeof value.message === "string" ); } function isPausePeriod(value: unknown): boolean { return isRecord(value) && isFiniteNumber(value.startedAt) && isOptionalFiniteNumber(value, "finishedAt"); } function isClusterState(value: unknown): value is ClusterState { if (!isRecord(value)) return false; const requireComplete = value.executionDataVersion === 1; return ( typeof value.runId === "string" && typeof value.goal === "string" && typeof value.cwd === "string" && (value.executionDataVersion === undefined || value.executionDataVersion === 1) && isClusterStatus(value.status) && Array.isArray(value.tasks) && value.tasks.every((task) => isTaskRuntime(task, requireComplete)) && Array.isArray(value.events) && value.events.every(isClusterEvent) && isFiniteNumber(value.startedAt) && isOptionalFiniteNumber(value, "finishedAt") && typeof value.paused === "boolean" && (value.pausePeriods === undefined || (Array.isArray(value.pausePeriods) && value.pausePeriods.every(isPausePeriod))) ); } function isClusterRunResult(value: unknown, requireComplete = false): value is ClusterRunResult { if (!isRecord(value)) return false; return ( typeof value.runId === "string" && isClusterStatus(value.status) && typeof value.summary === "string" && Array.isArray(value.tasks) && value.tasks.every((task) => isTaskRuntime(task, requireComplete)) && isUsageStats(value.usage, requireComplete) ); } export function isClusterSnapshot(value: unknown): value is ClusterSnapshot { if (!isRecord(value) || value.version !== 2 || !isClusterState(value.state)) return false; const requireComplete = isRecord(value.state) && value.state.executionDataVersion === 1; return value.result === undefined || isClusterRunResult(value.result, requireComplete); } function snapshotTime(snapshot: ClusterSnapshot): number { return snapshot.state.finishedAt ?? snapshot.state.startedAt; } function isNewerSnapshot( candidate: ClusterSnapshot, candidateEntry: string, existing: ClusterSnapshot, existingEntry: string, ): boolean { const candidateTime = snapshotTime(candidate); const existingTime = snapshotTime(existing); if (candidateTime !== existingTime) return candidateTime > existingTime; if ((candidate.state.finishedAt !== undefined) !== (existing.state.finishedAt !== undefined)) { return candidate.state.finishedAt !== undefined; } if ((candidate.result !== undefined) !== (existing.result !== undefined)) return candidate.result !== undefined; return candidateEntry > existingEntry; } export async function loadHistoricalSnapshots(cwd = process.cwd()): Promise { return loadHistoricalSnapshotsFromRunsDirectory(await runsRoot(cwd)); } export async function loadHistoricalSnapshotsFromRunsDirectory(runsDir: string): Promise { let entries: fs.Dirent[]; try { entries = await fs.promises.readdir(runsDir, { withFileTypes: true }); } catch { return []; } const snapshots = await Promise.all( entries.map(async (entry): Promise<{ entryName: string; snapshot?: ClusterSnapshot }> => { if (!entry.isDirectory()) return { entryName: entry.name }; const directory = path.join(runsDir, entry.name); try { const contents = await fs.promises.readFile(path.join(directory, "snapshot.json"), "utf8"); const value: unknown = JSON.parse(contents); if (isRecord(value) && value.version === 1) { await fs.promises.rm(directory, { recursive: true, force: true }); return { entryName: entry.name }; } return { entryName: entry.name, snapshot: isClusterSnapshot(value) ? value : undefined }; } catch { return { entryName: entry.name }; } }), ); const uniqueSnapshots = new Map(); for (const item of snapshots) { if (!item.snapshot) continue; const existing = uniqueSnapshots.get(item.snapshot.state.runId); if (!existing || isNewerSnapshot(item.snapshot, item.entryName, existing.snapshot, existing.entryName)) { uniqueSnapshots.set(item.snapshot.state.runId, { entryName: item.entryName, snapshot: item.snapshot }); } } return [...uniqueSnapshots.values()] .map((item) => item.snapshot) .sort((left, right) => { const timeDifference = snapshotTime(right) - snapshotTime(left); return timeDifference || right.state.runId.localeCompare(left.state.runId); }); } export function persistSnapshot(persistence: RunPersistence, state: ClusterState, result?: ClusterSnapshot["result"]): Promise { const snapshot: ClusterSnapshot = { version: 2, state: cloneState(state), result }; const write = async () => { const temporaryPath = `${persistence.snapshotPath}.${Date.now()}.tmp`; await fs.promises.writeFile(temporaryPath, JSON.stringify(snapshot, null, 2), "utf8"); await fs.promises.rename(temporaryPath, persistence.snapshotPath); }; persistence.writeQueue = (persistence.writeQueue ?? Promise.resolve()).then(write, write); return persistence.writeQueue; } export async function persistTaskOutput( persistence: RunPersistence, task: TaskRuntime, output: string, ): Promise { const outputPath = path.join(persistence.runDir, `${safeName(task.id)}.md`); await fs.promises.writeFile(outputPath, output, "utf8"); return outputPath; } export function summarizeState(state: ClusterState): { completedCount: number; failedCount: number; pausedCount: number; } { return { completedCount: state.tasks.filter((task) => task.status === "completed").length, failedCount: state.tasks.filter((task) => ["failed", "blocked", "cancelled", "timed_out"].includes(task.status)).length, pausedCount: state.tasks.filter((task) => task.status === "paused_for_user").length, }; } export function cloneState(state: ClusterState): ClusterState { return JSON.parse(JSON.stringify(state)) as ClusterState; } export function finalizeInterruptedSnapshot(snapshot: ClusterSnapshot, finishedAt = Date.now()): ClusterSnapshot { const state = cloneState(snapshot.state); const interrupted = state.status === "running" || state.status === "paused" || state.tasks.some((task) => !isTerminalStatus(task.status)); if (!interrupted) return snapshot; state.status = "cancelled"; if (state.paused) { const period = state.pausePeriods?.findLast((item) => item.finishedAt === undefined); if (period) period.finishedAt = finishedAt; } state.paused = false; state.finishedAt = finishedAt; for (const task of state.tasks) { if (isTerminalStatus(task.status)) continue; task.status = "cancelled"; task.finishedAt ??= finishedAt; } const result = snapshot.result ? { ...snapshot.result, status: "cancelled" as const, tasks: state.tasks, } : undefined; return { version: 2, state, result }; }