feat(coding-agent-eval): added runEvalAgent bridge for agent plan checks

- Added `agent()` in JS/Python preludes to call host bridge and parse returned text when schema is set.
- Added JS `parallel()` and `pipeline()` with bounded `__pool()` pools and concurrency normalization.
- Added `runEvalAgent` bridge logic with argument parsing plus plan-mode, allowlist, depth, and artifacts checks.
- Added tool routing and tests documenting new `agent/parallel/pipeline` behavior, defaults, and validation failures.
This commit is contained in:
can1357
2026-05-31 06:57:02 +02:00
parent d1bd14f020
commit cf621d0abf
10 changed files with 810 additions and 2 deletions
+3
View File
@@ -10,6 +10,7 @@ It covers tool behavior, runner lifecycle, environment handling, execution seman
- Subprocess kernel client: `src/eval/py/kernel.ts`
- Python wrapper / NDJSON server: `src/eval/py/runner.py`
- Prelude helpers loaded into every kernel: `src/eval/py/prelude.py`
- Host-side subagent helper bridge: `src/eval/agent-bridge.ts`
- MIME bundle renderer (text + structured outputs): `src/eval/py/display.ts`
- Interactive-mode renderer for user-triggered Python runs: `src/modes/components/eval-execution.ts`
- Runtime/env filtering and Python resolution: `src/eval/py/runtime.ts`
@@ -159,6 +160,8 @@ The runner additionally receives `PYTHONUNBUFFERED=1` and `PYTHONIOENCODING=utf-
If Python preflight fails and `eval.js` is enabled, `eval` remains available for `js` cells; `py` cells fail with a Python-backend availability error.
Python prelude helpers include `agent(prompt, *, agent_type="task", model=None, context=None, label=None, schema=None)`. It synchronously calls the host bridge, runs one subagent through the task executor, and returns the final text. When `schema` is supplied, the helper parses the subagent's JSON output and returns the object.
## Execution flow and cancellation/timeout
### Cell timeout
+25 -1
View File
@@ -9,6 +9,7 @@
- Model-facing prompt: `packages/coding-agent/src/prompts/tools/eval.md`
- Key collaborators:
- `packages/coding-agent/src/eval/backend.ts` — backend execution contract
- `packages/coding-agent/src/eval/agent-bridge.ts` — host-side `agent()` bridge into the subagent executor
- `packages/coding-agent/src/eval/js/executor.ts` — JS backend adapter
- `packages/coding-agent/src/eval/js/worker-core.ts` — JS execution, VM context, display/log capture
- `packages/coding-agent/src/eval/js/shared/prelude.txt` — JS global helper installer
@@ -131,11 +132,15 @@ Implemented in `packages/coding-agent/src/eval/js/worker-core.ts`, `packages/cod
- `read`, `write`, `append`, `sort`, `uniq`, `counter`, `diff`, `tree`, `env`, `output`
- `tool.<name>(args)` proxy for arbitrary session tool calls
- `llm(prompt, opts?)` for oneshot, stateless LLM calls (see _Oneshot LLM helper_ below)
- `agent(prompt, opts?)` for a single subagent call, plus JS-only `parallel()` / `pipeline()` bounded-pool helpers (see _Subagent helper_ below)
- JS helpers that touch the host/runtime boundary are async and `await`able; pure text helpers (`sort`, `uniq`, `counter`) return synchronously but may still be safely awaited.
- JS helper signatures use a trailing options object rather than Python keyword arguments:
- `await read(path, { offset?, limit? })`
- `await tree(path = ".", { maxDepth?, hidden? })`
- `sort(text, { reverse?, unique? })`, `uniq(text, { count? })`, `counter(items, { limit?, reverse? })`
- `await agent(prompt, { agentType?, model?, context?, label?, schema? })`
- `await parallel([() => agent("a"), () => agent("b")], { concurrency? })`
- `await pipeline(items, stage1, stage2, { concurrency? })`
- `display(value)` behavior:
- plain objects/arrays become JSON outputs
- `{ type: "image", data, mimeType }` becomes an image output
@@ -156,7 +161,7 @@ Implemented in `packages/coding-agent/src/eval/py/executor.ts`, `packages/coding
- initialize cwd / env / `sys.path`
- execute `PYTHON_PRELUDE`
- Python cells run in the runner's persistent asyncio event loop, so top-level `await` works; the prompt warns not to use `asyncio.run(...)`
- The Python prelude defines helpers with the same surface as JS where practical, including `tool.<name>(args)` through a per-run loopback bridge
- The Python prelude defines helpers with the same surface as JS where practical, including `tool.<name>(args)`, `llm(...)`, and `agent(...)` through a per-run loopback bridge
- Synchronous statement blocks run in the default executor with ContextVar state copied in; the GIL still serializes bytecode execution, but awaited regions can interleave with sibling cells
- Kernel `display_data` / `execute_result` messages map to:
- `application/x-omp-status` → status event
@@ -182,6 +187,21 @@ Both runtimes expose `llm()` — a single stateless completion against a model t
- `schema` (optional) is a plain JSON-Schema object. When present, the model is forced to call a single synthetic `respond` tool with that schema (loose, non-strict), and the helper returns the parsed object. When absent, the helper returns the completion string.
- Errors surface as exceptions: unresolved tier, missing API key, an `error`/`aborted` stop reason, or empty output each raise.
### Subagent helper (`agent`)
Both runtimes expose `agent()` — a single subagent invocation routed through `packages/coding-agent/src/eval/agent-bridge.ts` into the same `runSubprocess(...)` path used by the `task` tool. It uses the current eval session's spawn policy and inherits the parent eval executor id, so parent and subagent code share JS/Python runtime state.
- Signatures:
- JS: `await agent(prompt, { agentType?, model?, context?, label?, schema? })`
- Python: `agent(prompt, *, agent_type="task", model=None, context=None, label=None, schema=None)`
- `agentType` / `agent_type` defaults to the bundled `task` agent and resolves through normal agent discovery, so project and user agents work.
- `model` overrides the selected agent's model. Without it, normal per-agent settings and the agent frontmatter model apply.
- `context` supplies shared background; `label` controls the `agent://<id>` output label prefix.
- `schema` passes a JSON Schema to the subagent structured-output path. When present, the helper parses the final JSON text and returns an object.
- Spawn restrictions use `session.getSessionSpawns()` exactly like the `task` tool. Eval-driven subagent recursion is capped at depth 3.
- JS also exposes `parallel(thunks, { concurrency })` and `pipeline(items, ...stages, { concurrency })`; both use a bounded async pool with default concurrency 4, max 16, preserve item order, and propagate rejections.
- Errors surface as exceptions: unknown or disabled agent, disallowed spawn, recursion cap, subagent failure, or invalid structured output all fail the eval cell.
### Multi-language call behavior
A single tool call can mix Python and JS cells. Persistence is per language runtime:
@@ -202,12 +222,14 @@ A single tool call can mix Python and JS cells. Persistence is per language runt
- Subprocesses / native bindings
- Python availability check runs `<python> -c ...`.
- Python backend spawns one `python -u runner.py` subprocess per kernel; cancellation sends `SIGINT`. Details in `docs/python-repl.md`.
- `agent()` runs one in-process subagent via the task executor; that subagent may use its configured tools.
- Session state
- `session.assertEvalExecutionAllowed?.()` can block execution.
- `session.trackEvalExecution?.(...)` can register cancellable eval work.
- `session.getSessionFile?.()`, `session.getEvalSessionId?.()`, and `session.getEvalKernelOwnerId?.()` influence VM/kernel reuse and artifact lookup.
- JS VM contexts persist across eval calls until reset/disposal.
- Python retained kernels persist until reset, owner cleanup, or process exit.
- `agent()` allocates `agent://<id>` output artifacts and reuses the parent's eval executor id.
- User-visible prompts / interactive UI
- none; stdin requests are rejected programmatically
- Background work / cancellation
@@ -223,6 +245,8 @@ A single tool call can mix Python and JS cells. Persistence is per language runt
- Output truncation window: 50KB default (`DEFAULT_MAX_BYTES` in `packages/coding-agent/src/session/streaming-output.ts`)
- Output line cap inside truncation helpers: 3000 lines (`DEFAULT_MAX_LINES` in `packages/coding-agent/src/session/streaming-output.ts`)
- Streaming tail buffer for live updates: `DEFAULT_MAX_BYTES * 2` = 100KB (`packages/coding-agent/src/tools/eval.ts`)
- JS `parallel()` / `pipeline()` helper concurrency default: 4; maximum: 16
- Eval-driven `agent()` recursion cap: task depth 3 (`EVAL_AGENT_MAX_DEPTH`)
- Python retained kernel idle timeout: 5 minutes (`IDLE_TIMEOUT_MS` in `packages/coding-agent/src/eval/py/executor.ts`)
- Python retained kernel cap: 4 sessions (`MAX_KERNEL_SESSIONS` in `packages/coding-agent/src/eval/py/executor.ts`)
- Python retained kernel cleanup sweep: every 30s (`CLEANUP_INTERVAL_MS` in `packages/coding-agent/src/eval/py/executor.ts`)
+5
View File
@@ -1,8 +1,11 @@
# Changelog
## [Unreleased]
### Added
- Added `agent()` eval options `agent_type`/`agentType`, `model`, `context`, and `label`, and returned structured JSON when `schema` is provided in JS and Python eval cells
- Added `agent()` to the `eval` runtime so JS and Python cells can spawn one subagent through the existing task executor; JS eval also gained bounded `parallel()` and `pipeline()` helpers for orchestrating subagent calls.
- Added search support for virtual internal URLs (including `omp://` roots) by resolving and scanning in-memory internal resources as search targets alongside filesystem paths
- Added expansion of virtual internal URL search targets so `search` can match multiple internal documents when given `omp://`
- Added `/omfg <complaint>` slash command that drafts a TTSR rule from a complaint, validates it against the current conversation, saves it to project or `~/.omp/agent/rules`, and registers it live.
@@ -18,6 +21,8 @@
### Fixed
- Fixed `agent()` in eval to enforce plan-mode, spawn allowlist, and disabled-agent checks before launching subagents
- Fixed recursive `agent()` calls from eval by enforcing the existing max subagent depth limit
- Fixed runtime model switches (Ctrl+P cycling, `--model`, `/model`, model picker selections, and programmatic changes) so they no longer overwrite the persisted `modelRoles.default`; only the model picker's explicit "Set as default" action and settings changes persist the default.
- Fixed `search` to honor line-range suffixes on virtual internal URL targets so matches outside the requested ranges are no longer returned
- Fixed `search` to handle internal URLs without source files without incorrectly reporting `Path not found`, returning matches from virtual content instead
@@ -0,0 +1,342 @@
import { afterAll, afterEach, describe, expect, it, vi } from "bun:test";
import * as path from "node:path";
import { TempDir } from "@oh-my-pi/pi-utils";
import { Settings } from "../../config/settings";
import * as taskDiscovery from "../../task/discovery";
import type { PlanModeState } from "../../plan-mode/state";
import type { ExecutorOptions } from "../../task/executor";
import * as taskExecutor from "../../task/executor";
import { AgentOutputManager } from "../../task/output-manager";
import type { AgentDefinition, SingleResult } from "../../task/types";
import type { ToolSession } from "../../tools";
import { EVAL_AGENT_MAX_DEPTH, runEvalAgent } from "../agent-bridge";
import { disposeAllVmContexts } from "../js/context-manager";
import { executeJs } from "../js/executor";
import { disposeAllKernelSessions, executePython } from "../py/executor";
const taskAgent = {
name: "task",
description: "Task agent",
systemPrompt: "Run the task.",
source: "bundled",
spawns: "*",
model: ["pi/task"],
} satisfies AgentDefinition;
const reviewerAgent = {
name: "reviewer",
description: "Reviewer agent",
systemPrompt: "Review the task.",
source: "bundled",
model: ["pi/smol"],
} satisfies AgentDefinition;
interface SessionOptions {
cwd?: string;
sessionFile?: string | null;
artifactsDir?: string | null;
spawns?: string | null;
depth?: number;
activeModel?: string;
modelString?: string;
enableLsp?: boolean;
settings?: Settings;
outputManager?: AgentOutputManager;
planMode?: boolean;
}
function makeSession(options: SessionOptions = {}): ToolSession {
const settings =
options.settings ??
Settings.isolated({
"async.enabled": false,
"task.isolation.mode": "none",
"task.enableLsp": true,
});
const artifactsDir = options.artifactsDir ?? null;
return {
cwd: options.cwd ?? process.cwd(),
hasUI: false,
settings,
taskDepth: options.depth ?? 0,
enableLsp: options.enableLsp ?? true,
agentOutputManager: options.outputManager,
getSessionFile: () => options.sessionFile ?? null,
getSessionSpawns: () => options.spawns ?? "*",
getActiveModelString: () => options.activeModel ?? "p/active",
getModelString: () => options.modelString ?? "p/fallback",
getArtifactsDir: () => artifactsDir,
getSessionId: () => "test-session",
getEvalSessionId: () => "test-eval-session",
getPlanModeState: options.planMode
? () =>
({
enabled: true,
planFilePath: path.join(options.cwd ?? process.cwd(), "plan.md"),
}) satisfies PlanModeState
: undefined,
};
}
function mockAgents(agents: AgentDefinition[] = [taskAgent, reviewerAgent]): void {
vi.spyOn(taskDiscovery, "discoverAgents").mockResolvedValue({ agents, projectAgentsDir: null });
}
function singleResult(options: ExecutorOptions, overrides: Partial<SingleResult> = {}): SingleResult {
return {
index: options.index,
id: options.id,
agent: options.agent.name,
agentSource: options.agent.source,
task: options.task,
assignment: options.assignment,
description: options.description,
exitCode: 0,
output: "ok",
stderr: "",
truncated: false,
durationMs: 1,
tokens: 0,
...overrides,
};
}
function makeEvalSession(
tempDir: TempDir,
prefix: string,
): { session: ToolSession; sessionFile: string; sessionId: string } {
const sessionFile = path.join(tempDir.path(), "session.jsonl");
const artifactsDir = sessionFile.slice(0, -6);
const session = makeSession({
cwd: tempDir.path(),
sessionFile,
artifactsDir,
outputManager: new AgentOutputManager(() => artifactsDir),
});
return { session, sessionFile, sessionId: `${prefix}:${crypto.randomUUID()}` };
}
describe("runEvalAgent", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("resolves the default task agent and agentType overrides", async () => {
mockAgents();
const runSpy = vi.spyOn(taskExecutor, "runSubprocess").mockImplementation(async options =>
singleResult(options, {
output: options.agent.name,
}),
);
const session = makeSession();
const defaultResult = await runEvalAgent({ prompt: "hello" }, { session });
const overrideResult = await runEvalAgent({ prompt: "hello", agentType: "reviewer" }, { session });
expect(defaultResult.text).toBe("task");
expect(overrideResult.text).toBe("reviewer");
expect(runSpy.mock.calls[0]?.[0].agent.name).toBe("task");
expect(runSpy.mock.calls[1]?.[0].agent.name).toBe("reviewer");
});
it("throws for an unknown agent", async () => {
mockAgents([taskAgent]);
vi.spyOn(taskExecutor, "runSubprocess").mockImplementation(async options => singleResult(options));
await expect(runEvalAgent({ prompt: "hello", agentType: "missing" }, { session: makeSession() })).rejects.toThrow(
'Unknown agent "missing"',
);
});
it("enforces spawn restrictions and the eval recursion cap", async () => {
mockAgents();
const runSpy = vi.spyOn(taskExecutor, "runSubprocess").mockImplementation(async options => singleResult(options));
await expect(runEvalAgent({ prompt: "hello" }, { session: makeSession({ spawns: "" }) })).rejects.toThrow(
"spawns disabled",
);
await expect(runEvalAgent({ prompt: "hello" }, { session: makeSession({ spawns: "reviewer" }) })).rejects.toThrow(
"Allowed: reviewer",
);
await expect(
runEvalAgent({ prompt: "hello" }, { session: makeSession({ depth: EVAL_AGENT_MAX_DEPTH }) }),
).rejects.toThrow("maximum depth");
expect(runSpy).not.toHaveBeenCalled();
});
it("throws instead of spawning from plan mode", async () => {
mockAgents();
const runSpy = vi.spyOn(taskExecutor, "runSubprocess").mockImplementation(async options => singleResult(options));
await expect(runEvalAgent({ prompt: "hello" }, { session: makeSession({ planMode: true }) })).rejects.toThrow(
"unavailable in plan mode",
);
expect(runSpy).not.toHaveBeenCalled();
});
it("passes the parent execution context and only sets outputSchema when schema is supplied", async () => {
mockAgents();
const runSpy = vi.spyOn(taskExecutor, "runSubprocess").mockImplementation(async options => singleResult(options));
const abortController = new AbortController();
const schema = { type: "object", properties: { ok: { type: "boolean" } } };
const session = makeSession({ depth: 2, activeModel: "p/current", modelString: "p/fallback" });
await runEvalAgent(
{ prompt: " hello ", context: " context ", label: "My Agent", model: "p/override", schema },
{ session, signal: abortController.signal },
);
await runEvalAgent({ prompt: "plain" }, { session });
const firstOptions = runSpy.mock.calls[0]?.[0];
const secondOptions = runSpy.mock.calls[1]?.[0];
if (!firstOptions || !secondOptions) throw new Error("runSubprocess was not called");
expect(firstOptions.taskDepth).toBe(2);
expect(firstOptions.signal).toBe(abortController.signal);
expect(firstOptions.parentActiveModelPattern).toBe("p/current");
expect(firstOptions.outputSchema).toBe(schema);
expect(firstOptions.assignment).toBe("hello");
expect(firstOptions.context).toBe("context");
expect(firstOptions.description).toBe("My Agent");
expect(firstOptions.modelOverride).toEqual(["p/override"]);
expect(secondOptions.outputSchema).toBeUndefined();
});
it("maps successful and failed subagent results", async () => {
mockAgents();
const runSpy = vi.spyOn(taskExecutor, "runSubprocess");
runSpy.mockImplementationOnce(async options =>
singleResult(options, {
id: "0-EvalAgent",
output: "done",
resolvedModel: "p/model",
}),
);
runSpy.mockImplementationOnce(async options =>
singleResult(options, {
exitCode: 1,
output: "",
stderr: "stderr",
error: "boom",
}),
);
const result = await runEvalAgent({ prompt: "hello" }, { session: makeSession() });
expect(result).toEqual({
text: "done",
details: { agent: "task", id: "0-EvalAgent", model: "p/model", structured: false },
});
await expect(runEvalAgent({ prompt: "fail" }, { session: makeSession() })).rejects.toThrow("boom");
});
});
describe("agent() through eval runtimes", () => {
afterEach(() => {
vi.restoreAllMocks();
});
afterAll(async () => {
await disposeAllVmContexts();
await disposeAllKernelSessions();
});
it("exposes agent() in JavaScript and parses structured output", async () => {
using tempDir = TempDir.createSync("@omp-eval-agent-js-");
const { session, sessionFile, sessionId } = makeEvalSession(tempDir, "js-agent");
mockAgents();
vi.spyOn(taskExecutor, "runSubprocess").mockImplementation(async options =>
singleResult(options, {
output: options.outputSchema ? '{"ok":true,"n":3}' : "hello from agent",
}),
);
const result = await executeJs(
'const text = await agent("hi"); const data = await agent("json", { schema: { type: "object" } }); return JSON.stringify([text, data]);',
{ cwd: tempDir.path(), sessionId, session, sessionFile },
);
expect(result.exitCode).toBe(0);
expect(JSON.parse(result.output.trim())).toEqual(["hello from agent", { ok: true, n: 3 }]);
});
it("runs JavaScript parallel() with bounded concurrency while preserving order", async () => {
using tempDir = TempDir.createSync("@omp-eval-agent-js-parallel-");
const { session, sessionFile, sessionId } = makeEvalSession(tempDir, "js-agent-parallel");
mockAgents();
let inFlight = 0;
let maxInFlight = 0;
vi.spyOn(taskExecutor, "runSubprocess").mockImplementation(async options => {
inFlight++;
maxInFlight = Math.max(maxInFlight, inFlight);
try {
await Bun.sleep(options.assignment === "a" ? 30 : 10);
return singleResult(options, { output: options.assignment ?? "" });
} finally {
inFlight--;
}
});
const result = await executeJs(
'const values = await parallel(["a", "b", "c", "d"].map(name => () => agent(name)), { concurrency: 2 }); return JSON.stringify(values);',
{ cwd: tempDir.path(), sessionId, session, sessionFile },
);
expect(result.exitCode).toBe(0);
expect(JSON.parse(result.output.trim())).toEqual(["a", "b", "c", "d"]);
expect(maxInFlight).toBeGreaterThan(1);
expect(maxInFlight).toBeLessThanOrEqual(2);
});
it("propagates JavaScript parallel() rejections", async () => {
using tempDir = TempDir.createSync("@omp-eval-agent-js-reject-");
const { session, sessionFile, sessionId } = makeEvalSession(tempDir, "js-agent-reject");
mockAgents();
vi.spyOn(taskExecutor, "runSubprocess").mockImplementation(async options => {
if (options.assignment === "bad") {
return singleResult(options, { exitCode: 1, output: "", stderr: "boom", error: "boom" });
}
return singleResult(options, { output: options.assignment ?? "" });
});
const result = await executeJs('await parallel([() => agent("ok"), () => agent("bad")], { concurrency: 2 });', {
cwd: tempDir.path(),
sessionId,
session,
sessionFile,
});
expect(result.exitCode).toBe(1);
expect(result.output).toContain("boom");
});
it("exposes agent() in the Python runtime", async () => {
using tempDir = TempDir.createSync("@omp-eval-agent-py-");
const { session, sessionFile, sessionId } = makeEvalSession(tempDir, "py-agent");
mockAgents();
vi.spyOn(taskExecutor, "runSubprocess").mockImplementation(async options =>
singleResult(options, { output: "hello from python" }),
);
const probe = await executePython('print("probe")', {
cwd: tempDir.path(),
sessionId: `${sessionId}:probe`,
sessionFile,
kernelMode: "per-call",
});
if (probe.exitCode === undefined && probe.cancelled) {
expect(probe.output).toBe("");
return;
}
expect(probe.exitCode).toBe(0);
const result = await executePython('print(agent("hi"))', {
cwd: tempDir.path(),
sessionId,
sessionFile,
kernelMode: "per-call",
toolSession: session,
});
expect(result.exitCode).toBe(0);
expect(result.output.trim()).toBe("hello from python");
});
});
@@ -0,0 +1,281 @@
/**
* Host-side handler for the eval `agent()` helper.
*/
import * as fs from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import { prompt, Snowflake } from "@oh-my-pi/pi-utils";
import * as z from "zod/v4";
import { resolveAgentModelPatterns } from "../config/model-resolver";
import type { LocalProtocolOptions } from "../internal-urls";
import { MCPManager } from "../mcp/manager";
import subagentUserPromptTemplate from "../prompts/system/subagent-user-prompt.md" with { type: "text" };
import * as taskDiscovery from "../task/discovery";
import * as taskExecutor from "../task/executor";
import { AgentOutputManager } from "../task/output-manager";
import type { AgentDefinition, AgentProgress } from "../task/types";
import type { ToolSession } from "../tools";
import { ToolError } from "../tools/tool-errors";
import type { JsStatusEvent } from "./js/shared/types";
// Import review tools for side effects (registers subagent tool handlers).
import "../tools/review";
/** Synthetic bridge name reserved for the `agent()` helper across both runtimes. */
export const EVAL_AGENT_BRIDGE_NAME = "__agent__";
/** Hard recursion limit for eval-driven subagents. */
export const EVAL_AGENT_MAX_DEPTH = 3;
const DEFAULT_AGENT_TYPE = "task";
const DEFAULT_AGENT_LABEL = "EvalAgent";
const agentArgsSchema = z.object({
prompt: z.string().min(1, "prompt must be a non-empty string"),
agentType: z.string().min(1).optional(),
model: z.union([z.string().min(1), z.array(z.string().min(1)).min(1)]).optional(),
context: z.string().optional(),
label: z.string().optional(),
schema: z.unknown().optional(),
});
interface EvalAgentArgs {
prompt: string;
agentType?: string;
model?: string | string[];
context?: string;
label?: string;
schema?: unknown;
}
export interface EvalAgentBridgeOptions {
session: ToolSession;
signal?: AbortSignal;
emitStatus?: (event: JsStatusEvent) => void;
}
export interface EvalAgentResult {
text: string;
details: {
agent: string;
id: string;
model?: string | string[];
structured: boolean;
};
}
function parseAgentArgs(args: unknown): EvalAgentArgs {
const parsed = agentArgsSchema.safeParse(args);
if (!parsed.success) {
const issue = parsed.error.issues[0];
const where = issue?.path.length ? `${issue.path.join(".")}: ` : "";
throw new ToolError(`agent() received invalid arguments: ${where}${issue?.message ?? "bad input"}`);
}
return parsed.data;
}
function assertDepthAllowed(session: ToolSession): void {
const taskDepth = session.taskDepth ?? 0;
if (taskDepth >= EVAL_AGENT_MAX_DEPTH) {
throw new ToolError(
`agent() cannot spawn another agent at task depth ${taskDepth}; maximum depth is ${EVAL_AGENT_MAX_DEPTH}.`,
);
}
}
function assertSpawnAllowed(session: ToolSession, agentName: string): void {
const parentSpawns = session.getSessionSpawns() ?? "*";
if (parentSpawns === "*") return;
if (parentSpawns === "") {
throw new ToolError(`Cannot spawn '${agentName}'. Allowed: none (spawns disabled for this agent)`);
}
const allowedSpawns = parentSpawns.split(",").map(spawn => spawn.trim());
if (!allowedSpawns.includes(agentName)) {
throw new ToolError(`Cannot spawn '${agentName}'. Allowed: ${parentSpawns}`);
}
}
function assertAgentEnabled(session: ToolSession, agentName: string, agents: AgentDefinition[]): void {
const disabledAgents = session.settings.get("task.disabledAgents") as string[];
if (!disabledAgents.includes(agentName)) return;
const enabled = agents.filter(agent => !disabledAgents.includes(agent.name)).map(agent => agent.name);
throw new ToolError(
`Agent "${agentName}" is disabled in settings. Enable it via /agents, or use a different agent type.${enabled.length > 0 ? ` Available: ${enabled.join(", ")}` : ""}`,
);
}
function assertNotPlanMode(session: ToolSession): void {
if (session.getPlanModeState?.()?.enabled) {
throw new ToolError("agent() is unavailable in plan mode.");
}
}
function renderSubagentPrompt(assignment: string): string {
return prompt.render(subagentUserPromptTemplate, { assignment: assignment.trim(), independentMode: false });
}
function trimToUndefined(value: string | undefined): string | undefined {
const trimmed = value?.trim();
return trimmed ? trimmed : undefined;
}
function outputIdBase(label: string | undefined, agentName: string): string {
const source = trimToUndefined(label) ?? agentName ?? DEFAULT_AGENT_LABEL;
const sanitized = source.replace(/[^A-Za-z0-9_-]+/g, "").slice(0, 48);
return sanitized || DEFAULT_AGENT_LABEL;
}
function getOutputManager(session: ToolSession): AgentOutputManager {
if (session.agentOutputManager) return session.agentOutputManager;
const manager = new AgentOutputManager(session.getArtifactsDir ?? (() => null));
session.agentOutputManager = manager;
return manager;
}
async function getArtifacts(session: ToolSession): Promise<{
sessionFile: string | null;
artifactsDir: string;
contextFile?: string;
}> {
const sessionFile = session.getSessionFile();
const sessionArtifactsDir = sessionFile ? sessionFile.slice(0, -6) : null;
const artifactsDir = sessionArtifactsDir ?? path.join(os.tmpdir(), `omp-eval-agent-${Snowflake.next()}`);
await fs.mkdir(artifactsDir, { recursive: true });
const shouldWriteConversationContext = session.settings.get("irc.enabled") !== true;
const compactContext = shouldWriteConversationContext ? session.getCompactContext?.() : undefined;
if (!compactContext) return { sessionFile, artifactsDir };
const contextFile = path.join(artifactsDir, "context.md");
await Bun.write(contextFile, compactContext);
return { sessionFile, artifactsDir, contextFile };
}
function emitProgressStatus(emitStatus: ((event: JsStatusEvent) => void) | undefined, progress: AgentProgress): void {
emitStatus?.({
op: "agent",
agent: progress.agent,
id: progress.id,
status: progress.status,
lastIntent: progress.lastIntent,
toolCount: progress.toolCount,
durationMs: progress.durationMs,
model: progress.resolvedModel ?? progress.modelOverride,
});
}
/**
* Run a single subagent on behalf of an eval cell's `agent()` call.
*/
export async function runEvalAgent(args: unknown, options: EvalAgentBridgeOptions): Promise<EvalAgentResult> {
const parsed = parseAgentArgs(args);
const agentName = parsed.agentType ?? DEFAULT_AGENT_TYPE;
const structured = Object.hasOwn(parsed, "schema");
assertNotPlanMode(options.session);
assertDepthAllowed(options.session);
assertSpawnAllowed(options.session, agentName);
const { agents } = await taskDiscovery.discoverAgents(options.session.cwd);
const agent = taskDiscovery.getAgent(agents, agentName);
if (!agent) {
const available = agents.map(candidate => candidate.name).join(", ") || "none";
throw new ToolError(`Unknown agent "${agentName}". Available: ${available}`);
}
assertAgentEnabled(options.session, agentName, agents);
const effectiveAgent = agent;
const parentActiveModelPattern = options.session.getActiveModelString?.();
const agentModelOverrides = options.session.settings.get("task.agentModelOverrides");
const modelOverride = resolveAgentModelPatterns({
settingsOverride: parsed.model ?? agentModelOverrides[agentName],
agentModel: effectiveAgent.model,
settings: options.session.settings,
activeModelPattern: parentActiveModelPattern,
fallbackModelPattern: options.session.getModelString?.(),
});
const availableSkills = [...(options.session.skills ?? [])];
const resolvedAutoloadSkills =
effectiveAgent.autoloadSkills?.length && availableSkills.length > 0
? effectiveAgent.autoloadSkills
.map(name => availableSkills.find(skill => skill.name === name))
.filter((skill): skill is NonNullable<typeof skill> => skill !== undefined)
: [];
const contextFiles = options.session.contextFiles?.filter(
file => path.basename(file.path).toLowerCase() !== "agents.md",
);
const localProtocolOptions: LocalProtocolOptions = options.session.localProtocolOptions ?? {
getArtifactsDir: options.session.getArtifactsDir ?? (() => null),
getSessionId: options.session.getSessionId ?? (() => null),
};
const parentArtifactManager = options.session.getArtifactManager?.() ?? undefined;
const parentEvalSessionId = options.session.getEvalSessionId?.() ?? undefined;
const mcpManager = options.session.mcpManager ?? MCPManager.instance();
const { sessionFile, artifactsDir, contextFile } = await getArtifacts(options.session);
const outputManager = getOutputManager(options.session);
const id = await outputManager.allocate(outputIdBase(parsed.label, agentName));
const assignment = parsed.prompt.trim();
const context = trimToUndefined(parsed.context);
const result = await taskExecutor.runSubprocess({
cwd: options.session.cwd,
agent: effectiveAgent,
task: renderSubagentPrompt(assignment),
assignment,
context,
description: trimToUndefined(parsed.label),
index: 0,
id,
taskDepth: options.session.taskDepth ?? 0,
modelOverride,
parentActiveModelPattern,
thinkingLevel: effectiveAgent.thinkingLevel,
outputSchema: structured ? parsed.schema : undefined,
sessionFile,
persistArtifacts: Boolean(sessionFile),
artifactsDir,
contextFile,
enableLsp: (options.session.enableLsp ?? true) && options.session.settings.get("task.enableLsp"),
signal: options.signal,
eventBus: options.session.eventBus,
onProgress: progress => emitProgressStatus(options.emitStatus, progress),
authStorage: options.session.authStorage,
modelRegistry: options.session.modelRegistry,
settings: options.session.settings,
mcpManager,
contextFiles,
skills: availableSkills,
autoloadSkills: resolvedAutoloadSkills,
workspaceTree: options.session.workspaceTree,
promptTemplates: options.session.promptTemplates,
localProtocolOptions,
parentArtifactManager,
parentHindsightSessionState: options.session.getHindsightSessionState?.(),
parentMnemosyneSessionState: options.session.getMnemosyneSessionState?.(),
parentTelemetry: options.session.getTelemetry?.(),
parentEvalSessionId,
});
if (result.exitCode !== 0 || result.error) {
const failureMessage =
result.error ?? result.stderr ?? result.abortReason ?? `agent() subagent '${agentName}' failed.`;
throw new ToolError(failureMessage);
}
options.emitStatus?.({
op: "agent",
agent: result.agent,
id: result.id,
status: "completed",
chars: result.output.length,
model: result.resolvedModel ?? modelOverride,
});
return {
text: result.output,
details: {
agent: result.agent,
id: result.id,
model: result.resolvedModel ?? modelOverride,
structured,
},
};
}
@@ -39,11 +39,63 @@ if (!globalThis.__omp_js_prelude_loaded__) {
return values.length === 1 ? values[0] : values;
};
const hasOwn = (object, key) => Object.prototype.hasOwnProperty.call(object, key);
const llm = async (prompt, opts = {}) => {
const o = toOptions(opts);
const res = await globalThis.__omp_call_tool__("__llm__", { prompt, ...o });
const text = res && typeof res === "object" ? res.text : res;
return o.schema ? JSON.parse(text) : text;
return hasOwn(o, "schema") ? JSON.parse(text) : text;
};
const agent = async (prompt, opts = {}) => {
const o = toOptions(opts);
const res = await globalThis.__omp_call_tool__("__agent__", { prompt, ...o });
const text = res && typeof res === "object" ? res.text : res;
return hasOwn(o, "schema") ? JSON.parse(text) : text;
};
const normalizeConcurrency = value => {
const number = Number(value ?? 4);
if (!Number.isFinite(number)) return 4;
return Math.max(1, Math.min(16, Math.trunc(number)));
};
const __pool = async (items, limit, fn) => {
const list = Array.from(items ?? []);
const concurrency = Math.min(normalizeConcurrency(limit), list.length);
const results = new Array(list.length);
let next = 0;
const worker = async () => {
while (true) {
const index = next++;
if (index >= list.length) return;
results[index] = await fn(list[index], index);
}
};
await Promise.all(Array.from({ length: concurrency }, () => worker()));
return results;
};
const parallel = async (thunks, opts = {}) =>
__pool(thunks, toOptions(opts).concurrency, (thunk, index) => {
if (typeof thunk !== "function") throw new TypeError("parallel() expects an iterable of functions");
return thunk(index);
});
const pipeline = async (items, ...stagesAndOptions) => {
let opts = {};
const last = stagesAndOptions.at(-1);
if (last && typeof last === "object" && !Array.isArray(last)) {
opts = last;
stagesAndOptions = stagesAndOptions.slice(0, -1);
}
let current = Array.from(items ?? []);
for (const stage of stagesAndOptions) {
if (typeof stage !== "function") throw new TypeError("pipeline() stages must be functions");
current = await __pool(current, toOptions(opts).concurrency, stage);
}
return current;
};
const display = value => {
@@ -70,6 +122,10 @@ if (!globalThis.__omp_js_prelude_loaded__) {
globalThis.tool = tool;
globalThis.llm = llm;
globalThis.output = output;
globalThis.agent = agent;
globalThis.parallel = parallel;
globalThis.pipeline = pipeline;
globalThis.__pool = __pool;
globalThis.read = read;
globalThis.write = write;
globalThis.append = append;
@@ -1,6 +1,7 @@
import type { AgentTool, AgentToolResult } from "@oh-my-pi/pi-agent-core";
import type { ToolSession } from "../../tools";
import { ToolError } from "../../tools/tool-errors";
import { EVAL_AGENT_BRIDGE_NAME, runEvalAgent } from "../agent-bridge";
import { EVAL_LLM_BRIDGE_NAME, runEvalLlm } from "../llm-bridge";
import type { JsStatusEvent } from "./shared/types";
@@ -105,6 +106,9 @@ export async function callSessionTool(name: string, args: unknown, options: Tool
if (name === EVAL_LLM_BRIDGE_NAME) {
return await runEvalLlm(args, options);
}
if (name === EVAL_AGENT_BRIDGE_NAME) {
return await runEvalAgent(args, options);
}
const tool = getTool(options.session, name);
const normalizedArgs = normalizeArgs(args);
const toolCallId = `js-${name}-${crypto.randomUUID()}`;
@@ -479,3 +479,26 @@ if "__omp_prelude_loaded__" not in globals():
res = _bridge_call("__llm__", args)
text = res.get("text") if isinstance(res, dict) else res
return json.loads(text) if schema is not None else text
def agent(prompt, *, agent_type="task", model=None, context=None, label=None, schema=None):
"""Run a subagent and return its final output.
`agent_type` selects the subagent definition (default "task"). Pass
`model` to override that agent's model, `context` for shared background,
`label` for the output artifact id, and `schema` to request structured
JSON output; when `schema` is supplied the parsed object is returned.
"""
args = {"prompt": prompt}
if agent_type is not None:
args["agentType"] = agent_type
if model is not None:
args["model"] = model
if context is not None:
args["context"] = context
if label is not None:
args["label"] = label
if schema is not None:
args["schema"] = schema
res = _bridge_call("__agent__", args)
text = res.get("text") if isinstance(res, dict) else res
return json.loads(text) if schema is not None else text
@@ -46,6 +46,13 @@ tool.<name>(args) → unknown
Invoke any session tool by name. `args` is the tool's parameter object.
llm(prompt, model?="default", system?=None, schema?=None) → str | dict
Oneshot, stateless LLM call (no history, no tools). `model` picks a tier: "smol" (fast), "default" (this session's model), "slow" (most capable). Pass `system` for a system prompt. Pass a JSON-Schema `schema` to force structured output and get the parsed object back; otherwise returns the completion text.
agent(prompt, agent_type?="task", model?=None, context?=None, label?=None, schema?=None) → str | dict
Run a subagent and return its final output. Defaults to the bundled "task" agent; pass `agent_type`/`agentType` for another discovered agent. Pass a JSON-Schema `schema` to force structured output and get the parsed object back.
{{#if js}}parallel(thunks, { concurrency?=4 }) → list
Run async thunks through a bounded pool (default concurrency 4, max 16), preserving input order and propagating failures.
pipeline(items, ...stages, { concurrency?=4 }) → list
Apply async stage functions to each item with the same bounded-pool semantics; stage failures propagate.
{{/if}}
```
</prelude>
@@ -0,0 +1,63 @@
import { afterEach, describe, expect, it, vi } from "bun:test";
import { Settings } from "../../src/config/settings";
import { runEvalAgent } from "../../src/eval/agent-bridge";
import type { LocalProtocolOptions } from "../../src/internal-urls";
import type { MCPManager } from "../../src/mcp";
import * as taskDiscovery from "../../src/task/discovery";
import * as taskExecutor from "../../src/task/executor";
import type { AgentDefinition, SingleResult } from "../../src/task/types";
import type { ToolSession } from "../../src/tools";
function createResult(): SingleResult {
return {
index: 0,
id: "0-Task",
agent: "task",
agentSource: "bundled",
task: "do work",
exitCode: 0,
output: "done",
stderr: "",
truncated: false,
durationMs: 1,
tokens: 0,
};
}
describe("runEvalAgent", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("forwards session-scoped MCP and local protocol options", async () => {
const agent: AgentDefinition = {
name: "task",
description: "Task agent",
systemPrompt: "Handle task",
source: "bundled",
};
vi.spyOn(taskDiscovery, "discoverAgents").mockResolvedValue({ agents: [agent], projectAgentsDir: null });
const runSubprocessSpy = vi.spyOn(taskExecutor, "runSubprocess").mockResolvedValue(createResult());
const mcpManager = { sentinel: "mcp" } as unknown as MCPManager;
const localProtocolOptions: LocalProtocolOptions = {
getArtifactsDir: () => "/tmp/parent-artifacts",
getSessionId: () => "parent-session",
};
const session = {
cwd: "/tmp",
settings: Settings.isolated(),
getSessionSpawns: () => "*",
getSessionFile: () => null,
mcpManager,
localProtocolOptions,
} as unknown as ToolSession;
await runEvalAgent({ prompt: "do work", agentType: "task" }, { session });
expect(runSubprocessSpy).toHaveBeenCalledTimes(1);
const options = runSubprocessSpy.mock.calls[0]?.[0];
expect(options?.mcpManager).toBe(mcpManager);
expect(options?.localProtocolOptions).toBe(localProtocolOptions);
});
});