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
21 changes: 21 additions & 0 deletions docs/CMUX_MINI_RUNNER.md
Original file line number Diff line number Diff line change
Expand Up @@ -425,6 +425,27 @@ in at the producer's root. So one root job per root per mini:
detached holder tied to the job's Runner.Worker and released by job-completed. A re-take of a root the
same job already holds (a compile restoring its own product) is a no-op. It exits 1 when the root is
still busy after the wait, and 2 for a bad root or outside a runner job.
- A second, different root is refused (exit 2): two jobs taking two roots in opposite orders would deadlock.
`take ROOT --switch` swaps instead.
- It waits for ROOT while still holding the old root, and lets the old one go only once ROOT is held.
A timeout (exit 1) leaves the job on its old root, never without one.
- Two switchers after each other's root both time out, so callers keep `--wait` short.
- On success it writes the new `CMUX_CI_CANONICAL_ROOT` to `$GITHUB_ENV`.
- It only works when each held root has a live holder of its own: one taken with `take`, or the
admission root of a `ROOT_SWITCHERS` class (compile-gui, test-e2e's `build`), which the hook hands to
a separate holder at admission.
- A build that finds a product to reuse at another root switches to it. The caller must be done with
the old root.
- `glaeda-canonical-root take-gui [--wait S]` holds the mini's gui token from that step to the end of the
job, with a third holder (`<holder>-gui.pid`) that job-completed releases. It is a no-op when the job
already holds gui from admission. It exits 0 when held, 1 when still taken after the wait, and 2
outside a job.
- A job in take-gui holds a root, while a gui job waiting in take-root holds gui: opposite lock orders.
So every take-root waiter writes `capacity/root-k.want-<pid>`, containing `gui` when its job holds
the gui token. take-gui exits 3 at once when a gui holder is waiting for a root this job holds. The
caller then leaves its console-session work to another job and finishes, which frees the root.
Waiters rewrite their markers every second. take-gui removes a marker older than 10 s or with a dead
pid, so a killed waiter's leftover cannot make it give way.
- Jobs the hook does not know (seed-swiftpm-manifests, anything new) are pinned to root 1, because they use
/private/tmp/cmux-ci themselves. Ids that other workflows reuse (`build`, `test`, `lint`) are classed by
(workflow file, job id) from `GITHUB_WORKFLOW_REF`, so with `canonicalRoots` above 1 test-e2e's jobs are
Expand Down
20 changes: 13 additions & 7 deletions scripts/glaeda-cmux-runner
Original file line number Diff line number Diff line change
Expand Up @@ -470,15 +470,21 @@ def hook_wrappers(ctx: Context) -> dict[str, bytes]:
capacity = shlex.quote(os.fspath(Path(os.environ.get("GLAEDA_FLEET_DIR") or "/Users/Shared/cmux-build-fleet") / "capacity"))
home = ctx.home if isinstance(getattr(ctx, "home", None), Path) else Path.home()
state = shlex.quote(os.fspath(home / ".local/state/glaeda/cmux-runner"))
dirs = f" --capacity-dir {capacity} --state-dir {state}"
root = (
"#!/bin/sh\n"
"# glaeda-canonical-root take ROOT [--wait SECONDS] (generated, receipt-owned): hold canonical root ROOT\n"
"# (/private/tmp/cmux-ci[-N]) for the rest of this job; the job-started hook links it into the fleet bin.\n"
"[ \"${1:-}\" = take ] && [ -n \"${2:-}\" ] || "
"{ echo 'usage: glaeda-canonical-root take ROOT [--wait SECONDS]' >&2; exit 2; }\n"
"root=$2; shift 2\n"
f"exec {python} {hook} take-root --root \"$root\" --canonical-roots {roots_n}"
f" --capacity-dir {capacity} --state-dir {state} \"$@\"\n"
"# glaeda-canonical-root take ROOT [--wait SECONDS] [--switch] (generated, receipt-owned): hold canonical\n"
"# root ROOT (/private/tmp/cmux-ci[-N]) for the rest of this job; take-gui [--wait SECONDS] holds the\n"
"# mini's gui token (exit 3: gave way to a gui job waiting for this job's root). The job-started hook links\n"
"# it into the fleet bin.\n"
"usage() { echo 'usage: glaeda-canonical-root take ROOT [--wait SECONDS] [--switch]"
" | take-gui [--wait SECONDS]' >&2; exit 2; }\n"
"case \"${1:-}\" in\n"
"take) [ -n \"${2:-}\" ] || usage; root=$2; shift 2\n"
f" exec {python} {hook} take-root --root \"$root\" --canonical-roots {roots_n}{dirs} \"$@\" ;;\n"
f"take-gui) shift; exec {python} {hook} take-gui{dirs} \"$@\" ;;\n"
"*) usage ;;\n"
"esac\n"
)
return {"job-started.sh": started.encode(), "job-completed.sh": completed.encode(), "listen.sh": listen.encode(),
"glaeda-canonical-root": root.encode()}
Expand Down
172 changes: 153 additions & 19 deletions scripts/glaeda-cmux-runner-hook
Original file line number Diff line number Diff line change
Expand Up @@ -192,7 +192,17 @@ CANONICAL_ROOT_ENV = "CMUX_CI_CANONICAL_ROOT"
# With more than one root, these restore a producer's product into the producer's root (#filePath is baked
# in there), so they take that root from their restore step, not a token-chosen one at job start.
ROOT_CONSUMERS = ("gui", "product")
# Producers whose admission root sits in a holder of its own (root_holder_file), so a step can swap it for the
# root of a product the job reuses instead of compiling (glaeda-canonical-root take ROOT --switch): the old
# root is let go once the new one is held, so a switch that times out leaves the job on its old root.
ROOT_SWITCHERS = ("compile-gui",)
TAKE_ROOT_POLL_S = 1.0
# glaeda-canonical-root take-gui: a job takes the gui token only for the steps that use the console session.
# It holds a root while it waits, and a gui job (gui at admission) may be waiting in take-root for that root:
# the lock orders are opposite. So a take-root waiter leaves a marker, capacity/root-k.want-<pid> ("gui" when
# its job holds the gui token), and take-gui gives way (TAKE_GUI_GAVE_WAY) when a gui holder wants its root.
TAKE_GUI_GAVE_WAY = 3
WANT_FRESH_S = 10.0 # a waiter rewrites its marker every TAKE_ROOT_POLL_S; an older one is stale


def canonical_root(k: int) -> str:
Expand Down Expand Up @@ -634,10 +644,16 @@ def ensure_root_shim(bin_dir: Path) -> None:
os.replace(tmp, link)


def take_root(spec: str, wait_s: float, capacity_dir: Path, state_dir: Path, roots: int = 99) -> int:
def take_root(spec: str, wait_s: float, capacity_dir: Path, state_dir: Path, roots: int = 99,
switch: bool = False) -> int:
"""glaeda-canonical-root take ROOT: hold canonical root ROOT for the rest of this job (cmux#14338), waiting
up to wait_s for it. A root this job already holds is a no-op; a second, different root is refused (two
jobs taking two roots in opposite orders would deadlock). Prints the root's path."""
jobs taking two roots in opposite orders would deadlock), unless `switch` and every root the job holds
has a live holder of its own (ROOT_SWITCHERS, or an earlier take-root). A switch waits for ROOT while
still holding the old roots and lets them go only once ROOT is held, so a timeout leaves the job where
it was (exit 1) and never without a root; two switchers after each other's root both time out, so
callers keep --wait short. On success it points $GITHUB_ENV's CMUX_CI_CANONICAL_ROOT at ROOT. Prints
the root's path."""
k = root_number(spec, roots)
if k is None:
print(f"{PREFIX}: take-root: {spec!r} is not one of this mini's {roots} canonical root(s) "
Expand All @@ -653,29 +669,119 @@ def take_root(spec: str, wait_s: float, capacity_dir: Path, state_dir: Path, roo
if name in held:
print(path)
return 0
if held:
print(f"{PREFIX}: take-root: this job already holds {', '.join(held)}; one root per job", file=sys.stderr)
if held and not (switch and all(holder_alive(root_holder_file(state_dir, old)) for old in held)):
print(f"{PREFIX}: take-root: this job already holds {', '.join(held)}; one root per job"
+ ("" if switch else " (--switch lets go of a root with its own holder first)"), file=sys.stderr)
return 2
worker = runner_worker_pid()
if worker is None:
print(f"{PREFIX}: take-root: not inside a runner job (no Runner.Worker above this step)", file=sys.stderr)
return 2
deadline = time.monotonic() + max(0.0, wait_s)
want = capacity_dir / f"{name}.want-{os.getpid()}"
try:
while True:
fd = lock_file(capacity_dir / f"{name}.token", fcntl.LOCK_EX)
if fd is not None:
break
if time.monotonic() >= deadline:
print(f"{PREFIX}: take-root: {path} is still in use after {wait_s:g} s", file=sys.stderr)
return 1
# take-gui reads it: a gui holder waiting for a root must not be waited on. Rewritten every poll,
# so a marker a killed waiter left goes stale (WANT_FRESH_S) even if its pid comes back.
with contextlib.suppress(OSError):
want.write_text("gui\n" if gui_file(state_dir).exists() else "\n")
time.sleep(TAKE_ROOT_POLL_S)
finally:
with contextlib.suppress(OSError):
want.unlink()
ok, note = hold([fd], worker, state_dir, pid_file=root_holder_file(state_dir, name))
if not ok:
print(f"{PREFIX}: take-root: {note}", file=sys.stderr)
return 1
if held: # a switch: ROOT is held, so the old roots can go, after the record stops naming them
with contextlib.suppress(OSError):
roots_file(state_dir).write_text(name + "\n")
for old in held:
release_holder(root_holder_file(state_dir, old))
print(f"{PREFIX}: take-root: let go of {', '.join(held)}; {export_root(name)}", file=sys.stderr)
else:
with contextlib.suppress(OSError), open(roots_file(state_dir), "a", encoding="utf-8") as record:
record.write(name + "\n")
print(path)
return 0


def holder_alive(pid_file: Path) -> bool:
"""A holder pid file whose holder is still running (a killed one leaves its file behind)."""
try:
return pid_alive(int(pid_file.read_text().strip()))
except (OSError, ValueError):
return False


def gui_holder_wants(capacity_dir: Path, roots: list[str]) -> str | None:
"""A root in `roots` that a job holding the gui token is waiting for in take-root, from its marker."""
for name in roots:
for marker in capacity_dir.glob(f"{name}.want-*"):
try:
pid, mark = int(marker.name.rsplit("-", 1)[1]), marker.read_text()
except (OSError, ValueError):
continue
try:
fresh = time.time() - marker.stat().st_mtime < WANT_FRESH_S
except OSError:
continue
if not fresh or not pid_alive(pid):
with contextlib.suppress(OSError):
marker.unlink() # a waiter that was killed
elif mark.strip() == "gui":
return name
return None


def take_gui(wait_s: float, capacity_dir: Path, state_dir: Path) -> int:
"""glaeda-canonical-root take-gui: hold this mini's gui token (one console session) for the rest of this
job, waiting up to wait_s. A job that already holds it (from admission) is a no-op. Exits 0 when held, 1
when still taken after the wait, TAKE_GUI_GAVE_WAY when the holder waits for a root this job holds (see
TAKE_GUI_GAVE_WAY: the caller then leaves its console-session work to another job), 2 outside a job."""
if not os.environ.get("RUNNER_NAME"):
print(f"{PREFIX}: take-gui: RUNNER_NAME is not set (not a runner job step)", file=sys.stderr)
return 2
if gui_file(state_dir).exists():
return 0
worker = runner_worker_pid()
if worker is None:
print(f"{PREFIX}: take-gui: not inside a runner job (no Runner.Worker above this step)", file=sys.stderr)
return 2
deadline = time.monotonic() + max(0.0, wait_s)
while True:
fd = lock_file(capacity_dir / f"{name}.token", fcntl.LOCK_EX)
fd = lock_file(capacity_dir / "gui.token", fcntl.LOCK_EX)
if fd is not None:
break
held: list[str] = []
with contextlib.suppress(OSError):
held = roots_file(state_dir).read_text().split()
wanted = gui_holder_wants(capacity_dir, held)
if wanted:
print(f"{PREFIX}: take-gui: the gui token's holder is waiting for {wanted}, which this job holds; "
"giving way", file=sys.stderr)
return TAKE_GUI_GAVE_WAY
if time.monotonic() >= deadline:
print(f"{PREFIX}: take-root: {path} is still in use after {wait_s:g} s", file=sys.stderr)
print(f"{PREFIX}: take-gui: the gui token is still taken after {wait_s:g} s", file=sys.stderr)
return 1
time.sleep(TAKE_ROOT_POLL_S)
held, note = hold([fd], worker, state_dir, pid_file=root_holder_file(state_dir, name))
if not held:
print(f"{PREFIX}: take-root: {note}", file=sys.stderr)
main = holder_file(state_dir)
ok, note = hold([fd], worker, state_dir, pid_file=main.with_name(f"{main.stem}-gui.pid"))
if not ok:
print(f"{PREFIX}: take-gui: {note}", file=sys.stderr)
return 1
try:
gui_file(state_dir).write_text("gui\n")
except OSError as error: # unrecorded, a re-take would wait on itself and markers would miss it
release_holder(main.with_name(f"{main.stem}-gui.pid"))
print(f"{PREFIX}: take-gui: cannot record the gui token ({type(error).__name__})", file=sys.stderr)
return 1
with contextlib.suppress(OSError), open(roots_file(state_dir), "a", encoding="utf-8") as record:
record.write(name + "\n")
print(path)
return 0


Expand Down Expand Up @@ -812,6 +918,7 @@ def take_capacity(host_lock: str, capacity_dir: Path, total: int, job: str, watc
if waiting: # flock has no writer preference: stop admitting so the build worker gets in
return False, f"capacity: a fleet build is waiting for the host (pid {waiting[0]})"
taken: list[str] = []
named: dict[str, int] = {}
for token in tokens:
# persistent-dd comes in compile_slots copies once compile admission keeps its state per
# runner; copy 0 keeps the original file name. A mini whose runners disagree on the count runs
Expand All @@ -832,6 +939,7 @@ def take_capacity(host_lock: str, capacity_dir: Path, total: int, job: str, watc
if fd is not None:
fds.append(fd)
taken.append(name)
named[name] = fd
break
else:
what = "canonical root" if token == "root" else token
Expand All @@ -847,10 +955,22 @@ def take_capacity(host_lock: str, capacity_dir: Path, total: int, job: str, watc
got += 1
if got < units:
return False, (f"capacity: {got} of {total} units free, {job or 'an unknown job'} ({klass}) needs {units}")
held, note = hold(fds, watch_pid, state_dir, drop=(admission,))
fds = [] # the holder has them now, or hold closed them
what = "+".join([f"{units}/{total} units", *taken])
root = next((name for name in taken if name.startswith("root-")), None)
split = named[root] if root and klass in ROOT_SWITCHERS else None
if split is not None: # its own holder (ROOT_SWITCHERS); the main holder must not inherit it
fds.remove(split)
held, note = hold(fds, watch_pid, state_dir, drop=(admission,) if split is None else (admission, split))
fds = [] if split is None else [split] # the holder has them now, or hold closed them
if held and split is not None:
held, why = hold(fds, watch_pid, state_dir, drop=(admission,), pid_file=root_holder_file(state_dir, root))
fds = []
if not held:
release_holder(holder_file(state_dir))
note = why
what = "+".join([f"{units}/{total} units", *taken])
if held and "gui" in taken:
with contextlib.suppress(OSError):
gui_file(state_dir).write_text("gui\n")
if held and root:
note += f"; {export_root(root)}"
with contextlib.suppress(OSError):
Expand Down Expand Up @@ -932,6 +1052,12 @@ def roots_file(state_dir: Path) -> Path:
return main.with_name(main.stem + ".roots")


def gui_file(state_dir: Path) -> Path:
"""Present while this runner's job holds the gui token (from admission or take-gui)."""
main = holder_file(state_dir)
return main.with_name(main.stem + ".gui")


def root_holder_file(state_dir: Path, name: str) -> Path:
main = holder_file(state_dir)
return main.with_name(f"{main.stem}-{name}.pid")
Expand All @@ -941,10 +1067,13 @@ def release_host_lock(state_dir: Path) -> str:
"""job-completed: tell this job's holders to let go (the admission holder and any take-root holders),
and wait briefly for them."""
roots_file(state_dir).unlink(missing_ok=True)
gui_file(state_dir).unlink(missing_ok=True)
main = holder_file(state_dir)
extra = [release_holder(path) for path in sorted(state_dir.glob(f"{main.stem}-root-*.pid"))]
gui = [release_holder(path) for path in state_dir.glob(f"{main.stem}-gui.pid")]
note = release_holder(main)
return note + (f"; {len(extra)} canonical root holder(s) released" if extra else "")
return (note + (f"; {len(extra)} canonical root holder(s) released" if extra else "")
+ ("; the gui token holder released" if gui else ""))


def release_holder(pid_file: Path) -> str:
Expand Down Expand Up @@ -1825,12 +1954,14 @@ def disk_pressure(explicit: str | None) -> str:

def main(argv: list[str] | None = None) -> int:
p = argparse.ArgumentParser(prog=PREFIX)
p.add_argument("phase", choices=("job-started", "job-completed", "check", "listen", "take-root"),
p.add_argument("phase", choices=("job-started", "job-completed", "check", "listen", "take-root", "take-gui"),
help="check: the eligibility gate only, read-only (no lock, no rustup change); "
"listen: run the runner's run.sh under the listener gate (the LaunchAgent's program)")
p.add_argument("--runner-dir", type=Path, help="listen: the runner directory whose run.sh the gate runs")
p.add_argument("--root", help="take-root: the canonical root to hold, /private/tmp/cmux-ci[-N] or N")
p.add_argument("--wait", type=float, default=0, help="take-root: seconds to wait for the root")
p.add_argument("--wait", type=float, default=0, help="take-root, take-gui: seconds to wait for the token")
p.add_argument("--switch", action="store_true",
help="take-root: let go of the root this job holds (its own holder) and take ROOT instead")
p.add_argument("--no-waiters", action="store_true",
help="listen: only an exclusive holder or a reservation stops the listener (one runner per mini)")
p.add_argument("--adopt", type=int, default=0, help=argparse.SUPPRESS)
Expand Down Expand Up @@ -1886,7 +2017,9 @@ def main(argv: list[str] | None = None) -> int:
if args.phase == "take-root":
if not args.root:
p.error("take-root needs --root")
return take_root(args.root, args.wait, args.capacity_dir, args.state_dir, args.canonical_roots)
return take_root(args.root, args.wait, args.capacity_dir, args.state_dir, args.canonical_roots, args.switch)
if args.phase == "take-gui":
return take_gui(args.wait, args.capacity_dir, args.state_dir)
if args.phase == "check":
why, want = node_refusal(args.fleet_class)
if not why:
Expand Down Expand Up @@ -1918,6 +2051,7 @@ def main(argv: list[str] | None = None) -> int:
# The wrapper execs this script, so the parent is the job's Runner.Worker.
if args.capacity_units > 0:
roots_file(args.state_dir).unlink(missing_ok=True) # a crashed job's record must not look held
gui_file(args.state_dir).unlink(missing_ok=True)
ensure_root_shim(Path(FLEET_DIR) / "bin")
# the repository decide() admitted: the runner's variable, else the event payload's
repo = os.environ.get("GITHUB_REPOSITORY") or _repo_name((event or {}).get("repository")
Expand Down
Loading
Loading