import { randomUUID } from "node:crypto"; import { EventStore } from "../telemetry/event-store.js"; import { STATES, type ControllerState } from "../controller/state-machine.js"; import { POLICIES, TOPOLOGIES, type BaseEvent, type Policy, type TaskLedger, type Topology } from "../types.js"; import { isTaskLedger, loadLedger } from "./ledger.js"; const TERMINAL = new Set(["COMPLETED", "BLOCKED", "CANCELLED", "BUDGET_EXHAUSTED", "FAILED"]); const READ_ONLY_CHECKPOINTS = new Set(["RECEIVED", "PREFLIGHT", "ROUTED", "SCOUTING", "SWARMING", "DEEP_RUNNING", "WAR_ROOM"]); const MUTATION_MARKERS = new Set(["worktree.created", "workflow.started", "verification.completed", "review.completed"]); const COMMITTED_EFFECTS = new Set(["worktree.integrated", "destructive.command.completed", "external.side-effect.completed"]); export type RecoveryDecision = | { kind: "recoverable"; runId: string; state: ControllerState; topology: Topology; ledger: TaskLedger; sourceSessionId: string; sourceTaskId: string; sourcePolicy: Policy; sourcePolicyVersion: string; sourceConfigHash?: string; sourceExperimentId?: string } | { kind: "blocked"; runId?: string; state?: ControllerState; reason: "event-history-invalid" | "ledger-missing-or-invalid" | "config-version-mismatch" | "config-provenance-mismatch" | "committed-side-effect" | "unsafe-checkpoint" } | { kind: "terminal"; runId: string; state: ControllerState | "STOPPED" }; export interface RecoveredActiveMetadata { identity: { runId: string; taskId: string; sessionId: string; experimentId?: string }; recoveredFromRunId: string; state: "RECEIVED"; previousState: ControllerState; topology: Topology; ledger: TaskLedger; } export type RecoveryContinuation = | { kind: "scout"; task: string; acceptanceCommand: string } | { kind: "ledger-handoff"; facts: TaskLedger["acceptedFacts"] } | { kind: "blocked"; reason: "task-context-missing" }; function controllerState(value: unknown): ControllerState | undefined { return typeof value === "string" && (STATES as readonly string[]).includes(value) ? value as ControllerState : undefined; } function topology(value: unknown): Topology | undefined { return typeof value === "string" && (TOPOLOGIES as readonly string[]).includes(value) ? value as Topology : undefined; } function policy(value: unknown): Policy | undefined { return typeof value === "string" && (POLICIES as readonly string[]).includes(value) ? value as Policy : undefined; } function ordered(events: readonly BaseEvent[]): BaseEvent[] { return events.map((event, index) => ({ event, index })).sort((left, right) => left.event.timestamp.localeCompare(right.event.timestamp) || left.index - right.index).map(({ event }) => event); } function committedSideEffect(event: BaseEvent): boolean { return COMMITTED_EFFECTS.has(event.eventType) || event.destructive === true || event.externalSideEffect === true; } export function classifyInterruptedRun(input: readonly BaseEvent[], ledger: unknown): RecoveryDecision { const events = ordered(input); const first = events[0]; if (!first || !first.runId || events.some((event) => event.runId !== first.runId)) return { kind: "blocked", reason: "event-history-invalid" }; const sourceExperimentId = first.experimentId; if (events.some((event) => event.experimentId !== sourceExperimentId) || (sourceExperimentId !== undefined && !/^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$/.test(sourceExperimentId))) { return { kind: "blocked", runId: first.runId, reason: "event-history-invalid" }; } if (events.some((event) => event.eventType === "recovery.stopped" && typeof event.continuedAsRunId === "string")) return { kind: "terminal", runId: first.runId, state: "STOPPED" }; const terminalEvent = [...events].reverse().find((event) => event.eventType === "result.completed"); const lastTransition = [...events].reverse().find((event) => event.eventType === "route.changed"); const state = controllerState(terminalEvent?.terminalState) ?? controllerState(lastTransition?.state) ?? "RECEIVED"; if ((terminalEvent && !controllerState(terminalEvent.terminalState)) || (lastTransition && !controllerState(lastTransition.state))) return { kind: "blocked", runId: first.runId, reason: "event-history-invalid" }; if (TERMINAL.has(state)) return { kind: "terminal", runId: first.runId, state }; if (!isTaskLedger(ledger)) return { kind: "blocked", runId: first.runId, state, reason: "ledger-missing-or-invalid" }; if (events.some((event) => event.configVersion !== ledger.configVersion)) return { kind: "blocked", runId: first.runId, state, reason: "config-version-mismatch" }; const configHashes = events.flatMap((event) => event.configHash === undefined ? [] : [event.configHash]); const sourceConfigHash = configHashes[0]; if (!first.policyVersion || events.some((event) => event.policyVersion !== first.policyVersion) || configHashes.length === 0 || configHashes.some((hash) => typeof hash !== "string" || !/^[a-f0-9]{64}$/i.test(hash)) || new Set(configHashes).size > 1) { return { kind: "blocked", runId: first.runId, state, reason: "config-provenance-mismatch" }; } if (events.some(committedSideEffect)) return { kind: "blocked", runId: first.runId, state, reason: "committed-side-effect" }; const classifiedEvents = events.filter((event) => event.eventType === "request.classified"); const classified = classifiedEvents.at(-1); const runTopology = topology(classified?.topology); const sourcePolicy = policy(classified?.policy); if (!runTopology || !sourcePolicy || classifiedEvents.some((event) => topology(event.topology) !== runTopology || policy(event.policy) !== sourcePolicy)) { return { kind: "blocked", runId: first.runId, state, reason: "event-history-invalid" }; } const unsafe = !READ_ONLY_CHECKPOINTS.has(state) || runTopology === "direct" || events.some((event) => MUTATION_MARKERS.has(event.eventType)) || ledger.changedFiles.length > 0 || ledger.verificationState !== "pending"; if (unsafe) return { kind: "blocked", runId: first.runId, state, reason: "unsafe-checkpoint" }; return { kind: "recoverable", runId: first.runId, state, topology: runTopology, ledger: structuredClone(ledger), sourceSessionId: first.sessionId, sourceTaskId: first.taskId, sourcePolicy, sourcePolicyVersion: first.policyVersion, ...(typeof sourceConfigHash === "string" ? { sourceConfigHash } : {}), ...(sourceExperimentId ? { sourceExperimentId } : {}) }; } export function materializeRecovery(decision: Extract, identity: { runId: string; sessionId: string }): RecoveredActiveMetadata { if (!identity.runId.trim() || identity.runId === decision.runId) throw new Error("Recovery requires a new Ultra run id"); if (!identity.sessionId.trim()) throw new Error("Recovery requires a session id"); return { identity: { runId: identity.runId, taskId: identity.runId, sessionId: identity.sessionId, ...(decision.sourceExperimentId ? { experimentId: decision.sourceExperimentId } : {}) }, recoveredFromRunId: decision.runId, state: "RECEIVED", previousState: decision.state, topology: decision.topology, ledger: structuredClone(decision.ledger), }; } export function recoveryContinuation(decision: Extract, rawTask?: string, acceptanceCommand?: string): RecoveryContinuation { if (rawTask?.trim() && acceptanceCommand?.trim()) return { kind: "scout", task: rawTask, acceptanceCommand }; if (decision.ledger.acceptedFacts.length) return { kind: "ledger-handoff", facts: structuredClone(decision.ledger.acceptedFacts) }; return { kind: "blocked", reason: "task-context-missing" }; } export async function discoverInterruptedRuns(eventsDir: string, ledgersDir: string): Promise { const grouped = new Map(); for (const event of await new EventStore(eventsDir).all()) grouped.set(event.runId, [...(grouped.get(event.runId) ?? []), event]); const decisions = await Promise.all([...grouped].filter(([, events]) => events.some((event) => event.eventType === "request.received")).map(async ([runId, events]) => { let ledger: unknown; try { ledger = await loadLedger(ledgersDir, runId); } catch { ledger = undefined; } return classifyInterruptedRun(events, ledger); })); return decisions.filter((decision) => decision.kind !== "terminal"); } export async function markInterruptedRunsStopped(eventsDir: string, ledgersDir: string, now = new Date()): Promise { const store = new EventStore(eventsDir); const events = await store.all(); const decisions = await discoverInterruptedRuns(eventsDir, ledgersDir); for (const decision of decisions) { if (!decision.runId || events.some((event) => event.runId === decision.runId && event.eventType === "recovery.stopped" && event.stopReason === "pi-restart")) continue; const previous = [...events].reverse().find((event) => event.runId === decision.runId); if (!previous) continue; await store.append({ schemaVersion: 1, eventId: randomUUID(), timestamp: now.toISOString(), sessionId: previous.sessionId, taskId: previous.taskId, runId: previous.runId, spanId: randomUUID(), parentSpanId: previous.spanId, configVersion: previous.configVersion, policyVersion: previous.policyVersion, ...(previous.experimentId ? { experimentId: previous.experimentId } : {}), profile: previous.profile, eventType: "recovery.stopped", stopReason: "pi-restart", ...(decision.state ? { previousState: decision.state } : {}), }); } return decisions; }