import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; import { getCrewEnv } from "../config/env-vars.ts"; import { errors } from "../errors.ts"; import { atomicWriteFile, atomicWriteJson } from "../state/atomic-write.ts"; import { type ActiveRunRegistryEntry, activeRunEntries } from "../state/stores/active-run-registry.ts"; import { getCurrentPlanRecord } from "../state/stores/plan-store.ts"; import { loadManifestWithRecovery, loadTasksWithRecovery, saveRunManifest } from "../state/stores/state-store.ts"; import type { TeamRunManifest, TeamTaskState } from "../state/types.ts"; import { logInternalError } from "../utils/internal-error.ts"; import { findRunStateDir } from "../utils/paths.ts"; import { recordFromTask, upsertCrewAgent } from "./crew-agent-records.ts"; import { checkProcessLiveness } from "./process-status.ts"; /** Age threshold for orphaned temp directory cleanup: 1 hour. */ const ORPHAN_TEMP_DIR_AGE_THRESHOLD_MS = 60 * 60 * 1000; /** A `.cleanup-in-progress` sentinel older than this is abandoned — its owner * died mid-cleanup (session kill between sentinel-create and dir-delete) and * without reclamation the workspace is frozen out of the cleanup cycle * FOREVER (EEXIST → skip, no TTL). Live 2026-09-26: 42 frozen workspaces. */ const SENTINEL_STALE_RECLAIM_MS = 10 * 60 * 1000; /** Assumed production cadence of the temp-workspace reconcile tick (60s — * see observability.ts tempReconcileTimer). Only used to derive the * stateless round-robin batch index; a different real cadence still cycles * through every batch, just with different pacing. */ const ORPHAN_RECONCILE_TICK_MS = 60_000; /** Defense-in-depth: cap the number of /tmp/pi-crew-* entries processed per * reconcile tick. With a few thousand accumulated dirs, processing them all * synchronously can block the main thread for many seconds, causing the * terminal to appear hung. We process in batches; the rest are handled on * subsequent ticks. */ const ORPHAN_TEMP_SCAN_BATCH_SIZE = 50; /** * Result of reconciling a single stale run. */ export interface ReconcileResult { runId: string; /** What was found and what action was taken */ verdict: "healthy" | "blocked_awaiting_approval" | "waiting_answer" | "result_exists" | "pid_dead" | "pid_alive_stale" | "no_status"; /** Whether repair was applied */ repaired: boolean; /** Human-readable detail */ detail: string; /** Repaired task state, returned to a locked caller for persistence. */ repairedTasks?: TeamTaskState[]; /** NEW-1 (SDD 2026-09-30 WI-2): what the caller should PERSIST — the FULL * task array, including tasks the repair did not touch. Only the * individual-stale branch sets it (its `repairedTasks` is filtered down to * the stale subset, which must NOT be full-overwrite-persisted or healthy * tasks vanish from tasks.json). Callers fall back to `repairedTasks` * when absent — every other branch already returns the full array there. */ persistTasks?: TeamTaskState[]; } /** * Is this run intentionally waiting for human plan approval? * * Such runs are NOT stale even if their owning session died or their async PID * is no longer live — they are blocked on a human decision, not a crash. Crash * recovery and stale reconciliation must preserve them rather than mark them * failed or orphan-cancel them. See PR #32 (gustavo-pelissaro) for the * original analysis of this failure mode. */ export function isPlanApprovalPending(manifest: TeamRunManifest): boolean { return manifest.status === "blocked" && manifest.planApproval?.required === true && manifest.planApproval.status === "pending"; } /** T2/R4 (ADR-4 §2 reader migration): plan-record-first variant — the current * PlanRecord's approval (when present) is authoritative; pre-v2 runs fall * back to the manifest gate above. The `blocked + pending → protected` * invariant is preserved on BOTH paths. */ export function isPlanApprovalPendingEffective(manifest: TeamRunManifest): boolean { if (isPlanApprovalPending(manifest)) return true; if (manifest.status !== "blocked") return false; return getCurrentPlanRecord(manifest)?.approval?.status === "pending"; } /** WP-2/R2 (ADR-0 item 9, docs/decisions/2026-08-17-waiting-producer-ask.md): * waiting TTL — a manifest.waitState.askedAt younger than this marks an * intentional ask-wait; older marks a leaked park marker (worker died without * resolving the ask) and the run loses reconciler protection. Leak guard. */ const WAITING_TTL_MS = 24 * 60 * 60 * 1000; /** * Is this run parked in the ask tool awaiting a human answer (WP-2/R2)? * * True when waitState is set and its askedAt is within waitingTtl. The parked * worker is blocked polling the mailbox for a response — not crashed — and it * typically stops heartbeating while parked, so the run must be protected at * run level (before the no-PID/heartbeat staleness phases) exactly like * plan-approval runs. askedAt exactly at the TTL boundary still counts as * pending (<=, fail-safe toward preservation); a future askedAt (clock skew) * also passes. TTL-EXPIRED waitState returns false — the marker leaked and the * normal staleness repair paths apply (leak guard). */ export function isWaitAnswerPending(manifest: TeamRunManifest, now = Date.now()): boolean { const askedAtMs = manifest.waitState?.askedAt ? new Date(manifest.waitState.askedAt).getTime() : Number.NaN; return Number.isFinite(askedAtMs) && now - askedAtMs <= WAITING_TTL_MS; } /** * Generalized intentional-wait predicate (ADR-0 item 9, extended by ADR-4 §2): * the run is blocked on a human decision — plan approval (record-first with * manifest fallback — pre-v2 semantics preserved on the fallback path) or a * pending ask answer within waitingTtl. Such runs are not crashes and must not * be stale-repaired or cancelled by reconciliation. */ export function isIntentionalWait(manifest: TeamRunManifest, now = Date.now()): boolean { return isPlanApprovalPendingEffective(manifest) || isWaitAnswerPending(manifest, now); } const STALE_ALIVE_PID_MS = 24 * 60 * 60 * 1000; // 24 hours const ACTIVE_EVIDENCE_TTL_MS = 5 * 60 * 1000; /** For no-PID runs, repair when ALL running tasks have heartbeat stale beyond this threshold. */ const NO_PID_HEARTBEAT_STALE_MS = 5 * 60 * 1000; // 5 minutes — same as heartbeat-gradient deadMs /** * Phase 1: Check if a result file already exists for the run. * If so, the run completed but status wasn't updated — repair it. */ function checkResultFile(manifest: TeamRunManifest, tasks: TeamTaskState[]): { found: boolean; repaired: boolean } { // Check if all tasks already have terminal status (result was written but manifest wasn't updated) const allTerminal = tasks.length > 0 && tasks.every( (t) => t.status === "completed" || t.status === "failed" || t.status === "cancelled" || t.status === "skipped" || t.status === "needs_attention", ); if (allTerminal) { // All tasks are terminal but manifest status was not updated — repair it. // Derive the run status from task outcomes instead of blindly marking // "completed": a run where every task failed/cancelled is NOT a success. const hasFailed = tasks.some((t) => t.status === "failed" || t.status === "needs_attention"); const onlyCancelledOrSkipped = tasks.every((t) => t.status === "cancelled" || t.status === "skipped"); manifest.status = hasFailed ? "failed" : onlyCancelledOrSkipped ? "cancelled" : "completed"; // Persist manifest status change immediately to make checkResultFile self-contained. // NOTE (OPT-02): Remains sync — reconcileStaleRun → checkResultFile is called from // withRunLockSync in crash-recovery.ts and from sync reconcileOrphanedTempWorkspaces. // Converting to async would require cascading through all callers including sync UI hooks. saveRunManifest(manifest); // Sync agent records even when tasks are already terminal // (e.g., a previous reconcile fixed tasks but crashed before updating agents) for (const task of tasks) { try { upsertCrewAgent(manifest, recordFromTask(manifest, task, "scaffold")); } catch { /* non-critical */ } } return { found: true, repaired: false }; } return { found: false, repaired: false }; } /** * Get process start time in milliseconds since boot from /proc//stat. * Returns undefined if the process is gone or /proc is unavailable. * * Platform limitation: Relies on Linux's /proc filesystem. On macOS/Windows, * startTime will be undefined and PID reuse detection is weaker. * * The start time is in the 22nd field (index 21 after comm) of /proc//stat. */ function getProcessStartTime(pid: number): number | undefined { try { const stat = fs.readFileSync(`/proc/${pid}/stat`, "utf-8"); const lastParen = stat.lastIndexOf(")"); if (lastParen === -1) return undefined; const fieldsAfterComm = stat .slice(lastParen + 1) .trim() .split(/\s+/); // starttime is at index 19 (the 20th field after comm) const startTimeClockTicks = Number(fieldsAfterComm[19]); if (!Number.isFinite(startTimeClockTicks)) return undefined; // Convert clock ticks to ms using ~100 ticks/sec (CLK_TCK). // The absolute value matters less than uniqueness per PID lifecycle. return Math.floor(startTimeClockTicks * 10); } catch { return undefined; } } /** * Phase 2: Check PID liveness. * Uses process.kill(pid, 0) for the authoritative check, but also verifies * startTime before and after to detect PID recycling. If the OS recycles * the PID between the kill(0) call and the repair action, a newly spawned * unrelated process could be incorrectly identified as the async worker. * The heartbeat file is used as corroborating evidence when the process * appears dead per kill(0). */ function checkPidLiveness( pid: number | undefined, stateRoot?: string, ): { alive: boolean; detail: string; } { if (pid === undefined || !Number.isInteger(pid) || pid <= 0) { return { alive: false, detail: "no pid recorded" }; } // Capture startTime before kill(0) to detect PID recycling. const startTimeBefore = getProcessStartTime(pid); const liveness = checkProcessLiveness(pid); // Re-verify startTime after kill(0) to close the TOCTOU window. // If startTime changed, the PID was recycled to a different process. const startTimeAfter = getProcessStartTime(pid); if (startTimeBefore !== undefined && startTimeAfter !== undefined && startTimeBefore !== startTimeAfter) { // PID was recycled — treat as dead to avoid acting on wrong process. return { alive: false, detail: "pid_recycled" }; } // If process is alive per kill(0), we're done. if (liveness.alive) return { alive: true, detail: liveness.detail }; // Process is dead per kill(0). ESRCH is AUTHORITATIVE — the PID does not exist. // A frozen/zombie process still has a PID (kill(0) returns 0, or EPERM if we lack // permission to signal it), so reaching here with alive=false means the process // is genuinely gone. PID recycling is already handled above via the startTime // TOCTOU check, so the heartbeat file must NOT override the dead verdict. // // HISTORY: a <5min-old heartbeat.json previously OVERRRODE this verdict (returned // alive:true "process dead but heartbeat Xs old"), delaying repair by ~5min. // Combined with the 60s reconcile cadence this produced a ~300-360s detection lag // for dead workers (observed 353s in production). The heartbeat is now used only // to enrich the detail string, never to override the authoritative kill(0) verdict. let heartbeatNote = ""; if (stateRoot) { try { const heartbeatPath = path.join(stateRoot, "heartbeat.json"); if (fs.existsSync(heartbeatPath)) { const hb = JSON.parse(fs.readFileSync(heartbeatPath, "utf-8")) as { pid?: number; at?: number }; if (hb?.pid === pid && hb?.at) { heartbeatNote = ` (heartbeat was ${Math.round((Date.now() - hb.at) / 1000)}s old)`; } } } catch { /* ignore — best-effort */ } } return { alive: false, detail: `${liveness.detail}${heartbeatNote}` }; } /** * Phase 3: For dead PIDs, repair immediately. * For alive PIDs, only mark stale if status hasn't updated in STALE_ALIVE_PID_MS. */ function evaluateStaleness(manifest: TeamRunManifest, pidAlive: boolean, now: number): { stale: boolean; reason: string } { if (!pidAlive) { return { stale: true, reason: "pid_dead" }; } const updatedAt = new Date(manifest.updatedAt).getTime(); if (!Number.isFinite(updatedAt)) { return { stale: false, reason: "updated_at_invalid" }; } if (now - updatedAt > STALE_ALIVE_PID_MS) { return { stale: true, reason: `alive_but_stale_${Math.round((now - updatedAt) / 3600_000)}h`, }; } return { stale: false, reason: "alive_and_recent" }; } function hasRecentActiveEvidence(tasks: TeamTaskState[], now: number): boolean { return tasks.some((task) => { if (task.status !== "running" && task.status !== "waiting") return false; const heartbeatAt = task.heartbeat?.lastSeenAt ? new Date(task.heartbeat.lastSeenAt).getTime() : Number.NaN; if (task.heartbeat?.alive !== false && Number.isFinite(heartbeatAt) && now - heartbeatAt <= ACTIVE_EVIDENCE_TTL_MS) return true; const activityAt = task.agentProgress?.lastActivityAt ? new Date(task.agentProgress.lastActivityAt).getTime() : Number.NaN; return Number.isFinite(activityAt) && now - activityAt <= ACTIVE_EVIDENCE_TTL_MS; }); } /** * Check if a single task is stale (heartbeat AND activity both stale beyond * NO_PID_HEARTBEAT_STALE_MS). Tasks with no heartbeat AND no agent progress * are considered NOT stale (they may be newly spawned and haven't reported yet). */ function isTaskHeartbeatStale(task: TeamTaskState, now: number): boolean { const heartbeatAt = task.heartbeat?.lastSeenAt ? new Date(task.heartbeat.lastSeenAt).getTime() : Number.NaN; const activityAt = task.agentProgress?.lastActivityAt ? new Date(task.agentProgress.lastActivityAt).getTime() : Number.NaN; // If no heartbeat AND no activity, we can't determine staleness — not stale if (!Number.isFinite(heartbeatAt) && !Number.isFinite(activityAt)) return false; // Compute elapsed from both sources and use the fresher (minimum) one, // mirroring heartbeat-watcher logic: if either source has recent activity, // the task is not stale even if the other is stale. const heartbeatAge = Number.isFinite(heartbeatAt) ? now - heartbeatAt : Number.MAX_SAFE_INTEGER; const activityAge = Number.isFinite(activityAt) ? now - activityAt : Number.MAX_SAFE_INTEGER; const elapsed = Math.min(heartbeatAge, activityAge); if (elapsed <= NO_PID_HEARTBEAT_STALE_MS) return false; // F1 part 2 (battery 2026-09-10, live evidence run team_20260910171307): a // surface worker's recorder only flushes at turn boundaries — one slow LLM // turn (observed: 6.7min) freezes lastSeenAt/lastActivityAt while the worker // process is alive and mid-turn. The heartbeat-watcher already has a // PID-liveness gate for exactly this ("alive → downgrade dead→stale"); the // reconciler MUST apply the same gate before repairing, or it kills healthy // runs at the 5-minute mark mid-turn. A truly dead pid still repairs as before. const taskPid = task.heartbeat?.pid ?? task.checkpoint?.childPid; const pidAlive = taskPid ? checkProcessLiveness(taskPid).alive : false; if (taskPid && pidAlive) return false; if (getCrewEnv("PI_CREW_DEBUG_STALE") === "1") { // F1 forensic (battery 2026-09-10): sidecar log of every STALE verdict so // live repros can show exactly what the reconciler saw (pid present? alive? // elapsed?) — the reconciler may run in ANY host process, hence a fixed // sidecar path instead of console. try { fs.appendFileSync( "/tmp/pi-crew-f1-debug.log", `${JSON.stringify({ ts: new Date().toISOString(), taskId: task.id, status: task.status, taskPid, pidAlive, heartbeatAge, activityAge, hbLastSeen: task.heartbeat?.lastSeenAt })}\n`, ); } catch { /* best-effort forensic log */ } } return true; } /** * For no-PID runs: check if ALL running tasks have heartbeats stale beyond * the no-PID heartbeat threshold, and return the list of individually stale * task IDs. This detects zombie tasks where the worker process died but no * PID was recorded (e.g. live-session /tmp/ workspaces). * Merged from the former allRunningTasksHeartbeatStale and * findIndividuallyStaleTaskIds to avoid redundant duplicate checks. */ function getRunningTaskStaleness(tasks: TeamTaskState[], now: number): { allStale: boolean; staleTaskIds: string[] } { const runningTasks = tasks.filter((t) => t.status === "running" || t.status === "waiting"); if (runningTasks.length === 0) return { allStale: false, staleTaskIds: [] }; const staleTaskIds: string[] = []; for (const task of runningTasks) { if (isTaskHeartbeatStale(task, now)) { staleTaskIds.push(task.id); } } return { allStale: staleTaskIds.length === runningTasks.length, staleTaskIds, }; } /** * Repair a stale run by marking it as failed and cancelling running tasks. */ /** * E3/E1 (Round 15): Build a human-actionable error string for a stale-reconciled * task. Explains WHY the run was marked stale (the detected reason) and gives * concrete remediation, instead of the bare 'Stale run reconciled: '. * Now returns a structured CrewError (E012) so callers also get a machine- * readable code + help hint; `.message` carries the same rich text as before. */ function buildStaleReconcileError(task: TeamTaskState, reason: string): Error { const heartbeatAgeSeconds = task.heartbeat?.lastSeenAt ? Math.round((Date.now() - new Date(task.heartbeat.lastSeenAt).getTime()) / 1000) : undefined; return errors.runStale(reason, heartbeatAgeSeconds); } /** * G24 (SDD-4 WI-2): find a LIVE active-run-registry entry for this runId. * activeRunEntries() already applies the full liveness filter (terminal status, * dead async PID, >30min-stale non-async manifests, symlink safety), so an * entry found here means some live session or runner claims the run RIGHT NOW. * Read-only consult (register/unregister are never imported here). Fail-open: * on an unexpected registry error return undefined so normal staleness * reconciliation proceeds — a broken registry must not shield dead runs from * repair forever. */ function findLiveRegistryEntry(runId: string): ActiveRunRegistryEntry | undefined { try { return activeRunEntries().find((entry) => entry.runId === runId); } catch (err) { logInternalError("stale-reconciler", new Error(`active-run registry consult failed for ${runId}: ${err}`), undefined, "warn"); return undefined; } } function repairStaleRun(manifest: TeamRunManifest, tasks: TeamTaskState[], reason: string): TeamTaskState[] { const now = new Date().toISOString(); const repairedTasks = tasks.map((task) => { if (task.status === "running" || task.status === "queued" || task.status === "waiting") { return { ...task, status: "cancelled" as const, finishedAt: now, // E3/E1 (Round 15): structured CrewError (E012) with code + help hint. error: buildStaleReconcileError(task, reason).message, }; } return task; }); // Update agent records so widget sees cancelled status immediately for (const task of repairedTasks) { try { upsertCrewAgent(manifest, recordFromTask(manifest, task, "scaffold")); } catch { /* non-critical */ } } return repairedTasks; } /** * Three-phase stale run reconciliation. * * 1. Check if result already exists → use it * 2. Check PID liveness * 3. Dead PID → repair immediately; alive PID → only fail if stale > 24h * * NOTE: Callers must provide locking via withRunLock/withRunLockSync when * calling from contexts where concurrent reconciliation of the same runId * could occur (e.g., the auto-repair timer). The crash-recovery.ts caller * already provides this. The reconcileOrphanedTempWorkspaces caller handles * /tmp workspaces where concurrent access is a known benign race (separate * dirs, low consequence of redundant repair). */ export function reconcileStaleRun(manifest: TeamRunManifest, tasks: TeamTaskState[], now = Date.now()): ReconcileResult { const runId = manifest.runId; // Preserve runs intentionally blocked on human plan approval. These are not // crashes even if the owning PID is gone — they are waiting for a decision. // Must short-circuit before Phase 1 (result check) and Phase 2 (PID liveness). // ADR-4 §2: record-first predicate — a PlanRecord approval is authoritative, // the manifest field remains the pre-v2 fallback (never dropped). if (isPlanApprovalPendingEffective(manifest)) { return { runId, verdict: "blocked_awaiting_approval", repaired: false, detail: "Plan approval is pending; blocked run is intentionally waiting and must not be stale-repaired", }; } // WP-2/R2 (ADR-0 item 9): a run parked in the ask tool with a fresh // waitState.askedAt is intentionally waiting on a human answer — not stale. // Short-circuits before Phase 1/2 like plan-approval above: a parked worker // does not heartbeat, so the no-PID/heartbeat staleness phases would // otherwise cancel it long before its (server-clamped <=1h) ask deadline. // TTL-expired waitState (>24h) falls through — the park marker leaked // (worker died without resolving) and normal staleness repair applies. if (isWaitAnswerPending(manifest, now)) { return { runId, verdict: "waiting_answer", repaired: false, detail: "Ask answer pending (manifest.waitState.askedAt within waiting TTL); run is intentionally waiting and must not be stale-repaired", }; } // G24 (SDD-4 WI-2, plan §2 upgrade-2026-09-29): foreign-LIVE skip. The // active-run registry is the cross-session source of truth for "a live // session/runner owns this run right now" (registerActiveRun at dispatch; // entries self-expire via PID liveness + the 30min freshness horizon). Until // now ownership was only checked as "run of the CURRENT session" // (crash-recovery.ts ownerSessionId filter), so a child session — or the /tmp // orphan scan, which has no session identity at all — could repair (cancel) // a run whose workers were alive but heartbeat-frozen (G25 silent turn: // 12m27s measured on a live worker). Must run BEFORE Phase 1/2/3: no phase // below may act on a run a live owner still owns. Placement AFTER the // intentional-wait guards keeps their more precise verdicts (both already // return repaired:false, so no repair was ever at risk from them). const liveEntry = findLiveRegistryEntry(runId); if (liveEntry) { return { runId, verdict: "healthy", repaired: false, // Honest text (G24): says exactly what happened — skipped because a live // owner exists — and does not claim any repair. The "healthy" verdict // keeps callers passive: crash-recovery does not push it into the // notify result list, and reconcileOrphanedTempWorkspaces marks the // workspace hasRunning so the foreign run's state is preserved. detail: `Foreign-LIVE run: active-run registry entry is alive (cwd ${liveEntry.cwd}, updatedAt ${liveEntry.updatedAt}); reconcile skipped — a live session/runner owns this run, nothing was changed`, }; } // Phase 1: Check if results already exist const phase1 = checkResultFile(manifest, tasks); if (phase1.found) { // checkResultFile sets manifest.status='completed' and saves it, // but we re-save to ensure the completed status is persisted // before returning (avoids TOCTOU where caller might re-read stale data) // NOTE (OPT-02): Remains sync — reconcileStaleRun is called from withRunLockSync in // crash-recovery.ts and from sync reconcileOrphanedTempWorkspaces; cannot await here. saveRunManifest(manifest); return { runId, verdict: "result_exists", repaired: false, detail: "All tasks already terminal — no repair needed", }; } // Phase 2: Check PID liveness const pid = manifest.async?.pid; const pidStatus = checkPidLiveness(pid, manifest.stateRoot); if (pidStatus.detail === "no pid recorded") { // No async PID may be a foreground/live run. Preserve it if task heartbeat // or agent progress proves active work even when manifest.updatedAt is old. if (hasRecentActiveEvidence(tasks, now)) { return { runId, verdict: "no_status", repaired: false, detail: "No PID recorded, but recent task heartbeat/progress exists; not repairing", }; } // No PID and no recent activity. If ALL running tasks have stale heartbeats // (beyond NO_PID_HEARTBEAT_STALE_MS = 5min), repair immediately — the worker // process is dead but we have no PID to check. This handles /tmp/ live-session // workspaces where agents exit without calling submit_result. // Merged with individual stale task check to avoid redundant duplicate checks. const { allStale, staleTaskIds } = getRunningTaskStaleness(tasks, now); if (allStale) { const repaired = repairStaleRun(manifest, tasks, "no_pid_heartbeat_stale"); return { runId, verdict: "no_status", repaired: true, detail: `No PID; all running task heartbeats stale >${Math.round(NO_PID_HEARTBEAT_STALE_MS / 60_000)}min; repaired ${repaired.filter((t) => t.status === "cancelled").length} tasks`, repairedTasks: repaired, }; } // Check for individually stale tasks even when not all are stale. // This handles the case where task A is healthy but task B is a zombie. // We repair only the zombie tasks, not the whole run. // NOTE: Individual stale task repair (this branch) is intentionally // separate from run-level repair (the allStale branch above). Both // paths can coexist — the allStale path cancels the entire run when // ALL tasks are zombies, while this path repairs only the zombie // subset when some tasks are still healthy. if (staleTaskIds.length > 0) { const repaired = repairStaleRun(manifest, tasks, "no_pid_individual_stale_task"); // Only return the individually repaired tasks in detail return { runId, verdict: "no_status", repaired: true, detail: `No PID; ${staleTaskIds.length} individually stale task(s) repaired: ${staleTaskIds.join(", ")}`, repairedTasks: repaired.filter((t) => staleTaskIds.includes(t.id)), persistTasks: repaired, // NEW-1: persist FULL array (data-loss fix — SDD WI-2) }; } // Fall through: no recent activity but not all tasks stale enough yet. // Check the longer STALE_ALIVE_PID_MS threshold for very old runs. const updatedAt = new Date(manifest.updatedAt).getTime(); if (Number.isFinite(updatedAt) && now - updatedAt > STALE_ALIVE_PID_MS) { const repaired = repairStaleRun(manifest, tasks, "no_pid_stale"); return { runId, verdict: "no_status", repaired: true, detail: `No PID; stale ${Math.round((now - updatedAt) / 3600_000)}h; repaired ${repaired.filter((t) => t.status === "cancelled").length} tasks`, repairedTasks: repaired, }; } return { runId, verdict: "no_status", repaired: false, detail: "No PID recorded; not stale enough to repair", }; } // Phase 3: Evaluate staleness const staleness = evaluateStaleness(manifest, pidStatus.alive, now); if (!staleness.stale) { return { runId, verdict: "healthy", repaired: false, detail: `PID ${pid}: ${pidStatus.detail}, ${staleness.reason}`, }; } // Repair const repaired = repairStaleRun(manifest, tasks, staleness.reason); return { runId, verdict: pidStatus.alive ? "pid_alive_stale" : "pid_dead", repaired: true, detail: `PID ${pid}: ${pidStatus.detail}; ${staleness.reason}; repaired ${repaired.filter((t) => t.status === "cancelled").length} tasks`, repairedTasks: repaired, }; } /** * Result of orphaned temp workspace reconciliation. */ export interface OrphanReconcileResult { /** Number of runs repaired (manifests cancelled). */ repaired: number; /** Number of /tmp/pi-crew-* directories removed. */ cleanedDirs: number; } /** * Scan /tmp (os.tmpdir()) for orphaned pi-crew-* workspaces and reconcile * any stale runs found. This catches runs created by tests or crashed sessions * that the per-CWD auto-repair timer would miss. * * When `cleanupOrphanedTempDirs` is not explicitly set to `false`, directories * older than 1 hour with no remaining running manifests are deleted after * their runs are reconciled. * * @returns Number of runs repaired and directories cleaned. */ export function reconcileOrphanedTempWorkspaces( now = Date.now(), options?: { cleanupOrphanedTempDirs?: boolean; tmpDir?: string; scanBatchSize?: number; }, ): OrphanReconcileResult { // Injectable tmpDir + scanBatchSize for deterministic unit testing // (Round 19: tests must not depend on global /tmp cleanliness; the // production ORPHAN_TEMP_SCAN_BATCH_SIZE cap could exclude a test's dir // when leftover dirs accumulate). Defaults remain os.tmpdir() + the cap. const tmpDir = options?.tmpDir ?? getSafeTempDir(); if (!tmpDir) return { repaired: 0, cleanedDirs: 0 }; let repaired = 0; let cleanedDirs = 0; try { const entries = fs.readdirSync(tmpDir, { withFileTypes: true }); // Sort for deterministic order; cap to ORPHAN_TEMP_SCAN_BATCH_SIZE per // tick to avoid main-thread stalls when /tmp has thousands of // pi-crew-* dirs from past interrupted test runs. const scanBatch = options?.scanBatchSize ?? ORPHAN_TEMP_SCAN_BATCH_SIZE; const candidates = entries .filter((e) => e.isDirectory() && e.name.startsWith("pi-crew-")) .sort((a, b) => a.name.localeCompare(b.name)); // STATELESS ROUND-ROBIN across ticks (Bug B, live 2026-09-26): the old // `slice(0, scanBatch)` visited only the alphabetically-first batch, so a // cluster of uncleanable dirs (frozen sentinels / fresh sentinels / // waiting runs) starved everything behind it — 42 frozen sentinels in the // `agent-stale-wakeup-test-*` cluster (alphabetically first) blocked ~3.4k // dirs from EVER being scanned. Each tick now takes the NEXT slice, // derived deterministically from `now` (production cadence is 60s/tick — // see observability.ts tempReconcileTimer), so every batch is visited // within `batches` ticks without any persisted state. Deterministic under // the injected `now` that tests already use. const batches = Math.max(1, Math.ceil(candidates.length / scanBatch)); const batchIdx = Math.floor(now / ORPHAN_RECONCILE_TICK_MS) % batches; const selected = candidates.slice(batchIdx * scanBatch, batchIdx * scanBatch + scanBatch); for (const entry of selected) { if (!entry.isDirectory() || !entry.name.startsWith("pi-crew-")) continue; const workspaceDir = path.join(tmpDir, entry.name); // RR-020 Fix 3: accept BOTH supported run-state layouts — `/.crew/ // state/runs` and the `.pi`-based `/.pi/teams/state/runs`. Only the // `.crew` layout used to be seen here, so a temp workspace with live // `.pi/teams` run state was invisible to the reconciler (and then // deleted as debris by the legacy cleanup sweep). const stateRunsDir = findRunStateDir(workspaceDir); if (!stateRunsDir) continue; let hasRunning = false; try { // Cold-verify F-2: Dirent.isDirectory() is FALSE for a symlink-to-dir // (withFileTypes does not follow), so a planted `runs/` symlink // is skipped instead of being read/written out-of-tree. for (const entry of fs.readdirSync(stateRunsDir, { withFileTypes: true })) { if (!entry.isDirectory()) continue; const runDir = entry.name; const manifestPath = path.join(stateRunsDir, runDir, "manifest.json"); const tasksPath = path.join(stateRunsDir, runDir, "tasks.json"); if (!fs.existsSync(manifestPath) || !fs.existsSync(tasksPath)) continue; try { const manifest = loadManifestWithRecovery(manifestPath, runDir); if (!manifest) continue; // corrupt (now quarantined) or missing — skip if (manifest.status !== "running") continue; const tasks = loadTasksWithRecovery(tasksPath, path.join(stateRunsDir, runDir, "events.jsonl"), runDir); const result = reconcileStaleRun(manifest, tasks, now); if (result.repaired && result.repairedTasks) { // Persist repaired tasks (atomic — temp+rename to survive mid-write crash) atomicWriteJson(tasksPath, result.repairedTasks); // Update manifest status const updated = { ...manifest, status: "cancelled" as const, updatedAt: new Date(now).toISOString(), summary: `Stale run reconciled: ${result.detail}`, }; atomicWriteJson(manifestPath, updated); // Update agent records for (const task of result.repairedTasks) { try { upsertCrewAgent(updated, recordFromTask(updated, task, "scaffold")); } catch { /* non-critical */ } } repaired++; } // Persist result_exists manifest status update if (result.verdict === "result_exists") { // checkResultFile (called by reconcileStaleRun at line 518) // already saved manifest.status='completed' at line 66. // No additional manifest write needed here — avoids TOCTOU // window between checkResultFile save and re-write. // Sync agent records (checkResultFile also does this but // we sync here for consistency) for (const task of tasks) { try { upsertCrewAgent(manifest, recordFromTask(manifest, task, "scaffold")); } catch { /* non-critical */ } } } // If still running after reconciliation attempt, mark for dir-preserving if ( result.verdict === "healthy" || // WP-2/R2: a run intentionally waiting on an ask answer keeps the // temp workspace alive (its stateRoot + mailbox must survive for // the parked worker to receive the answer). result.verdict === "waiting_answer" || (result.verdict === "no_status" && !result.repaired) ) { hasRunning = true; } } catch (err) { // Log warning when skipping a directory due to error. // Note: manifestPath here refers to the variable defined in the // for loop at line 442 (outer scope of this catch block). const scanManifestPath = manifestPath; logInternalError( "stale-reconciler", new Error(`Skipping manifest due to parse error: ${scanManifestPath}: ${err}`), undefined, "warn", ); } } } catch (err) { // Cannot determine running state — treat as if running to prevent // premature cleanup of a potentially active workspace. hasRunning = true; logInternalError("stale-reconciler", new Error(`Skipping unreadable runs dir: ${stateRunsDir}: ${err}`), undefined, "warn"); } // Post-loop: check if this workspace dir can be cleaned up. // Eligible when cleanup is enabled, no running manifests remain, and // the directory is older than the age threshold. // Re-scan manifests to confirm no running runs remain BEFORE the // cleanup decision (fixes TOCTOU race where a manifest may have // transitioned from 'running' to 'completed' between the main loop // and this re-scan). // // TOCTOU fix: Create a sentinel file before the re-scan. New run // creation should check for this sentinel and refuse/wait if present. // Additionally, do a final check right before cleanup to catch any // new runs created between the re-scan and the cleanup decision. const sentinelPath = path.join(workspaceDir, ".cleanup-in-progress"); let canCleanup = !hasRunning; // FIX: Capture dirAge BEFORE creating the sentinel file, because writing // .cleanup-in-progress inside workspaceDir updates its mtime, making the // directory appear fresh and defeating the age check. const cleanupEnabled = options?.cleanupOrphanedTempDirs !== false; let dirAge = 0; if (cleanupEnabled && !hasRunning) { try { const stat = fs.statSync(workspaceDir); dirAge = now - stat.mtimeMs; } catch { // Directory disappeared } } if (canCleanup) { // Create sentinel file before re-scan to signal cleanup in progress. // O_EXCL requires workspaceDir to exist — ensure it before writing. try { if (!fs.existsSync(workspaceDir)) { fs.mkdirSync(workspaceDir, { recursive: true }); } // Atomic exclusive-create: reserve the sentinel (fails with EEXIST if // another cleanup is in progress), then fill it atomically. The brief // empty-file window is safe — the sentinel "exists" (correctly // signalling in-progress cleanup). const sentinelFd = fs.openSync(sentinelPath, "wx"); fs.closeSync(sentinelFd); atomicWriteFile(sentinelPath, JSON.stringify({ startedAt: now })); } catch { // Sentinel already exists. Either another cleanup is genuinely in // progress (it holds the sentinel for seconds at most) OR a previous // process died mid-cleanup and ABANDONED it — the frozen-workspace // bug (no TTL on the skip). Reclaim sentinels older than // SENTINEL_STALE_RECLAIM_MS and retry the exclusive create once; // losing the retry race means someone else reclaimed first — skip. let sentinelAgeMs = Number.NaN; try { sentinelAgeMs = now - fs.statSync(sentinelPath).mtimeMs; } catch { /* sentinel vanished — retry create below */ } if (Number.isFinite(sentinelAgeMs) && sentinelAgeMs <= SENTINEL_STALE_RECLAIM_MS) { canCleanup = false; // fresh sentinel — a live cleanup owns it } else { try { fs.unlinkSync(sentinelPath); // stale/absent — reclaim } catch { /* already gone */ } try { const retryFd = fs.openSync(sentinelPath, "wx"); fs.closeSync(retryFd); atomicWriteFile(sentinelPath, JSON.stringify({ startedAt: now, reclaimedFromStale: true })); } catch { canCleanup = false; // lost the reclaim race — skip } } } } if (canCleanup) { if (fs.existsSync(stateRunsDir)) { try { // Cold-verify round 5 (E1): a string readdir FOLLOWS a symlinked // `` entry, and loadManifestWithRecovery on a corrupt manifest // then QUARANTINES it — renameSync resolves through the symlinked // dir and renamed a file OUTSIDE the scanned tree (reproduced even // with cleanupOrphanedTempDirs:false, fired by the production // tempReconcileTimer every session). Dirent.isDirectory() is false // for a symlink-to-dir, so planted entries are skipped. Deletion // itself was already safe (rmSync unlinks symlinks, never follows). for (const entry of fs.readdirSync(stateRunsDir, { withFileTypes: true })) { if (!entry.isDirectory()) continue; const runDir = entry.name; const manifestPath = path.join(stateRunsDir, runDir, "manifest.json"); if (!fs.existsSync(manifestPath)) continue; const manifest = loadManifestWithRecovery(manifestPath, runDir); if (!manifest) continue; // corrupt (now quarantined) or missing — skip if (manifest?.status === "running") { canCleanup = false; break; } } } catch (err) { logInternalError( "stale-reconciler", new Error(`Skipping unreadable runs dir: ${stateRunsDir}: ${err}`), undefined, "warn", ); } } } if (canCleanup) { if (dirAge > ORPHAN_TEMP_DIR_AGE_THRESHOLD_MS) { // Final check: re-verify no running manifests appeared since re-scan. let stillClean = true; if (fs.existsSync(stateRunsDir)) { try { // Cold-verify round 5 (E1): same Dirent guard as the first gate — // quarantine-rename through a symlinked `` wrote outside. for (const entry of fs.readdirSync(stateRunsDir, { withFileTypes: true })) { if (!entry.isDirectory()) continue; const runDir = entry.name; const manifestPath = path.join(stateRunsDir, runDir, "manifest.json"); if (!fs.existsSync(manifestPath)) continue; const manifest = loadManifestWithRecovery(manifestPath, runDir); if (!manifest) continue; // corrupt (now quarantined) or missing — skip if (manifest.status === "running") { stillClean = false; break; } } } catch { stillClean = false; } } if (stillClean) { fs.rmSync(workspaceDir, { recursive: true, force: true, }); cleanedDirs++; } } // Clean up sentinel file regardless of outcome try { fs.unlinkSync(sentinelPath); } catch { /* ignore if already gone */ } } } } catch { /* skip if tmpdir unreadable */ } return { repaired, cleanedDirs }; } function getSafeTempDir(): string | undefined { try { return fs.existsSync(os.tmpdir()) ? os.tmpdir() : undefined; } catch { return undefined; } }