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
23 changes: 18 additions & 5 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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.

Expand Down
20 changes: 20 additions & 0 deletions docs/api/compute.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions docs/api/sandbox.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
46 changes: 34 additions & 12 deletions docs/user/cluster.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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

Expand Down
7 changes: 4 additions & 3 deletions src/lightcone/engine/compute/slurm.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
17 changes: 10 additions & 7 deletions src/lightcone/engine/compute/slurm_bootstrap.py
Original file line number Diff line number Diff line change
@@ -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

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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),
Expand All @@ -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()
Expand All @@ -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:
Expand Down
5 changes: 5 additions & 0 deletions src/lightcone/engine/sandbox/policy.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down
59 changes: 55 additions & 4 deletions tests/test_compute_slurm.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
import asyncio
import os
import shlex
import signal
import subprocess
import sys
import tempfile
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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))

Expand Down Expand Up @@ -1183,14 +1185,20 @@ 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:
processes.append(
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,
Expand All @@ -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())
Expand All @@ -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:
Expand Down
Loading
Loading