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>
122 lines
5.1 KiB
Python
122 lines
5.1 KiB
Python
"""Cost model for the evolution wall clock: measured cells, real schedules."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import pytest
|
|
|
|
from workflow_bench.measure_evolution_cost import (
|
|
CANDIDATE_ARM,
|
|
SHA_OVERHEAD_SECONDS,
|
|
DURATIONS_BY_ARM,
|
|
PROPOSER_SECONDS,
|
|
REVIEW_ARMS,
|
|
expected_task_seconds,
|
|
fed_makespan,
|
|
fed_pool_enabled,
|
|
generation_seconds,
|
|
graph_pipeline_enabled,
|
|
paid_arms,
|
|
task_cells,
|
|
wave_makespan,
|
|
)
|
|
|
|
|
|
def test_every_arm_has_its_own_unsorted_sample():
|
|
assert set(DURATIONS_BY_ARM) == set(REVIEW_ARMS)
|
|
for arm, sample in DURATIONS_BY_ARM.items():
|
|
assert len(sample) >= 10, arm
|
|
# Sorting would hand each task a uniform block and hide the variance
|
|
# the whole model exists to price.
|
|
assert list(sample) != sorted(sample), arm
|
|
assert PROPOSER_SECONDS > 0
|
|
assert SHA_OVERHEAD_SECONDS > 0
|
|
|
|
|
|
def test_weekly_reuse_pays_the_candidate_arm_only():
|
|
assert paid_arms(weekly=True, reuse_enabled=True) == (CANDIDATE_ARM,)
|
|
assert paid_arms(weekly=False, reuse_enabled=True) == REVIEW_ARMS
|
|
assert paid_arms(weekly=True, reuse_enabled=False) == REVIEW_ARMS
|
|
|
|
|
|
def test_cells_are_submitted_run_major_arm_minor():
|
|
# runner.py: [(run_idx, arm) for run_idx in range(runs) for arm in arms].
|
|
# At workers=3 that puts one cell of each arm in every wave.
|
|
cells = task_cells(2, REVIEW_ARMS, 0)
|
|
assert len(cells) == 6
|
|
expected = [DURATIONS_BY_ARM[arm][run] for run in range(2) for arm in REVIEW_ARMS]
|
|
assert cells == expected
|
|
|
|
|
|
def test_overhead_is_charged_per_sha_and_outside_the_pool():
|
|
# Two properties at once: the residual sits outside the schedule, where more
|
|
# workers cannot dissolve it, and it scales with SHAs rather than cells.
|
|
assert task_cells(1, (CANDIDATE_ARM,), 0) == [DURATIONS_BY_ARM[CANDIDATE_ARM][0]]
|
|
wide = generation_seconds(
|
|
task_count=1, runs=3, arms=REVIEW_ARMS, workers=9, fed_pool=True, unique_shas=5
|
|
)
|
|
assert wide >= PROPOSER_SECONDS + 5 * SHA_OVERHEAD_SECONDS
|
|
|
|
|
|
def test_sweep_overhead_does_not_shrink_with_the_arm_count():
|
|
"""The bias that made weekly look cheaper than it is.
|
|
|
|
A seeded weekly generation pays one arm instead of three but builds exactly
|
|
the same graphs. Charging the residual per cell billed it a third of a cost
|
|
the real sweep still pays; per SHA, the two attribute the same setup.
|
|
"""
|
|
|
|
kwargs = dict(task_count=6, runs=3, workers=3, fed_pool=False, unique_shas=5)
|
|
weekly = generation_seconds(arms=(CANDIDATE_ARM,), **kwargs)
|
|
cold = generation_seconds(arms=REVIEW_ARMS, **kwargs)
|
|
weekly_sessions = 6 * expected_task_seconds(3, (CANDIDATE_ARM,), 3, fed_pool=False)
|
|
cold_sessions = 6 * expected_task_seconds(3, REVIEW_ARMS, 3, fed_pool=False)
|
|
# Whatever each wall is, the non-session part is identical.
|
|
assert round(weekly - weekly_sessions) == round(cold - cold_sessions)
|
|
# Cycling wraps, so a task can ask for more runs than the sample holds.
|
|
long_sample = task_cells(len(DURATIONS_BY_ARM[CANDIDATE_ARM]) + 2, (CANDIDATE_ARM,), 0)
|
|
assert len(long_sample) == len(DURATIONS_BY_ARM[CANDIDATE_ARM]) + 2
|
|
|
|
|
|
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_task_seconds_is_alignment_averaged_and_deterministic():
|
|
waved = expected_task_seconds(3, REVIEW_ARMS, 3, fed_pool=False)
|
|
assert waved == expected_task_seconds(3, REVIEW_ARMS, 3, fed_pool=False)
|
|
assert expected_task_seconds(0, REVIEW_ARMS, 3, fed_pool=False) == 0.0
|
|
assert expected_task_seconds(3, (), 3, fed_pool=False) == 0.0
|
|
# The barrier can only cost time, never save it.
|
|
assert waved >= expected_task_seconds(3, REVIEW_ARMS, 3, fed_pool=True)
|
|
|
|
|
|
def test_a_generation_pays_one_proposer_session_on_top_of_its_tasks():
|
|
one = generation_seconds(
|
|
task_count=1, runs=3, arms=REVIEW_ARMS, workers=3, fed_pool=False, unique_shas=1
|
|
)
|
|
two = generation_seconds(
|
|
task_count=2, runs=3, arms=REVIEW_ARMS, workers=3, fed_pool=False, unique_shas=1
|
|
)
|
|
# Each extra task adds exactly one task's makespan. The proposer and the
|
|
# per-SHA sweep overhead are both paid once, not per task.
|
|
assert two - one == pytest.approx(
|
|
one - PROPOSER_SECONDS - SHA_OVERHEAD_SECONDS, abs=2.0
|
|
)
|
|
|
|
|
|
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_task_seconds(3, REVIEW_ARMS, 1, fed_pool=True)
|
|
assert expected_task_seconds(3, REVIEW_ARMS, workers, fed_pool=True) <= serial
|