GitNexus/eval/workflow_bench/comparator_reuse.py
Gergo Magyar b97ad89f38 fix(eval): address PR review feedback (#3207)
- aggregate: count admissible rows directly instead of subtracting the
  execution and evidence counters, which double-charged a row that is both
  a session error and invalid review evidence and could report UNUSABLE for
  an arm holding real measurements.
- run_proposer: bound the session timeout by what is left of
  --max-runtime-seconds, so clearing the sweep minimum cannot start a
  full-length session past the instance window.
- comparator reuse: hold one O_NOFOLLOW descriptor for the size check,
  digest and copy, and prove it is the inode that was checked, closing the
  swap window a concurrent writer of the reuse directory had.
- Drive the review-artifact mount assertion through run_arm and the
  clone-template assertion through run_cell, instead of rebuilding the
  expected values in the tests (also removes the CodeQL unnecessary lambda).
- Assert the workflow invokes run-evolution.sh rather than that its YAML
  mentions --max-runtime-seconds, which only appears in a comment.
- Correct the parse_review_output failure-mode claim: the fold was empty
  artifacts reported as "not valid UTF-8 JSON"; a never-created file raised
  FileNotFoundError.
- prettier: wrap the over-long readFileSync call flagged by PR autofix.

Note: pre-existing failure in tests/test_model_gateway.py::test_locked_litellm_translates_messages_to_offline_responses (local LiteLLM proxy never becomes ready in this environment) not addressed by this PR.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-07 17:02:37 +00:00

449 lines
18 KiB
Python

"""Reuse frozen comparator cells when the current sweep is still the same experiment.
Weekly skill evolution re-runs incumbent ``review`` / ``ce_review`` (and the
implementation incumbents) even when the model, effort, tasks, oracles,
incumbent skill bytes, and CE plugin have not changed. Those arms are the
baseline the gate compares a *new* candidate against — they are not the
thing being evolved. Replaying them burns two-thirds of a generation.
This module selects prior ``results.jsonl`` rows that are safe to carry
forward. Candidate arms are never reused. A mismatch on any bound field
falls through to a paid cell. Missing artifacts also fall through: a reused
row that the proposer cannot read is worse than spending the tokens again.
"""
from __future__ import annotations
import hashlib
import json
import os
import re
import stat
from collections.abc import Iterator, Mapping, Sequence
from contextlib import contextmanager
from dataclasses import dataclass
from datetime import UTC, datetime, timedelta
from pathlib import Path, PurePosixPath
from typing import Any
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(
{
"review",
"ce_review",
"workflow",
"workflow_direct",
"ce_workflow",
"ce_workflow_direct",
"baseline",
"baseline_nomcp",
}
)
# Must stay aligned with runner.EXCLUDED_ERROR_KINDS plus review-invalid.
# A reused row becomes promotion evidence; excluded kinds cannot enter that set.
REUSE_EXCLUDED_ERROR_KINDS = frozenset(
{
"session-error",
"infra-error",
"evidence-unverified",
"cleanup-failure",
"review-evidence-invalid",
"cancelled",
}
)
_TRANSCRIPT_NAME = re.compile(r"[A-Za-z0-9._-]{1,200}")
CellKey = tuple[str, str, int]
@dataclass(frozen=True)
class TaskReuseBinding:
"""Per-task identity the prior row must still match."""
task_base_sha: str
task_prompt_digest: str
oracle_digest: str
oracle_command_digest: str
oracle_manifest_digest: str
# The cell's environment is part of its identity: a comparator measured
# against different task assets or different sandbox dependencies is a
# measurement of a different machine, not a baseline for this sweep.
task_asset_manifest_digest: str | None = None
sandbox_dependency_manifest_digest: str | None = None
@dataclass(frozen=True)
class ComparatorReuseExpectation:
"""Sweep-wide lock for comparator reuse. Any drift pays for a fresh cell."""
model: str
effort: str
sandbox_backend: str
runtime_digest: str | None
now: datetime
max_age: timedelta
tasks: Mapping[str, TaskReuseBinding]
skill_digests: Mapping[str, str | None]
ce_plugin_version: str | None
ce_plugin_manifest_digest: str | None
def load_result_rows(path: Path) -> list[dict[str, Any]]:
"""Load ``results.jsonl``; skip malformed lines the same way evolve does."""
rows: list[dict[str, Any]] = []
for line in path.read_text().splitlines():
if not line.strip():
continue
try:
row = json.loads(line)
except json.JSONDecodeError:
continue
if isinstance(row, dict):
rows.append(row)
return rows
def current_runtime_digest() -> str | None:
"""Harness lockfile digest exported by ``run-evolution.sh``, if present."""
value = os.environ.get("RUNTIME_DIGEST", "").strip()
return value or None
def row_is_reusable_comparator(row: Mapping[str, Any], expected: ComparatorReuseExpectation) -> bool:
"""True when ``row`` is a complete, still-valid comparator measurement."""
arm = row.get("arm")
if not isinstance(arm, str) or arm in CANDIDATE_ARMS or arm not in REUSABLE_COMPARATOR_ARMS:
return False
if row.get("error_kind") in REUSE_EXCLUDED_ERROR_KINDS:
return False
if row.get("error_kind") not in (None, ""):
return False
if row.get("ok") is not True:
return False
if row.get("transcript_missing") is True:
return False
if row.get("candidate_overlay_digest") not in (None, ""):
return False
# Age against the ORIGINAL measurement, not the copy time: materialize_reused_row
# restamps recorded_at, so a chained row would otherwise refresh its own clock
# and never expire. Bound both directions - a future stamp is corrupt, not fresh.
recorded = _parse_recorded_at(row.get("reused_from_recorded_at") or row.get("recorded_at"))
if recorded is None:
return False
age = expected.now - recorded
if age > expected.max_age or age < timedelta(0):
return False
if row.get("model") != expected.model and row.get("benchmark_model") != expected.model:
return False
if row.get("effort") != expected.effort:
return False
if row.get("sandbox_backend") != expected.sandbox_backend:
return False
# Fail closed. A row with no runtime_digest was measured by a harness that
# did not record one, which is exactly the drift this lock exists to catch;
# treating the absence as agreement made every legacy row reusable forever.
prior_runtime = row.get("runtime_digest")
if not isinstance(prior_runtime, str) or not prior_runtime:
return False
if not expected.runtime_digest or prior_runtime != expected.runtime_digest:
return False
task_id = row.get("task")
binding = expected.tasks.get(task_id) if isinstance(task_id, str) else None
if binding is None:
return False
if row.get("task_base_sha") != binding.task_base_sha:
return False
if row.get("task_prompt_digest") != binding.task_prompt_digest:
return False
if row.get("oracle_digest") != binding.oracle_digest:
return False
if row.get("oracle_command_digest") != binding.oracle_command_digest:
return False
if row.get("oracle_manifest_digest") != binding.oracle_manifest_digest:
return False
# Fail closed on both sides, as the runtime digest does: an unbound
# expectation means this sweep could not determine its own environment, and
# a row without the field was measured before it was recorded.
for field, bound in (
("task_asset_manifest_digest", binding.task_asset_manifest_digest),
("sandbox_dependency_manifest_digest", binding.sandbox_dependency_manifest_digest),
):
prior = row.get(field)
if not isinstance(prior, str) or not prior or not bound or prior != bound:
return False
if arm in CE_ARMS:
if row.get("ce_plugin_version") != expected.ce_plugin_version:
return False
if row.get("ce_plugin_manifest_digest") != expected.ce_plugin_manifest_digest:
return False
else:
expected_skill = expected.skill_digests.get(arm)
if not expected_skill or row.get("skill_digest") != expected_skill:
return False
if arm in {"review", "ce_review"}:
if row.get("review_evidence_valid") is not True:
return False
if not isinstance(row.get("review_score"), dict):
return False
if row.get("review_weighted_f1") is None:
return False
artifacts = row.get("transcript_artifacts")
if not isinstance(artifacts, list) or not artifacts:
return False
try:
for artifact in artifacts:
_transcript_metadata(artifact)
except SandboxError:
return False
return True
def select_reusable_comparator_rows(
rows: Sequence[Mapping[str, Any]],
*,
expected: ComparatorReuseExpectation,
) -> dict[CellKey, dict[str, Any]]:
"""Index reusable rows by ``(task, arm, run)``. Conflicting duplicates drop the key."""
chosen: dict[CellKey, dict[str, Any]] = {}
blocked: set[CellKey] = set()
for row in rows:
if not row_is_reusable_comparator(row, expected):
continue
task_id = row["task"]
arm = row["arm"]
run = row.get("run")
if not isinstance(run, int) or isinstance(run, bool) or run < 0:
continue
key = (str(task_id), str(arm), run)
if key in blocked:
continue
previous = chosen.get(key)
if previous is None:
chosen[key] = dict(row)
continue
if _row_identity(previous) != _row_identity(row):
blocked.add(key)
chosen.pop(key, None)
return chosen
def materialize_reused_row(
row: Mapping[str, Any],
*,
source_dir: Path,
dest_dir: Path,
) -> dict[str, Any]:
"""Copy digest-bound artifacts into this sweep's evidence dir and stamp reuse."""
source = _resolved_directory(source_dir, label="reuse source")
dest = _resolved_directory(dest_dir, label="reuse destination")
if source == dest:
raise SandboxError("comparator reuse cannot read and write the same results directory")
materialized = dict(row)
materialized["reused"] = True
# Keep the FIRST measurement time across a chain. Overwriting it with the
# previous copy's stamp let a row refresh its own clock every generation and
# outlive the max_age bound entirely.
materialized["reused_from_recorded_at"] = row.get("reused_from_recorded_at") or row.get("recorded_at")
materialized["recorded_at"] = datetime.now(UTC).isoformat()
artifacts = row.get("transcript_artifacts")
if not isinstance(artifacts, list) or not artifacts:
raise SandboxError("reused row is missing transcript_artifacts")
copied_artifacts: list[dict[str, Any]] = []
for artifact in artifacts:
copied_artifacts.append(_copy_transcript_artifact(source, dest, artifact))
materialized["transcript_artifacts"] = copied_artifacts
review_name = row.get("review_artifact")
if isinstance(review_name, str) and review_name:
_copy_named_artifact(source, dest, review_name, label="review artifact")
task = row.get("task")
arm = row.get("arm")
run = row.get("run")
if isinstance(task, str) and isinstance(arm, str) and isinstance(run, int) and not isinstance(run, bool):
patch_name = f"{task}-{arm}-run{run}.patch"
patch = source / patch_name
if patch.is_file() and not patch.is_symlink():
_copy_named_artifact(source, dest, patch_name, label="patch artifact")
return materialized
def default_reuse_max_age() -> timedelta:
return timedelta(days=EVIDENCE_MAX_AGE_DAYS)
def _row_identity(row: Mapping[str, Any]) -> tuple[Any, ...]:
return (
row.get("skill_digest"),
row.get("oracle_digest"),
row.get("review_weighted_f1"),
row.get("ce_plugin_manifest_digest"),
row.get("recorded_at"),
)
def _parse_recorded_at(value: Any) -> datetime | None:
if not isinstance(value, str) or not value:
return None
try:
parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
except ValueError:
return None
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=UTC)
return parsed.astimezone(UTC)
def _transcript_metadata(metadata: Any) -> tuple[str, str, int]:
if not isinstance(metadata, dict) or set(metadata) != {"path", "sha256", "bytes", "source"}:
raise SandboxError("transcript artifact metadata must contain only path, sha256, bytes, and source")
relative = metadata["path"]
digest = metadata["sha256"]
size = metadata["bytes"]
if metadata["source"] != PARENT_EVENT_STREAM_SOURCE:
raise SandboxError("transcript artifact source is not the parent event stream")
if not isinstance(relative, str) or not isinstance(digest, str) or not re.fullmatch(r"[0-9a-f]{64}", digest):
raise SandboxError("transcript artifact metadata is malformed")
if not isinstance(size, int) or isinstance(size, bool) or size < 0 or size > MAX_TRANSCRIPT_BYTES:
raise SandboxError("transcript artifact byte count is out of range")
relative_path = PurePosixPath(relative)
if (
relative_path.is_absolute()
or len(relative_path.parts) != 2
or relative_path.parts[0] != "transcripts"
or any(part in {"", ".", ".."} for part in relative_path.parts)
or _TRANSCRIPT_NAME.fullmatch(relative_path.parts[1]) is None
):
raise SandboxError(f"unsafe transcript artifact path: {relative!r}")
return relative, digest, size
def _resolved_directory(path: Path, *, label: str) -> Path:
"""An existing, non-symlink directory, resolved through its parents.
Deliberately weaker than proposer_sandbox's same-shaped helper, which
refuses every symlink hop in the path. That one guards a MOUNT ROOT, where
a hop changes what an untrusted session is handed. This one guards a DATA
directory whose contents are validated individually anyway - every file
read goes through ``_regular_file`` (lstat, symlinks rejected) and every
write through ``O_NOFOLLOW`` - so a symlinked parent grants nothing those
guards do not already cover, while refusing one would reject ordinary
setups such as a symlinked artifacts directory or macOS's /var.
Separately named because they make different promises. Do not merge them
without first deciding which promise the reuse path should make.
"""
resolved = path.expanduser()
try:
metadata = resolved.lstat()
except OSError as exc:
raise SandboxError(f"{label} is unavailable: {resolved}: {exc}") from exc
if stat.S_ISLNK(metadata.st_mode) or not stat.S_ISDIR(metadata.st_mode):
raise SandboxError(f"{label} must be a real directory: {resolved}")
return resolved.resolve()
def _copy_transcript_artifact(source: Path, dest: Path, metadata: Mapping[str, Any]) -> dict[str, Any]:
relative, expected_digest, expected_size = _transcript_metadata(metadata)
dest_dir = dest / "transcripts"
dest_dir.mkdir(mode=0o700, exist_ok=True)
dest_dir.chmod(0o700)
destination = dest_dir / PurePosixPath(relative).name
# One descriptor for the size check, the digest and the copy. Re-opening the
# path between them is what let a concurrent writer swap the checked file
# for a symlink and have the copy follow it.
with _open_regular(source / Path(*PurePosixPath(relative).parts), label="transcript") as source_fd:
if os.fstat(source_fd).st_size != expected_size:
raise SandboxError(f"reused transcript size drifted: {relative}")
digest = _sha256_descriptor(source_fd)
if digest != expected_digest:
raise SandboxError(f"reused transcript digest drifted: {relative}")
_copy_owner_only(source_fd, destination)
return {"path": relative, "sha256": digest, "bytes": expected_size, "source": PARENT_EVENT_STREAM_SOURCE}
def _copy_named_artifact(source: Path, dest: Path, name: str, *, label: str) -> None:
relative = PurePosixPath(name)
if relative.is_absolute() or len(relative.parts) != 1 or relative.parts[0] in {"", ".", ".."}:
raise SandboxError(f"unsafe {label} path: {name!r}")
with _open_regular(source / name, label=label) as source_fd:
_copy_owner_only(source_fd, dest / name)
@contextmanager
def _open_regular(path: Path, *, label: str) -> Iterator[int]:
"""Open a regular non-symlink file and hold it open for every later read.
Checking the path and then re-opening it is a race the reuse directory is
exposed to: it is written by a previous sweep and read by this one, so a
concurrent writer can replace a validated file with a symlink in between.
O_NOFOLLOW refuses the leaf link and the fstat comparison proves the open
descriptor is the inode that was checked — the same guarantee
evolution._bounded_regular_bytes makes for evidence files.
"""
try:
before = path.lstat()
except OSError as exc:
raise SandboxError(f"{label} is missing: {path}: {exc}") from exc
if stat.S_ISLNK(before.st_mode) or not stat.S_ISREG(before.st_mode):
raise SandboxError(f"{label} must be a regular non-symlink file: {path}")
try:
descriptor = os.open(path, os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0))
except OSError as exc:
raise SandboxError(f"{label} is unreadable: {path}: {exc}") from exc
try:
opened = os.fstat(descriptor)
if not stat.S_ISREG(opened.st_mode) or (opened.st_dev, opened.st_ino) != (before.st_dev, before.st_ino):
raise SandboxError(f"{label} changed while opening: {path}")
yield descriptor
finally:
os.close(descriptor)
def _copy_owner_only(source: int, destination: Path) -> None:
# 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)
os.lseek(source, 0, os.SEEK_SET)
while True:
chunk = os.read(source, COPY_CHUNK_BYTES)
if not chunk:
break
_write_all(descriptor, chunk)
os.fsync(descriptor)
finally:
os.close(descriptor)
def _sha256_descriptor(descriptor: int) -> str:
os.lseek(descriptor, 0, os.SEEK_SET)
# dup so hashlib owns a file object it may close; the duplicate shares the
# offset, which is why every reader here seeks to 0 before it starts.
with os.fdopen(os.dup(descriptor), "rb") as handle:
return hashlib.file_digest(handle, "sha256").hexdigest()