From d8189a60713107ca60c819abd6d0c1af033515fa Mon Sep 17 00:00:00 2001 From: Leo Li Date: Thu, 24 Sep 2026 04:09:16 +0000 Subject: [PATCH 1/3] ci: expect test-e2e.yml to split auto across both macOS 26 pools A direct test-e2e.yml dispatch with runner=auto lands every commit on blacksmith-6vcpu-macos-26; only run-e2e.sh applies the commit-keyed split to blacksmith-12vcpu-macos-26. These tests expect the workflow to resolve the pool itself with the dispatcher's rule, run build/test and the Tart checks on that label, and honour a CI_E2E_LARGE_POOL_SPLIT=0 kill switch in both places. They fail until the workflow does. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_012pAcDGiHaibDaMU4CPvXAP --- tests/test_ci_e2e_compilation_cache.py | 2 +- tests/test_ci_self_hosted_guard.sh | 4 +- tests/test_run_e2e.py | 147 +++++++++++++++++++++++++ 3 files changed, 150 insertions(+), 3 deletions(-) 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..5d1a99a01208 100644 --- a/tests/test_run_e2e.py +++ b/tests/test_run_e2e.py @@ -477,6 +477,35 @@ 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 split switch too: without it the wrapper would pin a routed + # commit to the large pool after the workflow stopped splitting. + self.assertEqual( + wrapper["env"].get("CMUX_CI_E2E_LARGE_POOL_SPLIT"), + "${{ vars.CI_E2E_LARGE_POOL_SPLIT }}", + ) + + def test_the_split_switch_keeps_routed_commits_on_the_default_pool(self): + # REMOTE_HEAD ends in b and would be routed. With the switch off the + # workflow's auto stays on the 6vcpu pool, so the dispatcher must not + # pin the commit to the large pool behind its back. + result = self.launch( + "cmuxTests/ExampleTests", "--ref", "topic/fix", + LAUNCHER_VARIABLES=json.dumps([ + {"name": "CI_E2E_LARGE_POOL_SPLIT", "value": "0"}, + ]), + ) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertNotIn("runner", self.dispatch()) + + def test_a_workflow_job_passes_the_split_switch_it_cannot_list(self): + result = self.launch( + "cmuxTests/ExampleTests", "--ref", "topic/fix", + LAUNCHER_VARIABLES="not json", + CMUX_MACOS_RUNNER_TESTS="", CMUX_CI_E2E_LARGE_POOL_SPLIT="0", + ) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertNotIn("runner", self.dispatch()) + self.assertNotIn(["variable", "list"], [call[:2] for call in self.calls()]) 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 +642,124 @@ def test_ambiguous_dispatch_fails(self): self.dispatch.find_run(HEAD, "cmuxTests/Example", "mine") +class WorkflowRunnerPoolTests(unittest.TestCase): + """test-e2e.yml's `auto` splits commits across pools the way run-e2e.sh does. + + A direct dispatch with runner=auto used to land every commit on the 6vcpu + pool, which queued for hours while the 12vcpu pool sat idle. The workflow + now resolves the pool itself, with the dispatcher's rule, so the two agree + on the runner label that product reuse and in-flight matching key on. + """ + + SMALL = "blacksmith-6vcpu-macos-26" + LARGE = "blacksmith-12vcpu-macos-26" + # Every possible last hex digit, so both parities are covered. + 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 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="", split="", commit): + """Run the workflow's own step script with the values GitHub would pass.""" + step = self.pool_step() + env = dict(os.environ) + values = { + "${{ inputs.runner }}": requested, + "${{ vars.MACOS_RUNNER_TESTS }}": variable, + "${{ vars.CI_E2E_LARGE_POOL_SPLIT }}": split, + "${{ needs.resolve-ref.outputs.sha }}": commit, + } + 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) + 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"] + + def test_auto_splits_commits_across_both_pools(self): + for commit in self.COMMITS: + with self.subTest(commit=commit): + expected = self.LARGE if int(commit[-1], 16) % 2 else self.SMALL + self.assertEqual(self.pool.resolve("auto", "", "", commit), expected) + self.assertEqual(self.pool.resolve("", "", "", commit), expected) + self.assertEqual(self.run_pool_step(commit=commit), expected) + + def test_the_workflow_and_the_dispatcher_route_every_commit_alike(self): + for commit in self.COMMITS: + for variable in ("", self.SMALL, "blacksmith-6vcpu-macos-15"): + for split in ("", "1", "0"): + with self.subTest(commit=commit, variable=variable, split=split): + default = variable or self.SMALL + self.assertEqual( + self.run_pool_step(variable=variable, split=split, commit=commit), + self.dispatch.routed_runner( + commit, default, self.pool.split_enabled(split) + ), + ) + + def test_an_explicit_choice_or_admin_variable_is_never_rerouted(self): + odd = self.COMMITS[1] + self.assertEqual(self.pool.resolve(self.SMALL, "", "", odd), self.SMALL) + self.assertEqual(self.pool.resolve("tart-canary", "", "", odd), "tart-canary") + self.assertEqual(self.pool.resolve("auto", "blacksmith-6vcpu-macos-15", "", odd), + "blacksmith-6vcpu-macos-15") + self.assertEqual(self.pool.resolve("auto", "", "0", odd), self.SMALL) + self.assertEqual(self.run_pool_step(requested="tart-small", commit=odd), "tart-small") + + def test_an_unresolved_commit_is_refused(self): + for commit in ("", "main", "B" * 40, "b" * 39): + with self.subTest(commit=commit), self.assertRaises(ValueError): + self.pool.resolve("auto", "", "", commit) + + def test_the_pool_job_resolves_after_the_commit_on_linux(self): + job = self.jobs["runner"] + self.assertIn("resolve-ref", job["needs"]) + self.assertIn("ubuntu", job["runs-on"]) + self.assertEqual(job["outputs"]["label"], "${{ steps.pool.outputs.label }}") + self.assertEqual(self.pool_step()["id"], "pool") + checkout = next(step for step in job["steps"] if "actions/checkout" in step.get("uses", "")) + self.assertIn("scripts/ci/e2e_runner_pool.py", checkout["with"]["sparse-checkout"]) + self.assertIs(checkout["with"]["persist-credentials"], False) + + 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"] From 15652511fc0af38698dcc8bc9be7ef9d03bb90e3 Mon Sep 17 00:00:00 2001 From: Leo Li Date: Thu, 24 Sep 2026 04:09:16 +0000 Subject: [PATCH 2/3] ci: split test-e2e.yml's auto runner across both macOS 26 pools A dispatch with runner=auto resolved to blacksmith-6vcpu-macos-26 for every commit unless it came through run-e2e.sh, which alone routed odd commits to blacksmith-12vcpu-macos-26. E2E builds queued for hours on the 6vcpu pool while the 12vcpu pool sat idle. A new Linux `runner` job resolves the pool after resolve-ref, with the rule now in scripts/ci/e2e_runner_pool.py: an explicit runner wins; else MACOS_RUNNER_TESTS, else the 6vcpu pool; and when that is the 6vcpu pool, a resolved commit whose last hex digit is odd goes to the 12vcpu pool. CI_E2E_LARGE_POOL_SPLIT=0 turns the split off. build, test, their Tart identity checks and CMUX_PRODUCT_RUNNER read the job's output. dispatch-focused-test.py imports the same function and reads the same switch, and test-macos-suite.yml passes the switch to it. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_012pAcDGiHaibDaMU4CPvXAP --- .github/workflows/test-e2e.yml | 63 ++++++++++++++--- .github/workflows/test-macos-suite.yml | 1 + scripts/ci/dispatch-focused-test.py | 98 +++++++++++++++++--------- scripts/ci/e2e_runner_pool.py | 74 +++++++++++++++++++ 4 files changed, 194 insertions(+), 42 deletions(-) create mode 100644 scripts/ci/e2e_runner_pool.py diff --git a/.github/workflows/test-e2e.yml b/.github/workflows/test-e2e.yml index ac7c56b2496a..a3622dc03ede 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 splits commits across the 6vcpu and 12vcpu macOS 26 pools; 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 pre-split default even when the runner job routes the commit +# to the 12vcpu pool. Every auto dispatch of one ref and filter still shares a +# group, and still lands on one pool, because the split is keyed on the +# commit. run-e2e.sh names a routed pool 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,44 @@ jobs: with: ref: ${{ inputs.ref }} + runner: + # The macOS pool build and test run on. `auto` follows MACOS_RUNNER_TESTS + # and, on the default 6vcpu macOS 26 pool, sends odd commits to the 12vcpu + # pool (CI_E2E_LARGE_POOL_SPLIT=0 turns that off). run-e2e.sh applies the + # same rule from the same script, so product reuse and in-flight matching + # see one label per commit however the run was started. + needs: resolve-ref + runs-on: ${{ vars.LINUX_RUNNER || 'blacksmith-4vcpu-ubuntu-2404' }} + timeout-minutes: 5 + outputs: + label: ${{ steps.pool.outputs.label }} + steps: + # The helper comes from this workflow's revision, not the tested one, + # which may predate it. + - name: Checkout pool helper + uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2 + with: + sparse-checkout: scripts/ci/e2e_runner_pool.py + sparse-checkout-cone-mode: false + persist-credentials: false + + - name: Pick the macOS pool + id: pool + env: + REQUESTED_RUNNER: ${{ inputs.runner }} + RUNNER_VARIABLE: ${{ vars.MACOS_RUNNER_TESTS }} + LARGE_POOL_SPLIT: ${{ vars.CI_E2E_LARGE_POOL_SPLIT }} + COMMIT: ${{ needs.resolve-ref.outputs.sha }} + run: | + set -euo pipefail + label="$(python3 scripts/ci/e2e_runner_pool.py \ + --requested "$REQUESTED_RUNNER" \ + --variable "$RUNNER_VARIABLE" \ + --split "$LARGE_POOL_SPLIT" \ + --commit "$COMMIT")" + echo "label=$label" >> "$GITHUB_OUTPUT" + echo "Runner: $label (requested ${REQUESTED_RUNNER:-auto}) for $COMMIT" + 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 +233,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 +257,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 +656,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 +664,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 +682,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..b95435bffba9 100644 --- a/.github/workflows/test-macos-suite.yml +++ b/.github/workflows/test-macos-suite.yml @@ -92,6 +92,7 @@ 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_SPLIT: ${{ vars.CI_E2E_LARGE_POOL_SPLIT }} 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..61b97fc8798a 100644 --- a/scripts/ci/dispatch-focused-test.py +++ b/scripts/ci/dispatch-focused-test.py @@ -16,10 +16,16 @@ from urllib.parse import quote import uuid +sys.path.insert(0, str(Path(__file__).resolve().parent)) +from e2e_runner_pool import LARGE_RUNNER, SPLIT_VARIABLE, split_enabled +from e2e_runner_pool import routed_runner as pool_routed_runner + 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" +# `vars.CI_E2E_LARGE_POOL_SPLIT`, passed the same way; see large_pool_split(). +SPLIT_ENV = "CMUX_CI_E2E_LARGE_POOL_SPLIT" ROOT = Path(__file__).resolve().parents[2] RUN_DISCOVERY_ATTEMPTS = 12 RUN_DISCOVERY_TIMEOUT_SECONDS = 60.0 @@ -38,13 +44,9 @@ "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" +# Half of all commits compile on the large macOS 26 SKU. The rule lives in +# e2e_runner_pool.py, which test-e2e.yml runs too, so a run started here and +# one started from the Actions UI put the same commit on the same pool. # 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 +188,50 @@ 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 large_pool_split() -> bool: + """Whether `vars.CI_E2E_LARGE_POOL_SPLIT` leaves the split on. + + A workflow job passes the variable in CMUX_CI_E2E_LARGE_POOL_SPLIT. An + unreadable listing counts as on, the workflow's own default; the only + cost of guessing wrong is pinning a commit to the pool the workflow would + otherwise have split it onto. + """ + if SPLIT_ENV in os.environ: + return split_enabled(os.environ[SPLIT_ENV]) + if VARIABLE_ENV in os.environ: + # A job token cannot list variables; a caller that passed one variable + # but not the other predates the switch. + return True + return split_enabled((listed_variables() or {}).get(SPLIT_VARIABLE)) + + def default_runner() -> str | None: """The label `runner: auto` resolves to, or None when it cannot be known. @@ -207,22 +253,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): + variables = listed_variables() + if variables is None: return None - if not isinstance(variables, list): - 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,15 +269,9 @@ 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. - - Only the free default is split. A repository variable naming any other - pool is an admin decision, and it wins unchanged. - """ - if default == SMALL_RUNNER and int(commit[-1], 16) % 2: - return LARGE_RUNNER - return default +def routed_runner(commit: str, default: str | None, split: bool = True) -> str | None: + """The pool an unpinned dispatch at `commit` runs on; see e2e_runner_pool.""" + return pool_routed_runner(commit, default, split) def attempts( @@ -430,7 +460,9 @@ def main() -> int: # could not be established, and the in-flight guards below stay silent # rather than compare against a runner they guessed. pinned = args.runner not in (None, "auto") - runner = args.runner if pinned else routed_runner(commit, default_runner()) + runner = args.runner if pinned else routed_runner( + commit, default_runner(), large_pool_split() + ) # 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) @@ -517,6 +549,8 @@ def main() -> int: } if args.runner is not None: fields["runner"] = args.runner + # test-e2e.yml would route this commit to the same pool on its own. Name + # it anyway so the run title carries the pool the guards above match on. if not pinned and runner == LARGE_RUNNER: fields["runner"] = runner command = ["gh", "workflow", "run", WORKFLOW, "--repo", REPO] diff --git a/scripts/ci/e2e_runner_pool.py b/scripts/ci/e2e_runner_pool.py new file mode 100644 index 000000000000..e031bd9f90f2 --- /dev/null +++ b/scripts/ci/e2e_runner_pool.py @@ -0,0 +1,74 @@ +#!/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 puts +the same commit on the same pool. In-flight reuse, the failed-selector refusal +and the compiled-product contract all match on that runner label. + +`runner: auto` means `vars.MACOS_RUNNER_TESTS` when it names a pool, else +the free 6vcpu macOS 26 pool. When that resolves to the 6vcpu pool, half of +all commits run on the 12vcpu pool instead. The split is keyed on the commit +(odd last hex digit goes large), not drawn at random, so every dispatch at one +commit lands on one pool. `vars.CI_E2E_LARGE_POOL_SPLIT == '0'` turns the +split off. An explicit runner, or a variable naming any other pool, is never +rerouted. +""" +from __future__ import annotations + +import argparse +import re +import sys + +SMALL_RUNNER = "blacksmith-6vcpu-macos-26" +LARGE_RUNNER = "blacksmith-12vcpu-macos-26" +# The repository variable that turns the split off when set to "0". +SPLIT_VARIABLE = "CI_E2E_LARGE_POOL_SPLIT" +COMMIT = re.compile(r"[0-9a-f]{40}") + + +def split_enabled(value: str | None) -> bool: + """Whether the large-pool split is on. Unset or anything but "0" is on.""" + return (value or "").strip() != "0" + + +def routed_runner(commit: str, default: str | None, split: bool = True) -> str | None: + """The pool an unpinned run at `commit` lands on, given what auto means. + + Only the free default is split. A repository variable naming any other + pool is an admin decision, and it wins unchanged. None stays None: a + caller that could not establish the default must not act on a guess. + """ + if split and default == SMALL_RUNNER and int(commit[-1], 16) % 2: + return LARGE_RUNNER + return default + + +def resolve(requested: str | None, variable: str | None, split: str | None, commit: str) -> str: + """The runner label for a workflow run, from its inputs and variables.""" + if not COMMIT.fullmatch(commit): + raise ValueError(f"expected a full lowercase commit SHA, got {commit!r}") + requested = (requested or "").strip() + if requested and requested != "auto": + return requested + default = (variable or "").strip() or SMALL_RUNNER + return routed_runner(commit, default, split_enabled(split)) + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__.splitlines()[0]) + parser.add_argument("--commit", required=True, help="resolved 40-character commit SHA") + parser.add_argument("--requested", default="", help="the workflow's runner input") + parser.add_argument("--variable", default="", help="vars.MACOS_RUNNER_TESTS") + parser.add_argument("--split", default="", help=f"vars.{SPLIT_VARIABLE}") + args = parser.parse_args(argv) + try: + print(resolve(args.requested, args.variable, args.split, args.commit)) + except ValueError as error: + print(f"error: {error}", file=sys.stderr) + return 1 + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) From 31dd5a616d02c2f490cc925dcc0eab8f71826947 Mon Sep 17 00:00:00 2001 From: Leo Li Date: Thu, 24 Sep 2026 04:25:47 +0000 Subject: [PATCH 3/3] ci: overflow E2E to the 12vcpu macOS 26 pool only when 6vcpu is backed up The 12vcpu macOS 26 pool is reserved first for release and nightly builds. The parity split sent every odd commit's E2E run there however busy it was. `auto` now stays on the 6vcpu pool and overflows only when at least CI_E2E_OVERFLOW_MIN_QUEUED (default 4) other E2E runs are in flight on 6vcpu, no release or nightly run is in flight, nothing is queued on 12vcpu, and fewer than CI_E2E_OVERFLOW_MAX_LARGE_RUNNING (default 2) E2E runs are on it. Any API error, a possibly truncated listing, or an invalid threshold stays on 6vcpu, and CI_E2E_LARGE_POOL_OVERFLOW=0 (renamed from CI_E2E_LARGE_POOL_SPLIT) turns overflow off. The decision costs at most two GET requests (one page each of in-progress and queued runs), through the queue janitor's client, and attributes demand from run titles and workflow names rather than per-run job listings, so dozens of dispatches an hour stay well inside the shared GITHUB_TOKEN budget. The runner job gains actions: read. run-e2e.sh applies the same rule, names the pool it chose, and reads the queue only when it is about to dispatch. Because the choice is no longer a function of the commit, its in-flight reuse and overlap refusal now match a run on either macOS 26 pool for an unpinned dispatch. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_012pAcDGiHaibDaMU4CPvXAP --- .github/workflows/test-e2e.yml | 52 ++- .github/workflows/test-macos-suite.yml | 4 +- scripts/ci/dispatch-focused-test.py | 159 +++++--- scripts/ci/e2e_runner_pool.py | 303 +++++++++++++-- tests/test_run_e2e.py | 487 ++++++++++++++++++++----- 5 files changed, 810 insertions(+), 195 deletions(-) diff --git a/.github/workflows/test-e2e.yml b/.github/workflows/test-e2e.yml index a3622dc03ede..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 and splits commits across the 6vcpu and 12vcpu macOS 26 pools; 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 @@ -44,10 +44,10 @@ on: - tart-small # run-name and the concurrency group cannot read job outputs, so they spell -# `auto` as the pre-split default even when the runner job routes the commit -# to the 12vcpu pool. Every auto dispatch of one ref and filter still shares a -# group, and still lands on one pool, because the split is keyed on the -# commit. run-e2e.sh names a routed pool explicitly, so its titles are exact. +# `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 @@ -66,42 +66,58 @@ jobs: ref: ${{ inputs.ref }} runner: - # The macOS pool build and test run on. `auto` follows MACOS_RUNNER_TESTS - # and, on the default 6vcpu macOS 26 pool, sends odd commits to the 12vcpu - # pool (CI_E2E_LARGE_POOL_SPLIT=0 turns that off). run-e2e.sh applies the - # same rule from the same script, so product reuse and in-flight matching - # see one label per commit however the run was started. - needs: resolve-ref + # 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. + # 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 + 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_SPLIT: ${{ vars.CI_E2E_LARGE_POOL_SPLIT }} - COMMIT: ${{ needs.resolve-ref.outputs.sha }} + 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" \ - --split "$LARGE_POOL_SPLIT" \ - --commit "$COMMIT")" + --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}) for $COMMIT" + echo "Runner: $label (requested ${REQUESTED_RUNNER:-auto})" filter: # Seconds on Linux, and it rejects a malformed selector before either diff --git a/.github/workflows/test-macos-suite.yml b/.github/workflows/test-macos-suite.yml index b95435bffba9..7f6a23f5a088 100644 --- a/.github/workflows/test-macos-suite.yml +++ b/.github/workflows/test-macos-suite.yml @@ -92,7 +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_SPLIT: ${{ vars.CI_E2E_LARGE_POOL_SPLIT }} + 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 61b97fc8798a..ae08778de19b 100644 --- a/scripts/ci/dispatch-focused-test.py +++ b/scripts/ci/dispatch-focused-test.py @@ -17,15 +17,17 @@ import uuid sys.path.insert(0, str(Path(__file__).resolve().parent)) -from e2e_runner_pool import LARGE_RUNNER, SPLIT_VARIABLE, split_enabled -from e2e_runner_pool import routed_runner as pool_routed_runner +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" -# `vars.CI_E2E_LARGE_POOL_SPLIT`, passed the same way; see large_pool_split(). -SPLIT_ENV = "CMUX_CI_E2E_LARGE_POOL_SPLIT" +# 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 @@ -44,9 +46,12 @@ "tart-dual", "tart-small", ) -# Half of all commits compile on the large macOS 26 SKU. The rule lives in -# e2e_runner_pool.py, which test-e2e.yml runs too, so a run started here and -# one started from the Actions UI put the same commit on the same pool. +# 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 @@ -215,21 +220,41 @@ def listed_variables() -> dict[str, str] | None: return _listed # type: ignore[return-value] -def large_pool_split() -> bool: - """Whether `vars.CI_E2E_LARGE_POOL_SPLIT` leaves the split on. +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 passes the variable in CMUX_CI_E2E_LARGE_POOL_SPLIT. An - unreadable listing counts as on, the workflow's own default; the only - cost of guessing wrong is pinning a commit to the pool the workflow would - otherwise have split it onto. + 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 SPLIT_ENV in os.environ: - return split_enabled(os.environ[SPLIT_ENV]) + if env_name in os.environ: + return os.environ[env_name] if VARIABLE_ENV in os.environ: - # A job token cannot list variables; a caller that passed one variable - # but not the other predates the switch. - return True - return split_enabled((listed_variables() or {}).get(SPLIT_VARIABLE)) + 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: @@ -269,19 +294,51 @@ def default_runner() -> str | None: return literal.group(1) if literal else None -def routed_runner(commit: str, default: str | None, split: bool = True) -> str | None: - """The pool an unpinned dispatch at `commit` runs on; see e2e_runner_pool.""" - return pool_routed_runner(commit, default, split) +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. + + 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 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", ""))) @@ -290,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 @@ -316,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. @@ -328,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. """ @@ -338,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. @@ -456,16 +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(), large_pool_split() - ) + 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( @@ -476,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 @@ -488,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) @@ -509,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 " @@ -537,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 = { @@ -549,9 +614,9 @@ def main() -> int: } if args.runner is not None: fields["runner"] = args.runner - # test-e2e.yml would route this commit to the same pool on its own. Name - # it anyway so the run title carries the pool the guards above match on. - 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 index e031bd9f90f2..717844119926 100644 --- a/scripts/ci/e2e_runner_pool.py +++ b/scripts/ci/e2e_runner_pool.py @@ -2,71 +2,302 @@ """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 puts -the same commit on the same pool. In-flight reuse, the failed-selector refusal -and the compiled-product contract all match on that runner label. - -`runner: auto` means `vars.MACOS_RUNNER_TESTS` when it names a pool, else -the free 6vcpu macOS 26 pool. When that resolves to the 6vcpu pool, half of -all commits run on the 12vcpu pool instead. The split is keyed on the commit -(odd last hex digit goes large), not drawn at random, so every dispatch at one -commit lands on one pool. `vars.CI_E2E_LARGE_POOL_SPLIT == '0'` turns the -split off. An explicit runner, or a variable naming any other pool, is never -rerouted. +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" -# The repository variable that turns the split off when set to "0". -SPLIT_VARIABLE = "CI_E2E_LARGE_POOL_SPLIT" -COMMIT = re.compile(r"[0-9a-f]{40}") +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" -def split_enabled(value: str | None) -> bool: - """Whether the large-pool split is on. Unset or anything but "0" is on.""" +@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 routed_runner(commit: str, default: str | None, split: bool = True) -> str | None: - """The pool an unpinned run at `commit` lands on, given what auto means. +def thresholds(min_queued: str | None, max_large_running: str | None) -> Thresholds | None: + """Thresholds from repository variables; None when either is invalid. - Only the free default is split. A repository variable naming any other - pool is an admin decision, and it wins unchanged. None stays None: a - caller that could not establish the default must not act on a guess. + 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. """ - if split and default == SMALL_RUNNER and int(commit[-1], 16) % 2: - return LARGE_RUNNER - return default + 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 resolve(requested: str | None, variable: str | None, split: str | None, commit: str) -> str: +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.""" - if not COMMIT.fullmatch(commit): - raise ValueError(f"expected a full lowercase commit SHA, got {commit!r}") requested = (requested or "").strip() if requested and requested != "auto": return requested default = (variable or "").strip() or SMALL_RUNNER - return routed_runner(commit, default, split_enabled(split)) + 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: list[str] | None = None) -> int: +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("--commit", required=True, help="resolved 40-character commit SHA") parser.add_argument("--requested", default="", help="the workflow's runner input") parser.add_argument("--variable", default="", help="vars.MACOS_RUNNER_TESTS") - parser.add_argument("--split", default="", help=f"vars.{SPLIT_VARIABLE}") + 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) - try: - print(resolve(args.requested, args.variable, args.split, args.commit)) - except ValueError as error: - print(f"error: {error}", file=sys.stderr) - return 1 + + 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 diff --git a/tests/test_run_e2e.py b/tests/test_run_e2e.py index 5d1a99a01208..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_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_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_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,35 +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 split switch too: without it the wrapper would pin a routed - # commit to the large pool after the workflow stopped splitting. - self.assertEqual( - wrapper["env"].get("CMUX_CI_E2E_LARGE_POOL_SPLIT"), - "${{ vars.CI_E2E_LARGE_POOL_SPLIT }}", - ) - - def test_the_split_switch_keeps_routed_commits_on_the_default_pool(self): - # REMOTE_HEAD ends in b and would be routed. With the switch off the - # workflow's auto stays on the 6vcpu pool, so the dispatcher must not - # pin the commit to the large pool behind its back. + # 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", "--ref", "topic/fix", + "cmuxTests/ExampleTests", LAUNCHER_QUEUE=BACKED_UP, LAUNCHER_VARIABLES=json.dumps([ - {"name": "CI_E2E_LARGE_POOL_SPLIT", "value": "0"}, + {"name": "CI_E2E_LARGE_POOL_OVERFLOW", "value": "0"}, ]), ) self.assertEqual(result.returncode, 0, result.stderr) - self.assertNotIn("runner", self.dispatch()) + self.assertEqual(self.dispatch()["runner"], SMALL) + self.assertEqual(self.queue_reads(), []) - def test_a_workflow_job_passes_the_split_switch_it_cannot_list(self): + def test_a_workflow_job_passes_the_overflow_variables_it_cannot_list(self): result = self.launch( - "cmuxTests/ExampleTests", "--ref", "topic/fix", + "cmuxTests/ExampleTests", LAUNCHER_QUEUE=BACKED_UP, LAUNCHER_VARIABLES="not json", - CMUX_MACOS_RUNNER_TESTS="", CMUX_CI_E2E_LARGE_POOL_SPLIT="0", + CMUX_MACOS_RUNNER_TESTS="", CMUX_CI_E2E_LARGE_POOL_OVERFLOW="0", ) self.assertEqual(result.returncode, 0, result.stderr) - self.assertNotIn("runner", self.dispatch()) + 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 @@ -642,18 +768,31 @@ 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): - """test-e2e.yml's `auto` splits commits across pools the way run-e2e.sh does. + """E2E overflows to the 12vcpu pool only when 6vcpu is backed up. - A direct dispatch with runner=auto used to land every commit on the 6vcpu - pool, which queued for hours while the 12vcpu pool sat idle. The workflow - now resolves the pool itself, with the dispatcher's rule, so the two agree - on the runner label that product reuse and in-flight matching key on. + 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. """ - SMALL = "blacksmith-6vcpu-macos-26" - LARGE = "blacksmith-12vcpu-macos-26" - # Every possible last hex digit, so both parities are covered. COMMITS = ["0123456789abcdef0123456789abcdef0123456" + digit for digit in "0123456789abcdef"] @classmethod @@ -671,19 +810,188 @@ def setUpClass(cls): 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="", split="", commit): - """Run the workflow's own step script with the values GitHub would pass.""" + 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 = dict(os.environ) + 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_SPLIT }}": split, - "${{ needs.resolve-ref.outputs.sha }}": commit, + "${{ 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}") @@ -692,55 +1000,48 @@ def run_pool_step(self, *, requested="auto", variable="", split="", commit): output = Path(temp) / "output" output.write_text("") env["GITHUB_OUTPUT"] = str(output) - subprocess.run(["bash", "-e", "-c", step["run"]], cwd=ROOT, env=env, check=True, - capture_output=True, text=True) + 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"] - - def test_auto_splits_commits_across_both_pools(self): - for commit in self.COMMITS: - with self.subTest(commit=commit): - expected = self.LARGE if int(commit[-1], 16) % 2 else self.SMALL - self.assertEqual(self.pool.resolve("auto", "", "", commit), expected) - self.assertEqual(self.pool.resolve("", "", "", commit), expected) - self.assertEqual(self.run_pool_step(commit=commit), expected) - - def test_the_workflow_and_the_dispatcher_route_every_commit_alike(self): - for commit in self.COMMITS: - for variable in ("", self.SMALL, "blacksmith-6vcpu-macos-15"): - for split in ("", "1", "0"): - with self.subTest(commit=commit, variable=variable, split=split): - default = variable or self.SMALL - self.assertEqual( - self.run_pool_step(variable=variable, split=split, commit=commit), - self.dispatch.routed_runner( - commit, default, self.pool.split_enabled(split) - ), - ) - - def test_an_explicit_choice_or_admin_variable_is_never_rerouted(self): - odd = self.COMMITS[1] - self.assertEqual(self.pool.resolve(self.SMALL, "", "", odd), self.SMALL) - self.assertEqual(self.pool.resolve("tart-canary", "", "", odd), "tart-canary") - self.assertEqual(self.pool.resolve("auto", "blacksmith-6vcpu-macos-15", "", odd), + 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") - self.assertEqual(self.pool.resolve("auto", "", "0", odd), self.SMALL) - self.assertEqual(self.run_pool_step(requested="tart-small", commit=odd), "tart-small") - def test_an_unresolved_commit_is_refused(self): - for commit in ("", "main", "B" * 40, "b" * 39): - with self.subTest(commit=commit), self.assertRaises(ValueError): - self.pool.resolve("auto", "", "", commit) - - def test_the_pool_job_resolves_after_the_commit_on_linux(self): + def test_the_pool_job_reads_actions_and_nothing_else(self): + self.assertEqual(self.workflow["permissions"], {"contents": "read"}) job = self.jobs["runner"] - self.assertIn("resolve-ref", job["needs"]) + self.assertEqual(job["permissions"], {"contents": "read", "actions": "read"}) self.assertIn("ubuntu", job["runs-on"]) self.assertEqual(job["outputs"]["label"], "${{ steps.pool.outputs.label }}") - self.assertEqual(self.pool_step()["id"], "pool") + 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", "")) - self.assertIn("scripts/ci/e2e_runner_pool.py", checkout["with"]["sparse-checkout"]) + 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 }}"