refactor(coding-agent): extracted shared auto-backgrounding and job wait helpers

- Extract foreground wait, notice formatting, and settlement racing logic into a shared module.
- Update async job type definitions to support eval jobs alongside bash and task jobs.
- Add configuration settings for eval auto-background thresholds.
This commit is contained in:
can1357
2026-08-20 05:06:44 +02:00
parent 77930e467e
commit aeed1e6195
19 changed files with 847 additions and 373 deletions
+1 -1
View File
@@ -8,7 +8,7 @@
### Fixed
- Fixed the tool-argument repair layer applying lossy repairs on union-branch diagnoses: when a value failed every `anyOf`/`oneOf` variant, the first failing branch's issues were treated as authoritative, so object payloads got JSON-stringified into string-typed fields and unrecognized keys were silently deleted — corrupting subagent `yield` payloads (validation "passed" and parents received `summary: "{\"purge\":13,…}"` instead of a retryable error). Issues surfaced from a failed union branch are now marked at every depth and only receive lossless repairs (JSON-string parsing, boolean spellings, scalar coercion); container stringification, key deletion, and singleton-array wrapping require an authoritative non-union diagnosis.
- Fixed the tool-argument repair layer applying lossy repairs on union-branch diagnoses: when a value failed every `anyOf`/`oneOf` variant, the first failing branch's issues were treated as authoritative, so object payloads got JSON-stringified into string-typed fields and unrecognized keys were silently deleted — corrupting subagent `yield` payloads (validation "passed" and parents received `summary: "{\"purge\":13,…}"` instead of a retryable error). Issues from a failed union are now marked at every depth and only receive lossless repairs (JSON-string parsing, boolean spellings, scalar coercion) — unless a `const`/`enum` discriminator uniquely tag-selects the intended variant, whose diagnosis stays fully repairable (and is now the surfaced one, rather than the first failing branch's).
- Fixed local OpenAI-compatible servers with strict `chat_template_kwargs` whitelists (e.g. NInfer) failing every Qwen 3.8+ turn with `400 chat_template_kwargs.reasoning_effort is not supported` after the effort routing fix: the reasoning-effort fallback now recognizes a rejection of the kwargs spelling itself, retries with the kwarg stripped while keeping the effort on the standard top-level `reasoning_effort` field (hoisting it there for the kwargs-only vLLM dialect), and remembers the shape for the rest of the session. Value-level rejections and drops now also update the `chat_template_kwargs.reasoning_effort` twin instead of leaving a stale effort for kwargs-reading renderers, and unknown-parameter 400s naming `reasoning_effort` are recognized as effort rejections.
## [17.3.8] - 2026-08-19
+2
View File
@@ -11,6 +11,7 @@
- Added `extendedContext` setting (`/settings` → Context → General, default on). When off, models with a premium long-context price tier (OpenAI GPT-5.6 Sol/Terra/Luna bill 2x input / 1.5x output above 272K input tokens, on both the API and subscription Codex) are capped at the standard-pricing threshold — they appear as 272K again and compaction fires before a request crosses into premium billing. Toggling mid-session re-clamps or restores the active model's window immediately. Anthropic Claude 4.6+ serves its full 1M window at standard pricing, so no Anthropic model is affected.
- Added click-to-toggle and drag-to-reorder controls for list-valued `/settings` editors.
- Added `compaction.asyncEnabled` (Async Compaction, default on): when context enters the band just below the compaction threshold, maintenance speculatively summarizes in the background off a branch snapshot (first configured LLM-backed method — remote, handoff, or soft — isolated from the live turn by a side session id) and holds the armed result; crossing the threshold then splices it in instantly instead of blocking on a summarization round-trip. Armed results are invalidated by branch changes, reset boundaries, model switches that strand provider-native replay payloads, and context growth past `keepRecentTokens` (which re-speculates). The status line pulses the auto-compact icon while a speculation runs and holds it in accent once a result is armed.
- Added auto-backgrounding for eval cells, mirroring the bash policy: with `eval.autoBackground.enabled` (default off), a cell still running after `eval.autoBackground.thresholdMs` (default 60s) — or when a queued user/peer message steers mid-wait — converts into a background job that keeps executing on the kernel and delivers its result automatically. Both settings inherit the RPC-host default treatment alongside their bash counterparts.
### Changed
@@ -20,6 +21,7 @@
- `/settings` rows can now carry a risk note: a warning glyph on the row plus a warning-colored line above the description. `External Thinking` (`externalThinking`, `--external-thinking`) is the first user — providers have flagged the request shape it produces as abuse, up to account-level enforcement, so both the settings entry and `--help` now say so.
- The todo HUD header now draws a summed progress bar counting closed/total tasks across every stage. Once all tasks close, the bar smoothly collapses before the row disappears.
- The todo HUD now carries overall progress in the tree spine instead of a horizontal header bar: the top-level connector column runs unbroken from the `TODO` header down the left edge and closes with a short `└────` elbow tail under the block. The closed/total fraction (summed across every stage) fills that path in accent — down the spine, around the bend, out along the tail. Once all tasks close, the accent smoothly drains back up before the header row disappears.
- The todo HUD now carries overall progress in the tree spine instead of a horizontal header bar: the top-level connector column runs unbroken from the `TODO` header to the last row, and the closed/total fraction (summed across every stage) colors it top-down in accent. Once all tasks close, the accent smoothly drains out of the spine before the header row disappears.
- Token counting is now scoped to the model being billed rather than to a process-global tokenizer: session maintenance, stats, advisors, `/context`, snapcompact inline imaging, and `compress` each count through the owning agent's `Tokenizer` (`agent.tokenizer`). Message counting is `Tokenizer.countMessage`/`countMessages` (replacing the free `estimateTokens(message, tokenizer)` helper; the legacy shim keeps a compat `estimateTokens` export for legacy pi extensions). `estimateToolSchemaTokens`, `estimateSkillsTokens`, `computeNonMessageTokens`, and `computeNonMessageBreakdown` take an explicit tokenizer; standalone prompt inspection intentionally keeps the default estimate because it has no resolved catalog model.
- The advisor runtime's `maintainContext` hook now receives the pending update as a message instead of a pre-computed token count — sizing it needs the advisor model's tokenizer, which the host owns.
- Handoff no longer starts a new session: `/handoff` and the auto-maintenance `handoff` method now commit the generated document as a regular compaction entry on the current session (document becomes the summary, recent history is kept per `compaction.keepRecentTokens`, session id/transcript/cache key unchanged). The `session_before_switch`/`session_switch` extension events no longer fire with reason `"handoff"`, mid-turn maintenance no longer skips the handoff preference, and overflow recovery can now apply a pre-armed handoff result.
@@ -0,0 +1,78 @@
/**
* Shared foreground-wait helpers for tools that auto-background long-running
* work as {@link AsyncJobManager} jobs (bash commands, eval cells): the
* LLM-facing background notice, the threshold-vs-timeout wait budget, and the
* settlement race against abort/steering signals.
*/
/** Default foreground-wait threshold before a tool call auto-backgrounds. */
export const DEFAULT_AUTO_BACKGROUND_THRESHOLD_MS = 60_000;
/** LLM-facing footer appended when a tool call is converted into a background job. */
export function formatBackgroundNotice(jobId: string): string {
return `Backgrounded as job ${jobId}; result will be delivered automatically.`;
}
/**
* How long a tool foreground-waits before backgrounding. Bounded by the call's
* own timeout minus a small buffer so a deadline expiry resolves inline instead
* of backgrounding moments before it fires. `0` means background immediately.
*/
export function resolveAutoBackgroundWaitMs(thresholdMs: number, timeoutMs: number | undefined): number {
if (thresholdMs <= 0) return 0;
if (timeoutMs === undefined) return thresholdMs;
const timeoutBufferMs = 1_000;
return Math.max(0, Math.min(thresholdMs, timeoutMs - timeoutBufferMs));
}
/** Non-settled outcomes of {@link raceJobSettlement}. */
export type JobWaitInterrupt = { kind: "running" } | { kind: "steer" } | { kind: "aborted" };
/**
* Race a managed job's settlement against the auto-background threshold, the
* caller's abort signal, and the turn's steering signal. Returns the job's own
* completion when it settles first; otherwise reports why the wait ended:
* "running" = threshold elapsed (background it), "steer" = a queued message
* arrived mid-wait, "aborted" = the caller cancelled.
*/
export async function raceJobSettlement<C>(
completion: Promise<C>,
thresholdMs: number,
signal?: AbortSignal,
steeringSignal?: AbortSignal,
): Promise<C | JobWaitInterrupt> {
if (signal?.aborted) {
return { kind: "aborted" };
}
if (steeringSignal?.aborted) {
return { kind: "steer" };
}
// Cancellable threshold: a bare Bun.sleep(thresholdMs) leaves a live, ref'd
// timer for the full threshold after the job finishes (or abort/steer) wins
// the race first — delaying SDK/headless shutdown and accumulating timers
// under fast completion rates. Settle a withResolvers promise from
// setTimeout so the finally can clear it regardless of which waiter wins.
const { promise: thresholdPromise, resolve: resolveThreshold } = Promise.withResolvers<{ kind: "running" }>();
const thresholdTimer = setTimeout(() => resolveThreshold({ kind: "running" }), thresholdMs);
const waiters: Array<Promise<C | JobWaitInterrupt>> = [completion, thresholdPromise];
const { promise: abortedPromise, resolve: resolveAborted } = Promise.withResolvers<{ kind: "aborted" }>();
const onAbort = () => resolveAborted({ kind: "aborted" });
const { promise: steerPromise, resolve: resolveSteer } = Promise.withResolvers<{ kind: "steer" }>();
const onSteer = () => resolveSteer({ kind: "steer" });
if (signal) {
signal.addEventListener("abort", onAbort, { once: true });
waiters.push(abortedPromise);
}
if (steeringSignal) {
steeringSignal.addEventListener("abort", onSteer, { once: true });
waiters.push(steerPromise);
}
try {
return await Promise.race(waiters);
} finally {
clearTimeout(thresholdTimer);
signal?.removeEventListener("abort", onAbort);
steeringSignal?.removeEventListener("abort", onSteer);
}
}
+1
View File
@@ -1 +1,2 @@
export * from "./auto-background";
export * from "./job-manager";
@@ -29,9 +29,12 @@ interface PollEscalationState {
lastPollEndAt: number;
}
/** Kind of work a managed job runs; drives job-row badges and delivery labels. */
export type AsyncJobType = "bash" | "task" | "eval";
export interface AsyncJob {
id: string;
type: "bash" | "task";
type: AsyncJobType;
status: "running" | "completed" | "failed" | "cancelled";
startTime: number;
label: string;
@@ -183,7 +186,7 @@ export class AsyncJobManager {
}
register(
type: "bash" | "task",
type: AsyncJobType,
label: string,
run: (ctx: {
jobId: string;
@@ -3652,6 +3652,22 @@ export const SETTINGS_SCHEMA = {
},
},
"eval.autoBackground.enabled": {
type: "boolean",
default: false,
ui: {
tab: "shell",
group: "Eval & Runtimes",
label: "Eval Auto-Background",
description: "Automatically background long-running eval cells and deliver the result later",
},
},
"eval.autoBackground.thresholdMs": {
type: "number",
default: 60_000,
},
// Runtime knobs (consumed by eval backends and the /python slash command)
"python.kernelMode": {
type: "enum",
+6
View File
@@ -45,4 +45,10 @@ export interface EvalToolDetails {
languages?: EvalLanguage[];
/** Optional human-readable notice (e.g. fallback explanation). */
notice?: string;
/** Present when the cell was auto-backgrounded as an async job. */
async?: {
state: "running" | "completed" | "failed";
jobId: string;
type: "eval";
};
}
+2
View File
@@ -155,6 +155,8 @@ const RPC_BACKGROUND_DEFAULTED_SETTING_PATHS: SettingPath[] = [
"async.maxJobs",
"bash.autoBackground.enabled",
"bash.autoBackground.thresholdMs",
"eval.autoBackground.enabled",
"eval.autoBackground.thresholdMs",
];
// Protocol-mode hosts opt into a small set of paths whose host-default we
@@ -7,6 +7,7 @@
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { type Component, Text } from "@oh-my-pi/pi-tui";
import { formatBytes, formatDuration } from "@oh-my-pi/pi-utils";
import type { AsyncJobType } from "../../async";
import {
type CustomMessage,
type FileMentionMessage,
@@ -32,10 +33,10 @@ export function buildAsyncResultBlock(message: CustomOrHookMessage): ToolActivit
const details = (
message as CustomMessage<{
jobId?: string;
type?: "bash" | "task";
type?: AsyncJobType;
label?: string;
durationMs?: number;
jobs?: Array<{ jobId?: string; type?: "bash" | "task"; label?: string; durationMs?: number }>;
jobs?: Array<{ jobId?: string; type?: AsyncJobType; label?: string; durationMs?: number }>;
}>
).details;
const jobs =
@@ -43,3 +43,5 @@ Acyclic waves via `agent(…, handle=true)` + `pipeline`/`parallel`:
<critical>
Prior top-level names survive into the next cell — reuse; NEVER re-import/re-declare. Re-read only if file changed since last read.
</critical>
{{#if autoBackgroundEnabled}}Long-running cells may auto-background and deliver later; the kernel stays busy until the cell finishes. Need inline? Raise `timeout`.{{/if}}
@@ -9,7 +9,7 @@
* every completion — regardless of owner — into the first top-level session.
*/
import { prompt } from "@oh-my-pi/pi-utils";
import type { AsyncJob } from "../async";
import type { AsyncJob, AsyncJobType } from "../async";
import asyncResultTemplate from "../prompts/tools/async-result.md" with { type: "text" };
import type { CustomMessage } from "./messages";
@@ -41,7 +41,7 @@ export interface AsyncResultEntry {
type AsyncResultJobDetails = {
jobId: string;
type?: "bash" | "task";
type?: AsyncJobType;
label?: string;
durationMs?: number;
};
+1 -1
View File
@@ -888,7 +888,7 @@ export function createSubagentSettings(
return Settings.isolated(
{
...snapshot,
// Async jobs and bash auto-backgrounding are inherited from the parent:
// Async jobs and bash/eval auto-backgrounding are inherited from the parent:
// background jobs are owner-routed to the subagent's own session, and
// the run driver's quiescence barrier + teardown reap guarantee no
// owner job outlives the run, so worktree capture/cleanup stays
+9 -62
View File
@@ -10,6 +10,12 @@ import type {
import type { Component } from "@oh-my-pi/pi-tui";
import { ImageProtocol, TERMINAL } from "@oh-my-pi/pi-tui";
import { getProjectDir, isEnoent, logger, prompt } from "@oh-my-pi/pi-utils";
import {
DEFAULT_AUTO_BACKGROUND_THRESHOLD_MS,
formatBackgroundNotice,
raceJobSettlement,
resolveAutoBackgroundWaitMs,
} from "../async";
import type { Settings } from "../config/settings";
import { applyDirenvPreflight, type BashResult, executeBash } from "../exec/bash-executor";
import type { RenderResultOptions } from "../extensibility/custom-tools/types";
@@ -57,7 +63,6 @@ import { clampTimeout, TOOL_TIMEOUTS } from "./tool-timeouts";
export const BASH_DEFAULT_PREVIEW_LINES = DEFAULT_TERMINAL_PREVIEW_LINES;
const BASH_ENV_NAME_PATTERN = /^[A-Za-z_][A-Za-z0-9_]*$/;
const DEFAULT_AUTO_BACKGROUND_THRESHOLD_MS = 60_000;
const BASH_APPROVAL_SHELL_CONTROL_CHARS: Record<string, true> = {
"\n": true,
"\r": true,
@@ -505,10 +510,6 @@ function formatExitCodeNotice(exitCode: number): string {
return `Command exited with code ${exitCode}`;
}
function formatBackgroundNotice(jobId: string): string {
return `Backgrounded as job ${jobId}; result will be delivered automatically.`;
}
/**
* Strip the trailing occurrence of `notice` (plus a single surrounding newline
* on each side) so the TUI can echo the value via a styled footer label
@@ -894,60 +895,6 @@ export class BashTool implements AgentTool<typeof bashSchemaBase | typeof bashSc
};
}
async #waitForManagedBashJob(
job: ManagedBashJobHandle,
thresholdMs: number,
signal?: AbortSignal,
steeringSignal?: AbortSignal,
): Promise<ManagedBashJobCompletion | { kind: "running" } | { kind: "steer" } | { kind: "aborted" }> {
if (signal?.aborted) {
return { kind: "aborted" };
}
if (steeringSignal?.aborted) {
return { kind: "steer" };
}
// Cancellable threshold: a bare Bun.sleep(thresholdMs) leaves a live, ref'd
// timer for the full threshold after the command finishes (or abort/steer)
// wins the race first — delaying SDK/headless shutdown and accumulating
// timers under fast command rates. Settle a withResolvers promise from
// setTimeout so the finally can clear it regardless of which waiter wins.
const { promise: thresholdPromise, resolve: resolveThreshold } = Promise.withResolvers<{
kind: "running";
}>();
const thresholdTimer = setTimeout(() => resolveThreshold({ kind: "running" }), thresholdMs);
const waiters: Array<
Promise<ManagedBashJobCompletion | { kind: "running" } | { kind: "steer" } | { kind: "aborted" }>
> = [job.completion, thresholdPromise];
const { promise: abortedPromise, resolve: resolveAborted } = Promise.withResolvers<{ kind: "aborted" }>();
const onAbort = () => resolveAborted({ kind: "aborted" });
const { promise: steerPromise, resolve: resolveSteer } = Promise.withResolvers<{ kind: "steer" }>();
const onSteer = () => resolveSteer({ kind: "steer" });
if (signal) {
signal.addEventListener("abort", onAbort, { once: true });
waiters.push(abortedPromise);
}
if (steeringSignal) {
steeringSignal.addEventListener("abort", onSteer, { once: true });
waiters.push(steerPromise);
}
try {
return await Promise.race(waiters);
} finally {
clearTimeout(thresholdTimer);
signal?.removeEventListener("abort", onAbort);
steeringSignal?.removeEventListener("abort", onSteer);
}
}
#resolveAutoBackgroundWaitMs(timeoutMs: number | undefined): number {
if (this.#autoBackgroundThresholdMs <= 0) return 0;
if (timeoutMs === undefined) return this.#autoBackgroundThresholdMs;
const timeoutBufferMs = 1_000;
return Math.max(0, Math.min(this.#autoBackgroundThresholdMs, timeoutMs - timeoutBufferMs));
}
async execute(
_toolCallId: string,
{
@@ -1099,7 +1046,7 @@ export class BashTool implements AgentTool<typeof bashSchemaBase | typeof bashSc
autoBgManager &&
!autoBgManager.atCapacity
) {
const autoBackgroundWaitMs = this.#resolveAutoBackgroundWaitMs(timeoutMs);
const autoBackgroundWaitMs = resolveAutoBackgroundWaitMs(this.#autoBackgroundThresholdMs, timeoutMs);
const startBackgrounded = autoBackgroundWaitMs === 0;
const job = this.#startManagedBashJob({
command,
@@ -1123,8 +1070,8 @@ export class BashTool implements AgentTool<typeof bashSchemaBase | typeof bashSc
// 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,
const waitResult = await raceJobSettlement(
job.completion,
autoBackgroundWaitMs,
signal,
ctx?.toolCall?.steeringSignal,
+20 -4
View File
@@ -586,6 +586,10 @@ export const evalToolRenderer = {
warningLine = formatStyledTruncationWarning(details.meta, uiTheme) ?? undefined;
}
const noticeLine = details?.notice ? uiTheme.fg("dim", wrapBrackets(details.notice, uiTheme)) : undefined;
const asyncLine =
details?.async?.state === "running"
? uiTheme.fg("dim", wrapBrackets(`Backgrounded: ${details.async.jobId}`, uiTheme))
: undefined;
const cellResults = details?.cells;
if (cellResults && cellResults.length > 0) {
@@ -670,6 +674,9 @@ export const evalToolRenderer = {
if (noticeLine) {
lines.push(noticeLine);
}
if (asyncLine) {
lines.push(asyncLine);
}
if (warningLine) {
lines.push(warningLine);
}
@@ -693,14 +700,19 @@ export const evalToolRenderer = {
);
if (!combinedOutput && statusLines.length === 0) {
const lines = [timeoutLine, noticeLine, warningLine].filter(Boolean) as string[];
const lines = [timeoutLine, noticeLine, asyncLine, warningLine].filter(Boolean) as string[];
return new Text(lines.join("\n"), 0, 0);
}
if (!combinedOutput && statusLines.length > 0) {
const lines = [uiTheme.fg("dim", "Status"), ...statusLines, timeoutLine, noticeLine, warningLine].filter(
Boolean,
) as string[];
const lines = [
uiTheme.fg("dim", "Status"),
...statusLines,
timeoutLine,
noticeLine,
asyncLine,
warningLine,
].filter(Boolean) as string[];
return new Text(lines.join("\n"), 0, 0);
}
@@ -714,6 +726,7 @@ export const evalToolRenderer = {
...(statusLines.length > 0 ? [uiTheme.fg("dim", "Status"), ...statusLines] : []),
timeoutLine,
noticeLine,
asyncLine,
warningLine,
].filter(Boolean) as string[];
return new Text(lines.join("\n"), 0, 0);
@@ -765,6 +778,9 @@ export const evalToolRenderer = {
if (noticeLine) {
outputLines.push(truncateToWidth(noticeLine, width));
}
if (asyncLine) {
outputLines.push(truncateToWidth(asyncLine, width));
}
if (warningLine) {
outputLines.push(truncateToWidth(warningLine, width));
}
+478 -294
View File
@@ -2,6 +2,12 @@ import { type } from "@oh-my-pi/omptype";
import type { AgentTool, AgentToolContext, AgentToolResult, AgentToolUpdateCallback } from "@oh-my-pi/pi-agent-core";
import type { ImageContent, ToolExample } from "@oh-my-pi/pi-ai";
import { prompt } from "@oh-my-pi/pi-utils";
import {
DEFAULT_AUTO_BACKGROUND_THRESHOLD_MS,
formatBackgroundNotice,
raceJobSettlement,
resolveAutoBackgroundWaitMs,
} from "../async";
import { jsBackend, juliaBackend, pythonBackend, rubyBackend } from "../eval";
import type { ExecutorBackend, ExecutorBackendResult } from "../eval/backend";
import { EVAL_TIMEOUT_PAUSE_OP, EVAL_TIMEOUT_RESUME_OP } from "../eval/bridge-timeout";
@@ -169,6 +175,8 @@ export interface EvalToolDescriptionOptions {
* `false`/`""` hides `agent()`, and a comma list drives the advertised default.
*/
spawns?: boolean | string | null;
/** Advertise auto-backgrounding of long-running cells in the tool prompt. */
autoBackgroundEnabled?: boolean;
}
export function getEvalToolDescription(options: EvalToolDescriptionOptions = {}): string {
@@ -182,6 +190,7 @@ export function getEvalToolDescription(options: EvalToolDescriptionOptions = {})
js,
rb,
jl,
autoBackgroundEnabled: options.autoBackgroundEnabled ?? false,
spawns: spawnPolicy.enabled,
spawnDefaultAgent: spawnPolicy.defaultAgent,
spawnAllowedAgentsText: spawnPolicy.allowedPromptText,
@@ -206,6 +215,11 @@ interface ResolvedEvalCell {
resolved: ResolvedBackend;
}
/** Settlement handed from a managed eval job to its foreground waiter. */
type ManagedEvalJobCompletion =
| { kind: "completed"; result: AgentToolResult<EvalToolDetails | undefined> }
| { kind: "failed"; error: unknown };
function uniqueEvalLanguages(cells: ResolvedEvalCell[]): EvalLanguage[] {
return [...new Set(cells.map(cell => cell.resolved.backend.id))];
}
@@ -302,6 +316,7 @@ export class EvalTool implements AgentTool<typeof evalSchema> {
rb: backends.ruby,
jl: backends.julia,
spawns: sessionSpawns,
autoBackgroundEnabled: this.session.settings.get("eval.autoBackground.enabled"),
});
}
/** All reuse-chain examples; the `examples` getter filters by enabled languages. */
@@ -394,7 +409,7 @@ export class EvalTool implements AgentTool<typeof evalSchema> {
params: typeof evalSchema.infer,
signal?: AbortSignal,
onUpdate?: AgentToolUpdateCallback,
_ctx?: AgentToolContext,
ctx?: AgentToolContext,
): Promise<AgentToolResult<EvalToolDetails | undefined>> {
if (this.#proxyExecutor) {
return this.#proxyExecutor(params, signal);
@@ -428,6 +443,185 @@ export class EvalTool implements AgentTool<typeof evalSchema> {
const languages = uniqueEvalLanguages(cells);
const notice = detailsNotice(cells);
const sessionAbortController = new AbortController();
const emitToolUpdate = onUpdate
? (text: string, details: EvalToolDetails): void => {
onUpdate({ content: [{ type: "text", text }], details });
}
: undefined;
const run = (
runSignal: AbortSignal | undefined,
emitUpdate: ((text: string, details: EvalToolDetails) => void) | undefined,
): Promise<AgentToolResult<EvalToolDetails | undefined>> => {
const execution = this.#runCells({
session,
cells,
languages,
notice,
excludeWebP,
signal: runSignal,
sessionAbortController,
emitUpdate,
});
return session.trackEvalExecution?.(execution, sessionAbortController) ?? execution;
};
const autoBgManager = session.asyncJobManager;
// At the running-job cap, fall through to direct foreground execution
// instead of failing every eval call until a slot frees up.
if (!session.settings.get("eval.autoBackground.enabled") || !autoBgManager || autoBgManager.atCapacity) {
return await run(signal, emitToolUpdate);
}
const thresholdMs = Math.max(
0,
Math.floor(session.settings.get("eval.autoBackground.thresholdMs") ?? DEFAULT_AUTO_BACKGROUND_THRESHOLD_MS),
);
// The wait budget mirrors #runCells' clamped cell timeout. The cell budget
// is runtime work (it pauses across agent()/tool bridge calls), so a cell
// can legitimately outlive it in wall time — exactly the case
// backgrounding exists for.
const cellTimeoutMs =
cells[0].timeoutMs === 0
? undefined
: clampTimeout("eval", cells[0].timeoutMs / 1000, session.settings.get("tools.maxTimeout")) * 1000;
const autoBackgroundWaitMs = resolveAutoBackgroundWaitMs(thresholdMs, cellTimeoutMs);
const startBackgrounded = autoBackgroundWaitMs === 0;
const rawLabel = params.title?.trim() || params.code.trim().split("\n", 1)[0] || "eval cell";
const label = rawLabel.length > 120 ? `${rawLabel.slice(0, 117)}...` : rawLabel;
let latestText = "";
let latestDetails: EvalToolDetails | undefined;
let forwardUpdates = !startBackgrounded;
const completion = Promise.withResolvers<ManagedEvalJobCompletion>();
const jobId = autoBgManager.register(
"eval",
label,
async ({ jobId, signal: runSignal, reportProgress }) => {
try {
const result = await run(runSignal, (text, details) => {
latestText = text;
latestDetails = details;
void reportProgress(text, { async: { state: "running", jobId, type: "eval" } });
if (forwardUpdates) emitToolUpdate?.(text, details);
});
const finalText = result.content.find(block => block.type === "text")?.text ?? "";
latestText = finalText;
// Hand the full result (images included) to the foreground waiter
// before deciding the job's terminal state.
completion.resolve({ kind: "completed", result });
if (result.isError === true) {
// A failed, cancelled, or timed-out cell is a completed execution
// that errored. Re-enter the failure path so the job manager
// records it as failed and delivers the error text.
throw new ToolError(finalText || "Eval cell failed");
}
await reportProgress(finalText, { async: { state: "completed", jobId, type: "eval" } });
return finalText;
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
latestText = message;
completion.resolve({ kind: "failed", error });
await reportProgress(message, { async: { state: "failed", jobId, type: "eval" } });
throw error;
}
},
{ ownerId: session.getAgentId?.() ?? undefined },
);
if (startBackgrounded) {
return this.#buildBackgroundStartResult(jobId, cells, languages, notice, latestText, latestDetails);
}
// 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([jobId]);
const waitResult = await raceJobSettlement(
completion.promise,
autoBackgroundWaitMs,
signal,
ctx?.toolCall?.steeringSignal,
);
if (waitResult.kind === "completed") {
return waitResult.result;
}
if (waitResult.kind === "failed") {
throw waitResult.error;
}
if (waitResult.kind === "aborted") {
autoBgManager.cancel(jobId);
throw new ToolAbortError(latestText || "Eval cell aborted");
}
forwardUpdates = false;
autoBgManager.resumeDeliveries([jobId]);
// "steer": a queued user/peer message arrived mid-wait — background the
// cell (it keeps running) so the message injects promptly.
const steerNotice =
waitResult.kind === "steer"
? "Backgrounded early to handle an incoming message; the cell keeps running."
: undefined;
return this.#buildBackgroundStartResult(jobId, cells, languages, notice, latestText, latestDetails, steerNotice);
}
/**
* Tool result returned when a cell converts into a background job: the live
* output tail plus the background notice, with details carrying the running
* cell snapshot and the async job marker the transcript renderer keys on.
*/
#buildBackgroundStartResult(
jobId: string,
cells: ResolvedEvalCell[],
languages: EvalLanguage[],
notice: string | undefined,
previewText: string,
latestDetails: EvalToolDetails | undefined,
extraNotice?: string,
): AgentToolResult<EvalToolDetails> {
// latestDetails snapshots are per-update copies (buildUpdateDetails), so
// tagging the async marker on cannot leak into later job progress.
const details: EvalToolDetails = latestDetails ?? {
language: languages[0],
languages,
cells: cells.map(cell => ({
index: cell.index,
title: cell.title,
code: cell.code,
language: cell.resolved.backend.id,
output: previewText,
status: "running" as const,
})),
};
if (notice) details.notice ??= notice;
details.async = { state: "running", jobId, type: "eval" };
const lines: string[] = [];
const trimmedPreview = previewText.trimEnd();
if (trimmedPreview.length > 0) {
lines.push(trimmedPreview, "");
}
if (extraNotice) {
lines.push(extraNotice, "");
}
lines.push(formatBackgroundNotice(jobId));
return { content: [{ type: "text", text: lines.join("\n") }], details };
}
/**
* Execute the resolved cells against their backends, streaming tail/detail
* updates through `emitUpdate`. Runs identically in the foreground path and
* inside a managed background job (which passes the job's own signal).
*/
async #runCells(options: {
session: ToolSession;
cells: ResolvedEvalCell[];
languages: EvalLanguage[];
notice: string | undefined;
excludeWebP: boolean | undefined;
signal: AbortSignal | undefined;
sessionAbortController: AbortController;
emitUpdate?: (text: string, details: EvalToolDetails) => void;
}): Promise<AgentToolResult<EvalToolDetails | undefined>> {
const { session, cells, languages, notice, excludeWebP, signal, sessionAbortController, emitUpdate } = options;
let outputSink: OutputSink | undefined;
let outputSummary: OutputSummary | undefined;
let outputDumped = false;
@@ -437,309 +631,299 @@ export class EvalTool implements AgentTool<typeof evalSchema> {
outputDumped = true;
return outputSummary;
};
try {
if (signal?.aborted) {
throw new ToolAbortError();
}
session.assertEvalExecutionAllowed?.();
const execution = (async (): Promise<AgentToolResult<EvalToolDetails | undefined>> => {
try {
if (signal?.aborted) {
throw new ToolAbortError();
}
session.assertEvalExecutionAllowed?.();
const tailBuffer = new TailBuffer(DEFAULT_MAX_BYTES * 2);
const jsonOutputs: unknown[] = [];
const images: ImageContent[] = [];
const statusEvents: EvalStatusEvent[] = [];
const tailBuffer = new TailBuffer(DEFAULT_MAX_BYTES * 2);
const jsonOutputs: unknown[] = [];
const images: ImageContent[] = [];
const statusEvents: EvalStatusEvent[] = [];
const cellResults: EvalCellResult[] = cells.map(cell => ({
index: cell.index,
title: cell.title,
code: cell.code,
language: cell.resolved.backend.id,
output: "",
status: "pending",
}));
const cellOutputs: string[] = [];
// The cell currently inside backend.execute(). Streamed stdout is
// appended to its rendered `output` live so a long-running cell (e.g. a
// sleep loop) shows progress instead of nothing until it returns. A
// dedicated per-cell tail buffer keeps attribution correct and avoids
// double-counting against the aggregate `tailBuffer`; on completion the
// authoritative `cellResult.output` (below) overwrites this live tail.
let activeLiveCell: { result: EvalCellResult; buf: TailBuffer } | undefined;
const cellResults: EvalCellResult[] = cells.map(cell => ({
index: cell.index,
title: cell.title,
code: cell.code,
language: cell.resolved.backend.id,
output: "",
status: "pending",
}));
const cellOutputs: string[] = [];
// The cell currently inside backend.execute(). Streamed stdout is
// appended to its rendered `output` live so a long-running cell (e.g. a
// sleep loop) shows progress instead of nothing until it returns. A
// dedicated per-cell tail buffer keeps attribution correct and avoids
// double-counting against the aggregate `tailBuffer`; on completion the
// authoritative `cellResult.output` (below) overwrites this live tail.
let activeLiveCell: { result: EvalCellResult; buf: TailBuffer } | undefined;
const appendTail = (text: string) => {
tailBuffer.append(text);
};
const buildUpdateDetails = (): EvalToolDetails => {
const details: EvalToolDetails = {
language: languages[0],
languages,
cells: cellResults.map(cell => ({
...cell,
statusEvents: cell.statusEvents ? [...cell.statusEvents] : undefined,
})),
};
if (jsonOutputs.length > 0) {
details.jsonOutputs = jsonOutputs;
}
if (images.length > 0) {
details.images = images;
}
if (statusEvents.length > 0) {
details.statusEvents = statusEvents;
}
if (notice) {
details.notice = notice;
}
return details;
};
const pushUpdate = () => {
if (!onUpdate) return;
const tailText = tailBuffer.text();
onUpdate({
content: [{ type: "text", text: tailText }],
details: buildUpdateDetails(),
});
};
const sessionFile = session.getSessionFile?.() ?? undefined;
const kernelOwnerId = session.getEvalKernelOwnerId?.() ?? undefined;
const { path: artifactPath, id: artifactId } = (await session.allocateOutputArtifact?.("eval")) ?? {};
session.assertEvalExecutionAllowed?.();
outputSink = new OutputSink({
artifactPath,
artifactId,
headBytes: resolveOutputSinkHeadBytes(session.settings),
maxColumns: resolveOutputMaxColumns(session.settings),
onChunk: chunk => {
appendTail(chunk);
if (activeLiveCell) {
activeLiveCell.buf.append(chunk);
activeLiveCell.result.output = activeLiveCell.buf.text();
}
pushUpdate();
},
});
const sessionId = session.getEvalSessionId?.() ?? defaultEvalSessionId(session);
for (let i = 0; i < cells.length; i++) {
const cell = cells[i];
const backend = cell.resolved.backend;
// The per-cell `timeout` is a budget on the cell runtime's *own*
// work. Host-side `agent()`/`parallel()`/`completion()` bridge calls suspend
// that budget entirely and restart a fresh timeout window when control
// returns to the active backend runtime. Compute, stdout, `log()`/`phase()`, and
// ordinary tool calls all count against the budget. The watchdog drives
// `combinedSignal`; we pass no wall-clock deadline downstream so the
// backends never arm a competing fixed timer.
const idleTimeoutMs =
cell.timeoutMs === 0
? undefined
: clampTimeout("eval", cell.timeoutMs / 1000, session.settings.get("tools.maxTimeout")) * 1000;
const idle = idleTimeoutMs === undefined ? undefined : new IdleTimeout(idleTimeoutMs);
const combinedSignal =
signal && idle
? AbortSignal.any([signal, idle.signal, sessionAbortController.signal])
: signal
? AbortSignal.any([signal, sessionAbortController.signal])
: idle
? AbortSignal.any([idle.signal, sessionAbortController.signal])
: sessionAbortController.signal;
const cellResult = cellResults[i];
cellResult.status = "running";
cellResult.output = "";
cellResult.statusEvents = undefined;
cellResult.exitCode = undefined;
cellResult.durationMs = undefined;
activeLiveCell = { result: cellResult, buf: new TailBuffer(DEFAULT_MAX_BYTES * 2) };
pushUpdate();
const startTime = Date.now();
let result: ExecutorBackendResult;
try {
result = await backend.execute(cell.code, {
cwd: session.cwd,
sessionId,
sessionFile: sessionFile ?? undefined,
kernelOwnerId,
signal: combinedSignal,
session,
idleTimeoutMs,
reset: cell.reset,
onChunk: chunk => {
outputSink!.push(chunk);
},
onStatus: event => {
if (event.op === EVAL_TIMEOUT_PAUSE_OP) {
idle?.pause();
return;
}
if (event.op === EVAL_TIMEOUT_RESUME_OP) {
idle?.resume();
return;
}
cellResult.statusEvents ??= [];
upsertStatusEvent(cellResult.statusEvents, event);
pushUpdate();
},
});
} finally {
idle?.dispose();
activeLiveCell = undefined;
}
const durationMs = Date.now() - startTime;
const cellStatusEvents: EvalStatusEvent[] = [];
const cellDisplayOutputs: EvalDisplayOutput[] = [];
const cellImageNotes: string[] = [];
let cellHasMarkdown = false;
for (const output of result.displayOutputs) {
if (output.type === "json") {
jsonOutputs.push(output.data);
cellDisplayOutputs.push(output);
}
if (output.type === "image") {
const resized = await resizeImage(
{
type: "image",
data: output.data,
mimeType: output.mimeType,
},
{ excludeWebP },
);
const image: ImageContent = {
type: "image",
data: resized.data,
mimeType: resized.mimeType,
};
images.push(image);
cellDisplayOutputs.push({
type: "image",
data: image.data,
mimeType: image.mimeType,
});
const dimensionNote = formatDimensionNote(resized);
if (dimensionNote) {
cellImageNotes.push(`display image ${cellImageNotes.length + 1}: ${dimensionNote}`);
}
}
if (output.type === "status") {
upsertStatusEvent(statusEvents, output.event);
upsertStatusEvent(cellStatusEvents, output.event);
}
if (output.type === "markdown") {
cellHasMarkdown = true;
}
}
const stdoutTrimmed = result.output.trim();
const imageText = cellImageNotes.join("\n");
const displayText = formatDisplayOutputsForText(cellDisplayOutputs);
const visibleDisplayText =
displayText && imageText ? `${displayText}\n\n${imageText}` : displayText || imageText;
const cellOutput =
stdoutTrimmed && visibleDisplayText
? `${stdoutTrimmed}\n\n${visibleDisplayText}`
: stdoutTrimmed || visibleDisplayText;
cellResult.output = cellOutput;
cellResult.exitCode = result.exitCode;
cellResult.durationMs = durationMs;
cellResult.statusEvents = cellStatusEvents.length > 0 ? cellStatusEvents : undefined;
cellResult.hasMarkdown = cellHasMarkdown || undefined;
if (cellOutput) {
cellOutputs.push(cellOutput);
appendTail(cellOutput);
}
if (result.cancelled) {
cellResult.status = "error";
pushUpdate();
const errorMsg = result.output || "Command aborted";
const combinedOutput = cellOutputs.join("\n\n");
const outputText = combinedOutput || errorMsg;
const summaryForMeta = await summarizeFinal(combinedOutput, finalizeOutput);
const details: EvalToolDetails = {
language: languages[0],
languages,
cells: cellResults,
jsonOutputs: jsonOutputs.length > 0 ? jsonOutputs : undefined,
statusEvents: statusEvents.length > 0 ? statusEvents : undefined,
isError: true,
};
if (notice) details.notice = notice;
return toolResult(details)
.content([{ type: "text", text: outputText }, ...images])
.truncationFromSummary(summaryForMeta, { direction: "tail" })
.done();
}
if (result.exitCode !== 0 && result.exitCode !== undefined) {
cellResult.status = "error";
pushUpdate();
const combinedOutput = cellOutputs.join("\n\n");
const outputText = combinedOutput
? `${combinedOutput}\n\nCommand exited with code ${result.exitCode}`
: `Command exited with code ${result.exitCode}`;
const summaryForMeta = await summarizeFinal(combinedOutput, finalizeOutput);
const details: EvalToolDetails = {
language: languages[0],
languages,
cells: cellResults,
jsonOutputs: jsonOutputs.length > 0 ? jsonOutputs : undefined,
statusEvents: statusEvents.length > 0 ? statusEvents : undefined,
isError: true,
};
if (notice) details.notice = notice;
return toolResult(details)
.content([{ type: "text", text: outputText }, ...images])
.truncationFromSummary(summaryForMeta, { direction: "tail" })
.done();
}
cellResult.status = "complete";
pushUpdate();
}
const combinedOutput = cellOutputs.join("\n\n");
const hasImages = images.length > 0;
const outputText =
combinedOutput ||
(hasImages
? `(displayed ${images.length} image${images.length === 1 ? "" : "s"}; no text output)`
: "(no output)");
const summaryForMeta = await summarizeFinal(combinedOutput, finalizeOutput);
const appendTail = (text: string) => {
tailBuffer.append(text);
};
const buildUpdateDetails = (): EvalToolDetails => {
const details: EvalToolDetails = {
language: languages[0],
languages,
cells: cellResults,
jsonOutputs: jsonOutputs.length > 0 ? jsonOutputs : undefined,
statusEvents: statusEvents.length > 0 ? statusEvents : undefined,
cells: cellResults.map(cell => ({
...cell,
statusEvents: cell.statusEvents ? [...cell.statusEvents] : undefined,
})),
};
if (notice) details.notice = notice;
return toolResult(details)
.content([{ type: "text", text: outputText }, ...images])
.truncationFromSummary(summaryForMeta, { direction: "tail" })
.done();
} finally {
if (!outputDumped) {
try {
await finalizeOutput();
} catch {}
if (jsonOutputs.length > 0) {
details.jsonOutputs = jsonOutputs;
}
}
})();
if (images.length > 0) {
details.images = images;
}
if (statusEvents.length > 0) {
details.statusEvents = statusEvents;
}
if (notice) {
details.notice = notice;
}
return details;
};
return await (session.trackEvalExecution?.(execution, sessionAbortController) ?? execution);
const pushUpdate = () => {
emitUpdate?.(tailBuffer.text(), buildUpdateDetails());
};
const sessionFile = session.getSessionFile?.() ?? undefined;
const kernelOwnerId = session.getEvalKernelOwnerId?.() ?? undefined;
const { path: artifactPath, id: artifactId } = (await session.allocateOutputArtifact?.("eval")) ?? {};
session.assertEvalExecutionAllowed?.();
outputSink = new OutputSink({
artifactPath,
artifactId,
headBytes: resolveOutputSinkHeadBytes(session.settings),
maxColumns: resolveOutputMaxColumns(session.settings),
onChunk: chunk => {
appendTail(chunk);
if (activeLiveCell) {
activeLiveCell.buf.append(chunk);
activeLiveCell.result.output = activeLiveCell.buf.text();
}
pushUpdate();
},
});
const sessionId = session.getEvalSessionId?.() ?? defaultEvalSessionId(session);
for (let i = 0; i < cells.length; i++) {
const cell = cells[i];
const backend = cell.resolved.backend;
// The per-cell `timeout` is a budget on the cell runtime's *own*
// work. Host-side `agent()`/`parallel()`/`completion()` bridge calls suspend
// that budget entirely and restart a fresh timeout window when control
// returns to the active backend runtime. Compute, stdout, `log()`/`phase()`, and
// ordinary tool calls all count against the budget. The watchdog drives
// `combinedSignal`; we pass no wall-clock deadline downstream so the
// backends never arm a competing fixed timer.
const idleTimeoutMs =
cell.timeoutMs === 0
? undefined
: clampTimeout("eval", cell.timeoutMs / 1000, session.settings.get("tools.maxTimeout")) * 1000;
const idle = idleTimeoutMs === undefined ? undefined : new IdleTimeout(idleTimeoutMs);
const combinedSignal =
signal && idle
? AbortSignal.any([signal, idle.signal, sessionAbortController.signal])
: signal
? AbortSignal.any([signal, sessionAbortController.signal])
: idle
? AbortSignal.any([idle.signal, sessionAbortController.signal])
: sessionAbortController.signal;
const cellResult = cellResults[i];
cellResult.status = "running";
cellResult.output = "";
cellResult.statusEvents = undefined;
cellResult.exitCode = undefined;
cellResult.durationMs = undefined;
activeLiveCell = { result: cellResult, buf: new TailBuffer(DEFAULT_MAX_BYTES * 2) };
pushUpdate();
const startTime = Date.now();
let result: ExecutorBackendResult;
try {
result = await backend.execute(cell.code, {
cwd: session.cwd,
sessionId,
sessionFile: sessionFile ?? undefined,
kernelOwnerId,
signal: combinedSignal,
session,
idleTimeoutMs,
reset: cell.reset,
onChunk: chunk => {
outputSink!.push(chunk);
},
onStatus: event => {
if (event.op === EVAL_TIMEOUT_PAUSE_OP) {
idle?.pause();
return;
}
if (event.op === EVAL_TIMEOUT_RESUME_OP) {
idle?.resume();
return;
}
cellResult.statusEvents ??= [];
upsertStatusEvent(cellResult.statusEvents, event);
pushUpdate();
},
});
} finally {
idle?.dispose();
activeLiveCell = undefined;
}
const durationMs = Date.now() - startTime;
const cellStatusEvents: EvalStatusEvent[] = [];
const cellDisplayOutputs: EvalDisplayOutput[] = [];
const cellImageNotes: string[] = [];
let cellHasMarkdown = false;
for (const output of result.displayOutputs) {
if (output.type === "json") {
jsonOutputs.push(output.data);
cellDisplayOutputs.push(output);
}
if (output.type === "image") {
const resized = await resizeImage(
{
type: "image",
data: output.data,
mimeType: output.mimeType,
},
{ excludeWebP },
);
const image: ImageContent = {
type: "image",
data: resized.data,
mimeType: resized.mimeType,
};
images.push(image);
cellDisplayOutputs.push({
type: "image",
data: image.data,
mimeType: image.mimeType,
});
const dimensionNote = formatDimensionNote(resized);
if (dimensionNote) {
cellImageNotes.push(`display image ${cellImageNotes.length + 1}: ${dimensionNote}`);
}
}
if (output.type === "status") {
upsertStatusEvent(statusEvents, output.event);
upsertStatusEvent(cellStatusEvents, output.event);
}
if (output.type === "markdown") {
cellHasMarkdown = true;
}
}
const stdoutTrimmed = result.output.trim();
const imageText = cellImageNotes.join("\n");
const displayText = formatDisplayOutputsForText(cellDisplayOutputs);
const visibleDisplayText =
displayText && imageText ? `${displayText}\n\n${imageText}` : displayText || imageText;
const cellOutput =
stdoutTrimmed && visibleDisplayText
? `${stdoutTrimmed}\n\n${visibleDisplayText}`
: stdoutTrimmed || visibleDisplayText;
cellResult.output = cellOutput;
cellResult.exitCode = result.exitCode;
cellResult.durationMs = durationMs;
cellResult.statusEvents = cellStatusEvents.length > 0 ? cellStatusEvents : undefined;
cellResult.hasMarkdown = cellHasMarkdown || undefined;
if (cellOutput) {
cellOutputs.push(cellOutput);
appendTail(cellOutput);
}
if (result.cancelled) {
cellResult.status = "error";
pushUpdate();
const errorMsg = result.output || "Command aborted";
const combinedOutput = cellOutputs.join("\n\n");
const outputText = combinedOutput || errorMsg;
const summaryForMeta = await summarizeFinal(combinedOutput, finalizeOutput);
const details: EvalToolDetails = {
language: languages[0],
languages,
cells: cellResults,
jsonOutputs: jsonOutputs.length > 0 ? jsonOutputs : undefined,
statusEvents: statusEvents.length > 0 ? statusEvents : undefined,
isError: true,
};
if (notice) details.notice = notice;
return toolResult(details)
.content([{ type: "text", text: outputText }, ...images])
.truncationFromSummary(summaryForMeta, { direction: "tail" })
.done();
}
if (result.exitCode !== 0 && result.exitCode !== undefined) {
cellResult.status = "error";
pushUpdate();
const combinedOutput = cellOutputs.join("\n\n");
const outputText = combinedOutput
? `${combinedOutput}\n\nCommand exited with code ${result.exitCode}`
: `Command exited with code ${result.exitCode}`;
const summaryForMeta = await summarizeFinal(combinedOutput, finalizeOutput);
const details: EvalToolDetails = {
language: languages[0],
languages,
cells: cellResults,
jsonOutputs: jsonOutputs.length > 0 ? jsonOutputs : undefined,
statusEvents: statusEvents.length > 0 ? statusEvents : undefined,
isError: true,
};
if (notice) details.notice = notice;
return toolResult(details)
.content([{ type: "text", text: outputText }, ...images])
.truncationFromSummary(summaryForMeta, { direction: "tail" })
.done();
}
cellResult.status = "complete";
pushUpdate();
}
const combinedOutput = cellOutputs.join("\n\n");
const hasImages = images.length > 0;
const outputText =
combinedOutput ||
(hasImages
? `(displayed ${images.length} image${images.length === 1 ? "" : "s"}; no text output)`
: "(no output)");
const summaryForMeta = await summarizeFinal(combinedOutput, finalizeOutput);
const details: EvalToolDetails = {
language: languages[0],
languages,
cells: cellResults,
jsonOutputs: jsonOutputs.length > 0 ? jsonOutputs : undefined,
statusEvents: statusEvents.length > 0 ? statusEvents : undefined,
};
if (notice) details.notice = notice;
return toolResult(details)
.content([{ type: "text", text: outputText }, ...images])
.truncationFromSummary(summaryForMeta, { direction: "tail" })
.done();
} finally {
if (!outputDumped) {
try {
await finalizeOutput();
} catch {}
}
}
}
}
+2 -2
View File
@@ -7,7 +7,7 @@
import type { AgentToolResult } from "@oh-my-pi/pi-agent-core";
import type { Component } from "@oh-my-pi/pi-tui";
import { Text } from "@oh-my-pi/pi-tui";
import type { AsyncJob, AsyncJobManager } from "../../async";
import type { AsyncJob, AsyncJobManager, AsyncJobType } from "../../async";
import { settings } from "../../config/settings";
import type { RenderResultOptions } from "../../extensibility/custom-tools/types";
import { shimmerEnabled, shimmerText } from "../../modes/theme/shimmer";
@@ -145,7 +145,7 @@ function describeAgents(agents: AgentActivitySnapshot[]): string[] {
interface TrackedJobLike {
id: string;
type: "bash" | "task";
type: AsyncJobType;
status: string;
label: string;
startTime: number;
+2 -1
View File
@@ -5,6 +5,7 @@
*/
import type { AgentToolResult } from "@oh-my-pi/pi-agent-core";
import type { AsyncJobType } from "../../async";
import type { IrcDeliveryReceipt, IrcMessage } from "../../irc/bus";
import type { LaunchParams, LaunchToolDetails } from "./launch";
@@ -42,7 +43,7 @@ export interface HubPeerInfo {
/** Background-job row surfaced by `wait`/`cancel`/`jobs` results. */
export interface JobSnapshot {
id: string;
type: "bash" | "task";
type: AsyncJobType;
status: "running" | "completed" | "failed" | "cancelled";
label: string;
durationMs: number;
@@ -5,7 +5,7 @@ import {
ASIDE_MESSAGE_DISCARD,
type CommittableAsideMessage,
} from "@oh-my-pi/pi-agent-core";
import { type AsyncJob, AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async";
import { type AsyncJob, AsyncJobManager, type AsyncJobType } from "@oh-my-pi/pi-coding-agent/async";
import type { CustomMessage } from "@oh-my-pi/pi-coding-agent/session/messages";
import { YieldQueue } from "@oh-my-pi/pi-coding-agent/session/yield-queue";
import type { ToolSession } from "@oh-my-pi/pi-coding-agent/tools";
@@ -21,7 +21,7 @@ type AsyncEntry = {
type AsyncDetails = {
jobs: Array<{
jobId: string;
type?: "bash" | "task";
type?: AsyncJobType;
label?: string;
durationMs?: number;
}>;
@@ -0,0 +1,215 @@
import { afterEach, describe, expect, it, vi } from "bun:test";
import type { AgentToolContext } from "@oh-my-pi/pi-agent-core";
import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async";
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import * as evalIndex from "@oh-my-pi/pi-coding-agent/eval";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { ToolSession } from "@oh-my-pi/pi-coding-agent/tools";
import { EvalTool } from "@oh-my-pi/pi-coding-agent/tools/eval";
function makeSession(settings: Settings, asyncJobManager: AsyncJobManager): ToolSession {
return {
cwd: "/tmp/eval-test",
hasUI: false,
getSessionFile: () => null,
getSessionSpawns: () => null,
settings,
asyncJobManager,
};
}
function baseResult(overrides: Record<string, unknown> = {}) {
return {
output: "",
exitCode: 0,
cancelled: false,
truncated: false,
artifactId: undefined,
totalLines: 0,
totalBytes: 0,
outputLines: 0,
outputBytes: 0,
displayOutputs: [] as unknown[],
...overrides,
};
}
/**
* Mock the JS backend with a cell that streams one chunk immediately and then
* blocks until the returned `release()` gate opens — so backgrounding is decided
* by the tool's own threshold/steer race, never by a guessed sleep.
*/
function mockGatedCell(finalOutput: string): { release: () => void } {
const gate = Promise.withResolvers<void>();
vi.spyOn(evalIndex.jsBackend, "execute").mockImplementation((async (
_code: string,
options: { onChunk?: (chunk: string) => void },
) => {
options.onChunk?.("start\n");
await gate.promise;
return baseResult({ output: finalOutput });
}) as never);
return { release: gate.resolve };
}
function steeringContext(steeringSignal: AbortSignal): AgentToolContext {
return {
sessionManager: SessionManager.inMemory(),
modelRegistry: {
find: () => undefined,
getAll: () => [],
getApiKey: async () => undefined,
} as unknown as AgentToolContext["modelRegistry"],
model: undefined,
isIdle: () => true,
hasQueuedMessages: () => false,
abort: () => {},
toolNames: [],
toolCall: {
batchId: "batch-1",
index: 0,
total: 1,
toolCalls: [{ id: "call-steer", name: "eval" }],
steeringSignal,
},
} as AgentToolContext;
}
/**
* Defends the eval auto-background contract (mirror of bash's): a cell that
* finishes before the threshold resolves inline with no job leftovers, a cell
* that outlives the threshold converts into a running async job whose result is
* delivered later, and a steering interrupt backgrounds the cell immediately so
* the queued message can inject while the kernel keeps working.
*/
describe("EvalTool auto-background", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("keeps fast cells inline and suppresses their job delivery", async () => {
const deliveries: string[] = [];
const asyncJobManager = new AsyncJobManager({
onJobComplete: async (_jobId, text) => {
deliveries.push(text);
},
});
vi.spyOn(evalIndex.jsBackend, "execute").mockImplementation((async () =>
baseResult({ output: "quick\n" })) as never);
const tool = new EvalTool(
makeSession(
Settings.isolated({
"eval.autoBackground.enabled": true,
"eval.autoBackground.thresholdMs": 2_000,
}),
asyncJobManager,
),
);
const result = await tool.execute("call-inline", { language: "js", code: "print('quick')" });
const text = result.content.map(c => (c.type === "text" ? c.text : "")).join("\n");
expect(text).toContain("quick");
expect(result.details?.async).toBeUndefined();
expect(result.details?.cells?.[0]?.status).toBe("complete");
await asyncJobManager.drainDeliveries({ timeoutMs: 1 });
expect(deliveries).toEqual([]);
await asyncJobManager.dispose();
});
it("backgrounds a cell that outlives the threshold and delivers its result", async () => {
const deliveries: Array<{ jobId: string; text: string }> = [];
const updates: string[] = [];
const asyncJobManager = new AsyncJobManager({
onJobComplete: async (jobId, text) => {
deliveries.push({ jobId, text });
},
});
const cell = mockGatedCell("start\ndone\n");
const tool = new EvalTool(
makeSession(
Settings.isolated({
"eval.autoBackground.enabled": true,
"eval.autoBackground.thresholdMs": 10,
}),
asyncJobManager,
),
);
// The gated cell cannot finish on its own, so execute() returning proves
// the threshold path backgrounded it.
const result = await tool.execute(
"call-background",
{ language: "js", code: "print('start'); await work(); print('done')" },
undefined,
update => {
updates.push(update.content?.find(block => block.type === "text")?.text ?? "");
},
);
expect(result.details?.async?.state).toBe("running");
expect(result.details?.async?.type).toBe("eval");
const text = result.content.map(c => (c.type === "text" ? c.text : "")).join("\n");
expect(text).toContain("Backgrounded as job");
// The snapshot keeps the running cell (with its streamed tail) for the transcript.
expect(result.details?.cells?.[0]?.status).toBe("running");
const jobId = result.details?.async?.jobId;
if (!jobId) {
throw new Error("expected an auto-backgrounded job id");
}
const runningJob = asyncJobManager.getJob(jobId);
expect(runningJob?.status).toBe("running");
const updatesAtBackground = updates.slice();
cell.release();
await runningJob?.promise;
await asyncJobManager.drainDeliveries({ timeoutMs: 1 });
expect(deliveries).toHaveLength(1);
expect(deliveries[0]?.jobId).toBe(jobId);
expect(deliveries[0]?.text).toContain("done");
// Tool-call updates stop once the cell is backgrounded.
expect(updates).toEqual(updatesAtBackground);
await asyncJobManager.dispose();
});
it("backgrounds a running cell when the steering signal fires mid-wait", async () => {
const asyncJobManager = new AsyncJobManager({});
const cell = mockGatedCell("steered\n");
const tool = new EvalTool(
makeSession(
Settings.isolated({
"eval.autoBackground.enabled": true,
// High threshold: only the steering signal can background this.
"eval.autoBackground.thresholdMs": 60_000,
}),
asyncJobManager,
),
);
const steering = new AbortController();
steering.abort();
const result = await tool.execute(
"call-steer",
{ language: "js", code: "await work()" },
undefined,
undefined,
steeringContext(steering.signal),
);
// The steer backgrounds the cell instead of killing it: the call returns a
// running job and the cell finishes on its own.
expect(result.details?.async?.state).toBe("running");
const text = result.content.map(c => (c.type === "text" ? c.text : "")).join("\n");
expect(text).toContain("Backgrounded early to handle an incoming message");
const jobId = result.details?.async?.jobId;
if (!jobId) {
throw new Error("expected a steer-backgrounded job id");
}
const job = asyncJobManager.getJob(jobId);
expect(job?.status).toBe("running");
cell.release();
await job?.promise;
expect(asyncJobManager.getJob(jobId)?.status).toBe("completed");
await asyncJobManager.dispose();
});
});