diff --git a/.github/workflows/gitnexus-skill-evolution.yml b/.github/workflows/gitnexus-skill-evolution.yml index c004bb515..8d02917c1 100644 --- a/.github/workflows/gitnexus-skill-evolution.yml +++ b/.github/workflows/gitnexus-skill-evolution.yml @@ -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. diff --git a/eval/tests/test_comparator_reuse.py b/eval/tests/test_comparator_reuse.py new file mode 100644 index 000000000..bcb8a7b3d --- /dev/null +++ b/eval/tests/test_comparator_reuse.py @@ -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) diff --git a/eval/tests/test_evolve.py b/eval/tests/test_evolve.py index c135584de..bf790d1b6 100644 --- a/eval/tests/test_evolve.py +++ b/eval/tests/test_evolve.py @@ -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: diff --git a/eval/tests/test_review_corpus.py b/eval/tests/test_review_corpus.py index 42f18cb7c..cbaa2695e 100644 --- a/eval/tests/test_review_corpus.py +++ b/eval/tests/test_review_corpus.py @@ -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(): diff --git a/eval/tests/test_sanitized_graph.py b/eval/tests/test_sanitized_graph.py index 7b5dfc093..d7f73d756 100644 --- a/eval/tests/test_sanitized_graph.py +++ b/eval/tests/test_sanitized_graph.py @@ -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", diff --git a/eval/tests/test_session_progress.py b/eval/tests/test_session_progress.py index 9e59ba4bf..416c6206c 100644 --- a/eval/tests/test_session_progress.py +++ b/eval/tests/test_session_progress.py @@ -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 diff --git a/eval/tests/test_workflow_bench.py b/eval/tests/test_workflow_bench.py index 45d0b2f35..889cdee78 100644 --- a/eval/tests/test_workflow_bench.py +++ b/eval/tests/test_workflow_bench.py @@ -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) diff --git a/eval/tests/test_workflow_bench_sessions.py b/eval/tests/test_workflow_bench_sessions.py index 1f66f071f..9a7eb4e87 100644 --- a/eval/tests/test_workflow_bench_sessions.py +++ b/eval/tests/test_workflow_bench_sessions.py @@ -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() diff --git a/eval/workflow_bench/README.md b/eval/workflow_bench/README.md index 79bd4bfea..0b722d340 100644 --- a/eval/workflow_bench/README.md +++ b/eval/workflow_bench/README.md @@ -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) diff --git a/eval/workflow_bench/comparator_reuse.py b/eval/workflow_bench/comparator_reuse.py new file mode 100644 index 000000000..96a0f9a4a --- /dev/null +++ b/eval/workflow_bench/comparator_reuse.py @@ -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() diff --git a/eval/workflow_bench/evolve.py b/eval/workflow_bench/evolve.py index 71a2299cb..3e2774435 100644 --- a/eval/workflow_bench/evolve.py +++ b/eval/workflow_bench/evolve.py @@ -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): diff --git a/eval/workflow_bench/run-evolution.sh b/eval/workflow_bench/run-evolution.sh index c9d2a77a1..396fdac28 100755 --- a/eval/workflow_bench/run-evolution.sh +++ b/eval/workflow_bench/run-evolution.sh @@ -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[@]}" diff --git a/eval/workflow_bench/runner.py b/eval/workflow_bench/runner.py index 82185186c..f78394bf9 100644 --- a/eval/workflow_bench/runner.py +++ b/eval/workflow_bench/runner.py @@ -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 = [ diff --git a/eval/workflow_bench/runner_artifacts.py b/eval/workflow_bench/runner_artifacts.py index 7f05e14b1..02dbd8a38 100644 --- a/eval/workflow_bench/runner_artifacts.py +++ b/eval/workflow_bench/runner_artifacts.py @@ -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], diff --git a/eval/workflow_bench/runner_sessions.py b/eval/workflow_bench/runner_sessions.py index 9daef00b4..053d5e48d 100644 --- a/eval/workflow_bench/runner_sessions.py +++ b/eval/workflow_bench/runner_sessions.py @@ -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: diff --git a/eval/workflow_bench/sanitized_graph.py b/eval/workflow_bench/sanitized_graph.py index 3aea9e8c1..fc3d766a2 100644 --- a/eval/workflow_bench/sanitized_graph.py +++ b/eval/workflow_bench/sanitized_graph.py @@ -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( diff --git a/eval/workflow_bench/tasks.review.scenarios.yaml b/eval/workflow_bench/tasks.review.scenarios.yaml index 8ff323ed1..6e285a2bc 100644 --- a/eval/workflow_bench/tasks.review.scenarios.yaml +++ b/eval/workflow_bench/tasks.review.scenarios.yaml @@ -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 diff --git a/eval/workflow_bench/tasks.scenarios.yaml b/eval/workflow_bench/tasks.scenarios.yaml index af4755f88..59448ae4f 100644 --- a/eval/workflow_bench/tasks.scenarios.yaml +++ b/eval/workflow_bench/tasks.scenarios.yaml @@ -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 diff --git a/gitnexus/test/unit/skill-evolution-workflow.test.ts b/gitnexus/test/unit/skill-evolution-workflow.test.ts index b10313e42..5b3890c3a 100644 --- a/gitnexus/test/unit/skill-evolution-workflow.test.ts +++ b/gitnexus/test/unit/skill-evolution-workflow.test.ts @@ -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', () => {