From 513a85dd2ad695f02d60e926eb761b123d868f66 Mon Sep 17 00:00:00 2001 From: Xiwei Pan Date: Mon, 20 Jul 2026 04:31:00 +0800 Subject: [PATCH] feat: run self-selected Top50 episodes --- benchmark/agent_environment.py | 13 +- benchmark/process_control.py | 5 + benchmark/run_top50.py | 102 +++++ benchmark/tests/test_run_top50.py | 36 ++ benchmark/tests/test_top50_runner.py | 223 +++++++++++ benchmark/top50_config.yaml | 38 ++ benchmark/top50_runner.py | 556 +++++++++++++++++++++++++++ docker/Dockerfile | 4 +- 8 files changed, 974 insertions(+), 3 deletions(-) create mode 100644 benchmark/run_top50.py create mode 100644 benchmark/tests/test_run_top50.py create mode 100644 benchmark/tests/test_top50_runner.py create mode 100644 benchmark/top50_config.yaml create mode 100644 benchmark/top50_runner.py diff --git a/benchmark/agent_environment.py b/benchmark/agent_environment.py index c2ce4f5..a69354c 100644 --- a/benchmark/agent_environment.py +++ b/benchmark/agent_environment.py @@ -6,6 +6,17 @@ from benchmark.process_control import ProcessLimits, run_capped_process +_SAFE_ENV_KEYS = {"PATH", "PYTHONPATH", "LANG", "LC_ALL", "PAGER", "MANPAGER", "LESS", + "PRB_PRED_DIR", "PRB_SUBMIT_DIR", "PRB_ARTIFACT_DIR", + "PRB_PRED_TIMEOUT", "PRB_SUBMIT_TIMEOUT"} + + +def sanitized_agent_env(extra: dict[str, str] | None = None) -> dict[str, str]: + """Return the small non-secret environment visible to model-authored commands.""" + merged = dict(os.environ) + merged.update(extra or {}) + return {key: value for key, value in merged.items() if key in _SAFE_ENV_KEYS} + def run_as_agent( command: str, @@ -64,7 +75,7 @@ def execute(self, action: dict, cwd: str = "", *, timeout: int | None = None) -> result = run_as_agent( command, cwd=action_cwd, - env=os.environ | self.config.env, + env=sanitized_agent_env(self.config.env), timeout=timeout or self.config.timeout, uid=uid, gid=gid, diff --git a/benchmark/process_control.py b/benchmark/process_control.py index 76b55d9..3036bd2 100644 --- a/benchmark/process_control.py +++ b/benchmark/process_control.py @@ -132,6 +132,11 @@ def run_capped_process( timed_out = True terminate_process_group(process) process.wait() + else: + # A successful shell leader may leave background children behind. Kill the + # entire private process group before reading pipes to completion so no process + # can survive into a later rule episode or keep a pipe open indefinitely. + terminate_process_group(process) for reader in readers: reader.join() stdout, stderr = capture.result() diff --git a/benchmark/run_top50.py b/benchmark/run_top50.py new file mode 100644 index 0000000..1634334 --- /dev/null +++ b/benchmark/run_top50.py @@ -0,0 +1,102 @@ +#!/usr/bin/env python3 +"""Production entrypoint for the standardized self-selected Top50 model track.""" +from __future__ import annotations + +import argparse +import json +import os +from pathlib import Path + +from benchmark.env_setup import find_pred_binary, pinned_commit, verify_pred_version +from benchmark.evidence_budget import EvidenceBudget +from benchmark.run_mini import list_rules +from benchmark.top50_runner import ( + PhaseResult, + Top50Contract, + Top50Runner, + TriageBudget, + build_rankable_runner, +) + + +def pilot_contract() -> Top50Contract: + """Return explicitly provisional values; issue #68 freezes the public contract.""" + return Top50Contract( + triage=TriageBudget( + model_generations=int(os.environ.get("PRB_TRIAGE_GENERATIONS", "8")), + shell_actions=int(os.environ.get("PRB_TRIAGE_ACTIONS", "12"))), + episode=EvidenceBudget( + model_generations=int(os.environ.get("PRB_EPISODE_GENERATIONS", "10")), + shell_actions=int(os.environ.get("PRB_EPISODE_ACTIONS", "12")), + pred_calls=int(os.environ.get("PRB_PRED_CALLS", "24")), + solve_calls=int(os.environ.get("PRB_SOLVE_CALLS", "10")), + submit_attempts=2, + max_output_chars=int(os.environ.get("PRB_MAX_OUTPUT_CHARS", "10000")), + pred_timeout_seconds=int(os.environ.get("PRB_PRED_TIMEOUT_SECONDS", "300"))), + ) + + +class _FakeExecutor: + def run_triage(self, session, *, repo_path, inventory, model): + payload = session.workdir / "shortlist.json" + payload.write_text(json.dumps(list(inventory[:50])), encoding="utf-8") + session.commit_file(str(payload)) + return PhaseResult(messages=[]) + + def run_episode(self, session, **kwargs): + return PhaseResult(messages=[]) + + +def run(*, model: str, repo_dir: str | Path, output: str | Path, + fake: bool = False, api_base: str | None = None, api_key: str | None = None, + model_kwargs: dict | None = None) -> dict: + repo = Path(repo_dir).resolve() + inventory = list_rules(repo) + if len(inventory) < 50: + raise ValueError(f"canonical inventory has only {len(inventory)} runnable rules") + pred_binary = find_pred_binary() + verify_pred_version(pred_binary) + contract = pilot_contract() + if fake: + runner = Top50Runner( + executor=_FakeExecutor(), contract=contract, pred_binary=pred_binary) + else: + runner = build_rankable_runner( + contract=contract, pred_binary=pred_binary, + agent_uid=int(os.environ["PRB_AGENT_UID"]), + agent_gid=int(os.environ["PRB_AGENT_GID"]), + oracle_uid=int(os.environ["PRB_ORACLE_UID"]), + oracle_gid=int(os.environ["PRB_ORACLE_GID"]), + evidence_gid=int(os.environ["PRB_EVIDENCE_GID"]), + api_base=api_base, api_key=api_key, model_kwargs=model_kwargs) + result = runner.run(model=model, repo_path=repo, inventory=inventory, output=output) + result["library_commit"] = pinned_commit() + result["budget_contract_status"] = "pilot-unfrozen" + # Rewrite once with provenance added after the runner's final checkpoint. + Path(output).write_text(json.dumps(result, indent=2), encoding="utf-8") + return result + + +def main(argv: list[str] | None = None) -> None: + parser = argparse.ArgumentParser(description="Standardized Model API Top50 benchmark") + parser.add_argument("--model", default=os.environ.get("MODEL_NAME"), required=False) + parser.add_argument("--repo-dir", default=os.environ.get("REPO_DIR", "/app/pr-src")) + parser.add_argument("--output", default=os.environ.get("OUTPUT", "/out/submission.json")) + parser.add_argument("--api-base", default=os.environ.get("API_BASE")) + parser.add_argument("--api-key", default=os.environ.get("API_KEY")) + parser.add_argument("--model-kwargs", default=os.environ.get("MODEL_KWARGS")) + parser.add_argument("--fake", action="store_true", default=bool(os.environ.get("FAKE"))) + args = parser.parse_args(argv) + if not args.model: + parser.error("--model (or MODEL_NAME) is required") + kwargs = json.loads(args.model_kwargs) if args.model_kwargs else None + if kwargs is not None and not isinstance(kwargs, dict): + parser.error("--model-kwargs must be a JSON object") + result = run(model=args.model, repo_dir=args.repo_dir, output=args.output, + fake=args.fake, api_base=args.api_base, api_key=args.api_key, + model_kwargs=kwargs) + print(f"Top50 {result['status']} ({len(result['episodes'])}/50 episodes) → {args.output}") + + +if __name__ == "__main__": + main() diff --git a/benchmark/tests/test_run_top50.py b/benchmark/tests/test_run_top50.py new file mode 100644 index 0000000..a69211d --- /dev/null +++ b/benchmark/tests/test_run_top50.py @@ -0,0 +1,36 @@ +"""Entrypoint wiring tests for the standardized Top50 track.""" +from __future__ import annotations + +import json + +import pytest + +from benchmark import run_top50 + + +def test_fake_entrypoint_uses_50_isolated_episodes(tmp_path, monkeypatch): + repo = tmp_path / "repo" + rules = repo / "src" / "rules" + rules.mkdir(parents=True) + for index in range(55): + (rules / f"rule_{index:02d}.rs").write_text("// rule") + pred = tmp_path / "pred" + pred.write_text("#!/bin/sh\necho pred 0.6.0\n") + pred.chmod(0o700) + monkeypatch.setattr(run_top50, "find_pred_binary", lambda: pred) + monkeypatch.setattr(run_top50, "verify_pred_version", lambda binary: "0.6.0") + output = tmp_path / "top50.json" + + result = run_top50.run( + model="fake/model", repo_dir=repo, output=output, fake=True) + + assert len(result["shortlist"]) == 50 + assert len(result["episodes"]) == 50 + assert result["rankable"] is False + assert json.loads(output.read_text())["budget_contract_status"] == "pilot-unfrozen" + + +@pytest.mark.parametrize("forbidden", ["--backend", "--config", "--strategy-file"]) +def test_standard_entrypoint_rejects_custom_harness_options(forbidden): + with pytest.raises(SystemExit): + run_top50.main(["--model", "fake/model", "--fake", forbidden, "custom"]) diff --git a/benchmark/tests/test_top50_runner.py b/benchmark/tests/test_top50_runner.py new file mode 100644 index 0000000..50137a6 --- /dev/null +++ b/benchmark/tests/test_top50_runner.py @@ -0,0 +1,223 @@ +"""Deterministic end-to-end acceptance tests for the self-selected Top50 workflow.""" +from __future__ import annotations + +import json +import time +from pathlib import Path + +import pytest + +from benchmark.top50_runner import ( + PhaseResult, + ShortlistEntry, + Top50Contract, + Top50Runner, + TriageBudget, + build_rankable_runner, + format_status, +) +from benchmark.evidence_budget import EvidenceBudget +from benchmark.verify import Verdict + + +def _fake_pred(tmp_path: Path) -> Path: + script = tmp_path / "real-pred" + script.write_text("#!/bin/sh\necho ran:$*\n", encoding="utf-8") + script.chmod(0o700) + return script + + +def _contract() -> Top50Contract: + return Top50Contract( + triage=TriageBudget(model_generations=3, shell_actions=3), + episode=EvidenceBudget( + model_generations=2, shell_actions=2, pred_calls=1, solve_calls=0, + submit_attempts=2, max_output_chars=1024, pred_timeout_seconds=2), + ) + + +class FakeExecutor: + def __init__(self, shortlist_payload=None, *, fail_episode: int | None = None): + self.shortlist_payload = shortlist_payload + self.fail_episode = fail_episode + self.workspaces: list[Path] = [] + self.initial_statuses: list[dict] = [] + + def run_triage(self, session, *, repo_path, inventory, model): + session.record_model_generation() + session.admit_shell_action("commit-top50 shortlist.json") + payload = self.shortlist_payload + if payload is None: + payload = [{"rule": rule, "hypothesis": f"risk-{index}"} + for index, rule in enumerate(inventory[:50])] + path = session.workdir / "shortlist.json" + path.write_text(json.dumps(payload), encoding="utf-8") + accepted, _ = session.commit_file(str(path)) + return PhaseResult(messages=[{"role": "triage", "accepted": accepted}]) + + def run_episode(self, session, *, repo_path, entry: ShortlistEntry, index, total, model): + self.workspaces.append(session.workdir) + self.initial_statuses.append(session.status()) + assert not (session.workdir / "sentinel").exists() + if index == 1: + (session.workdir / "sentinel").write_text("private") + session.state.reserve("pred_calls") + session.submit._handle( + {"op": "submit", "certificate_text": json.dumps({ + "rule": entry.rule, "source": {}, "bundle": {"target": {"type": "T"}}})}) + session.record_model_generation() + session.admit_shell_action("pwd") + error = "provider unavailable" if self.fail_episode == index else None + return PhaseResult(messages=[{"role": "episode", "rule": entry.rule}], error=error) + + +def _runner(tmp_path: Path, executor) -> Top50Runner: + return Top50Runner( + executor=executor, + contract=_contract(), + pred_binary=_fake_pred(tmp_path), + verifier=lambda cert: Verdict(False, "not a bug"), + ) + + +def test_frozen_top50_runs_in_order_with_fresh_state(tmp_path): + inventory = [f"rule_{index:02d}" for index in range(60)] + executor = FakeExecutor() + result = _runner(tmp_path, executor).run( + model="fake/model", repo_path=tmp_path, inventory=inventory) + + # Injected executors are development-only even when they finish all 50 episodes. + assert result["rankable"] is False + assert [entry["rule"] for entry in result["shortlist"]] == inventory[:50] + assert [episode["rule"] for episode in result["episodes"]] == inventory[:50] + assert len({str(path) for path in executor.workspaces}) == 50 + assert all(status["pred_calls"]["used"] == 0 for status in executor.initial_statuses) + assert all(status["submit_attempts"]["used"] == 0 for status in executor.initial_statuses) + assert all(episode["messages"] == [{"role": "episode", "rule": episode["rule"]}] + for episode in result["episodes"]) + + +@pytest.mark.parametrize("payload", [ + [f"rule_{index:02d}" for index in range(49)], + [f"rule_{index:02d}" for index in range(51)], + ["rule_00"] * 50, + [*[f"rule_{index:02d}" for index in range(49)], "unknown"], + [*[f"rule_{index:02d}" for index in range(49)], + {"rule": "rule_49", "hypothesis": "x" * 501}], +]) +def test_invalid_shortlist_never_starts_an_episode(tmp_path, payload): + executor = FakeExecutor(payload) + result = _runner(tmp_path, executor).run( + model="fake/model", repo_path=tmp_path, + inventory=[f"rule_{index:02d}" for index in range(60)]) + + assert result["rankable"] is False + assert result["episodes"] == [] + assert "valid frozen Top50" in result["run_error"] + + +def test_episode_infrastructure_error_is_partial_and_unrankable(tmp_path): + executor = FakeExecutor(fail_episode=3) + result = _runner(tmp_path, executor).run( + model="fake/model", repo_path=tmp_path, + inventory=[f"rule_{index:02d}" for index in range(60)]) + + assert result["rankable"] is False + assert len(result["episodes"]) == 3 + assert result["episodes"][-1]["status"] == "run_error" + assert "provider unavailable" in result["run_error"] + + +def test_raised_episode_error_is_checkpointed(tmp_path): + class RaisingExecutor(FakeExecutor): + def run_episode(self, session, **kwargs): + if kwargs["index"] == 2: + raise RuntimeError("gateway vanished") + return super().run_episode(session, **kwargs) + + output = tmp_path / "partial.json" + result = _runner(tmp_path, RaisingExecutor()).run( + model="fake/model", repo_path=tmp_path, + inventory=[f"rule_{index:02d}" for index in range(60)], output=output) + + assert result["rankable"] is False + assert len(result["episodes"]) == 2 + assert "gateway vanished" in result["run_error"] + assert json.loads(output.read_text())["run_error"] == result["run_error"] + + +def test_second_shortlist_commit_is_rejected(tmp_path): + from benchmark.top50_runner import TriageSession + + inventory = tuple(f"rule_{index:02d}" for index in range(60)) + with TriageSession(inventory=inventory, budget=TriageBudget(2, 2)) as session: + path = session.workdir / "shortlist.json" + path.write_text(json.dumps(list(inventory[:50]))) + assert session.commit_file(str(path))[0] is True + assert session.commit_file(str(path)) == (False, "shortlist is already frozen") + + +def test_observation_status_reports_every_authoritative_counter(tmp_path): + executor = FakeExecutor() + runner = _runner(tmp_path, executor) + entry = ShortlistEntry("rule_00") + from benchmark.evidence_budget import EvidenceBudgetSession + + with EvidenceBudgetSession( + rule=entry.rule, budget=_contract().episode, pred_binary=runner.pred_binary, + verifier=lambda cert: Verdict(False, "no bug"), + ) as session: + text = format_status(session, index=1, total=50) + + for expected in ("rule 1/50: rule_00", "model generations: 0/2", + "shell actions: 0/2", "pred calls: 0/1", + "solve calls: 0/0", "submit attempts: 0/2"): + assert expected in text + + +def test_rankable_contract_is_fixed_to_model_api_surface(): + executor = FakeExecutor() + assert not hasattr(executor, "backend") + with pytest.raises(ValueError, match="exactly 50"): + Top50Contract(_contract().triage, _contract().episode, shortlist_size=49) + + +def test_only_standard_factory_can_mark_a_run_rankable(tmp_path, monkeypatch): + import benchmark.top50_runner as top50 + + monkeypatch.setattr(top50.os, "geteuid", lambda: 0) + runner = build_rankable_runner( + contract=_contract(), pred_binary=_fake_pred(tmp_path), + agent_uid=10001, agent_gid=10001, oracle_uid=10002, oracle_gid=10002, + evidence_gid=10003) + assert runner._rankable_contract is True + with pytest.raises(ValueError, match="distinct"): + build_rankable_runner( + contract=_contract(), pred_binary=_fake_pred(tmp_path), + agent_uid=10001, agent_gid=10001, oracle_uid=10001, oracle_gid=10002, + evidence_gid=10003) + + +def test_agent_environment_does_not_expose_provider_credentials(monkeypatch): + from benchmark.agent_environment import sanitized_agent_env + + monkeypatch.setenv("OPENAI_API_KEY", "secret") + monkeypatch.setenv("ANTHROPIC_API_KEY", "secret") + monkeypatch.setenv("PRB_PRED_DIR", "/gateway") + env = sanitized_agent_env() + assert "OPENAI_API_KEY" not in env and "ANTHROPIC_API_KEY" not in env + assert env["PRB_PRED_DIR"] == "/gateway" + + +def test_successful_shell_cannot_leave_background_process(tmp_path): + from benchmark.agent_environment import run_as_agent, sanitized_agent_env + import os + + sentinel = tmp_path / "late" + command = f"(sleep 0.2; touch {sentinel}) >/dev/null 2>&1 &" + result = run_as_agent( + command, cwd=str(tmp_path), env=sanitized_agent_env(), timeout=2, + uid=os.getuid(), gid=os.getgid()) + assert result.returncode == 0 + time.sleep(0.4) + assert not sentinel.exists() diff --git a/benchmark/top50_config.yaml b/benchmark/top50_config.yaml new file mode 100644 index 0000000..4a94829 --- /dev/null +++ b/benchmark/top50_config.yaml @@ -0,0 +1,38 @@ +agent: + system_template: | + You are evaluating reduction rules under a fixed logical evidence budget. Each response + must contain exactly one bash command. The authoritative remaining budget is injected into + every observation; elapsed time is not a score and unused budget never transfers. + + {% if phase == "triage" %} + This is source-only triage. Read {{repo_dir}}/src/rules but do not run pred or submit. + Canonical rule inventory: {{inventory}} + Write a JSON list of exactly 50 unique entries to the workspace. Each entry is either a + rule string or {"rule": "...", "hypothesis": "..."}. Freeze it with: + commit-top50 shortlist.json + {% else %} + Investigate only rule {{rule}} ({{rule_index}}/{{rule_total}}). + Triage hypothesis: {{hypothesis}} + Source: {{repo_dir}}/src/rules/{{rule}}.rs + Use pred for dynamic evidence. If you reproduce a bug, write a certificate under + $PRB_ARTIFACT_DIR and call submit immediately. There are exactly two submit attempts for + this rule. Never investigate another rule. + {% endif %} + + Initial authoritative budget: + {{budget_status}} + + Finish normally with: echo COMPLETE_TASK_AND_SUBMIT_FINAL_OUTPUT + instance_template: | + {{task}} + step_limit: 0 + cost_limit: 0 + max_consecutive_format_errors: 3 + +model: + observation_template: | + {% if output.exception_info %}{{output.exception_info}}{% endif %} + {{output.returncode}} + {{output.output}} + format_error_template: | + Format error: {{error}}. Provide exactly one bash command. diff --git a/benchmark/top50_runner.py b/benchmark/top50_runner.py new file mode 100644 index 0000000..3eb1499 --- /dev/null +++ b/benchmark/top50_runner.py @@ -0,0 +1,556 @@ +"""Self-selected Top50 workflow with frozen triage and isolated rule episodes.""" +from __future__ import annotations + +import copy +import json +import os +import shlex +import shutil +import tempfile +import threading +from dataclasses import asdict, dataclass +from pathlib import Path +from typing import Callable, Protocol + +from benchmark.agent_environment import make_agent_environment, run_as_agent, sanitized_agent_env +from benchmark.evidence_budget import EvidenceBudget, EvidenceBudgetSession, EvidenceBudgetState +from benchmark.run_mini import ( + DEFAULT_MAX_TOKENS, + _build_model, + _load_agent_config, + _message_text, + _session_usage, +) +from benchmark.usage import Usage, usage_as_dict +from benchmark.verify import Verdict, verify + +TOP50_SIZE = 50 +DEFAULT_HYPOTHESIS_CHARS = 500 +MAX_SHORTLIST_BYTES = 128 * 1024 +_RANKABLE_TOKEN = object() + + +@dataclass(frozen=True) +class TriageBudget: + model_generations: int + shell_actions: int + max_output_chars: int = 10_000 + command_timeout_seconds: int = 300 + + def __post_init__(self) -> None: + for name, value in asdict(self).items(): + if not isinstance(value, int) or isinstance(value, bool) or value <= 0: + raise ValueError(f"{name} must be a positive integer") + + +@dataclass(frozen=True) +class Top50Contract: + triage: TriageBudget + episode: EvidenceBudget + shortlist_size: int = TOP50_SIZE + hypothesis_chars: int = DEFAULT_HYPOTHESIS_CHARS + + def __post_init__(self) -> None: + if self.shortlist_size != TOP50_SIZE: + raise ValueError("the rankable contract requires exactly 50 rules") + if self.hypothesis_chars <= 0: + raise ValueError("hypothesis_chars must be positive") + + +@dataclass(frozen=True) +class ShortlistEntry: + rule: str + hypothesis: str = "" + + +@dataclass +class PhaseResult: + messages: list[dict] + tokens_k: float = 0.0 + usage: object | None = None + error: str | None = None + + +class PhaseExecutor(Protocol): + def run_triage(self, session: "TriageSession", *, repo_path: Path, + inventory: tuple[str, ...], model: str) -> PhaseResult: ... + + def run_episode(self, session: EvidenceBudgetSession, *, repo_path: Path, + entry: ShortlistEntry, index: int, total: int, + model: str) -> PhaseResult: ... + + +class TriageSession: + """Evaluation-owned source-only workspace and one-shot shortlist controller.""" + + def __init__(self, *, inventory: tuple[str, ...], budget: TriageBudget, + shortlist_size: int = TOP50_SIZE, + hypothesis_chars: int = DEFAULT_HYPOTHESIS_CHARS, + agent_uid: int | None = None, agent_gid: int | None = None): + if len(set(inventory)) != len(inventory): + raise ValueError("canonical inventory contains duplicates") + if (agent_uid is None) != (agent_gid is None): + raise ValueError("agent_uid and agent_gid must be provided together") + self.inventory = inventory + self.budget = budget + self.shortlist_size = shortlist_size + self.hypothesis_chars = hypothesis_chars + self.agent_uid = agent_uid + self.agent_gid = agent_gid + evidence = EvidenceBudget( + model_generations=budget.model_generations, + shell_actions=budget.shell_actions, + pred_calls=0, + solve_calls=0, + submit_attempts=2, + max_output_chars=budget.max_output_chars, + pred_timeout_seconds=budget.command_timeout_seconds, + ) + self.state = EvidenceBudgetState(evidence) + self._tmpdir: Path | None = None + self._workdir: Path | None = None + self._shortlist: tuple[ShortlistEntry, ...] | None = None + self._events: list[dict] = [] + self._lock = threading.RLock() + + @property + def workdir(self) -> Path: + if self._workdir is None: + raise RuntimeError("triage session is not active") + return self._workdir + + @property + def shortlist(self) -> tuple[ShortlistEntry, ...] | None: + with self._lock: + return copy.deepcopy(self._shortlist) + + def __enter__(self) -> "TriageSession": + self._tmpdir = Path(tempfile.mkdtemp(prefix="prb-triage-", dir="/tmp")).resolve() + self._workdir = self._tmpdir / "work" + self._workdir.mkdir(mode=0o700) + if self.agent_uid is not None and self.agent_gid is not None: + self._tmpdir.chmod(0o711) + os.chown(self._workdir, self.agent_uid, self.agent_gid) + self._workdir.chmod(0o700) + return self + + def __exit__(self, exc_type, exc, tb) -> None: + if self._tmpdir is not None: + shutil.rmtree(self._tmpdir, ignore_errors=True) + + def record_model_generation(self, *, outcome: str = "completed", + infrastructure_error: bool = False) -> bool: + reservation = None if infrastructure_error else self.state.reserve("model_generations") + admitted = infrastructure_error or reservation is not None + self._append_event("model_generation", reservation is not None, + outcome if admitted else "budget_exhausted") + return admitted + + def admit_shell_action(self, command: str) -> bool: + reservation = self.state.reserve("shell_actions") + self._append_event("shell_action", reservation is not None, + "admitted" if reservation is not None else "budget_exhausted", + command=command) + return reservation is not None + + def commit_file(self, path: str) -> tuple[bool, str]: + """Validate and atomically freeze a model-authored shortlist JSON file.""" + candidate = Path(path) + if not candidate.is_absolute(): + candidate = self.workdir / candidate + candidate = candidate.resolve() + if not candidate.is_relative_to(self.workdir): + return False, "shortlist file must be inside the triage workspace" + try: + raw = candidate.read_bytes() + except OSError as error: + return False, f"cannot read shortlist: {error}" + if len(raw) > MAX_SHORTLIST_BYTES: + return False, f"shortlist exceeds {MAX_SHORTLIST_BYTES} bytes" + try: + payload = json.loads(raw) + entries = self._validate(payload) + except (UnicodeError, json.JSONDecodeError, ValueError) as error: + return False, str(error) + with self._lock: + if self._shortlist is not None: + return False, "shortlist is already frozen" + self._shortlist = tuple(entries) + self._events.append({"type": "shortlist_commit", "accepted": True, + "rules": [entry.rule for entry in entries]}) + return True, f"frozen {len(entries)} rules" + + def status(self) -> dict: + return self.state.status() + + def ledger(self) -> dict: + with self._lock: + return {"budget": asdict(self.budget), "status": self.status(), + "events": copy.deepcopy(self._events), + "shortlist": ([asdict(entry) for entry in self._shortlist] + if self._shortlist is not None else None)} + + def _validate(self, payload: object) -> list[ShortlistEntry]: + if not isinstance(payload, list) or len(payload) != self.shortlist_size: + raise ValueError(f"shortlist must contain exactly {self.shortlist_size} entries") + entries: list[ShortlistEntry] = [] + for item in payload: + if isinstance(item, str): + rule, hypothesis = item, "" + elif isinstance(item, dict): + rule, hypothesis = item.get("rule"), item.get("hypothesis", "") + else: + raise ValueError("each shortlist entry must be a rule string or object") + if not isinstance(rule, str) or rule not in self.inventory: + raise ValueError(f"unknown rule in shortlist: {rule!r}") + if not isinstance(hypothesis, str) or len(hypothesis) > self.hypothesis_chars: + raise ValueError( + f"hypothesis for {rule!r} must be at most {self.hypothesis_chars} characters") + entries.append(ShortlistEntry(rule, hypothesis)) + rules = [entry.rule for entry in entries] + if len(set(rules)) != len(rules): + raise ValueError("shortlist rules must be unique") + return entries + + def _append_event(self, event_type: str, charged: bool, outcome: str, **extra) -> None: + with self._lock: + self._events.append({"sequence": len(self._events) + 1, "type": event_type, + "charged": charged, "outcome": outcome, + "budget": self.status(), **extra}) + + +class Top50Runner: + """Freeze one self-selected Top50, then execute 50 fresh sequential episodes.""" + + def __init__(self, *, executor: PhaseExecutor, contract: Top50Contract, + pred_binary: str | Path, verifier: Callable[[dict], Verdict] = verify, + agent_uid: int | None = None, agent_gid: int | None = None, + oracle_uid: int | None = None, oracle_gid: int | None = None, + evidence_gid: int | None = None, _rankable_token=None): + self.executor = executor + self.contract = contract + self.pred_binary = Path(pred_binary) + self.verifier = verifier + self.identities = {"agent_uid": agent_uid, "agent_gid": agent_gid, + "oracle_uid": oracle_uid, "oracle_gid": oracle_gid, + "evidence_gid": evidence_gid} + self._rankable_contract = _rankable_token is _RANKABLE_TOKEN + + def run(self, *, model: str, repo_path: str | Path, + inventory: list[str] | tuple[str, ...], output: str | Path | None = None) -> dict: + repo_path = Path(repo_path).resolve() + canonical = tuple(inventory) + with TriageSession( + inventory=canonical, + budget=self.contract.triage, + shortlist_size=self.contract.shortlist_size, + hypothesis_chars=self.contract.hypothesis_chars, + agent_uid=self.identities["agent_uid"], + agent_gid=self.identities["agent_gid"], + ) as triage: + try: + triage_result = self.executor.run_triage( + triage, repo_path=repo_path, inventory=canonical, model=model) + except Exception as error: + triage_result = PhaseResult( + messages=[], error=f"{type(error).__name__}: {error}") + shortlist = triage.shortlist + triage_ledger = triage.ledger() + if triage_result.error: + result = self._result(model, triage_ledger, shortlist, [], triage_result, + f"triage infrastructure error: {triage_result.error}") + return _persist(result, output) + if shortlist is None: + result = self._result(model, triage_ledger, None, [], triage_result, + "triage ended without a valid frozen Top50") + return _persist(result, output) + + episodes: list[dict] = [] + run_error = None + for index, entry in enumerate(shortlist, 1): + episode = None + try: + with EvidenceBudgetSession( + rule=entry.rule, + budget=self.contract.episode, + pred_binary=self.pred_binary, + verifier=self.verifier, + **self.identities, + ) as episode: + phase = self.executor.run_episode( + episode, repo_path=repo_path, entry=entry, + index=index, total=len(shortlist), model=model) + ledger = episode.ledger() + except Exception as error: + phase = PhaseResult(messages=[], error=f"{type(error).__name__}: {error}") + ledger = (episode.ledger() if episode is not None + and episode.submit is not None and episode.pred is not None else {}) + accepted = next((attempt for attempt in ledger.get("submit", []) + if attempt.get("accepted")), None) + record = { + "index": index, + "rule": entry.rule, + "hypothesis": entry.hypothesis, + "status": "run_error" if phase.error else ( + "bug_found" if accepted else "completed"), + "accepted_submit_attempt": accepted.get("attempt") if accepted else None, + "ledger": ledger, + "messages": copy.deepcopy(phase.messages), + "tokens_k": phase.tokens_k, + "usage": _usage_dict(phase.usage), + } + episodes.append(record) + if phase.error: + run_error = f"episode {index} ({entry.rule}) infrastructure error: {phase.error}" + break + checkpoint = self._result( + model, triage_ledger, shortlist, episodes, triage_result, + f"run incomplete after episode {index}/{len(shortlist)}") + _persist(checkpoint, output) + result = self._result( + model, triage_ledger, shortlist, episodes, triage_result, run_error) + return _persist(result, output) + + def _result(self, model: str, triage: dict, + shortlist: tuple[ShortlistEntry, ...] | None, episodes: list[dict], + triage_result: PhaseResult, run_error: str | None) -> dict: + result = { + "model": model, + "status": "run_error" if run_error else "completed", + "rankable": (self._rankable_contract and run_error is None + and len(episodes) == self.contract.shortlist_size), + "contract": {"triage": asdict(self.contract.triage), + "episode": asdict(self.contract.episode), + "shortlist_size": self.contract.shortlist_size, + "hypothesis_chars": self.contract.hypothesis_chars}, + "shortlist": ([asdict(entry) for entry in shortlist] if shortlist else None), + "triage": {"ledger": triage, "messages": copy.deepcopy(triage_result.messages), + "tokens_k": triage_result.tokens_k, + "usage": _usage_dict(triage_result.usage)}, + "episodes": episodes, + } + if run_error: + result["run_error"] = run_error + return result + + +def build_rankable_runner( + *, contract: Top50Contract, pred_binary: str | Path, + agent_uid: int, agent_gid: int, oracle_uid: int, oracle_gid: int, evidence_gid: int, + api_base: str | None = None, api_key: str | None = None, + max_tokens: int = DEFAULT_MAX_TOKENS, model_kwargs: dict | None = None, + verifier: Callable[[dict], Verdict] = verify, +) -> Top50Runner: + """Construct the sole rankable harness after proving the required OS boundary.""" + if os.geteuid() != 0: + raise RuntimeError("rankable Top50 runs require the root runner privilege boundary") + if len({agent_uid, oracle_uid}) != 2: + raise ValueError("agent and oracle must use distinct identities") + executor = MiniSwePhaseExecutor( + api_base=api_base, api_key=api_key, max_tokens=max_tokens, + model_kwargs=model_kwargs, agent_uid=agent_uid, agent_gid=agent_gid, + evidence_gid=evidence_gid) + return Top50Runner( + executor=executor, contract=contract, pred_binary=pred_binary, verifier=verifier, + agent_uid=agent_uid, agent_gid=agent_gid, oracle_uid=oracle_uid, + oracle_gid=oracle_gid, evidence_gid=evidence_gid, + _rankable_token=_RANKABLE_TOKEN) + + +def _usage_dict(usage: object | None) -> dict | None: + return usage_as_dict(usage) if isinstance(usage, Usage) else None + + +def _persist(result: dict, output: str | Path | None) -> dict: + if output is None: + return result + destination = Path(output) + destination.parent.mkdir(parents=True, exist_ok=True) + temporary = destination.with_name(f".{destination.name}.{os.getpid()}.tmp") + temporary.write_text(json.dumps(result, indent=2), encoding="utf-8") + os.replace(temporary, destination) + return result + + +class TriageEnvironment: + """mini-swe environment that exposes source shell actions and intercepted commit-top50.""" + + def __init__(self, session: TriageSession, *, uid: int, gid: int, + extra_groups: tuple[int, ...] = ()): + self.session = session + self.uid = uid + self.gid = gid + self.extra_groups = extra_groups + + def execute(self, action: dict, cwd: str = "", *, timeout: int | None = None) -> dict: + command = action.get("command", "") + if not self.session.admit_shell_action(command): + return {"output": "shell action budget exhausted\n", "returncode": 75, + "exception_info": ""} + try: + words = shlex.split(command) + except ValueError as error: + return {"output": str(error), "returncode": 2, "exception_info": ""} + if words and words[0] == "commit-top50": + if len(words) != 2: + return {"output": "usage: commit-top50 SHORTLIST.json\n", "returncode": 2, + "exception_info": ""} + accepted, message = self.session.commit_file(words[1]) + return {"output": message + "\n", "returncode": 0 if accepted else 2, + "exception_info": ""} + try: + result = run_as_agent( + command, cwd=cwd or str(self.session.workdir), env=sanitized_agent_env(), + timeout=timeout or self.session.budget.command_timeout_seconds, + uid=self.uid, gid=self.gid, extra_groups=self.extra_groups, + max_output_chars=self.session.budget.max_output_chars) + return {"output": result.stdout, "returncode": result.returncode, + "exception_info": ""} + except Exception as error: + return {"output": getattr(error, "output", "") or "", "returncode": -1, + "exception_info": f"{type(error).__name__}: {error}"} + + def get_template_vars(self, **kwargs) -> dict: + return {"cwd": str(self.session.workdir), **kwargs} + + def serialize(self) -> dict: + return {"info": {"environment_type": type(self).__name__}} + + +def format_status(session, *, index: int | None = None, total: int = TOP50_SIZE) -> str: + status = session.status() + lines = [] + if index is not None: + lines.append(f"rule {index}/{total}: {session.rule}") + for key, label in (("model_generations", "model generations"), + ("shell_actions", "shell actions"), + ("pred_calls", "pred calls"), ("solve_calls", "solve calls")): + if key in status: + counter = status[key] + lines.append(f"{label}: {counter['used']}/{counter['limit']}") + if "submit_attempts" in status: + counter = status["submit_attempts"] + lines.append(f"submit attempts: {counter['used']}/{counter['limit']}") + return "\n".join(lines) + + +class MiniSwePhaseExecutor: + """The sole rankable mini-swe/LiteLLM implementation of the phase protocol.""" + + def __init__(self, *, api_base: str | None = None, api_key: str | None = None, + max_tokens: int = DEFAULT_MAX_TOKENS, model_kwargs: dict | None = None, + agent_uid: int | None = None, agent_gid: int | None = None, + evidence_gid: int | None = None): + self.api_base = api_base + self.api_key = api_key + self.max_tokens = max_tokens + self.model_kwargs = model_kwargs + self.agent_uid = os.getuid() if agent_uid is None else agent_uid + self.agent_gid = os.getgid() if agent_gid is None else agent_gid + self.evidence_gid = evidence_gid + config_path = Path(__file__).with_name("top50_config.yaml") + self.agent_config, self.model_config, _ = _load_agent_config( + config_path, config_path, "") + self._models: dict[str, object] = {} + + def run_triage(self, session: TriageSession, *, repo_path: Path, + inventory: tuple[str, ...], model: str) -> PhaseResult: + environment = TriageEnvironment( + session, uid=self.agent_uid, gid=self.agent_gid, + extra_groups=((self.evidence_gid,) if self.evidence_gid is not None else ())) + return self._run_agent( + model, environment, session, + task="Select and commit exactly 50 high-risk reduction rules.", + template_vars={"repo_dir": str(repo_path), "inventory": json.dumps(inventory), + "phase": "triage"}, status=lambda: format_status(session)) + + def run_episode(self, session: EvidenceBudgetSession, *, repo_path: Path, + entry: ShortlistEntry, index: int, total: int, + model: str) -> PhaseResult: + environment = make_agent_environment( + session, uid=self.agent_uid, gid=self.agent_gid, + extra_groups=((self.evidence_gid,) if self.evidence_gid is not None else ())) + return self._run_agent( + model, environment, session, + task=f"Investigate only reduction rule {entry.rule}.", + template_vars={"repo_dir": str(repo_path), "rule": entry.rule, + "hypothesis": entry.hypothesis, "phase": "episode", + "rule_index": index, "rule_total": total}, + status=lambda: format_status(session, index=index, total=total)) + + def _run_agent(self, model_name: str, environment, session, *, task: str, + template_vars: dict, status: Callable[[], str]) -> PhaseResult: + from minisweagent.agents.default import DefaultAgent + from minisweagent.exceptions import FormatError, Submitted + + model = self._models.get(model_name) + if model is None: + model = _build_model( + model_name, self.api_base, self.max_tokens, + model_kwargs=self.model_kwargs, api_key=self.api_key, + observation_template=self.model_config.get("observation_template"), + format_error_template=self.model_config.get("format_error_template")) + self._models[model_name] = model + + class BudgetedAgent(DefaultAgent): + def query(self): + if session.status()["model_generations"]["remaining"] <= 0: + raise Submitted({"role": "exit", "content": "generation budget exhausted", + "extra": {"exit_status": "Submitted", "submission": ""}}) + try: + message = super().query() + except FormatError as error: + session.record_model_generation(outcome="format_error") + budget = status() + for format_message in error.messages: + format_message["content"] = ( + f"[authoritative budget]\n{budget}\n\n" + f"{format_message.get('content', '')}") + raise + except Exception: + session.record_model_generation( + outcome="provider_error", infrastructure_error=True) + raise + session.record_model_generation(outcome="completed") + return message + + def execute_actions(self, message): + actions = message.get("extra", {}).get("actions", []) + budget_text = status() + if len(actions) == 1: + outputs = [self.env.execute(actions[0])] + elif actions: + outputs = [{"output": "exactly one command is required", "returncode": 2, + "exception_info": ""} for _ in actions] + else: + observation = self.model.format_message( + role="user", content=(f"[authoritative budget]\n{budget_text}\n\n" + "Format error: exactly one command is required")) + return self.add_messages(observation) + for output in outputs: + output["output"] = f"[authoritative budget]\n{budget_text}\n\n{output['output']}" + messages = self.add_messages(*self.model.format_observation_messages( + message, outputs, self.get_template_vars())) + current = session.status() + hard_exhausted = any( + current[name]["used"] > 0 and current[name]["remaining"] == 0 + for name in ("model_generations", "shell_actions", "pred_calls") + if name in current) + submit_closed = bool(getattr(getattr(session, "submit", None), "closed", False)) + if hard_exhausted or submit_closed or getattr(session, "shortlist", None) is not None: + raise Submitted({"role": "exit", "content": "phase complete", + "extra": {"exit_status": "Submitted", "submission": ""}}) + return messages + + agent = BudgetedAgent(model, environment, **self.agent_config) + agent.extra_template_vars = template_vars | {"budget_status": status()} + error = None + try: + agent.run(task=task) + except Exception as exception: + error = f"{type(exception).__name__}: {exception}" + tokens_k, usage = _session_usage(agent) + messages = [{"role": message.get("role", ""), "content": _message_text(message)} + for message in agent.messages] + return PhaseResult(messages=messages, tokens_k=tokens_k, + usage=usage, error=error) diff --git a/docker/Dockerfile b/docker/Dockerfile index 0639dee..94478c0 100644 --- a/docker/Dockerfile +++ b/docker/Dockerfile @@ -124,7 +124,7 @@ ENV REPO_DIR=/app/pr-src \ OUTPUT=/out/submission.json VOLUME ["/out"] # Build-time sanity that the assembled runner imports and runs (no API/pred calls). -RUN python -m benchmark.run_submission --fake --model fake/check \ +RUN python -m benchmark.run_top50 --fake --model fake/check \ --output /tmp/buildcheck.json && rm -f /tmp/buildcheck.json -ENTRYPOINT ["python", "-m", "benchmark.run_submission"] +ENTRYPOINT ["python", "-m", "benchmark.run_top50"] CMD []