fix: sub-agent hang/crash stranding runs mid-phase with no save
Some checks failed
port-to-omp / port (push) Failing after 2s

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.
This commit is contained in:
2026-08-11 08:47:12 -04:00
parent 6d23b04ef6
commit 8768f9de97
7 changed files with 435 additions and 48 deletions

View File

@@ -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"`. 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 ## Commands
Every command accepts a `[path]` target (default: the current directory) and is Every command accepts a `[path]` target (default: the current directory) and is

View File

@@ -17,9 +17,32 @@
import { mkdir, writeFile } from "node:fs/promises"; import { mkdir, writeFile } from "node:fs/promises";
import { dirname, isAbsolute, join } from "node:path"; 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"; 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 { export interface AgentTaskOptions {
/** Absolute working directory for the sub-agent. */ /** Absolute working directory for the sub-agent. */
cwd: string; cwd: string;
@@ -37,6 +60,62 @@ export interface AgentTaskOptions {
* turns tool_execution_start/end + assistant turns into chat messages. * turns tool_execution_start/end + assistant turns into chat messages.
*/ */
onEvent?: (event: AgentSessionEvent) => void; 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<T>(
promise: Promise<T>,
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<typeof setTimeout> | 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 { export interface AgentRunResult {
@@ -54,8 +133,21 @@ export type AgentRunner = (opts: AgentTaskOptions) => Promise<AgentRunResult>;
let currentRunner: AgentRunner = defaultAgentRunner; let currentRunner: AgentRunner = defaultAgentRunner;
/** Entry point used by the check-runner. */ /** Entry point used by the check-runner. */
export function runAgentTask(opts: AgentTaskOptions): Promise<AgentRunResult> { export async function runAgentTask(
return currentRunner(opts); opts: AgentTaskOptions,
): Promise<AgentRunResult> {
// 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). */ /** Override the active agent runner (primarily for tests). */
@@ -121,43 +213,40 @@ export async function defaultAgentRunner(
resourceLoader: loader, 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<AgentRunResult> {
const acc: SessionEventAccumulator = { text: "" };
try { try {
let text = "";
let stopReason: string | undefined;
let errorMessage: string | undefined;
const unsubscribe = session.subscribe((event: AgentSessionEvent) => { const unsubscribe = session.subscribe((event: AgentSessionEvent) => {
if ( applySessionEvent(acc, event, opts.onEvent);
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);
}
}); });
await session.prompt(opts.task, { expandPromptTemplates: false }); await session.prompt(opts.task, { expandPromptTemplates: false });
// Ensure the agent has fully settled (tool calls may still be in-flight // 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. // Surface session errors that didn't throw but left no useful output.
// A session ending with stopReason "error" and no text means the model // A session ending with stopReason "error" and no text means the model
// call failed silently — treat that as a failed run, not ok:true. // call failed silently — treat that as a failed run, not ok:true.
if (errorMessage) { if (acc.errorMessage) {
return { ok: false, text, error: errorMessage }; return { ok: false, text: acc.text, error: acc.errorMessage };
} }
if (!text.trim() && stopReason === "error") { if (!acc.text.trim() && acc.stopReason === "error") {
return { return {
ok: false, ok: false,
text, text: acc.text,
error: "sub-agent session ended in error with no output.", error: "sub-agent session ended in error with no output.",
}; };
} }
return { ok: true, text }; return { ok: true, text: acc.text };
} catch (err) { } catch (err) {
return { return {
ok: false, 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. * Extract joined text from an assistant message's content blocks.
* Mirrors piolium's `extractAssistantText`. * Mirrors piolium's `extractAssistantText`.

View File

@@ -92,8 +92,11 @@ ${scopeRulesMarkdown()}
3. Apply the rubric below to each comment and classify it: RESTATE, VERBOSE, 3. Apply the rubric below to each comment and classify it: RESTATE, VERBOSE,
WHY, or OK. WHY, or OK.
4. Write a findings report to \`${findingsFile}\` with per-file line refs. 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 5. Return a ONE-LINE summary as your final message, e.g.
file). The host captures it as the analysis-phase findings. \`<count> comment smell(s) across <files> 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} ${RUBRIC}
@@ -110,7 +113,8 @@ ${RUBRIC}
\`\`\` \`\`\`
If no smells are found, write \`# comments — findings\n\n0 comment smell(s).\` 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). Write the report under \`${outDir}\` (create directories as needed).
`; `;
@@ -128,6 +132,7 @@ export function buildCommentsFixTask(
): string { ): string {
const outDir = commentsArtifactDir(scope); const outDir = commentsArtifactDir(scope);
const changesFile = changesPath(scope); const changesFile = changesPath(scope);
const findingsFile = findingsPath(scope);
return `# Task: comments hygiene fix return `# Task: comments hygiene fix
You are running the **comments** hygiene fix phase. You are running the **comments** hygiene fix phase.
@@ -136,6 +141,9 @@ You are running the **comments** hygiene fix phase.
- Fix target: \`${scope.target}\` - Fix target: \`${scope.target}\`
## Input: scan findings ## 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)"} ${findings.trim().length > 0 ? findings : "(no findings text provided)"}
## What to do ## What to do

View File

@@ -142,7 +142,12 @@ export async function runCheck(
): Promise<CheckRunOutcome> { ): Promise<CheckRunOutcome> {
const startMs = Date.now(); const startMs = Date.now();
const outcome = await runCheckImpl(opts); 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; return outcome;
} }
@@ -383,7 +388,16 @@ async function runCheckImplInner(
error = err instanceof Error ? err.message : String(err); error = err instanceof Error ? err.message : String(err);
markCheckStatus(state, check.name, "failed", error); markCheckStatus(state, check.name, "failed", error);
markRunStatus(state, reconcileRunStatus(state)); 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 }; return { status: "failed", error, findings, changes, state };
} finally { } finally {
strip.done(); strip.done();

137
tests/agent-runner.test.ts Normal file
View File

@@ -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");
});
});

View File

@@ -20,6 +20,8 @@ import {
setAgentRunner, setAgentRunner,
resetAgentRunner, resetAgentRunner,
fakeAgentRunner, fakeAgentRunner,
AGENT_TIMEOUT_ENV,
type AgentRunResult,
} from "../src/agent-runner.js"; } from "../src/agent-runner.js";
import { handleCheckCommand, type PygieniumCtx } from "../src/commands.js"; import { handleCheckCommand, type PygieniumCtx } from "../src/commands.js";
import { loadRunState, runStatePath } from "../src/run-state.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.status).toBe("skipped");
expect(state?.checks.gated.error).toBe("no source files matched"); 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<AgentRunResult>(() => {}));
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");
});
}); });

View File

@@ -30,6 +30,8 @@ import {
commentsCheck, commentsCheck,
findingsPath, findingsPath,
changesPath, changesPath,
buildCommentsScanTask,
buildCommentsFixTask,
} from "../src/checks/comments.js"; } from "../src/checks/comments.js";
import type { CheckScope } from "../src/checks/registry.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!.findings).toBeDefined();
expect(check!.changes).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");
});
}); });