From 9c48d16a6b0f375db9e966f718ba5054fd70d1c7 Mon Sep 17 00:00:00 2001 From: Teakowa <27560638+Teakowa@users.noreply.github.com> Date: Sat, 3 Oct 2026 22:49:49 +0800 Subject: [PATCH 1/4] feat(bench): add an mcp tool level beside bin for the agent benchmark MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds `level` to the condition grid so the same scenario, prompt, grader, model, and trial can run with Wright as CLI (bin, the default) or as native MCP tools (mcp). Under mcp the wright CLI stays off PATH and the harness hands the adapter BENCH_MCP_CMD — `wright serve --transport mcp` through the existing tracing shim — while the canary enforces that the CLI is unreachable. - agent_bench: LEVELS + --level flag, cell normalization/labels (wright-mcp+.../knowledge/network), wright-only level validation, BENCH_TOOL_LEVEL/BENCH_MCP_CMD env, mcp no-server invalid check, and result.toolCalls from the normalized transcript. - bench_trace: serve lines now carry their session id and transport; serve_request mirrors each transport's answer rule (blank lines and JSON-RPC notifications are silent, everything else is answered), so request/response pairing cannot desynchronize and concurrent sessions do not cross-pair. tools/call maps to the wright op its tool name carries; handshake methods keep mcp: names — excluded from uses but their response bytes (the tools/list schema payload) count as output. Structured refusals are refused, only parse-level errors are malformed. - direct adapter: stdio JSON-RPC client that initializes, lists tools, converts schemas to model tools, dispatches calls, and closes reliably; the module state resets per run and setup failures exit 1 (agent error), not provider-interrupted. - bench_report: paired mcp-vs-bin section on (scenario, agent, trial, knowledge, network, skills) with Wilson intervals on usable/passed, search/read counts, wright/bash/tool calls, turns, and tokens. Closes #474 --- benchmarks/agent/adapters/direct.py | 139 +++++++++++++--- benchmarks/agent/agent_bench.py | 52 ++++-- benchmarks/agent/bench_report.py | 72 ++++++++ benchmarks/agent/bench_trace.py | 237 +++++++++++++++++++++++---- benchmarks/agent/matrix.example.json | 7 + benchmarks/agent/test_adapters.py | 63 ++++++- benchmarks/agent/test_agent_bench.py | 157 +++++++++++++++++- docs/agent-benchmark.md | 43 +++-- 8 files changed, 685 insertions(+), 85 deletions(-) diff --git a/benchmarks/agent/adapters/direct.py b/benchmarks/agent/adapters/direct.py index 05e5c244..a0f75fb8 100644 --- a/benchmarks/agent/adapters/direct.py +++ b/benchmarks/agent/adapters/direct.py @@ -4,6 +4,8 @@ Honors the BENCH_* contract (docs/agent-benchmark.md). BENCH_MODEL is `anthropic/` (ANTHROPIC_API_KEY, optional ANTHROPIC_BASE_URL) or `openai/` (OPENAI_API_KEY, optional OPENAI_BASE_URL, so any OpenAI-compatible endpoint works). The model gets one `bash` tool that runs in the workspace on the shimmed PATH, plus `fetch` only when knowledge is `web`. +Under tool level `mcp` it also registers `BENCH_MCP_CMD` (`wright serve --transport mcp`) and exposes the server's tools +natively; the `wright` CLI is then not on PATH. The system prompt lists the installed skills by name and description, and the model reads their files itself. The loop, its limits, and the tool set are fixed here so every model faces the same protocol; the recorded harness commit identifies them. It does not sandbox the network: pair it with the harness --canary-cmd. Exit 75 marks a provider or infrastructure failure. @@ -14,6 +16,8 @@ import json import os import re +import selectors +import shlex import shutil import subprocess import sys @@ -122,6 +126,82 @@ def step(self) -> tuple[str, list[tuple[str, str, dict]], dict]: return message.get("content") or "", calls, usage_row((u.get("prompt_tokens") or 0) - cached, (u.get("completion_tokens") or 0) - reasoning, cached, 0, reasoning) +class Mcp: + """A `wright serve --transport mcp` server over stdio (ADR-0020): line-delimited JSON-RPC, one request at a time.""" + + def __init__(self, command: str, env: dict): + self.proc = subprocess.Popen(["/bin/sh", "-c", command], stdin=subprocess.PIPE, stdout=subprocess.PIPE, text=True, env=env) + self.next_id = 0 + + def request(self, method: str, params: dict | None = None) -> dict: + self.next_id += 1 + message = {"jsonrpc": "2.0", "id": self.next_id, "method": method} + if params is not None: + message["params"] = params + try: + self.proc.stdin.write(json.dumps(message) + "\n") + self.proc.stdin.flush() + for line in self._lines(): + reply = json.loads(line) + if isinstance(reply, dict) and reply.get("id") == self.next_id: + return reply + except (BrokenPipeError, json.JSONDecodeError, OSError): + pass + return {"error": {"code": -32000, "message": "the MCP server closed or timed out"}} + + def _lines(self): + """Reply lines, bounded to COMMAND_SECONDS where the OS can poll a pipe; a hung server then reads as a tool error.""" + if os.name != "posix": + yield from self.proc.stdout + return + deadline = time.monotonic() + COMMAND_SECONDS + selector = selectors.DefaultSelector() + with selector: + selector.register(self.proc.stdout, selectors.EVENT_READ) + while selector.select(max(0.0, deadline - time.monotonic())): + line = self.proc.stdout.readline() + if not line: + return + yield line + + def notify(self, method: str) -> None: + try: + self.proc.stdin.write(json.dumps({"jsonrpc": "2.0", "method": method}) + "\n") + self.proc.stdin.flush() + except (BrokenPipeError, OSError): + pass + + def close(self) -> None: + try: + self.proc.stdin.close() + self.proc.wait(timeout=10) + except Exception: + self.proc.kill() + self.proc.stdout.close() + + +MCP: Mcp | None = None + + +def mcp_tools(command: str, env: dict) -> list[dict]: + """Register the MCP server and return its tools in this loop's {name, description, schema} form.""" + global MCP + MCP = Mcp(command, env) + try: + initialized = MCP.request("initialize", {"protocolVersion": "2025-06-18", "capabilities": {}, "clientInfo": {"name": "wright-agent-bench", "version": "1"}}) + if "error" in initialized: + raise ProviderError(f"MCP initialize failed: {initialized['error'].get('message')}", False) + MCP.notify("notifications/initialized") + listed = (MCP.request("tools/list").get("result") or {}).get("tools") or [] + if not listed: + raise ProviderError("MCP tools/list returned no tools", False) + except ProviderError: + MCP.close() + MCP = None + raise + return [{"name": tool["name"], "description": tool.get("description") or "", "schema": tool.get("inputSchema") or {"type": "object"}} for tool in listed] + + def run_tool(name: str, args: dict, env: dict) -> str: try: if name == "bash": @@ -130,6 +210,14 @@ def run_tool(name: str, args: dict, env: dict) -> str: elif name == "fetch": with urllib.request.urlopen(args["url"], timeout=60) as response: out = response.read(OUTPUT_CHARS * 2).decode(errors="replace") + elif MCP: + reply = MCP.request("tools/call", {"name": name, "arguments": args}) + if "error" in reply: + return f"[mcp error {reply['error'].get('code')}: {reply['error'].get('message')}]" + result = reply.get("result") or {} + out = "\n".join(block.get("text", "") for block in result.get("content") or [] if isinstance(block, dict) and block.get("type") == "text") + if result.get("isError"): + out = f"[isError] {out}" else: return f"unknown tool {name}" except subprocess.TimeoutExpired: @@ -160,35 +248,48 @@ def main() -> int: return 2 web = env["BENCH_KNOWLEDGE"] == "web" tools = [BASH] + ([FETCH] if web else []) + child_env = {k: v for k, v in env.items() if k not in ("BENCH_HOST_PATH", "ANTHROPIC_API_KEY", "OPENAI_API_KEY")} + global MCP + MCP = None + if env.get("BENCH_MCP_CMD"): + try: + tools += mcp_tools(env["BENCH_MCP_CMD"], child_env) + except ProviderError as error: + print(f"mcp setup: {error}", file=sys.stderr) + return 1 listing, loaded = skill_listing([Path(p) for p in env.get("BENCH_SKILL_DIRS", "").split(os.pathsep) if p]) effort = env.get("BENCH_THINKING") if provider == "openai" else None system = SYSTEM + listing chat = Anthropic(model, tools, system) if provider == "anthropic" else OpenAI(model, tools, system, effort) Path(env["BENCH_AGENT_INFO"]).write_text(json.dumps({ "agent": "direct", "model": env["BENCH_MODEL"], "effort": effort, "tools": [t["name"] for t in tools], + "toolLevel": env.get("BENCH_TOOL_LEVEL", "bin"), "protocol": {"maxTurns": MAX_TURNS, "commandSeconds": COMMAND_SECONDS, "outputChars": OUTPUT_CHARS, "maxOutputTokens": MAX_OUTPUT_TOKENS}}, indent=2)) Path(env["BENCH_CONTEXT"]).write_text(json.dumps({"loaded": loaded})) - child_env = {k: v for k, v in env.items() if k not in ("BENCH_HOST_PATH", "ANTHROPIC_API_KEY", "OPENAI_API_KEY")} chat.user(sys.stdin.read()) final, code = "", 0 - with open(env["BENCH_USAGE"], "w", buffering=1) as usage, open(env["BENCH_TRANSCRIPT"], "w", buffering=1) as transcript: - transcript.write(json.dumps({"t": time.time(), "type": "system", "text": system}) + "\n") - for _ in range(MAX_TURNS): - try: - text, calls, row = chat.step() - except ProviderError as error: - print(f"provider error: {error}", file=sys.stderr) - code = INFRA_EXIT if error.transient else 1 - break - usage.write(json.dumps(row) + "\n") - transcript.write(json.dumps({"t": time.time(), "type": "assistant", "text": text, "calls": [{"name": n, "input": a} for _, n, a in calls]}) + "\n") - final = text - if not calls: - break - outputs = [(i, run_tool(n, a, child_env)) for i, n, a in calls] - for (_, n, a), (_, out) in zip(calls, outputs): - transcript.write(json.dumps({"t": time.time(), "type": "tool_result", "name": n, "output": out}) + "\n") - chat.results(outputs) + try: + with open(env["BENCH_USAGE"], "w", buffering=1) as usage, open(env["BENCH_TRANSCRIPT"], "w", buffering=1) as transcript: + transcript.write(json.dumps({"t": time.time(), "type": "system", "text": system}) + "\n") + for _ in range(MAX_TURNS): + try: + text, calls, row = chat.step() + except ProviderError as error: + print(f"provider error: {error}", file=sys.stderr) + code = INFRA_EXIT if error.transient else 1 + break + usage.write(json.dumps(row) + "\n") + transcript.write(json.dumps({"t": time.time(), "type": "assistant", "text": text, "calls": [{"name": n, "input": a} for _, n, a in calls]}) + "\n") + final = text + if not calls: + break + outputs = [(i, run_tool(n, a, child_env)) for i, n, a in calls] + for (_, n, a), (_, out) in zip(calls, outputs): + transcript.write(json.dumps({"t": time.time(), "type": "tool_result", "name": n, "output": out}) + "\n") + chat.results(outputs) + finally: + if MCP: + MCP.close() sys.stdout.write(final) return code diff --git a/benchmarks/agent/agent_bench.py b/benchmarks/agent/agent_bench.py index dac30b6d..29425c9d 100644 --- a/benchmarks/agent/agent_bench.py +++ b/benchmarks/agent/agent_bench.py @@ -37,6 +37,7 @@ SCENARIOS = HERE / "scenarios" RESULT_CONTRACT = "wright-agent-bench/v3" TOOLS = ("none", "wright", "overpy") +LEVELS = ("bin", "mcp") # how a `wright` tool reaches the agent: the CLI shim on PATH, or `wright serve --transport mcp` the adapter registers SKILLS = ("wright-skill", "workshop-skill", "opy-skill", "workshop-format-skill") SKILL_LANGUAGE = {"opy-skill": "opy", "workshop-format-skill": "workshop"} # skills that teach one language and apply to its scenarios only KNOWLEDGE_LEVELS = ("none", "wiki", "web") @@ -99,12 +100,13 @@ def baseline_path(path: str) -> str: def normalize_cell(raw: dict) -> dict: - return {"tool": raw["tool"], "skills": sorted(raw.get("skills") or []), "knowledge": raw["knowledge"], "network": raw["network"]} + return {"tool": raw["tool"], "level": raw.get("level") or "bin", "skills": sorted(raw.get("skills") or []), "knowledge": raw["knowledge"], "network": raw["network"]} def cell_label(cell: dict) -> str: - """`tool[+skill...]/knowledge/network`, for example `wright+wright-skill/none/off`; the baseline is `none/none/off`.""" - return f"{'+'.join([cell['tool'], *cell['skills']])}/{cell['knowledge']}/{cell['network']}" + """`tool[+skill...]/knowledge/network`, for example `wright+wright-skill/none/off`; a non-`bin` level suffixes the tool: `wright-mcp+.../none/off`.""" + tool = cell["tool"] if cell.get("level") in (None, "bin") else f"{cell['tool']}-{cell['level']}" + return f"{'+'.join([tool, *cell['skills']])}/{cell['knowledge']}/{cell['network']}" def applicable(scenario: dict, cell: dict) -> bool: @@ -144,8 +146,11 @@ def overpy_launcher(out: Path, host_path: str) -> Path: def check_cell(cell: dict, args: argparse.Namespace) -> None: - if cell["tool"] not in TOOLS or cell["knowledge"] not in KNOWLEDGE_LEVELS or cell["network"] not in ("off", "on") or any(s not in SKILLS for s in cell["skills"]): + level = cell.get("level") or "bin" # `normalize_cell` fills this; a hand-built cell defaults to `bin` the same way + if cell["tool"] not in TOOLS or level not in LEVELS or cell["knowledge"] not in KNOWLEDGE_LEVELS or cell["network"] not in ("off", "on") or any(s not in SKILLS for s in cell["skills"]): raise SystemExit(f"invalid condition {cell_label(cell)}") + if level != "bin" and cell["tool"] != "wright": + raise SystemExit(f"invalid condition {cell_label(cell)}: level '{level}' is a `wright` level") if cell["knowledge"] == "web" and cell["network"] != "on": raise SystemExit("knowledge 'web' requires network 'on'") for name in cell["skills"]: @@ -169,20 +174,24 @@ def build_env(cell: dict, args: argparse.Namespace, out: Path, workspace: Path) env.setdefault("HOME", str(home)) path = baseline_path(os.environ["PATH"]) if cell["tool"] != "none": - shim_dir = out / "bin" - shim_dir.mkdir() - shim = shim_dir / cell["tool"] - shim.write_text(f'#!/bin/sh\nexec "{sys.executable}" "{(HERE / "bench_trace.py").resolve()}" shim {cell["tool"]} "$@"\n') - shim.chmod(0o755) real = args.wright if cell["tool"] == "wright" else str(overpy_launcher(out, os.environ["PATH"])) env.update({f"BENCH_TOOL_REAL_{cell['tool'].upper()}": real, "BENCH_TOOL_TRACE": str(out / "tool-trace.jsonl"), "BENCH_TOOL_SIDECAR": str(out / "tool-calls")}) - path = f"{shim_dir}{os.pathsep}{path}" + if cell["tool"] == "wright" and cell.get("level") == "mcp": + # the adapter registers the traced MCP server; the `wright` CLI itself stays off PATH (the canary enforces that) + env["BENCH_MCP_CMD"] = shlex.join([sys.executable, str((HERE / "bench_trace.py").resolve()), "shim", "wright", "serve", "--transport", "mcp", str(workspace)]) + else: + shim_dir = out / "bin" + shim_dir.mkdir() + shim = shim_dir / cell["tool"] + shim.write_text(f'#!/bin/sh\nexec "{sys.executable}" "{(HERE / "bench_trace.py").resolve()}" shim {cell["tool"]} "$@"\n') + shim.chmod(0o755) + path = f"{shim_dir}{os.pathsep}{path}" env.update( PATH=path, BENCH_HOST_PATH=os.environ["PATH"], BENCH_WORKSPACE=str(workspace), BENCH_RUN_DIR=str(out), BENCH_AGENT_ID=args.agent_id, BENCH_USAGE=str(out / "usage.jsonl"), BENCH_TRANSCRIPT=str(out / "transcript.jsonl"), BENCH_CONTEXT=str(out / "context.json"), BENCH_AGENT_INFO=str(out / "agent-info.json"), - BENCH_KNOWLEDGE=cell["knowledge"], BENCH_NETWORK=cell["network"], BENCH_TOOL=cell["tool"], BENCH_SKILLS=",".join(cell["skills"]), + BENCH_KNOWLEDGE=cell["knowledge"], BENCH_NETWORK=cell["network"], BENCH_TOOL=cell["tool"], BENCH_TOOL_LEVEL=cell.get("level") or "bin", BENCH_SKILLS=",".join(cell["skills"]), ) if cell["skills"]: env["BENCH_SKILL_DIRS"] = os.pathsep.join(str(Path(args.skill_dirs[name]).resolve()) for name in cell["skills"]) @@ -202,8 +211,9 @@ def canaries(cell: dict, env: dict, workspace: Path, args: argparse.Namespace) - if args.check_ancestors and (found := ancestor_instructions(workspace)): return f"instruction files in ancestor directories of the workspace: {found}; use --out outside the repository" for tool in ("wright", "overpy"): - if cell["tool"] != tool and shutil.which(tool, path=env["PATH"]): - return f"{tool} reachable although the tool is '{cell['tool']}'" + shimmed = cell["tool"] == tool and cell.get("level", "bin") == "bin" # level `mcp` keeps the CLI off PATH on purpose + if not shimmed and shutil.which(tool, path=env["PATH"]): + return f"{tool} reachable although the tool is '{cell['tool']}' at level '{cell.get('level', 'bin')}'" if cell["network"] == "off" and args.canary_cmd: if subprocess.run(args.canary_cmd, shell=True, cwd=workspace, env=env, capture_output=True).returncode == 0: return "network reachable under network 'off'" @@ -426,15 +436,22 @@ def run_trial(scenario: dict, cell: dict, args: argparse.Namespace, out: Path) - result["fileReadEnforcement"] = "unrestricted" context = context_report(out, [skill_name(Path(args.skill_dirs[name])) for name in cell["skills"]]) result["context"] = context + reasons = [] if context.get("unexpected"): - result["invalid"] = f"unexpected loaded context: {context['unexpected']}" + reasons.append(f"unexpected loaded context: {context['unexpected']}") if (out / "agent-info.json").is_file(): result["agentInfo"] = json.loads((out / "agent-info.json").read_text()) events = bench_trace.read_events(out / "tool-trace.jsonl") result["toolUse"] = bench_trace.summarize_trace(events) + if cell.get("level") == "mcp" and agent_exit != INFRA_EXIT and not any(e.get("type") == "serve" for e in bench_trace.tool_events(events, "wright")): + reasons.append("level 'mcp' but the adapter never started the wright MCP server (BENCH_MCP_CMD)") + if calls := bench_trace.call_counts(out / "transcript.jsonl"): + result["toolCalls"] = calls result.update(bench_grade.grade(scenario, workspace, args.wright, out / "grading")) if any(c.get("unavailable") for c in result["checks"]): - result["invalid"] = "a required grader was unavailable" + reasons.append("a required grader was unavailable") + if reasons: + result["invalid"] = "; ".join(reasons) entry = workspace / scenario["entry"] final_sha = hashlib.sha256(entry.read_bytes()).hexdigest() if entry.is_file() else None result["friction"] = bench_trace.friction(events) @@ -520,7 +537,7 @@ def verdict_word(result: dict) -> str: def cmd_run(args: argparse.Namespace) -> int: scenario = load_scenario(args.scenario) - cell = normalize_cell({"tool": args.tool, "skills": args.skills, "knowledge": args.knowledge, "network": args.network}) + cell = normalize_cell({"tool": args.tool, "level": args.level, "skills": args.skills, "knowledge": args.knowledge, "network": args.network}) code = 0 for trial in range(1, args.trials + 1): result = run_trial(scenario, cell, args, trial_dir(args.out, args.scenario, args.agent_id, cell, trial)) @@ -822,11 +839,12 @@ def main() -> int: run.add_argument("--agent-cmd", required=True, help="shell command; the task prompt arrives on stdin, cwd is the workspace, BENCH_* describes the condition") run.add_argument("--agent-id", required=True, help="recorded agent/model/version label") run.add_argument("--tool", choices=TOOLS, default="wright") + run.add_argument("--level", choices=LEVELS, default="bin", help="how a `wright` tool reaches the agent: `bin` puts the CLI shim on PATH; `mcp` gives the adapter BENCH_MCP_CMD to register `wright serve --transport mcp`") run.add_argument("--skills", nargs="*", choices=SKILLS, default=[], help="skills installed in this condition") run.add_argument("--knowledge", choices=KNOWLEDGE_LEVELS, default="none") run.add_argument("--network", choices=("off", "on"), default="off") run.add_argument("--trials", type=int, default=1) - sub.choices["matrix"].add_argument("config", type=Path, help="JSON: agents[{id,cmd}], cells[{tool,skills,knowledge,network}], scenarios, trials, parallel, seed, options") + sub.choices["matrix"].add_argument("config", type=Path, help="JSON: agents[{id,cmd}], cells[{tool,level,skills,knowledge,network}], scenarios, trials, parallel, seed, options") ev = sub.choices["evaluate"] ev.add_argument("--adapter", choices=sorted(ADAPTERS), required=True, help="agent adapter; `direct` is the built-in loop that needs no agent harness") ev.add_argument("--model", required=True, help="BENCH_MODEL, in the form the adapter expects") diff --git a/benchmarks/agent/bench_report.py b/benchmarks/agent/bench_report.py index f37f2105..bfc630c9 100644 --- a/benchmarks/agent/bench_report.py +++ b/benchmarks/agent/bench_report.py @@ -81,6 +81,70 @@ def fmt(value, digits: int = 0) -> str: return "n/a" if value is None else f"{value:,.{digits}f}" +def mean_of(values: list) -> float | None: + values = [v for v in values if v is not None] + return mean(values) if values else None + + +def search_reads(result: dict) -> int: + """Information-gathering operations per run, the same definition at both levels: every wright invocation (CLI calls under + `bin`, `tools/call` under `mcp` — the shim counts both) plus `bash` calls that ran a search/read shell command. A `bash` + call that itself invokes the wright CLI is counted once, on the wright side.""" + import bench_trace + shell = bench_trace.shell_search_reads(result["_dir"] / "transcript.jsonl") + wright = ((result.get("toolUse") or {}).get("wright") or {}).get("invocations") or 0 + return shell + wright + + +def level_stats(runs: list[dict]) -> dict: + """One level's metrics inside a paired group: rates, search/read and tool-call counts, turns, tokens.""" + usage = [(r.get("usage") or {}) for r in runs] + return { + "n": len(runs), + "usable": sum(1 for r in runs if r.get("usable")), + "passed": sum(1 for r in runs if r.get("passed")), + "searchReads": mean_of([search_reads(r) for r in runs]), + "wrightCalls": mean_of([(r.get("toolUse") or {}).get("wright", {}).get("invocations") for r in runs]), + "bashCalls": mean_of([(r.get("toolCalls") or {}).get("bash") for r in runs]), + "toolCalls": mean_of([sum(calls.values()) if calls else None for calls in (r.get("toolCalls") for r in runs)]), + "turns": mean_of([u.get("turns") for u in usage]), + "tokens": mean_of([total_tokens(r) for r in runs]), + } + + +def levels(runs: list[dict]) -> tuple[list[str], list[dict]]: + """`mcp` vs `bin` paired on (scenario, agent, trial) within one cell: rates with Wilson intervals and efficiency (#474).""" + by_key: dict[tuple, dict[str, dict]] = defaultdict(dict) + for r in runs: + condition = r.get("condition") or {} + if condition.get("tool") != "wright": + continue + key = (r["scenario"], r["agent"]["id"], r["_trial"], condition.get("knowledge"), condition.get("network"), tuple(condition.get("skills") or [])) + by_key[key][condition.get("level", "bin")] = r + groups: dict[tuple, list[tuple]] = defaultdict(list) + for (scenario, agent, _trial, knowledge, network, skills), sides in by_key.items(): + if "bin" in sides and "mcp" in sides: + groups[(agent, skills, knowledge, network)].append((sides["bin"], sides["mcp"])) + lines, records = [], [] + for (agent, skills, knowledge, network), pairs in sorted(groups.items()): + cell = f"wright{'+'.join(['', *skills]) if skills else ''}/{knowledge}/{network}" + bins, mcps = [b for b, _ in pairs], [m for _, m in pairs] + stats = {level: level_stats(group) for level, group in (("bin", bins), ("mcp", mcps))} + gain = sum(1 for b, m in pairs if m.get("usable") and not b.get("usable")) + loss = sum(1 for b, m in pairs if b.get("usable") and not m.get("usable")) + both = [(total_tokens(b), total_tokens(m)) for b, m in pairs if b.get("usable") and m.get("usable") and total_tokens(b) and total_tokens(m)] + record = {"agent": agent, "cell": cell, "pairs": len(pairs), "stats": stats, "usableGained": gain, "usableLost": loss, + "tokenSaving": mean(1 - m / b for b, m in both) if both else None, "bothUsable": len(both)} + records.append(record) + for level in ("bin", "mcp"): + s = stats[level] + lines.append(f"| {agent} | {cell} | {level} | {s['n']} | {rate(s['usable'], s['n'])} | {rate(s['passed'], s['n'])} | " + f"{fmt(s['searchReads'], 1)} | {fmt(s['wrightCalls'], 1)} | {fmt(s['bashCalls'], 1)} | {fmt(s['toolCalls'], 1)} | {fmt(s['turns'], 1)} | {fmt(s['tokens'])} |") + saving = f"{record['tokenSaving']:+.0%} tokens (n={len(both)})" if both else "no both-usable pairs" + lines.append(f"| {agent} | {cell} | Δ paired | {len(pairs)} | +{gain} / -{loss} | — | — | — | — | — | — | {saving} |") + return lines, records + + def paired(runs: list[dict], reference: str = BASELINE) -> list[str]: by_key: dict[tuple, dict] = {(r["scenario"], r["agent"]["id"], r["_trial"], label(r)): r for r in runs} lines = [] @@ -200,6 +264,14 @@ def render(results: list[dict], regrade: list[str] | None = None, reference: str g = [r for r in runs if r.get("split") == split and label(r) == cell] if g: out.append(f"| {split} | {cell} | {rate_runs(g)} |") + level_lines, level_records = levels(runs) + if level_lines: + summary["levels"] = level_records + out += ["", "## Level comparison: `mcp` vs `bin` (paired on scenario, agent, trial)", "", + "`search/read` counts information-gathering operations: every wright invocation plus `bash` calls that ran a " + "search/read shell command (a `bash` call invoking the wright CLI counts once, as a wright call).", "", + "| agent | cell | level | runs | usable | passed | search/read | wright calls | bash calls | tool calls | turns | tokens/run |", + "| --- | --- | --- | --- | --- | --- | --- | --- | --- | --- | --- | --- |", *level_lines] pairs = paired(runs, reference) if pairs: out += ["", f"## Paired against `{reference}` (same scenario, agent, trial)", "", "| agent | comparison | pairs | usable gained/lost | tokens where both usable |", "| --- | --- | --- | --- | --- |", *pairs] diff --git a/benchmarks/agent/bench_trace.py b/benchmarks/agent/bench_trace.py index 9a50f8af..31752980 100644 --- a/benchmarks/agent/bench_trace.py +++ b/benchmarks/agent/bench_trace.py @@ -5,6 +5,7 @@ import hashlib import json import os +import re import subprocess import sys import threading @@ -79,13 +80,24 @@ def shim_main(argv: list[str]) -> int: return proc.returncode +def serve_transport(argv: list[str]) -> str: + """The transport a `wright serve` argv selected; stdio is the server default.""" + for i, a in enumerate(argv): + if a == "--transport" and i + 1 < len(argv): + return argv[i + 1] + if a.startswith("--transport="): + return a.split("=", 1)[1] + return "stdio" + + def serve_tee(real: str, argv: list[str], started: float) -> int: proc = subprocess.Popen([real, *argv], stdin=subprocess.PIPE, stdout=subprocess.PIPE) counts = {"req": 0, "res": 0} + session = {"session": os.getpid(), "transport": serve_transport(argv)} # separates this session's lines from a concurrent one def log(direction: str, line: bytes) -> None: counts[direction] += 1 - append_event({"type": "serve", "dir": direction, "t": time.time(), "line": line.decode(errors="replace")}) + append_event({"type": "serve", "dir": direction, "t": time.time(), "line": line.decode(errors="replace"), **session}) def pump() -> None: for line in proc.stdout: @@ -111,6 +123,145 @@ def read_events(path: Path) -> list[dict]: return [json.loads(line) for line in path.read_text().splitlines()] if path.is_file() else [] +def mcp_op(name) -> str | None: + """The Wright operation an MCP tool name carries: `wright_call_graph` -> `callGraph`.""" + if not isinstance(name, str): + return None + return re.sub(r"_([a-z])", lambda m: m.group(1).upper(), name.removeprefix("wright_")) + + +def serve_request(line: str, transport: str = "stdio") -> dict: + """One serve request line: `{op, args, expects}`. + + `expects` mirrors the server's answer rule for the session's transport (`serve.rs`/`mcp.rs`): a blank line is skipped + silently; under jsonrpc/mcp a notification — a JSON-RPC object with a string `method` and no `id` — is answered with + silence; every other line is answered, including unparseable or malformed input (a parse/invalid-request error). + `tools/call` maps to the Wright operation its tool name carries, and the jsonrpc transport's wright methods + (`compile`, `check`, `analyze`, `inspect`) map to their names; every other method keeps a `:` name so + transport traffic (handshake, `tools/list`) is distinguishable from Wright operations.""" + if not line.strip(): + return {"op": None, "args": None, "expects": False} + try: + message = json.loads(line) + except json.JSONDecodeError: + message = None + expects = transport == "stdio" or not ( + isinstance(message, dict) and message.get("jsonrpc") == "2.0" and isinstance(message.get("method"), str) and "id" not in message + ) + if not isinstance(message, dict): + return {"op": None, "args": None, "expects": expects} + if transport == "stdio": # a bare {op, ...} request per line; anything else is answered malformed-request + if "op" in message: + op = message["op"] + return {"op": op if isinstance(op, str) else "", "args": {k: v for k, v in message.items() if k != "op"}, "expects": expects} + return {"op": None, "args": None, "expects": expects} + method = message.get("method") + params = message.get("params") if isinstance(message.get("params"), dict) else {} + if transport == "mcp" and method == "tools/call": + return {"op": mcp_op(params.get("name")), "args": params.get("arguments"), "expects": expects} + if transport == "jsonrpc" and (method == "request" and (op := params.get("op")) or method in DECISION_COMMANDS and (op := method)): + return {"op": op if isinstance(op, str) else None, "args": params or None, "expects": expects} + return {"op": f"{transport}:{method}" if isinstance(method, str) else None, "args": None, "expects": expects} + + +def serve_pairs(events: list[dict]) -> list[tuple[dict, dict, dict | None]]: + """`(request event, parsed request, response event)` for every serve request that expects an answer, in request order. + + Serve lines carry their session's id and transport (the shim tags them); pairing is FIFO within a session because the + server answers its lines strictly in order. A request left without a response pairs with None.""" + sessions: dict = {} + for event in events: + if event.get("type") == "serve": + sessions.setdefault(event.get("session"), []).append(event) + pairs: list[tuple[dict, dict, dict | None]] = [] + for stream in sessions.values(): + pending: list[tuple[dict, dict]] = [] + for event in stream: + if event["dir"] == "req": + request = serve_request(event["line"], event.get("transport", "stdio")) + if request["expects"]: + pending.append((event, request)) + elif pending: + pairs.append((*pending.pop(0), event)) + pairs += [(*pair, None) for pair in pending] + return sorted(pairs, key=lambda pair: pair[0]["t"]) + + +def serve_error(line: str) -> str | None: + """Why a serve response failed: `malformed` when the request could not be understood at all, `refused` for a structured + rejection the server understood (unknown tool, bad params, or a Wright refusal), else None.""" + try: + message = json.loads(line) + except json.JSONDecodeError: + return None # counted by `unparsedServeResponses`, not here + if not isinstance(message, dict): + return None + error = message.get("error") + if isinstance(error, dict): + code = error.get("code") + return "malformed" if code in (-32700, -32600, "malformed-request") else "refused" + result = message.get("result") + if isinstance(result, dict) and result.get("isError") and "content" in result: # an MCP tool result carrying a refusal + return "refused" + return None + + +def serve_use(request: dict) -> bool: + """A serve request counts as a Wright use when it carried a Wright operation; `:` methods are transport traffic.""" + return bool(request["op"]) and ":" not in request["op"] + + +def request_key(request: dict) -> tuple: + """Identity for repeat/retry detection: the operation plus its serialized arguments.""" + return (request["op"], json.dumps(request["args"], sort_keys=True)) + + +def serve_ops(events: list[dict]) -> list[str]: + """The operation names a serve session carried, in order: Wright ops and `:` methods.""" + return [request["op"] for _, request, _ in serve_pairs(events) if request["op"]] + + +def transcript_events(path: Path): + """Parsed transcript events (dicts only); empty when the transcript does not exist.""" + if not path.is_file(): + return + for line in path.read_text().splitlines(): + try: + event = json.loads(line) + except json.JSONDecodeError: + continue + if isinstance(event, dict): + yield event + + +def call_counts(path: Path) -> dict[str, int] | None: + """Model tool calls by name from a normalized adapter transcript (`calls` on its events), None when it has none.""" + counts: dict[str, int] = {} + for event in transcript_events(path): + for call in event.get("calls") or []: + name = call.get("name") if isinstance(call, dict) else None + if name: + counts[name] = counts.get(name, 0) + 1 + return counts or None + + +SEARCH_COMMAND = re.compile(r"(? int: + """`bash` transcript calls that ran a search/read command; a call that invokes the wright CLI is excluded (`toolUse` counts it).""" + count = 0 + for event in transcript_events(path): + for call in event.get("calls") or []: + if not isinstance(call, dict) or call.get("name") != "bash": + continue + command = call.get("input").get("command") if isinstance(call.get("input"), dict) else None + if isinstance(command, str) and SEARCH_COMMAND.search(command) and not WRIGHT_IN_SHELL.search(command): + count += 1 + return count + + class Snapshots(threading.Thread): """Poll watched files during a run; keep a copy of every distinct content, with its time.""" @@ -158,19 +309,36 @@ def summarize_trace(events: list[dict]) -> dict: def summarize_tool(events: list[dict]) -> dict: + """Wright uses: every CLI invocation, plus each serve request's operation (`check`, `lint`, ...). + + A `serve` session is a container — its spawn is not itself a use unless it never carried a request — so a + `wright serve` session and an MCP `tools/call` count the same way. The `:` handshake is not a use + either, but its response bytes (`mcp:tools/list` carries the schemas) count as output.""" calls = [e for e in events if e["type"] == "call"] + cli = [c for c in calls if "requests" not in c] + sessions = [c for c in calls if "requests" in c] + pairs = serve_pairs(events) + uses = [(request, response) for _, request, response in pairs if serve_use(request)] by_command: dict[str, int] = {} for call in calls: by_command[command_of(call["argv"])] = by_command.get(command_of(call["argv"]), 0) + 1 + for _, request, _ in pairs: + if request["op"]: + by_command[request["op"]] = by_command.get(request["op"], 0) + 1 + output: dict[str, int] = {} + for call in calls: + command = command_of(call["argv"]) + if command: + output[command] = output.get(command, 0) + call.get("stdoutBytes", 0) + call.get("stderrBytes", 0) + for _, request, response in pairs: + if request["op"] and response: + output[request["op"]] = output.get(request["op"], 0) + len(response["line"]) return { - "invocations": len(calls), + "invocations": len(cli) + len(uses) + sum(1 for c in sessions if not c["requests"]), "byCommand": by_command, - "failedInvocations": sum(1 for c in calls if c["exit"] != 0), + "failedInvocations": sum(1 for c in calls if c["exit"] != 0) + sum(1 for request, response in uses if response and serve_error(response["line"])), "ownerOrEnvironmentGaps": [c["argv"] for c in calls if c["exit"] >= 3], - "outputTokensEstimate": { - cmd: sum(c.get("stdoutBytes", 0) + c.get("stderrBytes", 0) for c in calls if command_of(c["argv"]) == cmd) // TOKEN_BYTES - for cmd in by_command if cmd - }, + "outputTokensEstimate": {cmd: total // TOKEN_BYTES for cmd, total in output.items()}, } @@ -178,43 +346,35 @@ def friction(events: list[dict]) -> dict: """Wright friction; other tools are summarized by `summarize_trace` only.""" events = tool_events(events, "wright") calls = [e for e in events if e["type"] == "call"] - seen: list[tuple] = [] + pairs = serve_pairs(events) + requests = [(request, response) for _, request, response in pairs if serve_use(request)] + seen: list = [] repeats = 0 - for call in calls: - key = tuple(call["argv"]) + for key in [tuple(c["argv"]) for c in calls] + [request_key(request) for request, _ in requests]: repeats += key in seen seen.append(key) - serve_responses = [] + errors = [serve_error(e["line"]) for e in events if e["type"] == "serve" and e["dir"] == "res"] unparsed = 0 for event in events: if event["type"] == "serve" and event["dir"] == "res": try: - serve_responses.append(json.loads(event["line"])) + json.loads(event["line"]) except json.JSONDecodeError: unparsed += 1 + keys = [request_key(request) for request, _ in requests] return { "usageErrors": sum(1 for c in calls if c["exit"] == 2), "unknownSubcommands": sum(1 for c in calls if "unrecognized subcommand" in c.get("stderrHead", "")), "helpLookups": sum(1 for c in calls if any(a in ("--help", "-h", "help") for a in c["argv"])), - "retriesAfterUnsupported": sum(1 for i, c in enumerate(calls) if c["exit"] >= 3 and tuple(c["argv"]) in [tuple(x["argv"]) for x in calls[i + 1:]]), - "malformedServeRequests": sum(1 for r in serve_responses if r.get("error", {}).get("code") == "malformed-request"), + "retriesAfterUnsupported": sum(1 for i, c in enumerate(calls) if c["exit"] >= 3 and tuple(c["argv"]) in [tuple(x["argv"]) for x in calls[i + 1:]]) + + sum(1 for i, (request, response) in enumerate(requests) if response and serve_error(response["line"]) == "refused" and keys[i] in keys[i + 1:]), + "malformedServeRequests": errors.count("malformed"), "unparsedServeResponses": unparsed, "identicalRepeats": repeats, "callsToFirstSuccess": next((i + 1 for i, c in enumerate(calls) if c["exit"] == 0 and not any(a in ("--help", "-h", "--version") for a in c["argv"])), None), } -def serve_ops(events: list[dict]) -> list[str]: - ops = [] - for e in events: - if e["type"] == "serve" and e["dir"] == "req": - try: - ops.append(json.loads(e["line"]).get("op", "")) - except json.JSONDecodeError: - ops.append("") - return ops - - def expectation(status: str, detail: str = "") -> dict: return {"status": status, "detail": detail} @@ -223,24 +383,27 @@ def detect_expectations(events: list[dict], snapshots: list[dict], scenario: dic """SPEC-414 E01-E12 over the Wright trace. E05, E07, E09, E10 need the agent transcript: `unavailable`.""" events = tool_events(events, "wright") calls = [e for e in events if e["type"] == "call"] + pairs = serve_pairs(events) ops = serve_ops(events) + requests = [request for _, request, _ in pairs if serve_use(request)] + real_ops = [request["op"] for request in requests] used = bool(calls) unavailable = expectation("unavailable", "needs the normalized agent transcript") result = {k: unavailable for k in ("E05", "E07", "E09", "E10")} if not used: return {**result, **{k: expectation("na", "Wright not used") for k in ("E01", "E02", "E03", "E04", "E06", "E08", "E11", "E12")}} first = calls[0] - discovery = any(a in ("--help", "-h", "help", "--version") for a in first["argv"]) or "capabilities" in ops[:1] + discovery = any(a in ("--help", "-h", "help", "--version") for a in first["argv"]) or "capabilities" in ops[:1] or (ops[:1] and ":" in ops[0]) result["E01"] = expectation("pass" if discovery else "fail", f"first call: {' '.join(first['argv'])}") decision = [c for c in calls if command_of(c["argv"]) in DECISION_COMMANDS] structured = [c for c in decision if wants_json(c["argv"])] - if decision or ops: - rate = (len(structured) + len(ops)) / (len(decision) + len(ops)) + if decision or real_ops: + rate = (len(structured) + len(real_ops)) / (len(decision) + len(real_ops)) result["E02"] = expectation("pass" if rate >= 0.5 else "fail", f"structured share {rate:.2f}") else: result["E02"] = expectation("na", "no decision-driving calls") if scenario.get("stabilityRisk"): - stability = any(command_of(c["argv"]) in ("lint", "analyze") for c in calls) or any(o in ("lint", "findings", "analyze") for o in ops) + stability = any(command_of(c["argv"]) in ("lint", "analyze") for c in calls) or any(o in ("lint", "findings", "analyze") for o in real_ops) result["E03"] = expectation("pass" if stability else "fail", "lint or analyze run" if stability else "only check-level validation") else: result["E03"] = expectation("na", "scenario has no stability risk") @@ -249,6 +412,7 @@ def detect_expectations(events: list[dict], snapshots: list[dict], scenario: dic result["E04"] = expectation("na", "no edits observed") else: after = [c for c in calls if command_of(c["argv"]) in VALIDATING and c["t"] + c["seconds"] >= last_edit] + after += [e for e, request in ((e, p) for e, p, _ in pairs) if request["op"] in VALIDATING and e["t"] >= last_edit] matched = [c for c in after if c.get("envelope") and c["envelope"].get("inputIdentity") == final_sha256] result["E04"] = expectation("pass" if after else "fail", f"{len(after)} validation(s) after last edit; {len(matched)} match the final content") withheld = [c for c in calls if ((c.get("envelope") or {}).get("selection") or {}).get("withheld")] @@ -256,10 +420,15 @@ def detect_expectations(events: list[dict], snapshots: list[dict], scenario: dic result["E06"] = expectation("na" if not withheld and not flags else "info", f"{flags} selection-flag call(s), {len(withheld)} withheld result(s)") unsupported = [c for c in calls if c["exit"] >= 3] excess = [c for c in unsupported if sum(1 for x in calls if x["argv"] == c["argv"]) > 2] - result["E08"] = expectation("na" if not unsupported else ("fail" if excess else "pass"), f"{len(unsupported)} exit 3/4 call(s), {len(excess)} retried more than twice") + refused = [i for i, (_, request, response) in enumerate(pairs) if response and serve_error(response["line"]) == "refused"] + keys = [request_key(request) for _, request, _ in pairs] + retry_pairs = sum(1 for i in refused if keys.count(keys[i]) > 2) + result["E08"] = expectation("na" if not unsupported and not refused else ("fail" if excess or retry_pairs else "pass"), + f"{len(unsupported)} exit 3/4 call(s) and {len(refused)} refused request(s), {len(excess) + retry_pairs} retried more than twice") if ops: malformed = friction(events)["malformedServeRequests"] - result["E11"] = expectation("pass" if ops[0] == "capabilities" and not malformed else "fail", f"first op '{ops[0]}', {malformed} malformed") + discovery_op = ops[0] == "capabilities" or ":" in ops[0] # a transport handshake (`mcp:initialize`, ...) is the discovery + result["E11"] = expectation("pass" if discovery_op and not malformed else "fail", f"first op '{ops[0]}', {malformed} malformed") else: result["E11"] = expectation("na", "no serve session") edit_times = [s["t"] for s in snapshots] @@ -269,6 +438,12 @@ def detect_expectations(events: list[dict], snapshots: list[dict], scenario: dic if later["argv"] == c["argv"]: wasted += not any(c["t"] <= t <= later["t"] for t in edit_times) break + request_events = [(e, request_key(request)) for e, request, _ in pairs if serve_use(request)] + for i, (event, key) in enumerate(request_events): + for later, later_key in request_events[i + 1:]: + if later_key == key and not any(event["t"] <= t <= later["t"] for t in edit_times): + wasted += 1 + break result["E12"] = expectation("pass" if wasted == 0 else "fail", f"{wasted} identical repeat(s) with no edit between") return result diff --git a/benchmarks/agent/matrix.example.json b/benchmarks/agent/matrix.example.json index bd7f9da2..4be79d96 100644 --- a/benchmarks/agent/matrix.example.json +++ b/benchmarks/agent/matrix.example.json @@ -26,6 +26,13 @@ "knowledge": "none", "network": "off" }, + { + "tool": "wright", + "level": "mcp", + "skills": [], + "knowledge": "none", + "network": "off" + }, { "tool": "wright", "skills": [ diff --git a/benchmarks/agent/test_adapters.py b/benchmarks/agent/test_adapters.py index b5a5dc5f..0a68a842 100644 --- a/benchmarks/agent/test_adapters.py +++ b/benchmarks/agent/test_adapters.py @@ -170,11 +170,12 @@ def serve(self, replies): class Handler(BaseHTTPRequestHandler): def do_POST(self): - seen.append(json.loads(self.rfile.read(int(self.headers["content-length"])))) + request = json.loads(self.rfile.read(int(self.headers["content-length"]))) + seen.append(request) status, body = replies[min(len(seen), len(replies)) - 1] self.send_response(status) self.end_headers() - self.wfile.write(json.dumps(body).encode()) + self.wfile.write(json.dumps(body(request) if callable(body) else body).encode()) def log_message(self, *args): pass @@ -184,7 +185,7 @@ def log_message(self, *args): self.addCleanup(server.shutdown) return f"http://127.0.0.1:{server.server_port}", seen - def run_direct(self, model, base_var, replies): + def run_direct(self, model, base_var, replies, extra_env=None): base, seen = self.serve(replies) with tempfile.TemporaryDirectory(dir=Path(__file__).parent) as tmp: tmp = Path(tmp) @@ -192,7 +193,7 @@ def run_direct(self, model, base_var, replies): skill.mkdir() (skill / "SKILL.md").write_text("---\nname: demo\ndescription: A demo skill\n---\nbody\n") env = {"BENCH_MODEL": model, "BENCH_KNOWLEDGE": "none", "BENCH_SKILL_DIRS": str(skill), "ANTHROPIC_API_KEY": "k", "OPENAI_API_KEY": "k", base_var: base, - **{f"BENCH_{n}": str(tmp / n.lower()) for n in ("USAGE", "TRANSCRIPT", "CONTEXT", "AGENT_INFO")}, "PATH": os.environ["PATH"]} + **{f"BENCH_{n}": str(tmp / n.lower()) for n in ("USAGE", "TRANSCRIPT", "CONTEXT", "AGENT_INFO")}, "PATH": os.environ["PATH"], **(extra_env or {})} cwd = os.getcwd() os.makedirs(tmp / "work") os.chdir(tmp / "work") @@ -232,6 +233,60 @@ def test_a_malformed_payload_is_a_clean_provider_error_not_a_traceback(self): code, *_ = self.run_direct(model, base_var, replies) self.assertEqual(code, 1) + FAKE_MCP = '''\ +import json, os, sys +log = open(os.environ["MCP_CALLS"], "a") +for line in sys.stdin: + try: + msg = json.loads(line) + except json.JSONDecodeError: + continue + method = msg.get("method") + if method == "initialize": + result = {"protocolVersion": "2025-06-18", "capabilities": {"tools": {}}, "serverInfo": {"name": "fake", "version": "0"}} + elif method == "tools/list": + result = {"tools": [{"name": "wright_check", "description": "Check the project", "inputSchema": {"type": "object", "properties": {"file": {"type": "string"}}}}]} + elif method == "tools/call": + log.write(msg["params"]["name"] + "\\n") + log.flush() + result = {"content": [{"type": "text", "text": "{\\"ok\\": true}"}]} + else: + result = {} + if "id" in msg: + sys.stdout.write(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": result}) + "\\n") + sys.stdout.flush() +''' + + def test_mcp_tools_reach_the_model_dispatch_to_the_server_and_bill_schemas(self): + with tempfile.TemporaryDirectory(dir=Path(__file__).parent) as tmp: + tmp = Path(tmp) + server_py = tmp / "fake_mcp.py" + server_py.write_text(self.FAKE_MCP) + calls = tmp / "calls.txt" + use = {"content": [{"type": "tool_use", "id": "t1", "name": "wright_check", "input": {"file": "mode.ws"}}], + "usage": {"input_tokens": 1, "output_tokens": 1}} + billed = lambda request: {"content": [{"type": "text", "text": "done"}], + "usage": {"input_tokens": len(json.dumps(request["tools"])), "output_tokens": 1}} + code, final, seen, usage, context, info = self.run_direct( + "anthropic/m", "ANTHROPIC_BASE_URL", [(200, use), (200, billed)], + {"BENCH_MCP_CMD": f"{sys.executable} {server_py}", "MCP_CALLS": str(calls), "BENCH_TOOL_LEVEL": "mcp"}) + self.assertEqual((code, final), (0, "done")) + tools = {t["name"]: t for t in seen[0]["tools"]} + self.assertEqual([t["name"] for t in seen[0]["tools"]], ["bash", "wright_check"]) + self.assertEqual(tools["wright_check"]["input_schema"], {"type": "object", "properties": {"file": {"type": "string"}}}) + self.assertEqual(calls.read_text().strip(), "wright_check") + self.assertIn('"ok": true', seen[1]["messages"][-1]["content"][0]["content"]) + self.assertEqual(info["toolLevel"], "mcp") + # fixture: the billed input covers the whole serialized tool list, so the MCP tool schema counts toward mcp input tokens + self.assertEqual(usage[1]["input"], len(json.dumps(seen[1]["tools"]))) + + def test_mcp_setup_failure_is_not_an_infrastructure_exit(self): + with tempfile.TemporaryDirectory(dir=Path(__file__).parent) as tmp: + env = {"BENCH_MODEL": "anthropic/m", "BENCH_KNOWLEDGE": "none", "BENCH_MCP_CMD": "false", "ANTHROPIC_API_KEY": "k", + **{f"BENCH_{n}": str(Path(tmp) / n.lower()) for n in ("USAGE", "TRANSCRIPT", "CONTEXT", "AGENT_INFO")}, "PATH": os.environ["PATH"]} + with patch.dict(direct.os.environ, env, clear=True), patch.object(direct.sys, "stdin", io.StringIO("task")), patch.object(direct.sys, "stdout", io.StringIO()): + self.assertEqual(direct.main(), 1) + class PipeTest(unittest.TestCase): def test_drain_keeps_a_chatty_stderr_from_deadlocking_the_stdout_read(self): diff --git a/benchmarks/agent/test_agent_bench.py b/benchmarks/agent/test_agent_bench.py index 62443ede..8c4390ae 100644 --- a/benchmarks/agent/test_agent_bench.py +++ b/benchmarks/agent/test_agent_bench.py @@ -32,12 +32,12 @@ def setUp(self): self.out = Path(tempfile.mkdtemp(dir=agent_bench.ROOT / "target")).resolve() self.addCleanup(shutil.rmtree, self.out, True) - def trial(self, agent_cmd: str, tool: str = "wright", skills: tuple = (), knowledge: str = "none", scenario: str = SCENARIO, **options) -> dict: + def trial(self, agent_cmd: str, tool: str = "wright", skills: tuple = (), knowledge: str = "none", scenario: str = SCENARIO, level: str = "bin", **options) -> dict: args = argparse.Namespace(**{ "wright": str(Path(WRIGHT).resolve()), "agent_id": "fake", "agent_cmd": agent_cmd, "timeout": 60, "infra_retries": 2, "env_pass": [], "canary_cmd": None, "skill_dirs": {}, "wiki_dir": None, "check_ancestors": False, "out": self.out, **options, }) - cell = agent_bench.normalize_cell({"tool": tool, "skills": list(skills), "knowledge": knowledge, "network": "off"}) + cell = agent_bench.normalize_cell({"tool": tool, "level": level, "skills": list(skills), "knowledge": knowledge, "network": "off"}) return agent_bench.run_trial(agent_bench.load_scenario(scenario), cell, args, self.out / f"{scenario}-{tool}") def test_wiki_snapshot_is_copied_and_identified(self): @@ -89,6 +89,65 @@ def test_cell_labels_applicability_and_baseline(self): self.assertTrue(agent_bench.applicable({"language": "workshop"}, fmt)) self.assertFalse(agent_bench.applicable({"language": "opy"}, fmt)) + def test_mcp_level_label_and_cell_validation(self): + cell = agent_bench.normalize_cell({"tool": "wright", "level": "mcp", "skills": ["wright-skill"], "knowledge": "none", "network": "off"}) + self.assertEqual(agent_bench.cell_label(cell), "wright-mcp+wright-skill/none/off") + self.assertEqual(agent_bench.normalize_cell({"tool": "wright", "skills": [], "knowledge": "none", "network": "off"})["level"], "bin") + self.assertEqual(agent_bench.cell_label(agent_bench.normalize_cell({"tool": "wright", "skills": [], "knowledge": "none", "network": "off"})), "wright/none/off") + args = argparse.Namespace(skill_dirs={}, wiki_dir=None) + with self.assertRaisesRegex(SystemExit, "level 'mcp' is a `wright` level"): + agent_bench.check_cell(agent_bench.normalize_cell({"tool": "overpy", "level": "mcp", "skills": [], "knowledge": "none", "network": "off"}), args) + with self.assertRaisesRegex(SystemExit, "level 'mcp' is a `wright` level"): + agent_bench.check_cell(agent_bench.normalize_cell({"tool": "none", "level": "mcp", "skills": [], "knowledge": "none", "network": "off"}), args) + with self.assertRaisesRegex(SystemExit, "invalid condition"): + agent_bench.check_cell(agent_bench.normalize_cell({"tool": "wright", "level": "bogus", "skills": [], "knowledge": "none", "network": "off"}), args) + + def test_mcp_env_hides_the_cli_and_describes_the_server(self): + workspace = self.out / "ws" + workspace.mkdir() + args = argparse.Namespace(wright=str(Path(WRIGHT).resolve()), agent_id="fake", env_pass=[], skill_dirs={}, wiki_dir=None, check_ancestors=False, canary_cmd=None) + cell = agent_bench.normalize_cell({"tool": "wright", "level": "mcp", "skills": [], "knowledge": "none", "network": "off"}) + env = agent_bench.build_env(cell, args, self.out / "run", workspace) + self.assertEqual(env["BENCH_TOOL_LEVEL"], "mcp") + self.assertIn("serve --transport mcp", env["BENCH_MCP_CMD"]) + self.assertIn("bench_trace.py", env["BENCH_MCP_CMD"]) + self.assertIn(str(workspace), env["BENCH_MCP_CMD"]) + self.assertIsNone(shutil.which("wright", path=env["PATH"])) + self.assertFalse((self.out / "run/bin").exists()) + shims = self.out / "shims" + shims.mkdir() + (shims / "wright").write_text("#!/bin/sh\nexit 0\n") + (shims / "wright").chmod(0o755) + self.assertIn("reachable", agent_bench.canaries(cell, {**env, "PATH": str(shims)}, workspace, args)) + env_bin = agent_bench.build_env({**cell, "level": "bin"}, args, self.out / "run2", workspace) + self.assertTrue((self.out / "run2/bin").is_dir()) + self.assertTrue(str(shutil.which("wright", path=env_bin["PATH"])).startswith(str(self.out / "run2"))) + + def test_mcp_level_is_invalid_when_the_adapter_never_registers_the_server(self): + result = self.trial("true", level="mcp") + self.assertEqual(result["status"], "invalid") + self.assertIn("BENCH_MCP_CMD", result["invalid"]) + + def test_mcp_tools_calls_are_traced_as_wright_uses(self): + lines = [ + '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-06-18","capabilities":{}}}', + '{"jsonrpc":"2.0","method":"notifications/initialized"}', + '{"jsonrpc":"2.0","id":2,"method":"tools/list"}', + '{"jsonrpc":"2.0","id":3,"method":"tools/call","params":{"name":"wright_check","arguments":{}}}', + '{"jsonrpc":"2.0","id":4,"method":"tools/call","params":{"name":"wright_nope","arguments":{}}}', + ] + agent = "printf '%s\\n' " + " ".join(f"'{line}'" for line in lines) + ' | sh -c "$BENCH_MCP_CMD" >/dev/null' + result = self.trial(agent, level="mcp") + use = result["toolUse"]["wright"] + self.assertEqual((use["invocations"], use["byCommand"]["check"]), (2, 1)) + self.assertEqual(use["byCommand"]["nope"], 1) + self.assertEqual(use["failedInvocations"], 1) + self.assertNotIn("invalid", result) + self.assertEqual(result["expectations"]["E11"]["status"], "pass") + self.assertEqual(result["expectations"]["E01"]["status"], "pass") + self.assertEqual(result["expectations"]["E02"]["status"], "pass") + self.assertGreater(use["outputTokensEstimate"]["mcp:tools/list"], 0) + @unittest.skipUnless(bench_grade.oracle_available(), "run `agent_bench.py setup-oracle`") def test_overpy_tool_is_traced_and_wright_is_hidden(self): scenario = "repair-opy-runaway-loop" @@ -670,6 +729,69 @@ def test_friction_reports_unparseable_serve_responses_without_losing_errors(self self.assertEqual(result["malformedServeRequests"], 1) self.assertEqual(result["unparsedServeResponses"], 1) + def serve(self, line, direction="req", transport="stdio", session=7, t=0.0): + return {"tool": "wright", "type": "serve", "dir": direction, "t": t, "line": line, "transport": transport, "session": session} + + def test_serve_request_mirrors_the_servers_answer_rule(self): + notification = '{"jsonrpc":"2.0","method":"notifications/initialized"}' + self.assertFalse(bench_trace.serve_request("", "mcp")["expects"]) # blank lines are skipped silently + self.assertFalse(bench_trace.serve_request(" ", "stdio")["expects"]) + self.assertFalse(bench_trace.serve_request(notification, "mcp")["expects"]) + self.assertTrue(bench_trace.serve_request("not json", "stdio")["expects"]) # answered with a parse error + self.assertTrue(bench_trace.serve_request('{"foo":1}', "mcp")["expects"]) # no method: answered -32600 + self.assertTrue(bench_trace.serve_request(notification, "stdio")["expects"]) # stdio answers every non-blank line + request = bench_trace.serve_request('{"jsonrpc":"2.0","id":3,"method":"tools/call","params":{"name":"wright_call_graph","arguments":{"depth":2}}}', "mcp") + self.assertEqual((request["op"], request["args"], request["expects"]), ("callGraph", {"depth": 2}, True)) + + def test_serve_pairs_do_not_desynchronize_on_silent_lines(self): + res1, res2, res3 = (json.dumps({"jsonrpc": "2.0", "id": i, "result": {}}) for i in (1, 2, 3)) + events = [ + self.serve(""), # skipped silently by the server + self.serve('{"jsonrpc":"2.0","id":1,"method":"tools/call","params":{"name":"wright_check"}}', transport="mcp"), + self.serve('{"jsonrpc":"2.0","method":"notifications/initialized"}', transport="mcp"), # notification: no answer + self.serve('{"foo":1}', transport="mcp"), # answered -32600 + self.serve(res1, "res", "mcp"), self.serve('{"error":{"code":-32600}}', "res", "mcp"), + self.serve('{"jsonrpc":"2.0","id":2,"method":"tools/call","params":{"name":"wright_symbols"}}', transport="mcp"), + self.serve(res2, "res", "mcp"), + self.serve('{"jsonrpc":"2.0","id":3,"method":"tools/call","params":{"name":"wright_lint"}}', transport="mcp"), # unanswered + ] + pairs = bench_trace.serve_pairs(events) + self.assertEqual([r["op"] for _, r, _ in pairs], ["check", None, "symbols", "lint"]) + self.assertEqual(pairs[0][2]["line"], res1) + self.assertEqual(pairs[1][2]["line"], '{"error":{"code":-32600}}') + self.assertEqual(pairs[2][2]["line"], res2) + self.assertIsNone(pairs[3][2]) + + def test_serve_pairs_separate_concurrent_sessions(self): + events = [ + self.serve('{"op":"check"}', session=1, t=1), self.serve('{"op":"lint"}', session=2, t=2), + self.serve('{"r":1}', "res", session=1, t=3), self.serve('{"r":2}', "res", session=2, t=4), + ] + pairs = bench_trace.serve_pairs(events) + self.assertEqual([(r["op"], s["line"]) for _, r, s in pairs], [("check", '{"r":1}'), ("lint", '{"r":2}')]) + + def test_serve_error_classifies_refusals_and_malformed(self): + self.assertEqual(bench_trace.serve_error('{"error":{"code":-32700}}'), "malformed") + self.assertEqual(bench_trace.serve_error('{"error":{"code":"malformed-request"}}'), "malformed") + self.assertEqual(bench_trace.serve_error('{"error":{"code":-32602,"message":"unknown tool"}}'), "refused") + self.assertEqual(bench_trace.serve_error('{"result":{"isError":true,"content":[]}}'), "refused") + self.assertIsNone(bench_trace.serve_error('{"result":{"ok":true}}')) + + def test_call_counts_and_shell_search_reads(self): + (agent_bench.ROOT / "target").mkdir(exist_ok=True) + transcript = Path(tempfile.mkdtemp(dir=agent_bench.ROOT / "target")) / "transcript.jsonl" + self.addCleanup(shutil.rmtree, transcript.parent, True) + transcript.write_text("\n".join(json.dumps(e) for e in [ + {"type": "assistant", "calls": [{"name": "bash", "input": {"command": "rg foo && cat mode.ws"}}, + {"name": "wright_check", "input": {"file": "mode.ws"}}]}, + {"type": "assistant", "calls": [{"name": "bash", "input": {"command": "wright check mode.ws"}}, + {"name": "bash", "input": "a string input"}, + {"name": "bash", "input": {"command": "ls"}}]}, + "not json", + ])) + self.assertEqual(bench_trace.call_counts(transcript), {"bash": 4, "wright_check": 1}) + self.assertEqual(bench_trace.shell_search_reads(transcript), 2) # the wright-CLI call is counted by toolUse, not here + def test_serve_trace_preserves_large_json_responses(self): import io from types import SimpleNamespace @@ -743,6 +865,37 @@ def test_provider_failures_do_not_count_as_agent_failures(self): self.assertEqual(summary["infrastructureFailures"], 1) self.assertIn("excluded from outcome metrics", text) + def level_result(self, level, trial, usable, tokens, bash=(), wright_invocations=0, tool_calls=None, scenario="s"): + run = self.result("wright-mcp/none/off" if level == "mcp" else "wright/none/off", trial, usable, tokens, scenario=scenario) + run["condition"].update(tool="wright", level=level, knowledge="none", network="off", skills=[]) + run["toolUse"] = {"wright": {"invocations": wright_invocations}} + run["toolCalls"] = tool_calls or {} + (agent_bench.ROOT / "target").mkdir(exist_ok=True) + directory = Path(tempfile.mkdtemp(dir=agent_bench.ROOT / "target")) + self.addCleanup(shutil.rmtree, directory, True) + transcript = [{"type": "assistant", "calls": [{"name": "bash", "input": {"command": command}} for command in bash]}] + (directory / "transcript.jsonl").write_text("\n".join(json.dumps(e) for e in transcript)) + run["_dir"] = directory + return run + + def test_level_comparison_pairs_mcp_with_bin(self): + runs = [ + self.level_result("bin", 1, True, 1000, bash=("rg foo mode.ws", "wright check mode.ws", "cat mode.ws"), wright_invocations=2, tool_calls={"bash": 3}), + self.level_result("mcp", 1, True, 800, wright_invocations=3, tool_calls={"wright_check": 2, "wright_symbols": 1, "bash": 1}), + self.level_result("bin", 2, False, 900, bash=("ls",), wright_invocations=1, tool_calls={"bash": 1}), + self.level_result("mcp", 2, True, 700, wright_invocations=2, tool_calls={"wright_check": 2}), + ] + text, summary = bench_report.render(runs) + self.assertIn("Level comparison", text) + self.assertIn("+1 / -0", text) # mcp turned trial 2 usable + stats = {row["cell"]: row["stats"] for row in summary["levels"]} + self.assertEqual(stats["wright/none/off"]["mcp"]["searchReads"], 2.5) # wright invocations only; no shell search commands + self.assertEqual(stats["wright/none/off"]["bin"]["searchReads"], 3.0) # (2 wright + 2 shell) and (1 wright + 1 shell) + self.assertIn("[", text.split("Level comparison")[1]) # Wilson intervals on the rates + # trial keys do not cross-pair: bin trial 1 and mcp trial 2 are different trials + text, _ = bench_report.render([runs[0], runs[3]]) + self.assertNotIn("Level comparison", text) + if __name__ == "__main__": unittest.main() diff --git a/docs/agent-benchmark.md b/docs/agent-benchmark.md index ac4c410b..75e7a067 100644 --- a/docs/agent-benchmark.md +++ b/docs/agent-benchmark.md @@ -22,8 +22,9 @@ agent actually had (`Agent setup`). - The scenario workspace: the seed project only. - The scenario prompt, delivered on the agent's stdin. It states the requirement in user terms and never names Wright commands or Workshop APIs. -- In the `wright` tool condition, the released `wright` CLI and `wright serve` - session on `PATH`; in the `overpy` tool condition, the `overpy` compiler. +- In the `wright` tool condition, the released `wright` CLI on `PATH` (level + `bin`) or its MCP tools natively (level `mcp`), always through the tracing + shim; in the `overpy` tool condition, the `overpy` compiler. Not allowed in the primary condition: a Workshop/OverPy/OSTW system prompt or skill pack, a generated API reference, or task-specific hints. The agent, @@ -32,12 +33,14 @@ on a vendor. Other conditions below are experiments and are labeled as such. ## Conditions -A cell is `tool[+skill...]/knowledge/network`, for example `wright+wright-skill/none/off`; -the baseline is `none/none/off`. The task and prompt are identical in every cell. +A cell is `tool[-level][+skill...]/knowledge/network`, for example +`wright+wright-skill/none/off` or `wright-mcp/none/off`; the baseline is +`none/none/off`. The task and prompt are identical in every cell. | Factor | Levels | | --- | --- | -| `tool` | `none`: neither tool is reachable. `wright`: `wright` on `PATH` through a tracing shim. `overpy`: the pinned `overpy` compiler through a tracing shim (OverPy scenarios only). A tool that is not part of the condition is hidden from `PATH`, and a canary fails the run if it is reachable. | +| `tool` | `none`: neither tool is reachable. `wright`: Wright through a tracing shim; how it reaches the agent is set by `level`. `overpy`: the pinned `overpy` compiler through a tracing shim (OverPy scenarios only). A tool that is not part of the condition is hidden from `PATH`, and a canary fails the run if it is reachable. | +| `level` | `bin` (default): the `wright` CLI on `PATH` through the shim, so the agent uses Wright as shell + JSON. `mcp` (`wright` only): the `wright` CLI stays off `PATH` and the harness gives the adapter `BENCH_MCP_CMD`, the command that starts `wright serve --transport mcp` on the workspace through the same tracing shim; the adapter registers the server and exposes its tools natively. A run at level `mcp` where the adapter never starts the server is `invalid`. | | `skills` | Any of `wright-skill` (how to use Wright), `workshop-skill` (the progressive wiki knowledge skill, see below), `opy-skill` (how to write OverPy and use `overpy`; OverPy scenarios only), and `workshop-format-skill` (the raw Workshop source format; Workshop scenarios only, a local benchmark control). Each is given by `--skill-dir NAME=DIR` and installed through the agent's skill mechanism. | | `knowledge` | `none`; `wiki`: a pinned snapshot given by `--wiki-dir`, copied into the workspace as `./wiki` (a real copy, because tools such as `rg` do not follow symlinks; never counted as an edit); `web`: the adapter enables its web tools. | | `network` | `off` or `on`; `web` requires `on`. | @@ -90,7 +93,10 @@ elsewhere on disk, so run `none` in a clean environment when that matters. `--agent-cmd` is a shell command run in the workspace with the prompt on stdin. The harness describes the cell through environment variables, and the adapter -enforces it: `BENCH_TOOL`, `BENCH_SKILLS` (names), `BENCH_SKILL_DIRS` (one +enforces it: `BENCH_TOOL`, `BENCH_TOOL_LEVEL` (`bin`, or `mcp` when the cell is a +`wright` MCP cell), `BENCH_MCP_CMD` (level `mcp` only: the command that starts the +traced `wright serve --transport mcp` server for the adapter to register), +`BENCH_SKILLS` (names), `BENCH_SKILL_DIRS` (one directory per installed skill, separated by the path separator), `BENCH_KNOWLEDGE`, `BENCH_NETWORK`, `BENCH_HOST_PATH` (the unscrubbed `PATH`, for locating the agent binary itself; do not pass it to the agent), and @@ -129,7 +135,7 @@ keep host configuration out of the run, and it reports what loaded through | [`agy.py`](../benchmarks/agent/adapters/agy.py) | Antigravity CLI | `BENCH_MODEL` (for example `gemini-3.8-flash-high`) and `BENCH_THINKING`. Isolated `HOME` with only authentication files, workspace skills under `.agents/skills`, MCP/browser access denied and URL reads denied outside `web`. Observed built-in web tools under network `off` stop and invalidate the trial; URL permissions do not cover search. Streaming step usage is recorded. The CLI does not export observed loaded skills or context limits; these remain unreported. | | [`opencode.py`](../benchmarks/agent/adapters/opencode.py) | opencode | `BENCH_MODEL=provider/model` (as `opencode models` lists it) and `BENCH_THINKING` as the model variant. Isolated `HOME` and XDG directories holding only the credentials; `--pure`; Claude Code instructions and the real home's external skills are disabled by environment variable because opencode reads them regardless of `HOME`; skills installed under `.opencode/skills`; web tools denied outside `web`. Per-step usage comes from `step_finish` events; the tool list is what the agent used. | | [`grok.py`](../benchmarks/agent/adapters/grok.py) | Grok CLI | `BENCH_MODEL` (a `grok models` id) and `BENCH_THINKING` as reasoning effort. Isolated `GROK_HOME` holding only the login and the skills; prompt sent `--verbatim`; subagents disabled; web search disabled outside `web`. Tools, skills, and the context window come from the stream's init and result lines. | -| [`direct.py`](../benchmarks/agent/adapters/direct.py) | none (built-in loop) | See below. | +| [`direct.py`](../benchmarks/agent/adapters/direct.py) | none (built-in loop) | See below. Under level `mcp` it registers `BENCH_MCP_CMD` and exposes the server's tools to the model; only `direct` supports level `mcp` today. | The pi adapter also uses an isolated `HOME` containing only its authentication files. Explicit provider extensions remain referenced by path, not copied with @@ -272,16 +278,18 @@ with the workspace, `agent.log`, snapshots, and the Wright trace beside it. | `authorities`, `disagreement` | Validity per authority and any Wright/oracle disagreement | | `diagnostics`, `lintFindings`, `unsafeEdits` | Remaining `wright check` diagnostics, lint rule codes, files changed outside `writable` | | `unverifiedRuntimeClaims` | The scenario's `runtimeOnly` claims | -| `toolUse` | Per tool (`wright`, `overpy`): invocations by subcommand, failures, exits of 3 or 4 (candidate owner or environment gaps), and estimated output tokens per command | +| `toolUse` | Per tool (`wright`, `overpy`): invocations by subcommand, failures, exits of 3 or 4 (candidate owner or environment gaps), and estimated output tokens per command. Under level `mcp`, each `tools/call` counts as an invocation of the Wright operation its tool name carries; the `initialize`/`tools/list` handshake is not an invocation but its response bytes (the tool schemas) are included in `mcp:tools/list`'s estimated output tokens | +| `toolCalls` | Model-visible tool calls by name from the adapter's normalized transcript (`bash`, `fetch`, `wright_*` under `mcp`), when the adapter writes one | | `friction`, `expectations` | Usage errors, unknown subcommands, help lookups, retries, malformed `serve` requests, unparsed `serve` responses, identical repeats; expectation E01-E12 verdicts | | `snapshots` | Strict validity of each snapshot of the entry, first valid index, and valid-to-invalid regressions | | `usage`, `context` | Turns, tokens by kind, peak context (and its share of the limit), tokens to first valid; loaded context | | `invalid`, `infraRetries`, `fileReadEnforcement`, `fileWriteEnforcement`, `networkEnforcement` | Present when the run was excluded or retried; how file reads (`allow-list` with the hidden and allowed paths, or `unrestricted`), file writes (`trial-directory-only` or `unrestricted`), and network `off` (`canary-checked` or `declared-only`) were enforced | -`toolUse` is recorded per CLI invocation through the shim; `wright serve` sessions are teed -line by line into the trace. Comparing cells for the same scenario shows what -Wright adds. Exit 3 or 4 entries, disagreements, and failed `layer` values are -the input for owner Issues. +`toolUse` is recorded per CLI invocation through the shim; `wright serve` sessions — +stdio or MCP — are teed line by line into the trace, and each `tools/call` is +counted like one CLI invocation of the same operation. Comparing cells for the +same scenario shows what Wright adds. Exit 3 or 4 entries, disagreements, and +failed `layer` values are the input for owner Issues. Expectations E05, E07, E09, and E10 need the normalized agent transcript and report `unavailable` until an adapter provides it. Estimated tokens for Wright @@ -298,6 +306,17 @@ Diagnostics flag headroom (baseline usable rate of at least 95%), infrastructure failures, invalid runs, trial variance, and, with `--regrade`, a grader that gives different verdicts on the same stored workspace. +When a directory holds both `wright` and `wright-mcp` cells, a level-comparison +section pairs `mcp` with `bin` runs of the same scenario, agent, and trial and +reports each side's usable and passed rates with Wilson intervals, mean +search/read operations (every Wright invocation plus `bash` calls that ran a +search/read shell command), Wright calls, `bash` calls, total model tool calls, +turns, and tokens per run, plus the paired usable gained/lost and token delta. +The MCP tool schemas are part of the model request, so their context cost is +inside the provider-reported input tokens the usage accounting sums; the +`mcp:tools/list` estimate under `toolUse` shows the same payload on the Wright +side. + ## Score `agent_bench.py score` writes `score.json` and `score.txt`: one card per language From 25c55a297da4093473093b862ecee714cbff9032 Mon Sep 17 00:00:00 2001 From: Teakowa <27560638+Teakowa@users.noreply.github.com> Date: Sun, 4 Oct 2026 11:17:26 +0800 Subject: [PATCH 2/4] fix(bench): keep buffered MCP reply lines for the next request Mcp._lines polled select() on the pipe FD but read with TextIOWrapper.readline(), which buffers ahead past what select reports: a line behind the matching reply sat in the wrapper's buffer while the next request waited on an empty FD and timed out. Raw os.read into self.pending keeps every complete line either yielded or pending across requests; a line at the buffer boundary and a final line without a newline are handled the same way. --- benchmarks/agent/adapters/direct.py | 24 +++++++++++++++++++----- benchmarks/agent/test_adapters.py | 26 ++++++++++++++++++++++++++ 2 files changed, 45 insertions(+), 5 deletions(-) diff --git a/benchmarks/agent/adapters/direct.py b/benchmarks/agent/adapters/direct.py index a0f75fb8..3262c101 100644 --- a/benchmarks/agent/adapters/direct.py +++ b/benchmarks/agent/adapters/direct.py @@ -132,6 +132,7 @@ class Mcp: def __init__(self, command: str, env: dict): self.proc = subprocess.Popen(["/bin/sh", "-c", command], stdin=subprocess.PIPE, stdout=subprocess.PIPE, text=True, env=env) self.next_id = 0 + self.pending = b"" # bytes read past the last yielded line; keeps complete lines for the next request def request(self, method: str, params: dict | None = None) -> dict: self.next_id += 1 @@ -150,7 +151,10 @@ def request(self, method: str, params: dict | None = None) -> dict: return {"error": {"code": -32000, "message": "the MCP server closed or timed out"}} def _lines(self): - """Reply lines, bounded to COMMAND_SECONDS where the OS can poll a pipe; a hung server then reads as a tool error.""" + """Reply lines, bounded to COMMAND_SECONDS where the OS can poll a pipe; a hung server then reads as a tool error. + + readline() would buffer ahead past what select() reports, stranding a line the next request then waits for; + raw reads keep every complete line either yielded or in self.pending.""" if os.name != "posix": yield from self.proc.stdout return @@ -158,11 +162,21 @@ def _lines(self): selector = selectors.DefaultSelector() with selector: selector.register(self.proc.stdout, selectors.EVENT_READ) - while selector.select(max(0.0, deadline - time.monotonic())): - line = self.proc.stdout.readline() - if not line: + while True: + head, sep, rest = self.pending.partition(b"\n") + if sep: + self.pending = rest + yield head.decode("utf-8", "replace") + continue + if not selector.select(max(0.0, deadline - time.monotonic())): + return + chunk = os.read(self.proc.stdout.fileno(), 65536) + if not chunk: + if self.pending: + yield self.pending.decode("utf-8", "replace") + self.pending = b"" return - yield line + self.pending += chunk def notify(self, method: str) -> None: try: diff --git a/benchmarks/agent/test_adapters.py b/benchmarks/agent/test_adapters.py index 0a68a842..d999133d 100644 --- a/benchmarks/agent/test_adapters.py +++ b/benchmarks/agent/test_adapters.py @@ -287,6 +287,32 @@ def test_mcp_setup_failure_is_not_an_infrastructure_exit(self): with patch.dict(direct.os.environ, env, clear=True), patch.object(direct.sys, "stdin", io.StringIO("task")), patch.object(direct.sys, "stdout", io.StringIO()): self.assertEqual(direct.main(), 1) + NOISY_MCP = '''\ +import json, sys +for line in sys.stdin: + try: + msg = json.loads(line) + except json.JSONDecodeError: + continue + if "id" in msg: + sys.stdout.write(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}) + "\\n") + sys.stdout.write(json.dumps({"note": "trailing"}) + "\\n") + sys.stdout.flush() +''' + + @unittest.skipUnless(os.name == "posix", "the reply poll needs a selectable pipe") + def test_mcp_lines_behind_the_reply_are_not_lost_to_the_next_request(self): + with tempfile.TemporaryDirectory(dir=Path(__file__).parent) as tmp: + server_py = Path(tmp) / "noisy_mcp.py" + server_py.write_text(self.NOISY_MCP) + mcp = direct.Mcp(f"{sys.executable} {server_py}", dict(os.environ)) + try: + self.assertEqual(mcp.request("initialize")["id"], 1) # the buffered trailing line must not strand request 2 + self.assertEqual(mcp.request("tools/list")["id"], 2) + self.assertEqual(mcp.request("tools/call", {"name": "x", "arguments": {}})["id"], 3) + finally: + mcp.close() + class PipeTest(unittest.TestCase): def test_drain_keeps_a_chatty_stderr_from_deadlocking_the_stdout_read(self): From dbf02e1b5b1f0153e994d14547284dfc0e349f0a Mon Sep 17 00:00:00 2001 From: Teakowa <27560638+Teakowa@users.noreply.github.com> Date: Sun, 4 Oct 2026 12:42:04 +0800 Subject: [PATCH 3/4] fix(bench): map the jsonrpc transport's real methods and log requests before forwarding serve_request reused DECISION_COMMANDS for the jsonrpc method map, which has lint and lacks compile: a real compile call classified as transport traffic (invisible to toolUse and friction), while a lint the server can only answer MethodNotFound counted as a Wright use. serve.rs accepts compile/check/analyze/inspect, so those get their own constant. serve_tee also logged the req event after writing to the server, letting a fast response reach the trace first and desynchronizing serve_pairs; it now logs before the write. --- benchmarks/agent/bench_trace.py | 5 +++-- benchmarks/agent/test_agent_bench.py | 3 +++ 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/benchmarks/agent/bench_trace.py b/benchmarks/agent/bench_trace.py index 31752980..41e2df3e 100644 --- a/benchmarks/agent/bench_trace.py +++ b/benchmarks/agent/bench_trace.py @@ -14,6 +14,7 @@ TOKEN_BYTES = 4 # estimation only: bytes per token for Wright output attribution DECISION_COMMANDS = ("check", "lint", "analyze", "inspect") +JSONRPC_METHODS = ("compile", "check", "analyze", "inspect") # serve.rs's direct methods — `lint` exists only as a CLI command VALIDATING = ("check", "lint", "analyze", "compile") OUTPUT_FORMAT_FLAGS = ("--format", "-f") @@ -108,9 +109,9 @@ def pump() -> None: reader = threading.Thread(target=pump) reader.start() for line in sys.stdin.buffer: + log("req", line) # before the write: the pump can log a fast response before a req logged after the write proc.stdin.write(line) proc.stdin.flush() - log("req", line) proc.stdin.close() code = proc.wait() reader.join() @@ -159,7 +160,7 @@ def serve_request(line: str, transport: str = "stdio") -> dict: params = message.get("params") if isinstance(message.get("params"), dict) else {} if transport == "mcp" and method == "tools/call": return {"op": mcp_op(params.get("name")), "args": params.get("arguments"), "expects": expects} - if transport == "jsonrpc" and (method == "request" and (op := params.get("op")) or method in DECISION_COMMANDS and (op := method)): + if transport == "jsonrpc" and (method == "request" and (op := params.get("op")) or method in JSONRPC_METHODS and (op := method)): return {"op": op if isinstance(op, str) else None, "args": params or None, "expects": expects} return {"op": f"{transport}:{method}" if isinstance(method, str) else None, "args": None, "expects": expects} diff --git a/benchmarks/agent/test_agent_bench.py b/benchmarks/agent/test_agent_bench.py index 8c4390ae..69e157d9 100644 --- a/benchmarks/agent/test_agent_bench.py +++ b/benchmarks/agent/test_agent_bench.py @@ -742,6 +742,9 @@ def test_serve_request_mirrors_the_servers_answer_rule(self): self.assertTrue(bench_trace.serve_request(notification, "stdio")["expects"]) # stdio answers every non-blank line request = bench_trace.serve_request('{"jsonrpc":"2.0","id":3,"method":"tools/call","params":{"name":"wright_call_graph","arguments":{"depth":2}}}', "mcp") self.assertEqual((request["op"], request["args"], request["expects"]), ("callGraph", {"depth": 2}, True)) + # the jsonrpc transport's wright methods are compile/check/analyze/inspect — lint is a CLI op the server rejects + self.assertEqual(bench_trace.serve_request('{"jsonrpc":"2.0","id":4,"method":"compile"}', "jsonrpc")["op"], "compile") + self.assertEqual(bench_trace.serve_request('{"jsonrpc":"2.0","id":5,"method":"lint"}', "jsonrpc")["op"], "jsonrpc:lint") def test_serve_pairs_do_not_desynchronize_on_silent_lines(self): res1, res2, res3 = (json.dumps({"jsonrpc": "2.0", "id": i, "result": {}}) for i in (1, 2, 3)) From cc365dd3164b7da5c62fc4ea014f8326783c906d Mon Sep 17 00:00:00 2001 From: Teakowa <27560638+Teakowa@users.noreply.github.com> Date: Sun, 4 Oct 2026 14:11:47 +0800 Subject: [PATCH 4/4] fix(bench): tolerate partial artifacts everywhere a resume reads them MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit bench_report.load() also parsed result.json unguarded, so a truncated file left by a killed run still crashed evaluate's report step for trials the new run never attempted. The atomic writer moves to bench_report — the one module every artifact reader already imports — and now also covers score.json, leaderboard.json, and summary.json. compare() and load_entries() treat a partial score.json like a run whose scoring never finished instead of crashing on JSONDecodeError. --- benchmarks/agent/agent_bench.py | 11 ++--------- benchmarks/agent/bench_leaderboard.py | 7 +++++-- benchmarks/agent/bench_report.py | 15 +++++++++++++-- benchmarks/agent/bench_score.py | 11 ++++++++--- benchmarks/agent/test_agent_bench.py | 14 ++++++++++++++ 5 files changed, 42 insertions(+), 16 deletions(-) diff --git a/benchmarks/agent/agent_bench.py b/benchmarks/agent/agent_bench.py index 29425c9d..feba1920 100644 --- a/benchmarks/agent/agent_bench.py +++ b/benchmarks/agent/agent_bench.py @@ -410,7 +410,7 @@ def run_trial(scenario: dict, cell: dict, args: argparse.Namespace, out: Path) - if reason: result = base_result(scenario, cell, args, out, 0.0, None) result.update(invalid=reason, status="invalid") - write_json(out / "result.json", result) + bench_report.write_json(out / "result.json", result) return result snapshots = bench_trace.Snapshots(workspace, scenario.get("watch", [scenario["entry"]]), out / "snapshots") snapshots.start() @@ -460,7 +460,7 @@ def run_trial(scenario: dict, cell: dict, args: argparse.Namespace, out: Path) - first_valid = next((s["t"] for s in result["snapshots"]["series"] if s["valid"]), None) result["usage"] = bench_trace.usage_summary(out / "usage.jsonl", first_valid) result["status"] = run_status(result) - write_json(out / "result.json", result) + bench_report.write_json(out / "result.json", result) return result @@ -497,13 +497,6 @@ def file_sha256(path: str | Path) -> str: return hashlib.sha256(Path(path).read_bytes()).hexdigest() -def write_json(path: Path, data: dict) -> None: - """A kill mid-write leaves no partial JSON to crash the next resume: temp file in the same directory, then replace.""" - temporary = path.with_name(f".{path.name}.tmp") - temporary.write_text(json.dumps(data, indent=2) + "\n") - os.replace(temporary, path) - - def base_result(scenario: dict, cell: dict, args: argparse.Namespace, out: Path, seconds: float, agent_exit: int | None) -> dict: return { "contract": RESULT_CONTRACT, diff --git a/benchmarks/agent/bench_leaderboard.py b/benchmarks/agent/bench_leaderboard.py index 88681b72..f25c26ce 100644 --- a/benchmarks/agent/bench_leaderboard.py +++ b/benchmarks/agent/bench_leaderboard.py @@ -50,7 +50,10 @@ def load_entries(dirs: list[Path]) -> list[dict]: path = directory / "score.json" if not path.is_file(): continue - cards = {c["language"]: c for c in json.loads(path.read_text())["cards"] if "refused" not in c} + try: + cards = {c["language"]: c for c in json.loads(path.read_text())["cards"] if "refused" not in c} + except json.JSONDecodeError: + continue # a partial file means scoring never finished; the run has no entry yet if not cards: continue first = next(iter(cards.values())) @@ -183,7 +186,7 @@ def main(dirs: list[Path], out: Path) -> int: out.mkdir(parents=True, exist_ok=True) (out / "LEADERBOARD.md").write_text(markdown(data)) (out / "leaderboard.html").write_text(page(data)) - (out / "leaderboard.json").write_text(json.dumps(data, indent=2) + "\n") + bench_report.write_json(out / "leaderboard.json", data) print(markdown(data)) print(f"wrote {out / 'LEADERBOARD.md'}, {out / 'leaderboard.html'}, {out / 'leaderboard.json'}") return 0 if data["entries"] else 1 diff --git a/benchmarks/agent/bench_report.py b/benchmarks/agent/bench_report.py index bfc630c9..998e4a2e 100644 --- a/benchmarks/agent/bench_report.py +++ b/benchmarks/agent/bench_report.py @@ -4,6 +4,7 @@ import json import math +import os import re from collections import defaultdict from pathlib import Path @@ -13,6 +14,13 @@ HEADROOM = 0.95 +def write_json(path: Path, data: dict) -> None: + """A kill mid-write leaves no partial JSON to crash the next resume: temp file in the same directory, then replace.""" + temporary = path.with_name(f".{path.name}.tmp") + temporary.write_text(json.dumps(data, indent=2) + "\n") + os.replace(temporary, path) + + def wilson(k: int, n: int, z: float = 1.96) -> tuple[float, float]: if n == 0: return (0.0, 0.0) @@ -31,7 +39,10 @@ def load(dirs: list[Path]) -> list[dict]: results = [] for base in dirs: for path in sorted(base.rglob("result.json")): - result = json.loads(path.read_text()) + try: + result = json.loads(path.read_text()) + except json.JSONDecodeError: + continue # a partial file means the trial never finished; it will be retried if str(result.get("contract", "")).startswith("wright-agent-bench/") and all(k in result for k in ("status", "language", "condition", "scenario", "agent", "environment")): result["_dir"] = path.parent match = re.search(r"-(\d+)$", path.parent.name) @@ -321,6 +332,6 @@ def main(dirs: list[Path], wright: str, regrade: bool, load_scenario, reference: notes = regrade_notes([r for r in results if r["status"] != "invalid"], wright, load_scenario) if regrade else None text, summary = render(results, notes, reference) (dirs[0] / "report.md").write_text(text) - (dirs[0] / "summary.json").write_text(json.dumps(summary, indent=2) + "\n") + write_json(dirs[0] / "summary.json", summary) print(text) return 0 diff --git a/benchmarks/agent/bench_score.py b/benchmarks/agent/bench_score.py index aa10a2f7..733f3497 100644 --- a/benchmarks/agent/bench_score.py +++ b/benchmarks/agent/bench_score.py @@ -158,7 +158,12 @@ def compare(dirs: list[Path]) -> str: if not path.is_file(): rows.append(("", directory.name, "no score", "no score.json in this directory")) continue - for card_ in json.loads(path.read_text())["cards"]: + try: + cards_ = json.loads(path.read_text())["cards"] + except json.JSONDecodeError: + rows.append(("", directory.name, "no score", "score.json is incomplete — rerun report")) + continue + for card_ in cards_: if "refused" in card_: rows.append((card_["track"], directory.name, "no score", card_["refused"])) continue @@ -176,7 +181,7 @@ def compare(dirs: list[Path]) -> str: def main(dirs: list[Path], languages: list[str], expected_by_language: dict[str, list[str]], out: Path | None) -> int: - from bench_report import load + from bench_report import load, write_json results = load(dirs) status = 0 cards = [] @@ -186,6 +191,6 @@ def main(dirs: list[Path], languages: list[str], expected_by_language: dict[str, print(render(c)) status = max(status, 2 if "refused" in c else 0) target = out or dirs[0] - (target / "score.json").write_text(json.dumps({"contract": CONTRACT, "cards": cards}, indent=2) + "\n") + write_json(target / "score.json", {"contract": CONTRACT, "cards": cards}) (target / "score.txt").write_text("\n".join(render(c) for c in cards)) return status diff --git a/benchmarks/agent/test_agent_bench.py b/benchmarks/agent/test_agent_bench.py index 69e157d9..9082629e 100644 --- a/benchmarks/agent/test_agent_bench.py +++ b/benchmarks/agent/test_agent_bench.py @@ -818,6 +818,20 @@ def test_unused_wright_is_not_applicable(self): class ReportTest(unittest.TestCase): + def test_load_skips_a_partial_result_left_by_a_killed_run(self): + import tempfile + (agent_bench.ROOT / "target").mkdir(exist_ok=True) + root = Path(tempfile.mkdtemp(dir=agent_bench.ROOT / "target")).resolve() + self.addCleanup(shutil.rmtree, root, True) + finished = {k: v for k, v in {**self.result("wright/none/off", 2, True, 5), "environment": {}}.items() if not k.startswith("_")} # _dir/_trial are stamped on load, not stored + for name, content in (("partial-1", '{"contract":'), ("done-2", json.dumps(finished))): + directory = root / "s" / "a" / name + directory.mkdir(parents=True) + (directory / "result.json").write_text(content) + loaded = bench_report.load([root]) + self.assertEqual(len(loaded), 1) + self.assertEqual(loaded[0]["_trial"], 2) + def result(self, cell, trial, usable, tokens, split=None, scenario="s", language="opy", status="completed"): tool = cell.split("/")[0].split("+")[0] return {