mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-01 02:01:24 +00:00
* perf(eval): packed sweep scheduler and the harness that measured it
Extracted from the combined skill-evolution branch so it can be reviewed on its
own. Purely additive against main: no existing function changes behaviour, and
sweep_packed_cells has no production caller yet.
sweep_task_cells finishes one task before starting the next and drains a wave
before refilling it, so a task with fewer cells than workers leaves workers
idle and one slow cell stalls its whole wave. sweep_packed_cells feeds every
task's cells through a single pool instead, keeping the breaker's meaning: a
total submission order continued across task boundaries, a folder walking
results in that order, and consecutive systemic failures counted there, so a
doomed run aborts on the same cell it would have under waves.
simulate_sweep.py is what produced the numbers. It drives the real schedulers
with only the paid agent session stubbed, using the measured per-arm durations
in session_durations.json divided by a scale factor. The distribution's shape
is kept deliberately - median 826s against a 5400s ceiling - because that
spread is the entire reason a barrier costs anything, and uniform sleeps would
erase the effect under test. All schedulers consume one identical seeded plan.
Measured at workers=3 against the review corpus, packing is worth about 40% of
a cold sweep, and it is the only change that moves a seeded weekly run at all -
there a task is three cells and a wave is never full. The submission window is
a real trade, measured with failures injected at four positions:
window 3 -> -8% wall, overrun 2 (the wave scheduler's own bound)
window 6 -> -27% wall, overrun 4
window 12 -> -42% wall, overrun 9
window 54 -> -44% wall, overrun 11
Overrun is wasted paid sessions on an aborted sweep. The default multiplier is
2; the curve lives in the constant's comment so raising it is an informed
decision. Contention was measured separately by burning real CPU in
subprocesses under taskset: the advantage holds between -40% and -47% from 24
cores down to an oversubscribed 2, though packing erodes faster than waves do
because packing is what creates the concurrency.
measure_evolution_cost.py is the offline cost model, with no runtime caller. It
reports workers from the workflow's current default, which on this base is 1.
Limits worth stating: sleeping threads do not contend and the duration sample
was itself recorded at workers=1, so the speedups are upper bounds; the ordering
of the schedulers is trustworthy because they were compared under identical
conditions, the magnitudes are not.
562 eval tests pass at this base. The two test_model_gateway.py failures,
test_locked_litellm_translates_messages_to_offline_responses and
test_openai_gateway_never_leaves_proxy_output_on_an_undrained_pipe, fail
identically on origin/main in this environment.
* fix(eval): compare the shipped window and bound the overrun by it
Address PR review feedback (#3206).
run_faithful defaulted its submission window to `workers` while
runner.sweep_packed_cells defaults to `max(workers * PACKED_WINDOW_MULTIPLIER,
workers)`, so every run that named no window compared a prototype queued twice
as tightly as the shipped scheduler and presented it as the production
invariant. The faithful default now reads the same constant. Measured at
workers=3, faithful and production agreed on nothing before and agree exactly
now: breaker overrun 2/1/2 vs 2/4/3 becomes 2/4/3 vs 2/4/3 across the three
failure positions.
The contention sweep hard-coded `window=12` for faithful only, which the
production run never saw - masked at workers=6 where both are 12. Removed, and
the production measurement it was already paying for is now reported as
`production_s` instead of being discarded.
breaker_fidelity checked the overrun against `args.workers`. The bound the
producer actually enforces is `window - 1` cells past the fold pointer, which
is the wave scheduler's own `workers - 1` when window == workers; against the
shipped default of 6 the old predicate reported a failure for an in-bound run.
The window is now passed explicitly, reported in each row, and checked against
its own bound.
--window was parsed and never read. Wired into the schedulers that hold one.
Dropped two unused plan constructions CodeQL flagged, and the `skipped` set in
sweep_packed_cells that nothing reads - the None appended to `submitted` is the
skip representation the fold loop consumes.
Verification: 562 passed, 15 skipped, 2 failed (the two test_model_gateway.py
failures the PR description documents as reproducing on origin/main), ruff
clean.
* fix(eval): carry the cancellation scope into packed cells, reject the args that hang
Address PR review feedback (#3206).
sweep_packed_cells submits from a producer THREAD, and a new thread starts with
an empty context, so `copy_context()` there copied the producer's context rather
than the one cancellation_scope had just bound _CANCELLATION in. Every packed
cell therefore ran with no cancellation event, and run_managed falls back to
_CANCELLATION when none is passed - so a cancelled run's subprocesses would
never have learned about it. sweep_task_cells gets this right for free by
submitting from the thread that entered the scope. Reproduced directly: packed
workers observed [False, False], wave workers [True, True]. The caller's context
is now captured before the producer starts and copied per submission; the new
test fails without the fix.
Three CLI arguments were accepted and then wedged the run:
--scale 0 ZeroDivisionError before any scheduler starts
--graph-seconds -1 hangs: the builder thread dies on a negative
sleep, every scheduler waits on a readiness
event nobody sets
--window 0 (faithful) hangs: submitted - fold_pointer >= 0 holds
before the first submission, so the producer
and the consumer wait on each other
The first two are rejected at the parser, which is the only layer that runs
before a thread exists. run_faithful now enforces the same window >= workers
rule sweep_packed_cells already had, so the prototype rejects exactly what the
shipped function rejects. All three were confirmed to crash or hang first.
Verification: 563 passed, 15 skipped, 2 failed (the two test_model_gateway.py
failures the PR description documents as reproducing on origin/main), ruff
clean.
* chore(autofix): apply prettier + eslint fixes via /autofix command
* Address PR review feedback (#3206)
Preserve settled sibling rows when a packed cell raises. run_cell deliberately
lets unexpected harness exceptions propagate, and sweep_task_cells answers that
by folding every non-failing sibling before it re-raises - the cells already ran
and already spent their budget, so dropping their rows means paying for evidence
the sweep then discards. sweep_packed_cells called future.result() bare, so the
fold stopped at the failing index and every later cell that had already
completed was silently lost. It now folds forward over the settled futures
before re-raising. The failing index itself has no row, since execute() assigns
only on success, so folding forward cannot duplicate it.
Pinned by a regression test that fails without the fix: the later cell is made
to finish first, so there is real settled evidence to lose at the moment cell 0
raises.
Reject arguments that cannot produce a run, at the boundary rather than deep
inside a thread. NaN defeats every comparison it appears in, so the existing
"> 0" and ">= 0" checks admitted --scale nan and --graph-seconds nan; the NaN
then reached time.sleep in a worker or the graph thread, raised there, and left
every scheduler waiting forever on a readiness event nobody would set. Infinity
was worse than a crash: it scaled all durations to zero and the run reported a
sweep that took no time. Both flags now require a finite value.
The count flags are indexed or handed straight to a thread pool, so a zero
surfaced as an IndexError on plans[0], a median over an empty sequence, or
ThreadPoolExecutor's own error - none naming the flag responsible. --workers,
--repeat and --runs now require at least 1.
Two flags were not in the review but carry the same invariant and the same
one-line treatment, so they are fixed with the class rather than left to
resurface: --runs (same empty-plan path as --repeat) and --window, where zero
admits no cell at all because the producer waits for a fold pointer to move past
a cell it was never allowed to submit.
Verified each guard fires with its own message rather than a stack trace.
563 eval tests pass. Note: pre-existing failures in test_model_gateway.py not
addressed by this PR - litellm[proxy]'s console script is absent in this
environment, and neither test touches the files changed here.
* Address PR review feedback (#3206), round 2
Stop charging the fed baseline for overlap the wave scheduler gets free.
run_fed is documented as pricing the barrier alone, but it slept graph_seconds
serially before every task, while run_wave starts one background builder that
prepares task N+1 while task N's cells run. The fed-versus-wave delta therefore
mixed the loss of that overlap into what was reported as the price of the
barrier. run_fed now uses the same builder, started before the clock, so the
barrier is the only remaining difference.
This moved the numbers. On the weekly profile fed was 4.203s and is now 3.694s,
exactly equal to wave - which is the answer that profile should give. On cold,
fed was 5.995s and is now 5.487s, so the measured price of the barrier widens
from 1.844s to 2.352s: the old arrangement understated it by about a quarter.
No committed results file or PR-body figure quotes these, so there is nothing
stale to regenerate.
Enforce the window bound the schedulers actually hold. Last round's guard
required only >= 1, but run_faithful and sweep_packed_cells both refuse a window
below the worker count, so --scheduler faithful --workers 3 --window 1 passed
validation and then died on an uncaught ValueError. The check now uses the
worker count.
It also uses the LARGEST worker count the invocation will really use.
--contention-sweep runs its own counts irrespective of --workers, so validating
against --workers alone let the three-worker measurements finish and then raised
on the six-worker one, losing the run partway through. Those counts are now a
named constant the validator can see.
Verified: --scheduler faithful --workers 3 --window 1 is rejected naming 3, and
--workers 3 --window 3 --contention-sweep is rejected naming 6.
No regression test for the graph-overlap fix. Discriminating it from the old
behaviour requires cell work to overlap graph work, which makes the assertion a
timing comparison, and this project does not take non-deterministic tests. It is
verified by the before/after measurement above instead.
564 eval tests pass. Note: pre-existing failures in test_model_gateway.py not
addressed by this PR - litellm[proxy]'s console script is absent here.
---------
Co-authored-by: Gergo Magyar <gergomagyar0@gmail.com>
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com>
311 lines
12 KiB
Python
311 lines
12 KiB
Python
#!/usr/bin/env python3
|
|
"""Cheap cost model for the skill-evolution review generation.
|
|
|
|
This is the ce-optimize measurement harness. It does not start Claude and it
|
|
does not replay a run. It reads the review corpus, the evolve defaults and the
|
|
workflow's workers default, then schedules the measured cell durations in
|
|
``session_durations.json`` the way ``sweep_task_cells`` schedules real cells.
|
|
|
|
Everything priced here is measured. Cell durations and the proposer session
|
|
come from a real artifact, and the work outside the agent sessions comes from
|
|
that run's own step wall minus the time its sessions and proposer account for.
|
|
|
|
Weekly assumes a matching seed, so every reusable comparator cell is skipped
|
|
and only the candidate arm is paid. Cold assumes an empty seed.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import math
|
|
import re
|
|
import statistics as st
|
|
import subprocess
|
|
import sys
|
|
from pathlib import Path
|
|
|
|
REPO_ROOT = Path(__file__).resolve().parents[2]
|
|
EVAL_ROOT = REPO_ROOT / "eval"
|
|
REVIEW_TASKS = EVAL_ROOT / "workflow_bench" / "tasks.review.scenarios.yaml"
|
|
EVOLVE_PY = EVAL_ROOT / "workflow_bench" / "evolve.py"
|
|
RUNNER_PY = EVAL_ROOT / "workflow_bench" / "runner.py"
|
|
ARTIFACTS_PY = EVAL_ROOT / "workflow_bench" / "runner_artifacts.py"
|
|
REUSE_PY = EVAL_ROOT / "workflow_bench" / "comparator_reuse.py"
|
|
WORKFLOW = REPO_ROOT / ".github" / "workflows" / "gitnexus-skill-evolution.yml"
|
|
|
|
MEASURED = json.loads(
|
|
(Path(__file__).resolve().parent / "session_durations.json").read_text(encoding="utf-8")
|
|
)
|
|
# Per arm, because the arms are not interchangeable and the weekly lane pays
|
|
# only the candidate one. Cells are submitted run-major and arm-minor
|
|
# (runner.py ``planned``), so at workers=3 every wave holds one cell of each
|
|
# arm and the slowest arm sets the wave.
|
|
DURATIONS_BY_ARM: dict[str, tuple[float, ...]] = {
|
|
arm: tuple(values) for arm, values in MEASURED["cell_duration_s_by_arm"].items()
|
|
}
|
|
PROPOSER_SECONDS: float = MEASURED["proposer_duration_s"]
|
|
_RESIDUAL = MEASURED["residual"]
|
|
# Clone, graph build, sandbox, teardown: the sweep's own time, taken as that
|
|
# run's step wall minus what its sessions and proposer account for. Charged
|
|
# SERIALLY, outside the pool, and charged PER SHA rather than per cell. The
|
|
# residual mixes per-cell work with per-SHA graph setup and the artifact cannot
|
|
# separate them; per-SHA is the direction that refuses to credit a run for
|
|
# shrinking work it still performs, which per-cell did - a weekly generation
|
|
# pays one arm instead of three but builds exactly the same graphs. See
|
|
# session_durations.json residual._split_assumption.
|
|
SHA_OVERHEAD_SECONDS: float = _RESIDUAL["sha_overhead_s"]
|
|
|
|
# runner.py CANDIDATE_ARMS derives the candidate arm from its incumbent, and
|
|
# only an incumbent row can be reused from a prior generation.
|
|
CANDIDATE_ARM = "candidate_review"
|
|
REVIEW_ARMS = ("ce_review", "review", CANDIDATE_ARM)
|
|
|
|
SUITE_FILES = (
|
|
"tests/test_measure_evolution_cost.py",
|
|
"tests/test_comparator_reuse.py",
|
|
"tests/test_evolve.py",
|
|
"tests/test_sanitized_graph.py",
|
|
"tests/test_workflow_bench.py",
|
|
"tests/test_workflow_bench_sessions.py",
|
|
"tests/test_session_progress.py",
|
|
)
|
|
|
|
|
|
def _read(path: Path) -> str:
|
|
return path.read_text(encoding="utf-8")
|
|
|
|
|
|
def review_tasks(text: str) -> list[dict[str, str]]:
|
|
tasks: list[dict[str, str]] = []
|
|
current: dict[str, str] | None = None
|
|
for raw in text.splitlines():
|
|
line = raw.strip()
|
|
if line.startswith("id:"):
|
|
if current is not None:
|
|
tasks.append(current)
|
|
current = {"id": line.split(":", 1)[1].strip()}
|
|
elif line.startswith("ref:") and current is not None:
|
|
current["ref"] = line.split(":", 1)[1].strip()
|
|
if current is not None:
|
|
tasks.append(current)
|
|
return tasks
|
|
|
|
|
|
def evolve_default(name: str, text: str) -> int:
|
|
match = re.search(rf'add_argument\("--{re.escape(name)}".*?default=(\d+)', text, flags=re.S)
|
|
if match is None:
|
|
raise ValueError(f"evolve.py is missing --{name} default")
|
|
return int(match.group(1))
|
|
|
|
|
|
def workflow_dispatch_workers(text: str) -> int:
|
|
match = re.search(r"^\s+workers:\n(?:.*\n)*?^\s+default: '(\d+)'", text, flags=re.M)
|
|
if match is None:
|
|
raise ValueError("workflow_dispatch workers default is missing")
|
|
return int(match.group(1))
|
|
|
|
|
|
def feature_enabled() -> tuple[int, int]:
|
|
evolve = _read(EVOLVE_PY)
|
|
runner = _read(RUNNER_PY)
|
|
artifacts = _read(ARTIFACTS_PY)
|
|
reuse = int(
|
|
REUSE_PY.is_file()
|
|
and "--reuse-results" in evolve
|
|
and "select_reusable_comparator_rows" in runner
|
|
and "CANDIDATE" in _read(REUSE_PY)
|
|
)
|
|
templates = int("def copy_isolated_tree" in artifacts and "clone_templates" in runner)
|
|
return reuse, templates
|
|
|
|
|
|
def graph_pipeline_enabled(runner_text: str) -> int:
|
|
"""True when the runner prefetches the next SHA during paid sessions."""
|
|
|
|
return int("prefetch_next_graph" in runner_text or "GraphPrefetch" in runner_text)
|
|
|
|
|
|
def fed_pool_enabled(runner_text: str) -> int:
|
|
"""True when the sweep feeds a live pool instead of waiting on waves."""
|
|
|
|
return int("def _run_fed_pool" in runner_text)
|
|
|
|
|
|
def paid_arms(weekly: bool, reuse_enabled: bool) -> tuple[str, ...]:
|
|
"""Arms a generation actually pays for."""
|
|
|
|
if weekly and reuse_enabled:
|
|
return (CANDIDATE_ARM,)
|
|
return REVIEW_ARMS
|
|
|
|
|
|
def task_cells(runs: int, arms: tuple[str, ...], offset: int) -> list[float]:
|
|
"""One task's cell durations in submission order: run-major, arm-minor.
|
|
|
|
Each arm draws from its own measured sample, cycled from ``offset`` so the
|
|
caller can average over every alignment instead of trusting one.
|
|
"""
|
|
|
|
cells: list[float] = []
|
|
for run_idx in range(runs):
|
|
for arm in arms:
|
|
sample = DURATIONS_BY_ARM[arm]
|
|
cells.append(sample[(offset + run_idx) % len(sample)])
|
|
return cells
|
|
|
|
|
|
def wave_makespan(durations: list[float], workers: int) -> float:
|
|
"""Today's scheduler: fixed waves of ``workers``, with a barrier between."""
|
|
|
|
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 expected_task_seconds(
|
|
runs: int, arms: tuple[str, ...], workers: int, *, fed_pool: bool
|
|
) -> float:
|
|
"""Mean makespan of one task over every alignment of the measured samples.
|
|
|
|
One fixed alignment would let an accident of the source run - its slowest
|
|
cells happen to come first - decide the answer. Averaging keeps the real
|
|
multiset and the real ordering effects without that artifact, and stays
|
|
deterministic.
|
|
"""
|
|
|
|
if runs < 1 or not arms:
|
|
return 0.0
|
|
makespan = fed_makespan if fed_pool else wave_makespan
|
|
# lcm, not max: with samples of 13 and 14, max would wrap the shorter one
|
|
# and count its first entry twice.
|
|
alignments = math.lcm(*(len(DURATIONS_BY_ARM[arm]) for arm in arms))
|
|
return (
|
|
sum(makespan(task_cells(runs, arms, offset), workers) for offset in range(alignments))
|
|
/ alignments
|
|
)
|
|
|
|
|
|
def generation_seconds(
|
|
*,
|
|
task_count: int,
|
|
runs: int,
|
|
arms: tuple[str, ...],
|
|
workers: int,
|
|
fed_pool: bool,
|
|
unique_shas: int,
|
|
) -> int:
|
|
"""Whole generation: proposer, then the tasks back to back, plus overhead.
|
|
|
|
Prices a HEALTHY sweep. A run whose cells return unusable evidence does not
|
|
reach this wall at all: the outage breaker aborts after
|
|
``DEFAULT_OUTAGE_STREAK`` consecutive systemic failures, which for the
|
|
sample's own error sequence is cell 5 of 41.
|
|
|
|
Sweep overhead is charged per SHA, so it does not shrink with the arm count.
|
|
Weekly pays one arm instead of three but builds the same graphs, and billing
|
|
that per cell credited it for a saving the real run never makes.
|
|
"""
|
|
|
|
return round(
|
|
PROPOSER_SECONDS
|
|
+ task_count * expected_task_seconds(runs, arms, workers, fed_pool=fed_pool)
|
|
+ unique_shas * SHA_OVERHEAD_SECONDS
|
|
)
|
|
|
|
|
|
def _pytest_python() -> list[str]:
|
|
venv_python = EVAL_ROOT / ".venv" / "bin" / "python"
|
|
if venv_python.is_file():
|
|
return [str(venv_python)]
|
|
if (EVAL_ROOT / "uv.lock").is_file():
|
|
return ["uv", "run", "--locked", "--extra", "dev", "python"]
|
|
return [sys.executable]
|
|
|
|
|
|
def suite_passed() -> int:
|
|
files = [name for name in SUITE_FILES if (EVAL_ROOT / name).is_file()]
|
|
if not files:
|
|
return 0
|
|
cmd = [*_pytest_python(), "-m", "pytest", *files, "-q", "--tb=no", "--no-header"]
|
|
try:
|
|
completed = subprocess.run(
|
|
cmd, cwd=EVAL_ROOT, check=False, capture_output=True, text=True, timeout=240
|
|
)
|
|
except (OSError, subprocess.TimeoutExpired):
|
|
return 0
|
|
return int(completed.returncode == 0)
|
|
|
|
|
|
def main() -> int:
|
|
tasks = review_tasks(_read(REVIEW_TASKS))
|
|
evolve = _read(EVOLVE_PY)
|
|
runner = _read(RUNNER_PY)
|
|
runs = evolve_default("runs", evolve)
|
|
workers = workflow_dispatch_workers(_read(WORKFLOW))
|
|
reuse_enabled, clone_templates_enabled = feature_enabled()
|
|
fed_pool = fed_pool_enabled(runner)
|
|
|
|
# Both walls build the same graphs; the arm count does not change that.
|
|
unique_shas = len({t.get("ref", "") for t in tasks if t.get("ref")})
|
|
payload: dict[str, object] = {}
|
|
for label, weekly in (("weekly", True), ("cold", False)):
|
|
arms = paid_arms(weekly, bool(reuse_enabled))
|
|
payload[f"estimated_{label}_wall_seconds"] = generation_seconds(
|
|
task_count=len(tasks),
|
|
runs=runs,
|
|
arms=arms,
|
|
workers=workers,
|
|
fed_pool=bool(fed_pool),
|
|
unique_shas=unique_shas,
|
|
)
|
|
payload[f"paid_{label}_cells"] = len(tasks) * runs * len(arms)
|
|
# What the wave barrier costs: the same cells, continuously fed.
|
|
payload[f"fed_pool_{label}_wall_seconds"] = generation_seconds(
|
|
task_count=len(tasks),
|
|
runs=runs,
|
|
arms=arms,
|
|
workers=workers,
|
|
fed_pool=True,
|
|
unique_shas=unique_shas,
|
|
)
|
|
|
|
all_durations = [d for sample in DURATIONS_BY_ARM.values() for d in sample]
|
|
payload.update(
|
|
{
|
|
"suite_passed": suite_passed(),
|
|
"promotion_min_runs": evolve_default("promotion-min-runs", evolve),
|
|
"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": graph_pipeline_enabled(runner),
|
|
"fed_pool_enabled": fed_pool,
|
|
"measured_cell_count": len(all_durations),
|
|
"median_cell_seconds": round(st.median(all_durations)),
|
|
"mean_cell_seconds": round(st.mean(all_durations)),
|
|
"max_cell_seconds": round(max(all_durations)),
|
|
"median_candidate_cell_seconds": round(st.median(DURATIONS_BY_ARM[CANDIDATE_ARM])),
|
|
"mean_candidate_cell_seconds": round(st.mean(DURATIONS_BY_ARM[CANDIDATE_ARM])),
|
|
"proposer_seconds": round(PROPOSER_SECONDS),
|
|
"sha_overhead_seconds": round(SHA_OVERHEAD_SECONDS, 1),
|
|
}
|
|
)
|
|
json.dump(payload, sys.stdout, sort_keys=True)
|
|
sys.stdout.write("\n")
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|