import * as fs from "node:fs"; import * as path from "node:path"; import type { ExtensionContext } from "@earendil-works/pi-coding-agent"; import { DEFAULT_PATHS } from "../../config/defaults.ts"; import { appendHookEvent, executeHook } from "../../hooks/registry.ts"; import type { MetricRegistry } from "../../observability/metric-registry.ts"; import { discoverRunLockFiles, sweepStaleLocks, withRunLock, withRunLockSync } from "../../state/coordination/locks.ts"; import { appendEvent, scanSequence } from "../../state/event-log/event-log.ts"; import { readActiveRunRegistry, unregisterActiveRun } from "../../state/stores/active-run-registry.ts"; import { loadRunManifestById, loadRunManifestByIdAsync, saveRunManifest, saveRunTasks, updateRunStatus, } from "../../state/stores/state-store.ts"; import type { TeamTaskState } from "../../state/types.ts"; import { logInternalError } from "../../utils/internal-error.ts"; import { projectCrewRoot, userCrewRoot } from "../../utils/paths.ts"; import { resolveRealContainedPath } from "../../utils/safe-paths.ts"; import { sleepSync } from "../../utils/sleep.ts"; import { recordFromTask, upsertCrewAgent } from "../crew-agent-records.ts"; import { isWorkerHeartbeatStale } from "../heartbeat/worker-heartbeat.ts"; import { terminateLiveAgentsForRun } from "../live-session/live-agent-manager.ts"; import type { ManifestCache } from "../manifest-cache.ts"; import { checkProcessLiveness } from "../process-status.ts"; import { mapConcurrent } from "../scheduling/parallel-utils.ts"; import { isIntentionalWait, isPlanApprovalPendingEffective, type ReconcileResult, reconcileStaleRun } from "../stale-reconciler.ts"; export interface RecoveryPlan { runId: string; resumableTasks: string[]; preservedTasks: string[]; lastEventSeq: number; } function isTerminalTask(task: TeamTaskState): boolean { return ( task.status === "completed" || task.status === "failed" || task.status === "cancelled" || task.status === "skipped" || task.status === "needs_attention" ); } /** * Errno codes that represent transient filesystem conditions (e.g. Windows AV * scan temporarily locking the file, or a brief permission window). On these * errors the manifest is NOT corrupt — retrying the read after a short backoff * will succeed. */ const TRANSIENT_READ_ERRNO_CODES = new Set(["EBUSY", "EACCES", "EAGAIN", "EPERM", "ENFILE", "EMFILE"]); /** * Raw file-read function used by readManifestWithTransientRetry. Separated as * a module-level binding so tests can inject transient-error mocks without * needing to mutate the frozen ESM `fs` namespace. */ let _readManifestFileSync: (filePath: string) => string = (p) => fs.readFileSync(p, "utf-8"); /** * @internal — test seam for readManifestWithTransientRetry. Pass `null` to reset. */ export function _setReadManifestFileSyncForTest(fn: ((filePath: string) => string) | null): void { _readManifestFileSync = fn ?? ((p) => fs.readFileSync(p, "utf-8")); } /** * Read and JSON-parse a manifest file, retrying transient read errors * (EBUSY/EACCES/EAGAIN — common on Windows when AV software briefly locks the * file) with exponential backoff. * * ST-6: Previously, a single `catch` quarantined the manifest on ALL errors, * including transient I/O failures that left a perfectly healthy manifest * unparseable for a few milliseconds. This renamed the valid file to * `.corrupt-*`, making the run permanently unloadable. * * Behavior: * - Success → returns parsed object * - SyntaxError (genuinely corrupt JSON) → rethrown immediately (caller quarantines) * - Transient ErrnoException → retried up to `maxRetries` times; if still * failing, rethrown so the caller can skip WITHOUT quarantining * - Other ErrnoException (ENOENT, etc.) → rethrown immediately */ export function readManifestWithTransientRetry(manifestPath: string, maxRetries = 3, baseDelayMs = 100): Record { for (let attempt = 0; attempt <= maxRetries; attempt++) { try { return JSON.parse(_readManifestFileSync(manifestPath)) as Record; } catch (err) { // SyntaxError = genuinely corrupt JSON → rethrow immediately if (err instanceof SyntaxError) throw err; const code = (err as NodeJS.ErrnoException).code; // Transient read error → retry with exponential backoff if (code !== undefined && TRANSIENT_READ_ERRNO_CODES.has(code)) { if (attempt < maxRetries) { sleepSync(baseDelayMs * 2 ** attempt); continue; } } // Non-transient error, or retries exhausted → rethrow throw err; } } // Unreachable: the loop either returns or throws on every iteration. throw new Error(`unreachable: readManifestWithTransientRetry exhausted for ${manifestPath}`); } export function shouldRecoverTask(task: TeamTaskState, deadMs: number): boolean { if (task.status !== "running") return false; if (!task.heartbeat) return true; return task.heartbeat.alive === false || isWorkerHeartbeatStale(task.heartbeat, deadMs); } export function detectInterruptedRuns( cwd: string, manifestCache: ManifestCache, deadMs = 300_000, currentSessionId?: string, ): RecoveryPlan[] { const plans: RecoveryPlan[] = []; for (const manifest of manifestCache.list(50)) { if (manifest.status !== "running" && manifest.status !== "blocked") continue; // Preserve runs intentionally blocked on plan approval — not crashes. // B1 battery 2026-08-18 (case b): generalized to isIntentionalWait — a run // parked on an ask answer is equally intentional and must not be auto-resumed // out from under its parked worker. if (isIntentionalWait(manifest)) continue; if (manifest.async?.pid !== undefined && checkProcessLiveness(manifest.async.pid).alive) continue; // Skip runs owned by the current live session — a live session B must NOT // detect session A's still-running run as interrupted. if (currentSessionId && manifest.ownerSessionId && manifest.ownerSessionId === currentSessionId) continue; // NOTE: no withRunLock — best-effort only; concurrent writes may cause inconsistency const loaded = loadRunManifestById(cwd, manifest.runId); // NOTE: no withRunLock - best-effort only; concurrent writes may cause inconsistency if (!loaded) continue; const resumableTasks = loaded.tasks.filter((task) => shouldRecoverTask(task, deadMs)).map((task) => task.id); if (!resumableTasks.length) continue; plans.push({ runId: manifest.runId, resumableTasks, preservedTasks: loaded.tasks.filter(isTerminalTask).map((task) => task.id), lastEventSeq: scanSequence(loaded.manifest.eventsPath), }); } return plans; } export async function applyRecoveryPlan(plan: RecoveryPlan, ctx: Pick, registry?: MetricRegistry): Promise { const loaded = loadRunManifestById(ctx.cwd, plan.runId); // NOTE: no withRunLock - best-effort only; concurrent writes may cause inconsistency if (!loaded) throw new Error(`Run '${plan.runId}' not found.`); const hookReport = await executeHook("run_recovery", { runId: plan.runId, cwd: ctx.cwd, }); // R14-4: the hook await above is an arbitrary async gap — a concurrent writer // may have completed/cancelled the run in the meantime. Re-read inside the // lock and derive every write from the FRESH snapshot; never reset a run that // reached a terminal status. withRunLockSync(loaded.manifest, () => { const fresh = loadRunManifestById(ctx.cwd, plan.runId); // NOTE: inside withRunLockSync - consistent read if (!fresh) throw new Error(`Run '${plan.runId}' not found.`); if (fresh.manifest.status === "completed" || fresh.manifest.status === "failed" || fresh.manifest.status === "cancelled") { // Run reached a terminal status while the recovery hook was running — // do NOT reset it (no task reset, no status change). appendEvent(fresh.manifest.eventsPath, { type: "crew.run.recovery_skipped", runId: plan.runId, message: `Recovery skipped: run is already '${fresh.manifest.status}'`, data: { status: fresh.manifest.status }, }); return; } appendHookEvent(fresh.manifest, hookReport); if (hookReport.outcome === "block") { appendEvent(fresh.manifest.eventsPath, { type: "crew.run.recovery_blocked", runId: plan.runId, message: `Recovery blocked by hook: ${hookReport.reason ?? "run_recovery hook blocked the operation."}`, data: { hookOutcome: "block", reason: hookReport.reason }, }); return; } const reset = new Set(plan.resumableTasks); const tasks = fresh.tasks.map((task) => reset.has(task.id) ? { ...task, status: "queued" as const, startedAt: undefined, finishedAt: undefined, error: undefined, heartbeat: undefined, // WP-2/R2 (ADR-0 item 9): a requeued task must not resume parked // state — drop any stale ask-park marker along with the other // per-attempt transient fields. Tasks NOT reset (e.g. a task parked // in `waiting`) pass through untouched, marker intact; the marker's // deadline is then owned by the scheduler tick (dispatch-batch). waiting: undefined, } : task, ); saveRunTasks(fresh.manifest, tasks); appendEvent(fresh.manifest.eventsPath, { type: "crew.run.resumed", runId: plan.runId, message: `Recovered ${plan.resumableTasks.length} interrupted task(s).`, data: { recoveredFromSeq: plan.lastEventSeq, resumableTasks: plan.resumableTasks, }, }); registry?.counter("crew.run.count", "Total runs by status").inc({ status: "resumed" }); }); } export function declineRecoveryPlan(plan: RecoveryPlan, ctx: Pick): void { const loaded = loadRunManifestById(ctx.cwd, plan.runId); // NOTE: no withRunLock - best-effort only; concurrent writes may cause inconsistency if (!loaded) throw new Error(`Run '${plan.runId}' not found.`); // R14-4: re-read inside the lock so the decline decision applies to the // on-disk state, not the best-effort snapshot read above. Kept SYNC — // withRunLockSync is the sync lock variant. withRunLockSync(loaded.manifest, () => { const fresh = loadRunManifestById(ctx.cwd, plan.runId); // NOTE: inside withRunLockSync - consistent read if (!fresh) return; // Log the event first — if appendEvent fails, state remains consistent. appendEvent(fresh.manifest.eventsPath, { type: "crew.run.recovery_declined", runId: plan.runId, message: "Interrupted run was not resumed.", data: { recoveredFromSeq: plan.lastEventSeq }, }); updateRunStatus(fresh.manifest, "cancelled", "interrupted-not-resumed"); }); } /** * Run 3-phase stale reconciliation on all active runs. * Returns results for each reconciled run. */ /** * Auto-cancel orphaned runs whose owner session no longer exists. * * When a Pi session dies (crash, force-close, Ctrl+C), `session_shutdown` * does not fire and child workers are not terminated. The next Pi session * must detect these orphaned runs and cancel them. * * Criteria for orphan detection: * 1. Manifest status is "running" * 2. Manifest has an `ownerSessionId` that is NOT the current session * 3. The owner session's process is no longer alive (PID check) * 4. No recent heartbeat activity (task heartbeat or agent progress within threshold) * * Returns the number of runs cancelled. */ export function cancelOrphanedRuns( cwd: string, manifestCache: ManifestCache, currentSessionId: string, staleThresholdMs = 300_000, now = Date.now(), ): { cancelled: string[]; skipped: string[] } { const cancelled: string[] = []; const skipped: string[] = []; // Phase 1: Scan project-level manifests via manifestCache for (const manifest of manifestCache.list(50)) { if (manifest.status !== "running" && manifest.status !== "blocked") continue; // B1 battery 2026-08-18 (case b): preserve runs intentionally blocked on a // HUMAN DECISION — plan approval (pre-v2 semantics) OR a pending ask answer // within the waiting TTL (ADR-0 item 9: "not crashes and must not be // stale-repaired or cancelled by reconciliation"). Live proof: a resumed // run parked on ask (worker alive, waitState set, parked workers do NOT // heartbeat) was orphan-cancelled 5m18s into a 10-minute park by a THIRD // session's startup scan because the predicate only knew plan approval. if (isIntentionalWait(manifest)) { skipped.push(manifest.runId); continue; } // Only consider runs owned by a different session const ownerId = manifest.ownerSessionId; if (!ownerId || ownerId === currentSessionId) continue; // Check if the owner process is still alive const ownerPid = manifest.async?.pid; if (ownerPid !== undefined && checkProcessLiveness(ownerPid).alive) { skipped.push(manifest.runId); continue; } // Foreground (in-process) runs have NO async worker PID, so the liveness // check above cannot apply. They execute inside THIS pi process — a // recently-updated manifest conclusively proves the run is alive; only a // run that has gone quiet for the stale threshold can be an orphan (a // crashed process stops updating). Without this, a YOUNG foreground run // (first worker heartbeat not yet written) was orphan-cancelled by the // session that opened its agent view seconds after the run started — // killing the very run the user was looking at. if (!manifest.async) { const updatedMs = manifest.updatedAt ? new Date(manifest.updatedAt).getTime() : Number.NaN; if (Number.isFinite(updatedMs) && now - updatedMs <= staleThresholdMs) { skipped.push(manifest.runId); continue; } } // Check for recent heartbeat activity const loaded = loadRunManifestById(cwd, manifest.runId); // NOTE: no withRunLock - best-effort only; concurrent writes may cause inconsistency if (!loaded) continue; const hasRecentActivity = loaded.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 <= staleThresholdMs) return true; const activityAt = task.agentProgress?.lastActivityAt ? new Date(task.agentProgress.lastActivityAt).getTime() : Number.NaN; return Number.isFinite(activityAt) && now - activityAt <= staleThresholdMs; }); if (hasRecentActivity) { skipped.push(manifest.runId); continue; } // Orphan confirmed — mark durable state terminal before best-effort live-agent abort. // terminateLiveAgent unregisters handles before awaiting abort(), and live-executor's // isCurrent() checks durable terminal state before writing progress. // Orphan confirmed — cancel all running tasks let cancelledRun = false; withRunLockSync(loaded.manifest, () => { const fresh = loadRunManifestById(cwd, manifest.runId); // NOTE: inside withRunLockSync - consistent read if (!fresh) return; if (fresh.manifest.status !== "running" && fresh.manifest.status !== "blocked") { // Status changed between initial check (line 109) and acquiring the lock — normal concurrent update, not an orphan appendEvent(loaded.manifest.eventsPath, { type: "crew.run.orphan_skip", runId: manifest.runId, message: `Skipped orphan cancellation: status is '${fresh.manifest.status}' (was 'running'/'blocked' at initial scan)`, data: { currentStatus: fresh.manifest.status }, }); return; } const now_iso = new Date(now).toISOString(); const repairedTasks = fresh.tasks.map((task) => { if (task.status === "running" || task.status === "queued" || task.status === "waiting") { return { ...task, status: "cancelled" as const, finishedAt: now_iso, error: `Orphaned run: owner session ${ownerId} no longer exists`, }; } return task; }); saveRunTasks(fresh.manifest, repairedTasks); // B1 battery 2026-08-18 (case b leak fix): a cancelled run must not keep // its waitState pointer — observed live: orphan-cancel left waitState set, // shielding the dead run from staleness repair for the full 24h TTL. if (fresh.manifest.waitState) { const cleared = { ...fresh.manifest, waitState: undefined, updatedAt: new Date(now).toISOString() }; saveRunManifest(cleared); fresh.manifest = cleared; } for (const task of repairedTasks) { try { upsertCrewAgent(fresh.manifest, recordFromTask(fresh.manifest, task, "scaffold")); } catch { /* non-critical */ } } updateRunStatus(fresh.manifest, "cancelled", `Orphaned run: owner session ${ownerId} no longer exists`); appendEvent(fresh.manifest.eventsPath, { type: "crew.run.orphan_cancelled", runId: manifest.runId, message: `Auto-cancelled orphaned run (owner: ${ownerId})`, data: { ownerSessionId: ownerId, cancelledTasks: repairedTasks.filter((t) => t.status === "cancelled").length, }, }); cancelled.push(manifest.runId); cancelledRun = true; }); if (cancelledRun) void terminateLiveAgentsForRun(manifest.runId, "cancelled", appendEvent, loaded.manifest.eventsPath).catch((error) => logInternalError("crash-recovery.orphan.terminate", error, `runId=${manifest.runId}`, "warn"), ); } return { cancelled, skipped }; } /** * Purge the global active-run-index of entries whose manifest is no longer active. * * This scans every entry in active-run-index.json and removes any whose: * - manifest file no longer exists, OR * - manifest status is terminal (completed/failed/cancelled/blocked), OR * - manifest cwd directory no longer exists (e.g. temp test dirs) * * Also removes entries where the manifest is still "running" but: * - The cwd has been deleted (temp dir cleanup) * - The async worker PID is dead AND no heartbeat for > threshold * * This is the **global** cleanup that cancelOrphanedRuns (project-scoped) * cannot reach. */ /** * Best-effort removal of stateRoot and artifactsRoot directories for a purged run. * Uses resolveRealContainedPath to ensure we only delete paths that are safely * contained within a known crew root (project or user level). */ function tryRemoveRunDirectories(entry: { stateRoot: string; cwd: string }): void { const roots = [projectCrewRoot(entry.cwd), userCrewRoot()]; for (const root of roots) { try { resolveRealContainedPath(root, entry.stateRoot); // If we get here, stateRoot is safely contained — remove it fs.rmSync(entry.stateRoot, { recursive: true, force: true }); break; } catch { // Not contained in this root, try next } } // NOTE: artifactsRoot is shared across runs and cleaned up by pruneFinishedRuns/pruneUserLevelRuns — not deleted here. } /** * Age (ms) of the team-level heartbeat file for a run. The team-runner writes * `/heartbeat.json` periodically while a workflow is executing * (startTeamHeartbeat), so a fresh heartbeat is strong evidence the run is alive * even when its recorded PID check is inconclusive or its active-run-index * entry's `updatedAt` was frozen at registration. Returns Infinity when absent. */ function heartbeatAgeMs(entry: { stateRoot: string }, now: number): number { try { const mtime = fs.statSync(path.join(entry.stateRoot, "heartbeat.json")).mtimeMs; return Number.isFinite(mtime) ? now - mtime : Infinity; } catch { return Infinity; } } /** * True if there is recent evidence the run is (or was very recently) alive, so * it must NOT be purged. Any one of these signals is sufficient: * - on-disk `manifest.updatedAt` fresher than `staleThresholdMs` (rewritten on * every task transition / status change), and/or * - team-level `heartbeat.json` fresher than `staleThresholdMs`. * `entry.updatedAt` is intentionally NOT consulted: it is frozen at * registration and never refreshed during execution, which previously caused * long-running legitimate runs to be falsely purged — destroying their * stateRoot, and because saveRunTasks() silently no-ops once the state dir is * gone, hanging the workflow permanently at the current task with no * recoverable state ("Run not found"). */ function hasRecentLifeEvidence( entry: { stateRoot: string }, manifestUpdatedAt: string | undefined, now: number, staleThresholdMs: number, ): boolean { const manifestMs = manifestUpdatedAt ? new Date(manifestUpdatedAt).getTime() : NaN; if (Number.isFinite(manifestMs) && now - manifestMs <= staleThresholdMs) return true; const hbAge = heartbeatAgeMs(entry, now); if (Number.isFinite(hbAge) && hbAge <= staleThresholdMs) return true; return false; } /** * Purge the global active-run-index of entries whose manifest is no longer active. * * Note: This function only cleans user-level active run entries. * Project-level stale runs are handled by session_start auto-prune triggered during run creation. */ export function purgeStaleActiveRunIndex( staleThresholdMs = 300_000, now = Date.now(), currentSessionId?: string, ): { purged: string[]; kept: string[] } { const purged: string[] = []; const kept: string[] = []; const entries = readActiveRunRegistry(); for (const entry of entries) { // 1. Manifest file gone → stale, but PRESERVE stateRoot for manual recovery. // R-01: A single missing-manifest signal is insufficient to warrant // deleting all run state (events.jsonl, task data, artifacts). Require // corroboration — only remove stateRoot if it is ALREADY gone. if (!fs.existsSync(entry.manifestPath)) { if (!fs.existsSync(entry.stateRoot)) { tryRemoveRunDirectories(entry); } unregisterActiveRun(entry.runId); purged.push(entry.runId); continue; } // 2. CWD gone → temp dir cleaned up, but PRESERVE stateRoot for recovery. // R-01: Same dual-signal guard as step 1 — a single missing-CWD signal // must not delete all run state. Only clean up if stateRoot is already gone. if (!fs.existsSync(entry.cwd)) { if (!fs.existsSync(entry.stateRoot)) { tryRemoveRunDirectories(entry); } unregisterActiveRun(entry.runId); purged.push(entry.runId); continue; } // 3. Read manifest status let manifest: | { status?: string; updatedAt?: string; async?: { pid?: number }; ownerSessionId?: string; } | undefined; try { manifest = readManifestWithTransientRetry(entry.manifestPath) as typeof manifest; } catch (err) { // ST-6: Transient read errors (EBUSY/EACCES/EAGAIN from Windows AV // scans) must NOT quarantine a healthy manifest — the file is fine, // just temporarily locked. Skip this entry and try again next cycle. const code = (err as NodeJS.ErrnoException).code; const isTransient = code !== undefined && TRANSIENT_READ_ERRNO_CODES.has(code); if (isTransient) continue; // R-01: Quarantine the corrupt manifest instead of deleting all run // state. A single corrupted byte must NOT cause total data loss — // events.jsonl, task data, and artifacts are preserved for manual // recovery. SyntaxError (genuinely corrupt JSON) and unexpected // non-transient errors fall through to here. try { // Include pid for uniqueness across concurrent quarantines in the same ms. fs.renameSync(entry.manifestPath, entry.manifestPath + ".corrupt-" + Date.now() + "-" + process.pid); } catch { // Rename may fail if the path doesn't exist or permissions deny it — best effort. } unregisterActiveRun(entry.runId); purged.push(entry.runId); continue; } // 4. Terminal status → no longer active (just unregister, don't delete files) const terminalStatuses = new Set(["completed", "failed", "cancelled", "blocked"]); if (manifest && terminalStatuses.has(manifest.status ?? "")) { unregisterActiveRun(entry.runId); purged.push(entry.runId); continue; } // 4b. Skip runs owned by the current live session — a live session B must // NOT purge session A's still-running run from the active-run-index. if (currentSessionId && manifest?.ownerSessionId && manifest.ownerSessionId === currentSessionId) { kept.push(entry.runId); continue; } // 5. Still "running" with an async worker PID — only purge when the worker // is actually dead AND there is no recent evidence of life. We must NOT // rely solely on `entry.updatedAt` (frozen at registration) nor on a single // dead-PID reading: a long-running worker (e.g. a 15-minute explorer) // legitimately keeps the run "running" while periodically rewriting the // on-disk manifest.updatedAt and heartbeat.json. Falsely purging such a run // destroys its stateRoot, and because saveRunTasks() silently no-ops once // the state dir is gone, the workflow then hangs permanently at the // current task with no recoverable state ("Run not found"). When we do mark // a run cancelled here, we KEEP its stateRoot so the run stays queryable/ // resumable and its diagnostics survive; the finished-run pruner removes // the directory later on its normal schedule. if (manifest?.status === "running" && manifest.async?.pid !== undefined) { const pidAlive = checkProcessLiveness(manifest.async.pid).alive; if (!pidAlive && !hasRecentLifeEvidence(entry, manifest.updatedAt, now, staleThresholdMs)) { // Dead PID + no recent life evidence → cancel the manifest and unregister. // RT-F5: wrap load+modify+save in withRunLockSync so concurrent writers // (e.g. team-runner final saveRunTasks) cannot interleave between the // repaired-tasks write and the manifest-status flip. terminateLiveAgentsForRun // is fire-and-forget so it stays OUTSIDE the lock to avoid blocking the // critical section on async IPC. try { const fullLoaded = loadRunManifestById(entry.cwd, entry.runId); if (fullLoaded) { withRunLockSync(fullLoaded.manifest, () => { // Re-read under lock so the stale-best-effort load above // can't cause us to clobber a fresh status flip. const fresh = loadRunManifestById(entry.cwd, entry.runId); if (!fresh) return; if (fresh.manifest.status !== "running") return; const now_iso = new Date(now).toISOString(); const repairedTasks = fresh.tasks.map((task) => { if (task.status === "running" || task.status === "queued" || task.status === "waiting") { return { ...task, status: "cancelled" as const, finishedAt: now_iso, error: "Orphaned run: worker process dead and no recent activity", }; } return task; }); saveRunTasks(fresh.manifest, repairedTasks); for (const task of repairedTasks) { try { upsertCrewAgent(fresh.manifest, recordFromTask(fresh.manifest, task, "scaffold")); } catch { /* non-critical */ } } updateRunStatus(fresh.manifest, "cancelled", "Orphaned run: worker process dead and no recent activity"); }); void terminateLiveAgentsForRun( fullLoaded.manifest.runId, "cancelled", appendEvent, fullLoaded.manifest.eventsPath, ).catch((error) => logInternalError("crash-recovery.pid-dead.terminate", error, `runId=${fullLoaded.manifest.runId}`, "warn"), ); } } catch { // Best-effort manifest cleanup } unregisterActiveRun(entry.runId); purged.push(entry.runId); continue; } } // 6. "running" but no async worker PID — possible orphaned run where the // manifest was never updated to a terminal status after the worker exited. // Uses the same life-evidence corroboration as condition 5; the stateRoot is // kept on cancel so the run stays queryable/resumable with diagnostics. if (manifest?.status === "running" && manifest.async === undefined) { if (!hasRecentLifeEvidence(entry, manifest.updatedAt, now, staleThresholdMs)) { try { const fullLoaded = loadRunManifestById(entry.cwd, entry.runId); if (fullLoaded && fullLoaded.manifest.status === "running") { withRunLockSync(fullLoaded.manifest, () => { const fresh = loadRunManifestById(entry.cwd, entry.runId); if (fresh?.manifest.status !== "running") return; const now_iso = new Date(now).toISOString(); const repairedTasks = fresh.tasks.map((task) => { if (task.status === "running" || task.status === "queued" || task.status === "waiting") { return { ...task, status: "cancelled" as const, finishedAt: now_iso, error: "Orphaned run: workflow completed but manifest never updated to terminal status", }; } return task; }); saveRunTasks(fresh.manifest, repairedTasks); for (const task of repairedTasks) { try { upsertCrewAgent(fresh.manifest, recordFromTask(fresh.manifest, task, "scaffold")); } catch { /* non-critical */ } } updateRunStatus( fresh.manifest, "cancelled", "Orphaned run: no async worker and no manifest update in over " + Math.round(staleThresholdMs / 60000) + " minutes", ); }); void terminateLiveAgentsForRun( fullLoaded.manifest.runId, "cancelled", appendEvent, fullLoaded.manifest.eventsPath, ).catch((error) => logInternalError("crash-recovery.pid-dead.terminate", error, `runId=${fullLoaded.manifest.runId}`, "warn"), ); } } catch { // Best-effort } unregisterActiveRun(entry.runId); purged.push(entry.runId); continue; } } kept.push(entry.runId); } return { purged, kept }; } export async function reconcileAllStaleRuns( cwd: string, manifestCache: ManifestCache, now = Date.now(), currentSessionId?: string, ): Promise { // Capture runIds to reconcile BEFORE acquiring locks — avoids TOCTOU between cache iteration and lock acquisition. const runIds = manifestCache .list(50) .filter((m) => { if (m.status !== "running" && m.status !== "blocked") return false; // Skip runs owned by the current live session — a live session B must // NOT reconcile (mark failed) session A's still-running run. if (currentSessionId && m.ownerSessionId && m.ownerSessionId === currentSessionId) return false; return true; }) .map((m) => m.runId); // RR-021 WI-1.5 (audit C3): serial reconcile → mapConcurrent bound 4. // Each runId acquires its OWN run lock file (lockPath() is per-run), so the // acquisitions are independent; bound 4 caps fs/CPU contention. Results keep // input (snapshot) order — mapConcurrent indexes by item — and per-run // semantics (re-read inside lock, plan-approval re-check) are unchanged. const perRun = await mapConcurrent(runIds, 4, async (runId): Promise => { // RR-021 review remediation: per-run error isolation — one corrupt run // must not reject the whole sweep (serial baseline threw too, but the // SPEC's "independent per-run" intent requires continue-on-error). try { const cached = manifestCache.get(runId); if (!cached) return []; const loaded = await loadRunManifestByIdAsync(cwd, runId); // NOTE: best-effort only; concurrent writes may cause inconsistency if (!loaded) return []; const out: ReconcileResult[] = []; // Use lock to prevent race with cancel/status handlers modifying the same run. // Same v0.9.26 lock family as the old sync acquisition — interop-safe. await withRunLock(loaded.manifest, async () => { // Re-read inside lock to get freshest data const fresh = await loadRunManifestByIdAsync(cwd, runId); // NOTE: inside withRunLock - consistent read if (!fresh || (fresh.manifest.status !== "running" && fresh.manifest.status !== "blocked")) return; // Belt-and-suspenders: reconcileStaleRun itself guards this, but the run // may have flipped to blocked+plan-approval between cache-list and lock // acquisition — re-check the freshest manifest under the lock. if (isPlanApprovalPendingEffective(fresh.manifest)) { out.push({ runId, verdict: "blocked_awaiting_approval", repaired: false, detail: "Plan approval is pending; stale reconciliation skipped", }); return; } const result = reconcileStaleRun(fresh.manifest, fresh.tasks, now); if (result.repaired || result.verdict === "result_exists") { if (result.repairedTasks) { // NEW-1 (SDD 2026-09-30 WI-2): persist the FULL task array — the // individual-stale branch filters `repairedTasks` down to the stale // subset, and saveRunTasks full-overwrites tasks.json, which used to // drop every healthy task from disk. `persistTasks` carries the full // array; every other branch leaves it unset and falls back here. saveRunTasks(fresh.manifest, result.persistTasks ?? result.repairedTasks); for (const task of result.repairedTasks) { try { upsertCrewAgent(fresh.manifest, recordFromTask(fresh.manifest, task, "scaffold")); } catch { /* non-critical */ } } } updateRunStatus(fresh.manifest, "failed", `Stale run reconciled: ${result.detail}`); void terminateLiveAgentsForRun(fresh.manifest.runId, "failed", appendEvent, fresh.manifest.eventsPath).catch((error) => logInternalError("crash-recovery.reconcile.terminate", error, `runId=${fresh.manifest.runId}`, "warn"), ); appendEvent(fresh.manifest.eventsPath, { type: "crew.run.reconciled_stale", runId, message: result.detail, data: { verdict: result.verdict }, }); } if (result.verdict !== "healthy") { out.push(result); } }); return out; } catch (error) { logInternalError("crash-recovery.reconcileStaleRuns", error, `runId=${runId}`, "warn"); return []; } }); const results = perRun.flat(); // US-002 (2026-09-22): structured stale-lock sweep. Runs after reconciliation // so locks whose holder died (kill -9 / crash) do not linger. The sweep is // strictly safer than acquire-time steal (stale AND holder-dead only), so a // live holder is never disturbed. Best-effort. try { const runsRoots = [ path.join(projectCrewRoot(cwd), DEFAULT_PATHS.state.runsSubdir), path.join(userCrewRoot(), DEFAULT_PATHS.state.runsSubdir), ]; for (const runsRoot of runsRoots) { const { removed } = sweepStaleLocks(discoverRunLockFiles(runsRoot)); if (removed.length > 0) { logInternalError( "crash-recovery.sweep-stale-locks", new Error(`swept ${removed.length} stale run lock(s)`), removed.join(", "), "warn", ); } } } catch (error) { logInternalError("crash-recovery.sweep-stale-locks", error); } return results; }