Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
70 changes: 54 additions & 16 deletions scripts/ci/pr_runner_pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -1339,7 +1339,8 @@ def pick(load: Mapping[str, Mapping[str, int]], added: Mapping[str, int], usable
roots: Mapping[str, Mapping[str, int]] | None = None, root_jobs: int = 0,
queue_rounds: int = 0, taken: Mapping[str, int] | None = None,
taken_now: Mapping[str, int] | None = None, reserve: int = 0,
compared_jobs: int | None = None) -> Pick:
compared_jobs: int | None = None, root_taken: Mapping[str, int] | None = None,
root_taken_now: Mapping[str, int] | None = None) -> Pick:
"""The rule itself. `added` counts runs replayed since the snapshot on each pool.

With `queue_rounds` (CI_PR_POOL_QUEUE_ROUNDS) above 0: an owned pool, in
Expand All @@ -1351,7 +1352,10 @@ def pick(load: Mapping[str, Mapping[str, int]], added: Mapping[str, int], usable
pool with the least expected wait (expected_wait()), the earlier in
order on a tie. `taken` is the
peak of the runs since the snapshot that took each owned pool, by their
markers, and `taken_now` what they hold now (young_charge()). A
markers, and `taken_now` what they hold now (young_charge()).
`root_taken` and `root_taken_now` are the same for the root runners
(None: `taken` and `taken_now`); a newer run holds fewer root runners
than machines (choose()). A
replayed run counts REPLAYED_RUN_JOBS on an owned pool and one on its
root runners and on Blacksmith. `reserve` (main's full suite with
CI_OWNED_MAIN_RESERVE) is kept free on top of this run's jobs and root
Expand All @@ -1369,6 +1373,8 @@ def pick(load: Mapping[str, Mapping[str, int]], added: Mapping[str, int], usable
the fallback.
"""
roots, taken, taken_now = roots or {}, taken or {}, taken_now if taken_now is not None else taken or {}
root_taken = taken if root_taken is None else root_taken
root_taken_now = taken_now if root_taken_now is None else root_taken_now
blacksmith = [label for label in usable if not persistent(label)]
# Which Blacksmith pool: by its wait for this run's admission, since the
# shards may take another pool on the lane's Xcode (spread_shards()).
Expand All @@ -1391,7 +1397,8 @@ def pick(load: Mapping[str, Mapping[str, int]], added: Mapping[str, int], usable
limit = float(queue_rounds * job_minutes(label))
peak, now = taken.get(label, 0), taken_now.get(label, 0)
room = owned_room(label, load[label], added[label] * REPLAYED_RUN_JOBS, peak, now, queue_rounds, limit)
root_room = (owned_room(label, roots[label], added[label], peak, now, queue_rounds, limit)
root_room = (owned_room(label, roots[label], added[label], root_taken.get(label, 0),
root_taken_now.get(label, 0), queue_rounds, limit)
if label in roots else None)
rooms[label] = Pick(label, "owned", room, root_room, limit, whole if best and queue_rounds else None)
reserve = max(0, reserve)
Expand All @@ -1402,13 +1409,13 @@ def pick(load: Mapping[str, Mapping[str, int]], added: Mapping[str, int], usable
# An owned pool the run starts on now beats an earlier one it would
# queue on: with the rounds, std always fits by its queue places, so
# light sat idle while runs queued behind std's busy root runners.
def idle(counts: Mapping[str, int], added_jobs: int, label: str) -> int:
return counts["capacity"] - counts["running"] - counts["queued"] - taken_now.get(label, 0) - added_jobs
def idle(counts: Mapping[str, int], added_jobs: int, held: int) -> int:
return counts["capacity"] - counts["running"] - counts["queued"] - held - added_jobs

now = [label for label in fits
if idle(load[label], added[label] * REPLAYED_RUN_JOBS, label) >= jobs + reserve
if idle(load[label], added[label] * REPLAYED_RUN_JOBS, taken_now.get(label, 0)) >= jobs + reserve
and (label not in roots or root_jobs <= 0
or idle(roots[label], added[label], label) >= root_jobs + reserve)]
or idle(roots[label], added[label], root_taken_now.get(label, 0)) >= root_jobs + reserve)]
fits = now or fits
if split and not reserve and not fits and rooms and max(room.room for room in rooms.values()) >= 1:
# A pool with a root runner free first, when the run needs one.
Expand Down Expand Up @@ -1451,6 +1458,8 @@ def decide(
root_jobs: int = 0,
reserve: int = 0,
owned_now: Mapping[str, int] | None = None,
root_since: Mapping[str, int] | None = None,
root_now: Mapping[str, int] | None = None,
) -> Choice:
"""The preference rule over a janitor snapshot. Uncertainty keeps today's route.

Expand All @@ -1472,8 +1481,9 @@ def decide(
`split` lets this run take part of an owned pool (pick(), place()).
`root_jobs` is this run's peak on root runners, which an owned pool with
a root count must have free too. Its root runners are charged one per
replayed run (its admission), and a newer run's whole marker peak, since
a marker does not split it. `reserve` (main's full suite) is how many
replayed run (its admission), and `root_since` (`root_now` now) for the
runs in `owned_since`: the root runners they hold (choose()), None for
their whole peaks (`owned_since`, `owned_now`). `reserve` (main's full suite) is how many
machines, and root runners, an owned pool must keep free beyond this run.
`owned_now` is what the runs in `owned_since` hold now
(Routed.owned_now); None means their peaks.
Expand Down Expand Up @@ -1516,6 +1526,10 @@ def counted(label: str) -> bool:
taken = {label: max(0, int((owned_since or {}).get(label) or 0)) for label in usable if persistent(label)}
held = taken if owned_now is None else {
label: max(0, int(owned_now.get(label) or 0)) for label in usable if persistent(label)}
root_taken = taken if root_since is None else {
label: max(0, int(root_since.get(label) or 0)) for label in usable if persistent(label)}
root_held_now = held if root_now is None else {
label: max(0, int(root_now.get(label) or 0)) for label in usable if persistent(label)}
ephemeral = [label for label in usable if not persistent(label)]
queue_rounds = limits.queue_rounds
for _ in range(max(0, ephemeral_since) if ephemeral else 0):
Expand All @@ -1526,7 +1540,8 @@ def counted(label: str) -> bool:
for _ in range(max(0, routed_since)):
added[pick(load, added, usable, limits.max_queued, jobs=1, roots=roots, root_jobs=1,
queue_rounds=queue_rounds, taken=taken, taken_now=held,
compared_jobs=REPLAYED_RUN_JOBS).label] += 1
compared_jobs=REPLAYED_RUN_JOBS, root_taken=root_taken,
root_taken_now=root_held_now).label] += 1
if reserve:
# Main's full suite with a reserve (CI_OWNED_MAIN_RESERVE) takes an
# owned pool only while its peak and the reserve are free now: no
Expand All @@ -1535,7 +1550,7 @@ def counted(label: str) -> bool:
queue_rounds = 0
chosen = pick(load, added, candidates, limits.max_queued, jobs, split=split, roots=roots,
root_jobs=root_jobs, queue_rounds=queue_rounds, taken=taken, taken_now=held,
reserve=reserve)
reserve=reserve, root_taken=root_taken, root_taken_now=root_held_now)
label = chosen.label
if persistent(label) and chosen.how != "owned":
return Choice("", "", "every owned pool this run may take is busy, and no other pool is in the order")
Expand All @@ -1546,14 +1561,15 @@ def counted(label: str) -> bool:
for pool_label, count in taken.items() if count)
if chosen.how == "owned":

def idle(counts: Mapping[str, int], added_jobs: int) -> int:
def idle(counts: Mapping[str, int], added_jobs: int, peaks: Mapping[str, int],
now_held: Mapping[str, int]) -> int:
"""Free now: machines less running, queued and what newer runs hold (peaks with rounds 0)."""
if not queue_rounds:
return owned_room(label, counts, added_jobs, taken.get(label, 0), 0, 0, 0)
return counts["capacity"] - counts["running"] - counts["queued"] - held.get(label, 0) - added_jobs
return owned_room(label, counts, added_jobs, peaks.get(label, 0), 0, 0, 0)
return counts["capacity"] - counts["running"] - counts["queued"] - now_held.get(label, 0) - added_jobs

# Clamped for the text: an oversubscribed label has 0 free, not a negative count.
free_now = max(0, idle(load[label], added[label] * REPLAYED_RUN_JOBS))
free_now = max(0, idle(load[label], added[label] * REPLAYED_RUN_JOBS, taken, held))
places = max(0, chosen.room) - free_now
machines = f"{free_now} of {load[label]['capacity']} owned machines free"
if places > 0:
Expand All @@ -1562,7 +1578,7 @@ def idle(counts: Mapping[str, int], added_jobs: int) -> int:
machines += f" (Blacksmith's expected wait {chosen.blacksmith_wait:g} min)"
root, root_now = "", None
if chosen.root_room is not None:
root_now = max(0, idle(roots[label], added[label]))
root_now = max(0, idle(roots[label], added[label], root_taken, root_held_now))
root = f"; {root_now} of {roots[label]['capacity']} root runners free"
if chosen.root_room > root_now:
root += f" and {chosen.root_room - root_now} queue places"
Expand Down Expand Up @@ -1680,6 +1696,12 @@ def choose(

`queue_rounds` is CI_PR_POOL_QUEUE_ROUNDS as settings() reads it; a fork
run reads the janitor's copy instead.

A newer run that took an owned pool whose gui label has a count in
`owned_slots` holds one of its root runners, its admission: its gui-token
jobs take the gui label (gui_runner(), root_held()) and its side lanes
hold none. On a pool without one it may hold its whole marker peak there,
which a marker does not split, so that is its root charge.
"""
main = event == "workflow_dispatch" and ref == MAIN_REF
if event != "pull_request" and not main:
Expand Down Expand Up @@ -1814,11 +1836,27 @@ def choose(
# runner is idle; their later jobs are not charged (live_pools()).
older = {label: max(0, count - recent.runs().get(label, 0)) for label, count in before.runs().items()}
snapshot, owned_capacity = live_pools(snapshot, live_owned or {}, owned_capacity, older, live_online)
# On a pool with gui runners each newer run holds one root runner (its
# admission), not its whole peak: charging the peak left 0 of 15 root
# runners for a run while 3 newer runs held 3 (cmux run 36371179217,
# 2026-09-28). With the runners read live, the runs counted are the live
# window's; older runs' admissions show busy on the runners, or queued
# through `older` (live_pools()).
gui_slots = {} if fork else slots(owned_slots, xcode_pins.get(PR_XCODE_VARIABLE))
runs = routed.runs()

def root_charge(machines: Mapping[str, int]) -> dict[str, int]:
return {label: runs.get(label, 0) if gui_slots.get(gui_label(label), 0) > 0 else count
for label, count in machines.items()}

root_since = root_charge(routed.owned)
root_now = root_charge(routed.owned if routed.owned_now is None else routed.owned_now)
choice = decide(snapshot, limits, now=now, xcode_pins={} if fork else xcode_pins,
routed_since=routed.unknown, owned_since=routed.owned, ephemeral_since=routed.ephemeral,
auto_xcode=fork, owned_slots=owned_capacity, jobs=jobs,
split=(split or "").strip() == "1", shards=shards,
root_jobs=root_jobs, reserve=reserve, owned_now=routed.owned_now,
root_since=root_since, root_now=root_now,
# Main only ever takes an owned pool; the replay still
# spreads newer runs over the whole order.
choose_from=tuple(label for label in limits.order if persistent(label)) if main else None)
Expand Down
40 changes: 40 additions & 0 deletions tests/test_ci_pr_runner_pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -1341,6 +1341,7 @@ def test_main_names_the_owned_jobs_and_marks_their_peak(self):

ROOT_MINI = "glaeda-root-std-xcode-26.6"
SIDE_MINI = "glaeda-side-std-xcode-26.6"
GUI_MINI = "glaeda-gui-std-xcode-26.6"


class QueueBehindBusyRunners(unittest.TestCase):
Expand Down Expand Up @@ -3425,6 +3426,45 @@ def test_live_roots_ignore_committed_peaks(self):
self.assertGreaterEqual(live.root_budget, 9)
self.assertIn("0 of 19 root runners free", live.reason)

def test_newer_runs_hold_one_root_runner_each_with_gui_runners(self):
# The shape of run 36371179217 (2026-09-28): 15 root runners online and
# 12 busy, 2 runs replayed and 3 newer runs with a 32-machine peak on
# std. Charging that whole peak to the root runners left 0 for the
# run, so its admission went to Blacksmith. With gui runners each
# newer run holds one root runner, its admission.
routed = pool.Routed(unknown=2, owned={MINI: 32}, owned_now={MINI: 9}, owned_runs={MINI: 3},
live_now={MINI: 9}, live_runs={MINI: 3}, live_unknown=2)
kwargs = dict(machines=29, jobs=10, root_jobs=1, queue_rounds="2", split="1", routed=routed,
live_owned={MINI: 3, ROOT_MINI: 3}, live_online={MINI: 29, ROOT_MINI: 15})
fixed = owned_choice(fleet(busy=26), owned_slots=json.dumps({MINI: 29, ROOT_MINI: 19, GUI_MINI: 10}),
**kwargs)
self.assertEqual((fixed.runner, fixed.root_runner), (MINI, ROOT_MINI))
self.assertGreaterEqual(fixed.root_budget, 1)
# Its admission is placed owned.
self.assertIn("admission", pool.place(pool.FULL_RUN, fixed.owned_budget, root_budget=fixed.root_budget,
gui_runners=True)[0])
# Without gui runners on std (another pool's gui count does not count), a newer
# run's shards take root runners too: its whole peak is charged.
whole = owned_choice(fleet(busy=26), owned_slots=json.dumps({MINI: 29, ROOT_MINI: 19, LIGHT: 4,
"glaeda-gui-light-xcode-26.6": 4}), **kwargs)
self.assertEqual((whole.runner, whole.root_runner, whole.root_budget), (MINI, ROOT_MINI, 0))
self.assertIn("0 of 15 root runners free", whole.reason)

def test_decide_charges_the_root_runners_newer_runs_hold(self):
# Without `root_since` the root runners are charged the newer runs'
# machines; with it, only the root runners they hold.
snap = fleet(busy=0)
snap["pools"][ROOT_MINI] = {"queued": 0, "running": 0}
slots = {MINI: 11, ROOT_MINI: 10}
settings = pool.settings("", "", "", "1", PR_XCODE, "1")
peak = pool.decide(snap, settings, now=NOW, xcode_pins=OWNED_PINS, owned_slots=slots, jobs=1,
root_jobs=1, owned_since={MINI: 4}, owned_now={MINI: 4})
runs = pool.decide(snap, settings, now=NOW, xcode_pins=OWNED_PINS, owned_slots=slots, jobs=1,
root_jobs=1, owned_since={MINI: 4}, owned_now={MINI: 4},
root_since={MINI: 1}, root_now={MINI: 1})
self.assertIn("6 of 10 root runners free", peak.reason)
self.assertIn("9 of 10 root runners free", runs.reason)

def test_a_failed_snapshot_download_leaves_the_owned_pools_to_the_live_runners(self):
windows = []

Expand Down
Loading