From a88d9af818b26afe46330f18662d97e666d172cf Mon Sep 17 00:00:00 2001 From: Francois Lanusse Date: Tue, 29 Sep 2026 06:26:41 -0700 Subject: [PATCH] Supervise Slurm workers with stock Dask nannies --- CLAUDE.md | 23 ++++++-- docs/api/compute.md | 20 +++++++ docs/api/sandbox.md | 5 ++ docs/user/cluster.md | 46 +++++++++++---- src/lightcone/engine/compute/slurm.py | 7 ++- .../engine/compute/slurm_bootstrap.py | 17 +++--- src/lightcone/engine/sandbox/policy.py | 5 ++ tests/test_compute_slurm.py | 59 +++++++++++++++++-- tests/test_sandbox_oci.py | 28 +++++++++ 9 files changed, 179 insertions(+), 31 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index ee4aaced..23386042 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -187,7 +187,7 @@ src/lightcone/ # namespace — NO __init__.py │ ├── local.py # local provider: validated OS process identities │ ├── local_runtime.py # the detached LocalCluster owner │ ├── slurm.py # Slurm provider: native commands, JobName + Comment - │ ├── slurm_bootstrap.py # one stock Dask process per Slurm rank + │ ├── slurm_bootstrap.py # stock Dask Nanny per rank; rank zero hosts scheduler │ └── output.py # recipe bytes through Dask events ├── sandbox/ # the exec boundary │ ├── __init__.py # the public surface (detect, run, scope, the types) @@ -1748,6 +1748,19 @@ already is. `connection_root` defaults to `~/.lightcone/compute` is chosen by the bootstrap on each node (`tempfile.gettempdir()`), never frozen from the driver's temporary directory. +**Slurm supervises workers through stock Dask Nannies (2026-09).** Each rank +runs a Nanny and its separate worker process; rank zero also hosts the scheduler. +Nanny defaults `OMP_NUM_THREADS`, `MKL_NUM_THREADS`, and `OPENBLAS_NUM_THREADS` +to `1`, preserving explicit launch environment values, and sandbox policy passes +the effective values into recipe containers. The step sets +`--kill-on-bad-exit=0 --wait=0` so a rank exit alone does not terminate the others. +This is worker-process recovery, not complete recipe isolation: site OOM policy +may still kill a step or job; dead schedulers and Nannies are not restarted; +Dask retries are not fenced from surviving recipe subprocesses. Keep +`memory_limit=0` because Dask does not account for subprocess RSS. New commands +still require all expected workers. Serial per-output Git/annex commits remain +the driver's responsibility; changing that persistence/provenance model is deferred. + **Compute uses one shared Pydantic model family.** `Catalog` loads directly into the `Connection`, `Offer`, `Resources`, `TimeLimits`, and `Startup` objects used by providers; do not introduce parallel configuration classes. Memory units are explicit: @@ -1812,10 +1825,10 @@ not a visibility limit; it defaults to zero and does not select a model. Recipe memory retains ASTRA units (`8Gi` binary, `8GB` decimal, no bare quantities), independently of compute's SkyPilot units. Thread slots remain a separate concurrency cap. Reservations are cooperative, not -per-command OS CPU/RAM limits or BLAS thread counts. Local Nanny defaults keep -OMP/MKL/OPENBLAS threads at one unless the launch environment overrides them; Slurm -uses direct workers and the job environment. Unsupported disk/model requests and -fractional CPU/GPU counts fail explicitly. Exact bytes are shared in `units.py`; +per-command OS CPU/RAM limits or BLAS thread counts. Local and Slurm Nanny defaults +keep OMP/MKL/OPENBLAS threads at one unless the launch environment overrides them; +sandbox policy forwards the effective values into containers. Unsupported +disk/model requests and fractional CPU/GPU counts fail explicitly. Exact bytes are shared in `units.py`; allocation durations are parsed in `compute.model`, with error conversion only at the CLI request boundary. diff --git a/docs/api/compute.md b/docs/api/compute.md index b780bd2c..4f7dddc7 100644 --- a/docs/api/compute.md +++ b/docs/api/compute.md @@ -62,6 +62,22 @@ borrow clients through the common API. Provider settings stay behind that seam. identity checks. `local_runtime.py` and `slurm_bootstrap.py` compose stock Dask components; they do not define custom workers or membership protocols. +The Slurm bootstrap runs one stock `Nanny` per rank, each with a separate worker +process; rank zero also hosts the scheduler. The Nanny restarts an exited worker +and supplies Dask's default `OMP_NUM_THREADS`, `MKL_NUM_THREADS`, and +`OPENBLAS_NUM_THREADS` values of `1`, preserving explicit launch environment +values. Sandbox policy forwards those effective values into recipe containers. +`memory_limit=0` remains deliberate: Dask's worker memory accounting excludes +the external recipe subprocesses. The payload uses +`srun --kill-on-bad-exit=0 --wait=0` to avoid terminating healthy ranks merely +because another rank exited. Site OOM policy can still terminate the step or +allocation, and there is no recovery for a dead scheduler or Nanny. Dask can +reschedule lost tasks, but surviving recipe subprocesses are not fenced from +those retries. Connection readiness still requires every expected worker. +See [Dask's Nanny](https://distributed.dask.org/en/stable/worker.html#nanny), +[Dask resilience](https://distributed.dask.org/en/stable/resilience.html), and +[Slurm's `srun` options](https://slurm.schedmd.com/srun.html). + Configured connection and scratch roots are resolved before managed paths are appended, so filesystem aliases such as a symlinked home directory are supported. Managed directories and credential files retain strict symlink, ownership, and @@ -180,6 +196,10 @@ See [GPU deployment requirements](../user/cluster.md#gpu-allocations). ## Execution output and teardown +The driver still saves each output in its own Git/annex commit while Dask runs +the submitted graph. This serial storage work can dominate many short recipes; +adding workers does not accelerate it. + `output.py` transports byte chunks through standard Dask events so detached workers' output reaches the invoking CLI. It uses the borrowed client's event topic, which the schedulers lc launches drop as soon as the client disconnects diff --git a/docs/api/sandbox.md b/docs/api/sandbox.md index 813a1ac5..8709d6b7 100644 --- a/docs/api/sandbox.md +++ b/docs/api/sandbox.md @@ -38,6 +38,11 @@ runtimes receive a CPU policy and an explanatory note. CPU containers remain supported on all runtimes and set `NVIDIA_VISIBLE_DEVICES=void` to override image defaults. See [GPU allocations](../user/cluster.md#gpu-allocations). +Container recipes and probes explicitly receive the worker's effective +`OMP_NUM_THREADS`, `MKL_NUM_THREADS`, and `OPENBLAS_NUM_THREADS` values when set. +This preserves Dask Nanny's numerical-library thread settings across the OCI +boundary. Container environment variables remain an explicit allowlist. + ## What must stay true - **`wrap` stays pure** — no temp files, no FDs, no global state diff --git a/docs/user/cluster.md b/docs/user/cluster.md index 35cd38a1..54ab30ba 100644 --- a/docs/user/cluster.md +++ b/docs/user/cluster.md @@ -244,16 +244,37 @@ class; it does not guarantee a queue wait. Inspect the resolved plan before laun lc compute launch --cpus 32+ --memory 128+ --num-nodes 2 --time 1h --dry-run ``` -One allocation contains one `srun` step with one process per node, bound with -`--cpu-bind=threads` to exactly the hardware threads Slurm allocated. Rank zero -composes standard Dask `Scheduler` and `Worker` objects, and every other rank -starts a standard `Worker`. A one-node allocation has both scheduler and worker. +One allocation contains one `srun` step with one rank per node, bound with +`--cpu-bind=threads` to exactly the hardware threads Slurm allocated. Each rank +starts a standard Dask `Nanny`, which supervises a separate `Worker` process. +Rank zero also runs the scheduler, so a worker process exiting does not take the +scheduler with it. A one-node allocation has both scheduler and worker. Neither serves a dashboard or any other HTTP route. The scheduler consumes part of the offered resources; `task_slots_per_node` controls Dask task concurrency independently of the allocation's logical CPUs. -Dask memory management is disabled because recipes run in external subprocesses; -Slurm supplies allocation containment and memory enforcement. Lightcone always -requests a finite native `--time`. Actual termination follows Slurm's + +Dask's Nanny defaults `OMP_NUM_THREADS`, `MKL_NUM_THREADS`, and +`OPENBLAS_NUM_THREADS` to `1`, preventing each concurrent recipe from requesting +the whole node's numerical-library threads. Values explicitly set in the launch +environment take precedence and also reach containerized recipes. For recipes +that need more threads, set those values before launching and declare enough +recipe CPUs for each task; `task_slots_per_node` can further cap concurrency. See +[Dask's Nanny environment settings](https://distributed.dask.org/en/stable/worker.html#nanny). + +The Nanny restarts an exited worker, and Dask can reschedule its tasks. The step +uses `--kill-on-bad-exit=0 --wait=0` so an exited rank does not itself trigger +termination of the remaining ranks. This does not provide complete failure +isolation: site OOM policy can still kill the step or allocation, and the +scheduler and Nanny processes are not restarted if they die. Dask memory +management remains disabled because it does not account for the recipe +subprocesses' memory. Reduce task concurrency for memory-heavy recipes; there is +no per-recipe memory limit. A lost worker can also leave a recipe subprocess +running while Dask reschedules its task, so retries do not guarantee exclusive +access to output files. New commands still require every expected worker to be +connected. See [Dask's failure behavior](https://distributed.dask.org/en/stable/resilience.html) +and [Slurm's step termination settings](https://slurm.schedmd.com/srun.html). + +Lightcone always requests a finite native `--time`. Actual termination follows Slurm's `OverTimeLimit` and `KillWait` policy, which can permit an unlimited overrun. Lightcone does not impose an independent Slurm runtime deadline or require a preflight time-policy query. @@ -448,13 +469,14 @@ its request. Slurm enforces the overall allocation, while local execution uses cooperative budgets. Leave capacity for the scheduler, workers, and other overhead when declaring recipe requirements. -A CPU reservation does not set numerical-library thread counts. Local clusters -use Dask's Nanny defaults of `1` for `OMP_NUM_THREADS`, `MKL_NUM_THREADS`, and -`OPENBLAS_NUM_THREADS` when those variables are unset. Set the variables before -`lc compute launch`, or in the recipe command, to choose another value. See +A CPU reservation does not set numerical-library thread counts. Local and Slurm +clusters use Dask's Nanny defaults of `1` for `OMP_NUM_THREADS`, `MKL_NUM_THREADS`, +and `OPENBLAS_NUM_THREADS` when those variables are unset. Container recipes +receive the worker's effective values too. Set the variables before +`lc compute launch`, or in the recipe command, to choose another value, and +declare enough recipe CPUs for those threads. See [Dask's defaults](https://docs.dask.org/en/stable/configuration.html#distributed.nanny.pre-spawn-environ.OMP_NUM_THREADS) and [environment precedence](https://distributed.dask.org/en/stable/_modules/distributed/nanny.html). -Slurm workers run directly without a Nanny and inherit the job's thread settings. ## Execution requirements and limits diff --git a/src/lightcone/engine/compute/slurm.py b/src/lightcone/engine/compute/slurm.py index 7cad67aa..8bdb7484 100644 --- a/src/lightcone/engine/compute/slurm.py +++ b/src/lightcone/engine/compute/slurm.py @@ -303,10 +303,11 @@ def _payload(self, plan: LaunchPlan, token: str) -> list[str]: "--ntasks-per-node=1", f"--cpus-per-task={plan.resources.cpus}", f"--gres={details['gres']}", - # One process per node holds the whole allocation, so binding to - # exactly its allocated hardware threads is the only useful mask. + # Each rank and its worker inherit the node's allocated CPU mask. "--cpu-bind=threads", - "--kill-on-bad-exit=1", + # Losing a rank must not terminate healthy ranks and the scheduler. + "--kill-on-bad-exit=0", + "--wait=0", details["python"], "-P", "-m", diff --git a/src/lightcone/engine/compute/slurm_bootstrap.py b/src/lightcone/engine/compute/slurm_bootstrap.py index 2ef71a3c..904b63ef 100644 --- a/src/lightcone/engine/compute/slurm_bootstrap.py +++ b/src/lightcone/engine/compute/slurm_bootstrap.py @@ -1,4 +1,4 @@ -"""One stock Dask worker per Slurm rank, with a scheduler alongside rank zero.""" +"""One Dask nanny per Slurm rank, with a scheduler alongside rank zero.""" from __future__ import annotations @@ -79,8 +79,8 @@ def _allocation(args: argparse.Namespace) -> tuple[Identity, int, int]: async def run(args: argparse.Namespace) -> None: - """Run standard asynchronous Scheduler/Worker contexts for this native rank.""" - from distributed import Scheduler, Worker + """Run standard asynchronous Scheduler/Nanny contexts for this native rank.""" + from distributed import Nanny, Scheduler identity, restarts, rank = _allocation(args) if args.gpus: @@ -113,6 +113,7 @@ async def run(args: argparse.Namespace) -> None: worker_options = { **address, "nthreads": args.task_slots, + # Recipes run in subprocesses whose RSS Dask does not account for. "memory_limit": 0, "resources": {"CPU": args.cpus, "MEMORY": args.memory_bytes, "GPU": args.gpus}, "local_directory": str(scratch), @@ -139,7 +140,9 @@ async def run(args: argparse.Namespace) -> None: write_private_json( directory / "identity.json", {**identity_values, "scheduler_id": scheduler.id} ) - async with Worker(scheduler.address, security=security, **worker_options): + # The nanny keeps worker exits out of the scheduler process and + # supplies Dask's numerical-library thread defaults to recipes. + async with Nanny(scheduler.address, security=security, **worker_options): await scheduler.finished() # type: ignore[no-untyped-call] finally: await scheduler.close() @@ -160,10 +163,10 @@ async def run(args: argparse.Namespace) -> None: raise ComputeError("scheduler rendezvous belongs to another allocation attempt") break security = load_security(directory) - async with Worker( + async with Nanny( scheduler_file=str(directory / "scheduler.json"), security=security, **worker_options - ) as worker: - await worker.finished() + ) as nanny: + await nanny.finished() def main() -> None: diff --git a/src/lightcone/engine/sandbox/policy.py b/src/lightcone/engine/sandbox/policy.py index 6128af66..3b2a062c 100644 --- a/src/lightcone/engine/sandbox/policy.py +++ b/src/lightcone/engine/sandbox/policy.py @@ -348,6 +348,11 @@ def home_overlay(tmp_home: Path, env_dir: Path, *, containerized: bool = False) } if containerized: overlay["UV_PROJECT_ENVIRONMENT"] = str(env_dir) + # Dask's Nanny configures these on the worker. Containers do not + # inherit them, so carry the effective values across the boundary. + for name in ("OMP_NUM_THREADS", "MKL_NUM_THREADS", "OPENBLAS_NUM_THREADS"): + if name in os.environ: + overlay[name] = os.environ[name] return overlay diff --git a/tests/test_compute_slurm.py b/tests/test_compute_slurm.py index d51e40ec..cec8c871 100644 --- a/tests/test_compute_slurm.py +++ b/tests/test_compute_slurm.py @@ -6,6 +6,7 @@ import asyncio import os import shlex +import signal import subprocess import sys import tempfile @@ -360,7 +361,8 @@ def test_sbatch_launch_owns_payload_and_scrubs_ambient_overrides( assert script.startswith("#!/bin/bash\nset -euo pipefail\numask 077\n") payload = shlex.split(script.splitlines()[-1]) assert payload[:2] == ["exec", "srun"] - assert "--kill-on-bad-exit=1" in payload and "--overlap" not in payload + assert "--kill-on-bad-exit=0" in payload and "--overlap" not in payload + assert "--wait=0" in payload assert "--cpu-bind=threads" in payload assert "--ntasks=2" in payload python = payload.index(sys.executable) @@ -1100,7 +1102,7 @@ def test_gpu_worker_advertises_native_capacity_with_the_native_mask( worker.__aenter__ = AsyncMock(return_value=worker) worker.finished = AsyncMock() factory = MagicMock(return_value=worker) - monkeypatch.setattr(distributed, "Worker", factory) + monkeypatch.setattr(distributed, "Nanny", factory) asyncio.run(slurm_bootstrap.run(args)) @@ -1183,6 +1185,12 @@ def test_standard_bootstrap_starts_scheduler_and_worker_on_rank_zero_and_worker_ directory = slurm.attempt_directory(connection, IDENTITY, 0) processes = [] client = None + environment = dict(os.environ) + thread_variables = ("OMP_NUM_THREADS", "MKL_NUM_THREADS", "OPENBLAS_NUM_THREADS") + for variable in thread_variables: + environment.pop(variable, None) + # Explicit launch settings still take precedence over Dask's defaults. + environment["OMP_NUM_THREADS"] = "2" try: for rank in (1, 0): with (tmp_path / f"rank-{rank}.log").open("wb") as log: @@ -1190,7 +1198,7 @@ def test_standard_bootstrap_starts_scheduler_and_worker_on_rank_zero_and_worker_ subprocess.Popen( argv, cwd=tmp_path, - env={**os.environ, **_bootstrap_env(rank)}, + env={**environment, **_bootstrap_env(rank)}, stdin=subprocess.DEVNULL, stdout=log, stderr=subprocess.STDOUT, @@ -1206,7 +1214,11 @@ def test_standard_bootstrap_starts_scheduler_and_worker_on_rank_zero_and_worker_ time.sleep(0.1) metadata = read_private_json(directory / "identity.json") client = open_client(directory, metadata["scheduler_id"], timeout=5) - client.wait_for_workers(2, timeout=10) + try: + client.wait_for_workers(2, timeout=20) + except TimeoutError: + logs = "\n".join((tmp_path / f"rank-{rank}.log").read_text() for rank in (0, 1)) + pytest.fail(f"bootstrap failed to start workers: {logs}") workers = client.scheduler_info()["workers"] assert {worker["name"] for worker in workers.values()} == {"lightcone-0", "lightcone-1"} assert all(worker["nthreads"] == 1 for worker in workers.values()) @@ -1221,6 +1233,45 @@ def test_standard_bootstrap_starts_scheduler_and_worker_on_rank_zero_and_worker_ assert set(client.run(lambda dask_worker: dask_worker.http_server.address).values()) == { "127.0.0.1" } + assert all(worker["nanny"] for worker in workers.values()) + for variable in thread_variables: + expected = "2" if variable == "OMP_NUM_THREADS" else "1" + # A fresh recipe-style subprocess must inherit the worker's limits. + observed = client.run( + lambda key: subprocess.check_output( + [sys.executable, "-P", "-c", "import os; print(os.environ['" + key + "'])"], + text=True, + ).strip(), + variable, + ) + assert set(observed.values()) == {expected} + # Kill each worker, including rank zero's: its nanny and the scheduler + # must survive, the other worker must stay available, and capacity returns. + for rank in (0, 1): + workers = client.scheduler_info()["workers"] + address = next( + address for address, info in workers.items() if info["name"] == f"lightcone-{rank}" + ) + peer = next(address_ for address_ in workers if address_ != address) + pid = client.run(os.getpid, workers=[address])[address] + assert pid not in {process.pid for process in processes} + os.kill(pid, signal.SIGKILL) + assert client.submit(sum, [3, 4], workers=[peer], pure=False).result(timeout=5) == 7 + deadline = time.monotonic() + 20 + while True: + info = client.scheduler_info() + recovered = info["workers"] + if len(recovered) == 2 and address not in recovered: + break + assert time.monotonic() < deadline, "nanny did not replace the killed worker" + time.sleep(0.1) + assert info["id"] == metadata["scheduler_id"] + assert peer in recovered + assert all(process.poll() is None for process in processes) + replacement = next(address_ for address_ in recovered if address_ != peer) + assert recovered[replacement]["resources"] == workers[address]["resources"] + future = client.submit(sum, [4, 5], workers=[replacement], pure=False) + assert future.result(timeout=5) == 9 client.shutdown() client = None for process in processes: diff --git a/tests/test_sandbox_oci.py b/tests/test_sandbox_oci.py index de6efac3..8f7f7221 100644 --- a/tests/test_sandbox_oci.py +++ b/tests/test_sandbox_oci.py @@ -236,6 +236,34 @@ def test_the_environment_is_an_allowlist_never_ambient( assert "--env=LC_SANDBOX=podman" in argv +@pytest.mark.parametrize("recipe", [False, True], ids=["probe", "recipe"]) +@pytest.mark.parametrize("configured", [False, True], ids=["unset", "configured"]) +def test_container_commands_preserve_worker_numerical_thread_settings( + root: Path, monkeypatch: pytest.MonkeyPatch, recipe: bool, configured: bool, +) -> None: + values = {"OMP_NUM_THREADS": "1", "MKL_NUM_THREADS": "2", "OPENBLAS_NUM_THREADS": "3"} + for key, value in values.items(): + if configured: + monkeypatch.setenv(key, value) + else: + monkeypatch.delenv(key, raising=False) + monkeypatch.setenv("AWS_SECRET_ACCESS_KEY", "not-forwarded") + write_dir = root / "results" / "baseline" if recipe else None + if write_dir is not None: + write_dir.mkdir() + with boundary.scope( + exec_policy( + root, env_dir=root / ".lightcone" / "venv", + containerized=True, write_dir=write_dir, + ) + ) as built: + argv = _backend(root).wrap(built, ["true"]) + for key, value in values.items(): + forwarded = [arg for arg in argv if arg.startswith(f"--env={key}=")] + assert forwarded == ([f"--env={key}={value}"] if configured else []) + assert not any("AWS_SECRET_ACCESS_KEY" in arg for arg in argv) + + def test_no_host_resolved_env_binary_in_the_argv(root: Path, policy: Policy) -> None: """The overlay travels as `--env` flags: a host path for `env` (NixOS keeps it under /run/current-system/sw) need not exist in the image,