refactor(eval): consolidate duplicated harness logic after the review round

Simplification pass over the branch. Behavior-preserving throughout; three
reviewers, nine findings applied, two skipped.

The review-artifact block in _run_hidden_oracle was unreachable. That function
runs only in run_arm's non-review branch, while the directory it probes for is
created only in the review branch, and each sandbox serves exactly one arm - so
`review_artifact.parent.is_dir()` could never be true. It was added an hour
earlier to make the hidden oracle resolve the moved artifact; the oracle never
runs for review tasks, so the guard was dead on arrival. Deleting it also
removes the duplication it had with the verify-command wiring.

EXCLUDED_ERROR_KINDS is now one definition. runner.py and comparator_reuse.py
each carried the same six-member frozenset, kept in sync by a comment. Only one
direction is possible: runner already imports from comparator_reuse, so the
reverse import fails at module-init with a circular-import error. That is now
stated where the alias lives, so nobody tries it the other way.

ensure_task_graph and prefetch_next_graph shared ten keyword parameters, passed
through two call sites and forwarded whole between them. They now take a
GraphBuildEnv, mirroring TaskCellContext, which already bundles per-cell state
in this file. Its ready_keys() replaces an inline four-set union at the call
site.

Smaller consolidations: _sha256_file's hand-rolled chunk loop becomes
hashlib.file_digest (3.11+, already used in runner_artifacts); _copy_owner_only
reuses task_assets._write_all and COPY_CHUNK_BYTES instead of repeating the
short-write retry; its stat-then-open existence check becomes the O_EXCL failure
it was already relying on, which is atomic rather than merely narrow; and
runner_environment reads the digest through comparator_reuse.current_runtime_digest
instead of re-parsing the environment variable.

Three test docstrings summarised the branch's own history ("the branch's core
speedup", "the regression that produced fifteen runs") rather than the invariant
under test. Rewritten to state the constraint, which is what survives the merge.
Repaired the indentation left behind by the outage-streak edit and flattened the
prefetch dispatch from three nested conditionals to one.

Skipped: consolidating comparator_reuse._real_directory onto proposer_sandbox's
same-named helper - they differ, the sandbox one rejects any symlink in the
resolved path while this one checks only the leaf, so sharing it would tighten
behavior rather than preserve it. That needs a decision about which policy the
reuse path wants, not a simplification.

589 eval tests pass, ruff clean, 29 workflow contract tests pass. Unrelated and
pre-existing: two test_model_gateway.py failures, and
test_process_control.py::test_timeout_kills_term_ignoring_descendants_before_they_write,
which is a TERM-to-KILL timing flake (passes 2 of 3 in isolation) in a file this
branch does not touch.
This commit is contained in:
Gergo Magyar 2026-09-07 06:10:04 +00:00
parent 81f1b50770
commit 66d69194bd
7 changed files with 170 additions and 169 deletions

View file

@ -330,9 +330,9 @@ def test_clean_control_rewards_an_empty_approval_and_penalizes_noise():
def test_parse_review_output_names_the_actual_failure(tmp_path: Path):
"""One message per cause.
Folding these together is how a sandbox that made the artifact impossible
to write read for fifteen runs as an encoding fault: every cell reported
"not valid UTF-8 JSON" for a file the agent was never able to create.
Folding these together makes a sandbox that renders the artifact impossible
to write indistinguishable from an encoding fault: every cell reports "not
valid UTF-8 JSON" for a file the agent was never able to create.
"""
missing = tmp_path / "never-written.json"

View file

@ -1005,7 +1005,7 @@ def test_progress_line_reports_the_numbers_a_real_run_measured():
def test_review_artifact_is_mounted_as_a_writable_directory_outside_the_workspace(tmp_path):
"""The regression that produced fifteen runs of empty evidence.
"""A writable file inside a read-only directory is not a writable path.
A writable FILE inside a read-only directory is not writable to anything
that writes atomically. The Write tool creates `<target>.tmp.<n>.<hex>`

View file

@ -13,6 +13,7 @@ import yaml
from workflow_bench.runner import (
aggregate,
GraphBuildEnv,
broken_incumbent_arms,
build_parser,
infra_error_record,
@ -596,16 +597,18 @@ def test_prefetch_next_graph_runs_ensure_on_a_background_thread(monkeypatch):
task={"id": "review-b"},
binding={"repo_identity": "/repo", "resolved_sha": "bbb"},
graph_key=("/repo", "bbb"),
trees=Path("/tmp"),
task_asset_cache=None,
claude_bin="claude",
bwrap_bin="bwrap",
sandbox_backend="bwrap",
runtime_mounts=(),
clone_templates={},
clone_template_errors={},
graph_snapshots={},
graph_snapshot_errors={},
env=GraphBuildEnv(
trees=Path("/tmp"),
task_asset_cache=None,
claude_bin="claude",
bwrap_bin="bwrap",
sandbox_backend="bwrap",
runtime_mounts=(),
clone_templates={},
clone_template_errors={},
graph_snapshots={},
graph_snapshot_errors={},
),
cancel_event=cancel,
)
assert job.key == ("/repo", "bbb")
@ -632,3 +635,32 @@ def test_a_reused_resolution_does_not_count_as_this_sweeps_health():
mixed = aggregate([record(resolved=True, reused=True), record(resolved=True)])
assert mixed["resolved_fresh"] == 1
assert broken_incumbent_arms({"t": {"review": mixed}}, {"review"}) == []
def test_graph_build_env_ready_keys_covers_successes_and_failures():
"""A key that failed is attempted, not pending.
next_graph_prefetch_target skips keys already in ready_keys. If a failed
build were omitted, the sweep would prefetch it again every iteration and
pay a full clone and offline index each time for a build that cannot
succeed.
"""
env = GraphBuildEnv(
trees=Path("/tmp"),
task_asset_cache=None,
claude_bin="claude",
bwrap_bin="bwrap",
sandbox_backend="bwrap",
runtime_mounts=(),
clone_templates={("/repo", "aaa"): (Path("/tmp/a"), "aaa")},
clone_template_errors={("/repo", "bbb"): OSError("clone failed")},
graph_snapshots={("/repo", "ccc"): object()},
graph_snapshot_errors={("/repo", "ddd"): OSError("index failed")},
)
assert env.ready_keys() == {
("/repo", "aaa"),
("/repo", "bbb"),
("/repo", "ccc"),
("/repo", "ddd"),
}

View file

@ -1381,7 +1381,7 @@ def test_copy_isolated_tree_does_not_share_git_objects_or_refs(tmp_path):
def test_run_cell_uses_the_clone_template_instead_of_recloning(tmp_path, monkeypatch):
"""The branch's core speedup, which had no coverage at all.
"""run_cell must copy the template, never re-clone.
run_cell takes the clone-template branch on essentially every multi-cell
sweep: it copies a pre-sanitized template rather than paying `git clone

View file

@ -29,6 +29,7 @@ from .evolution import CANDIDATE_ARMS, EVIDENCE_MAX_AGE_DAYS
from .proposer_sandbox import SandboxError
from .runner_sessions import MAX_TRANSCRIPT_BYTES, PARENT_EVENT_STREAM_SOURCE
from .runtime_mounts import CE_ARMS
from .task_assets import COPY_CHUNK_BYTES, _write_all
REUSABLE_COMPARATOR_ARMS = frozenset(
{
@ -378,34 +379,29 @@ def _regular_file(path: Path, *, label: str) -> Path:
def _copy_owner_only(source: Path, destination: Path) -> None:
if destination.exists() or destination.is_symlink():
raise SandboxError(f"reuse destination already exists: {destination}")
descriptor = os.open(
destination,
os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_NOFOLLOW", 0),
0o600,
)
# O_CREAT|O_EXCL is the existence check, and unlike a stat beforehand it is
# atomic: a file appearing between check and open cannot slip through.
try:
descriptor = os.open(
destination,
os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_NOFOLLOW", 0),
0o600,
)
except FileExistsError as exc:
raise SandboxError(f"reuse destination already exists: {destination}") from exc
try:
os.fchmod(descriptor, 0o600)
with open(source, "rb") as handle:
while True:
chunk = handle.read(1024 * 1024)
chunk = handle.read(COPY_CHUNK_BYTES)
if not chunk:
break
view = memoryview(chunk)
while view:
written = os.write(descriptor, view)
if written <= 0:
raise OSError(f"short write while copying {source}")
view = view[written:]
_write_all(descriptor, chunk)
os.fsync(descriptor)
finally:
os.close(descriptor)
def _sha256_file(path: Path) -> str:
digest = hashlib.sha256()
with open(path, "rb") as handle:
for chunk in iter(lambda: handle.read(1024 * 1024), b""):
digest.update(chunk)
return digest.hexdigest()
return hashlib.file_digest(handle, "sha256").hexdigest()

View file

@ -48,6 +48,7 @@ import yaml
from . import runner
from . import runner_sessions
from .comparator_reuse import current_runtime_digest
from .model_gateway import (
ANTHROPIC_API_KEY_ENV,
attach_openai_gateway,
@ -1036,7 +1037,7 @@ def runner_environment(args: argparse.Namespace) -> dict[str, str]:
# process_control replaces the child environment wholesale, so a digest the
# workflow exported reaches the runner only if it is forwarded here. Without
# this the runner stamps no runtime_digest and the reuse lock never engages.
runtime_digest = os.environ.get("RUNTIME_DIGEST", "").strip()
runtime_digest = current_runtime_digest()
if runtime_digest:
env["RUNTIME_DIGEST"] = runtime_digest
if args.auth_token:

View file

@ -58,6 +58,7 @@ from typing import Any
import yaml
from .comparator_reuse import (
REUSE_EXCLUDED_ERROR_KINDS,
ComparatorReuseExpectation,
TaskReuseBinding,
current_runtime_digest,
@ -346,15 +347,6 @@ def _run_hidden_oracle(
# credited patch is captured.
oracle_mount = f"{SANDBOX_WORKSPACE}/{mount_name}"
oracle_env[ORACLE_ENV_VAR] = str(stage_root) if host_unsafe else oracle_mount
review_artifact = review_output_path(sandbox, REVIEW_OUTPUT)
review_mounts: tuple[ReadOnlyMount, ...] = ()
if review_artifact.parent.is_dir():
oracle_env[REVIEW_OUTPUT_ENV_VAR] = sandbox.host_text(
f"{SANDBOX_REVIEW_OUTPUT}/{REVIEW_OUTPUT}"
)
review_mounts = (
ReadOnlyMount(source=review_artifact.parent, target=SANDBOX_REVIEW_OUTPUT),
)
passed, _output = _verification_outcome(
run_verify(
snapshot.command,
@ -363,9 +355,9 @@ def _run_hidden_oracle(
command_prefix=sandbox.command_prefix_for(
read_only_workspace=True,
unshare_network=True,
extra_read_only_mounts=(*review_mounts,)
extra_read_only_mounts=()
if host_unsafe
else (ReadOnlyMount(source=stage_root, target=oracle_mount), *review_mounts),
else (ReadOnlyMount(source=stage_root, target=oracle_mount),),
),
env=oracle_env,
require_pid_namespace=getattr(sandbox, "require_pid_namespace", True),
@ -768,9 +760,11 @@ CHURN_FIELDS = ("diff_files", "diff_insertions", "diff_deletions")
# Rows where the session (or the harness) died carry no measured evidence and
# must not skew efficiency medians or resolve denominators. verify-failed and
# skill-not-invoked rows DO count: those sessions ran and spent real tokens.
EXCLUDED_ERROR_KINDS = frozenset(
{"session-error", "infra-error", "evidence-unverified", "cleanup-failure", "review-evidence-invalid", "cancelled"}
)
# One definition, in comparator_reuse: reuse eligibility and aggregate
# exclusion must never drift apart. The dependency only runs this way -
# comparator_reuse importing back from runner is a circular import.
EXCLUDED_ERROR_KINDS = REUSE_EXCLUDED_ERROR_KINDS
# A sustained upstream outage shows up as a run of session/infra/cleanup
# failures. (cleanup-failure overwrites the primary error_kind, so a
@ -1862,56 +1856,79 @@ def next_graph_prefetch_target(
return None
@dataclass(frozen=True)
class GraphBuildEnv:
"""Per-sweep state every graph build shares, and the caches it fills.
The four dicts are the sweep's memo of what has already been built, keyed by
(repo, sha). They are mutable by design and are written by both the sweep
thread and the prefetch thread, which is safe only because a build is
started for a key exactly once and joined before that key is read.
"""
trees: Path
task_asset_cache: TaskAssetCache
claude_bin: Path | str
bwrap_bin: Path | str
sandbox_backend: str
runtime_mounts: Sequence[ReadOnlyMount]
clone_templates: dict[tuple[str, str], tuple[Path, str]]
clone_template_errors: dict[tuple[str, str], BaseException]
graph_snapshots: dict[tuple[str, str], SanitizedGraphSnapshot]
graph_snapshot_errors: dict[tuple[str, str], BaseException]
def ready_keys(self) -> set[tuple[str, str]]:
"""Keys whose build has already been attempted, successfully or not."""
return (
set(self.clone_templates)
| set(self.clone_template_errors)
| set(self.graph_snapshots)
| set(self.graph_snapshot_errors)
)
def ensure_task_graph(
*,
task: Mapping[str, Any],
repo: Path,
task_sha: str,
graph_key: tuple[str, str],
trees: Path,
task_asset_cache: TaskAssetCache,
claude_bin: Path | str,
bwrap_bin: Path | str,
sandbox_backend: str,
runtime_mounts: Sequence[ReadOnlyMount],
clone_templates: dict[tuple[str, str], tuple[Path, str]],
clone_template_errors: dict[tuple[str, str], BaseException],
graph_snapshots: dict[tuple[str, str], SanitizedGraphSnapshot],
graph_snapshot_errors: dict[tuple[str, str], BaseException],
env: GraphBuildEnv,
) -> None:
"""Build one SHA's sanitized clone template and graph. Idempotent per key."""
if graph_key in graph_snapshots or graph_key in graph_snapshot_errors:
if graph_key in env.graph_snapshots or graph_key in env.graph_snapshot_errors:
return
try:
validate_no_prebuilt_graph_assets(task)
if graph_key not in clone_templates and graph_key not in clone_template_errors:
template = make_worktree(repo, task_sha, trees)
if graph_key not in env.clone_templates and graph_key not in env.clone_template_errors:
template = make_worktree(repo, task_sha, env.trees)
template_head = sanitize_clone_for_hidden_oracles(template)
clone_templates[graph_key] = (template, template_head)
env.clone_templates[graph_key] = (template, template_head)
clone_template: Path | None = None
template_head: str | None = None
if graph_key in clone_templates:
clone_template, template_head = clone_templates[graph_key]
if graph_key in clone_template_errors:
graph_snapshot_errors[graph_key] = clone_template_errors[graph_key]
if graph_key in env.clone_templates:
clone_template, template_head = env.clone_templates[graph_key]
if graph_key in env.clone_template_errors:
env.graph_snapshot_errors[graph_key] = env.clone_template_errors[graph_key]
return
graph_snapshots[graph_key] = prepare_sanitized_graph(
env.graph_snapshots[graph_key] = prepare_sanitized_graph(
task,
repo=repo,
resolved_sha=task_sha,
parent=trees,
cache=task_asset_cache,
claude_bin=claude_bin,
bwrap_bin=bwrap_bin,
sandbox_backend=sandbox_backend,
runtime_mounts=runtime_mounts,
parent=env.trees,
cache=env.task_asset_cache,
claude_bin=env.claude_bin,
bwrap_bin=env.bwrap_bin,
sandbox_backend=env.sandbox_backend,
runtime_mounts=env.runtime_mounts,
clone_template=clone_template,
sanitized_head=template_head,
)
except (ManagedProcessError, OSError, SandboxError, RuntimeError, ValueError) as exc:
graph_snapshot_errors[graph_key] = exc
clone_template_errors.setdefault(graph_key, exc)
env.graph_snapshot_errors[graph_key] = exc
env.clone_template_errors.setdefault(graph_key, exc)
@dataclass
@ -1930,16 +1947,7 @@ def prefetch_next_graph(
task: Mapping[str, Any],
binding: Mapping[str, Any],
graph_key: tuple[str, str],
trees: Path,
task_asset_cache: TaskAssetCache,
claude_bin: Path | str,
bwrap_bin: Path | str,
sandbox_backend: str,
runtime_mounts: Sequence[ReadOnlyMount],
clone_templates: dict[tuple[str, str], tuple[Path, str]],
clone_template_errors: dict[tuple[str, str], BaseException],
graph_snapshots: dict[tuple[str, str], SanitizedGraphSnapshot],
graph_snapshot_errors: dict[tuple[str, str], BaseException],
env: GraphBuildEnv,
cancel_event: threading.Event,
) -> GraphPrefetch:
"""Start clone+graph prep for the next unpaid SHA during paid sessions."""
@ -1951,22 +1959,7 @@ def prefetch_next_graph(
if cancel_event.is_set():
return
print(f"[prefetch_next_graph] clone+graph for {task_sha}")
ensure_task_graph(
task=task,
repo=repo,
task_sha=task_sha,
graph_key=graph_key,
trees=trees,
task_asset_cache=task_asset_cache,
claude_bin=claude_bin,
bwrap_bin=bwrap_bin,
sandbox_backend=sandbox_backend,
runtime_mounts=runtime_mounts,
clone_templates=clone_templates,
clone_template_errors=clone_template_errors,
graph_snapshots=graph_snapshots,
graph_snapshot_errors=graph_snapshot_errors,
)
ensure_task_graph(task=task, repo=repo, task_sha=task_sha, graph_key=graph_key, env=env)
# copy_context, as the worker pool already does at _run_wave: a plain Thread
# does not inherit ContextVars, so without this every run_managed inside the
@ -2091,10 +2084,18 @@ def _run_sweep(
f"reuse-results {reuse_source}: {len(reusable_rows)} comparator "
f"cell(s) match this sweep; candidate arms always run"
)
graph_snapshots: dict[tuple[str, str], SanitizedGraphSnapshot] = {}
graph_snapshot_errors: dict[tuple[str, str], BaseException] = {}
clone_templates: dict[tuple[str, str], tuple[Path, str]] = {}
clone_template_errors: dict[tuple[str, str], BaseException] = {}
graph_env = GraphBuildEnv(
trees=Path(trees),
task_asset_cache=task_asset_cache,
claude_bin=args.claude_bin,
bwrap_bin=bwrap_bin,
sandbox_backend=sandbox_backend,
runtime_mounts=runtime_mounts,
clone_templates={},
clone_template_errors={},
graph_snapshots={},
graph_snapshot_errors={},
)
graph_prefetch: GraphPrefetch | None = None
sweep_rows = list(zip(tasks, task_bindings, oracle_snapshots, strict=True))
@ -2139,26 +2140,13 @@ def _run_sweep(
graph_key = (str(repo), task_sha)
if graph_prefetch is not None and graph_prefetch.key == graph_key:
_join_graph_prefetch()
if paid_cells and graph_key not in graph_snapshots and graph_key not in graph_snapshot_errors:
if paid_cells:
ensure_task_graph(
task=task,
repo=repo,
task_sha=task_sha,
graph_key=graph_key,
trees=Path(trees),
task_asset_cache=task_asset_cache,
claude_bin=args.claude_bin,
bwrap_bin=bwrap_bin,
sandbox_backend=sandbox_backend,
runtime_mounts=runtime_mounts,
clone_templates=clone_templates,
clone_template_errors=clone_template_errors,
graph_snapshots=graph_snapshots,
graph_snapshot_errors=graph_snapshot_errors,
task=task, repo=repo, task_sha=task_sha, graph_key=graph_key, env=graph_env
)
graph_snapshot = graph_snapshots.get(graph_key)
graph_snapshot_error = graph_snapshot_errors.get(graph_key)
clone_template, template_head = clone_templates.get(graph_key, (None, None))
graph_snapshot = graph_env.graph_snapshots.get(graph_key)
graph_snapshot_error = graph_env.graph_snapshot_errors.get(graph_key)
clone_template, template_head = graph_env.clone_templates.get(graph_key, (None, None))
cell_context = TaskCellContext(
task=task,
oracle_snapshot=oracle_snapshot,
@ -2209,56 +2197,40 @@ def _run_sweep(
f"({started_cells}/{total_cells}, {(time.monotonic() - sweep_started) / 60:.0f}m elapsed)"
)
keep(run_idx, arm, record)
if paid_cells:
if graph_prefetch is None and not cancel_event.is_set():
ready_keys = (
set(clone_templates)
| set(clone_template_errors)
| set(graph_snapshots)
| set(graph_snapshot_errors)
if paid_cells and graph_prefetch is None and not cancel_event.is_set():
target = next_graph_prefetch_target(
[(later_task, later_binding) for later_task, later_binding, _ in sweep_rows[index + 1 :]],
arms=args.arms,
runs=args.runs,
reusable_rows=reusable_rows,
reuse_source=reuse_source,
ready_keys=graph_env.ready_keys(),
)
if target is not None:
later_task, later_binding, later_key = target
graph_prefetch = prefetch_next_graph(
task=later_task,
binding=later_binding,
graph_key=later_key,
env=graph_env,
cancel_event=cancel_event,
)
target = next_graph_prefetch_target(
[(later_task, later_binding) for later_task, later_binding, _ in sweep_rows[index + 1 :]],
arms=args.arms,
runs=args.runs,
reusable_rows=reusable_rows,
reuse_source=reuse_source,
ready_keys=ready_keys,
)
if target is not None:
later_task, later_binding, later_key = target
graph_prefetch = prefetch_next_graph(
task=later_task,
binding=later_binding,
graph_key=later_key,
trees=Path(trees),
task_asset_cache=task_asset_cache,
claude_bin=args.claude_bin,
bwrap_bin=bwrap_bin,
sandbox_backend=sandbox_backend,
runtime_mounts=runtime_mounts,
clone_templates=clone_templates,
clone_template_errors=clone_template_errors,
graph_snapshots=graph_snapshots,
graph_snapshot_errors=graph_snapshot_errors,
cancel_event=cancel_event,
)
# A reused success is evidence the pipeline can produce a good
# row, so it resets the consecutive-failure count the same way a
# paid success does. Leaving reused rows out let a streak carry
# across them and trip on stale history.
# A reused success is evidence the pipeline can produce a good row,
# so it resets the consecutive-failure count the same way a paid
# success does. Leaving reused rows out let a streak carry across
# them and trip on stale history.
for _run_idx, _arm, _record in reused_records:
outage_streak = systemic_outage_streak(_record.get("error_kind"), outage_streak)
outage_streak, outage_tripped = sweep_task_cells(
paid_cells,
workers=args.workers,
run=partial(run_cell, cell_context),
on_start=announce,
on_record=keep,
outage_streak=outage_streak,
outage_limit=args.outage_streak,
cancel_event=cancel_event,
)
paid_cells,
workers=args.workers,
run=partial(run_cell, cell_context),
on_start=announce,
on_record=keep,
outage_streak=outage_streak,
outage_limit=args.outage_streak,
cancel_event=cancel_event,
)
results[task["id"]] = {a: aggregate(rs) for a, rs in per_arm.items() if rs}
finally:
_join_graph_prefetch()