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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 34 additions & 8 deletions scripts/ci/persistent_compile_fleet.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import json
import os
import platform
import select
import shutil
import subprocess
import sys
Expand Down Expand Up @@ -614,11 +615,35 @@ def default_runner_name(local: LocalState) -> str:
return f"{node}-persistent-compile" if node else f"{platform.node().split('.')[0]}-persistent-compile"


def worker_running(directory: Path) -> bool:
def worker_pids(directory: Path) -> list[int]:
"""A Runner.Worker process exists only while this runner holds a job."""
result = subprocess.run(["pgrep", "-f", os.fspath(directory / "bin" / "Runner.Worker")],
capture_output=True, check=False)
return result.returncode == 0
capture_output=True, text=True, check=False)
return [int(pid) for pid in result.stdout.split()] if result.returncode == 0 else []


def wait_for_workers(directory: Path, deadline: float) -> None:
"""Block until no job is running here or the deadline passes, woken by the worker's exit.

kqueue's NOTE_EXIT fires the moment the process ends, which keeps the window in
which GitHub can assign another job before the stop as short as it can be.
"""
while (pids := worker_pids(directory)) and time.monotonic() < deadline:
queue = select.kqueue()
try:
watched = 0
for pid in pids:
try:
queue.control([select.kevent(pid, filter=select.KQ_FILTER_PROC,
flags=select.KQ_EV_ADD | select.KQ_EV_ONESHOT,
fflags=select.KQ_NOTE_EXIT)], 0)
watched += 1
except ProcessLookupError:
pass # already gone
if watched:
queue.control(None, 1, max(0.0, deadline - time.monotonic()))
finally:
queue.close()


def register_runner(name: str, token: str | None) -> None:
Expand Down Expand Up @@ -825,6 +850,10 @@ def glaeda_transition(local: LocalState, target: str, reason: str | None = None)
print("no Glaeda enrollment on this mini; skipping the Glaeda state change")
return
if local.enrollment.get("state") == target:
kept = local.enrollment.get("quarantineReason")
if reason is not None and reason != kept:
# Glaeda has no quarantined -> quarantined transition to record a new reason.
print(f"Glaeda: {local.enrollment.get('nodeId')} is already {target} for {kept}; that reason stays")
return
if local.glaeda is None:
raise Failure("cannot find the Glaeda checkout; set GLAEDA_ROOT or pass --glaeda-root")
Expand Down Expand Up @@ -857,12 +886,9 @@ def stop_taking_jobs(args: argparse.Namespace, target: str, reason: str | None =
except Failure as error:
glaeda_error = error
if local.runner_configured:
deadline = time.monotonic() + (0 if args.now else DRAIN_WAIT_SECONDS)
if worker_running(directory) and not args.now:
if not args.now and worker_pids(directory):
print(f"{name} is running a job; stopping as soon as it finishes (--now stops it immediately)")
# Poll tightly: between the job ending and the stop, GitHub can assign another.
while worker_running(directory) and time.monotonic() < deadline:
time.sleep(2)
wait_for_workers(directory, time.monotonic() + DRAIN_WAIT_SECONDS)
stop_service(directory)
print(f"{name} is stopped and stays stopped across reboots.")
elif glaeda_error is None:
Expand Down
64 changes: 56 additions & 8 deletions tests/test_ci_persistent_compile_fleet.py
Original file line number Diff line number Diff line change
Expand Up @@ -428,12 +428,9 @@ def test_doctor_says_turn_it_off_when_routing_into_an_empty_fleet(self) -> None:

def test_doctor_flags_a_pilot_with_an_empty_cohort(self) -> None:
github = fleet_state(runners=[runner()], variables={fleet.SELECTOR_VARIABLE: "pilot"})
lines = dict(fleet.doctor_lines(github, None)[0])["GitHub"]
self.assertIn(fleet.Line(False, "routing: pilot for (empty cohort: nothing routes)"), lines)

def test_doctor_points_an_empty_pilot_at_a_pr(self) -> None:
github = fleet_state(runners=[runner()], variables={fleet.SELECTOR_VARIABLE: "pilot"})
self.assertIn("scripts/persistent-compile pilot", fleet.doctor_lines(github, None)[1])
sections, nxt = fleet.doctor_lines(github, None)
self.assertIn(fleet.Line(False, "routing: pilot for (empty cohort: nothing routes)"), dict(sections)["GitHub"])
self.assertIn("scripts/persistent-compile pilot", nxt)


class XcodePin(unittest.TestCase):
Expand Down Expand Up @@ -510,7 +507,7 @@ def test_quarantine_also_stops_the_runner(self) -> None:
with mock.patch.object(fleet, "require_mac"), \
mock.patch.object(fleet, "read_local", return_value=mini()), \
mock.patch.object(fleet, "glaeda_transition", side_effect=lambda l, t, r=None: transitions.append((t, r))), \
mock.patch.object(fleet, "worker_running", return_value=False), \
mock.patch.object(fleet, "worker_pids", return_value=[]), \
mock.patch.object(fleet, "stop_service", side_effect=stopped.append), mock.patch("builtins.print"):
self.assertEqual(fleet.main(["quarantine", "disk_pressure"]), 0)
self.assertEqual(transitions, [("quarantined", "disk_pressure")])
Expand All @@ -521,7 +518,7 @@ def test_the_runner_stops_even_when_glaeda_refuses(self) -> None:
with mock.patch.object(fleet, "require_mac"), \
mock.patch.object(fleet, "read_local", return_value=mini()), \
mock.patch.object(fleet, "glaeda_transition", side_effect=fleet.Failure("unsupported")), \
mock.patch.object(fleet, "worker_running", return_value=False), \
mock.patch.object(fleet, "worker_pids", return_value=[]), \
mock.patch.object(fleet, "stop_service", side_effect=stopped.append), \
mock.patch("builtins.print"), mock.patch("sys.stderr"):
self.assertEqual(fleet.main(["quarantine", "hardware_failure"]), 1)
Expand All @@ -538,6 +535,57 @@ def test_a_refused_quarantine_without_a_runner_does_not_claim_one_stopped(self)
self.assertNotIn("stopped", str(caught.exception))
self.assertIn("no runner is configured", str(caught.exception))

def test_drain_waits_on_the_worker_exit_not_a_timer(self) -> None:
# Two looks: a job is running, then the worker has exited.
pids = iter([[4242], []])
registered, waits = [], []

class Queue:
def control(self, changes, max_events, timeout=None):
if changes:
registered.extend(event.ident for event in changes)
else:
waits.append(timeout)
return []

def close(self):
pass

fake_select = mock.Mock(kqueue=Queue, KQ_FILTER_PROC=-5, KQ_EV_ADD=1, KQ_EV_ONESHOT=16, KQ_NOTE_EXIT=1,
kevent=lambda ident, **_: mock.Mock(ident=ident))
with mock.patch.object(fleet, "select", fake_select), \
mock.patch.object(fleet, "worker_pids", side_effect=lambda _: next(pids)):
fleet.wait_for_workers(Path("/r"), fleet.time.monotonic() + 60)
self.assertEqual(registered, [4242])
self.assertEqual(len(waits), 1)
self.assertGreater(waits[0], 0)

def test_a_worker_gone_before_registration_is_skipped(self) -> None:
pids = iter([[4242], []])

class Queue:
def control(self, changes, max_events, timeout=None):
if changes:
raise ProcessLookupError
raise AssertionError("waited on a process that was already gone")

def close(self):
pass

fake_select = mock.Mock(kqueue=Queue, KQ_FILTER_PROC=-5, KQ_EV_ADD=1, KQ_EV_ONESHOT=16, KQ_NOTE_EXIT=1,
kevent=lambda ident, **_: ident)
with mock.patch.object(fleet, "select", fake_select), \
mock.patch.object(fleet, "worker_pids", side_effect=lambda _: next(pids)):
fleet.wait_for_workers(Path("/r"), fleet.time.monotonic() + 60)

def test_requarantine_says_the_old_reason_is_kept(self) -> None:
# Glaeda has no quarantined -> quarantined transition, so a new reason is not recorded.
local = mini(enrollment={"nodeId": "n", "state": "quarantined", "quarantineReason": "disk_pressure"})
with mock.patch("builtins.print") as printed, mock.patch.object(fleet, "run_checked") as run:
fleet.glaeda_transition(local, "quarantined", "hardware_failure")
run.assert_not_called()
self.assertIn("disk_pressure", " ".join(str(c.args[0]) for c in printed.call_args_list))

def test_resume_points_a_quarantined_mini_at_up(self) -> None:
local = mini(enrollment={"nodeId": "n", "state": "quarantined", "quarantineReason": "disk_pressure"})
with mock.patch.object(fleet, "require_mac"), mock.patch.object(fleet, "read_local", return_value=local), \
Expand Down