ci: extract traffic-controller into a testable Python module (#342)
The priority-based runner orchestration ("traffic-control") lived as a large
inline-bash step in ci.yml — a two-pass preemption + hold-back script that was
effectively untestable in YAML. Move it into .github/scripts/traffic_control.py,
structured as a pure decision CORE + a thin gh-I/O SHELL:
* Pure functions (no network/clock/subprocess), unit-testable in isolation:
- effective_priority(pr): lowest-numbered P0-P9, default P5; broken OR draft => 10.
- runs_to_cancel(this_pr, all_prs, self_run_id): PASS 1 — run ids to cancel,
empty unless THIS PR is P0; only strictly-lower running/queued runs; never self
(by number or run id), never equal-or-higher.
- wait_blockers(this_pr, all_prs): PASS 2 — yield to any strictly-higher PR with
an active/queued run, and to same-level peers ordered ahead (running-first,
then oldest createdAt). Empty => proceed.
* Shell (run_live): gathers the snapshot via gh, applies cancels, runs the bounded
hold-back poll loop; always exits 0. --dry-run feeds the core a snapshot JSON and
prints decisions with zero network.
* No jq/bash dependency (cross-platform, per the repo's Python-stdlib convention).
ci.yml's traffic-control job now checks out the repo and runs the module. Job
permissions gain `contents: read` (for checkout) alongside the existing
`actions: write` / `pull-requests: read`; env and downstream `needs:` wiring
unchanged; step stays `continue-on-error`.
33 stdlib unittest cases cover priority resolution, P0-only preemption, the
self/main/equal-or-higher invariants, and the same-level running-first/oldest
ordering.
Behaviour is preserved except the ticket's refinements: (1) drafts now count as
P10 (bottom); (2) an explicit same-level running-first-then-oldest tiebreaker; and
(3) per the ticket's order-of-operations, ONLY P0 preempts — the old bash also let
any higher-priority PR cancel a `broken` target's run, which no longer happens
(a broken/draft run is only cancelled by a P0, via the same strictly-lower rule).
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,249 @@
|
||||
#!/usr/bin/env python3
|
||||
# SPDX-License-Identifier: GPL-3.0-or-later
|
||||
"""Unit tests for the pure decision core of traffic_control.py (no network).
|
||||
|
||||
Covers: priority resolution (P-label / broken / draft / default P5), PASS 1
|
||||
preemption (only P0 cancels strictly-lower; P1-P9 never bump; self / main /
|
||||
equal-or-higher never cancelled), PASS 2 hold-back (yield to strictly-higher with
|
||||
an active run; same-level running-first then oldest-first), and a few end-to-end
|
||||
decision scenarios."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import unittest
|
||||
|
||||
import traffic_control as tc
|
||||
from traffic_control import PullRequest
|
||||
|
||||
|
||||
def pr(number, labels=(), *, draft=False, created_at="", status=tc.NONE, run_ids=()):
|
||||
"""Terse PullRequest builder for tests."""
|
||||
return PullRequest(
|
||||
number=number,
|
||||
labels=tuple(labels),
|
||||
is_draft=draft,
|
||||
created_at=created_at,
|
||||
run_status=status,
|
||||
run_ids=tuple(run_ids),
|
||||
)
|
||||
|
||||
|
||||
class EffectivePriorityTests(unittest.TestCase):
|
||||
def test_no_labels_defaults_to_p5(self):
|
||||
self.assertEqual(tc.effective_priority(pr(1)), 5)
|
||||
|
||||
def test_non_priority_labels_ignored_default_p5(self):
|
||||
self.assertEqual(tc.effective_priority(pr(1, ["bug", "enhancement"])), 5)
|
||||
|
||||
def test_single_p_label(self):
|
||||
self.assertEqual(tc.effective_priority(pr(1, ["P3"])), 3)
|
||||
self.assertEqual(tc.effective_priority(pr(1, ["P0"])), 0)
|
||||
|
||||
def test_lowest_numbered_p_label_wins(self):
|
||||
self.assertEqual(tc.effective_priority(pr(1, ["P4", "P1", "P7"])), 1)
|
||||
|
||||
def test_broken_is_bottom_p10_overriding_p0(self):
|
||||
self.assertEqual(tc.effective_priority(pr(1, ["broken", "P0"])), 10)
|
||||
|
||||
def test_draft_is_bottom_p10_overriding_p0(self):
|
||||
self.assertEqual(tc.effective_priority(pr(1, ["P0"], draft=True)), 10)
|
||||
|
||||
def test_draft_and_broken_still_p10(self):
|
||||
self.assertEqual(tc.effective_priority(pr(1, ["broken"], draft=True)), 10)
|
||||
|
||||
def test_double_digit_pseudo_label_is_not_a_priority(self):
|
||||
# Only P0-P9 count (regex ^P[0-9]$); "P10" is not a valid priority label.
|
||||
self.assertEqual(tc.effective_priority(pr(1, ["P10"])), 5)
|
||||
|
||||
def test_priority_label_text(self):
|
||||
self.assertEqual(tc.priority_label(0), "P0")
|
||||
self.assertEqual(tc.priority_label(5), "P5")
|
||||
self.assertIn("bottom", tc.priority_label(10))
|
||||
|
||||
|
||||
class RunsToCancelTests(unittest.TestCase):
|
||||
def test_non_p0_self_cancels_nothing(self):
|
||||
me = pr(1, ["P1"])
|
||||
others = [pr(2, ["P5"], status=tc.RUNNING, run_ids=[200])]
|
||||
self.assertEqual(tc.runs_to_cancel(me, [me, *others]), [])
|
||||
|
||||
def test_non_p0_self_never_cancels_even_broken(self):
|
||||
# Behaviour change vs the old bash (see module/PR notes): P1-P9 no longer
|
||||
# reclaim a broken target's runner — only P0 preempts.
|
||||
me = pr(1, ["P3"])
|
||||
broken = pr(2, ["broken"], status=tc.RUNNING, run_ids=[200])
|
||||
self.assertEqual(tc.runs_to_cancel(me, [me, broken]), [])
|
||||
|
||||
def test_p0_cancels_strictly_lower_active_runs(self):
|
||||
me = pr(1, ["P0"])
|
||||
low = pr(2, ["P5"], status=tc.RUNNING, run_ids=[200])
|
||||
queued = pr(3, ["P9"], status=tc.QUEUED, run_ids=[300])
|
||||
self.assertEqual(
|
||||
sorted(tc.runs_to_cancel(me, [me, low, queued])), [200, 300])
|
||||
|
||||
def test_p0_cancels_broken_and_draft_lower_runs(self):
|
||||
me = pr(1, ["P0"])
|
||||
broken = pr(2, ["broken"], status=tc.RUNNING, run_ids=[200])
|
||||
draft = pr(3, ["P2"], draft=True, status=tc.RUNNING, run_ids=[300])
|
||||
self.assertEqual(
|
||||
sorted(tc.runs_to_cancel(me, [me, broken, draft])), [200, 300])
|
||||
|
||||
def test_p0_never_cancels_equal_priority_p0(self):
|
||||
me = pr(1, ["P0"])
|
||||
peer = pr(2, ["P0"], status=tc.RUNNING, run_ids=[200])
|
||||
self.assertEqual(tc.runs_to_cancel(me, [me, peer]), [])
|
||||
|
||||
def test_p0_never_cancels_self(self):
|
||||
me = pr(1, ["P0"], status=tc.RUNNING, run_ids=[100])
|
||||
self.assertEqual(tc.runs_to_cancel(me, [me]), [])
|
||||
|
||||
def test_p0_excludes_own_run_id_defensively(self):
|
||||
me = pr(1, ["P0"], status=tc.RUNNING, run_ids=[100])
|
||||
# A lower PR that somehow reports our own run id must not be cancelled.
|
||||
low = pr(2, ["P5"], status=tc.RUNNING, run_ids=[100, 200])
|
||||
self.assertEqual(
|
||||
tc.runs_to_cancel(me, [me, low], self_run_id=100), [200])
|
||||
|
||||
def test_p0_skips_lower_with_no_active_run(self):
|
||||
me = pr(1, ["P0"])
|
||||
idle = pr(2, ["P5"], status=tc.NONE, run_ids=[])
|
||||
self.assertEqual(tc.runs_to_cancel(me, [me, idle]), [])
|
||||
|
||||
|
||||
class WaitBlockersTests(unittest.TestCase):
|
||||
def test_p0_never_waits(self):
|
||||
me = pr(1, ["P0"])
|
||||
higher = pr(2, ["P0"], status=tc.RUNNING) # nothing outranks P0 anyway
|
||||
self.assertEqual(tc.wait_blockers(me, [me, higher]), [])
|
||||
|
||||
def test_yields_to_strictly_higher_with_active_run(self):
|
||||
me = pr(2, ["P5"], status=tc.RUNNING)
|
||||
higher = pr(1, ["P2"], status=tc.RUNNING)
|
||||
blockers = tc.wait_blockers(me, [me, higher])
|
||||
self.assertEqual([b.number for b in blockers], [1])
|
||||
self.assertEqual(blockers[0].kind, "higher-priority")
|
||||
|
||||
def test_does_not_yield_to_higher_without_active_run(self):
|
||||
me = pr(2, ["P5"], status=tc.RUNNING)
|
||||
higher_idle = pr(1, ["P2"], status=tc.NONE)
|
||||
self.assertEqual(tc.wait_blockers(me, [me, higher_idle]), [])
|
||||
|
||||
def test_does_not_yield_to_lower_priority(self):
|
||||
me = pr(1, ["P2"], status=tc.RUNNING)
|
||||
lower = pr(2, ["P5"], status=tc.RUNNING)
|
||||
self.assertEqual(tc.wait_blockers(me, [me, lower]), [])
|
||||
|
||||
def test_same_level_oldest_running_proceeds(self):
|
||||
me = pr(1, ["P5"], created_at="2026-07-01T00:00:00Z", status=tc.RUNNING)
|
||||
newer = pr(2, ["P5"], created_at="2026-07-02T00:00:00Z", status=tc.RUNNING)
|
||||
self.assertEqual(tc.wait_blockers(me, [me, newer]), [])
|
||||
|
||||
def test_same_level_newer_running_yields_to_older(self):
|
||||
# Coordinator clarification: within a level, older createdAt goes first.
|
||||
older = pr(1, ["P5"], created_at="2026-07-01T00:00:00Z", status=tc.RUNNING)
|
||||
me = pr(2, ["P5"], created_at="2026-07-02T00:00:00Z", status=tc.RUNNING)
|
||||
blockers = tc.wait_blockers(me, [me, older])
|
||||
self.assertEqual([b.number for b in blockers], [1])
|
||||
self.assertEqual(blockers[0].kind, "same-level-ahead")
|
||||
|
||||
def test_same_level_running_first_beats_older_waiting(self):
|
||||
# An in-flight peer keeps its place; a not-yet-running OLDER peer does not
|
||||
# jump ahead of us while we are the one already running.
|
||||
me = pr(2, ["P5"], created_at="2026-07-02T00:00:00Z", status=tc.RUNNING)
|
||||
older_waiting = pr(1, ["P5"], created_at="2026-07-01T00:00:00Z", status=tc.NONE)
|
||||
self.assertEqual(tc.wait_blockers(me, [me, older_waiting]), [])
|
||||
|
||||
def test_same_level_waiting_orders_oldest_before_newer(self):
|
||||
# Neither running: strictly oldest-first among the waiting bucket.
|
||||
oldest = pr(1, ["P5"], created_at="2026-07-01T00:00:00Z", status=tc.NONE)
|
||||
middle = pr(2, ["P5"], created_at="2026-07-02T00:00:00Z", status=tc.NONE)
|
||||
me = pr(3, ["P5"], created_at="2026-07-03T00:00:00Z", status=tc.NONE)
|
||||
blockers = tc.wait_blockers(me, [oldest, middle, me])
|
||||
self.assertEqual([b.number for b in blockers], [1, 2])
|
||||
|
||||
def test_same_level_queued_counts_as_waiting_ordered_by_age(self):
|
||||
# A queued peer is "waiting to start", not in-flight: ordered purely by age.
|
||||
me = pr(1, ["P5"], created_at="2026-07-01T00:00:00Z", status=tc.NONE)
|
||||
newer_queued = pr(2, ["P5"], created_at="2026-07-02T00:00:00Z", status=tc.QUEUED)
|
||||
self.assertEqual(tc.wait_blockers(me, [me, newer_queued]), [])
|
||||
|
||||
def test_broken_self_yields_to_everyone_active(self):
|
||||
me = pr(1, ["broken"], status=tc.RUNNING) # effective P10
|
||||
normal = pr(2, ["P5"], status=tc.RUNNING)
|
||||
blockers = tc.wait_blockers(me, [me, normal])
|
||||
self.assertEqual([b.number for b in blockers], [2])
|
||||
self.assertEqual(blockers[0].kind, "higher-priority")
|
||||
|
||||
|
||||
class EndToEndDecisionTests(unittest.TestCase):
|
||||
def test_p0_emergency_cancels_lower_and_proceeds(self):
|
||||
me = pr(10, ["P0"], status=tc.RUNNING, run_ids=[1000])
|
||||
prs = [
|
||||
me,
|
||||
pr(11, ["P2"], status=tc.RUNNING, run_ids=[1100]),
|
||||
pr(12, ["P5"], status=tc.QUEUED, run_ids=[1200]),
|
||||
pr(13, ["P0"], status=tc.RUNNING, run_ids=[1300]), # equal — spared
|
||||
]
|
||||
dec = tc.decide(me, prs, self_run_id=1000)
|
||||
self.assertEqual(sorted(dec.cancel_run_ids), [1100, 1200])
|
||||
self.assertTrue(dec.proceed)
|
||||
|
||||
def test_p5_waits_behind_running_higher(self):
|
||||
me = pr(20, ["P5"], status=tc.RUNNING, run_ids=[2000])
|
||||
higher = pr(21, ["P2"], status=tc.RUNNING, run_ids=[2100])
|
||||
dec = tc.decide(me, [me, higher])
|
||||
self.assertEqual(dec.cancel_run_ids, ()) # not P0 — cancels nothing
|
||||
self.assertFalse(dec.proceed)
|
||||
self.assertEqual([b.number for b in dec.blockers], [21])
|
||||
|
||||
def test_lone_p5_proceeds(self):
|
||||
me = pr(30, ["P5"], status=tc.RUNNING, run_ids=[3000])
|
||||
dec = tc.decide(me, [me])
|
||||
self.assertEqual(dec.cancel_run_ids, ())
|
||||
self.assertTrue(dec.proceed)
|
||||
|
||||
|
||||
class SnapshotParsingTests(unittest.TestCase):
|
||||
def test_from_json_label_objects_and_fields(self):
|
||||
obj = {
|
||||
"number": 7,
|
||||
"labels": [{"name": "P3"}, {"name": "bug"}],
|
||||
"isDraft": True,
|
||||
"createdAt": "2026-07-01T00:00:00Z",
|
||||
"runStatus": "running",
|
||||
"runIds": [42, 43],
|
||||
}
|
||||
p = PullRequest.from_json(obj)
|
||||
self.assertEqual(p.number, 7)
|
||||
self.assertEqual(p.labels, ("P3", "bug"))
|
||||
self.assertTrue(p.is_draft)
|
||||
self.assertEqual(p.run_status, tc.RUNNING)
|
||||
self.assertEqual(p.run_ids, (42, 43))
|
||||
self.assertEqual(tc.effective_priority(p), 10) # draft => bottom
|
||||
|
||||
def test_from_json_plain_string_labels_and_unknown_status(self):
|
||||
p = PullRequest.from_json(
|
||||
{"number": 8, "labels": ["P1"], "runStatus": "bogus"})
|
||||
self.assertEqual(p.labels, ("P1",))
|
||||
self.assertEqual(p.run_status, tc.NONE) # unknown -> none
|
||||
|
||||
def test_load_snapshot_roundtrip(self):
|
||||
text = json.dumps({
|
||||
"self": 2,
|
||||
"self_run_id": 222,
|
||||
"prs": [
|
||||
{"number": 1, "labels": ["P2"], "runStatus": "running",
|
||||
"runIds": [111]},
|
||||
{"number": 2, "labels": ["P5"], "runStatus": "running",
|
||||
"runIds": [222]},
|
||||
],
|
||||
})
|
||||
this_pr, all_prs, self_run_id = tc._load_snapshot(text)
|
||||
self.assertEqual(this_pr.number, 2)
|
||||
self.assertEqual(len(all_prs), 2)
|
||||
self.assertEqual(self_run_id, 222)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,491 @@
|
||||
#!/usr/bin/env python3
|
||||
# SPDX-License-Identifier: GPL-3.0-or-later
|
||||
"""CI traffic-controller: priority-based runner orchestration for LibreMail's
|
||||
`ci.yml`. Extracted out of the old inline-bash `traffic-control` step into a
|
||||
Python module so the decision logic is developer-legible and, above all, unit-
|
||||
testable (see `test_traffic_control.py`).
|
||||
|
||||
DESIGN: pure decision CORE + thin gh-I/O SHELL
|
||||
----------------------------------------------
|
||||
The decisions ("who do we cancel?", "do we proceed or wait?") are pure functions
|
||||
over plain `PullRequest` snapshots — no network, no clock, no subprocess — so they
|
||||
can be exercised exhaustively in unit tests. The SHELL (`run_live`) is the only part
|
||||
that touches `gh`: it gathers the snapshot, applies the cancellations, and runs the
|
||||
bounded hold-back poll loop. Feed the core a snapshot JSON (`--dry-run`) to see its
|
||||
decisions with zero network.
|
||||
|
||||
ORDER OF OPERATIONS (issue #342)
|
||||
--------------------------------
|
||||
1. Effective priority orders everything: the lowest-numbered `P0`-`P9` label present
|
||||
(P0 = highest), default `P5` if none. A `broken` OR `draft` PR is effectively P10
|
||||
(bottom, below P9), overriding any P0-P9 label.
|
||||
2. PASS 1 - preemption: **only a P0 (emergency) preempts.** A P0 cancels the
|
||||
in-progress / queued CI runs of ALL strictly-lower OTHER open PRs to reclaim their
|
||||
runners. P1-P9 never bump a lower run mid-flight.
|
||||
3. PASS 2 - bounded hold-back: a non-P0 PR yields (cancels nothing) to any strictly-
|
||||
higher-priority OTHER PR that has an active/queued run, and — among its OWN
|
||||
priority level — to any PR ordered ahead of it (running-first, then oldest by
|
||||
`createdAt`). It proceeds the moment it is at the front, or when the wait budget
|
||||
elapses (a PR never blocks itself).
|
||||
|
||||
SAFETY INVARIANTS (preserved from the original step)
|
||||
----------------------------------------------------
|
||||
* never cancel a run on `main` / a push event — the shell's `gh run list` query
|
||||
filters `--event pull_request` and drops `headBranch == main`, so only PR-event
|
||||
runs ever reach the core;
|
||||
* never cancel THIS PR's own run — skipped by PR number AND by run id;
|
||||
* never cancel an equal-or-higher-priority PR — only strictly-lower (prio > self).
|
||||
|
||||
This is deliberately NOT a merge gate: every gh call is guarded, the shell always
|
||||
exits 0, and the ci.yml step stays `continue-on-error`, so a hiccup (API error,
|
||||
missing permission, fork PR) can never fail CI.
|
||||
|
||||
Pure standard library, cross-platform (the primary dev box is Windows, where the old
|
||||
bash + `jq` pipeline had no clean equivalent).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
from dataclasses import dataclass, field
|
||||
|
||||
# ── Priority model ───────────────────────────────────────────────────────────
|
||||
BROKEN_LABEL = "broken"
|
||||
DEFAULT_PRIORITY = 5 # a PR with no P0-P9 label
|
||||
BOTTOM_PRIORITY = 10 # broken OR draft — below P9
|
||||
TOP_PRIORITY = 0 # P0, the only priority that preempts
|
||||
_P_LABEL = re.compile(r"^P([0-9])$") # single digit only, matching the old jq `^P[0-9]$`
|
||||
|
||||
# ── Run-status model (normalised from gh's raw run statuses) ─────────────────
|
||||
RUNNING = "running" # gh status in_progress
|
||||
QUEUED = "queued" # gh status queued / waiting / requested / pending
|
||||
NONE = "none" # no active run (completed or absent)
|
||||
ACTIVE = frozenset({RUNNING, QUEUED})
|
||||
# For aggregating a branch's overall status from its runs: running beats queued
|
||||
# beats none (most-active wins). NB: the same-level ORDER (see _ordering_key) is a
|
||||
# coarser two-bucket split — in-flight (running) vs everything-else-by-age.
|
||||
_STATUS_RANK = {RUNNING: 0, QUEUED: 1, NONE: 2}
|
||||
# createdAt sentinel so a PR with an unknown timestamp sorts LAST (never wrongly
|
||||
# "oldest"/front, so it yields rather than preempts another PR's front slot).
|
||||
_FAR_FUTURE = "9999-12-31T23:59:59Z"
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class PullRequest:
|
||||
"""A snapshot of one open PR. The pure decision core consumes only these — no
|
||||
network. `run_ids` are the PR's active (non-completed) CI run ids, already
|
||||
filtered to pull_request events on a non-main head by the shell that built them."""
|
||||
|
||||
number: int
|
||||
labels: tuple[str, ...] = ()
|
||||
is_draft: bool = False
|
||||
created_at: str = ""
|
||||
run_status: str = NONE
|
||||
run_ids: tuple[int, ...] = ()
|
||||
|
||||
@classmethod
|
||||
def from_json(cls, obj: dict) -> "PullRequest":
|
||||
"""Build from a snapshot dict. `labels` may be a list of names or of gh's
|
||||
label objects (`{"name": ...}`)."""
|
||||
raw_labels = obj.get("labels") or []
|
||||
names: list[str] = []
|
||||
for lab in raw_labels:
|
||||
if isinstance(lab, dict):
|
||||
name = lab.get("name")
|
||||
else:
|
||||
name = lab
|
||||
if name:
|
||||
names.append(str(name))
|
||||
status = (obj.get("runStatus") or obj.get("run_status") or NONE).lower()
|
||||
if status not in (RUNNING, QUEUED, NONE):
|
||||
status = NONE
|
||||
raw_ids = obj.get("runIds") or obj.get("run_ids") or ()
|
||||
return cls(
|
||||
number=int(obj["number"]),
|
||||
labels=tuple(names),
|
||||
is_draft=bool(obj.get("isDraft") or obj.get("is_draft")
|
||||
or obj.get("draft") or False),
|
||||
created_at=str(obj.get("createdAt") or obj.get("created_at") or ""),
|
||||
run_status=status,
|
||||
run_ids=tuple(int(r) for r in raw_ids),
|
||||
)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Blocker:
|
||||
"""A PR that THIS PR must yield to in PASS 2 (purely informational for logging)."""
|
||||
|
||||
number: int
|
||||
priority: int
|
||||
kind: str # "higher-priority" | "same-level-ahead"
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Decision:
|
||||
"""The full point-in-time decision for THIS PR (used by --dry-run and tests)."""
|
||||
|
||||
self_number: int
|
||||
self_priority: int
|
||||
cancel_run_ids: tuple[int, ...] = ()
|
||||
blockers: tuple[Blocker, ...] = ()
|
||||
proceed: bool = True
|
||||
|
||||
|
||||
# ── Pure decision core (no network / clock / subprocess) ─────────────────────
|
||||
def effective_priority(pr: PullRequest) -> int:
|
||||
"""Effective priority: `broken` OR `draft` => 10 (bottom, overriding any P0-P9);
|
||||
else the lowest-numbered P0-P9 label present; else the default P5."""
|
||||
if pr.is_draft or BROKEN_LABEL in pr.labels:
|
||||
return BOTTOM_PRIORITY
|
||||
nums = [int(m.group(1)) for name in pr.labels if (m := _P_LABEL.match(name))]
|
||||
return min(nums) if nums else DEFAULT_PRIORITY
|
||||
|
||||
|
||||
def priority_label(prio: int) -> str:
|
||||
"""Human-readable priority for logs."""
|
||||
if prio >= BOTTOM_PRIORITY:
|
||||
return f"P{BOTTOM_PRIORITY} (broken/draft — bottom, below P9)"
|
||||
return f"P{prio}"
|
||||
|
||||
|
||||
def _ordering_key(pr: PullRequest) -> tuple[int, str, int]:
|
||||
"""Same-level ordering (issue #342 rule 3): an in-flight (RUNNING) run keeps its
|
||||
place at the front — a same-level peer never reorders it — then, among the PRs
|
||||
still waiting to start (QUEUED or no run yet), OLDEST createdAt first (ascending),
|
||||
then PR number as a stable final tiebreak so the order is fully deterministic."""
|
||||
in_flight = 0 if pr.run_status == RUNNING else 1
|
||||
return (in_flight, pr.created_at or _FAR_FUTURE, pr.number)
|
||||
|
||||
|
||||
def runs_to_cancel(
|
||||
this_pr: PullRequest,
|
||||
all_prs: list[PullRequest],
|
||||
*,
|
||||
self_run_id: int | None = None,
|
||||
) -> list[int]:
|
||||
"""PASS 1. Run ids to cancel — **empty unless THIS PR is P0**. A P0 preempts the
|
||||
active (running/queued) runs of every strictly-lower OTHER PR. Invariants: never
|
||||
cancel self (by number or run id), never cancel an equal-or-higher-priority PR."""
|
||||
if effective_priority(this_pr) != TOP_PRIORITY:
|
||||
return [] # only P0 preempts; P1-P9 never bump
|
||||
to_cancel: list[int] = []
|
||||
seen: set[int] = set()
|
||||
for pr in all_prs:
|
||||
if pr.number == this_pr.number:
|
||||
continue # never cancel self
|
||||
if effective_priority(pr) <= TOP_PRIORITY:
|
||||
continue # only strictly-lower (skip other P0s)
|
||||
if pr.run_status not in ACTIVE:
|
||||
continue # nothing running/queued to cancel
|
||||
for rid in pr.run_ids:
|
||||
if self_run_id is not None and rid == self_run_id:
|
||||
continue # never cancel our own run
|
||||
if rid in seen:
|
||||
continue
|
||||
seen.add(rid)
|
||||
to_cancel.append(rid)
|
||||
return to_cancel
|
||||
|
||||
|
||||
def wait_blockers(this_pr: PullRequest, all_prs: list[PullRequest]) -> list[Blocker]:
|
||||
"""PASS 2. The PRs THIS PR must yield to right now (empty => proceed). P0 never
|
||||
yields. Otherwise yield to (a) any strictly-higher-priority OTHER PR with an
|
||||
active/queued run, and (b) any SAME-priority PR ordered ahead of THIS PR
|
||||
(running-first, then oldest createdAt)."""
|
||||
self_prio = effective_priority(this_pr)
|
||||
if self_prio == TOP_PRIORITY:
|
||||
return [] # P0 outranks everything — never wait
|
||||
|
||||
others = [pr for pr in all_prs if pr.number != this_pr.number]
|
||||
blockers: list[Blocker] = []
|
||||
|
||||
# (a) strictly-higher-priority PRs that actually have an active/queued run.
|
||||
for pr in others:
|
||||
p = effective_priority(pr)
|
||||
if p < self_prio and pr.run_status in ACTIVE:
|
||||
blockers.append(Blocker(pr.number, p, "higher-priority"))
|
||||
|
||||
# (b) same-level ordering: THIS PR proceeds only when it is at the front.
|
||||
same_level = [pr for pr in others if effective_priority(pr) == self_prio]
|
||||
same_level.append(this_pr) # this_pr appears exactly once
|
||||
for pr in sorted(same_level, key=_ordering_key):
|
||||
if pr.number == this_pr.number:
|
||||
break # reached self => nobody ahead remains
|
||||
blockers.append(Blocker(pr.number, self_prio, "same-level-ahead"))
|
||||
|
||||
return blockers
|
||||
|
||||
|
||||
def decide(
|
||||
this_pr: PullRequest,
|
||||
all_prs: list[PullRequest],
|
||||
*,
|
||||
self_run_id: int | None = None,
|
||||
) -> Decision:
|
||||
"""Convenience: the full point-in-time decision (both passes) for THIS PR."""
|
||||
cancels = runs_to_cancel(this_pr, all_prs, self_run_id=self_run_id)
|
||||
blockers = wait_blockers(this_pr, all_prs)
|
||||
return Decision(
|
||||
self_number=this_pr.number,
|
||||
self_priority=effective_priority(this_pr),
|
||||
cancel_run_ids=tuple(cancels),
|
||||
blockers=tuple(blockers),
|
||||
proceed=not blockers,
|
||||
)
|
||||
|
||||
|
||||
# ── gh I/O shell (the only part that touches the network) ────────────────────
|
||||
def _log(msg: str) -> None:
|
||||
print(msg, flush=True)
|
||||
|
||||
|
||||
def _gh_json(args: list[str]) -> list | dict | None:
|
||||
"""Run `gh <args> --json ...` and parse stdout as JSON. Returns None (never
|
||||
raises) on any failure — the caller fails open."""
|
||||
try:
|
||||
proc = subprocess.run(
|
||||
["gh", *args],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
check=False,
|
||||
)
|
||||
except (OSError, ValueError) as exc:
|
||||
_log(f"::warning::gh invocation failed ({' '.join(args[:2])}): {exc}")
|
||||
return None
|
||||
if proc.returncode != 0:
|
||||
_log(f"::warning::gh exited {proc.returncode} ({' '.join(args[:2])}): "
|
||||
f"{proc.stderr.strip()}")
|
||||
return None
|
||||
try:
|
||||
return json.loads(proc.stdout or "null")
|
||||
except json.JSONDecodeError as exc:
|
||||
_log(f"::warning::could not parse gh JSON ({' '.join(args[:2])}): {exc}")
|
||||
return None
|
||||
|
||||
|
||||
def _normalise_status(raw: str) -> str:
|
||||
"""Map a gh run status onto our RUNNING / QUEUED / NONE model."""
|
||||
if raw == "in_progress":
|
||||
return RUNNING
|
||||
if raw == "completed":
|
||||
return NONE
|
||||
return QUEUED # queued / waiting / requested / pending
|
||||
|
||||
|
||||
def _runs_by_head(limit: int = 300) -> dict[str, dict]:
|
||||
"""One bulk `gh run list` -> {headBranch: {"status", "ids"}} for active PR-event
|
||||
runs. Enforces the 'never cancel main/push' invariant at the source: only
|
||||
`event == pull_request`, non-`main`, non-completed runs are kept. Active runs are
|
||||
the most recent, so `limit` most-recent runs comfortably covers them."""
|
||||
rows = _gh_json([
|
||||
"run", "list", "--workflow", "ci.yml", "--event", "pull_request",
|
||||
"--limit", str(limit),
|
||||
"--json", "databaseId,status,headBranch,event",
|
||||
])
|
||||
by_head: dict[str, dict] = {}
|
||||
for row in rows or []:
|
||||
if row.get("event") != "pull_request":
|
||||
continue
|
||||
head = row.get("headBranch")
|
||||
if not head or head == "main":
|
||||
continue
|
||||
if row.get("status") == "completed":
|
||||
continue
|
||||
entry = by_head.setdefault(head, {"status": NONE, "ids": []})
|
||||
entry["ids"].append(int(row["databaseId"]))
|
||||
status = _normalise_status(row.get("status", ""))
|
||||
# running beats queued beats none for the branch's aggregate status.
|
||||
if _STATUS_RANK[status] < _STATUS_RANK[entry["status"]]:
|
||||
entry["status"] = status
|
||||
return by_head
|
||||
|
||||
|
||||
def gather_snapshot(self_pr_number: int) -> tuple[PullRequest | None, list[PullRequest]]:
|
||||
"""Build (this_pr, all_prs) from live gh data. this_pr is forced to RUNNING —
|
||||
by definition our own run is in progress while this job executes."""
|
||||
prs = _gh_json([
|
||||
"pr", "list", "--state", "open", "--limit", "300",
|
||||
"--json", "number,headRefName,labels,isDraft,createdAt",
|
||||
])
|
||||
if prs is None:
|
||||
return None, []
|
||||
runs = _runs_by_head()
|
||||
all_prs: list[PullRequest] = []
|
||||
this_pr: PullRequest | None = None
|
||||
for obj in prs:
|
||||
head = obj.get("headRefName") or ""
|
||||
run_info = runs.get(head, {"status": NONE, "ids": []})
|
||||
number = int(obj["number"])
|
||||
is_self = number == self_pr_number
|
||||
pr = PullRequest.from_json({
|
||||
**obj,
|
||||
# self is definitionally running (this job is in progress).
|
||||
"runStatus": RUNNING if is_self else run_info["status"],
|
||||
"runIds": run_info["ids"],
|
||||
})
|
||||
all_prs.append(pr)
|
||||
if is_self:
|
||||
this_pr = pr
|
||||
return this_pr, all_prs
|
||||
|
||||
|
||||
def _cancel_run(run_id: int) -> bool:
|
||||
try:
|
||||
proc = subprocess.run(
|
||||
["gh", "run", "cancel", str(run_id)],
|
||||
capture_output=True, text=True, check=False,
|
||||
)
|
||||
except (OSError, ValueError) as exc:
|
||||
_log(f"::warning::could not cancel run {run_id}: {exc}")
|
||||
return False
|
||||
if proc.returncode == 0:
|
||||
return True
|
||||
_log(f"::warning::could not cancel run {run_id} — likely already finished. "
|
||||
f"{proc.stderr.strip()}")
|
||||
return False
|
||||
|
||||
|
||||
def _positive_int(env_name: str, default: int) -> int:
|
||||
raw = os.environ.get(env_name, "")
|
||||
return int(raw) if raw.isdigit() and int(raw) > 0 else default
|
||||
|
||||
|
||||
def run_live() -> int:
|
||||
"""The gh-driven shell: gather, PASS 1 (cancel), PASS 2 (bounded hold-back). Always
|
||||
returns 0 — the traffic-controller must never fail CI."""
|
||||
event = os.environ.get("GITHUB_EVENT_NAME", "")
|
||||
self_raw = os.environ.get("SELF_PR", "")
|
||||
if event != "pull_request" or not self_raw.isdigit():
|
||||
_log("Not a pull_request event (or no PR number) — nothing to do.")
|
||||
return 0
|
||||
self_number = int(self_raw)
|
||||
self_run_id = int(os.environ["GITHUB_RUN_ID"]) if os.environ.get(
|
||||
"GITHUB_RUN_ID", "").isdigit() else None
|
||||
|
||||
this_pr, all_prs = gather_snapshot(self_number)
|
||||
if this_pr is None:
|
||||
_log("::warning::Could not resolve THIS PR from the open-PR list — skipping.")
|
||||
return 0
|
||||
|
||||
self_prio = effective_priority(this_pr)
|
||||
_log(f"This PR #{self_number} effective priority: {priority_label(self_prio)} "
|
||||
"(P0 = highest/emergency, P9 = lowest, broken/draft = bottom).")
|
||||
|
||||
# ── PASS 1: PREEMPTION (only a P0 self) ──────────────────────────────────
|
||||
to_cancel = runs_to_cancel(this_pr, all_prs, self_run_id=self_run_id)
|
||||
if not to_cancel:
|
||||
if self_prio == TOP_PRIORITY:
|
||||
_log("P0 emergency — no strictly-lower active runs to cancel.")
|
||||
else:
|
||||
_log("Not P0 — no preemption (P1-P9 never cancel a lower run mid-flight).")
|
||||
else:
|
||||
cancelled = 0
|
||||
for rid in to_cancel:
|
||||
if _cancel_run(rid):
|
||||
_log(f" cancelled run {rid} (freed its runner).")
|
||||
cancelled += 1
|
||||
_log(f"P0 preemption complete — cancelled {cancelled}/{len(to_cancel)} run(s).")
|
||||
|
||||
# ── PASS 2: BOUNDED HOLD-BACK (yield to higher / same-level-ahead) ────────
|
||||
if self_prio == TOP_PRIORITY:
|
||||
_log("P0 emergency — not yielding; proceeding immediately.")
|
||||
return 0
|
||||
|
||||
budget = _positive_int("HOLD_BACK_BUDGET_SECONDS", 180)
|
||||
poll = _positive_int("HOLD_BACK_POLL_SECONDS", 15)
|
||||
deadline = time.monotonic() + budget
|
||||
_log(f"{priority_label(self_prio)} — holding back up to {budget}s for higher / "
|
||||
"earlier same-level PRs (no cancellation).")
|
||||
|
||||
while True:
|
||||
remaining = deadline - time.monotonic()
|
||||
if remaining <= 0:
|
||||
_log("Hold-back budget elapsed — proceeding; higher-priority PRs got their "
|
||||
"head start.")
|
||||
break
|
||||
# Refresh OTHER PRs so newly-opened higher-priority PRs are seen mid-wait;
|
||||
# THIS PR's own identity/priority stays fixed (matching the original).
|
||||
_, fresh = gather_snapshot(self_number)
|
||||
if not fresh:
|
||||
_log("::warning::Could not refresh open PRs — proceeding.")
|
||||
break
|
||||
blockers = wait_blockers(this_pr, fresh)
|
||||
if not blockers:
|
||||
_log("No higher-priority or earlier same-level PR is ahead — proceeding.")
|
||||
break
|
||||
tags = " ".join(f"#{b.number}(P{b.priority},{b.kind})" for b in blockers)
|
||||
sleep_s = min(poll, int(remaining)) if remaining >= 1 else 0
|
||||
_log(f"Yielding to: {tags} — re-checking in {sleep_s}s "
|
||||
f"({int(remaining)}s budget left).")
|
||||
if sleep_s > 0:
|
||||
time.sleep(sleep_s)
|
||||
|
||||
_log("Hold-back complete — this PR's heavy jobs may now start.")
|
||||
return 0
|
||||
|
||||
|
||||
# ── --dry-run: feed the pure core a snapshot JSON, print its decisions ───────
|
||||
def _load_snapshot(text: str) -> tuple[PullRequest, list[PullRequest], int | None]:
|
||||
data = json.loads(text)
|
||||
all_prs = [PullRequest.from_json(o) for o in data.get("prs", [])]
|
||||
self_number = int(data["self"])
|
||||
self_run_id = data.get("self_run_id")
|
||||
self_run_id = int(self_run_id) if self_run_id is not None else None
|
||||
this_pr = next((p for p in all_prs if p.number == self_number), None)
|
||||
if this_pr is None:
|
||||
raise ValueError(f"self #{self_number} not present in prs[]")
|
||||
return this_pr, all_prs, self_run_id
|
||||
|
||||
|
||||
def run_dry(text: str) -> int:
|
||||
this_pr, all_prs, self_run_id = _load_snapshot(text)
|
||||
dec = decide(this_pr, all_prs, self_run_id=self_run_id)
|
||||
_log(f"This PR #{dec.self_number} effective priority: "
|
||||
f"{priority_label(dec.self_priority)}")
|
||||
if dec.self_priority == TOP_PRIORITY:
|
||||
_log(f"PASS 1 (preemption): P0 — cancel run ids: "
|
||||
f"{list(dec.cancel_run_ids) or '(none active)'}")
|
||||
else:
|
||||
_log("PASS 1 (preemption): not P0 — no cancellations.")
|
||||
if dec.proceed:
|
||||
_log("PASS 2 (hold-back): PROCEED — no blockers.")
|
||||
else:
|
||||
tags = ", ".join(f"#{b.number}(P{b.priority}, {b.kind})" for b in dec.blockers)
|
||||
_log(f"PASS 2 (hold-back): WAIT — yielding to: {tags}")
|
||||
return 0
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
# Emit UTF-8 regardless of the host console so the log typography is stable on
|
||||
# the UTF-8 CI runners (and never raises on a legacy Windows code page).
|
||||
try:
|
||||
sys.stdout.reconfigure(encoding="utf-8", errors="replace")
|
||||
except (AttributeError, ValueError):
|
||||
pass
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument(
|
||||
"--dry-run", action="store_true",
|
||||
help="read a snapshot JSON (from --input or stdin), print decisions, no network.")
|
||||
parser.add_argument(
|
||||
"--input", help="snapshot JSON file for --dry-run (default: stdin).")
|
||||
args = parser.parse_args(argv)
|
||||
|
||||
if args.dry_run:
|
||||
text = (open(args.input, encoding="utf-8").read() if args.input
|
||||
else sys.stdin.read())
|
||||
return run_dry(text)
|
||||
return run_live()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
try:
|
||||
sys.exit(main())
|
||||
except Exception as exc: # never let the controller fail CI
|
||||
_log(f"::warning::traffic_control crashed, proceeding fail-open: {exc}")
|
||||
sys.exit(0)
|
||||
+52
-229
@@ -20,46 +20,43 @@ env:
|
||||
ANDROID_BUILD_TOOLS: "build-tools;37.0.0"
|
||||
|
||||
jobs:
|
||||
# ── Priority-based runner orchestration ──────────────────────────────────────
|
||||
# ── Priority-based runner orchestration ─────────────────────────────
|
||||
# Runs FIRST (the heavy jobs below all `needs: traffic-control`). It reads THIS
|
||||
# PR's P0–P9 label (and the special `broken` label) to order runner access.
|
||||
# Effective priority: `broken` => 10 (BOTTOM, below P9), overriding any P0–P9;
|
||||
# else the lowest-numbered P0–P9 label present (P0 = highest); else default P5.
|
||||
# PR's P0–P9 label, `broken` label, and draft state to order runner access. The
|
||||
# decision logic lives in .github/scripts/traffic_control.py — a pure, unit-tested
|
||||
# core (see .github/scripts/test_traffic_control.py) plus a thin gh-I/O shell; this
|
||||
# step just checks out the repo and runs it.
|
||||
#
|
||||
# Effective priority: a `broken` OR `draft` PR => 10 (BOTTOM, below P9), overriding
|
||||
# any P0–P9; else the lowest-numbered P0–P9 label present (P0 = highest); else P5.
|
||||
#
|
||||
# • P0 = EMERGENCY ONLY (app broken in production / emergency security update).
|
||||
# P0 PREEMPTS: it cancels the in-progress / queued CI runs of ALL strictly-
|
||||
# LOWER-priority OTHER open PRs to grab their runners immediately. A preempted
|
||||
# PR simply re-runs on its next push / autoupdate rebase.
|
||||
# PR simply re-runs on its next push / autoupdate rebase. P0 is the ONLY
|
||||
# priority that preempts — P1–P9 never cancel a lower run mid-flight.
|
||||
#
|
||||
# • `broken` = STUCK/FAILING PR — a MANUALLY-applied signal (maintainer / repo
|
||||
# owner only) meaning "deprioritise to the bottom so others aren't blocked
|
||||
# behind it while it's being fixed." Its effective priority is 10, so it NEVER
|
||||
# preempts (even if it's also labelled P0 — `broken` wins; a stuck PR can't be
|
||||
# an emergency merge) and ALWAYS yields: every other PR, even lower P-levels,
|
||||
# advances ahead of it. And because a broken PR's run is wasted (it can't
|
||||
# merge), ANY higher-priority PR — not just P0 — MAY cancel its in-progress run
|
||||
# to reclaim the runner. Removing the label restores its normal P-priority.
|
||||
# (Example: a P3 PR with failing CI was making lower-priority PRs wait behind
|
||||
# it; marking it `broken` lets them proceed — and reclaim its runner.)
|
||||
# • P1–P9 = YIELD WITHOUT BUMPING. They NEVER cancel a lower-priority run that is
|
||||
# already going — a higher-priority PR does not evict it, it just takes the next
|
||||
# free slot. Mechanism: a bounded hold-back. This job defers (up to
|
||||
# HOLD_BACK_BUDGET_SECONDS, kept well under timeout-minutes) while any strictly-
|
||||
# higher-priority OTHER open PR still has an active/queued CI run, and — within
|
||||
# its OWN priority level — while any peer is ordered ahead of it (an in-flight
|
||||
# run keeps its place; then oldest createdAt first). It proceeds the moment it
|
||||
# is at the front, or when the budget elapses (a PR never blocks itself).
|
||||
#
|
||||
# • P1–P9 = YIELD WITHOUT BUMPING. They NEVER cancel a NON-broken lower-priority
|
||||
# run that is already going — a higher-priority PR does not evict it, it just
|
||||
# takes the next free slot. Mechanism: a bounded hold-back. This job polls and
|
||||
# defers (up to HOLD_BACK_BUDGET_SECONDS, kept well under timeout-minutes) while
|
||||
# any strictly-higher-priority OTHER open PR still has an active/queued CI run,
|
||||
# so that PR's heavy jobs reach the runner queue ahead of this PR's. When the
|
||||
# budget elapses it proceeds anyway (a PR never blocks itself).
|
||||
# • `broken` / `draft` = BOTTOM (effective P10). Always yields, never preempts.
|
||||
# A maintainer applies `broken` to a stuck/failing PR to deprioritise it below
|
||||
# everything so others aren't blocked behind it; a draft isn't merge-ready, so
|
||||
# it likewise waits behind every ready PR. Only a P0 may cancel a bottom PR's
|
||||
# run (the same strictly-lower rule as any other target).
|
||||
#
|
||||
# Net preemption rule — a strictly-lower-priority OTHER PR's active run is cancelled
|
||||
# iff (THIS PR is P0) OR (that PR is `broken`); otherwise it is left to run and we
|
||||
# yield. So: P0 preempts ALL lower runs; ANY PR preempts lower `broken` runs; P1–P9
|
||||
# never preempt a non-broken run.
|
||||
#
|
||||
# Hard safety rules, all enforced in the script below:
|
||||
# • never cancels a run on main / a push event (filters --event pull_request);
|
||||
# Hard safety invariants, enforced in the script:
|
||||
# • never cancels a run on main / a push event (the gh query filters
|
||||
# --event pull_request and drops headBranch == main);
|
||||
# • never cancels THIS PR's own run (skips self by PR number + run id);
|
||||
# • never cancels an equal-or-higher-priority PR (only strictly-lower, prio > self);
|
||||
# • P1–P9 cancel NO non-broken run — they only wait (bounded), then proceed.
|
||||
# • P1–P9 cancel NOTHING — they only wait (bounded), then proceed.
|
||||
#
|
||||
# Honest limitation: GitHub Actions has no native priority queue and assigns
|
||||
# runners roughly FIFO, so the hold-back is a BEST-EFFORT head-start, not a hard
|
||||
@@ -69,215 +66,41 @@ jobs:
|
||||
#
|
||||
# It is deliberately NOT a merge-gate check: it is absent from `ci-passed`'s
|
||||
# needs, every API call is guarded, the script always exits 0, and the step is
|
||||
# `continue-on-error` — so a hiccup (API error, missing permission, fork PR)
|
||||
# can never fail or block CI. The heavy jobs only *order* after it via `needs`;
|
||||
# if it were ever skipped/failed they'd be skipped, which `ci-passed` now treats
|
||||
# as a gate failure (fail-safe: blocks merge, never spuriously passes).
|
||||
# `continue-on-error` — so a hiccup (API error, missing permission, fork PR) can
|
||||
# never fail or block CI. The heavy jobs only *order* after it via `needs`; if it
|
||||
# were ever skipped/failed they'd be skipped, which `ci-passed` treats as a gate
|
||||
# failure (fail-safe: blocks merge, never spuriously passes).
|
||||
traffic-control:
|
||||
name: Traffic control (runner priority)
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 6 # hard backstop; the P1–P9 hold-back budget below stays well under this
|
||||
timeout-minutes: 6 # hard backstop; the P1–P9 hold-back budget stays well under this
|
||||
permissions:
|
||||
actions: write # cancel lower-priority runs (P0 emergencies + broken targets)
|
||||
pull-requests: read # read PR P0–P9 labels
|
||||
contents: read # check out .github/scripts/traffic_control.py
|
||||
actions: write # cancel lower-priority runs (P0 emergencies)
|
||||
pull-requests: read # read PR P0–P9 labels + draft state
|
||||
env:
|
||||
GH_TOKEN: ${{ github.token }}
|
||||
GH_REPO: ${{ github.repository }}
|
||||
SELF_PR: ${{ github.event.pull_request.number }}
|
||||
# P1–P9 bounded hold-back knobs. BUDGET must stay comfortably below
|
||||
# timeout-minutes so the poll loop always exits 0 before the hard job timeout
|
||||
# fires — a timed-out job would skip the heavy jobs and fail `ci-passed`.
|
||||
# P1–P9 bounded hold-back knobs, read by traffic_control.py. BUDGET must stay
|
||||
# comfortably below timeout-minutes so the poll loop always exits 0 before the
|
||||
# hard job timeout fires — a timed-out job would skip the heavy jobs and fail
|
||||
# `ci-passed`.
|
||||
HOLD_BACK_BUDGET_SECONDS: "180"
|
||||
HOLD_BACK_POLL_SECONDS: "15"
|
||||
steps:
|
||||
# No checkout: this job only calls the gh CLI (auto-configured from GH_TOKEN /
|
||||
# GH_REPO), so it needs neither the repo contents nor the default contents:read.
|
||||
- name: Apply runner priority (P0/broken preempt; P1–P9 hold back)
|
||||
continue-on-error: true # belt-and-suspenders: never let this fail the run
|
||||
run: |
|
||||
# GitHub invokes run steps with `bash -eo pipefail`. Disable errexit so a
|
||||
# single failed API call can't abort the step; we guard every call and
|
||||
# always exit 0. Attacker-influenced values (branch names, labels) are only
|
||||
# ever read via env / gh JSON into shell vars — never interpolated as code.
|
||||
set +e
|
||||
|
||||
if [ "${GITHUB_EVENT_NAME:-}" != "pull_request" ] || [ -z "${SELF_PR:-}" ]; then
|
||||
echo "Not a pull_request event (or no PR number) — nothing to do."
|
||||
exit 0
|
||||
fi
|
||||
|
||||
# Effective priority of a labels JSON array read on stdin: a `broken` label
|
||||
# => 10 (bottom, below P9), overriding any P0–P9; else the highest-priority
|
||||
# (lowest-numbered) P0–P9 label present; else 5.
|
||||
prio_of() {
|
||||
jq -r 'if any(.[]; .name == "broken") then 10
|
||||
else ([ .[] | .name | select(test("^P[0-9]$")) | ltrimstr("P") | tonumber ]
|
||||
| if length == 0 then 5 else min end) end' 2>/dev/null
|
||||
}
|
||||
|
||||
# Snapshot of every open PR (number, head branch, labels) to stdout.
|
||||
list_open_prs() {
|
||||
gh pr list --state open --limit 300 --json number,headRefName,labels
|
||||
}
|
||||
|
||||
# Active (non-completed) CI run ids on head branch $1 — PR events only, never
|
||||
# main. Shared by the preemption pass (ids to cancel) and the hold-back
|
||||
# (presence => keep waiting).
|
||||
active_run_ids_for_head() {
|
||||
gh run list --workflow ci.yml --branch "$1" --event pull_request \
|
||||
--limit 100 --json databaseId,status,headBranch,event 2>/dev/null \
|
||||
| jq -r '.[]
|
||||
| select(.event == "pull_request")
|
||||
| select(.headBranch != "main")
|
||||
| select(.status != "completed")
|
||||
| .databaseId' 2>/dev/null
|
||||
}
|
||||
|
||||
# One initial snapshot, used to read THIS PR's own priority.
|
||||
if ! list_open_prs > open_prs.json 2>err.txt; then
|
||||
echo "::warning::Could not list open PRs — skipping. $(cat err.txt 2>/dev/null)"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
self_labels=$(jq -c --argjson pr "$SELF_PR" \
|
||||
'([ .[] | select(.number == $pr) | .labels ] | .[0]) // []' open_prs.json 2>/dev/null)
|
||||
self_prio=$(printf '%s' "${self_labels:-[]}" | prio_of)
|
||||
case "$self_prio" in ''|*[!0-9]*) self_prio=5 ;; esac
|
||||
if [ "$self_prio" -ge 10 ]; then prio_label="broken (below P9, bottom)"; else prio_label="P$self_prio"; fi
|
||||
echo "This PR #$SELF_PR effective priority: $prio_label (P0 = highest/emergency, P9 = lowest, 'broken' = bottom)."
|
||||
|
||||
# ── PASS 1: PREEMPTION — cancel a strictly-lower OTHER PR's active runs ──
|
||||
# A strictly-lower-priority (prio > self) OTHER PR's active CI run is
|
||||
# cancelled iff keeping it running is wasteful, i.e. EITHER:
|
||||
# • THIS PR is P0 (emergency — reclaim every lower runner now), OR
|
||||
# • that OTHER PR is `broken` (its run can't merge, so ANY higher-priority
|
||||
# PR — not just P0 — may reclaim its runner).
|
||||
# Otherwise (we're P1–P9 and the target isn't broken) we DON'T cancel; we
|
||||
# only yield to genuinely-higher-priority PRs in PASS 2.
|
||||
# Emit "number<TAB>head<TAB>prio<TAB>broken(0|1)" for every OTHER open PR
|
||||
# (broken => effective prio 10, the bottom, overriding any P0–P9 label).
|
||||
jq -r --argjson self "$SELF_PR" '
|
||||
.[] | select(.number != $self)
|
||||
| (any(.labels[]; .name == "broken")) as $b
|
||||
| [ .number, .headRefName,
|
||||
(if $b then 10
|
||||
else ([ .labels[] | .name | select(test("^P[0-9]$")) | ltrimstr("P") | tonumber ]
|
||||
| if length == 0 then 5 else min end) end),
|
||||
(if $b then 1 else 0 end) ]
|
||||
| @tsv' open_prs.json 2>/dev/null > others.tsv
|
||||
|
||||
cancelled_total=0
|
||||
while IFS=$'\t' read -r num head prio isbroken; do
|
||||
[ -n "${num:-}" ] || continue
|
||||
case "$prio" in ''|*[!0-9]*) prio=5 ;; esac
|
||||
[ "$isbroken" = "1" ] || isbroken=0
|
||||
|
||||
# Never touch an equal-or-higher-priority PR — only strictly lower.
|
||||
if [ "$prio" -le "$self_prio" ]; then
|
||||
echo "· PR #$num (P$prio): equal-or-higher priority — left untouched."
|
||||
continue
|
||||
fi
|
||||
|
||||
# Strictly lower, but only a P0 self OR a broken target is preemptible.
|
||||
if [ "$self_prio" -ne 0 ] && [ "$isbroken" != "1" ]; then
|
||||
echo "· PR #$num (P$prio): strictly lower, not broken, and we're not P0 — left to run (we don't cancel it)."
|
||||
continue
|
||||
fi
|
||||
|
||||
if [ "$self_prio" -eq 0 ]; then reason="P0 emergency"; else reason="target is 'broken'"; fi
|
||||
tag="P$prio"; [ "$isbroken" = "1" ] && tag="broken"
|
||||
echo "· PR #$num ($tag, head '$head'): preemptible ($reason) — checking for active CI runs."
|
||||
run_ids=$(active_run_ids_for_head "$head")
|
||||
|
||||
if [ -z "$run_ids" ]; then
|
||||
echo " no active CI runs."
|
||||
continue
|
||||
fi
|
||||
|
||||
while IFS= read -r run_id; do
|
||||
[ -n "$run_id" ] || continue
|
||||
[ "$run_id" = "${GITHUB_RUN_ID:-}" ] && continue # never cancel our own run
|
||||
if gh run cancel "$run_id" 2>err.txt; then
|
||||
echo " cancelled run $run_id (freed its runner)."
|
||||
cancelled_total=$((cancelled_total + 1))
|
||||
else
|
||||
echo "::warning::could not cancel run $run_id — likely already finished. $(cat err.txt 2>/dev/null)"
|
||||
fi
|
||||
done <<< "$run_ids"
|
||||
done < others.tsv
|
||||
echo "Preemption pass complete — cancelled $cancelled_total run(s)."
|
||||
|
||||
# ── PASS 2: BOUNDED HOLD-BACK — yield to strictly-higher, cancel NOTHING ──
|
||||
# P0 is top priority: nothing outranks an emergency, so it never yields.
|
||||
if [ "$self_prio" -eq 0 ]; then
|
||||
echo "P0 emergency — not yielding; proceeding immediately."
|
||||
exit 0
|
||||
fi
|
||||
|
||||
# Defer this PR's heavy jobs (which `needs: traffic-control`) while any
|
||||
# strictly-higher-priority OTHER open PR still has an active/queued CI run,
|
||||
# so those heavy jobs reach the runner queue first. Cancel NOTHING here.
|
||||
# Bounded by budget; on expiry proceed regardless (never block ourselves,
|
||||
# never hit the hard job timeout). Fail-open: any API hiccup => stop waiting.
|
||||
case "$HOLD_BACK_BUDGET_SECONDS" in ''|*[!0-9]*) HOLD_BACK_BUDGET_SECONDS=180 ;; esac
|
||||
case "$HOLD_BACK_POLL_SECONDS" in ''|*[!0-9]*) HOLD_BACK_POLL_SECONDS=15 ;; esac
|
||||
deadline=$(( $(date +%s) + HOLD_BACK_BUDGET_SECONDS ))
|
||||
echo "$prio_label — holding back up to ${HOLD_BACK_BUDGET_SECONDS}s for strictly-higher-priority PRs (no cancellation)."
|
||||
|
||||
while :; do
|
||||
remaining=$(( deadline - $(date +%s) ))
|
||||
if [ "$remaining" -le 0 ]; then
|
||||
echo "Hold-back budget elapsed — proceeding; higher-priority PRs got their head start."
|
||||
break
|
||||
fi
|
||||
|
||||
# Refresh so newly opened higher-priority PRs are seen mid-wait.
|
||||
if ! list_open_prs > open_prs.json 2>err.txt; then
|
||||
echo "::warning::Could not refresh open PRs — proceeding. $(cat err.txt 2>/dev/null)"
|
||||
break
|
||||
fi
|
||||
|
||||
# Strictly-higher-priority OTHER PRs (broken => 10, so a broken PR is never
|
||||
# higher than a non-broken one): "number<TAB>head<TAB>prio".
|
||||
jq -r --argjson self "$SELF_PR" --argjson me "$self_prio" '
|
||||
.[] | select(.number != $self)
|
||||
| { n: .number, h: .headRefName,
|
||||
p: (if any(.labels[]; .name == "broken") then 10
|
||||
else ([ .labels[] | .name | select(test("^P[0-9]$")) | ltrimstr("P") | tonumber ]
|
||||
| if length == 0 then 5 else min end) end) }
|
||||
| select(.p < $me)
|
||||
| [ .n, .h, .p ] | @tsv' open_prs.json 2>/dev/null > higher.tsv
|
||||
|
||||
if [ ! -s higher.tsv ]; then
|
||||
echo "No strictly-higher-priority open PRs — proceeding."
|
||||
break
|
||||
fi
|
||||
|
||||
blockers=""
|
||||
while IFS=$'\t' read -r num head prio; do
|
||||
[ -n "${num:-}" ] || continue
|
||||
if [ -n "$(active_run_ids_for_head "$head")" ]; then
|
||||
blockers="$blockers #$num(P$prio)"
|
||||
fi
|
||||
done < higher.tsv
|
||||
|
||||
if [ -z "$blockers" ]; then
|
||||
echo "No strictly-higher-priority PR has active CI runs — proceeding."
|
||||
break
|
||||
fi
|
||||
|
||||
sleep_s="$HOLD_BACK_POLL_SECONDS"
|
||||
[ "$remaining" -lt "$sleep_s" ] && sleep_s="$remaining"
|
||||
echo "Yielding to strictly-higher-priority PR(s) with active CI:${blockers} — re-checking in ${sleep_s}s (${remaining}s budget left)."
|
||||
[ "$sleep_s" -gt 0 ] && sleep "$sleep_s"
|
||||
done
|
||||
|
||||
echo "Hold-back complete — this PR's heavy jobs may now start."
|
||||
exit 0
|
||||
- name: Check out source
|
||||
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
|
||||
# gh is auto-configured from GH_TOKEN / GH_REPO; python3 is preinstalled on the
|
||||
# runner. The script guards every API call and always exits 0 (belt-and-braces
|
||||
# with continue-on-error), so it can never fail or block CI.
|
||||
- name: Apply runner priority (P0 preempts; P1–P9 hold back)
|
||||
continue-on-error: true
|
||||
run: python3 .github/scripts/traffic_control.py
|
||||
|
||||
debug-build:
|
||||
name: Debug build
|
||||
needs: traffic-control # order after runner-priority orchestration (P0/broken preempt; P1–P9 hold-back)
|
||||
needs: traffic-control # order after runner-priority orchestration (P0 preempts; P1–P9 hold-back)
|
||||
# x86_64: Linux-arm64 runners can't set up this SDK — android-actions/setup-android's sdkmanager
|
||||
# fails (exit 1) on the android-37.0 preview platform, and the emulator package has no arm64-Linux
|
||||
# build. Build/unit-test results are host-arch-independent anyway (R8/AGP/JVM); real arm64
|
||||
@@ -314,7 +137,7 @@ jobs:
|
||||
|
||||
unit-tests:
|
||||
name: Unit tests
|
||||
needs: traffic-control # order after runner-priority orchestration (P0/broken preempt; P1–P9 hold-back)
|
||||
needs: traffic-control # order after runner-priority orchestration (P0 preempts; P1–P9 hold-back)
|
||||
runs-on: ubuntu-latest
|
||||
steps:
|
||||
- name: Check out source
|
||||
@@ -371,7 +194,7 @@ jobs:
|
||||
|
||||
static-analysis:
|
||||
name: Static analysis
|
||||
needs: traffic-control # order after runner-priority orchestration (P0/broken preempt; P1–P9 hold-back)
|
||||
needs: traffic-control # order after runner-priority orchestration (P0 preempts; P1–P9 hold-back)
|
||||
runs-on: ubuntu-latest
|
||||
steps:
|
||||
- name: Check out source
|
||||
@@ -409,7 +232,7 @@ jobs:
|
||||
|
||||
e2e:
|
||||
name: E2E
|
||||
needs: traffic-control # order after runner-priority orchestration (P0/broken preempt; P1–P9 hold-back)
|
||||
needs: traffic-control # order after runner-priority orchestration (P0 preempts; P1–P9 hold-back)
|
||||
runs-on: ubuntu-latest
|
||||
strategy:
|
||||
fail-fast: false
|
||||
@@ -520,7 +343,7 @@ jobs:
|
||||
# 37 into the main `e2e` matrix and delete this job.
|
||||
e2e-preview:
|
||||
name: E2E (API 37 preview)
|
||||
needs: traffic-control # order after runner-priority orchestration (P0/broken preempt; P1–P9 hold-back)
|
||||
needs: traffic-control # order after runner-priority orchestration (P0 preempts; P1–P9 hold-back)
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 35
|
||||
env:
|
||||
|
||||
@@ -42,5 +42,9 @@ captures/
|
||||
# Kotlin
|
||||
.kotlin/
|
||||
|
||||
# Python (dev/CI helper scripts under .github/scripts, .claude/…)
|
||||
__pycache__/
|
||||
*.pyc
|
||||
|
||||
# Claude Code — personal settings (the shared settings.json IS committed)
|
||||
.claude/settings.local.json
|
||||
|
||||
Reference in New Issue
Block a user