fix(eval): cut skill-evolution wall clock without shrinking the gate

Reuse matching incumbent/CE cells, sanitize each SHA once, and default
dispatch workers to 3 so weekly review generations finish inside the
EventBridge window. Cap the sweep from leftover instance uptime so a
Friday dispatch still uploads evidence.

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Gergo Magyar 2026-09-06 06:48:27 +00:00
parent f48bf81256
commit 97afe84cef
19 changed files with 1228 additions and 72 deletions

View file

@ -63,13 +63,14 @@
# uploads, and a promotion (if any) opens a well-formed PR. Run
# 29907431284 (2026-07-22) went green end to end in 14h45m and reached a
# gate decision (`insufficient_evidence`, no promotion).
# [ ] After resizing the runner, prove a manual workers=3 run has zero excluded
# runs and does not stretch the 48-minute serial mean toward the session
# ceiling; then set GITNEXUS_EVOLUTION_WORKERS=3 and
# [ ] Confirm a workers=3 dispatch has zero excluded runs (review sessions in
# 33962002890 averaged ~19m serial, well under the 90m session ceiling).
# Then set GITNEXUS_EVOLUTION_WORKERS=3 and
# GITNEXUS_EVOLUTION_ENABLED=true for scheduled runs. Scheduled runs
# require both values, so leaving workers unset/1 is an immediate rollback;
# workflow_dispatch remains available for the proof and bills real API
# usage on GITNEXUS_BENCH_ANTHROPIC_API_KEY or GITNEXUS_BENCH_OPENAI_API_KEY.
# require both values, so leaving the var unset is an immediate rollback.
# Dispatch defaults to 3; pass workers=1 only to debug a contended host.
# Weekly generations reuse matching incumbent/CE cells from the previous
# artifact so the paid matrix is the new candidate, not a 54-cell replay.
name: GitNexus skill evolution
on:
@ -92,9 +93,9 @@ on:
default: '3'
type: string
workers:
description: 'Benchmark cells of one task to run at once — raise only to match the runner’s vCPUs'
description: 'Benchmark cells of one task to run at once — 3 fits the evolution box; drop to 1 only if siblings hit the session ceiling'
required: false
default: '1'
default: '3'
type: string
model:
description: 'Model for the benchmark arms (match the model your skill users run)'
@ -172,6 +173,12 @@ jobs:
# stops the runner just disappears mid-step. Scheduled runs can start well
# after the cron (the 2026-08-01 run was queued 65min late), so the job
# budget has to absorb that delay and still land inside the uptime window.
# A Friday workflow_dispatch on a box that already booted for Saturday's
# cron inherits leftover uptime, not a fresh 24h. Run 33962002890 started
# Friday 10:57 UTC and vanished at the Saturday 03:00 stop — 51 finished
# sessions never uploaded. run-evolution.sh therefore passes
# --max-runtime-seconds from /proc/uptime so the sweep fails in-process
# and this always() upload still runs.
timeout-minutes: 1260
permissions:
contents: read # The promotion PR uses a short-lived App token minted below.

View file

@ -0,0 +1,177 @@
"""Comparator-row reuse: skip unchanged incumbent/CE cells, never candidates."""
from __future__ import annotations
import hashlib
from datetime import UTC, datetime, timedelta
from pathlib import Path
import pytest
from workflow_bench.comparator_reuse import (
ComparatorReuseExpectation,
TaskReuseBinding,
materialize_reused_row,
row_is_reusable_comparator,
select_reusable_comparator_rows,
)
from workflow_bench.proposer_sandbox import SandboxError
from workflow_bench.runner_sessions import PARENT_EVENT_STREAM_SOURCE
def _digest(text: str = "blob") -> str:
return hashlib.sha256(text.encode()).hexdigest()
def _artifact(name: str = "session-1.jsonl", payload: bytes = b'{"type":"ok"}\n') -> dict:
return {
"path": f"transcripts/{name}",
"sha256": hashlib.sha256(payload).hexdigest(),
"bytes": len(payload),
"source": PARENT_EVENT_STREAM_SOURCE,
}
def _row(**overrides) -> dict:
base = {
"task": "review-pr-2718-defect",
"arm": "review",
"run": 0,
"ok": True,
"error_kind": None,
"model": "gpt-5.6-sol",
"benchmark_model": "gpt-5.6-sol",
"effort": "xhigh",
"sandbox_backend": "bwrap",
"task_base_sha": "a" * 40,
"task_prompt_digest": _digest("prompt"),
"oracle_digest": _digest("oracle"),
"oracle_command_digest": _digest("oracle-cmd"),
"oracle_manifest_digest": _digest("oracle-man"),
"skill_digest": _digest("skill"),
"candidate_overlay_digest": None,
"review_evidence_valid": True,
"review_score": {"weighted_f1": 0.4},
"review_weighted_f1": 0.4,
"transcript_missing": False,
"transcript_artifacts": [_artifact()],
"recorded_at": datetime.now(UTC).isoformat(),
}
base.update(overrides)
return base
def _expected(**overrides) -> ComparatorReuseExpectation:
now = datetime.now(UTC)
values = dict(
model="gpt-5.6-sol",
effort="xhigh",
sandbox_backend="bwrap",
runtime_digest=None,
now=now,
max_age=timedelta(days=90),
tasks={
"review-pr-2718-defect": TaskReuseBinding(
task_base_sha="a" * 40,
task_prompt_digest=_digest("prompt"),
oracle_digest=_digest("oracle"),
oracle_command_digest=_digest("oracle-cmd"),
oracle_manifest_digest=_digest("oracle-man"),
)
},
skill_digests={"review": _digest("skill"), "ce_review": None},
ce_plugin_version="3.24.0",
ce_plugin_manifest_digest=_digest("ce"),
)
values.update(overrides)
return ComparatorReuseExpectation(**values)
def test_matching_incumbent_review_row_is_reusable() -> None:
assert row_is_reusable_comparator(_row(), _expected()) is True
def test_candidate_rows_are_never_reusable() -> None:
assert row_is_reusable_comparator(_row(arm="candidate_review"), _expected()) is False
def test_skill_digest_drift_rejects_reuse() -> None:
assert row_is_reusable_comparator(_row(), _expected(skill_digests={"review": _digest("other")})) is False
def test_excluded_or_failed_rows_are_not_reusable() -> None:
expected = _expected()
assert row_is_reusable_comparator(_row(error_kind="session-error", ok=False), expected) is False
assert row_is_reusable_comparator(_row(ok=False), expected) is False
assert row_is_reusable_comparator(_row(review_evidence_valid=False), expected) is False
assert row_is_reusable_comparator(_row(recorded_at=(datetime.now(UTC) - timedelta(days=91)).isoformat()), expected) is False
def test_runtime_digest_mismatch_rejects_when_both_sides_are_bound() -> None:
row = _row(runtime_digest=_digest("old-cli"))
assert row_is_reusable_comparator(row, _expected(runtime_digest=_digest("new-cli"))) is False
assert row_is_reusable_comparator(row, _expected(runtime_digest=_digest("old-cli"))) is True
assert row_is_reusable_comparator(_row(), _expected(runtime_digest=_digest("new-cli"))) is True
def test_ce_review_matches_plugin_digest_not_repo_skill() -> None:
row = _row(
arm="ce_review",
skill_digest=None,
ce_plugin_version="3.24.0",
ce_plugin_manifest_digest=_digest("ce"),
)
assert row_is_reusable_comparator(row, _expected()) is True
assert (
row_is_reusable_comparator(row, _expected(ce_plugin_manifest_digest=_digest("other")))
is False
)
def test_select_drops_conflicting_duplicates() -> None:
first = _row(review_weighted_f1=0.4)
second = _row(review_weighted_f1=0.9, recorded_at=datetime.now(UTC).isoformat())
selected = select_reusable_comparator_rows([first, second], expected=_expected())
assert selected == {}
same = select_reusable_comparator_rows([first, dict(first)], expected=_expected())
assert ("review-pr-2718-defect", "review", 0) in same
def test_materialize_copies_transcript_and_review_artifacts(tmp_path: Path) -> None:
payload = b'{"type":"result"}\n'
source = tmp_path / "prior"
dest = tmp_path / "fresh"
(source / "transcripts").mkdir(parents=True)
dest.mkdir()
transcript = source / "transcripts" / "session-1.jsonl"
transcript.write_bytes(payload)
transcript.chmod(0o600)
review = source / "review-pr-2718-defect-review-run0.review.json"
review.write_text('{"verdict":"comment"}\n')
patch = source / "review-pr-2718-defect-review-run0.patch"
patch.write_text("diff\n")
row = _row(
review_artifact=review.name,
transcript_artifacts=[_artifact(payload=payload)],
)
copied = materialize_reused_row(row, source_dir=source, dest_dir=dest)
assert copied["reused"] is True
assert copied["reused_from_recorded_at"] == row["recorded_at"]
assert (dest / "transcripts" / "session-1.jsonl").read_bytes() == payload
assert (dest / review.name).read_text() == review.read_text()
assert (dest / patch.name).read_text() == "diff\n"
assert copied["transcript_artifacts"][0]["sha256"] == hashlib.sha256(payload).hexdigest()
def test_materialize_rejects_same_directory_and_missing_transcript(tmp_path: Path) -> None:
source = tmp_path / "prior"
source.mkdir()
row = _row()
with pytest.raises(SandboxError, match="same results directory"):
materialize_reused_row(row, source_dir=source, dest_dir=source)
dest = tmp_path / "fresh"
dest.mkdir()
with pytest.raises(SandboxError, match="missing"):
materialize_reused_row(row, source_dir=source, dest_dir=dest)

View file

@ -15,13 +15,18 @@ import pytest
from workflow_bench import evolve, evolution
from workflow_bench.runner_sessions import PARENT_EVENT_STREAM_SOURCE
from workflow_bench.evolve import (
MIN_INSTANCE_SWEEP_SECONDS,
build_parser,
build_proposer_prompt,
capped_timeout_seconds,
executed_benchmark_arms,
generation_timeout_seconds,
instance_window_budget_from_proc,
instance_window_budget_seconds,
load_jsonl,
proposer_evidence_entries,
read_learnings,
remaining_runtime_seconds,
resolve_incumbent_arms,
runner_argv,
select_evidence,
@ -922,6 +927,27 @@ def test_runner_argv_inserts_ce_review_for_review_overlay(tmp_path):
)
arms = argv[argv.index("--arms") + 1 : argv.index("--promotion-metric")]
assert arms == ["ce_review", "review", "candidate_review"]
assert "--reuse-results" not in argv
def test_runner_argv_forwards_prior_results_for_comparator_reuse(tmp_path):
args = build_parser().parse_args(
["--tasks", "t.yaml", "--model", "pinned", "--arms", "review"]
)
overlay = tmp_path / "overlay"
skill = overlay / ".claude" / "skills" / "gitnexus-review" / "SKILL.md"
skill.parent.mkdir(parents=True)
skill.write_text("candidate")
prior = tmp_path / "prior-bench"
argv = runner_argv(
args,
tmp_path / "bench",
overlay,
task_bindings=[{"id": "task"}],
target_base_digests={},
reuse_results=prior,
)
assert argv[argv.index("--reuse-results") + 1] == str(prior)
def test_runner_argv_omits_proposer_for_manual_overlay(tmp_path):
@ -1118,6 +1144,52 @@ def test_generation_timeout_rejects_unknown_arm() -> None:
)
def test_instance_window_budget_leaves_upload_reserve() -> None:
# Friday 10:57 on a box that booted 02:45 Saturday-window: ~8.2h uptime.
leftover = instance_window_budget_seconds(8.2 * 3600)
assert leftover == int(86_400 - 8.2 * 3600 - 5_400)
assert leftover >= MIN_INSTANCE_SWEEP_SECONDS
with pytest.raises(ValueError, match="only .*s left"):
instance_window_budget_seconds(23.5 * 3600)
with pytest.raises(ValueError, match="uptime must be"):
instance_window_budget_seconds(float("nan"))
def test_instance_window_budget_from_proc_reads_uptime_and_env(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
uptime = tmp_path / "uptime"
uptime.write_text("3600.00 8000.00\n")
monkeypatch.setenv("EVENTBRIDGE_INSTANCE_WINDOW_SECONDS", "20000")
monkeypatch.setenv("EVENTBRIDGE_STOP_RESERVE_SECONDS", "1000")
assert instance_window_budget_from_proc(uptime) == 20000 - 3600 - 1000
with pytest.raises(ValueError, match="cannot read instance uptime"):
instance_window_budget_from_proc(tmp_path / "missing")
def test_capped_timeout_clamps_to_leftover_window() -> None:
assert capped_timeout_seconds(10_000, None) == 10_000
assert capped_timeout_seconds(10_000, 90) == 90
with pytest.raises(ValueError, match="no time remains"):
capped_timeout_seconds(10_000, 0)
started = time.monotonic() - 40
assert remaining_runtime_seconds(max_runtime_seconds=None, started_monotonic=started) is None
leftover = remaining_runtime_seconds(max_runtime_seconds=100, started_monotonic=started)
assert leftover is not None
assert 50 <= leftover <= 60
def test_parser_rejects_non_positive_max_runtime() -> None:
with pytest.raises(SystemExit):
build_parser().parse_args(
["--tasks", "t.yaml", "--model", "pinned", "--max-runtime-seconds", "0"]
)
args = build_parser().parse_args(
["--tasks", "t.yaml", "--model", "pinned", "--max-runtime-seconds", "7200"]
)
assert args.max_runtime_seconds == 7200
@pytest.mark.skipif(sys.platform != "linux", reason="Bubblewrap PID namespaces require Linux")
def test_outer_runner_pid_namespace_kills_setsid_descendant(tmp_path):
try:

View file

@ -32,6 +32,10 @@ def test_review_corpus_is_immutable_and_task_bound():
assert task["ref"] == case["base_sha"]
assert task["sandbox_copy"] == [f"eval/workflow_bench/review_cases/{patch.name}"]
assert task["setup"] == review_case_setup_command(patch.name)
assert any(
dep.get("source") == "gitnexus-shared/dist" and dep.get("target") == "gitnexus-shared/dist"
for dep in task["sandbox_dependencies"]
)
def test_hidden_labels_are_not_recoverable_from_visible_task_input():

View file

@ -212,6 +212,22 @@ def test_prepare_sanitized_graph_builds_once_from_parentless_tree_and_caches_onl
assert removed == [seed]
def test_prepare_sanitized_graph_requires_head_when_given_a_template(tmp_path: Path):
with pytest.raises(SandboxError, match="sanitized HEAD"):
sanitized_graph.prepare_sanitized_graph(
{},
repo=tmp_path,
resolved_sha="b" * 40,
parent=tmp_path,
cache=SimpleNamespace(), # type: ignore[arg-type]
claude_bin="claude",
bwrap_bin="bwrap",
runtime_mounts=(),
clone_template=tmp_path,
sanitized_head=None,
)
def test_graph_snapshot_rejects_arm_sanitization_identity_drift(tmp_path: Path):
assets = SimpleNamespace(
digest="digest",

View file

@ -11,7 +11,7 @@ import io
import json
import time
from workflow_bench.runner_sessions import SessionProgress
from workflow_bench.runner_sessions import SessionProgress, neutralize_ci_log_text
def _drain_lines(stream: io.StringIO) -> list[str]:
@ -325,3 +325,48 @@ def test_cell_failure_detail_line_bounds_a_huge_detail() -> None:
assert line is not None
assert "truncated" in line
assert len(line) < MAX_CELL_DETAIL_CHARS + 200
def test_progress_neutralizes_github_actions_annotation_forms() -> None:
rewritten = neutralize_ci_log_text(
"gitnexus/src/cli/optional-grammars.ts(18,36): error TS2307: Cannot find module "
"'gitnexus-shared'\n::error::Composite projects may not disable incremental compilation.\n"
"##[error]tsc failed"
)
assert "): error TS2307" not in rewritten
assert "): compiler-error TS2307" in rewritten
assert "::error::" not in rewritten
assert "[:]error::" in rewritten
assert "##[error]" not in rewritten
assert "# [error]tsc failed" in rewritten
stream = io.StringIO()
progress = SessionProgress("review-pr-2718-defect-ce_review-run0", stream=stream, heartbeat_s=3600)
events = [
{
"type": "assistant",
"message": {
"content": [{"type": "tool_use", "id": "b1", "name": "Bash", "input": {"command": "npx tsc --noEmit"}}]
},
},
{
"type": "user",
"message": {
"content": [
{
"type": "tool_result",
"tool_use_id": "b1",
"is_error": True,
"content": "gitnexus/src/cli/optional-grammars.ts(18,36): error TS2307: Cannot find module 'gitnexus-shared'",
}
]
},
},
]
for event in events:
_observe(progress, (json.dumps(event) + "\n").encode())
output = stream.getvalue()
assert "): error TS2307" not in output
assert "): compiler-error TS2307" in output
assert "result=error" in output

View file

@ -245,6 +245,9 @@ def test_shipped_scenarios_opt_out_the_cross_module_cell_and_rebuild_graph_asset
assert skipped == ["cross-module-parse-retry"]
assert all(not task.get("sandbox_copy") for task in tasks)
assert all(task["sandbox_dependencies"] for task in tasks)
assert all(
any(dep.get("source") == "gitnexus-shared/dist" for dep in task["sandbox_dependencies"]) for task in tasks
)
assert all(task["oracle"]["command"] and task["oracle"]["files"] for task in tasks)
assert all("./node_modules/.bin/vitest run" in task["oracle"]["command"] for task in tasks)
assert all("npx vitest" not in task["oracle"]["command"] for task in tasks)

View file

@ -1348,3 +1348,26 @@ def test_make_worktree_clone_has_no_tags_but_keeps_all_branches(tmp_path):
current = _git(target, "rev-parse", "HEAD").stdout.strip()
assert current == other_sha
def test_copy_isolated_tree_does_not_share_git_objects_or_refs(tmp_path):
repo = tmp_path / "repo"
repo.mkdir()
_git(repo, "init", "--quiet")
_git(repo, "checkout", "--quiet", "-b", "main")
sha = _git_commit(repo, "base")
clones = tmp_path / "clones"
clones.mkdir()
template = runner.make_worktree(repo, sha, clones)
(template / "marker.txt").write_text("template\n")
copy = runner.copy_isolated_tree(template, clones)
assert copy != template
assert (copy / "marker.txt").read_text() == "template\n"
(copy / "marker.txt").write_text("copy\n")
assert (template / "marker.txt").read_text() == "template\n"
copy_head = _git(copy, "rev-parse", "HEAD").stdout.strip()
template_head = _git(template, "rev-parse", "HEAD").stdout.strip()
assert copy_head == template_head == sha
alternates = copy / ".git" / "objects" / "info" / "alternates"
assert not alternates.exists()

View file

@ -168,6 +168,34 @@ router thresholds as an incumbent policy, not permanent truth. Candidate
changes run offline in the same throwaway clones as the incumbent; production
skills never rewrite themselves from a live task.
On the self-hosted evolution box, `run-evolution.sh` caps the sweep with
`--max-runtime-seconds` derived from `/proc/uptime` (24h EventBridge window
minus a 90-minute upload reserve). A `workflow_dispatch` that lands on an
already-running instance therefore exits in-process instead of vanishing when
the box stops — a cancelled GitHub job skips even `if: always()`, which is
how run 33962002890 lost 51 finished sessions. Local runs are uncapped.
A review generation is 6 tasks × 3 arms × 3 runs. Serial workers=1 at ~19
minutes per session is a 16-hour job (run 33962002890). Two harness changes
cut that without shrinking the gate:
- **Comparator reuse.** `evolve.py` forwards the seed / prior generation as
`--reuse-results`. Incumbent `review` and `ce_review` rows are copied into
the new `results.jsonl` when model, effort, task SHA, prompt digest, oracle
bytes, incumbent skill digest, CE plugin digest, and sandbox backend still
match. Candidate arms always run. A weekly generation with an unchanged
incumbent therefore pays 18 sessions, not 54. A promotion, model change,
task-corpus change, or harness `RUNTIME_DIGEST` change invalidates the
lock and re-runs the comparators.
- **Sanitized clone templates.** Each unique task SHA is cloned and
sanitized once. Cells copy that parentless snapshot (reflink when the
filesystem allows) instead of `git clone --no-local` plus repack/prune/fsck
54 times. Isolation is a private `.git`, not a second copy of full history.
Dispatch defaults to `--workers 3` so those 18 paid cells can overlap. Size
workers to the host: a cell that loses CPU and hits the session ceiling is
an excluded run the gate refuses.
Build an overlay that mirrors only the canonical repo-local skill paths:
```text
@ -368,10 +396,10 @@ paired benchmark as any other candidate.
For ad-hoc use, run the driver on the existing re-evaluation triggers
(model/harness change or 90-day staleness). The repository workflow runs a
deliberate weekly drift check: scheduled concurrency stays serial unless
`GITNEXUS_EVOLUTION_WORKERS` is raised after a funded host-sized proof, and
`--workers` is bounded to 1–8 before paid work starts. `--generations` remains
the only loop bound.
deliberate weekly drift check: dispatch defaults to three concurrent cells
of one task; scheduled concurrency still requires
`GITNEXUS_EVOLUTION_WORKERS=3` after a clean proof. `--workers` is bounded
to 1–8 before paid work starts. `--generations` remains the only loop bound.
## Free-model setup (no paid tokens)

View file

@ -0,0 +1,387 @@
"""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 Mapping, Sequence
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
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
@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
recorded = _parse_recorded_at(row.get("recorded_at"))
if recorded is None or expected.now - recorded > expected.max_age:
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
prior_runtime = row.get("runtime_digest")
if (
isinstance(prior_runtime, str)
and prior_runtime
and expected.runtime_digest
and 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
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 = _real_directory(source_dir, label="reuse source")
dest = _real_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
materialized["reused_from_recorded_at"] = 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 _real_directory(path: Path, *, label: str) -> Path:
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)
source_file = _regular_file(source / Path(*PurePosixPath(relative).parts), label="transcript")
actual_size = source_file.stat().st_size
if actual_size != expected_size:
raise SandboxError(f"reused transcript size drifted: {relative}")
digest = _sha256_file(source_file)
if digest != expected_digest:
raise SandboxError(f"reused transcript digest drifted: {relative}")
dest_dir = dest / "transcripts"
dest_dir.mkdir(mode=0o700, exist_ok=True)
dest_dir.chmod(0o700)
destination = dest_dir / PurePosixPath(relative).name
_copy_owner_only(source_file, 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}")
source_file = _regular_file(source / name, label=label)
_copy_owner_only(source_file, dest / name)
def _regular_file(path: Path, *, label: str) -> Path:
try:
metadata = path.lstat()
except OSError as exc:
raise SandboxError(f"{label} is missing: {path}: {exc}") from exc
if stat.S_ISLNK(metadata.st_mode) or not stat.S_ISREG(metadata.st_mode):
raise SandboxError(f"{label} must be a regular non-symlink file: {path}")
return 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,
)
try:
os.fchmod(descriptor, 0o600)
with open(source, "rb") as handle:
while True:
chunk = handle.read(1024 * 1024)
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:]
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()

View file

@ -823,6 +823,92 @@ def _timeout_arm_key(arm: str) -> str:
return CANDIDATE_ARMS.get(arm, arm)
EVENTBRIDGE_INSTANCE_WINDOW_SECONDS = 86_400
EVENTBRIDGE_STOP_RESERVE_SECONDS = 5_400
MIN_INSTANCE_SWEEP_SECONDS = 600
def instance_window_budget_seconds(
uptime_seconds: float,
*,
window_seconds: int = EVENTBRIDGE_INSTANCE_WINDOW_SECONDS,
reserve_seconds: int = EVENTBRIDGE_STOP_RESERVE_SECONDS,
min_seconds: int = MIN_INSTANCE_SWEEP_SECONDS,
) -> int:
"""Seconds a sweep may run before an EventBridge 24h instance stop.
The dedicated evolution box is started ~15 minutes before the Saturday
cron and stopped 24h later. A ``workflow_dispatch`` that lands on an
already-running box inherits the leftover uptime, not a fresh day.
Run 33962002890 dispatched Friday 10:57 UTC and was still on its last
review cell when the Saturday 03:00 stop cancelled the runner — 51
finished sessions never uploaded because a cancelled job skips even
``if: always()``. Capping the in-process sweep so it *fails* (instead
of vanishing) leaves the reserve for the upload step.
"""
if window_seconds < 1 or reserve_seconds < 0 or min_seconds < 1:
raise ValueError("instance window and minimum must be positive; reserve must be non-negative")
if not math.isfinite(uptime_seconds) or uptime_seconds < 0:
raise ValueError("uptime must be a finite non-negative number")
leftover = int(window_seconds - uptime_seconds - reserve_seconds)
if leftover < min_seconds:
raise ValueError(
f"instance window has only {leftover}s left after a {reserve_seconds}s "
f"upload reserve (uptime {uptime_seconds:.0f}s of {window_seconds}s); "
f"need at least {min_seconds}s"
)
return leftover
def instance_window_budget_from_proc(
uptime_path: Path = Path("/proc/uptime"),
*,
window_seconds: int | None = None,
reserve_seconds: int | None = None,
) -> int:
"""Read host uptime and apply the EventBridge window env overrides."""
window = (
window_seconds
if window_seconds is not None
else int(os.environ.get("EVENTBRIDGE_INSTANCE_WINDOW_SECONDS", str(EVENTBRIDGE_INSTANCE_WINDOW_SECONDS)))
)
reserve = (
reserve_seconds
if reserve_seconds is not None
else int(os.environ.get("EVENTBRIDGE_STOP_RESERVE_SECONDS", str(EVENTBRIDGE_STOP_RESERVE_SECONDS)))
)
try:
uptime = float(uptime_path.read_text().split()[0])
except (OSError, IndexError, ValueError) as exc:
raise ValueError(f"cannot read instance uptime from {uptime_path}: {exc}") from exc
return instance_window_budget_seconds(uptime, window_seconds=window, reserve_seconds=reserve)
def remaining_runtime_seconds(*, max_runtime_seconds: int | None, started_monotonic: float) -> int | None:
"""Seconds left in an optional wall-clock cap, or None when uncapped."""
if max_runtime_seconds is None:
return None
if max_runtime_seconds < 1:
raise ValueError("max runtime must be positive")
leftover = max_runtime_seconds - (time.monotonic() - started_monotonic)
return max(0, int(leftover))
def capped_timeout_seconds(requested: int, remaining: int | None) -> int:
"""Clamp one managed-process timeout to the leftover instance window."""
if requested < 1:
raise ValueError("requested timeout must be positive")
if remaining is None:
return requested
if remaining < 1:
raise ValueError("no time remains in the instance window")
return min(requested, remaining)
def generation_timeout_seconds(
*,
task_count: int,
@ -877,6 +963,7 @@ def runner_argv(
task_bindings: list[dict[str, Any]],
target_base_digests: dict[str, str],
proposer_model: str | None = None,
reuse_results: Path | None = None,
) -> list[str]:
incumbent_arms = resolve_incumbent_arms(overlay_dir, args.arms)
paired_arms = executed_benchmark_arms(incumbent_arms)
@ -927,6 +1014,8 @@ def runner_argv(
argv += ["--ce-plugin-dir", str(args.ce_plugin_dir), "--ce-plugin-version", args.ce_plugin_version]
if args.unsafe_no_bwrap:
argv.append("--unsafe-no-bwrap")
if reuse_results is not None:
argv += ["--reuse-results", str(reuse_results)]
return argv
@ -1199,6 +1288,13 @@ def _require_finite_metric(value: Any, name: str, *, nullable: bool = False, max
raise ValueError(f"promotion has invalid {name}")
def _positive_int(value: str) -> int:
parsed = int(value)
if parsed < 1:
raise argparse.ArgumentTypeError(f"{value} is not a positive integer")
return parsed
def build_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--tasks", required=True, type=Path)
@ -1237,7 +1333,8 @@ def build_parser() -> argparse.ArgumentParser:
"--seed-results",
type=Path,
default=None,
help="prior wfbench results dir used as generation-0 proposer evidence",
help="prior wfbench results dir used as generation-0 proposer evidence "
"and as --reuse-results for unchanged incumbent/CE cells",
)
parser.add_argument(
"--initial-overlay",
@ -1271,6 +1368,13 @@ def build_parser() -> argparse.ArgumentParser:
default=runner_sessions.SESSION_TIMEOUT_SECONDS,
help="per session, seconds",
)
parser.add_argument(
"--max-runtime-seconds",
type=_positive_int,
default=None,
help="wall-clock cap for the whole evolve process (CI sets this from "
"instance uptime so the sweep exits before EventBridge stops the box)",
)
parser.add_argument("--base-url", default=None)
parser.add_argument(
"--anthropic-api-key",
@ -1404,6 +1508,7 @@ def _run_generations(
) -> int:
out_root = args.out_root or Path("results") / time.strftime("wfevolve-%Y%m%d-%H%M%S")
out_root.mkdir(parents=True, exist_ok=True)
started_monotonic = time.monotonic()
evidence_dir: Path | None = args.seed_results
# Only a proposal this driver wrote in this run is stageable: a
# --seed-results tree is an operator-supplied path, and its sibling
@ -1513,6 +1618,16 @@ def _run_generations(
print(f"[gen {generation}] promotion targets contain uncommitted or drifted bytes")
return 1
print(f"[gen {generation}] benchmarking candidate…")
leftover = remaining_runtime_seconds(
max_runtime_seconds=args.max_runtime_seconds,
started_monotonic=started_monotonic,
)
if leftover is not None and leftover < MIN_INSTANCE_SWEEP_SECONDS:
print(
f"[gen {generation}] stopping with {leftover}s left before the "
f"instance window ends; partial evidence is in {out_root}/"
)
return 1
benchmark_argv = runner_argv(
args,
bench_dir,
@ -1520,20 +1635,30 @@ def _run_generations(
task_bindings=selected_tasks,
target_base_digests=target_base_digests,
proposer_model=generation_proposer_model,
reuse_results=evidence_dir,
)
benchmark_command = (
benchmark_argv
if sandbox_backend == "host-unsafe"
else pid_namespace_command(benchmark_argv, bwrap_bin=bwrap_bin)
)
bench = run_managed(
benchmark_command,
timeout=generation_timeout_seconds(
sweep_timeout = capped_timeout_seconds(
generation_timeout_seconds(
task_count=len(selected_task_rows),
runs=args.runs,
session_timeout=args.timeout,
incumbent_arms=incumbent_arms,
),
leftover,
)
if leftover is not None:
print(
f"[gen {generation}] sweep timeout {sweep_timeout}s "
f"(instance window leftover {leftover}s)"
)
bench = run_managed(
benchmark_command,
timeout=sweep_timeout,
env=runner_environment(args),
require_pid_namespace=sandbox_backend == "bwrap",
# The sweep is the multi-hour phase; without this its per-run
@ -1545,7 +1670,13 @@ def _run_generations(
# The sweep runs with GITNEXUS_BENCH_ANTHROPIC_API_KEY in its environment,
# so its detail/stderr tail is a token-bearing sink like any other.
detail = redacted_failure(args, str(bench.detail or bench.stderr_tail[-1000:]))
print(f"[gen {generation}] benchmark run failed ({bench.state}, exit {bench.returncode}): {detail}")
if leftover is not None and bench.state == "timeout":
print(
f"[gen {generation}] benchmark hit the instance-window budget "
f"({sweep_timeout}s); partial evidence is in {bench_dir}: {detail}"
)
else:
print(f"[gen {generation}] benchmark run failed ({bench.state}, exit {bench.returncode}): {detail}")
return 1
promotion = json.loads((bench_dir / "promotion.json").read_text())
for line in summarize_gate(promotion):

View file

@ -15,6 +15,8 @@
# MODEL PROPOSER_MODEL EFFORT GENERATIONS RUNS WORKERS PROVIDER
# EVOLUTION_PROFILE CE_PLUGIN_DIR CE_PLUGIN_VERSION
# INCLUDE_EXPENSIVE SEED_RESULTS CLAUDE_BIN OUT_ROOT
# CI (caps --max-runtime-seconds from /proc/uptime)
# EVENTBRIDGE_INSTANCE_WINDOW_SECONDS EVENTBRIDGE_STOP_RESERVE_SECONDS
# UNSAFE_NO_BWRAP=1 (local review diagnostics only)
# GITNEXUS_BENCH_ANTHROPIC_API_KEY (legacy GITNEXUS_BENCH_AUTH_TOKEN)
# GITNEXUS_BENCH_OPENAI_API_KEY
@ -199,6 +201,20 @@ if ((dry_run)); then
exit 0
fi
# A cancelled GitHub job skips even `if: always()`, so evidence dies with the
# runner. The evolution box is EventBridge-stopped 24h after boot; a Friday
# dispatch inherits leftover uptime. Cap the sweep so it fails in-process and
# the upload step still runs (run 33962002890).
if [[ -n "${CI:-}" && -r /proc/uptime ]]; then
remaining="$(
cd "${eval_dir}"
uv run --locked --extra dev python -c \
'from workflow_bench.evolve import instance_window_budget_from_proc; print(instance_window_budget_from_proc())'
)"
cmd+=(--max-runtime-seconds "${remaining}")
echo "Capping the sweep to ${remaining}s so the instance-window reserve can upload evidence." >&2
fi
mkdir -p "${out_root}"
source_sha="$(git -C "${eval_dir}/.." rev-parse HEAD)"
runtime_digest="$(
@ -220,5 +236,9 @@ SOURCE_SHA="${source_sha}" RUNTIME_DIGEST="${runtime_digest}" SANDBOX_BACKEND="$
}, null, 2) + "\n")' "${out_root}/runtime-provenance.json"
export PYTHONUNBUFFERED=1
# The runner stamps this on every results.jsonl row and refuses to reuse a
# comparator cell when a prior row's digest disagrees. Keep it on the evolve
# process, not only in the provenance JSON sidecar.
export RUNTIME_DIGEST="${runtime_digest}"
cd "${eval_dir}"
exec "${cmd[@]}"

View file

@ -57,6 +57,15 @@ from typing import Any
import yaml
from .comparator_reuse import (
ComparatorReuseExpectation,
TaskReuseBinding,
current_runtime_digest,
default_reuse_max_age,
load_result_rows,
materialize_reused_row,
select_reusable_comparator_rows,
)
from .evolution import (
CANDIDATE_ARMS,
EVIDENCE_MAX_AGE_DAYS,
@ -121,6 +130,7 @@ from .runner_artifacts import (
enforce_phase_workspace,
enforce_work_evidence,
implementation_diff_digest,
copy_isolated_tree,
make_worktree,
new_plan_doc,
parse_shortstat as parse_shortstat,
@ -884,6 +894,8 @@ class TaskCellContext:
candidate_overlay: Path | None
overlay_digest: str | None
sandbox_backend: str = "bwrap"
clone_template: Path | None = None
sanitized_head: str | None = None
def run_cell(ctx: TaskCellContext, run_idx: int, arm: str) -> dict[str, Any]:
@ -908,8 +920,14 @@ def run_cell(ctx: TaskCellContext, run_idx: int, arm: str) -> dict[str, Any]:
raise RuntimeError("sanitized graph snapshot is unavailable")
if ctx.asset_snapshot is None:
raise RuntimeError("task asset snapshot is unavailable")
worktree = make_worktree(ctx.repo, ctx.task_sha, ctx.trees_dir)
sanitized_head = sanitize_clone_for_hidden_oracles(worktree)
if ctx.clone_template is not None:
if not ctx.sanitized_head:
raise RuntimeError("clone template is missing its sanitized HEAD")
worktree = copy_isolated_tree(ctx.clone_template, ctx.trees_dir)
sanitized_head = ctx.sanitized_head
else:
worktree = make_worktree(ctx.repo, ctx.task_sha, ctx.trees_dir)
sanitized_head = sanitize_clone_for_hidden_oracles(worktree)
ctx.graph_snapshot.materialize(worktree, sanitized_head=sanitized_head)
dependency_mounts = stage_task_assets(
task,
@ -1057,6 +1075,7 @@ def run_cell(ctx: TaskCellContext, run_idx: int, arm: str) -> dict[str, Any]:
"task_prompt_digest": hashlib.sha256(task["prompt"].encode()).hexdigest(),
"skill_digest": expected_skill_digest,
"candidate_overlay_digest": (ctx.overlay_digest if arm in CANDIDATE_ARMS else None),
"runtime_digest": current_runtime_digest(),
"recorded_at": datetime.now(UTC).isoformat(),
}
)
@ -1557,6 +1576,15 @@ def build_parser() -> argparse.ArgumentParser:
parser.add_argument("--task-bindings-json", default=None, help=argparse.SUPPRESS)
parser.add_argument("--promotion-target-bases-json", default=None, help=argparse.SUPPRESS)
parser.add_argument("--unsafe-no-bwrap", action="store_true", help=argparse.SUPPRESS)
parser.add_argument(
"--reuse-results",
type=Path,
default=None,
help="prior wfbench results dir whose incumbent/CE rows may be reused "
"when model, effort, tasks, oracles, skill bytes, and CE plugin still "
"match. Candidate arms always run. Used by evolve.py so a weekly "
"generation does not re-pay for an unchanged comparator.",
)
return parser
@ -1684,6 +1712,49 @@ def main() -> None:
gateway.__exit__(None, None, None)
def _comparator_reuse_expectation(
*,
args: argparse.Namespace,
tasks: Sequence[Any],
task_bindings: Sequence[Mapping[str, Any]],
oracle_snapshots: Sequence[Any],
sandbox_backend: str,
ce_plugin_snapshot: CePluginSnapshot | None,
) -> ComparatorReuseExpectation:
"""Bind this sweep's immutable identity for comparator-row reuse."""
skill_digests: dict[str, str | None] = {}
for arm in args.arms:
execution = CANDIDATE_ARMS.get(arm, arm)
if execution in EVALUATED_ARM_SKILLS:
skill_digests[arm] = skill_fingerprint(HARNESS_ROOT, execution)
else:
skill_digests[arm] = None
task_locks: dict[str, TaskReuseBinding] = {}
for task, binding, oracle in zip(tasks, task_bindings, oracle_snapshots, strict=True):
task_locks[str(task["id"])] = TaskReuseBinding(
task_base_sha=str(binding["resolved_sha"]),
task_prompt_digest=hashlib.sha256(str(task["prompt"]).encode()).hexdigest(),
oracle_digest=oracle.digest,
oracle_command_digest=oracle.command_digest,
oracle_manifest_digest=oracle.manifest_digest,
)
return ComparatorReuseExpectation(
model=args.model,
effort=args.effort,
sandbox_backend=sandbox_backend,
runtime_digest=current_runtime_digest(),
now=datetime.now(UTC),
max_age=default_reuse_max_age(),
tasks=task_locks,
skill_digests=skill_digests,
ce_plugin_version=ce_plugin_snapshot.version if ce_plugin_snapshot is not None else None,
ce_plugin_manifest_digest=(
ce_plugin_snapshot.manifest_digest if ce_plugin_snapshot is not None else None
),
)
def _run_sweep(
args: argparse.Namespace,
*,
@ -1740,8 +1811,35 @@ def _run_sweep(
except (OSError, SandboxError, ValueError) as exc:
parser.error(str(exc))
raise AssertionError("ArgumentParser.error() returned unexpectedly")
reuse_source = args.reuse_results.expanduser().resolve() if args.reuse_results is not None else None
reusable_rows: dict[tuple[str, str, int], dict[str, Any]] = {}
if reuse_source is not None:
if reuse_source == out_dir.resolve():
parser.error("--reuse-results cannot be this sweep's --out directory")
raise AssertionError("ArgumentParser.error() returned unexpectedly")
results_file = reuse_source / "results.jsonl"
if results_file.is_symlink() or not results_file.is_file():
parser.error("--reuse-results must contain a regular results.jsonl")
raise AssertionError("ArgumentParser.error() returned unexpectedly")
reusable_rows = select_reusable_comparator_rows(
load_result_rows(results_file),
expected=_comparator_reuse_expectation(
args=args,
tasks=tasks,
task_bindings=task_bindings,
oracle_snapshots=oracle_snapshots,
sandbox_backend=sandbox_backend,
ce_plugin_snapshot=ce_plugin_snapshot,
),
)
print(
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] = {}
for task, task_binding, oracle_snapshot in zip(
tasks,
task_bindings,
@ -1752,41 +1850,79 @@ def _run_sweep(
break
repo = Path(task_binding["repo_identity"])
task_sha = task_binding["resolved_sha"]
per_arm: dict[str, list[dict[str, Any]]] = {a: [] for a in args.arms}
planned = [(run_idx, arm) for run_idx in range(args.runs) for arm in args.arms]
reused_records: list[tuple[int, str, dict[str, Any]]] = []
paid_cells: list[tuple[int, str]] = []
for run_idx, arm in planned:
prior = reusable_rows.get((task["id"], arm, run_idx))
if prior is None or reuse_source is None:
paid_cells.append((run_idx, arm))
continue
try:
reused_records.append(
(
run_idx,
arm,
materialize_reused_row(prior, source_dir=reuse_source, dest_dir=out_dir),
)
)
except (OSError, SandboxError, ValueError) as exc:
print(
f"[{task['id']}][{arm}][run {run_idx}] comparator reuse "
f"failed ({exc}); running a paid cell"
)
paid_cells.append((run_idx, arm))
asset_snapshot: TaskAssetSnapshot | None = None
asset_snapshot_error: BaseException | None = None
graph_key = (str(repo), task_sha)
graph_snapshot: SanitizedGraphSnapshot | None = graph_snapshots.get(graph_key)
graph_snapshot_error: BaseException | None = graph_snapshot_errors.get(graph_key)
try:
validate_no_prebuilt_graph_assets(task)
if graph_snapshot is None and graph_snapshot_error is None:
graph_snapshot = prepare_sanitized_graph(
clone_template: Path | None = None
template_head: str | None = None
if paid_cells:
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, Path(trees))
template_head = sanitize_clone_for_hidden_oracles(template)
clone_templates[graph_key] = (template, template_head)
if graph_key in clone_templates:
clone_template, template_head = clone_templates[graph_key]
if graph_key in clone_template_errors and graph_snapshot_error is None:
graph_snapshot_error = clone_template_errors[graph_key]
if graph_snapshot is None and graph_snapshot_error is None:
graph_snapshot = prepare_sanitized_graph(
task,
repo=repo,
resolved_sha=task_sha,
parent=Path(trees),
cache=task_asset_cache,
claude_bin=args.claude_bin,
bwrap_bin=bwrap_bin,
sandbox_backend=sandbox_backend,
runtime_mounts=runtime_mounts,
clone_template=clone_template,
sanitized_head=template_head,
)
graph_snapshots[graph_key] = graph_snapshot
except (ManagedProcessError, OSError, SandboxError, RuntimeError, ValueError) as exc:
graph_snapshot_error = exc
graph_snapshot_errors[graph_key] = exc
clone_template_errors.setdefault(graph_key, exc)
try:
# Prepared here, once, rather than lazily inside the first cell:
# TaskAssetCache is a plain dict, so a lazy build would be a
# read-then-write race the moment cells stop running serially.
asset_snapshot = task_asset_cache.prepare(
task,
repo=repo,
resolved_sha=task_sha,
parent=Path(trees),
cache=task_asset_cache,
claude_bin=args.claude_bin,
bwrap_bin=bwrap_bin,
sandbox_backend=sandbox_backend,
runtime_mounts=runtime_mounts,
expected_dependency_binding=task_binding,
)
graph_snapshots[graph_key] = graph_snapshot
except (ManagedProcessError, OSError, SandboxError, RuntimeError, ValueError) as exc:
graph_snapshot_error = exc
graph_snapshot_errors[graph_key] = exc
try:
# Prepared here, once, rather than lazily inside the first cell:
# TaskAssetCache is a plain dict, so a lazy build would be a
# read-then-write race the moment cells stop running serially.
asset_snapshot = task_asset_cache.prepare(
task,
repo=repo,
resolved_sha=task_sha,
expected_dependency_binding=task_binding,
)
except (OSError, SandboxError, ValueError) as exc:
asset_snapshot_error = exc
except (OSError, SandboxError, ValueError) as exc:
asset_snapshot_error = exc
cell_context = TaskCellContext(
task=task,
oracle_snapshot=oracle_snapshot,
@ -1805,9 +1941,9 @@ def _run_sweep(
runtime_mounts=runtime_mounts,
candidate_overlay=candidate_overlay,
overlay_digest=overlay_digest,
clone_template=clone_template,
sanitized_head=template_head,
)
per_arm: dict[str, list[dict[str, Any]]] = {a: [] for a in args.arms}
cells = [(run_idx, arm) for run_idx in range(args.runs) for arm in args.arms]
def announce(run_idx: int, arm: str) -> None:
nonlocal started_cells
@ -1830,16 +1966,24 @@ def _run_sweep(
if failure:
print(failure)
outage_streak, outage_tripped = sweep_task_cells(
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,
)
for run_idx, arm, record in reused_records:
started_cells += 1
print(
f"[{task['id']}][{arm}][run {run_idx}] reused comparator "
f"({started_cells}/{total_cells}, {(time.monotonic() - sweep_started) / 60:.0f}m elapsed)"
)
keep(run_idx, arm, record)
if paid_cells:
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,
)
results[task["id"]] = {a: aggregate(rs) for a, rs in per_arm.items() if rs}
selection_report = [

View file

@ -325,6 +325,58 @@ def new_plan_doc(worktree: Path, before: dict[Path, str]) -> Path:
return changed[0]
def _assert_self_contained_git_objects(clone: Path) -> None:
"""Refuse clones that share pack/object bytes with another repository."""
alternates = clone / ".git" / "objects" / "info" / "alternates"
if alternates.exists():
raise RuntimeError(f"clone unexpectedly has an external object alternate: {alternates}")
objects = clone / ".git" / "objects"
if not objects.is_dir():
raise RuntimeError(f"clone is missing a git object store: {clone}")
for obj in objects.rglob("*"):
if obj.is_file() and obj.stat().st_nlink > 1:
raise RuntimeError(f"clone object is hardlinked to host storage: {obj}")
def copy_isolated_tree(source: Path, parent: Path) -> Path:
"""Copy a sanitized clone without sharing git objects or a ref namespace.
``git clone --no-local`` of GitNexus plus ``sanitize_clone_for_hidden_oracles``
(repack/prune/fsck) is minutes per cell. After sanitization the snapshot is
one parentless commit; copying that tree is the isolation boundary the
contamination bug actually required (a private ``.git``), not a second
fetch of full history. Prefer ``cp --reflink=auto`` so XFS/btrfs pay COW;
fall back to a full copy on filesystems that cannot reflink.
"""
try:
source_meta = source.expanduser().lstat()
except OSError as exc:
raise RuntimeError(f"clone template is unavailable: {source}: {exc}") from exc
if stat.S_ISLNK(source_meta.st_mode) or not stat.S_ISDIR(source_meta.st_mode):
raise RuntimeError(f"clone template must be a real directory: {source}")
source = source.expanduser().resolve()
target = Path(tempfile.mkdtemp(prefix="wfbench-", dir=parent))
target.rmdir()
try:
copied = run_managed(
["cp", "-a", "--reflink=auto", str(source), str(target)],
timeout=600,
)
if not copied.ok:
shutil.copytree(source, target, symlinks=True, copy_function=shutil.copy2)
_assert_self_contained_git_objects(target)
return target
except BaseException as primary:
if target.exists():
try:
shutil.rmtree(target)
except OSError as cleanup:
primary.add_note(f"clone copy cleanup also failed: {type(cleanup).__name__}: {cleanup}")
raise
def make_worktree(repo: Path, ref: str, parent: Path) -> Path:
"""Create a self-contained clone per benchmark arm."""
@ -344,12 +396,7 @@ def make_worktree(repo: Path, ref: str, parent: Path) -> Path:
],
timeout=600,
)
alternates = target / ".git" / "objects" / "info" / "alternates"
if alternates.exists():
raise RuntimeError(f"clone unexpectedly has an external object alternate: {alternates}")
for obj in (target / ".git" / "objects").rglob("*"):
if obj.is_file() and obj.stat().st_nlink > 1:
raise RuntimeError(f"clone object is hardlinked to host storage: {obj}")
_assert_self_contained_git_objects(target)
for candidate in (ref, f"origin/{ref}"):
proc = run_managed(
["git", "-C", str(target), "checkout", "--detach", "--quiet", candidate],

View file

@ -57,6 +57,24 @@ MAX_PROGRESS_PENDING = 256
MAX_PROGRESS_TOOL_ID_CHARS = 256
MAX_TOOL_PREVIEW_CHARS = 800
_SAFE_TOOL_NAME = re.compile(r"[A-Za-z0-9._:-]{1,64}")
_GHA_WORKFLOW_COMMAND = re.compile(r"(^|[\n\r])::")
_GHA_HASH_COMMAND = re.compile(r"##\[")
_GHA_COMPILER_ANNOTATION = re.compile(r"\((\d+),(\d+)\):\s+error\b", re.IGNORECASE)
def neutralize_ci_log_text(text: str) -> str:
"""Stop GitHub Actions from promoting tool output into check annotations.
Run 33962002890 logged in-sandbox ``tsc`` failures as
``file.ts(line,col): error TS2307``, which Actions parsed as workflow
annotations on ``.github``. The same parser treats ``::error::`` and
``##[error]`` as commands. Progress previews are evidence, not CI
signaling, so rewrite those forms before they hit the job log.
"""
text = _GHA_WORKFLOW_COMMAND.sub(r"\1[:]", text)
text = _GHA_HASH_COMMAND.sub("# [", text)
return _GHA_COMPILER_ANNOTATION.sub(r"(\1,\2): compiler-error", text)
def _safe_tool_name(value: Any) -> str:
@ -177,7 +195,9 @@ class SessionProgress:
def _say(self, message: str) -> None:
# Queue only: the stdout drain thread calls observe() and must not
# block on a full log pipe (process_control.stdout_observer contract).
self._pending_messages.append(f"[{self.label} {self._elapsed()}] {message}")
self._pending_messages.append(
neutralize_ci_log_text(f"[{self.label} {self._elapsed()}] {message}")
)
self._last_spoke = time.monotonic()
def _emit_pending(self) -> None:

View file

@ -24,7 +24,7 @@ from .proposer_sandbox import (
build_sandbox_environment,
prepare_sandbox,
)
from .runner_artifacts import make_worktree, remove_clone
from .runner_artifacts import copy_isolated_tree, make_worktree, remove_clone
from .task_assets import TaskAssetCache, TaskAssetSnapshot, _is_harness_sandbox_copy
GRAPH_ASSET_PATHS = (
@ -360,14 +360,29 @@ def prepare_sanitized_graph(
bwrap_bin: Path | str,
runtime_mounts: Sequence[ReadOnlyMount],
sandbox_backend: str = "bwrap",
clone_template: Path | None = None,
sanitized_head: str | None = None,
) -> SanitizedGraphSnapshot:
"""Sanitize, index offline once, scrub, and freeze graph assets for all arms."""
"""Sanitize, index offline once, scrub, and freeze graph assets for all arms.
When ``clone_template`` is an already-sanitized snapshot, this copies it
(the copy is scrubbed and indexed) so the template stays a clean cell
seed. Callers that already paid for ``make_worktree`` + sanitization
should pass that template rather than cloning GitNexus again.
"""
validate_no_prebuilt_graph_assets(task)
seed = make_worktree(repo, resolved_sha, parent)
if clone_template is not None:
if not isinstance(sanitized_head, str) or not sanitized_head:
raise SandboxError("clone template requires the sanitized HEAD")
seed = copy_isolated_tree(clone_template, parent)
else:
seed = make_worktree(repo, resolved_sha, parent)
sanitized_head = None
primary: BaseException | None = None
try:
sanitized_head = sanitize_clone_for_hidden_oracles(seed)
if sanitized_head is None:
sanitized_head = sanitize_clone_for_hidden_oracles(seed)
_scrub_source_references(seed)
_neutralize_target_index_inputs(seed)
with prepare_sandbox(

View file

@ -17,6 +17,9 @@ tasks:
- { source: node_modules, target: node_modules }
- { source: gitnexus/node_modules, target: gitnexus/node_modules }
- { source: gitnexus-shared/node_modules, target: gitnexus-shared/node_modules }
# Host-built types/JS. Historical clones have no dist/, so `tsc` in the
# read-only workspace otherwise reports TS2307/TS6379 (run 33962002890).
- { source: gitnexus-shared/dist, target: gitnexus-shared/dist }
- <<: *review_case
id: review-pr-2794-defect

View file

@ -44,6 +44,11 @@ tasks:
target: gitnexus/node_modules
- source: gitnexus-shared/node_modules
target: gitnexus-shared/node_modules
# Host-built types/JS. The clone has no dist/, and node_modules/gitnexus-shared
# is a relative symlink into that unbuilt tree — without this mount, in-sandbox
# `tsc --noEmit` / vitest fail with TS2307 / TS6379 (run 33962002890).
- source: gitnexus-shared/dist
target: gitnexus-shared/dist
prompt: >
Add -j as a short alias for --json on the gitnexus status command
(gitnexus/src/cli/index.ts), and cover the alias with a unit test in

View file

@ -270,12 +270,13 @@ describe('gitnexus skill-evolution workflow contract', () => {
});
it('passes the cell concurrency through to the benchmark', () => {
// The lane is serial unless told otherwise: concurrency only pays off when
// the runner has the vCPUs for it, and a cell starved of CPU drifts toward
// its session timeout, which the gate counts as an excluded run.
// Dispatch defaults to 3. Scheduled runs still fall back to serial unless
// GITNEXUS_EVOLUTION_WORKERS is set — a cell starved of CPU that hits the
// session ceiling is an excluded run the gate refuses.
expect(evolveJob?.env?.WORKERS).toBe(
"${{ inputs.workers || vars.GITNEXUS_EVOLUTION_WORKERS || '1' }}",
);
expect(workflow).toMatch(/workers:\n(?:[^\n]*\n){0,4} default: '3'/);
});
it('seeds from the newest usable completed main run, including failed runs', () => {
@ -572,6 +573,14 @@ exit 1`);
// uploads. The job must finish inside that window even when the schedule
// fires late (the 2026-08-01 run was queued 65 minutes after the cron).
expect(jobBudget as number).toBeLessThanOrEqual(21 * 60);
// A Friday dispatch inherits leftover uptime. The shared entrypoint — not
// the workflow YAML — must cap the sweep so it fails in-process and the
// always() upload still runs (run 33962002890).
const script = readFileSync(path.join(REPO_ROOT, 'eval/workflow_bench/run-evolution.sh'), 'utf8');
expect(script).toContain('--max-runtime-seconds');
expect(script).toContain('instance_window_budget_from_proc');
expect(script).toContain('export RUNTIME_DIGEST');
expect(workflow).toContain('--max-runtime-seconds');
});
it('uploads benchmark evidence unconditionally, on a path it addresses itself', () => {