From b1d8688dbbbdbd4bfacb50bd4ba5c424e19801be Mon Sep 17 00:00:00 2001 From: Gergo Magyar Date: Sun, 6 Sep 2026 10:10:28 +0000 Subject: [PATCH] perf(eval): price the benchmark against measured cell durations MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The cost model assumed every cell runs the 1140s mean. Cells are not uniform: the 41 rows in Actions run 33912693948's artifact are 826s at the median, 1262s at the mean, 2976s at p90, with two pinned at the 5400s session ceiling. A wave waits for its slowest cell, so a mean understates every concurrent schedule — the previous model called workers=3 cold 5.99h when the same schedule against real durations is 10.33h. session_durations.json carries the sample in submission order with its provenance and its caveat: every cell in that run returned unusable evidence, so the durations are real but a clean run may sit lower. It is the only live artifact; the 2026-07-22 green run's has expired. The model now simulates the schedule cell by cell rather than multiplying a mean by a wave count, averaged over all 41 rotations of the sample so no single alignment between sample order and cell index decides the answer. It prices today's barrier (wave_makespan) against a continuously fed pool (fed_makespan) and reports both, and it charges the proposer session — one per generation, measured at 344.7s — which it had been omitting entirely. Measurement only; no runtime behaviour changes. Co-Authored-By: Claude Opus 5 (1M context) Co-Authored-By: Claude Opus 5 (1M context) --- eval/tests/test_measure_evolution_cost.py | 89 +++--- eval/workflow_bench/measure_evolution_cost.py | 261 +++++++++++------- eval/workflow_bench/session_durations.json | 50 ++++ 3 files changed, 258 insertions(+), 142 deletions(-) create mode 100644 eval/workflow_bench/session_durations.json diff --git a/eval/tests/test_measure_evolution_cost.py b/eval/tests/test_measure_evolution_cost.py index 07f81baa3..21e7445b3 100644 --- a/eval/tests/test_measure_evolution_cost.py +++ b/eval/tests/test_measure_evolution_cost.py @@ -1,19 +1,24 @@ -"""Cost model for the evolution wall clock: waves, not cells; overlap, not wishes.""" +"""Cost model for the evolution wall clock: measured cells, real schedules.""" from __future__ import annotations +import pytest + from workflow_bench.measure_evolution_cost import ( CELL_COPY_SECONDS, + CELL_DURATIONS, GRAPH_ANALYZE_SECONDS, - MEAN_SESSION_SECONDS, + PROPOSER_SECONDS, TEMPLATE_SANITIZE_SECONDS, - cell_setup_seconds, + cell_durations, + expected_makespan, + fed_makespan, + fed_pool_enabled, graph_pipeline_enabled, paid_cells_per_task, - pipelined_wall_seconds, - session_wall_seconds, setup_wall_seconds, unique_paid_shas, + wave_makespan, ) TASKS = [ @@ -23,48 +28,62 @@ TASKS = [ ] +def test_measured_sample_is_present_and_unsorted(): + # Sorting would hand each task a uniform block and hide the variance the + # whole model exists to price. + assert len(CELL_DURATIONS) >= 20 + assert list(CELL_DURATIONS) != sorted(CELL_DURATIONS) + assert PROPOSER_SECONDS > 0 + + def test_weekly_reuse_pays_only_candidate_cells(): assert paid_cells_per_task(TASKS, runs=3, weekly=True, reuse_enabled=True) == [3, 3, 3] assert paid_cells_per_task(TASKS, runs=3, weekly=False, reuse_enabled=True) == [9, 9, 9] assert paid_cells_per_task(TASKS, runs=3, weekly=True, reuse_enabled=False) == [9, 9, 9] -def test_cell_clones_are_charged_once_per_wave(): - # run_cell clones inside its own pool worker: three siblings clone at once. - assert cell_setup_seconds(9, 3, True) == 3 * CELL_COPY_SECONDS - assert cell_setup_seconds(9, 1, True) == 9 * CELL_COPY_SECONDS - assert cell_setup_seconds(0, 3, True) == 0 - # Without templates the cell pays a full clone+sanitize, still per wave. - assert cell_setup_seconds(9, 3, False) == 3 * TEMPLATE_SANITIZE_SECONDS +def test_every_cell_carries_its_own_clone(): + durations = cell_durations(3, True) + assert durations == [d + CELL_COPY_SECONDS for d in CELL_DURATIONS[:3]] + assert cell_durations(2, False) == [d + TEMPLATE_SANITIZE_SECONDS for d in CELL_DURATIONS[:2]] + # Cycling wraps, so a task can be longer than the sample. + assert len(cell_durations(len(CELL_DURATIONS) + 5, True)) == len(CELL_DURATIONS) + 5 -def test_setup_counts_each_sha_once_and_each_task_wave_once(): +def test_a_wave_costs_its_slowest_cell_and_a_fed_pool_does_not(): + slow = [10.0, 1.0, 1.0, 10.0, 1.0, 1.0] + assert wave_makespan(slow, 3) == 20.0 + # Fed: one worker takes the first 10; the second 10 lands on a worker that + # has already cleared a 1, and the remaining 1s fill the third. + assert fed_makespan(slow, 3) == 11.0 + assert fed_makespan(slow, 1) == wave_makespan(slow, 1) == 24.0 + + +def test_expected_makespan_is_rotation_averaged_and_deterministic(): + first = expected_makespan(9, 3, fed_pool=False, clone_templates_enabled=True) + assert first == expected_makespan(9, 3, fed_pool=False, clone_templates_enabled=True) + assert expected_makespan(0, 3, fed_pool=False, clone_templates_enabled=True) == 0.0 + # The barrier can only cost time, never save it. + assert first >= expected_makespan(9, 3, fed_pool=True, clone_templates_enabled=True) + + +def test_setup_counts_each_sha_once_and_no_longer_charges_cells(): assert unique_paid_shas(TASKS, [9, 9, 9]) == 2 assert unique_paid_shas(TASKS, [9, 0, 0]) == 1 assert setup_wall_seconds( - unique_shas=2, - paid_per_task=[9, 9, 9], - workers=3, - clone_templates_enabled=True, - ) == 2 * (TEMPLATE_SANITIZE_SECONDS + GRAPH_ANALYZE_SECONDS) + 9 * CELL_COPY_SECONDS + unique_shas=2, paid_per_task=[9, 9, 9], clone_templates_enabled=True + ) == 2 * (TEMPLATE_SANITIZE_SECONDS + GRAPH_ANALYZE_SECONDS) + assert setup_wall_seconds(unique_shas=2, paid_per_task=[0], clone_templates_enabled=True) == 0 -def test_pipelining_hides_every_sha_but_the_first(): - paid = [9, 9, 9] - waves = [3 * MEAN_SESSION_SECONDS] * 3 - serial = session_wall_seconds(paid, 3) + setup_wall_seconds( - unique_shas=2, - paid_per_task=paid, - workers=3, - clone_templates_enabled=True, - ) - pipelined = pipelined_wall_seconds( - TASKS, paid, waves, workers=3, clone_templates_enabled=True - ) - # sha-2's setup fits inside task a's session wave; sha-1 is already ready. - assert serial - pipelined == TEMPLATE_SANITIZE_SECONDS + GRAPH_ANALYZE_SECONDS - - -def test_pipeline_flag_reads_the_runner_not_the_wish(): +def test_feature_flags_read_the_runner_not_the_wish(): assert graph_pipeline_enabled("def _run_sweep(): pass") == 0 assert graph_pipeline_enabled("graph_prefetch = GraphPrefetch(...)") == 1 + assert fed_pool_enabled("def _run_wave(): pass") == 0 + assert fed_pool_enabled("def _run_fed_pool(): pass") == 1 + + +@pytest.mark.parametrize("workers", [1, 3, 8]) +def test_more_workers_never_lengthen_a_task(workers): + serial = expected_makespan(9, 1, fed_pool=True, clone_templates_enabled=True) + assert expected_makespan(9, workers, fed_pool=True, clone_templates_enabled=True) <= serial diff --git a/eval/workflow_bench/measure_evolution_cost.py b/eval/workflow_bench/measure_evolution_cost.py index 7bd9daa20..5448bdc58 100644 --- a/eval/workflow_bench/measure_evolution_cost.py +++ b/eval/workflow_bench/measure_evolution_cost.py @@ -16,13 +16,22 @@ from __future__ import annotations import json import math import re +import statistics as st import subprocess import sys from pathlib import Path -# Agent-session mean from Actions run 33962002890 (review profile). The -# runner builds graphs before cells start, so setup is added separately. -MEAN_SESSION_SECONDS = 1140 +# Measured cell durations, not a mean: a wave waits for its SLOWEST cell, and +# these run 826s at the median against a 5400s session ceiling, so a model +# built on an average understates every concurrent schedule. See +# session_durations.json for provenance and its sampling caveat. +DURATIONS = json.loads( + (Path(__file__).resolve().parent / "session_durations.json").read_text(encoding="utf-8") +) +CELL_DURATIONS = tuple(DURATIONS["cell_duration_s"]) +# One proposer session per generation, ahead of the benchmark and unavoidably +# on the critical path. +PROPOSER_SECONDS = DURATIONS["proposer_duration_s"] # Conservative mean for `analyze --pdg --index-only` of a sanitized # GitNexus snapshot on the evolution box. The hard timeout is 3600s. GRAPH_ANALYZE_SECONDS = 600 @@ -151,36 +160,17 @@ def sha_setup_seconds(clone_templates_enabled: bool) -> int: return GRAPH_ANALYZE_SECONDS -def cell_setup_seconds(paid_cells: int, workers: int, clone_templates_enabled: bool) -> int: - """Per-cell clone cost, charged once per wave rather than once per cell. - - ``run_cell`` clones inside its own pool worker, so the siblings in a wave - clone concurrently and only one clone sits on the critical path per wave. - """ - - if paid_cells < 1: - return 0 - waves = math.ceil(paid_cells / workers) - if clone_templates_enabled: - return waves * CELL_COPY_SECONDS - return waves * TEMPLATE_SANITIZE_SECONDS - - def setup_wall_seconds( *, unique_shas: int, paid_per_task: list[int], - workers: int, clone_templates_enabled: bool, ) -> int: - """Fully serial SHA setup plus per-wave clone work, task by task.""" + """Serial per-SHA graph setup. Per-cell clones ride inside the cell.""" if unique_shas < 1 or sum(paid_per_task) < 1: return 0 - cells = sum( - cell_setup_seconds(paid, workers, clone_templates_enabled) for paid in paid_per_task - ) - return unique_shas * sha_setup_seconds(clone_templates_enabled) + cells + return unique_shas * sha_setup_seconds(clone_templates_enabled) def graph_pipeline_enabled(runner_text: str) -> int: @@ -199,9 +189,8 @@ def graph_pipeline_enabled(runner_text: str) -> int: def pipelined_wall_seconds( tasks: list[dict[str, str]], paid_per_task: list[int], - session_per_task: list[int], + session_per_task: list[float], *, - workers: int, clone_templates_enabled: bool, ) -> int: """Overlap the next unseen SHA's setup with the current task's session wave. @@ -212,7 +201,7 @@ def pipelined_wall_seconds( """ setup = sha_setup_seconds(clone_templates_enabled) - elapsed = 0 + elapsed = PROPOSER_SECONDS ready: set[str] = set() in_flight: tuple[str, int] | None = None for index, (task, paid, session) in enumerate(zip(tasks, paid_per_task, session_per_task, strict=True)): @@ -229,7 +218,7 @@ def pipelined_wall_seconds( if sha: ready.add(sha) session_start = elapsed - elapsed += session + cell_setup_seconds(paid, workers, clone_templates_enabled) + elapsed += session if in_flight is None: for later_task, later_paid in zip(tasks[index + 1 :], paid_per_task[index + 1 :], strict=True): later_sha = later_task.get("ref", "") @@ -239,18 +228,87 @@ def pipelined_wall_seconds( elif in_flight[1] <= elapsed: ready.add(in_flight[0]) in_flight = None - return elapsed + return round(elapsed) -def session_wall_seconds(paid_per_task: list[int], workers: int) -> int: +def cell_durations(count: int, clone_templates_enabled: bool, offset: int = 0) -> list[float]: + """Deterministic per-cell durations: the measured sample, cycled from ``offset``. + + Cycled rather than sampled so every scheduler is scored against the same + cells and the harness stays reproducible. Each cell carries its own clone, + which ``run_cell`` pays inside its worker. + """ + + clone = CELL_COPY_SECONDS if clone_templates_enabled else TEMPLATE_SANITIZE_SECONDS + size = len(CELL_DURATIONS) + return [CELL_DURATIONS[(offset + i) % size] + clone for i in range(count)] + + +def expected_makespan( + count: int, + workers: int, + *, + fed_pool: bool, + clone_templates_enabled: bool, +) -> float: + """Mean makespan over every rotation of the measured sample. + + One fixed alignment would let an accident of the source run — its slowest + cells happen to come first — decide the answer. Averaging all rotations + keeps the real multiset and the real ordering effects while removing that + alignment artifact, and stays deterministic. + """ + + if count <= 0: + return 0.0 + makespan = fed_makespan if fed_pool else wave_makespan + totals = [ + makespan(cell_durations(count, clone_templates_enabled, offset), workers) + for offset in range(len(CELL_DURATIONS)) + ] + return sum(totals) / len(totals) + + +def wave_makespan(durations: list[float], workers: int) -> float: + """Current scheduler: fixed waves of ``workers``, barrier between them.""" + + return sum( + max(durations[start : start + workers]) for start in range(0, len(durations), workers) + ) + + +def fed_makespan(durations: list[float], workers: int) -> float: + """Continuously fed pool: a free worker takes the next cell immediately.""" + + busy_until = [0.0] * workers + for duration in durations: + first = min(range(workers), key=busy_until.__getitem__) + busy_until[first] += duration + return max(busy_until) + + +def session_wall_seconds( + paid_per_task: list[int], + workers: int, + *, + fed_pool: bool, + clone_templates_enabled: bool, +) -> int: if workers < 1: raise ValueError("workers must be at least 1") - total = 0 + makespan = fed_makespan if fed_pool else wave_makespan + total = 0.0 for paid in paid_per_task: if paid <= 0: continue - total += math.ceil(paid / workers) * MEAN_SESSION_SECONDS - return total + total += makespan(cell_durations(paid, clone_templates_enabled), workers) + return round(total) + + +def fed_pool_enabled(runner_text: str) -> int: + """True only when the sweep feeds a live pool instead of waiting on waves.""" + + return int("def _run_fed_pool" in runner_text) def _pytest_python() -> list[str]: @@ -292,86 +350,75 @@ def suite_passed() -> int: def main() -> int: tasks = review_tasks(_read(REVIEW_TASKS)) evolve = _read(EVOLVE_PY) + runner = _read(RUNNER_PY) runs = evolve_default("runs", evolve) promotion_min_runs = evolve_default("promotion-min-runs", evolve) workers = workflow_dispatch_workers(_read(WORKFLOW)) reuse_enabled, clone_templates_enabled = feature_enabled() - pipeline_enabled = graph_pipeline_enabled(_read(RUNNER_PY)) - weekly_paid = paid_cells_per_task( - tasks, - runs=runs, - weekly=True, - reuse_enabled=bool(reuse_enabled), - ) - cold_paid = paid_cells_per_task( - tasks, - runs=runs, - weekly=False, - reuse_enabled=bool(reuse_enabled), - ) - weekly_sessions = session_wall_seconds(weekly_paid, workers) - cold_sessions = session_wall_seconds(cold_paid, workers) - weekly_session_waves = [ - 0 if paid <= 0 else math.ceil(paid / workers) * MEAN_SESSION_SECONDS for paid in weekly_paid - ] - cold_session_waves = [ - 0 if paid <= 0 else math.ceil(paid / workers) * MEAN_SESSION_SECONDS for paid in cold_paid - ] - weekly_setup = setup_wall_seconds( - unique_shas=unique_paid_shas(tasks, weekly_paid), - paid_per_task=weekly_paid, - workers=workers, - clone_templates_enabled=bool(clone_templates_enabled), - ) - cold_setup = setup_wall_seconds( - unique_shas=unique_paid_shas(tasks, cold_paid), - paid_per_task=cold_paid, - workers=workers, - clone_templates_enabled=bool(clone_templates_enabled), - ) - if pipeline_enabled: - weekly_wall = pipelined_wall_seconds( - tasks, - weekly_paid, - weekly_session_waves, - workers=workers, - clone_templates_enabled=bool(clone_templates_enabled), + pipeline_enabled = graph_pipeline_enabled(runner) + fed_pool = fed_pool_enabled(runner) + templates = bool(clone_templates_enabled) + + payload: dict[str, object] = {} + for label, weekly in (("weekly", True), ("cold", False)): + paid = paid_cells_per_task( + tasks, runs=runs, weekly=weekly, reuse_enabled=bool(reuse_enabled) ) - cold_wall = pipelined_wall_seconds( - tasks, - cold_paid, - cold_session_waves, - workers=workers, - clone_templates_enabled=bool(clone_templates_enabled), + per_task = [ + expected_makespan( + count, workers, fed_pool=bool(fed_pool), clone_templates_enabled=templates + ) + for count in paid + ] + sessions = round(sum(per_task)) + setup = setup_wall_seconds( + unique_shas=unique_paid_shas(tasks, paid), + paid_per_task=paid, + clone_templates_enabled=templates, ) - weekly_setup = weekly_wall - weekly_sessions - cold_setup = cold_wall - cold_sessions - else: - weekly_wall = weekly_sessions + weekly_setup - cold_wall = cold_sessions + cold_setup - payload = { - "estimated_weekly_wall_seconds": weekly_wall, - "estimated_cold_wall_seconds": cold_wall, - "suite_passed": suite_passed(), - "promotion_min_runs": promotion_min_runs, - "review_task_count": len(tasks), - "candidate_cells": len(tasks) * runs, - "paid_weekly_cells": sum(weekly_paid), - "paid_cold_cells": sum(cold_paid), - "workers": workers, - "unique_task_shas": len({task.get("ref", "") for task in tasks if task.get("ref")}), - "reuse_enabled": reuse_enabled, - "clone_templates_enabled": clone_templates_enabled, - "graph_pipeline_enabled": pipeline_enabled, - "mean_session_seconds": MEAN_SESSION_SECONDS, - "session_weekly_seconds": weekly_sessions, - "session_cold_seconds": cold_sessions, - "setup_weekly_seconds": weekly_setup, - "setup_cold_seconds": cold_setup, - "graph_analyze_seconds": GRAPH_ANALYZE_SECONDS, - "template_sanitize_seconds": TEMPLATE_SANITIZE_SECONDS, - } - json.dump(payload, sys.stdout, separators=(",", ":")) + if pipeline_enabled: + wall = pipelined_wall_seconds( + tasks, paid, per_task, clone_templates_enabled=templates + ) + setup = wall - sessions + else: + wall = round(sessions + setup + PROPOSER_SECONDS) + payload[f"estimated_{label}_wall_seconds"] = wall + payload[f"paid_{label}_cells"] = sum(paid) + payload[f"session_{label}_seconds"] = sessions + payload[f"setup_{label}_seconds"] = setup + # What the barrier costs: the same cells under a continuously fed pool. + payload[f"fed_pool_{label}_session_seconds"] = round( + sum( + expected_makespan( + count, workers, fed_pool=True, clone_templates_enabled=templates + ) + for count in paid + ) + ) + + payload.update( + { + "suite_passed": suite_passed(), + "promotion_min_runs": promotion_min_runs, + "review_task_count": len(tasks), + "candidate_cells": len(tasks) * runs, + "workers": workers, + "unique_task_shas": len({t.get("ref", "") for t in tasks if t.get("ref")}), + "reuse_enabled": reuse_enabled, + "clone_templates_enabled": clone_templates_enabled, + "graph_pipeline_enabled": pipeline_enabled, + "fed_pool_enabled": fed_pool, + "measured_cell_count": len(CELL_DURATIONS), + "median_cell_seconds": round(st.median(CELL_DURATIONS)), + "mean_cell_seconds": round(st.mean(CELL_DURATIONS)), + "max_cell_seconds": round(max(CELL_DURATIONS)), + "proposer_seconds": round(PROPOSER_SECONDS), + "graph_analyze_seconds": GRAPH_ANALYZE_SECONDS, + "template_sanitize_seconds": TEMPLATE_SANITIZE_SECONDS, + } + ) + json.dump(payload, sys.stdout, sort_keys=True) sys.stdout.write("\n") return 0 diff --git a/eval/workflow_bench/session_durations.json b/eval/workflow_bench/session_durations.json new file mode 100644 index 000000000..838643655 --- /dev/null +++ b/eval/workflow_bench/session_durations.json @@ -0,0 +1,50 @@ +{ + "_provenance": "Actions run 33912693948 (2026-09-04), review profile, workers=1, gen-0. Artifact gitnexus-evolution-33912693948-1, gen-0/bench/results.jsonl.", + "_caveat": "Every cell in that run was unusable evidence (32 review-evidence-invalid, 6 session-error, 3 skill-not-invoked) and two hit the 5400s session ceiling. Durations are real; a run that resolves cleanly may sit lower.", + "session_ceiling_s": 5400, + "proposer_duration_s": 344.7, + "cell_duration_s": [ + 3744.6, + 5400.0, + 2485.6, + 2140.4, + 2976.4, + 1338.0, + 1418.4, + 1162.8, + 3075.2, + 436.1, + 627.1, + 702.5, + 653.3, + 901.3, + 1240.4, + 1022.3, + 436.9, + 5400.0, + 851.2, + 704.3, + 762.1, + 1191.1, + 963.4, + 342.9, + 902.1, + 631.8, + 826.3, + 847.2, + 627.9, + 489.7, + 502.6, + 663.5, + 675.4, + 991.3, + 662.4, + 337.8, + 1222.7, + 361.3, + 734.3, + 544.0, + 741.1 + ], + "_order": "Submission order from the source run, deliberately unsorted: the model cycles this list across cells, so sorting it would hand every task a uniform block and hide exactly the variance being measured." +}