diff --git a/.github/workflows/test-e2e.yml b/.github/workflows/test-e2e.yml index ac7c56b2496a..d285f3da0825 100644 --- a/.github/workflows/test-e2e.yml +++ b/.github/workflows/test-e2e.yml @@ -29,7 +29,7 @@ on: default: true type: boolean runner: - description: "Runner OS (auto follows MACOS_RUNNER_TESTS; tart-* choices use isolated VMs)" + description: "Runner OS (auto follows MACOS_RUNNER_TESTS and overflows to the 12vcpu macOS 26 pool only when the 6vcpu pool is backed up and 12vcpu is idle; tart-* choices use isolated VMs)" required: false default: "auto" type: choice @@ -43,6 +43,11 @@ on: - tart-dual - tart-small +# run-name and the concurrency group cannot read job outputs, so they spell +# `auto` as the 6vcpu default even when the runner job overflows the run to +# the 12vcpu pool. Every auto dispatch of one ref and filter still shares a +# group, so a newer one still cancels the older wherever it landed. run-e2e.sh +# names the pool it chose explicitly, so its titles are exact. concurrency: group: e2e-${{ (!inputs.runner || inputs.runner == 'auto') && (vars.MACOS_RUNNER_TESTS || 'blacksmith-6vcpu-macos-26') || inputs.runner }}-${{ inputs.ref || github.ref_name }}-${{ inputs.test_filter }} cancel-in-progress: true @@ -60,6 +65,60 @@ jobs: with: ref: ${{ inputs.ref }} + runner: + # The macOS pool build and test run on. `auto` follows MACOS_RUNNER_TESTS. + # The 12vcpu macOS 26 pool is reserved first for release and nightly + # builds, so on the 6vcpu default an E2E run overflows to 12vcpu only when + # the 6vcpu pool is backed up (at least CI_E2E_OVERFLOW_MIN_QUEUED other + # E2E runs in flight on it, default 4) and 12vcpu has spare room (no + # release or nightly run in flight, nothing queued, and fewer than + # CI_E2E_OVERFLOW_MAX_LARGE_RUNNING E2E runs on it, default 2). Any error + # reading the queue stays on 6vcpu, and CI_E2E_LARGE_POOL_OVERFLOW=0 turns + # overflow off. run-e2e.sh applies the same rule from the same script. + runs-on: ${{ vars.LINUX_RUNNER || 'blacksmith-4vcpu-ubuntu-2404' }} + timeout-minutes: 5 + # Reads one page each of in-progress and queued runs (two API calls at + # most; see e2e_runner_pool.py for the budget) to estimate pool demand. + permissions: + contents: read + actions: read + outputs: + label: ${{ steps.pool.outputs.label }} + steps: + # The helper comes from this workflow's revision, not the tested one, + # which may predate it. The workflow files tell it which release and + # nightly workflows never use macOS. + - name: Checkout pool helper + uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2 + with: + sparse-checkout: | + scripts/ci/e2e_runner_pool.py + scripts/ci/queue_janitor.py + .github/workflows/ + sparse-checkout-cone-mode: false + persist-credentials: false + + - name: Pick the macOS pool + id: pool + env: + GH_TOKEN: ${{ github.token }} + GH_REPO: ${{ github.repository }} + REQUESTED_RUNNER: ${{ inputs.runner }} + RUNNER_VARIABLE: ${{ vars.MACOS_RUNNER_TESTS }} + LARGE_POOL_OVERFLOW: ${{ vars.CI_E2E_LARGE_POOL_OVERFLOW }} + OVERFLOW_MIN_QUEUED: ${{ vars.CI_E2E_OVERFLOW_MIN_QUEUED }} + OVERFLOW_MAX_LARGE_RUNNING: ${{ vars.CI_E2E_OVERFLOW_MAX_LARGE_RUNNING }} + run: | + set -euo pipefail + label="$(python3 scripts/ci/e2e_runner_pool.py \ + --requested "$REQUESTED_RUNNER" \ + --variable "$RUNNER_VARIABLE" \ + --overflow "$LARGE_POOL_OVERFLOW" \ + --min-queued "$OVERFLOW_MIN_QUEUED" \ + --max-large-running "$OVERFLOW_MAX_LARGE_RUNNING")" + echo "label=$label" >> "$GITHUB_OUTPUT" + echo "Runner: $label (requested ${REQUESTED_RUNNER:-auto})" + filter: # Seconds on Linux, and it rejects a malformed selector before either # macOS job is scheduled. After the build/test split a bad filter would @@ -190,8 +249,8 @@ jobs: echo "Resolved selectors ($count): $selectors" build: - needs: [resolve-ref, filter] - runs-on: ${{ (!inputs.runner || inputs.runner == 'auto') && (vars.MACOS_RUNNER_TESTS || 'blacksmith-6vcpu-macos-26') || inputs.runner }} + needs: [resolve-ref, filter, runner] + runs-on: ${{ needs.runner.outputs.label }} # Reuse lists this contract's artifacts and downloads one from an earlier # run. Nothing else in this lane reads the Actions API, and nothing writes. permissions: @@ -214,13 +273,13 @@ jobs: # same Xcode build. The intra-run handoff uses the artifact id, not the # name, so nothing here depends on it. CMUX_SKIP_ZIG_BUILD: "1" - CMUX_PRODUCT_RUNNER: ${{ (!inputs.runner || inputs.runner == 'auto') && (vars.MACOS_RUNNER_TESTS || 'blacksmith-6vcpu-macos-26') || inputs.runner }} + CMUX_PRODUCT_RUNNER: ${{ needs.runner.outputs.label }} steps: - name: Validate Tart canary identity - if: ${{ startsWith((!inputs.runner || inputs.runner == 'auto') && (vars.MACOS_RUNNER_TESTS || 'blacksmith-6vcpu-macos-26') || inputs.runner, 'tart-') }} + if: ${{ startsWith(needs.runner.outputs.label, 'tart-') }} env: - REQUESTED_RUNNER: ${{ (!inputs.runner || inputs.runner == 'auto') && (vars.MACOS_RUNNER_TESTS || 'blacksmith-6vcpu-macos-26') || inputs.runner }} + REQUESTED_RUNNER: ${{ needs.runner.outputs.label }} RUNNER_CONTEXT_NAME: ${{ runner.name }} run: | set -euo pipefail @@ -613,7 +672,7 @@ jobs: esac test: - needs: [resolve-ref, filter, build] + needs: [resolve-ref, filter, runner, build] # scripts/ci/parallel_artifact_download.py reads this run's artifact # metadata through the Actions API; the workflow-level token is # contents-only. Narrow and additive, mirroring ci-macos.yml's admission @@ -621,7 +680,7 @@ jobs: permissions: contents: read actions: read - runs-on: ${{ (!inputs.runner || inputs.runner == 'auto') && (vars.MACOS_RUNNER_TESTS || 'blacksmith-6vcpu-macos-26') || inputs.runner }} + runs-on: ${{ needs.runner.outputs.label }} timeout-minutes: ${{ fromJSON(inputs.job_timeout || '45') }} env: CMUX_CI_MAX_MACOS_SDK_MAJOR: "26" @@ -639,9 +698,9 @@ jobs: steps: - name: Validate Tart canary identity - if: ${{ startsWith((!inputs.runner || inputs.runner == 'auto') && (vars.MACOS_RUNNER_TESTS || 'blacksmith-6vcpu-macos-26') || inputs.runner, 'tart-') }} + if: ${{ startsWith(needs.runner.outputs.label, 'tart-') }} env: - REQUESTED_RUNNER: ${{ (!inputs.runner || inputs.runner == 'auto') && (vars.MACOS_RUNNER_TESTS || 'blacksmith-6vcpu-macos-26') || inputs.runner }} + REQUESTED_RUNNER: ${{ needs.runner.outputs.label }} RUNNER_CONTEXT_NAME: ${{ runner.name }} run: | set -euo pipefail diff --git a/.github/workflows/test-macos-suite.yml b/.github/workflows/test-macos-suite.yml index 4675bac686a0..7f6a23f5a088 100644 --- a/.github/workflows/test-macos-suite.yml +++ b/.github/workflows/test-macos-suite.yml @@ -92,6 +92,9 @@ jobs: # The job token cannot list repository variables, and without the # runner the wrapper cannot see an identical run already in flight. CMUX_MACOS_RUNNER_TESTS: ${{ vars.MACOS_RUNNER_TESTS }} + CMUX_CI_E2E_LARGE_POOL_OVERFLOW: ${{ vars.CI_E2E_LARGE_POOL_OVERFLOW }} + CMUX_CI_E2E_OVERFLOW_MIN_QUEUED: ${{ vars.CI_E2E_OVERFLOW_MIN_QUEUED }} + CMUX_CI_E2E_OVERFLOW_MAX_LARGE_RUNNING: ${{ vars.CI_E2E_OVERFLOW_MAX_LARGE_RUNNING }} UNIT_TEST_SUITES: ${{ inputs.unit_test_suites }} TEST_REF: ${{ inputs.ref || github.sha }} TEST_TIMEOUT: ${{ inputs.test_timeout }} diff --git a/scripts/ci/dispatch-focused-test.py b/scripts/ci/dispatch-focused-test.py index aed9805c6440..ae08778de19b 100644 --- a/scripts/ci/dispatch-focused-test.py +++ b/scripts/ci/dispatch-focused-test.py @@ -16,10 +16,18 @@ from urllib.parse import quote import uuid +sys.path.insert(0, str(Path(__file__).resolve().parent)) +import e2e_runner_pool as pool # noqa: E402 +from e2e_runner_pool import LARGE_RUNNER, SMALL_RUNNER # noqa: E402 + REPO = "manaflow-ai/cmux" WORKFLOW = "test-e2e.yml" # `vars.MACOS_RUNNER_TESTS`, as passed by a workflow job; see default_runner(). VARIABLE_ENV = "CMUX_MACOS_RUNNER_TESTS" +# The overflow variables, passed the same way; see repository_variable(). +OVERFLOW_ENV = "CMUX_" + pool.OVERFLOW_VARIABLE +MIN_QUEUED_ENV = "CMUX_" + pool.MIN_QUEUED_VARIABLE +MAX_LARGE_RUNNING_ENV = "CMUX_" + pool.MAX_LARGE_RUNNING_VARIABLE ROOT = Path(__file__).resolve().parents[2] RUN_DISCOVERY_ATTEMPTS = 12 RUN_DISCOVERY_TIMEOUT_SECONDS = 60.0 @@ -38,13 +46,12 @@ "tart-dual", "tart-small", ) -# Half of all commits compile on the large macOS 26 SKU, so the two sizes are -# compared on real focused-run traffic rather than one benchmark. The split is -# keyed on the commit, not drawn at random: every dispatch at one commit lands -# on one pool, which is what in-flight reuse, the failed-selector refusal and -# the product contract all match on. -SMALL_RUNNER = "blacksmith-6vcpu-macos-26" -LARGE_RUNNER = "blacksmith-12vcpu-macos-26" +# An unpinned run overflows to the 12vcpu macOS 26 pool only when the 6vcpu +# pool is backed up and the 12vcpu pool, reserved first for release and +# nightly builds, has room. The rule lives in e2e_runner_pool.py, which +# test-e2e.yml runs too. Because the choice depends on the queue at dispatch +# time, not on the commit, the in-flight guards below look on both pools. +OVERFLOW_POOLS = (SMALL_RUNNER, LARGE_RUNNER) # GitHub rejects a concurrency group longer than this as a workflow file # issue: the run is created with no jobs and no message saying why. MAX_CONCURRENCY_GROUP = 400 @@ -186,6 +193,70 @@ def parse_run_name(title: str) -> tuple[list[str], str, str] | None: return [part.strip() for part in head.split(",")], runner.strip(), ref +_UNLISTED = object() +_listed: object = _UNLISTED + + +def listed_variables() -> dict[str, str] | None: + """Repository variables by name, read once, or None when unreadable.""" + global _listed + if _listed is _UNLISTED: + try: + payload = output( + "gh", "variable", "list", "--repo", REPO, "--json", "name,value", + timeout=PRIOR_ATTEMPT_TIMEOUT_SECONDS, + ) + variables = json.loads(payload) + except (subprocess.SubprocessError, OSError, ValueError, json.JSONDecodeError): + variables = None + if isinstance(variables, list): + _listed = { + str(entry["name"]): str(entry.get("value", "")) + for entry in variables + if isinstance(entry, dict) and "name" in entry + } + else: + _listed = None + return _listed # type: ignore[return-value] + + +def repository_variable(name: str, env_name: str) -> str | None: + """An overflow variable's value; None or empty means unset (the default). + + A workflow job cannot list variables and passes them in CMUX_* instead. + A job that passed MACOS_RUNNER_TESTS but not this one predates it, so it + gets the default. Elsewhere an unreadable listing also means the default: + overflow is still bounded by the queue it reads, and fails to 6vcpu. + """ + if env_name in os.environ: + return os.environ[env_name] + if VARIABLE_ENV in os.environ: + return None + return (listed_variables() or {}).get(name) + + +class GhApi(pool.queue_janitor.GitHub): + """The queue janitor's GitHub client, speaking through `gh api`. + + `gh` carries the caller's own credentials, locally or in a workflow job, + so this needs no token handling of its own. + """ + + def __init__(self) -> None: + super().__init__("", REPO) + + def request(self, method: str, path: str, body=None): + self.calls += 1 + try: + payload = output( + "gh", "api", "--method", method, path.lstrip("/"), + timeout=PRIOR_ATTEMPT_TIMEOUT_SECONDS, + ) + return json.loads(payload) if payload else {} + except (subprocess.SubprocessError, OSError, ValueError) as error: + raise RuntimeError(f"{method} {path.split('?')[0]} failed") from error + + def default_runner() -> str | None: """The label `runner: auto` resolves to, or None when it cannot be known. @@ -207,22 +278,12 @@ def default_runner() -> str | None: if value: return value else: - try: - payload = output( - "gh", "variable", "list", "--repo", REPO, "--json", "name,value", - timeout=PRIOR_ATTEMPT_TIMEOUT_SECONDS, - ) - variables = json.loads(payload) - except (subprocess.SubprocessError, OSError, ValueError, json.JSONDecodeError): - return None - if not isinstance(variables, list): + variables = listed_variables() + if variables is None: return None - for entry in variables: - if isinstance(entry, dict) and entry.get("name") == "MACOS_RUNNER_TESTS": - value = str(entry.get("value", "")).strip() - if value: - return value - break + value = variables.get("MACOS_RUNNER_TESTS", "").strip() + if value: + return value try: workflow = (ROOT / ".github/workflows" / WORKFLOW).read_text() except OSError: @@ -233,25 +294,51 @@ def default_runner() -> str | None: return literal.group(1) if literal else None -def routed_runner(commit: str, default: str | None) -> str | None: - """The pool an unpinned dispatch at `commit` runs on. +def routed_runner(default: str | None) -> str | None: + """The pool an unpinned dispatch runs on now; see e2e_runner_pool. + + Only called when a dispatch is about to happen, so a run reused from the + history spends no API calls on the queue. + """ + return pool.auto_runner( + default, + enabled=pool.overflow_enabled( + repository_variable(pool.OVERFLOW_VARIABLE, OVERFLOW_ENV)), + limits=pool.thresholds( + repository_variable(pool.MIN_QUEUED_VARIABLE, MIN_QUEUED_ENV), + repository_variable(pool.MAX_LARGE_RUNNING_VARIABLE, MAX_LARGE_RUNNING_ENV), + ), + measure=lambda limits: pool.measure_load( + GhApi(), REPO, limits, workflows_dir=ROOT / ".github" / "workflows"), + log=lambda message: print(f"Runner pool: {message}", file=sys.stderr, flush=True), + ) + + +def candidate_runners(runner: str | None, pinned: bool) -> tuple[str, ...]: + """Every pool a dispatch with this runner could land on. - Only the free default is split. A repository variable naming any other - pool is an admin decision, and it wins unchanged. + A pinned runner is exact. An unpinned dispatch on the 6vcpu default may + overflow to the 12vcpu pool, so a run on either one already answers it. + Empty means the default could not be established. """ - if default == SMALL_RUNNER and int(commit[-1], 16) % 2: - return LARGE_RUNNER - return default + if runner is None: + return () + if not pinned and runner == SMALL_RUNNER: + return OVERFLOW_POOLS + return (runner,) def attempts( - runs: list[dict], commit: str, selector: str, runner: str | None = None + runs: list[dict], commit: str, selector: str, + runner: str | tuple[str, ...] | None = None, ) -> list[dict]: """Runs of this selector at this exact commit, newest first. - `runner` narrows to one pool. None means every pool, which is what the - repeat guard wants: a red result is usually a property of the commit. + `runner` narrows to one pool, or to any of several. None means every + pool, which is what the repeat guard wants: a red result is usually a + property of the commit. """ + runners = (runner,) if isinstance(runner, str) else runner found = [] for run in runs: parsed = parse_run_name(str(run.get("displayTitle", ""))) @@ -260,7 +347,7 @@ def attempts( selectors, run_runner, ref = parsed if ref != commit or selector not in selectors: continue - if runner is not None and run_runner != runner: + if runners is not None and run_runner not in runners: continue found.append(run) return found @@ -286,7 +373,7 @@ def prior_attempts( def live_attempts( - runs: list[dict], commit: str, selector: str, runner: str + runs: list[dict], commit: str, selector: str, runner: str | tuple[str, ...] ) -> list[dict]: """Attempts GitHub has accepted that have not reported a conclusion yet. @@ -298,7 +385,8 @@ def live_attempts( not collide, and instead pays a second full compile of identical source to answer a question already in flight. - `runner` is required and exact. A run on another pool shares neither the + `runner` is required and exact: one pool, or the pools an unpinned + dispatch could overflow between. A run on another pool shares neither the concurrency group nor the question: reusing its result would report macOS 15's answer to someone who asked about macOS 26. """ @@ -308,6 +396,11 @@ def live_attempts( ] +def parsed_runner(run: dict) -> str: + parsed = parse_run_name(str(run.get("displayTitle", ""))) + return parsed[1] if parsed else "an unknown runner" + + def watchable(run: dict) -> bool: """Whether this history entry carries enough to point a caller at the run. @@ -426,14 +519,16 @@ def main() -> int: if args.ref is None and commit != requested_ref: raise ValueError("GitHub revision differs from local HEAD; push the intended commit first") - # Which pool this dispatch will actually land on. None means the answer - # could not be established, and the in-flight guards below stay silent - # rather than compare against a runner they guessed. + # Which pools this dispatch could land on. Empty means the answer could + # not be established, and the in-flight guards below stay silent rather + # than compare against a runner they guessed. The queue is read only once + # the guards have decided to dispatch. pinned = args.runner not in (None, "auto") - runner = args.runner if pinned else routed_runner(commit, default_runner()) + default = args.runner if pinned else default_runner() + pools = candidate_runners(default, pinned) # test-e2e.yml groups on "e2e---". When the pool # is unknown, measure against the longest label in the runner dropdown. - label = runner or max(RUNNERS, key=len) + label = max(pools or RUNNERS, key=len) group_length = len(f"e2e-{label}-{commit}-{test_filter}") if group_length > MAX_CONCURRENCY_GROUP: parser.error( @@ -444,11 +539,12 @@ def main() -> int: if not args.force: history = recent_dispatches() - if runner is not None: + if pools: # An identical dispatch is already answering this exact question on - # this exact pool. Attach to it instead of cancelling it: the - # concurrency group keyed on runner/ref/test_filter would kill the - # run mid-compile and start the same compile again from cold. + # a pool this one could land on. Attach to it instead of cancelling + # it or paying a second compile on the other macOS 26 pool: the + # concurrency group keyed on runner/ref/test_filter would kill a + # same-pool run mid-compile and start the compile again from cold. requested = set(args.test_filter) running = [ run for run in history @@ -456,14 +552,14 @@ def main() -> int: and watchable(run) and (parsed := parse_run_name(str(run.get("displayTitle", "")))) is not None and parsed[2] == commit - and parsed[1] == runner + and parsed[1] in pools and set(parsed[0]) == requested ] if running: live = running[0] print( f"{test_filter} is already {live['status']} at {commit} " - f"on {runner}; reusing that run instead of dispatching.", + f"on {parsed_runner(live)}; reusing that run instead of dispatching.", flush=True, ) print(f"Run: {live['url']}", flush=True) @@ -477,12 +573,12 @@ def main() -> int: # Refuse per entry: one already-red selector makes the whole batch a # reprint of a known failure, and the compile it would pay for is shared. for entry in args.test_filter: - live = [run for run in live_attempts(history, commit, entry, runner) - if watchable(run)] if runner is not None else [] + live = [run for run in live_attempts(history, commit, entry, pools) + if watchable(run)] if pools else [] if live: raise ValueError( f"{entry} is already {live[0]['status']} at {commit} on " - f"{runner}, in {live[0]['url']}, under a different set of " + f"{parsed_runner(live[0])}, in {live[0]['url']}, under a different set of " "selectors. Dispatching now would compile identical source " "a second time to answer a question already in flight. Wait " "for that run, dispatch the remaining selectors on their " @@ -505,6 +601,7 @@ def main() -> int: "the new commit. Pass --force to dispatch anyway." ) + runner = args.runner if pinned else routed_runner(default) dispatch_id = uuid.uuid4().hex video = not args.no_video and test_target != "cmuxTests" fields = { @@ -517,7 +614,9 @@ def main() -> int: } if args.runner is not None: fields["runner"] = args.runner - if not pinned and runner == LARGE_RUNNER: + # Name the pool chosen here, so the run title carries the pool the guards + # above match on and test-e2e.yml does not read the queue a second time. + if not pinned and runner in OVERFLOW_POOLS: fields["runner"] = runner command = ["gh", "workflow", "run", WORKFLOW, "--repo", REPO] if args.workflow_ref: diff --git a/scripts/ci/e2e_runner_pool.py b/scripts/ci/e2e_runner_pool.py new file mode 100644 index 000000000000..717844119926 --- /dev/null +++ b/scripts/ci/e2e_runner_pool.py @@ -0,0 +1,305 @@ +#!/usr/bin/env python3 +"""Pick the macOS pool an E2E run lands on. + +test-e2e.yml and scripts/ci/dispatch-focused-test.py both call this, so a +workflow started from the Actions UI, `gh workflow run`, or run-e2e.sh applies +one rule. + +`runner: auto` means `vars.MACOS_RUNNER_TESTS` when it names a pool, else the +6vcpu macOS 26 pool. The 12vcpu macOS 26 pool ("macOS large") is reserved +first for release and nightly builds (nightly.yml's build job, the release +workflows), so E2E must never be the reason it is backed up. An `auto` run on +the 6vcpu default therefore overflows to the 12vcpu pool only when the 6vcpu +pool is backed up and the 12vcpu pool has spare room: + + queued(6vcpu) >= vars.CI_E2E_OVERFLOW_MIN_QUEUED (default 4) + queued(12vcpu) == 0 + running(12vcpu) < vars.CI_E2E_OVERFLOW_MAX_LARGE_RUNNING (default 2) + +Anything else stays on the 6vcpu pool: any error reading the queue, a listing +that may be truncated, or an invalid threshold. `vars.CI_E2E_LARGE_POOL_OVERFLOW +== '0'` turns overflow off. An explicit runner, or a variable naming any other +pool, is never rerouted. + +API budget: at most two requests per decision, never retried or polled. The +GITHUB_TOKEN allows about 1000 requests an hour for the whole repository and +E2E dispatches can run to dozens an hour, so the queue janitor's per-run job +listings (one request per in-flight run) are out of reach. The decision reads +one page of in-progress runs and one page of queued runs and attributes pool +demand from run metadata alone: + + * an E2E run's title names its pool (" on @ "), so an + in-flight E2E run titled with the 12vcpu pool counts as running there, or + queued there while the run itself is queued; + * an in-flight release or nightly run (the queue janitor's reserved + workflow names, minus workflows that never use macOS) counts as queued on + the 12vcpu pool, so E2E yields to it whether or not its macOS job has + started; + * queued(6vcpu) is estimated as the other E2E runs in flight on the 6vcpu + pool. That is E2E's own demand on the pool, not the pool's whole job + queue, which only job listings can show. + +The in-progress page is read first, and when it already rules overflow out +the queued page is never requested. A full page (100 runs) may hide more, so +it counts as unknown. A 6vcpu `auto` run started from the Actions UI that +overflowed is still titled 6vcpu (run-name cannot read job outputs), so it +counts as 6vcpu demand; run-e2e.sh names its pool, so its titles are exact. +""" +from __future__ import annotations + +import argparse +import dataclasses +import os +from collections.abc import Callable, Mapping, Sequence +from pathlib import Path +import re +import sys +from typing import Any, Protocol + +sys.path.insert(0, str(Path(__file__).resolve().parent)) +import queue_janitor # noqa: E402 + +SMALL_RUNNER = "blacksmith-6vcpu-macos-26" +LARGE_RUNNER = "blacksmith-12vcpu-macos-26" +E2E_WORKFLOW_PATH = ".github/workflows/test-e2e.yml" + +# Repository variables. The kill switch turns overflow off when set to "0". +OVERFLOW_VARIABLE = "CI_E2E_LARGE_POOL_OVERFLOW" +MIN_QUEUED_VARIABLE = "CI_E2E_OVERFLOW_MIN_QUEUED" +MAX_LARGE_RUNNING_VARIABLE = "CI_E2E_OVERFLOW_MAX_LARGE_RUNNING" +DEFAULT_MIN_QUEUED = 4 +DEFAULT_MAX_LARGE_RUNNING = 2 + +# The whole API budget of one decision; see the module docstring. +MAX_API_CALLS = 2 +PAGE_SIZE = 100 +# Workflows whose macOS jobs have first claim on the 12vcpu pool. +RESERVED_WORKFLOW = re.compile(r"release|nightly", re.IGNORECASE) +TITLE_RUNNER = re.compile(r" on (?P\S+) @ ") + +WORKFLOWS_DIR = Path(__file__).resolve().parents[2] / ".github" / "workflows" + + +@dataclasses.dataclass(frozen=True) +class Thresholds: + min_queued: int = DEFAULT_MIN_QUEUED + max_large_running: int = DEFAULT_MAX_LARGE_RUNNING + + +@dataclasses.dataclass(frozen=True) +class PoolLoad: + small_queued: int = 0 + large_queued: int = 0 + large_running: int = 0 + + +class ApiClient(Protocol): + """queue_janitor.GitHub's request(), or anything with the same shape.""" + + def request(self, method: str, path: str) -> Any: ... + + +def overflow_enabled(value: str | None) -> bool: + """Whether overflow is on. Unset or anything but "0" is on.""" + return (value or "").strip() != "0" + + +def thresholds(min_queued: str | None, max_large_running: str | None) -> Thresholds | None: + """Thresholds from repository variables; None when either is invalid. + + Blank means the default. A minimum below one would send every run to the + 12vcpu pool whenever it is idle, so it counts as invalid, and an invalid + value keeps E2E off the 12vcpu pool rather than guessing. + """ + try: + queued = int(min_queued) if (min_queued or "").strip() else DEFAULT_MIN_QUEUED + running = (int(max_large_running) if (max_large_running or "").strip() + else DEFAULT_MAX_LARGE_RUNNING) + except ValueError: + return None + if queued < 1 or running < 0: + return None + return Thresholds(queued, running) + + +def overflows(load: PoolLoad | None, limits: Thresholds) -> bool: + """The overflow rule itself. An unknown load never overflows.""" + return ( + load is not None + and load.small_queued >= limits.min_queued + and load.large_queued == 0 + and load.large_running < limits.max_large_running + ) + + +def title_runner(run: Mapping[str, Any]) -> str | None: + """The pool an E2E run's title names, or None for any other run.""" + if str(run.get("path") or "").split("@", 1)[0] != E2E_WORKFLOW_PATH: + return None + match = TITLE_RUNNER.search(str(run.get("display_title") or "")) + return match.group("runner") if match else None + + +def is_reserved(run: Mapping[str, Any], linux_only: frozenset[str]) -> bool: + """A release or nightly run that may want the 12vcpu pool.""" + path = str(run.get("path") or "").split("@", 1)[0] + if path in linux_only: + return False + return bool(RESERVED_WORKFLOW.search(f"{run.get('name') or ''} {path}")) + + +def add_runs( + load: PoolLoad, + runs: Sequence[Mapping[str, Any]], + *, + queued: bool, + linux_only: frozenset[str], + exclude_run_id: int | None, +) -> PoolLoad: + small, large_queued, large_running = load.small_queued, load.large_queued, load.large_running + for run in runs: + if run.get("id") == exclude_run_id: + continue + if is_reserved(run, linux_only): + large_queued += 1 + continue + runner = title_runner(run) + if runner == LARGE_RUNNER: + if queued: + large_queued += 1 + else: + large_running += 1 + elif runner == SMALL_RUNNER: + small += 1 + return PoolLoad(small, large_queued, large_running) + + +def measure_load( + client: ApiClient, + repo: str, + limits: Thresholds, + *, + workflows_dir: Path = WORKFLOWS_DIR, + exclude_run_id: int | None = None, +) -> PoolLoad | None: + """Pool demand from at most MAX_API_CALLS requests, or None when unknown. + + Raises RuntimeError (from the client) on an API failure. + """ + linux_only = queue_janitor.linux_only_workflow_paths(workflows_dir) + load = PoolLoad() + for status in ("in_progress", "queued"): + payload = client.request("GET", f"/repos/{repo}/actions/runs?status={status}&per_page={PAGE_SIZE}") + runs = payload.get("workflow_runs") if isinstance(payload, Mapping) else None + if not isinstance(runs, list): + raise RuntimeError(f"unexpected {status} runs payload") + if len(runs) >= PAGE_SIZE: + return None + load = add_runs(load, [run for run in runs if isinstance(run, Mapping)], + queued=status == "queued", linux_only=linux_only, + exclude_run_id=exclude_run_id) + if load.large_queued or load.large_running >= limits.max_large_running: + # Already ruled out; the second page cannot change that. + break + return load + + +def auto_runner( + default: str | None, + *, + enabled: bool, + limits: Thresholds | None, + measure: Callable[[Thresholds], PoolLoad | None], + log: Callable[[str], None] = lambda message: None, +) -> str | None: + """The pool an unpinned run lands on, given what `auto` means. + + Only the 6vcpu default overflows. None stays None: a caller that could + not establish the default must not act on a guess. `measure` is called + only when overflow is possible, and any error it raises stays on 6vcpu. + """ + if default != SMALL_RUNNER: + return default + if not enabled: + log(f"{OVERFLOW_VARIABLE}=0; staying on {SMALL_RUNNER}") + return default + if limits is None: + log(f"invalid {MIN_QUEUED_VARIABLE} or {MAX_LARGE_RUNNING_VARIABLE}; staying on {SMALL_RUNNER}") + return default + try: + load = measure(limits) + except Exception as error: # noqa: BLE001 - every failure is fail-safe + log(f"could not read the runner queue ({error}); staying on {SMALL_RUNNER}") + return default + if load is None: + log(f"too many in-flight runs to read in one page; staying on {SMALL_RUNNER}") + return default + chosen = LARGE_RUNNER if overflows(load, limits) else default + log( + f"E2E waiting on {SMALL_RUNNER}: {load.small_queued} (overflow at >= {limits.min_queued}); " + f"{LARGE_RUNNER} queued {load.large_queued}, running {load.large_running} " + f"(max {limits.max_large_running}) -> {chosen}" + ) + return chosen + + +def resolve( + requested: str | None, + variable: str | None, + *, + overflow: str | None, + min_queued: str | None, + max_large_running: str | None, + measure: Callable[[Thresholds], PoolLoad | None], + log: Callable[[str], None] = lambda message: None, +) -> str: + """The runner label for a workflow run, from its inputs and variables.""" + requested = (requested or "").strip() + if requested and requested != "auto": + return requested + default = (variable or "").strip() or SMALL_RUNNER + return auto_runner( + default, + enabled=overflow_enabled(overflow), + limits=thresholds(min_queued, max_large_running), + measure=measure, + log=log, + ) or SMALL_RUNNER + + +def main(argv: Sequence[str] | None = None, env: Mapping[str, str] | None = None) -> int: + env = os.environ if env is None else env + parser = argparse.ArgumentParser(description=__doc__.splitlines()[0]) + parser.add_argument("--requested", default="", help="the workflow's runner input") + parser.add_argument("--variable", default="", help="vars.MACOS_RUNNER_TESTS") + parser.add_argument("--overflow", default="", help=f"vars.{OVERFLOW_VARIABLE}") + parser.add_argument("--min-queued", default="", help=f"vars.{MIN_QUEUED_VARIABLE}") + parser.add_argument("--max-large-running", default="", help=f"vars.{MAX_LARGE_RUNNING_VARIABLE}") + parser.add_argument("--workflows-dir", type=Path, default=WORKFLOWS_DIR) + args = parser.parse_args(argv) + + repo = env.get("GH_REPO") or env.get("GITHUB_REPOSITORY") or "" + token = env.get("GH_TOKEN") or env.get("GITHUB_TOKEN") + run_id = (env.get("GITHUB_RUN_ID") or "").strip() + + def measure(limits: Thresholds) -> PoolLoad | None: + if not token or not repo: + raise RuntimeError("GH_TOKEN and GH_REPO are required") + return measure_load( + queue_janitor.GitHub(token, repo), repo, limits, + workflows_dir=args.workflows_dir, + exclude_run_id=int(run_id) if run_id.isdigit() else None, + ) + + print(resolve( + args.requested, args.variable, + overflow=args.overflow, min_queued=args.min_queued, + max_large_running=args.max_large_running, + measure=measure, + log=lambda message: print(message, file=sys.stderr), + )) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_ci_e2e_compilation_cache.py b/tests/test_ci_e2e_compilation_cache.py index 00edb21d08c1..9aabdd55444e 100644 --- a/tests/test_ci_e2e_compilation_cache.py +++ b/tests/test_ci_e2e_compilation_cache.py @@ -173,7 +173,7 @@ def test_one_build_serves_the_test_job_and_every_retry(self): outputs = WORKFLOW['jobs']['build']['outputs'] self.assertEqual(outputs['artifact_id'], '${{ steps.upload-product.outputs.artifact-id }}') self.assertEqual(outputs['sha256'], '${{ steps.package.outputs.sha256 }}') - self.assertEqual(WORKFLOW['jobs']['test']['needs'], ['resolve-ref', 'filter', 'build']) + self.assertEqual(WORKFLOW['jobs']['test']['needs'], ['resolve-ref', 'filter', 'runner', 'build']) def test_the_test_job_verifies_the_product_before_using_it(self): # A transport is allowed to miss; it is not allowed to hand over diff --git a/tests/test_ci_self_hosted_guard.sh b/tests/test_ci_self_hosted_guard.sh index bacd6a66309d..a7857baeaeff 100755 --- a/tests/test_ci_self_hosted_guard.sh +++ b/tests/test_ci_self_hosted_guard.sh @@ -172,8 +172,8 @@ check_e2e_runner_fallbacks() { if ! awk ' /^[[:space:]]*- name: Validate Tart canary identity$/ { in_tart_step=1; next } in_tart_step && /^ - / { in_tart_step=0; in_runner_reject=0; in_marker_reject=0 } - in_tart_step && /startsWith\(\(!inputs\.runner \|\| inputs\.runner == '\''auto'\''\) && \(vars\.MACOS_RUNNER_[A-Z0-9_]+ \|\| '\''blacksmith-6vcpu-macos-[0-9]+'\''\) \|\| inputs\.runner, '\''tart-'\''\)/ { saw_effective_runner=1 } - in_tart_step && /REQUESTED_RUNNER:.*inputs\.runner/ { saw_requested_runner=1 } + in_tart_step && /startsWith\(needs\.runner\.outputs\.label, '\''tart-'\''\)/ { saw_effective_runner=1 } + in_tart_step && /REQUESTED_RUNNER: \$\{\{ needs\.runner\.outputs\.label \}\}/ { saw_requested_runner=1 } in_tart_step && /RUNNER_CONTEXT_NAME: \$\{\{ runner\.name \}\}/ { saw_runner_context=1 } in_tart_step && /tart-cmux-\*/ { saw_runner_pattern=1 } in_tart_step && /^[[:space:]]*\*\)$/ { in_runner_reject=1 } diff --git a/tests/test_run_e2e.py b/tests/test_run_e2e.py index 15cf0c9a4c07..63d053760746 100644 --- a/tests/test_run_e2e.py +++ b/tests/test_run_e2e.py @@ -20,13 +20,47 @@ ).group(1) HEAD = "a" * 40 REMOTE_HEAD = "b" * 40 -FAKE_GH = r'''#!/usr/bin/env python3 +SMALL = "blacksmith-6vcpu-macos-26" +LARGE = "blacksmith-12vcpu-macos-26" + + +def e2e_run(runner, run_id, *, status="in_progress"): + """An in-flight test-e2e.yml run as the Actions runs listing returns it.""" + return { + "id": run_id, "status": status, "name": "E2E test with video recording", + "path": ".github/workflows/test-e2e.yml", "event": "workflow_dispatch", + "display_title": f"cmuxTests/Other{run_id} on {runner} @ {'c' * 40} [x{run_id}]", + } + + +def queue(*, small=0, large_running=0, large_queued=0, reserved=0): + """A runs listing, keyed by status, with this much demand per pool.""" + in_progress = [e2e_run(SMALL, 100 + n) for n in range(small)] + in_progress += [e2e_run(LARGE, 200 + n) for n in range(large_running)] + in_progress += [{ + "id": 300 + n, "status": "in_progress", "name": "Nightly", + "path": ".github/workflows/nightly.yml", "event": "schedule", + "display_title": "Nightly", + } for n in range(reserved)] + queued = [e2e_run(LARGE, 400 + n, status="queued") for n in range(large_queued)] + return {"in_progress": in_progress, "queued": queued} + + +BACKED_UP = json.dumps(queue(small=4)) +FAKE_GH =r'''#!/usr/bin/env python3 import json, os, pathlib, sys args = sys.argv[1:] root = pathlib.Path(os.environ["LAUNCHER_TEST_DIR"]) with (root / "calls.jsonl").open("a") as f: f.write(json.dumps(args) + "\n") -if args[0] == "api": +if args[0] == "api" and "/actions/runs?" in args[-1]: + # The runner-pool decision's queue read: one page per run status. + if os.environ.get("LAUNCHER_QUEUE_FAIL"): + sys.exit(1) + status = args[-1].split("status=", 1)[1].split("&", 1)[0] + queue = json.loads(os.environ.get("LAUNCHER_QUEUE", "{}")) + print(json.dumps({"workflow_runs": queue.get(status, [])})) +elif args[0] == "api": if os.environ.get("LAUNCHER_MISSING_COMMIT"): sys.exit(1) print(json.dumps({"sha": "b" * 40 if "topic%2Ffix" in args[1] else "a" * 40})) @@ -110,47 +144,133 @@ def test_explicit_runner_reaches_the_workflow_dispatch(self): with self.subTest(runner=runner): result = self.launch("cmuxTests/ExampleTests", "--runner", runner) self.assertEqual(result.returncode, 0, result.stderr) - self.assertEqual(self.dispatch()["runner"], runner) + # `auto` is decided here and named, so the title is exact. + expected = SMALL if runner == "auto" else runner + self.assertEqual(self.dispatch()["runner"], expected) + + def queue_reads(self): + return [call for call in self.calls() + if call[:1] == ["api"] and "/actions/runs?" in call[-1]] + + def test_an_idle_queue_keeps_the_6vcpu_pool_whatever_the_commit(self): + # The large pool is reserved first for release and nightly builds, so + # no commit goes there by default: REMOTE_HEAD ends in b, which the + # earlier parity split sent to the 12vcpu pool. + for ref in ("topic/fix", "main"): + with self.subTest(ref=ref): + self.setUp() + result = self.launch("cmuxTests/ExampleTests", "--ref", ref) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertEqual(self.dispatch()["runner"], SMALL) - def test_default_runner_keeps_the_workflow_default(self): - result = self.launch("cmuxTests/ExampleTests") + def test_a_backed_up_6vcpu_pool_overflows_to_an_idle_12vcpu_pool(self): + result = self.launch("cmuxTests/ExampleTests", LAUNCHER_QUEUE=BACKED_UP) self.assertEqual(result.returncode, 0, result.stderr) - self.assertNotIn("runner", self.dispatch()) + self.assertEqual(self.dispatch()["runner"], LARGE) + self.assertLessEqual(len(self.queue_reads()), 2) + + def test_a_busy_12vcpu_pool_keeps_e2e_on_6vcpu(self): + for busy in ( + queue(small=9, large_queued=1), + queue(small=9, large_running=2), + queue(small=9, reserved=1), + ): + with self.subTest(queue=busy): + self.setUp() + result = self.launch("cmuxTests/ExampleTests", LAUNCHER_QUEUE=json.dumps(busy)) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertEqual(self.dispatch()["runner"], SMALL) + self.assertLessEqual(len(self.queue_reads()), 2) + + def test_an_unreadable_queue_keeps_e2e_on_6vcpu(self): + result = self.launch( + "cmuxTests/ExampleTests", LAUNCHER_QUEUE=BACKED_UP, LAUNCHER_QUEUE_FAIL="1", + ) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertEqual(self.dispatch()["runner"], SMALL) + self.assertIn("staying on", result.stderr) - def test_odd_commits_compile_on_the_large_sku(self): - # REMOTE_HEAD ends in b, so it is routed; HEAD ends in a, so it is not. - result = self.launch("cmuxTests/ExampleTests", "--ref", "topic/fix") + def test_the_overflow_thresholds_are_repository_variables(self): + result = self.launch( + "cmuxTests/ExampleTests", + LAUNCHER_QUEUE=json.dumps(queue(small=2)), + LAUNCHER_VARIABLES=json.dumps([ + {"name": "CI_E2E_OVERFLOW_MIN_QUEUED", "value": "2"}, + ]), + ) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertEqual(self.dispatch()["runner"], LARGE) + self.setUp() + result = self.launch( + "cmuxTests/ExampleTests", LAUNCHER_QUEUE=BACKED_UP, + LAUNCHER_VARIABLES=json.dumps([ + {"name": "CI_E2E_OVERFLOW_MAX_LARGE_RUNNING", "value": "0"}, + ]), + ) self.assertEqual(result.returncode, 0, result.stderr) - self.assertEqual(self.dispatch()["runner"], "blacksmith-12vcpu-macos-26") + self.assertEqual(self.dispatch()["runner"], SMALL) def test_an_explicit_runner_is_never_rerouted(self): result = self.launch( "cmuxTests/ExampleTests", "--ref", "topic/fix", - "--runner", "blacksmith-6vcpu-macos-26", + "--runner", SMALL, LAUNCHER_QUEUE=BACKED_UP, ) self.assertEqual(result.returncode, 0, result.stderr) - self.assertEqual(self.dispatch()["runner"], "blacksmith-6vcpu-macos-26") + self.assertEqual(self.dispatch()["runner"], SMALL) + self.assertEqual(self.queue_reads(), []) - def test_an_admin_runner_variable_is_never_split(self): + def test_an_admin_runner_variable_is_never_overflowed(self): result = self.launch( - "cmuxTests/ExampleTests", "--ref", "topic/fix", + "cmuxTests/ExampleTests", "--ref", "topic/fix", LAUNCHER_QUEUE=BACKED_UP, LAUNCHER_VARIABLES=json.dumps([ {"name": "MACOS_RUNNER_TESTS", "value": "blacksmith-6vcpu-macos-15"}, ]), ) self.assertEqual(result.returncode, 0, result.stderr) self.assertNotIn("runner", self.dispatch()) + self.assertEqual(self.queue_reads(), []) + + def test_an_unpinned_dispatch_reuses_an_in_flight_run_on_either_macos_26_pool(self): + # Where auto lands depends on the queue at dispatch time, so the same + # commit and filter may already be running on the other pool. Reusing + # it costs no compile and no queue read. + for runner in (SMALL, LARGE): + for queue_state in ("{}", BACKED_UP): + with self.subTest(runner=runner, queue=queue_state): + self.setUp() + result = self.launch( + "cmuxTests/ExampleTests", + LAUNCHER_PRIOR_RUNS=self._live(runner=runner), + LAUNCHER_QUEUE=queue_state, + ) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertIn("reusing that run", result.stdout) + self.assertIn(f"on {runner}", result.stdout) + self.assertFalse((self.root / "dispatch.json").exists()) + self.assertEqual(self.queue_reads(), []) + + def test_an_overlapping_run_on_the_other_macos_26_pool_is_refused(self): + live = self._live(selector="cmuxTests/ExampleTests,cmuxTests/OtherTests", runner=LARGE) + result = self.launch("cmuxTests/ExampleTests", LAUNCHER_PRIOR_RUNS=live) + self.assertNotEqual(result.returncode, 0) + self.assertIn(f"already in_progress at {HEAD} on {LARGE}", result.stderr) + self.assertFalse((self.root / "dispatch.json").exists(), "must not dispatch") - def test_a_routed_commit_reuses_its_in_flight_run_on_the_large_sku(self): + def test_a_failure_on_the_12vcpu_pool_refuses_an_unpinned_repeat(self): result = self.launch( - "cmuxTests/ExampleTests", "--ref", "topic/fix", - LAUNCHER_PRIOR_RUNS=self._live( - commit=REMOTE_HEAD, runner="blacksmith-12vcpu-macos-26", - ), + "cmuxTests/ExampleTests", LAUNCHER_PRIOR_RUNS=self._prior("failure", runner=LARGE), + ) + self.assertNotEqual(result.returncode, 0) + self.assertIn("already failed", result.stderr) + self.assertFalse((self.root / "dispatch.json").exists(), "must not dispatch") + + def test_a_pinned_pool_ignores_a_run_on_the_other_macos_26_pool(self): + result = self.launch( + "cmuxTests/ExampleTests", "--runner", SMALL, + LAUNCHER_PRIOR_RUNS=self._live(runner=LARGE), ) self.assertEqual(result.returncode, 0, result.stderr) - self.assertIn("reusing that run", result.stdout) - self.assertFalse((self.root / "dispatch.json").exists()) + self.assertEqual(self.dispatch()["runner"], SMALL) def test_invalid_runner_is_rejected_before_github_access(self): result = self.launch("cmuxTests/ExampleTests", "--runner", "macos-15") @@ -477,6 +597,41 @@ def test_the_focused_suite_job_passes_the_runner_variable(self): self.assertEqual( wrapper["env"].get("CMUX_MACOS_RUNNER_TESTS"), "${{ vars.MACOS_RUNNER_TESTS }}" ) + # The overflow switch and thresholds too: without them the wrapper + # would overflow on defaults after an admin turned overflow off. + for name in ("CI_E2E_LARGE_POOL_OVERFLOW", "CI_E2E_OVERFLOW_MIN_QUEUED", + "CI_E2E_OVERFLOW_MAX_LARGE_RUNNING"): + self.assertEqual(wrapper["env"].get("CMUX_" + name), "${{ vars.%s }}" % name) + self.assertNotIn("SPLIT", json.dumps(wrapper["env"])) + + def test_the_overflow_switch_keeps_e2e_on_6vcpu_without_reading_the_queue(self): + result = self.launch( + "cmuxTests/ExampleTests", LAUNCHER_QUEUE=BACKED_UP, + LAUNCHER_VARIABLES=json.dumps([ + {"name": "CI_E2E_LARGE_POOL_OVERFLOW", "value": "0"}, + ]), + ) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertEqual(self.dispatch()["runner"], SMALL) + self.assertEqual(self.queue_reads(), []) + + def test_a_workflow_job_passes_the_overflow_variables_it_cannot_list(self): + result = self.launch( + "cmuxTests/ExampleTests", LAUNCHER_QUEUE=BACKED_UP, + LAUNCHER_VARIABLES="not json", + CMUX_MACOS_RUNNER_TESTS="", CMUX_CI_E2E_LARGE_POOL_OVERFLOW="0", + ) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertEqual(self.dispatch()["runner"], SMALL) + self.assertNotIn(["variable", "list"], [call[:2] for call in self.calls()]) + self.setUp() + result = self.launch( + "cmuxTests/ExampleTests", LAUNCHER_QUEUE=json.dumps(queue(small=1)), + LAUNCHER_VARIABLES="not json", + CMUX_MACOS_RUNNER_TESTS="", CMUX_CI_E2E_OVERFLOW_MIN_QUEUED="1", + ) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertEqual(self.dispatch()["runner"], LARGE) def test_a_run_without_a_dispatch_id_is_still_seen(self): # A run started from the GitHub UI shares the concurrency group and its @@ -613,6 +768,299 @@ def test_ambiguous_dispatch_fails(self): self.dispatch.find_run(HEAD, "cmuxTests/Example", "mine") +class FakeActions: + """Serves one runs page per status and counts every request.""" + + def __init__(self, pages=None, *, fail=False): + self.pages = pages or {} + self.fail = fail + self.paths = [] + + def request(self, method, path, body=None): + self.paths.append((method, path)) + if self.fail: + raise RuntimeError("GET /repos/x/actions/runs failed (503)") + status = path.split("status=", 1)[1].split("&", 1)[0] + return {"workflow_runs": self.pages.get(status, [])} + + +class WorkflowRunnerPoolTests(unittest.TestCase): + """E2E overflows to the 12vcpu pool only when 6vcpu is backed up. + + The 12vcpu macOS 26 pool is reserved first for release and nightly + builds. The earlier parity split sent half of all commits there however + busy it was; `auto` now stays on 6vcpu unless the 6vcpu pool is backed up + and the 12vcpu pool has spare room, and fails safe to 6vcpu. + """ + + COMMITS = ["0123456789abcdef0123456789abcdef0123456" + digit for digit in "0123456789abcdef"] + + @classmethod + def setUpClass(cls): + cls.workflow = yaml.safe_load((ROOT / ".github/workflows/test-e2e.yml").read_text()) + cls.jobs = cls.workflow["jobs"] + spec = importlib.util.spec_from_file_location( + "e2e_runner_pool", ROOT / "scripts/ci/e2e_runner_pool.py" + ) + cls.pool = importlib.util.module_from_spec(spec) + spec.loader.exec_module(cls.pool) + spec = importlib.util.spec_from_file_location( + "focused_dispatch_pool", ROOT / "scripts/ci/dispatch-focused-test.py" + ) + cls.dispatch = importlib.util.module_from_spec(spec) + spec.loader.exec_module(cls.dispatch) + + def load(self, small=0, large_queued=0, large_running=0): + return self.pool.PoolLoad(small, large_queued, large_running) + + def decide(self, load, *, variable="", overflow="", min_queued="", max_large_running="", + requested="auto"): + calls = [] + + def measure(limits): + calls.append(limits) + if isinstance(load, Exception): + raise load + return load + + label = self.pool.resolve( + requested, variable, overflow=overflow, min_queued=min_queued, + max_large_running=max_large_running, measure=measure, + ) + return label, calls + + def test_the_rule(self): + defaults = self.pool.Thresholds() + self.assertEqual((defaults.min_queued, defaults.max_large_running), (4, 2)) + cases = [ + (self.load(4, 0, 0), LARGE), + (self.load(4, 0, 1), LARGE), + (self.load(9, 0, 1), LARGE), + (self.load(3, 0, 0), SMALL), # 6vcpu not backed up + (self.load(9, 1, 0), SMALL), # something already waits on 12vcpu + (self.load(9, 0, 2), SMALL), # 12vcpu already has its share of E2E + (None, SMALL), # unknown load + ] + for load, expected in cases: + with self.subTest(load=load): + self.assertEqual(self.decide(load)[0], expected) + + def test_the_thresholds_are_variables_and_invalid_values_fail_safe(self): + self.assertEqual(self.decide(self.load(2), min_queued="2")[0], LARGE) + self.assertEqual(self.decide(self.load(9, 0, 3), max_large_running="4")[0], LARGE) + self.assertEqual(self.decide(self.load(9), max_large_running="0")[0], SMALL) + for min_queued, max_running in (("x", ""), ("0", ""), ("-1", ""), ("", "-1"), ("", "two")): + with self.subTest(min_queued=min_queued, max_running=max_running): + label, calls = self.decide(self.load(99), min_queued=min_queued, + max_large_running=max_running) + self.assertEqual(label, SMALL) + self.assertEqual(calls, [], "an invalid threshold must not read the queue") + + def test_the_kill_switch_never_reads_the_queue(self): + label, calls = self.decide(self.load(99), overflow="0") + self.assertEqual((label, calls), (SMALL, [])) + for value in ("", "1", "yes"): + with self.subTest(value=value): + self.assertEqual(self.decide(self.load(99), overflow=value)[0], LARGE) + + def test_any_measurement_error_fails_safe(self): + for error in (RuntimeError("503"), ValueError("bad json"), KeyError("workflow_runs")): + with self.subTest(error=error): + self.assertEqual(self.decide(error)[0], SMALL) + + def test_an_explicit_choice_or_admin_variable_is_never_rerouted(self): + for requested in (SMALL, LARGE, "tart-canary"): + with self.subTest(requested=requested): + self.assertEqual(self.decide(self.load(99), requested=requested), (requested, [])) + self.assertEqual(self.decide(self.load(99), variable="blacksmith-6vcpu-macos-15"), + ("blacksmith-6vcpu-macos-15", [])) + + def test_the_commit_no_longer_decides(self): + # Every commit gets the same answer for the same queue; the parity + # split sent odd commits to 12vcpu even when it was busy. + for commit in self.COMMITS: + with self.subTest(commit=commit): + self.assertEqual(self.decide(self.load(0))[0], SMALL) + self.assertEqual(self.decide(self.load(9, 1))[0], SMALL) + + def measure(self, pages, *, workflows_dir=None, exclude_run_id=None): + client = FakeActions(pages) + load = self.pool.measure_load( + client, "manaflow-ai/cmux", self.pool.Thresholds(), + workflows_dir=workflows_dir or ROOT / ".github/workflows", + exclude_run_id=exclude_run_id, + ) + return load, client + + def test_measurement_costs_at_most_two_api_calls(self): + for pages in ( + {}, + queue(small=4), + queue(small=60, large_queued=3), + {"in_progress": [e2e_run(SMALL, n) for n in range(99)]}, + ): + with self.subTest(pages=len(pages.get("in_progress", []))): + _, client = self.measure(pages) + self.assertLessEqual(len(client.paths), self.pool.MAX_API_CALLS) + self.assertLessEqual(self.pool.MAX_API_CALLS, 2) + for method, path in client.paths: + self.assertEqual(method, "GET") + self.assertRegex( + path, r"^/repos/manaflow-ai/cmux/actions/runs\?status=(in_progress|queued)&per_page=100$") + + def test_a_busy_12vcpu_pool_stops_after_one_call(self): + for pages in (queue(small=9, large_running=2), queue(small=9, reserved=1)): + with self.subTest(pages=pages): + load, client = self.measure(pages) + self.assertEqual(len(client.paths), 1) + self.assertFalse(self.pool.overflows(load, self.pool.Thresholds())) + + def test_measurement_attributes_runs_to_pools(self): + pages = queue(small=5, large_running=1, large_queued=1) + pages["in_progress"].append({"id": 9, "status": "in_progress", "name": "CI", + "path": ".github/workflows/ci.yml", "display_title": "fix"}) + load, _ = self.measure(pages) + self.assertEqual(load, self.load(5, 1, 1)) + # The run deciding is not its own demand. + load, _ = self.measure(queue(small=4), exclude_run_id=100) + self.assertEqual(load.small_queued, 3) + + def test_a_full_page_is_unknown(self): + load, client = self.measure({"in_progress": [e2e_run(SMALL, n) for n in range(100)]}) + self.assertIsNone(load) + self.assertEqual(len(client.paths), 1) + + def test_reserved_workflows_that_never_use_macos_do_not_block(self): + with tempfile.TemporaryDirectory() as temp: + workflows = Path(temp) + (workflows / "nightly.yml").write_text("jobs:\n b:\n runs-on: blacksmith-12vcpu-macos-26\n") + (workflows / "release-notes.yml").write_text("jobs:\n b:\n runs-on: ubuntu-latest\n") + notes = {"id": 7, "status": "in_progress", "name": "Release notes", + "path": ".github/workflows/release-notes.yml", "display_title": "notes"} + load, _ = self.measure({"in_progress": [notes] + queue(small=4)["in_progress"]}, + workflows_dir=workflows) + self.assertEqual(load, self.load(4)) + load, _ = self.measure(queue(small=4, reserved=1), workflows_dir=workflows) + self.assertEqual(load.large_queued, 1) + + def test_an_api_error_fails_safe_end_to_end(self): + client = FakeActions(fail=True) + label = self.pool.resolve( + "auto", "", overflow="", min_queued="", max_large_running="", + measure=lambda limits: self.pool.measure_load(client, "manaflow-ai/cmux", limits), + ) + self.assertEqual(label, SMALL) + self.assertEqual(len(client.paths), 1) + + def test_the_dispatcher_reads_the_queue_through_the_janitor_client(self): + # One rule, one client shape: run-e2e.sh subclasses the queue + # janitor's GitHub client and only swaps its transport for `gh api`. + self.assertTrue(issubclass(self.dispatch.GhApi, self.pool.queue_janitor.GitHub)) + # A workflow job passes the variables, so only queue reads remain. + variables = {"CMUX_MACOS_RUNNER_TESTS": "", "CMUX_CI_E2E_LARGE_POOL_OVERFLOW": ""} + with mock.patch.dict(os.environ, variables), mock.patch.object( + self.dispatch, "output", return_value=json.dumps( + {"workflow_runs": queue(small=4)["in_progress"]})) as output: + label = self.dispatch.routed_runner(SMALL) + self.assertEqual(label, LARGE) + self.assertLessEqual(output.call_count, 2) + for call in output.call_args_list: + self.assertEqual(call.args[:4], ("gh", "api", "--method", "GET")) + with mock.patch.dict(os.environ, variables), mock.patch.object( + self.dispatch, "output", side_effect=subprocess.CalledProcessError(1, "gh")): + self.assertEqual(self.dispatch.routed_runner(SMALL), SMALL) + + # Workflow wiring ------------------------------------------------------ + + def pool_step(self): + steps = self.jobs["runner"]["steps"] + return next(step for step in steps if "e2e_runner_pool.py" in step.get("run", "")) + + def run_pool_step(self, *, requested="auto", variable="", overflow="", min_queued="", + max_large_running=""): + """Run the workflow's own step script with the values GitHub would pass. + + No token reaches it, so a decision that reads the queue fails safe. + """ + step = self.pool_step() + env = {k: v for k, v in os.environ.items() if k not in ("GH_TOKEN", "GITHUB_TOKEN")} + values = { + "${{ github.token }}": "", + "${{ github.repository }}": "manaflow-ai/cmux", + "${{ inputs.runner }}": requested, + "${{ vars.MACOS_RUNNER_TESTS }}": variable, + "${{ vars.CI_E2E_LARGE_POOL_OVERFLOW }}": overflow, + "${{ vars.CI_E2E_OVERFLOW_MIN_QUEUED }}": min_queued, + "${{ vars.CI_E2E_OVERFLOW_MAX_LARGE_RUNNING }}": max_large_running, + } + for name, expression in step["env"].items(): + self.assertIn(expression, values, f"unexpected input {name}: {expression}") + env[name] = values[expression] + with tempfile.TemporaryDirectory() as temp: + output = Path(temp) / "output" + output.write_text("") + env["GITHUB_OUTPUT"] = str(output) + result = subprocess.run(["bash", "-e", "-c", step["run"]], cwd=ROOT, env=env, check=True, + capture_output=True, text=True) + lines = dict(line.split("=", 1) for line in output.read_text().splitlines() if "=" in line) + return lines["label"], result.stderr + + def test_the_workflow_step_resolves_through_the_rule(self): + self.assertEqual(self.run_pool_step()[0], SMALL) + label, stderr = self.run_pool_step() + self.assertIn("could not read the runner queue", stderr) + self.assertEqual(self.run_pool_step(overflow="0")[0], SMALL) + self.assertEqual(self.run_pool_step(requested="tart-small")[0], "tart-small") + self.assertEqual(self.run_pool_step(requested=LARGE)[0], LARGE) + self.assertEqual(self.run_pool_step(variable="blacksmith-6vcpu-macos-15")[0], + "blacksmith-6vcpu-macos-15") + + def test_the_pool_job_reads_actions_and_nothing_else(self): + self.assertEqual(self.workflow["permissions"], {"contents": "read"}) + job = self.jobs["runner"] + self.assertEqual(job["permissions"], {"contents": "read", "actions": "read"}) + self.assertIn("ubuntu", job["runs-on"]) + self.assertEqual(job["outputs"]["label"], "${{ steps.pool.outputs.label }}") + step = self.pool_step() + self.assertEqual(step["id"], "pool") + self.assertEqual(step["env"]["GH_TOKEN"], "${{ github.token }}") + self.assertNotIn("SPLIT", yaml.safe_dump(job)) + checkout = next(step for step in job["steps"] if "actions/checkout" in step.get("uses", "")) + paths = checkout["with"]["sparse-checkout"].split() + for path in ("scripts/ci/e2e_runner_pool.py", "scripts/ci/queue_janitor.py", ".github/workflows/"): + self.assertIn(path, paths) + self.assertIs(checkout["with"]["persist-credentials"], False) + # No other job gained write access for this. + for name, other in self.jobs.items(): + for scope, level in (other.get("permissions") or {}).items(): + with self.subTest(job=name, scope=scope): + self.assertEqual(level, "read") + + def test_the_pool_helper_explains_the_release_priority(self): + source = (ROOT / "scripts/ci/e2e_runner_pool.py").read_text() + for phrase in ("reserved", "release and nightly", "at most two requests"): + self.assertIn(phrase, source) + comment = (ROOT / ".github/workflows/test-e2e.yml").read_text() + self.assertIn("reserved first for release and nightly", comment) + + def test_macos_jobs_run_on_the_resolved_pool(self): + label = "${{ needs.runner.outputs.label }}" + for name in ("build", "test"): + with self.subTest(job=name): + job = self.jobs[name] + self.assertIn("runner", job["needs"]) + self.assertEqual(job["runs-on"], label) + tart = next(step for step in job["steps"] + if step.get("name") == "Validate Tart canary identity") + self.assertEqual(tart["if"], "${{ startsWith(needs.runner.outputs.label, 'tart-') }}") + self.assertEqual(tart["env"]["REQUESTED_RUNNER"], label) + # Nothing in a macOS job may resolve the pool a second way. + text = yaml.safe_dump(job) + self.assertNotIn("inputs.runner", text) + self.assertNotIn("vars.MACOS_RUNNER_TESTS", text) + self.assertEqual(self.jobs["build"]["env"]["CMUX_PRODUCT_RUNNER"], label) + + class SuiteWorkflowForwardsFocusedRuns(unittest.TestCase): def test_focused_selectors_never_compile_in_the_suite_workflow(self): jobs = yaml.safe_load((ROOT / ".github/workflows/test-macos-suite.yml").read_text())["jobs"]