From 2ab88b72177bd8bc3b04949e7080cf826882ea45 Mon Sep 17 00:00:00 2001 From: Gergo Magyar Date: Mon, 7 Sep 2026 06:49:50 +0000 Subject: [PATCH] test(eval): measure the sweep scheduler instead of modelling it measure_evolution_cost predicts wall clock from a model of what sweep_task_cells does. This runs the real thing - real threads, the real wave barrier, the real outage breaker - with only the paid agent session replaced by a sleep, and times it. Durations are the measured per-arm samples divided by 5000, so a 1416s cell takes ~0.28s. The shape is kept on purpose: the median cell is 826s against a 5400s ceiling, and that spread is the entire reason a barrier costs anything. Uniform random sleeps would erase the effect under test. All schedulers consume one identical seeded plan, so a comparison cannot be an artifact of one of them drawing luckier cells. The model survives contact: it tracks real execution within about 10%, and workers=1 - which runs without a pool at all - sits at 0.95, so the residual above 1.0 at higher worker counts is per-wave thread overhead rather than a modelling error. Two structural claims that were arithmetic are now observed. Weekly is flat from workers=3: 3.59, 3.59, 3.59, 3.60, 3.59, 3.59 across w=3..8. workers=4 buys nothing over workers=3 on cold, 7.68 against 7.78. Two prototype schedulers are measured beside it, deliberately before any production code exists. A continuously fed pool per task is worth more than the model claimed on cold, -27.3% against a predicted -17.9%, and exactly nothing on weekly, +0.0%, because a weekly task is one wave with nothing to feed. One pool across all tasks beats both: -40.7% weekly and -42.9% cold at workers=3, rising to -65.7% and -63.9% at workers=8. It also subsumes the fed pool, since packing across tasks is a fed pool. That reorders the backlog. Cross-task packing moves from second to first: it dominates on both profiles, and it is the only thing that moves weekly at all. Raising the worker count is worth nothing until it lands - under the barrier weekly does not improve from w=3 to w=8, and speedup against serial is 1.58x for three workers and only 2.40x for eight. The bound on all of it: sleeping threads do not contend. Real sandboxed sessions compete for CPU, page cache and disk, and the duration sample was itself measured at workers=1, so it carries no contention either. These speedups are upper bounds. The ordering is trustworthy because the schedulers were compared under identical conditions; the magnitudes are not. The packed prototype is also a bare ThreadPoolExecutor with no breaker folding, no per-task graph lifecycle and no reuse binding - which is the actual cost of building it, and is not measured here. --- eval/workflow_bench/simulate_sweep.py | 193 ++++++++++++++++++++++++++ 1 file changed, 193 insertions(+) create mode 100644 eval/workflow_bench/simulate_sweep.py diff --git a/eval/workflow_bench/simulate_sweep.py b/eval/workflow_bench/simulate_sweep.py new file mode 100644 index 000000000..dbd251f8f --- /dev/null +++ b/eval/workflow_bench/simulate_sweep.py @@ -0,0 +1,193 @@ +#!/usr/bin/env python3 +"""Run the real sweep scheduler against stub sessions and time it. + +``measure_evolution_cost`` is arithmetic: it predicts wall clock from a model of +what ``sweep_task_cells`` does. This runs the actual function - real threads, +the real wave barrier, the real outage breaker - and replaces only the paid +agent session with a sleep. If the two disagree, the model is wrong. + +Durations are the measured per-arm samples from ``session_durations.json`` +divided by ``--scale``, so a cell that really took 1416s takes ~0.28s here. The +shape is preserved deliberately: the median cell is 826s against a 5400s +ceiling, and that spread is the whole reason a barrier costs anything. Uniform +random sleeps would erase the effect under test. + +Three schedulers run against an IDENTICAL seeded duration sequence: + +``wave`` the shipped ``sweep_task_cells`` - fixed waves of ``workers``, a + barrier between them, one task at a time. +``fed`` a continuously fed pool per task: a free worker takes the next cell + immediately instead of waiting for its wave to drain (H1). +``packed`` one pool across every task, so a task's leftover capacity is filled + by the next task's cells (H2). + +``fed`` and ``packed`` are measured here as prototypes, deliberately, before any +production code is written - the point is to find out whether the idea is worth +the invariants it would cost. + + python3 -m workflow_bench.simulate_sweep --workers 3 + python3 -m workflow_bench.simulate_sweep --compare --repeat 5 +""" + +from __future__ import annotations + +import argparse +import json +import random +import statistics +import threading +import time +from concurrent.futures import ThreadPoolExecutor +from typing import Any + +from . import runner +from .measure_evolution_cost import ( + CANDIDATE_ARM, + DURATIONS_BY_ARM, + REVIEW_ARMS, + REVIEW_TASKS, + _read, + expected_task_seconds, + review_tasks, +) + +DEFAULT_SCALE = 5000.0 +Cell = tuple[int, str, float] + + +def build_plan( + *, task_count: int, runs: int, arms: tuple[str, ...], scale: float, seed: int +) -> list[list[Cell]]: + """Per-task cells in submission order (run-major, arm-minor) with durations. + + Generated once and shared by every scheduler so a comparison cannot be an + artifact of one of them drawing luckier cells. + """ + + rng = random.Random(seed) + plan: list[list[Cell]] = [] + for _task in range(task_count): + cells: list[Cell] = [] + for run_idx in range(runs): + for arm in arms: + sample = DURATIONS_BY_ARM[arm] + cells.append((run_idx, arm, sample[rng.randrange(len(sample))] / scale)) + plan.append(cells) + return plan + + +def _record(run_idx: int, arm: str) -> dict[str, Any]: + return { + "run": run_idx, + "arm": arm, + "ok": True, + "resolved": True, + "error_kind": None, + "review_evidence_valid": True, + } + + +def run_wave(plan: list[list[Cell]], workers: int) -> float: + """The shipped scheduler, driven for real.""" + + started = time.monotonic() + for cells in plan: + by_key = {(run_idx, arm): seconds for run_idx, arm, seconds in cells} + + def fake_run(run_idx: int, arm: str) -> dict[str, Any]: + time.sleep(by_key[(run_idx, arm)]) + return _record(run_idx, arm) + + _streak, tripped = runner.sweep_task_cells( + [(run_idx, arm) for run_idx, arm, _ in cells], + workers=workers, + run=fake_run, + on_start=lambda *_: None, + on_record=lambda *_: None, + outage_streak=0, + outage_limit=0, + ) + assert not tripped + return time.monotonic() - started + + +def _drain(cells: list[Cell], workers: int) -> None: + lock = threading.Lock() + order: list[tuple[int, str]] = [] + + def work(cell: Cell) -> None: + run_idx, arm, seconds = cell + time.sleep(seconds) + with lock: + order.append((run_idx, arm)) + + with ThreadPoolExecutor(max_workers=workers) as pool: + list(pool.map(work, cells)) + + +def run_fed(plan: list[list[Cell]], workers: int) -> float: + """H1: continuously fed pool, still one task at a time.""" + + started = time.monotonic() + for cells in plan: + _drain(cells, workers) + return time.monotonic() - started + + +def run_packed(plan: list[list[Cell]], workers: int) -> float: + """H2: one pool across every task.""" + + started = time.monotonic() + _drain([cell for cells in plan for cell in cells], workers) + return time.monotonic() - started + + +SCHEDULERS = {"wave": run_wave, "fed": run_fed, "packed": run_packed} + + +def main() -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--workers", type=int, default=3) + parser.add_argument("--scale", type=float, default=DEFAULT_SCALE) + parser.add_argument("--seed", type=int, default=1729) + parser.add_argument("--repeat", type=int, default=1) + parser.add_argument("--runs", type=int, default=3) + parser.add_argument("--scheduler", choices=sorted(SCHEDULERS), default="wave") + parser.add_argument("--compare", action="store_true", help="all schedulers, both profiles") + args = parser.parse_args() + + task_count = len(review_tasks(_read(REVIEW_TASKS))) + names = sorted(SCHEDULERS) if args.compare else [args.scheduler] + rows: list[dict[str, Any]] = [] + for label, weekly in (("weekly", True), ("cold", False)): + arms = (CANDIDATE_ARM,) if weekly else REVIEW_ARMS + plans = [ + build_plan( + task_count=task_count, runs=args.runs, arms=arms, scale=args.scale, seed=args.seed + i + ) + for i in range(args.repeat) + ] + serial = statistics.median(sum(c[2] for cells in p for c in cells) for p in plans) + predicted = task_count * expected_task_seconds( + args.runs, arms, args.workers, fed_pool=False + ) / args.scale + for name in names: + observed = statistics.median(SCHEDULERS[name](p, args.workers) for p in plans) + rows.append( + { + "profile": label, + "scheduler": name, + "workers": args.workers, + "observed_s": round(observed, 3), + "wave_model_s": round(predicted, 3), + "serial_s": round(serial, 3), + "speedup_vs_serial": round(serial / observed, 3) if observed else None, + "cells": sum(len(cells) for cells in plans[0]), + } + ) + print(json.dumps({"scale": args.scale, "repeat": args.repeat, "rows": rows}, indent=2)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main())