From 8768f9de976af5dfd67d5fd5cf8af2cda1ede17e Mon Sep 17 00:00:00 2001 From: Michael Freno Date: Tue, 11 Aug 2026 08:47:12 -0400 Subject: [PATCH] fix: sub-agent hang/crash stranding runs mid-phase with no save MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A stalled sub-agent (provider stream never settling after the final tool call, or a throwing session-event listener recursing through the SDK's run-failure path to stack overflow) left the phase stuck in_progress with no state save, no error, and no completion message — observed twice in freno-dev, both times after comments wrote findings.md. - agent-runner: 60-min settle watchdog on every agent phase (PYGIENIUM_AGENT_TIMEOUT_MS, env-tunable); timeout disposes the session and fails the phase loudly with a /pygienium-resume hint instead of hanging the check-runner's await forever. - agent-runner: applySessionEvent — the session.subscribe listener can no longer throw into the SDK event pipeline (guards for partial/malformed events, optional-chained message_update, throwing chat forwarder). - comments: scan no longer regenerates the full report as its final message (findings.md is the source of truth); fix phase reads findings.md with embedded fallback. - check-runner: failing state-save inside the catch path can't double- fault or escape as an unhandled rejection; completion posting is best-effort. - tests: watchdog timeout test, comments task-text regression tests, applySessionEvent malformed-shape unit tests. 146 pass, tsc clean. - README: document PYGIENIUM_AGENT_TIMEOUT_MS. --- README.md | 6 + src/agent-runner.ts | 252 ++++++++++++++++++++++++++++++------- src/checks/comments.ts | 14 ++- src/modes/check-runner.ts | 18 ++- tests/agent-runner.test.ts | 137 ++++++++++++++++++++ tests/check-runner.test.ts | 27 ++++ tests/comments.test.ts | 29 +++++ 7 files changed, 435 insertions(+), 48 deletions(-) create mode 100644 tests/agent-runner.test.ts diff --git a/README.md b/README.md index cb0392d..47e7c91 100644 --- a/README.md +++ b/README.md @@ -36,6 +36,12 @@ Pygienium reads its chat-rendering style from pi's `settings.json` (`~/.pi/agent No entry means `"verbose"` (the default). An unreadable or missing `settings.json` also falls back to `"verbose"`. +Environment variable (read at run time): + +| Env var | Default | Description | +| --- | --- | --- | +| `PYGIENIUM_AGENT_TIMEOUT_MS` | `3600000` (60 min) | Per-agent-phase settle deadline. A sub-agent session that never settles (stalled provider stream, hung retry/auto-compaction after its last tool call) is aborted and the phase fails with a visible error instead of hanging the run mid-transition with no state update and no completion message. Lower it (e.g. `1200000`) to fail fast on flaky providers; raise it for very large repos. + ## Commands Every command accepts a `[path]` target (default: the current directory) and is diff --git a/src/agent-runner.ts b/src/agent-runner.ts index 4c0ed32..2a182ae 100644 --- a/src/agent-runner.ts +++ b/src/agent-runner.ts @@ -17,9 +17,32 @@ import { mkdir, writeFile } from "node:fs/promises"; import { dirname, isAbsolute, join } from "node:path"; -import type { AgentSessionEvent } from "@earendil-works/pi-coding-agent"; +import type { + AgentSession, + AgentSessionEvent, +} from "@earendil-works/pi-coding-agent"; import { loadAgents, extensionRoot, type AgentDef } from "./agents.js"; +/** + * Env var overriding the per-agent settle timeout (ms). Guards against a + * sub-agent session that never settles (stalled provider stream, hung retry / + * auto-compaction after the final tool call), which otherwise leaves the run + * stuck mid-phase with no state save, no error, and no completion message. + */ +export const AGENT_TIMEOUT_ENV = "PYGIENIUM_AGENT_TIMEOUT_MS"; + +/** Default settle timeout per agent phase: 60 minutes. */ +const DEFAULT_AGENT_TIMEOUT_MS = 60 * 60_000; + +/** Resolve the per-agent settle timeout, honouring the env override. */ +export function agentTimeoutMs(): number { + const raw = process.env[AGENT_TIMEOUT_ENV]; + if (raw && /^\d+$/.test(raw.trim()) && Number(raw.trim()) > 0) { + return Number(raw.trim()); + } + return DEFAULT_AGENT_TIMEOUT_MS; +} + export interface AgentTaskOptions { /** Absolute working directory for the sub-agent. */ cwd: string; @@ -37,6 +60,62 @@ export interface AgentTaskOptions { * turns tool_execution_start/end + assistant turns into chat messages. */ onEvent?: (event: AgentSessionEvent) => void; + /** + * Maximum time the agent run may take before it is aborted and the phase + * fails loudly (default {@link agentTimeoutMs}). A session that never + * settles — stalled provider stream, hung retry/compaction after its last + * tool call — would otherwise hang the check run silently with no state + * update and no completion message. + */ + timeoutMs?: number; +} + +/** + * Race `promise` against a settle deadline. Returns `{ value }` on success or + * `{ error }` when the deadline elapsed first (the caller aborts the work). + * `promise` is still awaited-then-ignored afterwards so late rejections can + * never surface as unhandled. + */ +export async function withSettleTimeout( + promise: Promise, + timeoutMs: number, + label: string, +): Promise<{ value: T } | { error: string }> { + if (timeoutMs <= 0) { + try { + return { value: await promise }; + } catch (err) { + return { error: err instanceof Error ? err.message : String(err) }; + } + } + const settled = promise.then( + (value) => ({ value }) as { value: T }, + (err) => + ({ error: err instanceof Error ? err.message : String(err) }) as { + error: string; + }, + ); + let timer: ReturnType | undefined; + const deadline = new Promise<{ error: string }>((resolve) => { + timer = setTimeout(() => { + // ~`1m 30s` / `45s` for the message (timeoutMs is ms; tests use small values). + const totalSec = Math.round(timeoutMs / 1000); + const duration = + totalSec >= 60 + ? `${Math.floor(totalSec / 60)}m${totalSec % 60 ? ` ${totalSec % 60}s` : ""}` + : `${totalSec}s`; + resolve({ + error: `${label} did not settle within ${duration}; aborted. Check model/provider connectivity, then resume with /pygienium-resume.`, + }); + }, timeoutMs); + // Never hold the process open just because a deadline is pending. + timer.unref?.(); + }); + try { + return await Promise.race([settled, deadline]); + } finally { + if (timer) clearTimeout(timer); + } } export interface AgentRunResult { @@ -54,8 +133,21 @@ export type AgentRunner = (opts: AgentTaskOptions) => Promise; let currentRunner: AgentRunner = defaultAgentRunner; /** Entry point used by the check-runner. */ -export function runAgentTask(opts: AgentTaskOptions): Promise { - return currentRunner(opts); +export async function runAgentTask( + opts: AgentTaskOptions, +): Promise { + // Backstop: even a third-party/custom runner must not be able to hang the + // run forever. The default runner additionally aborts its session on + // timeout (see defaultAgentRunner); this race covers every other runner. + const settled = await withSettleTimeout( + currentRunner(opts), + opts.timeoutMs ?? agentTimeoutMs(), + `sub-agent "${opts.agentName}"`, + ); + if ("error" in settled) { + return { ok: false, text: "", error: settled.error }; + } + return settled.value; } /** Override the active agent runner (primarily for tests). */ @@ -121,43 +213,40 @@ export async function defaultAgentRunner( resourceLoader: loader, }); + // A session that never settles (stalled provider stream, hung retry / + // auto-compaction after its last tool call) must fail the phase loudly + // instead of hanging the run mid-transition with zero diagnostics. On + // timeout the session is disposed, which aborts the in-flight run. + const settled = await withSettleTimeout( + runSessionToCompletion(session, opts), + opts.timeoutMs ?? agentTimeoutMs(), + `sub-agent "${opts.agentName}"`, + ); + if ("error" in settled) { + // Abort the still-running session so it can't keep burning provider + // calls; the background settle path then finishes and disposes too. + try { + session.dispose(); + } catch { + /* ignore dispose errors */ + } + return { ok: false, text: "", error: settled.error }; + } + return settled.value; +} + +/** + * Run an in-memory agent session to completion, streaming tool events to the + * chat forwarder, and return the agent's final text. Owns session cleanup. + */ +async function runSessionToCompletion( + session: AgentSession, + opts: AgentTaskOptions, +): Promise { + const acc: SessionEventAccumulator = { text: "" }; try { - let text = ""; - let stopReason: string | undefined; - let errorMessage: string | undefined; const unsubscribe = session.subscribe((event: AgentSessionEvent) => { - if ( - event.type === "message_update" && - event.assistantMessageEvent.type === "text_delta" - ) { - text += event.assistantMessageEvent.delta; - } - if (event.type === "message_end") { - // Capture the full assistant text from the finalized message — - // models that don't stream text_delta (or truncate) still surface - // their output here. Prefer the streamed text when non-empty. - const message = event.message as { - role?: string; - content?: unknown; - stopReason?: string; - errorMessage?: string; - }; - if (message.stopReason) stopReason = message.stopReason; - if (message.errorMessage) errorMessage = message.errorMessage; - if (message.role === "assistant") { - const full = extractAssistantText(message.content).trim(); - if (full && !text.trim()) text = full; - } - } - // Forward the stream-driving events to the chat forwarder; it turns - // each into its own `pygienium-stream` message (see index.ts). - if ( - event.type === "tool_execution_start" || - event.type === "tool_execution_end" || - event.type === "message_end" - ) { - opts.onEvent?.(event); - } + applySessionEvent(acc, event, opts.onEvent); }); await session.prompt(opts.task, { expandPromptTemplates: false }); // Ensure the agent has fully settled (tool calls may still be in-flight @@ -168,17 +257,17 @@ export async function defaultAgentRunner( // Surface session errors that didn't throw but left no useful output. // A session ending with stopReason "error" and no text means the model // call failed silently — treat that as a failed run, not ok:true. - if (errorMessage) { - return { ok: false, text, error: errorMessage }; + if (acc.errorMessage) { + return { ok: false, text: acc.text, error: acc.errorMessage }; } - if (!text.trim() && stopReason === "error") { + if (!acc.text.trim() && acc.stopReason === "error") { return { ok: false, - text, + text: acc.text, error: "sub-agent session ended in error with no output.", }; } - return { ok: true, text }; + return { ok: true, text: acc.text }; } catch (err) { return { ok: false, @@ -194,6 +283,83 @@ export async function defaultAgentRunner( } } +/** Running capture state while a session streams assistant output. */ +export interface SessionEventAccumulator { + /** Joined assistant text seen so far (text_delta stream). */ + text: string; + /** stopReason of the final assistant message, when reported. */ + stopReason?: string; + /** errorMessage of the final assistant message, when reported. */ + errorMessage?: string; +} + +/** + * Interpret one session event into the running accumulator and forward the + * stream-driving events to the chat. + * + * MUST never throw: it runs synchronously inside the SDK's event pipeline + * (`session.subscribe` listeners are invoked from `processEvents`/`_emit`). + * A throw there is not contained — the SDK's run-failure path re-emits + * failure events through the same callback, so a callback that always throws + * recurses until stack overflow and crashes the host, stranding run-state at + * the phase boundary with no save and no completion. Guard every shape (the + * runtime events are partial/streaming and weaker than their types) and never + * let the cosmetic chat forwarder take the run down. + */ +export function applySessionEvent( + acc: SessionEventAccumulator, + event: AgentSessionEvent, + forward?: (event: AgentSessionEvent) => void, +): void { + if (!event) return; + try { + if ( + event.type === "message_update" && + event.assistantMessageEvent?.type === "text_delta" + ) { + acc.text += event.assistantMessageEvent.delta ?? ""; + } + if (event.type === "message_end") { + // Capture the full assistant text from the finalized message — + // models that don't stream text_delta (or truncate) still surface + // their output here. Prefer the streamed text when non-empty. + const message = event.message as + | { + role?: string; + content?: unknown; + stopReason?: string; + errorMessage?: string; + } + | undefined; + if (message) { + if (message.stopReason) acc.stopReason = message.stopReason; + if (message.errorMessage) acc.errorMessage = message.errorMessage; + if (message.role === "assistant") { + const full = extractAssistantText(message.content).trim(); + if (full && !acc.text.trim()) acc.text = full; + } + } + } + } catch (err) { + // A malformed event must not crash the SDK event pipeline; log and skip. + console.error("pygienium: error processing sub-agent event", err); + } + // Forward the stream-driving events to the chat forwarder; it turns each + // into its own `pygienium-stream` message (see index.ts). Cosmetic chat + // UI: a broken forwarder must not fail or hang the agent run. + if ( + event.type === "tool_execution_start" || + event.type === "tool_execution_end" || + event.type === "message_end" + ) { + try { + forward?.(event); + } catch (err) { + console.error("pygienium: error forwarding sub-agent event", err); + } + } +} + /** * Extract joined text from an assistant message's content blocks. * Mirrors piolium's `extractAssistantText`. diff --git a/src/checks/comments.ts b/src/checks/comments.ts index 9a735f1..a9339de 100644 --- a/src/checks/comments.ts +++ b/src/checks/comments.ts @@ -92,8 +92,11 @@ ${scopeRulesMarkdown()} 3. Apply the rubric below to each comment and classify it: RESTATE, VERBOSE, WHY, or OK. 4. Write a findings report to \`${findingsFile}\` with per-file line refs. -5. Return the findings report text as your final message (same content as the - file). The host captures it as the analysis-phase findings. +5. Return a ONE-LINE summary as your final message, e.g. + \` comment smell(s) across file(s).\` The full report lives + in findings.md — do NOT regenerate the report text in your final message + (regenerating a large report doubles the fragile output right after the + file write and can stall the run). ${RUBRIC} @@ -110,7 +113,8 @@ ${RUBRIC} \`\`\` If no smells are found, write \`# comments — findings\n\n0 comment smell(s).\` -and return that text. Always create findings.md so the run has an artifact. +and return \`0 comment smell(s).\` as your final message. Always create +findings.md so the run has an artifact. Write the report under \`${outDir}\` (create directories as needed). `; @@ -128,6 +132,7 @@ export function buildCommentsFixTask( ): string { const outDir = commentsArtifactDir(scope); const changesFile = changesPath(scope); + const findingsFile = findingsPath(scope); return `# Task: comments hygiene fix You are running the **comments** hygiene fix phase. @@ -136,6 +141,9 @@ You are running the **comments** hygiene fix phase. - Fix target: \`${scope.target}\` ## Input: scan findings +Read the detailed per-file findings from \`${findingsFile}\` (created by the +scan phase). If that file is missing or unreadable, fall back to the findings +text below: ${findings.trim().length > 0 ? findings : "(no findings text provided)"} ## What to do diff --git a/src/modes/check-runner.ts b/src/modes/check-runner.ts index 5658ea4..55221fd 100644 --- a/src/modes/check-runner.ts +++ b/src/modes/check-runner.ts @@ -142,7 +142,12 @@ export async function runCheck( ): Promise { const startMs = Date.now(); const outcome = await runCheckImpl(opts); - postCheckCompletion(opts, outcome, Date.now() - startMs); + try { + postCheckCompletion(opts, outcome, Date.now() - startMs); + } catch { + // Completion posting is best-effort chat UI: a renderer or send + // failure must never reject the run after its state was persisted. + } return outcome; } @@ -383,7 +388,16 @@ async function runCheckImplInner( error = err instanceof Error ? err.message : String(err); markCheckStatus(state, check.name, "failed", error); markRunStatus(state, reconcileRunStatus(state)); - await saveRunState(state); + try { + await saveRunState(state); + } catch (saveErr) { + // A failing state save inside the error path must not mask the + // original error or escape as an unhandled rejection (which would + // kill the run with nothing persisted or reported). + error += `; (also failed to persist run-state: ${ + saveErr instanceof Error ? saveErr.message : String(saveErr) + })`; + } return { status: "failed", error, findings, changes, state }; } finally { strip.done(); diff --git a/tests/agent-runner.test.ts b/tests/agent-runner.test.ts new file mode 100644 index 0000000..45251b3 --- /dev/null +++ b/tests/agent-runner.test.ts @@ -0,0 +1,137 @@ +/** + * agent-runner.test.ts — unit tests for the session event accumulator. + * + * `applySessionEvent` runs synchronously inside the SDK's event pipeline; a + * throw there crashes the host (the SDK's run-failure path re-emits events + * through the same callback → stack overflow), stranding run-state at the + * phase boundary. These tests pin the handler to never throw on the partial / + * malformed event shapes streaming sessions actually emit. + */ +import { describe, expect, it } from "bun:test"; +import { + applySessionEvent, + type SessionEventAccumulator, +} from "../src/agent-runner.js"; +import type { AgentSessionEvent } from "@earendil-works/pi-coding-agent"; + +function fresh(): SessionEventAccumulator { + return { text: "" }; +} + +/** Build a typed-as-unknown event so malformed shapes compile in tests. */ +function event(shape: unknown): AgentSessionEvent { + return shape as AgentSessionEvent; +} + +describe("applySessionEvent", () => { + it("accumulates text_delta stream events in order", () => { + const acc = fresh(); + applySessionEvent( + acc, + event({ + type: "message_update", + assistantMessageEvent: { type: "text_delta", delta: "foo" }, + }), + ); + applySessionEvent( + acc, + event({ + type: "message_update", + assistantMessageEvent: { type: "text_delta", delta: "bar" }, + }), + ); + expect(acc.text).toBe("foobar"); + }); + + it("does not throw on a message_update with no assistantMessageEvent", () => { + const acc = fresh(); + expect(() => + applySessionEvent(acc, event({ type: "message_update" })), + ).not.toThrow(); + expect(acc.text).toBe(""); + }); + + it("does not throw on a message_update with an unknown event shape", () => { + const acc = fresh(); + expect(() => applySessionEvent(acc, event({ type: "bogus_event" }))).not.toThrow(); + expect(() => applySessionEvent(acc, event(null))).not.toThrow(); + expect(() => applySessionEvent(acc, event(undefined))).not.toThrow(); + expect(acc.text).toBe(""); + }); + + it("captures full text + stopReason from message_end when nothing streamed", () => { + const acc = fresh(); + applySessionEvent( + acc, + event({ + type: "message_end", + message: { + role: "assistant", + stopReason: "stop", + content: [{ type: "text", text: "full report" }], + }, + }), + ); + expect(acc.text).toBe("full report"); + expect(acc.stopReason).toBe("stop"); + }); + + it("does not throw on a message_end with no message", () => { + const acc = fresh(); + expect(() => applySessionEvent(acc, event({ type: "message_end" }))).not.toThrow(); + expect(acc.text).toBe(""); + }); + + it("records an errorMessage surfaced on the final message", () => { + const acc = fresh(); + applySessionEvent( + acc, + event({ + type: "message_end", + message: { role: "assistant", errorMessage: "upstream 529" }, + }), + ); + expect(acc.errorMessage).toBe("upstream 529"); + }); + + it("forwards stream-driving events and swallows a throwing forwarder", () => { + const originalError = console.error; + console.error = () => {}; + try { + const forwarded: string[] = []; + const forward = (ev: AgentSessionEvent) => { + forwarded.push(ev.type); + if (ev.type === "tool_execution_end") throw new Error("renderer boom"); + }; + const acc = fresh(); + applySessionEvent( + acc, + event({ type: "tool_execution_start", toolName: "bash" }), + forward, + ); + applySessionEvent( + acc, + event({ type: "tool_execution_end", toolName: "bash" }), + forward, + ); + expect(forwarded).toEqual(["tool_execution_start", "tool_execution_end"]); + } finally { + console.error = originalError; + } + }); + + it("does not forward non-stream-driving events", () => { + const forwarded: string[] = []; + const acc = fresh(); + applySessionEvent( + acc, + event({ type: "message_update", assistantMessageEvent: { type: "text_delta", delta: "x" } }), + (ev) => forwarded.push(ev.type), + ); + applySessionEvent(acc, event({ type: "agent_settled" }), (ev) => + forwarded.push(ev.type), + ); + expect(forwarded).toEqual([]); + expect(acc.text).toBe("x"); + }); +}); diff --git a/tests/check-runner.test.ts b/tests/check-runner.test.ts index a02c2e1..a38a519 100644 --- a/tests/check-runner.test.ts +++ b/tests/check-runner.test.ts @@ -20,6 +20,8 @@ import { setAgentRunner, resetAgentRunner, fakeAgentRunner, + AGENT_TIMEOUT_ENV, + type AgentRunResult, } from "../src/agent-runner.js"; import { handleCheckCommand, type PygieniumCtx } from "../src/commands.js"; import { loadRunState, runStatePath } from "../src/run-state.js"; @@ -128,4 +130,29 @@ describe("check-runner integration", () => { expect(state?.checks.gated.status).toBe("skipped"); expect(state?.checks.gated.error).toBe("no source files matched"); }); + + it("fails the check loudly when the agent never settles", async () => { + // Models the silent hang: the sub-agent session never resolves after + // its last tool call (stalled provider stream / hung retry), leaving + // the run stuck mid-phase with no state save and no completion. The + // watchdog must abort the phase with a visible error instead. + setAgentRunner(() => new Promise(() => {})); + const prev = process.env[AGENT_TIMEOUT_ENV]; + process.env[AGENT_TIMEOUT_ENV] = "50"; + try { + await handleCheckCommand(smokeCheck(), "", stubCtx(cwd)); + } finally { + if (prev === undefined) delete process.env[AGENT_TIMEOUT_ENV]; + else process.env[AGENT_TIMEOUT_ENV] = prev; + } + + const state = await loadRunState(cwd); + expect(state?.checks.smoke.status).toBe("failed"); + expect(state?.checks.smoke.error).toContain("did not settle within"); + expect(state?.checks.smoke.error).toContain("/pygienium-resume"); + expect( + state?.checks.smoke.phases.find((p) => p.id === "analysis")?.status, + ).toBe("failed"); + expect(state?.status).toBe("failed"); + }); }); diff --git a/tests/comments.test.ts b/tests/comments.test.ts index 1a70c4a..ba058ea 100644 --- a/tests/comments.test.ts +++ b/tests/comments.test.ts @@ -30,6 +30,8 @@ import { commentsCheck, findingsPath, changesPath, + buildCommentsScanTask, + buildCommentsFixTask, } from "../src/checks/comments.js"; import type { CheckScope } from "../src/checks/registry.js"; @@ -256,4 +258,31 @@ describe("comments check (end-to-end)", () => { expect(check!.findings).toBeDefined(); expect(check!.changes).toBeDefined(); }); + + it("scan task keeps the full report in findings.md, not the final message", () => { + // Regression: the scan task used to demand the agent regenerate the + // whole report as its final message right after writing findings.md — + // a second huge output that stalled the phase transition (observed in + // freno-dev twice). The final message must stay a one-line summary. + const scope: CheckScope = { + cwd, + target, + fix: false, + rest: [], + }; + const task = buildCommentsScanTask(cwd, scope); + expect(task).toContain("ONE-LINE summary"); + expect(task).toContain("do NOT regenerate the report text"); + expect(task).not.toContain( + "Return the findings report text as your final message", + ); + expect(task).not.toContain("same content as the file"); + + // The fix phase must read the detailed findings from the artifact so a + // one-line scan summary can't starve it. + const fixTask = buildCommentsFixTask(cwd, scope, "fallback-findings-text"); + expect(fixTask).toContain("Read the detailed per-file findings"); + expect(fixTask).toContain(findingsPath(scope)); + expect(fixTask).toContain("fallback-findings-text"); + }); });