fix(coding-agent): fixed bash output integrity, job lifecycle, and interception

artifact spill now includes the head-retained bytes (full capture was missing first ~20KB); chunk throttle coalesces instead of dropping; cd-prefix extraction defers shell-expanded paths; interceptor rule is quote-aware and catches clobber and variable targets; completed async jobs release their Shell; at job cap commands degrade to foreground; PTY mode drops the non-interactive env and notes silent downgrades; timeout/abort annotations always appended; removed dead idle-timeout-watchdog.
This commit is contained in:
can1357
2026-06-10 01:27:16 +02:00
parent c3054d5cea
commit a13e9827f4
10 changed files with 248 additions and 147 deletions
+57 -3
View File
@@ -23,6 +23,12 @@ export interface AsyncJob {
* supply an id (e.g. legacy tests, SDK consumers without an agent context).
*/
ownerId?: string;
/**
* Job is registered but parked behind a caller-managed gate (e.g. a task
* batch semaphore). Queued jobs do not count toward the running-job limit
* until the caller invokes `markRunning()` from the run context.
*/
queued?: boolean;
}
export interface AsyncJobManagerOptions {
@@ -53,6 +59,8 @@ export interface AsyncJobRegisterOptions {
/** Registry id of the agent that owns this job; used to scope cancelAll. */
ownerId?: string;
onProgress?: (text: string, details?: Record<string, unknown>) => void | Promise<void>;
/** Register the job in queued state; see {@link AsyncJob.queued}. */
queued?: boolean;
}
/**
@@ -110,6 +118,17 @@ export class AsyncJobManager {
this.#retentionMs = Math.max(0, Math.floor(options.retentionMs ?? DEFAULT_RETENTION_MS));
}
/** True when the running-job count has reached the configured cap. */
get atCapacity(): boolean {
if (this.#disposed) return true;
// Mirror register(): queued jobs hold no execution slot.
let activeCount = 0;
for (const job of this.#jobs.values()) {
if (job.status === "running" && !job.queued) activeCount++;
}
return activeCount >= this.#maxRunningJobs;
}
register(
type: "bash" | "task",
label: string,
@@ -117,14 +136,21 @@ export class AsyncJobManager {
jobId: string;
signal: AbortSignal;
reportProgress: (text: string, details?: Record<string, unknown>) => Promise<void>;
/** Clear the queued flag once the job actually starts executing. */
markRunning: () => void;
}) => Promise<string>,
options?: AsyncJobRegisterOptions,
): string {
if (this.#disposed) {
throw new Error("Async job manager is disposed");
}
const runningCount = this.getRunningJobs().length;
if (runningCount >= this.#maxRunningJobs) {
// Queued jobs hold no execution slot yet — only count jobs that are
// actually running so a large parked batch cannot starve registration.
let activeCount = 0;
for (const existing of this.#jobs.values()) {
if (existing.status === "running" && !existing.queued) activeCount++;
}
if (activeCount >= this.#maxRunningJobs) {
throw new Error(
`Background job limit reached (${this.#maxRunningJobs}). Wait for running jobs to finish or cancel one.`,
);
@@ -144,6 +170,7 @@ export class AsyncJobManager {
abortController,
promise: Promise.resolve(),
ownerId: options?.ownerId,
queued: options?.queued === true,
};
const reportProgress = async (text: string, details?: Record<string, unknown>): Promise<void> => {
@@ -159,7 +186,14 @@ export class AsyncJobManager {
};
job.promise = (async () => {
try {
const text = await run({ jobId: id, signal: abortController.signal, reportProgress });
const text = await run({
jobId: id,
signal: abortController.signal,
reportProgress,
markRunning: () => {
job.queued = false;
},
});
if (job.status === "cancelled") {
job.resultText = text;
this.#scheduleEviction(id);
@@ -278,6 +312,26 @@ export class AsyncJobManager {
return before - this.#deliveries.length;
}
/**
* Lift a foreground-wait suppression set via `acknowledgeDeliveries`. If the
* job already finished while suppressed (its delivery enqueue was skipped),
* re-enqueue the completion so the result is still delivered exactly once.
*/
resumeDeliveries(jobIds: string[]): void {
for (const rawId of jobIds) {
const jobId = rawId.trim();
if (!jobId) continue;
if (!this.#suppressedDeliveries.delete(jobId)) continue;
const job = this.#jobs.get(jobId);
if (!job || (job.status !== "completed" && job.status !== "failed")) continue;
const queued =
this.#deliveries.some(delivery => delivery.jobId === jobId) ||
this.#inFlightDeliveries.some(delivery => delivery.jobId === jobId);
if (queued) continue;
this.#enqueueDelivery(jobId, job.status === "completed" ? (job.resultText ?? "") : (job.errorText ?? ""));
}
}
/**
* Cancel running jobs. With `filter.ownerId` set, cancels only jobs the
* matching agent registered; with no filter, cancels every running job
@@ -246,7 +246,11 @@ export const DEFAULT_BASH_INTERCEPTOR_RULES: BashInterceptorRule[] = [
message: "Use the `edit` tool instead of awk -i inplace. It provides diff preview and fuzzy matching.",
},
{
pattern: "^\\s*(echo|printf|cat\\s*<<)\\s+.*[^|]>\\s*\\S",
// `>` must sit outside quoted regions (so `echo "a -> b"` passes) and be
// followed by a plausible filename — including `$VAR` targets; `>|`
// (clobber) counts as a redirect; `>&2`/`2>&1` style fd duplication is
// not matched.
pattern: "^\\s*(echo|printf|cat\\s*<<)\\s+(?:[^\"'>]|\"[^\"]*\"|'[^']*')*(?<!\\|)>{1,2}\\|?\\s*[$\\w./~\"'-]",
tool: "write",
message: "Use the `write` tool instead of echo/cat redirection. It handles encoding and provides confirmation.",
},
@@ -314,7 +314,9 @@ export async function executeBash(command: string, options?: BashExecutorOptions
if (userSignal) {
userSignal.removeEventListener("abort", abortHandler);
}
if (resetSession) {
if (resetSession || options?.sessionKey?.includes(":async:")) {
// `:async:` keys are per-job (jobId is unique), so the Shell would
// otherwise stay in the process-global map forever after completion.
shellSessions.delete(sessionKey);
}
}
@@ -1,126 +0,0 @@
export type ExecutionAbortReason = "idle-timeout" | "signal";
export interface IdleTimeoutWatchdogOptions {
timeoutMs?: number;
signal?: AbortSignal;
hardTimeoutGraceMs: number;
onAbort?: (reason: ExecutionAbortReason) => void;
}
export class IdleTimeoutWatchdog {
#abortController = new AbortController();
#abortReason?: ExecutionAbortReason;
#hardTimeoutDeferred = Promise.withResolvers<"hard-timeout">();
#hardTimeoutGraceMs: number;
#hardTimeoutTimer?: NodeJS.Timeout;
#idleTimer?: NodeJS.Timeout;
#onAbort?: (reason: ExecutionAbortReason) => void;
#signal?: AbortSignal;
#signalAbortHandler?: () => void;
#timeoutMs?: number;
constructor(options: IdleTimeoutWatchdogOptions) {
this.#timeoutMs = options.timeoutMs;
this.#hardTimeoutGraceMs = options.hardTimeoutGraceMs;
this.#onAbort = options.onAbort;
this.#signal = options.signal;
if (this.#signal) {
if (this.#signal.aborted) {
this.#abort("signal");
return;
}
this.#signalAbortHandler = () => {
this.#abort("signal");
};
this.#signal.addEventListener("abort", this.#signalAbortHandler, { once: true });
}
this.touch();
}
get abortedBySignal(): boolean {
return this.#abortReason === "signal";
}
get hardTimeoutPromise(): Promise<"hard-timeout"> {
return this.#hardTimeoutDeferred.promise;
}
get signal(): AbortSignal {
return this.#abortController.signal;
}
get timedOut(): boolean {
return this.#abortReason === "idle-timeout";
}
touch(): void {
if (this.#abortReason || this.#timeoutMs === undefined || this.#timeoutMs <= 0) {
return;
}
if (this.#idleTimer) {
clearTimeout(this.#idleTimer);
}
this.#idleTimer = setTimeout(() => {
this.#abort("idle-timeout");
}, this.#timeoutMs);
}
dispose(): void {
if (this.#idleTimer) {
clearTimeout(this.#idleTimer);
this.#idleTimer = undefined;
}
if (this.#hardTimeoutTimer) {
clearTimeout(this.#hardTimeoutTimer);
this.#hardTimeoutTimer = undefined;
}
if (this.#signal && this.#signalAbortHandler) {
this.#signal.removeEventListener("abort", this.#signalAbortHandler);
this.#signalAbortHandler = undefined;
}
}
#abort(reason: ExecutionAbortReason): void {
if (this.#abortReason) {
return;
}
this.#abortReason = reason;
if (this.#idleTimer) {
clearTimeout(this.#idleTimer);
this.#idleTimer = undefined;
}
if (!this.#abortController.signal.aborted) {
this.#abortController.abort(reason);
}
this.#onAbort?.(reason);
this.#armHardTimeout();
}
#armHardTimeout(): void {
if (this.#hardTimeoutTimer || this.#hardTimeoutGraceMs <= 0) {
return;
}
this.#hardTimeoutTimer = setTimeout(() => {
this.#hardTimeoutDeferred.resolve("hard-timeout");
}, this.#hardTimeoutGraceMs);
}
}
export function formatIdleTimeoutMessage(timeoutMs?: number): string {
if (timeoutMs === undefined) {
return "Command timed out without output";
}
const seconds = Math.max(1, Math.round(timeoutMs / 1000));
return `Command timed out after ${seconds} seconds without output`;
}
@@ -650,6 +650,7 @@ export class OutputSink {
#sawData = false;
#truncated = false;
#lastChunkTime = 0;
#pendingChunk = "";
// Per-line column cap streaming state (persists across `push` calls so a
// long line split across chunks still trips the same trigger).
@@ -701,14 +702,20 @@ export class OutputSink {
push(chunk: string): void {
chunk = sanitizeWithOptionalSixelPassthrough(chunk, sanitizeText);
// Throttled onChunk: only call the callback when enough time has passed.
// Throttled onChunk: coalesce chunks arriving inside the throttle window
// and flush the buffered concatenation on the next eligible tick (plus a
// final flush in dump()) so the preview never has silent gaps.
// Live preview gets the raw (pre-cap) chunk so the TUI never lags behind
// what reached the sink — the column cap is for the persisted LLM view.
if (this.#onChunk) {
const now = Date.now();
if (now - this.#lastChunkTime >= this.#chunkThrottleMs) {
this.#lastChunkTime = now;
this.#onChunk(chunk);
const merged = this.#pendingChunk + chunk;
this.#pendingChunk = "";
this.#onChunk(merged);
} else {
this.#pendingChunk += chunk;
}
}
@@ -880,6 +887,11 @@ export class OutputSink {
const sink = Bun.file(this.#artifactPath).writer();
this.#file = { path: this.#artifactPath, artifactId: this.#artifactId, sink };
// Head-retained bytes precede the rolling tail buffer in the capture.
if (this.#head.length > 0) {
sink.write(this.#head);
}
// Flush existing buffer to file BEFORE it gets trimmed further.
if (this.#buffer.length > 0) {
sink.write(this.#buffer);
@@ -946,10 +958,19 @@ export class OutputSink {
this.#columnEllipsisAdded = false;
this.#columnDroppedBytes = 0;
this.#columnTruncatedLines = 0;
this.#pendingChunk = "";
}
async dump(notice?: string): Promise<OutputSummary> {
const noticeLine = notice ? `[${notice}]\n` : "";
// Flush any chunk still held back by the throttle so the live preview
// ends with the complete stream.
if (this.#onChunk && this.#pendingChunk.length > 0) {
const pending = this.#pendingChunk;
this.#pendingChunk = "";
this.#onChunk(pending);
}
const totalLines = this.#sawData ? this.#totalLines + 1 : 0;
if (this.#file) await this.#file.sink.end();
@@ -14,7 +14,6 @@ import { sanitizeText } from "@oh-my-pi/pi-utils";
import type { Terminal as XtermTerminalType } from "@xterm/headless";
import xterm from "@xterm/headless";
import { Settings } from "../config/settings";
import { NON_INTERACTIVE_ENV } from "../exec/non-interactive-env";
import type { Theme } from "../modes/theme/theme";
import { OutputSink, type OutputSummary } from "../session/streaming-output";
import { sanitizeWithOptionalSixelPassthrough } from "../utils/sixel";
@@ -358,8 +357,11 @@ export async function runInteractiveBashPty(
command: options.command,
cwd: options.cwd,
timeoutMs: options.timeoutMs,
// Interactive PTY: inherit the user's environment (the Rust side
// applies these as overrides), with a real TERM so editors,
// pagers, and TUIs behave like a normal terminal.
env: {
...NON_INTERACTIVE_ENV,
TERM: "xterm-256color",
...options.env,
},
signal: options.signal,
+52 -11
View File
@@ -410,10 +410,19 @@ export class BashTool implements AgentTool<BashToolSchema, BashToolDetails> {
*/
#throwIfUnfinished(result: BashResult | BashInteractiveResult, timeoutSec: number, outputText: string): void {
if (result.cancelled) {
throw new ToolError(normalizeResultOutput(result) || "Command aborted");
// executeBash output already carries a `[Command cancelled]` notice from
// the sink; PTY/bridge interactive output does not, so annotate it here.
const out = normalizeResultOutput(result);
const annotated = isInteractiveResult(result) && out ? `${out}\n\n[Command aborted]` : out;
throw new ToolError(annotated || "Command aborted");
}
if (isInteractiveResult(result) && result.timedOut) {
throw new ToolError(normalizeResultOutput(result) || `Command timed out after ${timeoutSec} seconds`);
const out = normalizeResultOutput(result);
throw new ToolError(
out
? `${out}\n\n[Command timed out after ${timeoutSec} seconds]`
: `Command timed out after ${timeoutSec} seconds`,
);
}
if (result.exitCode === undefined) {
throw new ToolError(`${outputText}\n\nCommand failed: missing exit status`);
@@ -669,7 +678,10 @@ export class BashTool implements AgentTool<BashToolSchema, BashToolDetails> {
// script can't pull the entire script into the "cwd" capture.
if (!cwd) {
const cdMatch = command.match(/^cd[ \t]+((?:[^&\\\n\r]|\\.)+?)[ \t]*&&[ \t]*/);
if (cdMatch) {
// Skip extraction when the path needs shell expansion ($VAR, $(...),
// backticks) — resolveToCwd only expands `~`, so routing those through
// cwd would reject commands the shell itself handles fine.
if (cdMatch && !/[$`(]/.test(cdMatch[1])) {
cwd = cdMatch[1].trim().replace(/^["']|["']$/g, "");
command = command.slice(cdMatch[0].length);
}
@@ -771,8 +783,24 @@ export class BashTool implements AgentTool<BashToolSchema, BashToolDetails> {
});
}
// The client-bridge terminal provides a live terminal card in the editor;
// when available it wins over auto-backgrounding (both are opt-in, and
// auto-background would otherwise silently disable the terminal route).
const clientBridge = this.session.getClientBridge?.();
const bridgeTerminalAvailable = Boolean(
clientBridge?.capabilities.terminal && clientBridge.createTerminal && !pty,
);
const autoBgManager = this.session.asyncJobManager;
if (this.#autoBackgroundEnabled && !pty && autoBgManager) {
// At the running-job cap, fall through to direct foreground execution
// instead of failing every bash call until a slot frees up.
if (
this.#autoBackgroundEnabled &&
!pty &&
!bridgeTerminalAvailable &&
autoBgManager &&
!autoBgManager.atCapacity
) {
const autoBackgroundWaitMs = this.#resolveAutoBackgroundWaitMs(timeoutMs);
const startBackgrounded = autoBackgroundWaitMs === 0;
const job = this.#startManagedBashJob({
@@ -793,21 +821,23 @@ export class BashTool implements AgentTool<BashToolSchema, BashToolDetails> {
notices: pendingNotices,
});
}
// Suppress the completion delivery up front so a job finishing while we
// foreground-wait cannot also be injected by the delivery loop. Lifted
// via resumeDeliveries() if we end up backgrounding after all.
autoBgManager.acknowledgeDeliveries([job.jobId]);
const waitResult = await this.#waitForManagedBashJob(job, autoBackgroundWaitMs, signal);
if (waitResult.kind === "completed") {
autoBgManager.acknowledgeDeliveries([job.jobId]);
return waitResult.result;
}
if (waitResult.kind === "failed") {
autoBgManager.acknowledgeDeliveries([job.jobId]);
throw waitResult.error;
}
if (waitResult.kind === "aborted") {
autoBgManager.cancel(job.jobId);
autoBgManager.acknowledgeDeliveries([job.jobId]);
throw new ToolAbortError(job.getLatestText() || "Command aborted");
}
job.setBackgrounded(true);
autoBgManager.resumeDeliveries([job.jobId]);
return this.#buildBackgroundStartResult(job.jobId, job.label, job.getLatestText(), timeoutSec, {
requestedTimeoutSec,
notices: pendingNotices,
@@ -816,7 +846,6 @@ export class BashTool implements AgentTool<BashToolSchema, BashToolDetails> {
// Route through the client terminal when the client advertises the terminal capability.
// Skip when pty=true (PTY needs the local terminal UI).
const clientBridge = this.session.getClientBridge?.();
if (clientBridge?.capabilities.terminal && clientBridge.createTerminal && !pty) {
const bridgeWallTimeStart = performance.now();
const handle = await clientBridge.createTerminal({
@@ -993,6 +1022,9 @@ export class BashTool implements AgentTool<BashToolSchema, BashToolDetails> {
const { path: artifactPath, id: artifactId } = (await this.session.allocateOutputArtifact?.("bash")) ?? {};
const interactiveUi = canUseInteractiveBashPty(pty, ctx) ? ctx?.ui : undefined;
if (pty && !interactiveUi) {
pendingNotices.push("pty requested but unavailable in this environment; ran without a terminal");
}
const wallTimeStart = performance.now();
const result: BashResult | BashInteractiveResult = interactiveUi
? await runInteractiveBashPty(interactiveUi, {
@@ -1017,13 +1049,22 @@ export class BashTool implements AgentTool<BashToolSchema, BashToolDetails> {
});
const wallTimeMs = performance.now() - wallTimeStart;
if (result.cancelled) {
const out = normalizeResultOutput(result);
// PTY output carries no cancel/timeout notice of its own; annotate so
// the model can tell an abort from a plain failure.
const message = isInteractiveResult(result) && out ? `${out}\n\n[Command aborted]` : out || "Command aborted";
if (signal?.aborted) {
throw new ToolAbortError(normalizeResultOutput(result) || "Command aborted");
throw new ToolAbortError(message);
}
throw new ToolError(normalizeResultOutput(result) || "Command aborted");
throw new ToolError(message);
}
if (isInteractiveResult(result) && result.timedOut) {
throw new ToolError(normalizeResultOutput(result) || `Command timed out after ${timeoutSec} seconds`);
const out = normalizeResultOutput(result);
throw new ToolError(
out
? `${out}\n\n[Command timed out after ${timeoutSec} seconds]`
: `Command timed out after ${timeoutSec} seconds`,
);
}
return this.#buildCompletedResult(result, timeoutSec, {
requestedTimeoutSec,
@@ -135,6 +135,47 @@ describe("AsyncJobManager", () => {
manager.cancel(firstJobId);
});
test("queued jobs do not count toward the cap until markRunning", async () => {
const manager = new AsyncJobManager({
maxRunningJobs: 1,
onJobComplete: async () => {},
});
const gate = Promise.withResolvers<void>();
const started = Promise.withResolvers<void>();
const release = Promise.withResolvers<void>();
const queuedJobId = manager.register(
"task",
"queued",
async ({ markRunning }) => {
await gate.promise;
markRunning();
started.resolve();
await release.promise;
return "queued done";
},
{ queued: true },
);
// Queued job holds no slot: another job registers fine at cap 1.
const runningJobId = manager.register("bash", "running", async ({ signal }) => {
await new Promise<void>(resolve => {
signal.addEventListener("abort", () => resolve(), { once: true });
});
return "done";
});
// Free the slot, then let the queued job start: it now occupies the slot.
manager.cancel(runningJobId);
gate.resolve();
await started.promise;
expect(() => manager.register("bash", "third", async () => "third")).toThrow(/Background job limit reached/);
release.resolve();
await manager.waitForAll();
expect(manager.getJob(queuedJobId)?.status).toBe("completed");
});
test("evicts completed jobs after retention period", async () => {
const manager = new AsyncJobManager({
retentionMs: 25,
@@ -280,6 +280,42 @@ describe("OutputSink", () => {
expect(dumped.output).toBe("bcdef");
});
test("artifact file includes head-retained bytes when head retention is enabled", async () => {
const dir = await createTempDir();
const artifactPath = path.join(dir, "output.log");
const sink = new OutputSink({
artifactPath,
artifactId: "artifact-2",
spillThreshold: 5,
headBytes: 4,
});
// First chunk lands fully in the head window; later chunks overflow the
// tail budget and trigger the artifact spill.
sink.push("head");
sink.push("abc");
sink.push("defgh");
const dumped = await sink.dump();
const artifactText = await Bun.file(artifactPath).text();
expect(dumped.truncated).toBe(true);
expect(artifactText).toBe("headabcdefgh");
});
test("throttled onChunk coalesces held-back chunks instead of dropping them", async () => {
const chunks: string[] = [];
const sink = new OutputSink({ onChunk: chunk => chunks.push(chunk), chunkThrottleMs: 60_000 });
sink.push("a");
// Inside the throttle window: buffered, not dropped.
sink.push("b");
sink.push("c");
const dumped = await sink.dump();
// First push fires immediately; dump flushes the coalesced remainder.
expect(chunks).toEqual(["a", "bc"]);
expect(dumped.output).toBe("abc");
});
test("createInput decodes streamed UTF-8 chunks correctly", async () => {
const sink = new OutputSink();
const writer = sink.createInput().getWriter();
@@ -1,9 +1,13 @@
import { describe, expect, it } from "bun:test";
import type { AgentToolContext } from "@oh-my-pi/pi-agent-core";
import { validateToolArguments } from "@oh-my-pi/pi-ai/utils/validation";
import type { BashInterceptorRule } from "@oh-my-pi/pi-coding-agent/config/settings-schema";
import {
type BashInterceptorRule,
DEFAULT_BASH_INTERCEPTOR_RULES,
} from "@oh-my-pi/pi-coding-agent/config/settings-schema";
import type { ToolSession } from "@oh-my-pi/pi-coding-agent/tools";
import { BashTool, type BashToolInput } from "@oh-my-pi/pi-coding-agent/tools/bash";
import { checkBashInterception } from "@oh-my-pi/pi-coding-agent/tools/bash-interceptor";
function createBashTool(rules: BashInterceptorRule[]): BashTool {
const session = {
@@ -58,6 +62,28 @@ describe("BashTool interception", () => {
});
});
describe("default echo/printf redirect rule", () => {
const tools = ["write"];
it("blocks unquoted redirects to files", () => {
expect(checkBashInterception("echo hi > out.txt", tools, DEFAULT_BASH_INTERCEPTOR_RULES).block).toBe(true);
expect(checkBashInterception("echo hi >> out.txt", tools, DEFAULT_BASH_INTERCEPTOR_RULES).block).toBe(true);
expect(checkBashInterception('printf "%s" foo > /tmp/x', tools, DEFAULT_BASH_INTERCEPTOR_RULES).block).toBe(true);
});
it("blocks clobber and variable-target redirects", () => {
expect(checkBashInterception("echo hi >| out.txt", tools, DEFAULT_BASH_INTERCEPTOR_RULES).block).toBe(true);
expect(checkBashInterception("echo hi > $OUT", tools, DEFAULT_BASH_INTERCEPTOR_RULES).block).toBe(true);
});
it("does not block `>` inside quoted text or fd duplication", () => {
expect(checkBashInterception('echo "a -> b"', tools, DEFAULT_BASH_INTERCEPTOR_RULES).block).toBe(false);
expect(checkBashInterception('echo "<p>hi</p>"', tools, DEFAULT_BASH_INTERCEPTOR_RULES).block).toBe(false);
expect(checkBashInterception("printf 'use 2>&1'", tools, DEFAULT_BASH_INTERCEPTOR_RULES).block).toBe(false);
expect(checkBashInterception('echo "err" >&2', tools, DEFAULT_BASH_INTERCEPTOR_RULES).block).toBe(false);
});
});
describe("BashTool argument validation", () => {
it("preserves async requests so disabled async mode returns the explicit error", async () => {
const tool = createBashTool([]);