import assert from "node:assert/strict"; import { spawn } from "node:child_process"; import { once } from "node:events"; import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; import { describe, it } from "node:test"; import { writeAsyncResultFile, writePendingAsyncResultFile } from "../../src/runs/background/result-files.ts"; import { listAsyncRuns } from "../../src/runs/background/async-status.ts"; import { finalizeProcessTerminal, initializeProcessTerminal, readProcessTerminal } from "../../src/runs/background/process-terminal.ts"; import { checkPidLiveness, reconcileAsyncRun } from "../../src/runs/background/stale-run-reconciler.ts"; import { summarizeAsyncStatus } from "../../src/runs/background/async-status.ts"; function tempRoot(prefix: string): string { return fs.mkdtempSync(path.join(os.tmpdir(), prefix)); } function writeStatus(asyncDir: string, status: Record): void { fs.mkdirSync(asyncDir, { recursive: true }); fs.writeFileSync(path.join(asyncDir, "status.json"), JSON.stringify(status, null, 2), "utf-8"); } function errno(code: string): NodeJS.ErrnoException { const error = new Error(code) as NodeJS.ErrnoException; error.code = code; return error; } describe("async stale-run reconciliation", () => { it("classifies pid liveness without treating EPERM as dead", () => { assert.equal(checkPidLiveness(process.pid), "alive"); assert.equal(checkPidLiveness(2_147_483_647, () => true), "alive"); assert.equal(checkPidLiveness(123, () => { throw errno("ESRCH"); }), "dead"); assert.equal(checkPidLiveness(123, () => { throw errno("EPERM"); }), "unknown"); assert.equal(checkPidLiveness(123, () => { throw new Error("boom"); }), "unknown"); }); it("repairs a recent Linux zombie only with a verified PID namespace", { skip: process.platform !== "linux", timeout: 10_000 }, async () => { const root = tempRoot("pi-stale-run-zombie-"); const script = String.raw` const { spawn } = require("node:child_process"); const fs = require("node:fs"); const child = spawn(process.execPath, ["-e", 'process.title = "pi) zombie"; process.exit(0);'], { stdio: "ignore" }); process.stdout.write(String(child.pid) + "\n"); // Keep libuv from reaping the child until the test releases stdin. fs.readSync(0, Buffer.alloc(1), 0, 1, null); `; const parent = spawn(process.execPath, ["--eval", script], { stdio: ["pipe", "pipe", "pipe"] }); const closed = once(parent, "close"); parent.stdin.on("error", () => {}); try { const [output] = await once(parent.stdout, "data", { signal: AbortSignal.timeout(5_000) }); const pid = Number(String(output).trim()); assert.ok(Number.isSafeInteger(pid) && pid > 0); let stat = ""; const deadline = Date.now() + 5_000; while (Date.now() < deadline) { stat = fs.readFileSync(`/proc/${pid}/stat`, "utf-8"); if (stat[stat.lastIndexOf(") ") + 2] === "Z") break; await new Promise((resolve) => setTimeout(resolve, 10)); } assert.match(stat, /\(pi\) zombie\) Z /); assert.equal(process.kill(pid, 0), true); assert.equal(checkPidLiveness(pid, process.kill, true), "dead"); const scope = fs.readlinkSync("/proc/self/ns/pid", "utf-8").trim(); const now = Date.now(); const resultsDir = path.join(root, "results"); for (const scenario of [ { name: "matching", recorded: scope, observed: scope, immediate: true }, { name: "missing-recorded", recorded: undefined, observed: scope, immediate: false }, { name: "different", recorded: "pid:[other]", observed: scope, immediate: false }, { name: "missing-observed", recorded: scope, observed: undefined, immediate: false }, { name: "both-missing", recorded: undefined, observed: undefined, immediate: false }, ]) { const runId = `run-zombie-${scenario.name}`; const asyncDir = path.join(root, runId); writeStatus(asyncDir, { runId, sessionId: "zombie-session", mode: "single", state: "running", pid, ...(scenario.recorded !== undefined ? { pidNamespaceScope: scenario.recorded } : {}), startedAt: now - 100, lastUpdate: now - 100, steps: [{ agent: "worker", status: "running", startedAt: now - 100 }], }); const options = { resultsDir, pidNamespaceScope: () => scenario.observed, now: () => now }; const repaired = reconcileAsyncRun(asyncDir, options); assert.equal(repaired.repaired, scenario.immediate, scenario.name); assert.equal(repaired.status?.state, scenario.immediate ? "failed" : "running", scenario.name); const resultPath = path.join(resultsDir, `${runId}.json`); assert.equal(fs.existsSync(resultPath), scenario.immediate, scenario.name); if (!scenario.immediate) { const stale = reconcileAsyncRun(asyncDir, { ...options, staleAlivePidMs: 1000, now: () => now + 2000 }); assert.equal(stale.repaired, true, scenario.name); assert.equal(stale.status?.state, "failed", scenario.name); assert.match(stale.message ?? "", /status has not updated for 2100ms.*PID ownership is unverified/); } const receipt = fs.readFileSync(resultPath, "utf-8"); assert.equal(JSON.parse(receipt).state, "failed"); assert.equal(reconcileAsyncRun(asyncDir, options).repaired, false); assert.equal(fs.readFileSync(resultPath, "utf-8"), receipt); } assert.equal(checkPidLiveness(pid), "alive", "unscoped callers retain signal-only probing"); } finally { parent.stdin.end("reap"); await closed; fs.rmSync(root, { recursive: true, force: true }); } }); it("repairs a run whose PID belongs to another namespace only after status goes stale", () => { const root = tempRoot("pi-stale-run-pid-namespace-"); try { const asyncDir = path.join(root, "run-live"); writeStatus(asyncDir, { runId: "run-live", mode: "single", state: "running", pid: 1404, pidNamespaceScope: "pid:[runner]", startedAt: 1000, lastUpdate: 1000, steps: [{ agent: "reviewer", status: "running", startedAt: 1000 }], }); const result = reconcileAsyncRun(asyncDir, { kill: () => { throw errno("ESRCH"); }, pidNamespaceScope: () => "pid:[observer]", now: () => 2000, }); assert.equal(result.repaired, false); assert.equal(result.status?.state, "running"); assert.equal(JSON.parse(fs.readFileSync(path.join(asyncDir, "status.json"), "utf-8")).state, "running"); const stale = reconcileAsyncRun(asyncDir, { kill: () => { throw errno("ESRCH"); }, pidNamespaceScope: () => "pid:[observer]", staleAlivePidMs: 5000, now: () => 7000, }); assert.equal(stale.repaired, true); assert.equal(stale.status?.state, "failed"); assert.match(stale.message ?? "", /PID 1404 cannot be probed from this process; status has not updated for 6000ms/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("does not repair a recent namespaced run from an observer without a PID namespace", () => { const root = tempRoot("pi-stale-run-pid-namespace-unknown-"); try { const asyncDir = path.join(root, "run-dead"); writeStatus(asyncDir, { runId: "run-dead", mode: "single", state: "running", pid: 1404, pidNamespaceScope: "pid:[runner]", startedAt: 1000, lastUpdate: 1000, steps: [{ agent: "reviewer", status: "running", startedAt: 1000 }], }); const result = reconcileAsyncRun(asyncDir, { kill: () => { throw errno("ESRCH"); }, pidNamespaceScope: () => undefined, now: () => 2000, }); assert.equal(result.repaired, false); assert.equal(result.status?.state, "running"); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("marks a dead runner failed and writes exactly one completion result", () => { const root = tempRoot("pi-stale-run-"); try { const asyncDir = path.join(root, "run-dead"); const resultsDir = path.join(root, "results"); writeStatus(asyncDir, { lifecycleArtifactVersion: 3, runId: "run-dead", sessionId: "session-current", completionOwnerId: "owner-current", mode: "single", state: "running", pid: 12345, pidNamespaceScope: "pid:[same]", processTerminal: { version: 1, state: "pending", runId: "run-dead", runnerProcessInstanceId: "runner-dead" }, startedAt: 1000, lastUpdate: 1000, currentStep: 0, steps: [{ agent: "scout", status: "running", startedAt: 1000, contextOverflow: true }], }); const result = reconcileAsyncRun(asyncDir, { resultsDir, kill: () => { throw errno("ESRCH"); }, pidNamespaceScope: () => "pid:[same]", now: () => 2000, }); assert.equal(result.repaired, true); assert.equal(result.status?.state, "failed"); assert.match(result.message ?? "", /process 12345 exited or disappeared/); const status = JSON.parse(fs.readFileSync(path.join(asyncDir, "status.json"), "utf-8")); assert.equal(status.state, "failed"); assert.equal(status.sessionId, "session-current"); assert.equal(status.completionOwnerId, "owner-current"); assert.equal(status.steps[0].status, "failed"); assert.match(status.steps[0].error, /process 12345 exited or disappeared/); const resultJson = JSON.parse(fs.readFileSync(path.join(resultsDir, "run-dead.json"), "utf-8")); assert.equal(resultJson.success, false); assert.equal(resultJson.sessionId, "session-current"); assert.equal(resultJson.completionOwnerId, "owner-current"); assert.equal(resultJson.state, "failed"); assert.equal(resultJson.exitCode, 1); assert.equal(resultJson.results[0].contextOverflow, true); assert.match(resultJson.summary, /process 12345 exited or disappeared/); assert.match(fs.readFileSync(path.join(asyncDir, "events.jsonl"), "utf-8"), /subagent\.run\.repaired_stale/); const resultText = fs.readFileSync(path.join(resultsDir, "run-dead.json"), "utf-8"); const second = reconcileAsyncRun(asyncDir, { resultsDir, kill: () => { throw errno("ESRCH"); }, now: () => 3000, }); assert.equal(second.repaired, false); assert.equal(fs.readFileSync(path.join(resultsDir, "run-dead.json"), "utf-8"), resultText); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("reports the observed runner exit code and signal when a killed runner left no result", () => { const root = tempRoot("pi-stale-run-killed-"); try { const asyncDir = path.join(root, "run-killed"); const resultsDir = path.join(root, "results"); writeStatus(asyncDir, { lifecycleArtifactVersion: 3, runId: "run-killed", sessionId: "session-current", mode: "single", state: "running", pid: 4242, startedAt: 1000, lastUpdate: 1000, currentStep: 0, steps: [{ agent: "worker", status: "running", startedAt: 1000 }], }); initializeProcessTerminal(asyncDir, "run-killed", "runner-killed"); const proof = finalizeProcessTerminal(asyncDir, "run-killed", { processInstanceId: "runner-killed", closeObservedAt: 1500, exitCode: null, signal: "SIGKILL" }); assert.equal(proof.state, "unknown"); reconcileAsyncRun(asyncDir, { resultsDir, kill: () => { throw errno("ESRCH"); }, now: () => 2000 }); const resultJson = JSON.parse(fs.readFileSync(path.join(resultsDir, "run-killed.json"), "utf-8")); assert.equal(resultJson.state, "failed"); assert.match(resultJson.summary, /Async runner process 4242 exited with code none \(signal SIGKILL\) before writing a result/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("repairs sessionless stale status without writing an unindexed result", () => { const root = tempRoot("pi-stale-run-sessionless-"); try { const asyncDir = path.join(root, "run-sessionless"); const resultsDir = path.join(root, "results"); writeStatus(asyncDir, { runId: "run-sessionless", mode: "single", state: "running", pid: 12345, startedAt: 1000, lastUpdate: 1000, currentStep: 0, steps: [{ agent: "scout", status: "running", startedAt: 1000 }], }); const result = reconcileAsyncRun(asyncDir, { resultsDir, kill: () => { throw errno("ESRCH"); }, now: () => 2000, }); assert.equal(result.repaired, true); assert.equal(result.resultPath, undefined); assert.equal(JSON.parse(fs.readFileSync(path.join(asyncDir, "status.json"), "utf-8")).state, "failed"); assert.equal(fs.existsSync(path.join(resultsDir, "run-sessionless.json")), false); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("includes runner stderr diagnostics when repairing a stale startup crash", () => { const root = tempRoot("pi-stale-run-stderr-"); try { const asyncDir = path.join(root, "run-dead-stderr"); const resultsDir = path.join(root, "results"); writeStatus(asyncDir, { runId: "run-dead-stderr", sessionId: "session-current", mode: "single", state: "running", pid: 12345, startedAt: 1000, lastUpdate: 1000, currentStep: 0, steps: [{ agent: "scout", status: "running", startedAt: 1000 }], }); fs.writeFileSync(path.join(asyncDir, "runner.stderr.log"), "startup failed\nmissing peer package\n", "utf-8"); const result = reconcileAsyncRun(asyncDir, { resultsDir, kill: () => { throw errno("ESRCH"); }, now: () => 2000, }); assert.equal(result.repaired, true); assert.match(result.message ?? "", /Runner stderr tail:/); assert.match(result.message ?? "", /missing peer package/); const status = JSON.parse(fs.readFileSync(path.join(asyncDir, "status.json"), "utf-8")); assert.match(status.steps[0].error, /missing peer package/); const resultJson = JSON.parse(fs.readFileSync(path.join(resultsDir, "run-dead-stderr.json"), "utf-8")); assert.match(resultJson.summary, /missing peer package/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("keeps stale repair successful when the event log cannot be appended", () => { const root = tempRoot("pi-stale-event-log-collision-"); try { const asyncDir = path.join(root, "run-dead-events-dir"); const resultsDir = path.join(root, "results"); writeStatus(asyncDir, { runId: "run-dead-events-dir", sessionId: "session-current", mode: "single", state: "running", pid: 12345, startedAt: 1000, lastUpdate: 1000, currentStep: 0, steps: [{ agent: "worker", status: "running", startedAt: 1000 }], }); fs.mkdirSync(path.join(asyncDir, "events.jsonl")); const result = reconcileAsyncRun(asyncDir, { resultsDir, kill: () => { throw errno("ESRCH"); }, now: () => 2000, }); assert.equal(result.repaired, true); assert.equal(result.status?.state, "failed"); assert.equal(JSON.parse(fs.readFileSync(path.join(asyncDir, "status.json"), "utf-8")).state, "failed"); assert.equal(JSON.parse(fs.readFileSync(path.join(resultsDir, "run-dead-events-dir.json"), "utf-8")).success, false); assert.equal(fs.statSync(path.join(asyncDir, "events.jsonl")).isDirectory(), true); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("repairs stale status with per-child result outcomes", () => { const root = tempRoot("pi-stale-mixed-result-"); try { const asyncDir = path.join(root, "run-mixed"); const resultsDir = path.join(root, "results"); fs.mkdirSync(resultsDir, { recursive: true }); writeStatus(asyncDir, { runId: "run-mixed", mode: "chain", state: "running", pid: 12345, startedAt: 1000, lastUpdate: 1000, steps: [ { agent: "scout", status: "running", startedAt: 1000, model: "planned-scout" }, { agent: "worker", status: "running", startedAt: 1100, model: "planned-worker" }, ], }); const scoutSession = path.join(root, "scout.jsonl"); const workerSession = path.join(root, "worker.jsonl"); fs.writeFileSync(path.join(resultsDir, "run-mixed.json"), JSON.stringify({ id: "run-mixed", success: false, state: "failed", results: [ { agent: "scout", success: true, sessionFile: scoutSession, model: "fast" }, { agent: "worker", success: false, error: "boom", sessionFile: workerSession, model: "careful", contextOverflow: true }, ], }, null, 2), "utf-8"); const result = reconcileAsyncRun(asyncDir, { resultsDir, kill: () => { throw errno("ESRCH"); }, now: () => 2000, }); assert.equal(result.repaired, true); assert.equal(result.status?.state, "failed"); assert.equal(result.status?.steps?.[0]?.status, "complete"); assert.equal(result.status?.steps?.[0]?.exitCode, 0); assert.equal(result.status?.steps?.[0]?.model, "fast"); assert.equal(result.status?.steps?.[0]?.sessionFile, scoutSession); assert.equal(result.status?.steps?.[1]?.status, "failed"); assert.equal(result.status?.steps?.[1]?.exitCode, 1); assert.equal(result.status?.steps?.[1]?.error, "boom"); assert.equal(result.status?.steps?.[1]?.model, "careful"); assert.equal(result.status?.steps?.[1]?.contextOverflow, true); assert.equal(result.status?.steps?.[1]?.sessionFile, workerSession); assert.equal(result.status?.error, "boom"); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("retains a root failure when stale result repair writes after runner finalization", () => { const root = tempRoot("pi-stale-result-race-"); try { const runId = "run-result-race"; const staleRunnerProcessInstanceId = "runner-stale-result-race"; const finalRunnerProcessInstanceId = "runner-final-result-race"; const asyncDir = path.join(root, runId); const resultsDir = path.join(root, "results"); initializeProcessTerminal(asyncDir, runId, staleRunnerProcessInstanceId); const staleStatus = { lifecycleArtifactVersion: 3, runId, mode: "chain", state: "running", pid: 12345, startedAt: 1000, lastUpdate: 1500, processTerminal: { version: 1, state: "pending", runId, runnerProcessInstanceId: staleRunnerProcessInstanceId }, steps: [ { agent: "scout", status: "complete", startedAt: 1000, endedAt: 1200, exitCode: 0, processTerminal: { version: 1, state: "pending", runId, childIndex: 0, runnerProcessInstanceId: "runner-stale-step-0" } }, { agent: "worker", status: "failed", startedAt: 1000, endedAt: 1400, exitCode: 1, error: "child diagnostic", processTerminal: { version: 1, state: "not-started", runId, childIndex: 1, runnerProcessInstanceId: "runner-stale-step-1" } }, ], }; writeStatus(asyncDir, staleStatus); fs.mkdirSync(resultsDir, { recursive: true }); const resultPath = path.join(resultsDir, `${runId}.json`); fs.writeFileSync(resultPath, JSON.stringify({ id: runId, success: false, state: "failed", results: [ { agent: "scout", success: true }, { agent: "worker", success: false, error: "child diagnostic" }, ], }, null, 2), "utf-8"); const observedTerminal = { version: 1, state: "observed", runId, runnerProcessInstanceId: finalRunnerProcessInstanceId, observedAt: 1900, instances: [{ kind: "runner", processInstanceId: finalRunnerProcessInstanceId, closeObservedAt: 1900, exitCode: 0, signal: null }], } as const; const observedStepTerminal = { version: 1, state: "observed", runId, childIndex: 0, runnerProcessInstanceId: "runner-current-step-0", observedAt: 1901, resumeDisposition: "non-resumable", instances: [{ kind: "runner", processInstanceId: "runner-current-step-0", closeObservedAt: 1901, exitCode: 0, signal: null }], } as const; const unknownStepTerminal = { version: 1, state: "unknown", runId, childIndex: 1, runnerProcessInstanceId: "runner-current-step-1", reason: "process-tree-unverified", resumeDisposition: "non-resumable", } as const; const repaired = reconcileAsyncRun(asyncDir, { resultsDir, now: () => 2000 }, () => { fs.writeFileSync(path.join(asyncDir, "process-terminal.json"), JSON.stringify(observedTerminal), "utf-8"); writeStatus(asyncDir, { ...staleStatus, state: "failed", error: "ROOT DIAGNOSTIC", processTerminal: observedTerminal, steps: [ { ...staleStatus.steps[0], processTerminal: observedStepTerminal }, { ...staleStatus.steps[1], processTerminal: unknownStepTerminal }, ], lastUpdate: 1900, endedAt: 1900, }); }); assert.equal(repaired.repaired, true); assert.equal(repaired.status?.state, "failed"); assert.equal(repaired.status?.error, "ROOT DIAGNOSTIC"); assert.equal(repaired.status?.steps?.[0]?.status, "complete"); assert.equal(repaired.status?.steps?.[1]?.status, "failed"); assert.equal(repaired.status?.steps?.[1]?.error, "child diagnostic"); assert.deepEqual(repaired.status?.processTerminal, observedTerminal); assert.deepEqual(repaired.status?.steps?.[0]?.processTerminal, observedStepTerminal); assert.deepEqual(repaired.status?.steps?.[1]?.processTerminal, unknownStepTerminal); const persistedStatus = JSON.parse(fs.readFileSync(path.join(asyncDir, "status.json"), "utf-8")); assert.equal(persistedStatus.error, "ROOT DIAGNOSTIC"); assert.deepEqual(persistedStatus.processTerminal, observedTerminal); assert.deepEqual(persistedStatus.steps[0].processTerminal, observedStepTerminal); assert.deepEqual(persistedStatus.steps[1].processTerminal, unknownStepTerminal); const publicStatus = summarizeAsyncStatus(asyncDir, persistedStatus); assert.deepEqual(publicStatus.processTerminal, observedTerminal); assert.deepEqual(publicStatus.steps[0]?.processTerminal, observedStepTerminal); assert.deepEqual(publicStatus.steps[1]?.processTerminal, unknownStepTerminal); assert.deepEqual(readProcessTerminal(asyncDir, { runId, runnerProcessInstanceId: finalRunnerProcessInstanceId }), observedTerminal); const staleLookup = readProcessTerminal(asyncDir, { runId, runnerProcessInstanceId: staleRunnerProcessInstanceId }); assert.equal(staleLookup?.state, "unknown"); assert.equal(staleLookup?.reason, "proof-write-failed"); const proof = finalizeProcessTerminal(asyncDir, runId, { processInstanceId: finalRunnerProcessInstanceId, closeObservedAt: 2100, exitCode: 0, signal: null, }); const finalStatus = JSON.parse(fs.readFileSync(path.join(asyncDir, "status.json"), "utf-8")); assert.deepEqual(proof, observedTerminal); assert.deepEqual(finalStatus.processTerminal, proof); assert.deepEqual(finalStatus.steps[0].processTerminal, observedStepTerminal); assert.deepEqual(finalStatus.steps[1].processTerminal, unknownStepTerminal); assert.equal(finalStatus.state, "failed"); assert.equal(finalStatus.error, "ROOT DIAGNOSTIC"); assert.equal(JSON.parse(fs.readFileSync(resultPath, "utf-8")).results[1].error, "child diagnostic"); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("preserves explicit root failure precedence during result repair", () => { const root = tempRoot("pi-stale-result-error-precedence-"); try { const resultsDir = path.join(root, "results"); fs.mkdirSync(resultsDir, { recursive: true }); for (const [index, fixture] of [ { statusError: "status failure", resultError: "result failure", expected: "status failure" }, { statusError: undefined, resultError: "result failure", expected: "result failure" }, ].entries()) { const runId = `run-error-precedence-${index}`; const asyncDir = path.join(root, runId); writeStatus(asyncDir, { runId, mode: "single", state: "running", pid: 12345, startedAt: 1000, lastUpdate: 1000, ...(fixture.statusError ? { error: fixture.statusError } : {}), steps: [{ agent: "worker", status: "running", startedAt: 1000 }], }); fs.writeFileSync(path.join(resultsDir, `${runId}.json`), JSON.stringify({ id: runId, success: false, state: "failed", error: fixture.resultError, results: [{ agent: "worker", success: false, error: "child failure" }], }), "utf-8"); const repaired = reconcileAsyncRun(asyncDir, { resultsDir, now: () => 2000 }); assert.equal(repaired.status?.error, fixture.expected); } } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("repairs stale running status from an indexed pending result", () => { const root = tempRoot("pi-stale-indexed-pending-result-"); const originalError = console.error; try { console.error = () => {}; const asyncDir = path.join(root, "run-pending-result"); const resultsDir = path.join(root, "results"); writeStatus(asyncDir, { runId: "run-pending-result", sessionId: "session-current", mode: "single", state: "running", pid: 12345, startedAt: 1000, lastUpdate: 1000, currentStep: 0, steps: [{ agent: "worker", status: "running", startedAt: 1000 }], }); const publicResultPath = path.join(resultsDir, "run-pending-result.json"); fs.mkdirSync(publicResultPath, { recursive: true }); assert.deepEqual(writeAsyncResultFile(publicResultPath, { id: "run-pending-result", runId: "run-pending-result", sessionId: "session-current", success: true, state: "complete", summary: "already done", }), { state: "pending" }); const result = reconcileAsyncRun(asyncDir, { resultsDir, kill: () => { throw errno("ESRCH"); }, now: () => 2000, }); assert.equal(result.repaired, true); assert.equal(result.status?.state, "complete"); assert.equal(result.status?.steps?.[0]?.status, "complete"); assert.equal(JSON.parse(fs.readFileSync(path.join(asyncDir, "status.json"), "utf-8")).state, "complete"); assert.equal(fs.statSync(publicResultPath).isDirectory(), true); assert.match(fs.readFileSync(result.resultPath!, "utf-8"), /already done/); } finally { console.error = originalError; fs.rmSync(root, { recursive: true, force: true }); } }); it("repairs stale rejected status from an indexed pending status snapshot", () => { const root = tempRoot("pi-stale-rejected-pending-result-"); const originalError = console.error; try { console.error = () => {}; const asyncDir = path.join(root, "run-rejected"); const resultsDir = path.join(root, "results"); writeStatus(asyncDir, { runId: "run-rejected", sessionId: "session-current", mode: "chain", state: "rejected", pid: 12345, startedAt: 1000, lastUpdate: 1000, currentStep: 0, steps: [{ agent: "worker", status: "rejected", startedAt: 1000, error: "Run rejected." }], }); const publicResultPath = path.join(resultsDir, "run-rejected.json"); fs.mkdirSync(publicResultPath, { recursive: true }); writePendingAsyncResultFile(publicResultPath, { id: "run-rejected", runId: "run-rejected", sessionId: "session-current", success: false, state: "rejected", summary: "Run rejected.", results: [{ agent: "worker", success: false, error: "Run rejected." }], }); const result = reconcileAsyncRun(asyncDir, { resultsDir, kill: () => { throw errno("ESRCH"); }, now: () => 2000, }); assert.equal(result.repaired, false); assert.equal(result.status?.state, "rejected"); assert.equal(result.status?.steps?.[0]?.status, "rejected"); assert.match(fs.readFileSync(result.resultPath!, "utf-8"), /Run rejected/); } finally { console.error = originalError; fs.rmSync(root, { recursive: true, force: true }); } }); it("fails a stale run when a live pid has not updated beyond the stale threshold", () => { const root = tempRoot("pi-stale-live-pid-"); try { const asyncDir = path.join(root, "run-reused-pid"); const resultsDir = path.join(root, "results"); writeStatus(asyncDir, { runId: "run-reused-pid", mode: "single", state: "running", pid: 12345, startedAt: 1000, lastUpdate: 1000, steps: [{ agent: "worker", status: "running", startedAt: 1000 }], }); const result = reconcileAsyncRun(asyncDir, { resultsDir, kill: () => true, now: () => 5000, staleAlivePidMs: 1000, }); assert.equal(result.repaired, true); assert.equal(result.status?.state, "failed"); assert.match(result.message ?? "", /PID .* is still live; status has not updated/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("preserves an existing result instead of overwriting it with stale-run failure", () => { const root = tempRoot("pi-stale-existing-result-"); try { const asyncDir = path.join(root, "run-result"); const resultsDir = path.join(root, "results"); fs.mkdirSync(resultsDir, { recursive: true }); writeStatus(asyncDir, { runId: "run-result", mode: "single", state: "running", pid: 12345, startedAt: 1000, lastUpdate: 1000, steps: [{ agent: "worker", status: "running", startedAt: 1000 }], }); const resultPath = path.join(resultsDir, "run-result.json"); fs.writeFileSync(resultPath, JSON.stringify({ id: "run-result", success: true, state: "complete", summary: "already done" }, null, 2), "utf-8"); const result = reconcileAsyncRun(asyncDir, { resultsDir, kill: () => { throw errno("ESRCH"); }, now: () => 2000, }); assert.equal(result.repaired, true); assert.equal(result.status?.state, "complete"); assert.equal(JSON.parse(fs.readFileSync(resultPath, "utf-8")).summary, "already done"); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("does not abort the whole async-run list when one run's reconcile write throws", () => { const root = tempRoot("pi-list-resilience-"); try { const asyncDirRoot = path.join(root, "async-subagent-runs"); const asyncDir = path.join(asyncDirRoot, "run-stale"); fs.mkdirSync(path.join(asyncDirRoot, ".active-runs"), { recursive: true }); fs.writeFileSync(path.join(asyncDirRoot, ".active-runs", "run-stale"), ""); writeStatus(asyncDir, { runId: "run-stale", sessionId: "session-current", state: "running", pid: 12345, startedAt: 1000, lastUpdate: 1000, steps: [{ agent: "worker", status: "running", startedAt: 1000 }], }); // Place the results dir under a regular file so the stale run's // repair-write lookup throws ENOTDIR mid-reconcile, mirroring the // unaddressable-run failures that used to abort the whole restore. const blocker = path.join(root, "blocked"); fs.writeFileSync(blocker, "not a dir"); const badResultsDir = path.join(blocker, "async-subagent-results"); const runs = listAsyncRuns(asyncDirRoot, { states: ["queued", "running"], sessionId: "session-current", resultsDir: badResultsDir, now: () => 2000, }); // The throwing run is still surfaced via the non-repairing status read; // critically, no error escaped to abort restoring other runs. assert.equal(runs.length, 1); assert.equal(runs[0].id, "run-stale"); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); });