From 85b8d8dc79cc4053b00874f82951e68373df1d4e Mon Sep 17 00:00:00 2001 From: Francois Lanusse Date: Tue, 29 Sep 2026 08:25:17 -0700 Subject: [PATCH 1/4] Streamline local compute launch and enforce local policy --- README.md | 3 +- docs/api/compute.md | 24 ++- docs/architecture.md | 6 +- docs/cli/compute.md | 58 ++++-- docs/user/cluster.md | 89 +++++--- docs/user/getting-started.md | 5 +- src/lightcone/cli/compute.py | 65 ++++-- src/lightcone/engine/compute/__init__.py | 101 ++++++--- src/lightcone/engine/compute/catalog.py | 72 +++++-- src/lightcone/engine/compute/local.py | 47 +++++ src/lightcone/engine/compute/local_runtime.py | 3 + tests/test_compute.py | 191 +++++++++++++++++- tests/test_compute_local.py | 99 +++++++-- 13 files changed, 620 insertions(+), 143 deletions(-) diff --git a/README.md b/README.md index fd5a9679..04bccffc 100644 --- a/README.md +++ b/README.md @@ -32,8 +32,7 @@ cd my-analysis uv add numpy # When you are done with your edits, commit: git add -A && git commit -m "First analysis" -CLUSTER=$(lc compute launch --cpus 1 --memory 1) -lc compute status "$CLUSTER" --wait +CLUSTER=$(lc compute launch --wait) # Generate outputs with full provenance tracking lc materialize "$CLUSTER" lc compute down "$CLUSTER" diff --git a/docs/api/compute.md b/docs/api/compute.md index 4f7dddc7..c5434843 100644 --- a/docs/api/compute.md +++ b/docs/api/compute.md @@ -7,7 +7,7 @@ It owns no service, registry, or saved current-cluster selection. | Symbol | Contract | |---|---| | `Request.parse(...)` | Exact/minimum CPU and memory requests, exact accelerator type/count, node count, walltime, startup class. | -| `Catalog.load(path)` | Ordered fixed shapes and stable connection namespaces; use the built-in local catalog only when the implicit default file is absent. | +| `Catalog.load(path)` | Ordered fixed shapes and stable connection namespaces; apply local defaults or disable policy alongside configured offers. | | `Compute.plan(request, *, name=None)` | Select an eligible offer and freeze its native launch settings and optional name without allocation. | | `Compute.launch(plan)` | Check names across native authorities, generate one if omitted, submit once, and return a self-contained `Identity`. | | `Compute.discover()` | Snapshots and per-connection errors, querying each authority once. | @@ -16,15 +16,22 @@ It owns no service, registry, or saved current-cluster selection. | `connect(cluster_id, timeout=10)` | Resolve a name or full ID; borrow a standard Dask client, closing the client but never the allocation. | | `Provider` | `plan`, `launch`, `discover`, `inspect`, `connect`, `terminate`. | -`Catalog.load()` defaults to `~/.lightcone/compute.yaml`. When that implicit file -is absent, the built-in catalog exposes a `local` offer: one CPU, 1 GiB, one node, -fast startup, 30-minute default and two-hour maximum lifetime. GPU offers require -an explicit catalog. It creates no configuration file or allocation. Configured -catalogs replace it completely. +`Catalog.load()` defaults to `~/.lightcone/compute.yaml`. The built-in `local` +offer uses detected usable CPUs and RAM, one node, fast startup, and a 30-minute +default/two-hour maximum lifetime. `local.resources` overrides its CPU/RAM budget; +`local.enabled: false` blocks local launch and execution while retaining connections +for inspection and termination. Remote catalogs retain the implicit local offer +unless disabled. Explicit local connections supply their own offers instead and +cannot be combined with `local.resources`. GPU offers require explicit configuration. +Loading creates no configuration file or allocation. Missing paths selected through an argument or `LC_COMPUTE_CONFIG`, unreadable files, and invalid catalogs remain errors. Stable connection namespaces let separate invocations discover and attach to the same local allocations. +`Compute.plan_local()` selects only local offers and defaults the name to `local`. +CLI `launch --wait` waits through `Compute.status` using the accepted immutable ID; +errors retain that ID without resubmission or termination. + `model.py` defines the shared Pydantic models: `Connection`, `Offer`, `Resources`, `Accelerator`, `TimeLimits`, `Startup`, `Request`, `Identity`, `LaunchPlan`, and `Snapshot`. `Catalog` validates YAML directly into these objects, which providers also use. @@ -208,6 +215,11 @@ command. A driver that exits before every task reports says so with `UNSTOPPED`: closing a client cannot prove that a remote subprocess stopped. Probes preserve both streams; materialization sends recipe output to stderr to leave stdout for its report. +A host-local OS file lock limits local allocations to one per user, independent of +connection roots and namespaces. The launcher acquires it before spawning and passes +the descriptor to the detached owner, which retains it for its lifetime. The lock +file is never unlinked; OS process exit releases the lock without stale-lock cleanup. + Local teardown drains the allocation's validated process group rather than assuming the owner's exit proves every child stopped. Boot UUID, UID, process session and the exact command containing a random allocation token establish diff --git a/docs/architecture.md b/docs/architecture.md index 4c47acf4..f4d72608 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -169,8 +169,10 @@ config-blob id, never a tag. `engine.compute` owns allocation lifecycle through a small provider protocol. A YAML catalog supplies ordered resource offers and stable native service namespaces. When the implicit default file is absent, a built-in local catalog -provides one CPU and 1 GiB without setup. GPU offers need an explicit catalog, -which replaces the built-in defaults. Missing explicit paths and invalid files +provides all detected usable CPUs and RAM without setup. The `local` policy can +override that budget or disable local compute. Remote catalogs retain the default +local offer; explicit local connections supply their own offers. GPU offers need +an explicit catalog. Missing explicit paths and invalid files remain errors. No catalog is written and no allocation starts until `compute launch` resolves resources and submits once. Slurm queries and validated local OS identities are authoritative for allocations; Dask is the diff --git a/docs/cli/compute.md b/docs/cli/compute.md index b082434b..799cddcf 100644 --- a/docs/cli/compute.md +++ b/docs/cli/compute.md @@ -5,45 +5,48 @@ No project is required for these commands. ```text lc compute resources [--json] -lc compute launch --cpus VALUE --memory VALUE [--gpus NAME[:COUNT]|0] - [--name NAME] [--num-nodes N] [--time DURATION] [--startup fast] [--dry-run] [--json] +lc compute launch [--cpus VALUE --memory VALUE] [--gpus NAME[:COUNT]|0] + [--name NAME] [--num-nodes N] [--time DURATION] [--startup fast] [--dry-run] [--wait] [--timeout SECONDS] [--json] lc compute status [CLUSTER] [--wait] [--timeout SECONDS] [--json] lc compute down CLUSTER [--json] ``` -Without configuration, `resources` exposes a built-in `local` offer: one CPU, -1 GiB, one node, fast startup, and a 30-minute default lifetime (two-hour maximum). -Launch it with `lc compute launch --cpus 1 --memory 1`; execution still requires -the returned cluster name or its full immutable ID. -GPU offers require an explicit catalog. On Linux, local GPU launches also require -an externally configured `CUDA_VISIBLE_DEVICES` mask; Lightcone does not discover -GPU hardware. See [GPU allocations](../user/cluster.md#gpu-allocations). - -`~/.lightcone/compute.yaml`, when present, replaces this built-in catalog. -`LC_COMPUTE_CONFIG` selects another file for both compute and execution commands, -so an allocation launched from a catalog can be found by `lc run` and -`lc materialize` too. Missing explicit paths and invalid catalogs are errors. Only a missing implicit default file enables the -built-in catalog, without writing a file or starting any compute. See the -[local and Slurm setup](../user/cluster.md) for examples. +With no resource flags, `lc compute launch` starts the default local CPU offer, +names the cluster `local`, and uses all detected usable logical CPUs and RAM. +The built-in offer has one node, fast startup, a 30-minute default lifetime, and +a two-hour maximum. `--name` and `--time` override the name and lifetime. +CPU and RAM are cooperative scheduling budgets, not exclusive reservations. +GPUs still require explicit offers and, locally, a `CUDA_VISIBLE_DEVICES` mask. + +`~/.lightcone/compute.yaml` configures resource offers; `LC_COMPUTE_CONFIG` selects +another file for all compute and execution commands. The top-level `local` block +can override the built-in CPU/RAM budget or disable local compute. Without an +explicit local connection, the built-in local offer is appended after configured +offers. Catalogs with explicit local connections use their own offers instead; +the shortcut chooses the first eligible local offer. See +[local configuration](../user/cluster.md#customize-resource-offers). +Missing explicit paths and invalid catalogs are errors. Loading a catalog or +planning with `--dry-run` creates no files or compute. | Command | Behavior | |---|---| | `resources` | Ordered available offers, per-node shape, node limit, default/maximum time, and startup class. Free capacity remains unknown. | | `launch` | Resolve one resource request and submit exactly once; print only the cluster name to stdout on acceptance. | +| `launch --wait` | Submit once, then wait for all expected workers. `--timeout` sets the readiness deadline (default 300 seconds). | | `launch --dry-run` | Show the resolved shape and native launch parameters without allocation. | | `status` | List one `name: status` line per current allocation across every configured connection; report the connections that could not be queried. | | `status CLUSTER` | Resolve a name or full ID, inspect native state, and probe Dask readiness separately. | | `status CLUSTER --wait` | Wait for readiness: the allocation is active and every node's worker is connected. The default deadline is 300 seconds, and queries grow less frequent as the wait goes on (up to every 30 seconds). Exits 1 on timeout, or at once if the allocation is ending or has ended; the allocation is left unchanged. | | `down CLUSTER` | Request native termination even if the scheduler is unavailable. An allocation that has already ended is a successful no-op when addressed by full ID; its name no longer resolves. A Slurm job that has left the queue is refused unless accounting confirms it ended. | -Choose a name with `--name analysis`, or omit it to generate `lc-` followed by -12 random hexadecimal characters. Names contain 1โ€“63 lowercase ASCII letters, +The local shortcut defaults to `local`. For explicit CPU/memory requests, omitting +`--name` generates `lc-` followed by 12 random hexadecimal characters. +Use `--name analysis` to choose either name yourself. Names contain 1โ€“63 lowercase ASCII letters, digits, or hyphens; they start with a letter and end with a letter or digit. Launch guidance goes to stderr, so the default output can be captured directly: ```bash -CLUSTER=$(lc compute launch --cpus 1 --memory 1) -lc compute status "$CLUSTER" --wait +CLUSTER=$(lc compute launch --wait) lc compute down "$CLUSTER" ``` @@ -56,6 +59,9 @@ durable reference to that allocation. Use the full immutable `id` from launch or status JSON to address one allocation directly, including when unrelated connections are unavailable. No name registry is maintained. +Supply both `--cpus` and `--memory`, or omit both for the local shortcut. +The shortcut never selects a remote offer, including when local compute is disabled. + Resource quantities are **per node**, and `--num-nodes` defaults to one. CPU and memory requests follow SkyPilot's exact/minimum convention: `4` is exact and `4+` means at least four. Compute memory uses binary units: `16`, `16GB`, and `16GiB` @@ -86,14 +92,22 @@ submission is an error, with no automatic resubmission elsewhere. An uncertain s includes its token and any known cluster ID. Inspect existing allocations before retrying it. -`--wait` requires a CLUSTER, and `--timeout` requires `--wait`. +On `status`, `--wait` requires a CLUSTER. On both commands, `--timeout` requires +`--wait`. Launch rejects `--wait --dry-run`. A waiting launch prints its cluster +name only once ready; JSON adds `ready: true`. A timeout or startup failure exits +1 and includes the accepted immutable ID in the error. Waiting never resubmits +or terminates the accepted allocation. + +Only one local cluster may run per user on a machine, across names, namespaces, +and configured roots. Concurrent launches are serialized by an OS lifetime lock; +a second launch fails until the existing cluster ends. `--json` emits versioned (`schema_version: 1`), allowlisted data without scheduler credentials: | Command | Keys | |---|---| -| `launch` | `plan`, `id`, `name`, `accepted` (only `plan` with `--dry-run`) | +| `launch` | `plan`, `id`, `name`, `accepted` (`ready: true` after `--wait`; only `plan` with `--dry-run`) | | `status CLUSTER` | `id`, `name`, `phase`, `allocation`, `dask`, `reason`, `native_state` | | `status` | `clusters` (a list of the above) and `errors` (by connection name) | | `down` | `id`, `name`, `termination_requested` | diff --git a/docs/user/cluster.md b/docs/user/cluster.md index 54ab30ba..f3bdc6dc 100644 --- a/docs/user/cluster.md +++ b/docs/user/cluster.md @@ -7,26 +7,24 @@ present. `lc materialize --check` and `lc status` remain local project inspectio ## Start locally -No configuration is needed on a fresh installation. When -`~/.lightcone/compute.yaml` is absent, Lightcone exposes a built-in `local` CPU offer: -one logical CPU, 1 GiB, one node, and fast startup. Its default lifetime is -30 minutes, with a maximum of two hours. This creates no catalog file and starts -no processes until you launch a cluster. -Use a custom catalog for larger CPU or RAM budgets and for GPU offers. -See [GPU allocations](#gpu-allocations). +No configuration is needed on a fresh installation. `lc compute launch` uses all +detected usable logical CPUs and RAM on this machine and names the cluster `local`. +The default lifetime is 30 minutes, with a maximum of two hours. Use `--time` to +change the lifetime. GPUs require explicit offers; see [GPU allocations](#gpu-allocations). ```bash lc compute resources -lc compute launch --cpus 1 --memory 1 --dry-run -CLUSTER=$(lc compute launch --cpus 1 --memory 1) -lc compute status "$CLUSTER" --wait +lc compute launch --dry-run +CLUSTER=$(lc compute launch --wait) lc run "$CLUSTER" -- python -c 'print("hello from the cluster")' lc materialize "$CLUSTER" lc compute down "$CLUSTER" ``` Run the execution commands from your project root. A launch returns when native -allocation is accepted; `status --wait` waits for Dask readiness. Execution never +allocation is accepted; `launch --wait` or `status --wait` waits for Dask readiness. +Both accept `--timeout SECONDS` (default 300). Waiting failures retain the accepted +cluster ID and leave the allocation unchanged. Execution never waits: `lc run` and `lc materialize` refuse a cluster that is not active with every expected worker connected, for example: @@ -43,6 +41,11 @@ Use `lc compute status NAME` for resource details and Dask readiness. ## Local allocations Local resources are cooperative limits, not an exclusive CPU/RAM reservation. +Only one local cluster can run per user on each machine. An OS lock held by the +detached owner prevents concurrent launches, including through different names, +catalogs, namespaces, or connection roots. End the existing cluster before +launching another. The lock releases when its owner exits, including on failure +or walltime expiry. An allocation owns a detached process session and standard `LocalCluster`: one worker process with `task_slots_per_node` threads, and a scheduler that listens on `127.0.0.1` over TLS. Its own logs are discarded; a startup failure is kept @@ -57,7 +60,7 @@ launched it. Other boot sessions are excluded from discovery, and an explicit ID from one is refused rather than reported as stopped. Once an allocation has ended, its credentials and scratch directory are removed; its full ID still reports `ended`. -Local compute is available wherever the catalog exposes a valid local offer; +Local compute is available when `local.enabled` is true and the catalog exposes a valid local offer; Lightcone does not infer permission from login-node names or site environment variables. Allocation choices are explicit and native permissions still apply. @@ -69,12 +72,11 @@ the hostname it belongs to. Local offers take no `config`. ## Cluster names -Choose a name at launch, or omit `--name` to generate a short name such as -`lc-a1b2c3d4e5f6`: +The local shortcut defaults to `local`. Explicit CPU/memory requests generate a +short name such as `lc-a1b2c3d4e5f6`. Override either with `--name`: ```bash -lc compute launch --name analysis --cpus 1 --memory 1 -lc compute status analysis --wait +lc compute launch --name analysis --wait lc compute down analysis ``` @@ -97,12 +99,45 @@ separate name registry. Native state still decides whether an allocation exists. ## Customize resource offers -Create `~/.lightcone/compute.yaml` to expose other resource shapes or services. -A configured catalog replaces the built-in catalog completely; no extra local -offer is added to it. The namespace is a stable UUID identifying a connection; -keep it unchanged while that connection's clusters exist. +Create `~/.lightcone/compute.yaml`, or select a file with `LC_COMPUTE_CONFIG`. +For a smaller default local budget: -For example, this catalog exposes a larger local allocation: +```yaml +version: 1 +local: + resources: {cpus: 4, memory: 8GiB} +``` + +Both CPU and memory are required in `local.resources`; omit that block to use +detected capacity. The same capacity validation applies to configured budgets. +This controls the default offer, not hard OS resource limits. + +On a login node, disable local compute and configure Slurm offers: + +```yaml +version: 1 +local: + enabled: false +# Add Slurm connections and offers as shown below. +``` + +This blocks local launches, including explicit local offers, and new execution +commands on local clusters. Existing local allocations remain inspectable and +stoppable; disabling does not kill them. Bare launch reports that local compute is +disabled; supply CPU/memory requirements to select a Slurm allocation. Select this +catalog on login nodes through `LC_COMPUTE_CONFIG`. This is Lightcone configuration +policy; native site permissions enforce machine-wide restrictions. + +By default, a built-in `local` connection and offer accompany remote offers, with +configured offers taking selection priority. If the catalog already defines local +connections, those offers replace the implicit local offer; omit `local.resources` +and size those offers directly. The no-resource shortcut selects the first eligible +local offer and defaults its cluster name to `local`. + +The namespace is a stable UUID identifying a connection; keep it unchanged while +that connection's clusters exist. + +For example, this catalog explicitly defines a local allocation: ```yaml version: 1 @@ -121,16 +156,16 @@ offers: Set `LC_COMPUTE_CONFIG` to choose another file for all commands, including `lc run` and `lc materialize`, which find clusters through the same catalog. A -missing explicit path or an invalid catalog is an error; only an absent implicit default file -selects the built-in offer. Stop existing built-in allocations before replacing -their connection with your own catalog. +missing explicit path or an invalid catalog is an error. An absent implicit default +file uses the default local policy. Stop existing built-in allocations before replacing +their connection namespace with your own. This example keeps the built-in connection's namespace, so allocations launched from the built-in offer stay visible and can still be stopped after the file exists. -A catalog has `version: 1`, a `connections` mapping, and an ordered `offers` -list: +A catalog has `version: 1`, an optional `local` policy, a `connections` mapping, +and an ordered `offers` list. Connections and offers default to empty: - A connection has a `namespace` (a UUID), a `provider` (`local` or `slurm`), an optional `context`, and optional provider `launch` settings. Namespaces must @@ -178,6 +213,8 @@ resource sizing. It has not been validated by submitting a job at NERSC: ```yaml version: 1 +local: + enabled: false connections: perlmutter: namespace: 9d0c0fc5-9be8-407a-a3ec-f17c4110b162 diff --git a/docs/user/getting-started.md b/docs/user/getting-started.md index 80d71b9f..9843d4c4 100644 --- a/docs/user/getting-started.md +++ b/docs/user/getting-started.md @@ -250,13 +250,12 @@ the record of what your results were computed with. ## 5. Materialize Launch the built-in local offer; no compute configuration is needed. It provides -one CPU and 1 GiB for 30 minutes. Keep the returned ID in `CLUSTER` for this +all usable CPUs and RAM for 30 minutes. Keep the returned ID in `CLUSTER` for this walkthrough. If you already have a compute catalog, its offers replace that default; see [Running on a Cluster](cluster.md). ```bash -CLUSTER=$(lc compute launch --cpus 1 --memory 1 --json | python -c 'import json,sys; print(json.load(sys.stdin)["id"])') -lc compute status "$CLUSTER" --wait +CLUSTER=$(lc compute launch --wait --json | python -c 'import json,sys; print(json.load(sys.stdin)["id"])') ``` Execution always requires this cluster ID. `lc materialize --check` can inspect diff --git a/src/lightcone/cli/compute.py b/src/lightcone/cli/compute.py index aa139375..1f745524 100644 --- a/src/lightcone/cli/compute.py +++ b/src/lightcone/cli/compute.py @@ -3,6 +3,7 @@ from __future__ import annotations import json +import math from collections.abc import Iterator from contextlib import contextmanager from typing import Any @@ -96,9 +97,9 @@ def resources(as_json: bool) -> None: @compute.command() -@click.option("--name", help="Cluster name; defaults to a generated short name.") -@click.option("--cpus", required=True, help="Logical CPUs per node; suffix + requests a minimum.") -@click.option("--memory", required=True, +@click.option("--name", help="Cluster name; defaults to local for the shortcut, else a short name.") +@click.option("--cpus", help="Logical CPUs per node; suffix + requests a minimum.") +@click.option("--memory", help="Memory per node, e.g. 16 or 16GB; suffix + requests a minimum.") @click.option("--gpus", default="0", show_default=True, help="Accelerator NAME[:COUNT] per node, e.g. A100:4 or GPU:1; 0 requests CPU only.") @@ -110,35 +111,48 @@ def resources(as_json: bool) -> None: "--startup", type=click.Choice(["fast"]), help="Require a fast startup service class." ) @click.option("--dry-run", is_flag=True, help="Resolve the launch without allocating compute.") +@click.option("--wait", is_flag=True, help="Wait until the new cluster is ready for execution.") +@click.option("--timeout", type=click.FloatRange(min=0, min_open=True), + help="Readiness deadline in seconds (default: 300); requires --wait.") @click.option("--json", "as_json", is_flag=True, help="Emit structured output.") def launch( name: str | None, - cpus: str, - memory: str, + cpus: str | None, + memory: str | None, gpus: str, num_nodes: int, walltime: str | None, startup: str | None, dry_run: bool, + wait: bool, + timeout: float | None, as_json: bool, ) -> None: - """Create one new cluster from a provider-independent resource request.""" + """Create a cluster; omit CPU and memory to use the default local offer.""" from lightcone.engine.compute import Compute - from lightcone.engine.compute.model import Request + from lightcone.engine.compute.model import ComputeError, Request with _errors(as_json): + if timeout is not None and not wait: + raise ComputeError("--timeout requires --wait") + if timeout is not None and not math.isfinite(timeout): + raise ComputeError("timeout must be finite and positive") + if dry_run and wait: + raise ComputeError("--wait cannot be combined with --dry-run") + if (cpus is None) != (memory is None): + raise ComputeError("supply both --cpus and --memory, or omit both for local compute") service = Compute() - plan = service.plan( - Request.parse( - cpus, - memory, - gpus=gpus, - num_nodes=num_nodes, - time=walltime, - startup=startup, - ), - name=name, - ) + if cpus is None or memory is None: + plan = service.plan_local( + name=name, time=walltime, gpus=gpus, num_nodes=num_nodes, startup=startup, + ) + else: + plan = service.plan( + Request.parse( + cpus, memory, gpus=gpus, num_nodes=num_nodes, time=walltime, startup=startup, + ), + name=name, + ) data: dict[str, Any] = {"schema_version": 1, "plan": plan.as_dict()} if dry_run: if as_json: @@ -148,11 +162,26 @@ def launch( return identity = service.launch(plan) data.update(id=identity.encode(), name=identity.name, accepted=True) + if wait: + if not as_json: + click.echo( + f"Allocation accepted: {identity.name}. Waiting for readiness.", err=True, + ) + try: + snapshot = service.status(identity.encode(), wait=True, timeout=timeout or 300) + if not snapshot.ready: + raise ComputeError( + f"cluster {identity.name} is {snapshot.phase}: {snapshot.reason}", + ) + except ComputeError as exc: + raise ComputeError(str(exc), cluster_id=identity.encode()) from exc + data["ready"] = True if as_json: click.echo(json.dumps(data)) else: click.echo(identity.name) click.echo( + f"Cluster {identity.name} is ready." if wait else f"Allocation accepted. Use lc compute status {identity.name} --wait for readiness.", err=True, ) diff --git a/src/lightcone/engine/compute/__init__.py b/src/lightcone/engine/compute/__init__.py index 2139ca9b..f299cfeb 100644 --- a/src/lightcone/engine/compute/__init__.py +++ b/src/lightcone/engine/compute/__init__.py @@ -15,6 +15,7 @@ Connection, Identity, LaunchPlan, + Offer, Provider, ProviderFactory, Request, @@ -127,34 +128,10 @@ def plan(self, request: Request, *, name: str | None = None) -> LaunchPlan: validate_name(name) unavailable: list[str] = [] for offer in self.catalog.offers: - if request.num_nodes > offer.max_nodes: - continue - if request.startup is not None and request.startup != offer.startup.class_: - continue - if request.seconds is not None and request.seconds > offer.time.max_seconds: - continue - if ( - offer.resources.cpus < request.cpus - if request.min_cpus - else offer.resources.cpus != request.cpus - ): - continue - if ( - offer.resources.memory_bytes < request.memory_bytes - if request.min_memory - else offer.resources.memory_bytes != request.memory_bytes - ): - continue - if request.accelerators is None: - matches_accelerators = offer.resources.accelerators is None - else: - matches_accelerators = request.accelerators.matches(offer.resources.accelerators) - if not matches_accelerators: - continue try: - provider = self.provider(self.catalog.connections[offer.connection]) - plan = provider.plan(offer, request) - return plan.replace(name=name) + plan = self._plan_offer(offer, request, name) + if plan is not None: + return plan except UnavailableOfferError as exc: unavailable.append(f"{offer.name}: {exc}") raise ComputeError( @@ -162,8 +139,70 @@ def plan(self, request: Request, *, name: str | None = None) -> LaunchPlan: + ("; " + "; ".join(unavailable) if unavailable else "") ) + def _plan_offer( + self, offer: Offer, request: Request, name: str | None, + ) -> LaunchPlan | None: + """Match one shape and validate its provider without allocating anything.""" + connection = self.catalog.connections[offer.connection] + if connection.provider == "local" and not self.catalog.local.enabled: + return None + if request.num_nodes > offer.max_nodes: + return None + if request.startup is not None and request.startup != offer.startup.class_: + return None + if request.seconds is not None and request.seconds > offer.time.max_seconds: + return None + if ( + offer.resources.cpus < request.cpus + if request.min_cpus else offer.resources.cpus != request.cpus + ): + return None + if ( + offer.resources.memory_bytes < request.memory_bytes + if request.min_memory else offer.resources.memory_bytes != request.memory_bytes + ): + return None + if request.accelerators is None: + matches = offer.resources.accelerators is None + else: + matches = request.accelerators.matches(offer.resources.accelerators) + if not matches: + return None + return self.provider(connection).plan(offer, request).replace(name=name) + + def plan_local( + self, *, name: str | None = None, time: str | None = None, + gpus: str = "0", num_nodes: int = 1, startup: str | None = None, + ) -> LaunchPlan: + """Plan the first usable local offer, without considering remote backends.""" + if not self.catalog.local.enabled: + raise ComputeError("local compute is disabled by the compute configuration") + name = "local" if name is None else name + validate_name(name) + unavailable: list[str] = [] + for offer in self.catalog.offers: + connection = self.catalog.connections[offer.connection] + if connection.provider != "local": + continue + request = Request.parse( + str(offer.resources.cpus), f"{offer.resources.memory_bytes}B", + gpus=gpus, num_nodes=num_nodes, time=time, startup=startup, + ) + try: + plan = self._plan_offer(offer, request, name) + if plan is not None: + return plan + except UnavailableOfferError as exc: + unavailable.append(f"{offer.name}: {exc}") + raise ComputeError( + "no local offer matches; configure local resources or supply --cpus and --memory " + "for a remote allocation" + ("; " + "; ".join(unavailable) if unavailable else "") + ) + def launch(self, plan: LaunchPlan) -> Identity: """Choose an unused name from native observations, then submit exactly once.""" + if plan.connection.provider == "local" and not self.catalog.local.enabled: + raise ComputeError("local compute is disabled by the compute configuration") if plan.name is not None: validate_name(plan.name) snapshots, errors = self.discover() @@ -252,7 +291,13 @@ def connect(cluster_id: str, *, timeout: float = 10) -> Iterator[Any]: """Borrow a validated standard Dask client; detach without closing its cluster.""" if not math.isfinite(timeout) or timeout <= 0: raise ComputeError("timeout must be finite and positive") - provider, identity = Compute().resolve(cluster_id) + service = Compute() + provider, identity = service.resolve(cluster_id) + connection = service.catalog.connection_for(identity.namespace) + if connection.provider == "local" and not service.catalog.local.enabled: + raise ComputeError( + "local compute is disabled by the compute configuration", cluster_id=identity.encode(), + ) snapshot = provider.inspect(identity) if snapshot.phase != "active": raise ComputeError( diff --git a/src/lightcone/engine/compute/catalog.py b/src/lightcone/engine/compute/catalog.py index a9b92aa8..c9bc964a 100644 --- a/src/lightcone/engine/compute/catalog.py +++ b/src/lightcone/engine/compute/catalog.py @@ -3,7 +3,6 @@ from __future__ import annotations import os -from decimal import Decimal from pathlib import Path from typing import Annotated, Any, Self @@ -43,12 +42,28 @@ def _unique_mapping(loader: yaml.SafeLoader, node: yaml.MappingNode) -> dict[str _UniqueLoader.add_constructor(yaml.resolver.BaseResolver.DEFAULT_MAPPING_TAG, _unique_mapping) +class LocalSettings(ComputeModel): + """Policy for local launches, including the implicit workstation offer.""" + + enabled: bool = True + resources: Resources | None = None + + @model_validator(mode="after") + def cpu_only(self) -> Self: + if self.resources is not None and self.resources.gpus: + raise ValueError( + "local.resources supports CPUs and memory; configure GPU offers explicitly" + ) + return self + + class Catalog(ComputeModel): """Configuration for new requests, never a registry of live clusters.""" version: Annotated[int, Field(ge=1, le=1)] - connections: dict[Name, Connection] - offers: list[Offer] + connections: dict[Name, Connection] = Field(default_factory=dict) + offers: list[Offer] = Field(default_factory=list) + local: LocalSettings = Field(default_factory=LocalSettings) @model_validator(mode="after") def relationships(self) -> Self: @@ -73,7 +88,7 @@ def relationships(self) -> Self: @classmethod def load(cls, path: Path | None = None) -> Catalog: - """Load configured offers, or a small local offer if the default file is absent. + """Load configured offers and apply the local-compute policy. Explicit paths and existing catalogs must be readable and valid. Loading the built-in offer writes no catalog and allocates no compute. @@ -88,25 +103,52 @@ def load(cls, path: Path | None = None) -> Catalog: except FileNotFoundError as exc: if configured or path.is_symlink(): raise ComputeError(f"cannot read compute catalog {path}: {exc}") from exc - return cls( - version=1, - connections={"local": Connection(namespace=_LOCAL_NAMESPACE, provider="local")}, - offers=[Offer( - name="local", connection="local", - resources=Resources(cpus=1, memory_gib=Decimal(1)), - max_nodes=1, time=TimeLimits(default="30m", max="2h"), - startup=Startup(class_="fast"), - )], - ) + raw = {"version": 1} except (OSError, UnicodeError, yaml.YAMLError) as exc: raise ComputeError(f"cannot read compute catalog {path}: {exc}") from exc try: - return cls.model_validate(raw) + catalog = cls.model_validate(raw) + return catalog._with_local() except ValidationError as exc: raise ComputeError( f"invalid compute catalog {path}:\n{validation_message(exc)}" ) from exc + def _with_local(self) -> Catalog: + """Keep explicit local connections, or add the stable built-in connection.""" + connections = dict(self.connections) + offers = list(self.offers) + explicit = any(connection.provider == "local" for connection in connections.values()) + if explicit and self.local.resources is not None: + raise ComputeError( + "local.resources cannot be combined with explicit local connections; " + "set their offer resources instead" + ) + if not explicit: + if "local" in connections: + raise ComputeError( + "connection name 'local' is reserved for the built-in local backend" + ) + # Retain this authority even when disabled so existing allocations can be stopped. + connections["local"] = Connection(namespace=_LOCAL_NAMESPACE, provider="local") + if self.local.enabled: + from dask.system import CPU_COUNT + from distributed.system import MEMORY_LIMIT + + resources = self.local.resources or Resources.from_bytes( + cpus=CPU_COUNT, memory_bytes=MEMORY_LIMIT, + ) + offers.append(Offer( + name="local", connection="local", resources=resources, + max_nodes=1, time=TimeLimits(default="30m", max="2h"), + startup=Startup(class_="fast"), + )) + if not self.local.enabled: + offers = [ + offer for offer in offers if connections[offer.connection].provider != "local" + ] + return self.replace(connections=connections, offers=offers) + def connection_for(self, namespace: str) -> Connection: """Find the configured authority without relying on current offers.""" for connection in self.connections.values(): diff --git a/src/lightcone/engine/compute/local.py b/src/lightcone/engine/compute/local.py index 85b05272..56c44380 100644 --- a/src/lightcone/engine/compute/local.py +++ b/src/lightcone/engine/compute/local.py @@ -7,6 +7,7 @@ import shutil import signal import socket +import stat import subprocess import sys import tempfile @@ -47,6 +48,39 @@ _RETIRED = "ended.json" +@contextmanager +def _allocation_lock() -> Iterator[int]: + """Hold a per-user, host-local lock across catalogs and detached owner lifetime.""" + import fcntl + + # Deliberately independent of TMPDIR, connection namespaces and configured roots. + # Resolve /tmp for macOS, where it is an alias of /private/tmp. + root = private_directory(Path("/tmp").resolve() / f"lightcone-local-{os.getuid()}", create=True) + descriptor: int | None = None + try: + descriptor = os.open( + root / "allocation.lock", os.O_CREAT | os.O_RDWR | os.O_NOFOLLOW, 0o600, + ) + info = os.fstat(descriptor) + if (not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid() + or stat.S_IMODE(info.st_mode) != 0o600 or info.st_nlink != 1): + raise ComputeError("local allocation lock must be a private file owned by you") + try: + fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError as exc: + raise ComputeError( + "a local cluster is already running or starting for this user on this machine; " + "stop it with lc compute down before launching another" + ) from exc + yield descriptor + except OSError as exc: + raise ComputeError(f"cannot acquire the local allocation lock: {exc}") from exc + finally: + if descriptor is not None: + # Do not unlock: the detached owner inherits this same open file description. + os.close(descriptor) + + def _boot_identity() -> str: try: if sys.platform == "linux": @@ -145,6 +179,16 @@ def launch(self, plan: LaunchPlan) -> Identity: raise ComputeError("local launch plan belongs to a different connection or node count") if plan.name is not None: validate_name(plan.name) + with _allocation_lock() as lock_fd: + # Also recognize allocations launched before the lifetime lock existed. + if self.discover(): + raise ComputeError( + "a local cluster is already running; stop it before launching another" + ) + return self._launch_locked(plan, lock_fd) + + def _launch_locked(self, plan: LaunchPlan, lock_fd: int) -> Identity: + """Transfer the singleton lock to the allocation owner before releasing our copy.""" boot = _boot_identity() token = uuid4().hex directory = private_directory(self.root / token, create=True) @@ -157,6 +201,7 @@ def launch(self, plan: LaunchPlan) -> Identity: "task_slots": plan.details["task_slots_per_node"], "scratch": str(scratch), "identity": "", + "lock_fd": lock_fd, }, ) process: subprocess.Popen[bytes] | None = None @@ -176,6 +221,7 @@ def launch(self, plan: LaunchPlan) -> Identity: stderr=subprocess.DEVNULL, start_new_session=True, close_fds=True, + pass_fds=(lock_fd,), env=environment, ) identity = Identity( @@ -203,6 +249,7 @@ def launch(self, plan: LaunchPlan) -> Identity: "task_slots": plan.details["task_slots_per_node"], "scratch": str(scratch), "identity": identity.encode(), + "lock_fd": lock_fd, }, ) return identity diff --git a/src/lightcone/engine/compute/local_runtime.py b/src/lightcone/engine/compute/local_runtime.py index 1ca6d552..b6e3785e 100644 --- a/src/lightcone/engine/compute/local_runtime.py +++ b/src/lightcone/engine/compute/local_runtime.py @@ -25,6 +25,9 @@ def main() -> None: os.umask(0o077) directory = private_directory(Path(sys.argv[1])) launch = read_private_json(directory / "launch.json") + if "lock_fd" in launch: + # Keep the inherited lock until process exit, but exclude worker execs. + os.set_inheritable(int(launch["lock_fd"]), False) if os.getsid(0) != os.getpid() or os.getpgrp() != os.getpid(): raise RuntimeError("the local allocation owner must lead its own process session") stopped = threading.Event() diff --git a/tests/test_compute.py b/tests/test_compute.py index 29dfead0..31be98ce 100644 --- a/tests/test_compute.py +++ b/tests/test_compute.py @@ -57,12 +57,13 @@ def expand(path: Path) -> Path: @pytest.fixture -def catalog(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Path: +def catalog(default_home: Path, tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Path: path = tmp_path / "compute.yaml" path.write_text( yaml.safe_dump( { "version": 1, + "local": {"enabled": False}, "connections": {"test": {"namespace": NAMESPACE, "provider": "fake"}}, "offers": [ { @@ -251,6 +252,8 @@ def test_incomplete_discovery_cannot_establish_names_but_full_ids_still_work( def test_missing_default_catalog_exposes_stable_local_resources_without_writing_files( default_home: Path, monkeypatch: pytest.MonkeyPatch, ) -> None: + monkeypatch.setattr("dask.system.CPU_COUNT", 1) + monkeypatch.setattr("distributed.system.MEMORY_LIMIT", GIB) first, second = Catalog.load(), Catalog.load() assert first == second assert set(first.connections) == {"local"} @@ -291,7 +294,7 @@ def test_missing_default_catalog_exposes_stable_local_resources_without_writing_ assert list(default_home.iterdir()) == [] -def test_configured_catalogs_replace_the_builtin_and_obey_path_precedence( +def test_configured_catalogs_can_disable_local_and_obey_path_precedence( catalog: Path, default_home: Path, monkeypatch: pytest.MonkeyPatch, ) -> None: monkeypatch.delenv("LC_COMPUTE_CONFIG", raising=False) @@ -299,11 +302,10 @@ def test_configured_catalogs_replace_the_builtin_and_obey_path_precedence( default.parent.mkdir() default.write_text(catalog.read_text()) configured = Catalog.load() - assert set(configured.connections) == {"test"} + assert set(configured.connections) == {"test", "local"} assert [offer.name for offer in configured.offers] == ["quick", "large"] - # An empty configured catalog explicitly exposes nothing; the builtin is - # never merged into it, whether selected by default, environment, or option. - default.write_text("version: 1\nconnections: {}\noffers: []\n") + # Disabled local connections remain available for inspection and termination. + default.write_text("version: 1\nlocal: {enabled: false}\nconnections: {}\noffers: []\n") assert Catalog.load().offers == [] monkeypatch.setenv("LC_COMPUTE_CONFIG", str(catalog)) assert Catalog.load() == configured @@ -475,6 +477,127 @@ def test_builtin_catalog_stays_cpu_only_without_probing_native_gpus( assert not list(default_home.iterdir()) +def test_local_shortcut_uses_detected_capacity_without_writing_files( + default_home: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setattr("dask.system.CPU_COUNT", 6) + monkeypatch.setattr("distributed.system.MEMORY_LIMIT", 7 * GIB + 123) + result = CliRunner().invoke(main, ["compute", "launch", "--dry-run", "--json"]) + assert result.exit_code == 0, result.output + plan = compute.Compute().plan_local() + assert plan.name == "local" + assert plan.resources.cpus == 6 + assert plan.resources.memory_bytes == 7 * GIB + 123 + assert plan.details["task_slots_per_node"] == 6 + assert json.loads(result.stdout)["plan"] == plan.as_dict() + assert not list(default_home.iterdir()) + + +def test_local_config_overrides_default_and_survives_disabling( + default_home: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setattr("dask.system.CPU_COUNT", 8) + monkeypatch.setattr("distributed.system.MEMORY_LIMIT", 16 * GIB) + path = default_home / "compute.yaml" + path.write_text("version: 1\nlocal:\n resources: {cpus: 2, memory: 3}\n") + monkeypatch.setenv("LC_COMPUTE_CONFIG", str(path)) + service = compute.Compute() + plan = service.plan_local(name="sandbox", time="1h") + assert (plan.name, plan.resources.cpus, plan.resources.memory_bytes) == ("sandbox", 2, 3 * GIB) + assert plan.seconds == 3600 + connection = plan.connection + path.write_text("version: 1\nlocal: {enabled: false}\n") + disabled = compute.Compute() + assert not disabled.resources()["offers"] + assert disabled.catalog.connection_for(connection.namespace) == connection + with pytest.raises(ComputeError, match="disabled"): + disabled.plan_local() + with pytest.raises(ComputeError, match="disabled"): + disabled.launch(plan) + with pytest.raises(ComputeError, match="no configured offer"): + disabled.plan(Request.parse("2", "3")) + + +def test_catalog_adds_local_after_remote_offers_and_shortcut_never_selects_remote( + catalog: Path, provider: MagicMock, monkeypatch: pytest.MonkeyPatch, +) -> None: + data = yaml.safe_load(catalog.read_text()) + data.pop("local") + catalog.write_text(yaml.safe_dump(data)) + monkeypatch.setattr("dask.system.CPU_COUNT", 4) + monkeypatch.setattr("distributed.system.MEMORY_LIMIT", 8 * GIB) + service = compute.Compute() + assert [offer.name for offer in service.catalog.offers] == ["quick", "large", "local"] + assert service.plan(Request.parse("4", "8")).offer.name == "quick" + assert service.plan_local().connection.provider == "local" + with pytest.raises(ComputeError, match="no local offer"): + service.plan_local(num_nodes=2) + provider.launch.assert_not_called() + + +def test_explicit_local_offers_keep_their_sizes_and_replace_the_implicit_offer( + catalog: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + data = yaml.safe_load(catalog.read_text()) + data["local"] = {"enabled": True} + data["connections"]["test"]["provider"] = "local" + catalog.write_text(yaml.safe_dump(data)) + monkeypatch.setattr("dask.system.CPU_COUNT", 8) + monkeypatch.setattr("distributed.system.MEMORY_LIMIT", 16 * GIB) + service = compute.Compute() + assert set(service.catalog.connections) == {"test"} + plan = service.plan_local() + assert (plan.offer.name, plan.name, plan.resources.cpus) == ("quick", "local", 4) + assert plan.resources.memory_bytes == 8 * GIB + data["local"]["resources"] = {"cpus": 2, "memory": 2} + catalog.write_text(yaml.safe_dump(data)) + with pytest.raises(ComputeError, match="cannot be combined with explicit local"): + Catalog.load() + + +def test_configured_local_budget_still_must_fit_host_capacity( + catalog: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + data = yaml.safe_load(catalog.read_text()) + data["local"] = {"resources": {"cpus": 8, "memory": 16}} + catalog.write_text(yaml.safe_dump(data)) + monkeypatch.setattr("dask.system.CPU_COUNT", 4) + monkeypatch.setattr("distributed.system.MEMORY_LIMIT", 8 * GIB) + result = CliRunner().invoke(main, ["compute", "launch", "--dry-run", "--json"]) + assert result.exit_code == 1 + assert "exceeds this host's CPU or RAM capacity" in json.loads(result.stdout)["error"] + + +@pytest.mark.parametrize("settings, message", [ + ({"enabled": "false"}, "local.enabled"), + ({"resources": {"cpus": 0, "memory": 1}}, "local.resources.cpus"), + ({"resources": {"cpus": 1, "memory": 1, "accelerators": "GPU:1"}}, "GPU offers"), +]) +def test_local_policy_validation(catalog: Path, settings: object, message: str) -> None: + data = yaml.safe_load(catalog.read_text()) + data["local"] = settings + catalog.write_text(yaml.safe_dump(data)) + with pytest.raises(ComputeError, match=message): + Catalog.load() + + +def test_disabled_policy_blocks_explicit_local_offers_and_execution_but_allows_down( + catalog: Path, provider: MagicMock, monkeypatch: pytest.MonkeyPatch, +) -> None: + data = yaml.safe_load(catalog.read_text()) + data["connections"]["test"]["provider"] = "local" + catalog.write_text(yaml.safe_dump(data)) + monkeypatch.setitem(compute.PROVIDERS, "local", lambda connection: provider) + service = compute.Compute() + assert service.catalog.offers == [] + with pytest.raises(ComputeError, match="disabled"): + with compute.connect(IDENTITY.encode()): + pytest.fail("borrowed disabled local compute") + provider.connect.assert_not_called() + service.down(IDENTITY.encode()) + provider.terminate.assert_called_once_with(IDENTITY) + + def test_configured_catalogs_do_not_probe_local_gpus( catalog: Path, monkeypatch: pytest.MonkeyPatch, ) -> None: @@ -851,6 +974,7 @@ def test_invalid_catalog_encoding_is_a_structured_error(catalog: Path) -> None: @pytest.mark.parametrize("setting", ["connection_root", "scratch_root", "python"]) def test_local_catalog_path_types_fail_without_a_traceback(catalog: Path, setting: str) -> None: data = yaml.safe_load(catalog.read_text()) + data["local"] = {"enabled": True} data["connections"]["test"]["provider"] = "local" data["connections"]["test"]["launch"] = {setting: None} catalog.write_text(yaml.safe_dump(data)) @@ -939,6 +1063,61 @@ def test_cli_launch_name_output_can_be_captured_without_json( provider.terminate.assert_called_once_with(identity) +@pytest.mark.parametrize("as_json", [False, True]) +def test_cli_launch_wait_returns_only_after_ready( + catalog: Path, provider: MagicMock, as_json: bool, +) -> None: + result = CliRunner().invoke(main, [ + "compute", "launch", "--cpus", "4", "--memory", "8", "--wait", + *(["--json"] if as_json else []), + ]) + assert result.exit_code == 0, result.output + provider.launch.assert_called_once() + provider.inspect.assert_called_once_with(IDENTITY) + provider.connect.assert_called_once() + if as_json: + data = json.loads(result.stdout) + assert data["accepted"] and data["ready"] + assert data["id"] == IDENTITY.encode() + else: + assert result.stdout == IDENTITY.name + "\n" + assert "is ready" in result.stderr + + +@pytest.mark.parametrize("phase", ["pending", "stopping", "ended"]) +def test_cli_launch_wait_failure_keeps_accepted_id_and_never_resubmits( + catalog: Path, provider: MagicMock, phase: str, +) -> None: + provider.inspect.return_value.phase = phase + result = CliRunner().invoke(main, [ + "compute", "launch", "--cpus", "4", "--memory", "8", + "--wait", "--timeout", "0.01", "--json", + ]) + assert result.exit_code == 1, result.output + data = json.loads(result.stdout) + assert data["id"] == IDENTITY.encode() + assert "error" in data + provider.launch.assert_called_once() + provider.terminate.assert_not_called() + + +@pytest.mark.parametrize("flags, message", [ + (["--timeout", "1"], "requires --wait"), + (["--wait", "--timeout", "inf"], "finite"), + (["--wait", "--timeout", "nan"], "finite"), + (["--wait", "--dry-run"], "cannot be combined"), + (["--cpus", "1"], "both --cpus and --memory"), + (["--memory", "1"], "both --cpus and --memory"), +]) +def test_cli_launch_rejects_invalid_flags_before_submission( + catalog: Path, provider: MagicMock, flags: list[str], message: str, +) -> None: + result = CliRunner().invoke(main, ["compute", "launch", *flags, "--json"]) + assert result.exit_code == 1, result.output + assert message in json.loads(result.stdout)["error"] + provider.launch.assert_not_called() + + def test_cli_partial_failure_and_ambiguous_submit(catalog: Path, provider: MagicMock) -> None: provider.discover.side_effect = ComputeError("unavailable") runner = CliRunner() diff --git a/tests/test_compute_local.py b/tests/test_compute_local.py index c50ecf17..70fcdbdb 100644 --- a/tests/test_compute_local.py +++ b/tests/test_compute_local.py @@ -20,7 +20,9 @@ import psutil import pytest +from click.testing import CliRunner +from lightcone.cli.commands import main from lightcone.engine.compute import Compute, local, local_runtime from lightcone.engine.compute.catalog import Catalog from lightcone.engine.compute.local import LocalProvider @@ -66,6 +68,68 @@ def _launch(provider: LocalProvider, *, seconds: int = 60) -> Identity: ) +def test_singleton_survives_cli_exit_and_spans_names_and_connection_roots( + provider: LocalProvider, tmp_path: Path, +) -> None: + # Independent launchers compete for the same host/user lock, with distinct catalogs. + script = """ +import json, sys +from lightcone.engine.compute.local import LocalProvider +from lightcone.engine.compute.model import ( + Connection, ComputeError, Offer, Request, Resources, TimeLimits, +) +p = LocalProvider(Connection.model_validate_json(sys.argv[1])) +offer = Offer(name='small', connection='workstation', resources=Resources(cpus=1, memory_gib=0.5), + max_nodes=1, time=TimeLimits(default='1m', max='1m')) +plan = p.plan(offer, Request(cpus=1, memory_bytes=512 * 1024**2)).replace(name=sys.argv[2]) +input() +try: + print(json.dumps({'id': p.launch(plan).encode()})) +except ComputeError as exc: + print(json.dumps({'error': str(exc)})) +""" + other = LocalProvider(provider.connection.replace( + namespace=str(uuid4()), launch={"connection_root": str(tmp_path / "other")}, + )) + providers = [provider, other] + processes = [subprocess.Popen( + [sys.executable, "-c", script, item.connection.model_dump_json(), f"local-{index}"], + stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, + ) for index, item in enumerate(providers)] + identities: list[tuple[LocalProvider, Identity]] = [] + try: + for process in processes: + assert process.stdin is not None + process.stdin.write("go\n") + process.stdin.flush() + results = [] + for item, process in zip(providers, processes): + stdout, stderr = process.communicate(timeout=20) + assert process.returncode == 0, stderr + result = json.loads(stdout) + results.append(result) + if "id" in result: + identities.append((item, Identity.decode(result["id"]))) + assert len(identities) == 1, results + assert "already running or starting" in next(r["error"] for r in results if "error" in r) + owner, identity = identities[0] + _ready(owner, identity) + with pytest.raises(ComputeError, match="already running"): + _launch(other) + owner.terminate(identity) + _ended(owner, identity) + replacement = _launch(other) + identities.append((other, replacement)) + _ready(other, replacement) + finally: + for process in processes: + if process.poll() is None: + process.kill() + process.wait(timeout=5) + for item, identity in identities: + item.terminate(identity) + + def _ready(provider: LocalProvider, identity: Identity) -> dict[str, object]: deadline = time.monotonic() + 20 while time.monotonic() < deadline: @@ -195,9 +259,16 @@ def expand(path: Path) -> Path: monkeypatch.setattr(Path, "expanduser", expand) monkeypatch.delenv("LC_COMPUTE_CONFIG", raising=False) - service = Compute() - identity = service.launch(service.plan(Request(cpus=1, memory_bytes=GIB, seconds=60))) + monkeypatch.setattr("dask.system.CPU_COUNT", 1) + monkeypatch.setattr("distributed.system.MEMORY_LIMIT", GIB) + result = CliRunner().invoke(main, [ + "compute", "launch", "--wait", "--time", "1m", "--timeout", "20", "--json", + ]) + data = json.loads(result.stdout) + identity = Identity.decode(data["id"]) try: + assert result.exit_code == 0, result.output + assert data["ready"] and data["name"] == "local" script = """ import sys from pathlib import Path @@ -500,22 +571,20 @@ def wait(*, timeout: float) -> None: def test_failed_spawn_and_unpublished_launch_do_not_hide_healthy_allocations( provider: LocalProvider, monkeypatch: pytest.MonkeyPatch, ) -> None: + popen = subprocess.Popen + with monkeypatch.context() as patch: + def fail(argv: list[str], **kwargs: Any) -> subprocess.Popen[bytes]: + if "lightcone.engine.compute.local_runtime" in argv: + raise OSError("configured interpreter cannot execute") + return popen(argv, **kwargs) + + patch.setattr("lightcone.engine.compute.local.subprocess.Popen", fail) + with pytest.raises(ComputeError, match="cannot execute"): + _launch(provider) + # A failed spawn must release the singleton lock as well as its private files. identity = _launch(provider) try: _ready(provider, identity) - before = set(provider.root.iterdir()) - popen = subprocess.Popen - with monkeypatch.context() as patch: - def fail(argv: list[str], **kwargs: Any) -> subprocess.Popen[bytes]: - if "lightcone.engine.compute.local_runtime" in argv: - raise OSError("configured interpreter cannot execute") - # macOS also uses Popen for its native boot-identity query. - return popen(argv, **kwargs) - - patch.setattr("lightcone.engine.compute.local.subprocess.Popen", fail) - with pytest.raises(ComputeError, match="cannot execute"): - _launch(provider) - assert set(provider.root.iterdir()) == before # A launcher interrupted before identity publication can also leave a # directory. This is not a published allocation or a discovery error. interrupted = private_directory(provider.root / uuid4().hex, create=True) From 016e522f6f89ec3822d73e372e96e8d56809f45e Mon Sep 17 00:00:00 2001 From: Francois Lanusse Date: Tue, 29 Sep 2026 08:29:43 -0700 Subject: [PATCH 2/4] Update eval prompt for default local compute and launch waiting --- evals/prompt.md | 36 ++++++++++++++++++++++++++---------- 1 file changed, 26 insertions(+), 10 deletions(-) diff --git a/evals/prompt.md b/evals/prompt.md index 6b0c9e3c..e9c4b7da 100644 --- a/evals/prompt.md +++ b/evals/prompt.md @@ -38,16 +38,32 @@ This project is driven by two CLIs โ€” use them rather than improvising: - When a recipe fails, `lc materialize` reports which output failed and why; fix the script or the spec, commit, and re-run. -Allocate compute before running commands or recipes. A fresh installation exposes -a built-in local offer with no setup: `lc compute launch --cpus 1 --memory 1`. -Add `--name analysis` to choose a name, or omit it to receive a generated short name. -Replace `` in these commands with the returned name or a full immutable -ID from the launch or status JSON output (`--json`). -Wait for readiness with `lc compute status --wait`, then reuse the -cluster with `run` and `materialize`; neither creates compute automatically. -If the 30-minute allocation expires, launch another and use its new name. -A configured catalog replaces the default offers; inspect `lc compute resources` -if the request does not match. +Allocate compute before running commands or recipes: `lc compute launch --wait`. +With no CPU/memory flags, this starts a cluster named `local` using all detected +usable CPUs and RAM on this machine and waits until it is ready. No configuration +is needed. GPUs still require explicit offers. Add `--name analysis` to choose +another name. Replace `` in these commands with the returned name +(normally `local`) or a full immutable ID from launch or status JSON (`--json`). +Reuse the cluster with `run` and `materialize`; neither creates compute automatically. + +Only one local cluster can run per user on this machine, even with different names +or catalogs. If one already exists, inspect `lc compute status` and reuse it rather +than launching another. `lc compute status --wait` waits for an +existing cluster. Launch's `--wait` defaults to a 300-second readiness timeout; +`--timeout SECONDS` overrides it. A waiting launch that fails reports the accepted +cluster ID and leaves the allocation unchanged: inspect it before retrying. +The default local lifetime is 30 minutes; `--time` overrides it up to two hours +for the built-in offer. After expiry, launch again; the name `local` can be reused, +but the immutable ID changes. + +Compute configuration is `~/.lightcone/compute.yaml`, or the file selected by +`LC_COMPUTE_CONFIG`. Its `local.resources` mapping can override the built-in CPU/RAM +budget (for example, `{cpus: 4, memory: 8GiB}`). Explicit local connections use +their own offers instead. Remote offers otherwise coexist with the default local +offer. If `local.enabled: false` is configured, respect that policy: local launch +and execution are disabled. Inspect `lc compute resources` and supply both +`--cpus` and `--memory` to select a configured remote allocation; `--wait` works +there too. The no-resource shortcut never selects remote compute automatically. ## Recipe template grammar From 4854e61f9b1b907892b000fc62da27e2346ae33f Mon Sep 17 00:00:00 2001 From: Francois Lanusse Date: Tue, 29 Sep 2026 09:44:09 -0700 Subject: [PATCH 3/4] Harden local allocation startup and guard NERSC login nodes Make detached local startup transactional with a parent-child pipe, clean unpublished failures, and preserve published identities. Protect the per-user lock, retain owner diagnostics, and handle closed stdio and the brief kernel teardown window after expiry. Disable local launch and execution on recognized NERSC login nodes while permitting interactive compute nodes and retaining Slurm access, inspection, and termination. Clarify catalog errors, isolate lifecycle tests, and update the documentation. Validation: 249 compute, local lifecycle, and output tests passed; Ruff, mypy, and git diff --check passed. Earlier Slurm and CLI validation passed with one bootstrap test passing on retry. No live NERSC allocation was tested; the documented home-filesystem flock limitation remains. --- CLAUDE.md | 51 ++- docs/api/compute.md | 23 +- docs/architecture.md | 5 +- docs/cli/compute.md | 6 + docs/user/cluster.md | 35 ++- docs/user/getting-started.md | 6 +- evals/prompt.md | 8 +- src/lightcone/cli/compute.py | 6 +- src/lightcone/engine/compute/__init__.py | 22 +- src/lightcone/engine/compute/catalog.py | 42 ++- src/lightcone/engine/compute/local.py | 201 ++++++++---- src/lightcone/engine/compute/local_runtime.py | 21 +- tests/conftest.py | 17 + tests/test_compute.py | 116 ++++++- tests/test_compute_local.py | 292 +++++++++++++++--- tests/test_compute_output.py | 2 + 16 files changed, 697 insertions(+), 156 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 23386042..482b2768 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -182,7 +182,7 @@ src/lightcone/ # namespace โ€” NO __init__.py โ”œโ”€โ”€ compute/ # explicit allocations and borrowed Dask clients โ”‚ โ”œโ”€โ”€ __init__.py # Compute: catalog, resolve, launch, status, down; connect() โ”‚ โ”œโ”€โ”€ model.py # the shared Pydantic models and the Provider protocol - โ”‚ โ”œโ”€โ”€ catalog.py # compute.yaml, or the built-in local offer + โ”‚ โ”œโ”€โ”€ catalog.py # compute.yaml with local defaults and policy โ”‚ โ”œโ”€โ”€ runtime.py # private files, TLS material, the scheduler config โ”‚ โ”œโ”€โ”€ local.py # local provider: validated OS process identities โ”‚ โ”œโ”€โ”€ local_runtime.py # the detached LocalCluster owner @@ -1710,7 +1710,8 @@ identities are the allocation authority; standard Dask supplies execution state. No Lightcone server, lifecycle database, custom Dask worker, or implicit allocation. **Names are native labels, not a registry.** `launch --name analysis` chooses a name; -otherwise launch generates `lc-` plus 12 random hexadecimal characters. Plain stdout +otherwise the local shortcut uses `local`, and explicit resource requests generate +`lc-` plus 12 random hexadecimal characters. Plain stdout contains only the name for shell capture; JSON retains the full immutable ID too. Resolve names through fresh discovery and refuse missing, ambiguous, or incomplete observations. Check existing names before submission, but do not claim atomic global @@ -1724,15 +1725,25 @@ comment retention requires Slurm's `AccountingStoreFlags` to include `job_commen Resolve the Slurm command user's UID through `id -u` on the same command runner, and use it for every native ownership check and filter. -**Local compute needs no setup.** An absent implicit `~/.lightcone/compute.yaml` -selects a built-in local catalog: one CPU, 1 GiB, one node, fast startup, 30-minute -default and two-hour maximum lifetime. It writes no catalog and starts no cluster. +**Local compute needs no setup.** The built-in local offer provides detected usable +logical CPUs and RAM, one node, fast startup, a 30-minute default and two-hour +maximum lifetime. Loading the catalog writes no catalog and starts no cluster. +`lc compute launch` without CPU/memory flags selects only local offers and defaults +the name to `local`. `--wait` returns when the accepted allocation is ready; timeout +or startup failure retains its ID without resubmitting or terminating it. +Configured remote offers precede the built-in local offer in selection order. +Explicit local connections supply their own offers instead. `local.resources` +overrides the built-in CPU/RAM budget and cannot accompany explicit local connections. +`local.enabled: false` blocks local launch and execution while preserving inspection +and termination. Recognized NERSC login nodes disable local compute automatically; +other sites can disable it in their catalogs. Native permissions remain the +enforcement boundary. GPU offers require an explicit catalog and, for local launches, a nonempty `CUDA_VISIBLE_DEVICES` mask on Linux. No GPU auto-discovery. Local GPU capacity and model labels are configured, not hardware-verified; allocations do not reserve devices exclusively against other host programs or allocations. -Configured catalogs replace it completely; missing explicit paths and invalid -files are errors. Execution still requires an explicitly launched cluster's name or ID. +Missing explicit paths and invalid files are errors. Execution still requires an +explicitly launched cluster's name or ID. **A Slurm connection needs no launch settings (2026-09).** Every `launch` key defaults, and the defaults assume a home directory shared by login and compute @@ -1869,9 +1880,12 @@ ended. Every scheduler lc launches runs under `runtime.SCHEDULER_CONFIG`, whose zero `events-cleanup-delay` drops a departed client's forwarded recipe output instead of holding it for Dask's default hour. -**No login-node guard or venue module (PR #226 review).** Explicit catalog -selection and native backend permissions determine allocation. Do not infer -permission from hostnames, NERSC_HOST, or inherited Slurm job variables. Both +**Block local compute on recognized NERSC login nodes (2026-09).** A nonempty +`NERSC_HOST` plus a short hostname matching `login[0-9]+` disables local launch +and execution, including explicit local offers and `local.enabled: true`. +Interactive compute nodes remain eligible; inherited `SLURM_JOB_ID` never exempts +a login node. Apply this policy at runtime without writing a configuration file. +Keep inspection, termination, and Slurm execution available. Both execution commands submit ordinary tasks to the Dask scheduler without worker restrictions; existing task runtime gates and sandbox checks remain. Do not add a per-worker validation framework around ordinary Dask task submission. Shared project storage @@ -2112,6 +2126,23 @@ unlinks before writing; a new tampering test should too. ### Recorded decisions +- **Guard NERSC login nodes by default (2026-09).** This reverses the earlier + no-login-node-guard decision: an unconfigured first launch must not allocate a + whole shared login node. Detect the documented NERSC environment marker and + login hostname together, without DNS queries or scheduler probes. A runtime + restriction also covers copied or incomplete catalogs; no generated file or + override flag is needed. Local compute in interactive compute-node sessions + remains available, subject to the configured local policy. + +- **Local compute accompanies remote catalogs (2026-09).** The built-in offer + uses the host's usable CPU/RAM capacity and follows configured offers, replacing + the previous one-CPU/1-GiB fallback that disappeared when a catalog existed. + `local.resources` sets a smaller budget; explicit local connections use their + own offers. Login-node catalogs disable local launch and execution through + `local.enabled: false`. Inspection and termination stay available. One local + allocation per user per machine is enforced across catalogs and connection + roots by an OS lock held throughout the detached owner's lifetime. + - **The engine is the host's uv tool, never a project dependency** (2026-08, reversing spec ยง2's engine-in-lock rule and deleting layer 3). `lc init` scaffolds no `lightcone-cli` dependency, and there is no diff --git a/docs/api/compute.md b/docs/api/compute.md index c5434843..83a2d3e8 100644 --- a/docs/api/compute.md +++ b/docs/api/compute.md @@ -24,6 +24,11 @@ for inspection and termination. Remote catalogs retain the implicit local offer unless disabled. Explicit local connections supply their own offers instead and cannot be combined with `local.resources`. GPU offers require explicit configuration. Loading creates no configuration file or allocation. +The effective local policy also disables local offers on recognized NERSC login +nodes: nonempty `NERSC_HOST` and a short hostname matching `login[0-9]+`. +Explicit enablement and Slurm job environment variables do not override this +guard; interactive compute nodes remain eligible. Local planning, launch, and +execution check the same policy, while status and termination remain available. Missing paths selected through an argument or `LC_COMPUTE_CONFIG`, unreadable files, and invalid catalogs remain errors. Stable connection namespaces let separate invocations discover and attach to the same local allocations. @@ -143,7 +148,8 @@ See [Slurm's accounting field documentation](https://slurm.schedmd.com/sacct.htm Execution submits ordinary tasks through the borrowed client's `submit` method. Dask chooses the workers and handles dependencies; invocation-specific keys prevent unintended reuse across commands. There is no worker-selection layer, per-worker -preflight orchestration, source fingerprinting, or login-node guard. Driver-side +preflight orchestration, or source fingerprinting. The local login-node guard +does not restrict remote Slurm execution from a login node. Driver-side preparation and the existing task runtime/sandbox checks remain in their owners. Workers advertise standard Dask `CPU`, `MEMORY`, and `GPU` resources; memory is measured @@ -218,7 +224,20 @@ materialization sends recipe output to stderr to leave stdout for its report. A host-local OS file lock limits local allocations to one per user, independent of connection roots and namespaces. The launcher acquires it before spawning and passes the descriptor to the detached owner, which retains it for its lifetime. The lock -file is never unlinked; OS process exit releases the lock without stale-lock cleanup. +resides under the account's home directory at +`.lightcone/local-locks//allocation.lock`, independently of `HOME`, +`TMPDIR`, and catalog paths. Its boot UUID keeps hosts with a shared home separate. +The file is never unlinked; OS process exit releases the lock without stale-lock +cleanup. Advisory metadata identifies the owning cluster, catalog, and connection +root when a competing launch is refused. +The account-home filesystem must support `flock`; the login-node guard alone +does not establish filesystem support on an interactive compute node. + +A startup pipe lets the owner proceed only after the launcher publishes its +identity and launch records. If the launcher dies before completing publication, +the pipe closes and the owner exits, releasing the lock. Failures before identity +publication remove the launcher's private files; a published identity remains +inspectable after a startup failure. Local teardown drains the allocation's validated process group rather than assuming the owner's exit proves every child stopped. Boot UUID, UID, process diff --git a/docs/architecture.md b/docs/architecture.md index f4d72608..f5fdd193 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -171,8 +171,9 @@ A YAML catalog supplies ordered resource offers and stable native service namespaces. When the implicit default file is absent, a built-in local catalog provides all detected usable CPUs and RAM without setup. The `local` policy can override that budget or disable local compute. Remote catalogs retain the default -local offer; explicit local connections supply their own offers. GPU offers need -an explicit catalog. Missing explicit paths and invalid files +local offer; explicit local connections supply their own offers. A runtime guard +disables local compute on recognized NERSC login nodes while permitting interactive +compute nodes. GPU offers need an explicit catalog. Missing explicit paths and invalid files remain errors. No catalog is written and no allocation starts until `compute launch` resolves resources and submits once. Slurm queries and validated local OS identities are authoritative for allocations; Dask is the diff --git a/docs/cli/compute.md b/docs/cli/compute.md index 799cddcf..1a67187e 100644 --- a/docs/cli/compute.md +++ b/docs/cli/compute.md @@ -28,6 +28,12 @@ the shortcut chooses the first eligible local offer. See Missing explicit paths and invalid catalogs are errors. Loading a catalog or planning with `--dry-run` creates no files or compute. +Local launch and execution are automatically disabled on recognized NERSC login +nodes, including the first run without a catalog. Interactive compute nodes remain +eligible. This runtime guard creates no configuration file and cannot be overridden +by `local.enabled: true`; Slurm, status, and termination remain available. +See [local allocations](../user/cluster.md#local-allocations) for detection details. + | Command | Behavior | |---|---| | `resources` | Ordered available offers, per-node shape, node limit, default/maximum time, and startup class. Free capacity remains unknown. | diff --git a/docs/user/cluster.md b/docs/user/cluster.md index f3bdc6dc..4c1c3636 100644 --- a/docs/user/cluster.md +++ b/docs/user/cluster.md @@ -45,7 +45,9 @@ Only one local cluster can run per user on each machine. An OS lock held by the detached owner prevents concurrent launches, including through different names, catalogs, namespaces, or connection roots. End the existing cluster before launching another. The lock releases when its owner exits, including on failure -or walltime expiry. +or walltime expiry. A refusal identifies the owning cluster and its original +catalog and connection root. Use that catalog to inspect or stop the cluster if +the current catalog no longer includes its connection. An allocation owns a detached process session and standard `LocalCluster`: one worker process with `task_slots_per_node` threads, and a scheduler that listens on `127.0.0.1` over TLS. Its own logs are discarded; a startup failure is kept @@ -60,9 +62,19 @@ launched it. Other boot sessions are excluded from discovery, and an explicit ID from one is refused rather than reported as stopped. Once an allocation has ended, its credentials and scratch directory are removed; its full ID still reports `ended`. -Local compute is available when `local.enabled` is true and the catalog exposes a valid local offer; -Lightcone does not infer permission from login-node names or site environment -variables. Allocation choices are explicit and native permissions still apply. +Local compute requires an enabled local policy and a valid local offer. +On NERSC login nodes it is disabled automatically, even with no catalog or with +`local.enabled: true`. The guard recognizes a nonempty `NERSC_HOST` and a short +hostname matching `login[0-9]+`; it does not perform DNS or scheduler queries. +The guard permits compute nodes such as `nid200021`, including interactive +sessions. A `SLURM_JOB_ID` variable does not exempt a login node. +See NERSC's [environment conventions](https://docs.nersc.gov/environment/) and +[interactive sessions](https://docs.nersc.gov/connect/vscode/). + +Local startup also requires `flock` support in the account home, where the +singleton lock lives. NERSC documents that its compute-node home mounts do not +support `flock`, so this separate requirement can prevent local startup there; +the Slurm provider does not use that lock. A local connection's optional `launch` settings are `connection_root` (default `~/.lightcone/compute`), `scratch_root` (default: the temporary directory), @@ -92,6 +104,10 @@ succeeding. Concurrent launches can still choose the same name, so lookup also refuses ambiguous names or incomplete discovery. Use a full ID to select a known allocation directly when another connection cannot be queried. +This includes local connections when `local.enabled: false`: existing local +allocations remain visible and can still have conflicting names. Repair a +connection's discovery error before launching another cluster or resolving names. + A name may be reused once its allocation has ended. Keep the full ID when you need a durable reference to one allocation; a later cluster with the same name has a different ID. Names are discovered from allocation metadata, without a @@ -112,7 +128,8 @@ Both CPU and memory are required in `local.resources`; omit that block to use detected capacity. The same capacity validation applies to configured budgets. This controls the default offer, not hard OS resource limits. -On a login node, disable local compute and configure Slurm offers: +NERSC login nodes are guarded without setup. For other sites, or to disable local +compute on every node using the catalog, set: ```yaml version: 1 @@ -127,6 +144,9 @@ stoppable; disabling does not kill them. Bare launch reports that local compute disabled; supply CPU/memory requirements to select a Slurm allocation. Select this catalog on login nodes through `LC_COMPUTE_CONFIG`. This is Lightcone configuration policy; native site permissions enforce machine-wide restrictions. +Leave `local.enabled` at its default to use local compute inside a NERSC +interactive compute-node session. The login-node guard still applies, and it +creates no configuration file. By default, a built-in `local` connection and offer accompany remote offers, with configured offers taking selection priority. If the catalog already defines local @@ -134,6 +154,11 @@ connections, those offers replace the implicit local offer; omit `local.resource and size those offers directly. The no-resource shortcut selects the first eligible local offer and defaults its cluster name to `local`. +Resource requests can select the built-in local offer when no earlier remote +offer is eligible. Set `local.enabled: false` for catalogs that must use only remote +compute. Without an explicit local connection, the connection name `local` is +reserved for the built-in backend; its offer name is also reserved while enabled. + The namespace is a stable UUID identifying a connection; keep it unchanged while that connection's clusters exist. diff --git a/docs/user/getting-started.md b/docs/user/getting-started.md index 9843d4c4..7e2e5df6 100644 --- a/docs/user/getting-started.md +++ b/docs/user/getting-started.md @@ -251,8 +251,10 @@ the record of what your results were computed with. Launch the built-in local offer; no compute configuration is needed. It provides all usable CPUs and RAM for 30 minutes. Keep the returned ID in `CLUSTER` for this -walkthrough. If you already have a compute catalog, its offers replace that default; -see [Running on a Cluster](cluster.md). +walkthrough. Configured remote offers coexist with that default. A catalog can +override the local budget, disable local compute, or provide explicit local offers; +see [Running on a Cluster](cluster.md). NERSC login nodes automatically refuse local +compute; use a compute node in an interactive allocation or a configured Slurm offer. ```bash CLUSTER=$(lc compute launch --wait --json | python -c 'import json,sys; print(json.load(sys.stdin)["id"])') diff --git a/evals/prompt.md b/evals/prompt.md index e9c4b7da..d905996e 100644 --- a/evals/prompt.md +++ b/evals/prompt.md @@ -47,8 +47,9 @@ another name. Replace `` in these commands with the returned name Reuse the cluster with `run` and `materialize`; neither creates compute automatically. Only one local cluster can run per user on this machine, even with different names -or catalogs. If one already exists, inspect `lc compute status` and reuse it rather -than launching another. `lc compute status --wait` waits for an +or catalogs. If one already exists, use the catalog identified in the refusal to +inspect `lc compute status` and reuse it rather than launching another. +`lc compute status --wait` waits for an existing cluster. Launch's `--wait` defaults to a 300-second readiness timeout; `--timeout SECONDS` overrides it. A waiting launch that fails reports the accepted cluster ID and leaves the allocation unchanged: inspect it before retrying. @@ -64,6 +65,9 @@ offer. If `local.enabled: false` is configured, respect that policy: local launc and execution are disabled. Inspect `lc compute resources` and supply both `--cpus` and `--memory` to select a configured remote allocation; `--wait` works there too. The no-resource shortcut never selects remote compute automatically. +Recognized NERSC login nodes refuse local compute automatically, even without a +catalog. Use Slurm or an interactive compute-node session; local compute remains +eligible on those compute nodes. Do not try to override the login-node guard. ## Recipe template grammar diff --git a/src/lightcone/cli/compute.py b/src/lightcone/cli/compute.py index 1f745524..01d221d4 100644 --- a/src/lightcone/cli/compute.py +++ b/src/lightcone/cli/compute.py @@ -59,8 +59,10 @@ def _duration(seconds: int) -> str: def compute() -> None: """Allocate resources, inspect clusters, and end allocations. - The catalog is LC_COMPUTE_CONFIG, else ~/.lightcone/compute.yaml, else a - built-in local offer. + Read LC_COMPUTE_CONFIG or ~/.lightcone/compute.yaml. A built-in local + offer follows configured offers unless local compute is disabled or + explicit local connections supply their own offers. Local compute is + automatically disabled on recognized NERSC login nodes. """ diff --git a/src/lightcone/engine/compute/__init__.py b/src/lightcone/engine/compute/__init__.py index f299cfeb..694b7998 100644 --- a/src/lightcone/engine/compute/__init__.py +++ b/src/lightcone/engine/compute/__init__.py @@ -9,7 +9,7 @@ from typing import Any from uuid import uuid4 -from .catalog import Catalog +from .catalog import Catalog, local_disabled_reason from .model import ( ComputeError, Connection, @@ -144,8 +144,6 @@ def _plan_offer( ) -> LaunchPlan | None: """Match one shape and validate its provider without allocating anything.""" connection = self.catalog.connections[offer.connection] - if connection.provider == "local" and not self.catalog.local.enabled: - return None if request.num_nodes > offer.max_nodes: return None if request.startup is not None and request.startup != offer.startup.class_: @@ -175,8 +173,8 @@ def plan_local( gpus: str = "0", num_nodes: int = 1, startup: str | None = None, ) -> LaunchPlan: """Plan the first usable local offer, without considering remote backends.""" - if not self.catalog.local.enabled: - raise ComputeError("local compute is disabled by the compute configuration") + if reason := local_disabled_reason(self.catalog.local.enabled): + raise ComputeError(reason) name = "local" if name is None else name validate_name(name) unavailable: list[str] = [] @@ -201,8 +199,10 @@ def plan_local( def launch(self, plan: LaunchPlan) -> Identity: """Choose an unused name from native observations, then submit exactly once.""" - if plan.connection.provider == "local" and not self.catalog.local.enabled: - raise ComputeError("local compute is disabled by the compute configuration") + if plan.connection.provider == "local" and ( + reason := local_disabled_reason(self.catalog.local.enabled) + ): + raise ComputeError(reason) if plan.name is not None: validate_name(plan.name) snapshots, errors = self.discover() @@ -294,10 +294,10 @@ def connect(cluster_id: str, *, timeout: float = 10) -> Iterator[Any]: service = Compute() provider, identity = service.resolve(cluster_id) connection = service.catalog.connection_for(identity.namespace) - if connection.provider == "local" and not service.catalog.local.enabled: - raise ComputeError( - "local compute is disabled by the compute configuration", cluster_id=identity.encode(), - ) + if connection.provider == "local" and ( + reason := local_disabled_reason(service.catalog.local.enabled) + ): + raise ComputeError(reason, cluster_id=identity.encode()) snapshot = provider.inspect(identity) if snapshot.phase != "active": raise ComputeError( diff --git a/src/lightcone/engine/compute/catalog.py b/src/lightcone/engine/compute/catalog.py index c9bc964a..efd23fb8 100644 --- a/src/lightcone/engine/compute/catalog.py +++ b/src/lightcone/engine/compute/catalog.py @@ -3,6 +3,8 @@ from __future__ import annotations import os +import re +import socket from pathlib import Path from typing import Annotated, Any, Self @@ -25,6 +27,20 @@ _LOCAL_NAMESPACE = "22c84e48-2f0a-4cd2-90a2-30ce2e909bd1" +def local_disabled_reason(enabled: bool = True) -> str | None: + """Explain a local-compute refusal while retaining inspection and termination.""" + if os.environ.get("NERSC_HOST") and re.fullmatch( + r"login[0-9]+", socket.gethostname().split(".", 1)[0].lower(), + ): + return ( + "local compute is disabled on NERSC login nodes; use an interactive compute node " + "or configure a Slurm offer and launch with --cpus and --memory" + ) + if not enabled: + return "local compute is disabled by the compute configuration" + return None + + class _UniqueLoader(yaml.SafeLoader): """Do not silently replace a connection or limit through duplicate YAML keys.""" @@ -109,17 +125,21 @@ def load(cls, path: Path | None = None) -> Catalog: try: catalog = cls.model_validate(raw) return catalog._with_local() - except ValidationError as exc: + except (ValidationError, ComputeError) as exc: + detail = validation_message(exc) if isinstance(exc, ValidationError) else str(exc) raise ComputeError( - f"invalid compute catalog {path}:\n{validation_message(exc)}" + f"invalid compute catalog {path}:\n{detail}" ) from exc def _with_local(self) -> Catalog: """Keep explicit local connections, or add the stable built-in connection.""" connections = dict(self.connections) offers = list(self.offers) + local = self.local + if local_disabled_reason(local.enabled) is not None: + local = local.replace(enabled=False) explicit = any(connection.provider == "local" for connection in connections.values()) - if explicit and self.local.resources is not None: + if explicit and local.resources is not None: raise ComputeError( "local.resources cannot be combined with explicit local connections; " "set their offer resources instead" @@ -127,15 +147,21 @@ def _with_local(self) -> Catalog: if not explicit: if "local" in connections: raise ComputeError( - "connection name 'local' is reserved for the built-in local backend" + "connection name 'local' is reserved for the built-in local backend; " + "rename the configured connection and its offer references" ) # Retain this authority even when disabled so existing allocations can be stopped. connections["local"] = Connection(namespace=_LOCAL_NAMESPACE, provider="local") - if self.local.enabled: + if local.enabled: + if any(offer.name == "local" for offer in offers): + raise ComputeError( + "offer name 'local' is reserved for the built-in local backend; " + "rename the configured offer" + ) from dask.system import CPU_COUNT from distributed.system import MEMORY_LIMIT - resources = self.local.resources or Resources.from_bytes( + resources = local.resources or Resources.from_bytes( cpus=CPU_COUNT, memory_bytes=MEMORY_LIMIT, ) offers.append(Offer( @@ -143,11 +169,11 @@ def _with_local(self) -> Catalog: max_nodes=1, time=TimeLimits(default="30m", max="2h"), startup=Startup(class_="fast"), )) - if not self.local.enabled: + if not local.enabled: offers = [ offer for offer in offers if connections[offer.connection].provider != "local" ] - return self.replace(connections=connections, offers=offers) + return self.replace(connections=connections, offers=offers, local=local) def connection_for(self, namespace: str) -> Connection: """Find the configured authority without relying on current offers.""" diff --git a/src/lightcone/engine/compute/local.py b/src/lightcone/engine/compute/local.py index 56c44380..a5902757 100644 --- a/src/lightcone/engine/compute/local.py +++ b/src/lightcone/engine/compute/local.py @@ -2,8 +2,12 @@ from __future__ import annotations +import fcntl +import json import os +import pwd import re +import shlex import shutil import signal import socket @@ -20,6 +24,7 @@ import psutil +from lightcone.engine.compute.catalog import local_disabled_reason from lightcone.engine.compute.model import ( ComputeError, Connection, @@ -46,35 +51,85 @@ _OWNER_MODULE = "lightcone.engine.compute.local_runtime" _STOP_GRACE = 3.0 _RETIRED = "ended.json" +# The account home is stable even when a command overrides HOME. A boot UUID +# separates machines sharing that home without depending on a mutable hostname. +_LOCK_ROOT = Path(pwd.getpwuid(os.getuid()).pw_dir) / ".lightcone" / "local-locks" + + +def _above_stdio(descriptor: int) -> int: + """Keep inherited control descriptors clear of Popen's stdio redirections.""" + if descriptor < 3: + duplicate = fcntl.fcntl(descriptor, fcntl.F_DUPFD_CLOEXEC, 3) + os.close(descriptor) + return duplicate + return descriptor @contextmanager def _allocation_lock() -> Iterator[int]: - """Hold a per-user, host-local lock across catalogs and detached owner lifetime.""" - import fcntl - - # Deliberately independent of TMPDIR, connection namespaces and configured roots. - # Resolve /tmp for macOS, where it is an alias of /private/tmp. - root = private_directory(Path("/tmp").resolve() / f"lightcone-local-{os.getuid()}", create=True) + """Hold a per-user, per-boot lock across catalogs and detached owner lifetime.""" + root = private_directory(configured_directory(_LOCK_ROOT) / _boot_identity(), create=True) + path = root / "allocation.lock" descriptor: int | None = None try: - descriptor = os.open( - root / "allocation.lock", os.O_CREAT | os.O_RDWR | os.O_NOFOLLOW, 0o600, - ) - info = os.fstat(descriptor) - if (not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid() - or stat.S_IMODE(info.st_mode) != 0o600 or info.st_nlink != 1): - raise ComputeError("local allocation lock must be a private file owned by you") try: - fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB) - except BlockingIOError as exc: - raise ComputeError( - "a local cluster is already running or starting for this user on this machine; " - "stop it with lc compute down before launching another" - ) from exc + descriptor = os.open(path, os.O_CREAT | os.O_RDWR | os.O_NOFOLLOW, 0o600) + descriptor = _above_stdio(descriptor) + info = os.fstat(descriptor) + if (not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid() + or stat.S_IMODE(info.st_mode) != 0o600 or info.st_nlink != 1): + raise ComputeError("local allocation lock must be a private file owned by you") + deadline = time.monotonic() + 0.1 + while True: + try: + fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB) + break + except BlockingIOError as exc: + # A zombie owner can still have exiting threads whose + # descriptor references the kernel has not released yet. + remaining = deadline - time.monotonic() + if remaining > 0: + time.sleep(min(0.01, remaining)) + continue + message = ( + "a local cluster is already running or starting " + "for this user on this machine" + ) + try: + owner = json.loads(os.pread(descriptor, 65536, 0)) + except (OSError, ValueError): + owner = {} + cluster_id = owner.get("identity") if isinstance(owner, dict) else None + if isinstance(cluster_id, str) and cluster_id: + command = "lc compute down " + shlex.quote(cluster_id) + if owner.get("explicit_catalog") is True and isinstance( + owner.get("catalog"), str, + ): + command = ( + "LC_COMPUTE_CONFIG=" + shlex.quote(owner["catalog"]) + " " + command + ) + else: + command = "env -u LC_COMPUTE_CONFIG " + command + message += ( + f"; allocation {cluster_id} uses local namespace " + f"{owner.get('namespace')!r} with connection_root " + f"{owner.get('connection_root')!r} (locator {owner.get('locator')!r}); " + "using that launch catalog, " + f"stop it with `{command}`" + ) + else: + message += ( + "; its launcher is publishing the allocation identity; retry shortly" + ) + raise ComputeError(message) from exc + current = path.stat(follow_symlinks=False) + if (current.st_dev, current.st_ino) != (info.st_dev, info.st_ino): + raise ComputeError("the local allocation lock changed while it was being acquired") + # Stale diagnostics have no authority once the kernel grants a new lock. + os.ftruncate(descriptor, 0) + except OSError as exc: + raise ComputeError(f"cannot acquire the local allocation lock: {exc}") from exc yield descriptor - except OSError as exc: - raise ComputeError(f"cannot acquire the local allocation lock: {exc}") from exc finally: if descriptor is not None: # Do not unlock: the detached owner inherits this same open file description. @@ -109,6 +164,8 @@ def __init__(self, connection: Connection) -> None: def plan(self, offer: Offer, request: Request) -> LaunchPlan: """Validate a one-node local offer without creating allocation files.""" + if reason := local_disabled_reason(): + raise UnavailableOfferError(reason) if os.name != "posix": raise ComputeError("local allocations require POSIX process sessions and signals") if request.num_nodes != 1: @@ -175,43 +232,49 @@ def plan(self, offer: Offer, request: Request) -> LaunchPlan: def launch(self, plan: LaunchPlan) -> Identity: """Start a detached allocation owner and retain its immutable OS identity.""" + if reason := local_disabled_reason(): + raise ComputeError(reason) if plan.connection != self.connection or plan.num_nodes != 1: raise ComputeError("local launch plan belongs to a different connection or node count") if plan.name is not None: validate_name(plan.name) with _allocation_lock() as lock_fd: - # Also recognize allocations launched before the lifetime lock existed. - if self.discover(): - raise ComputeError( - "a local cluster is already running; stop it before launching another" - ) return self._launch_locked(plan, lock_fd) def _launch_locked(self, plan: LaunchPlan, lock_fd: int) -> Identity: """Transfer the singleton lock to the allocation owner before releasing our copy.""" boot = _boot_identity() token = uuid4().hex - directory = private_directory(self.root / token, create=True) - scratch = private_directory(Path(plan.details["scratch_root"]) / f"lc-{token}", create=True) + directory: Path | None = None + scratch: Path | None = None started = time.monotonic() - write_private_json( - directory / "launch.json", - { - "deadline": started + plan.seconds, - "task_slots": plan.details["task_slots_per_node"], - "scratch": str(scratch), - "identity": "", - "lock_fd": lock_fd, - }, - ) process: subprocess.Popen[bytes] | None = None identity: Identity | None = None + published = False + startup_read: int | None = None + startup_write: int | None = None environment = {**os.environ, "CUDA_VISIBLE_DEVICES": plan.details["cuda_visible_devices"]} if plan.details["cuda_device_order"] is None: environment.pop("CUDA_DEVICE_ORDER", None) else: environment["CUDA_DEVICE_ORDER"] = plan.details["cuda_device_order"] try: + directory = private_directory(self.root / token, create=True) + scratch = private_directory( + Path(plan.details["scratch_root"]) / f"lc-{token}", create=True, + ) + startup_read, startup_write = os.pipe() + startup_read = _above_stdio(startup_read) + startup_write = _above_stdio(startup_write) + launch = { + "deadline": started + plan.seconds, + "task_slots": plan.details["task_slots_per_node"], + "scratch": str(scratch), + "identity": "", + "lock_fd": lock_fd, + "startup_fd": startup_read, + } + write_private_json(directory / "launch.json", launch) # This allocation outlives a command; the ordinary run-to-completion # subprocess seam cannot own it. Logs are discarded rather than grow. process = subprocess.Popen( @@ -221,9 +284,11 @@ def _launch_locked(self, plan: LaunchPlan, lock_fd: int) -> Identity: stderr=subprocess.DEVNULL, start_new_session=True, close_fds=True, - pass_fds=(lock_fd,), + pass_fds=(lock_fd, startup_read), env=environment, ) + os.close(startup_read) + startup_read = None identity = Identity( namespace=self.connection.namespace, native_id=str(process.pid), token=token, host=socket.gethostname(), @@ -241,19 +306,26 @@ def _launch_locked(self, plan: LaunchPlan, lock_fd: int) -> Identity: "accelerator_name": plan.resources.accelerator_name or "GPU", } write_private_json(directory / "identity.json", record) - # The child waits for this file before publishing its TLS connection. - write_private_json( - directory / "launch.json", - { - "deadline": started + plan.seconds, - "task_slots": plan.details["task_slots_per_node"], - "scratch": str(scratch), - "identity": identity.encode(), - "lock_fd": lock_fd, - }, - ) + published = True + launch["identity"] = identity.encode() + write_private_json(directory / "launch.json", launch) + owner = json.dumps({ + "identity": identity.encode(), + "namespace": self.connection.namespace, + "connection_root": str(self.root.parent), + "locator": str(self.root), + "explicit_catalog": "LC_COMPUTE_CONFIG" in os.environ, + "catalog": str(Path(os.environ.get( + "LC_COMPUTE_CONFIG", "~/.lightcone/compute.yaml", + )).expanduser().absolute()), + }).encode() + if os.pwrite(lock_fd, owner, 0) != len(owner): + raise OSError("cannot publish the local allocation lock owner") + # EOF without this byte means the launcher died before publication, + # including SIGKILL, which no Python exception handler can clean up. + os.write(startup_write, b"1") return identity - except Exception as exc: + except BaseException as exc: if process is not None: # The unreturned child is still our direct Popen child; no file # lookup or stale PID is needed to identify this failed launch. @@ -262,21 +334,24 @@ def _launch_locked(self, plan: LaunchPlan, lock_fd: int) -> Identity: except ProcessLookupError: pass process.wait(timeout=_STOP_GRACE) - if identity is None: - # Publication is the discovery boundary. A failed spawn has no - # allocation to retain, and these are the only files we wrote. - try: - (directory / "launch.json").unlink(missing_ok=True) - directory.rmdir() - scratch.rmdir() - except OSError: - # Discovery ignores unpublished directories even if cleanup - # is interrupted or the filesystem becomes unavailable. - pass + if published and directory is not None: + self._retire(directory) + else: + # No published identity exists to retain or inspect. + if directory is not None: + shutil.rmtree(directory, ignore_errors=True) + if scratch is not None: + shutil.rmtree(scratch, ignore_errors=True) + if not isinstance(exc, Exception): + raise raise ComputeError( f"cannot start the local allocation: {exc}", - cluster_id=identity.encode() if identity is not None else None, + cluster_id=identity.encode() if published and identity is not None else None, ) from exc + finally: + for descriptor in (startup_read, startup_write): + if descriptor is not None: + os.close(descriptor) def _directory(self, identity: Identity) -> Path: if ( diff --git a/src/lightcone/engine/compute/local_runtime.py b/src/lightcone/engine/compute/local_runtime.py index b6e3785e..cf643f06 100644 --- a/src/lightcone/engine/compute/local_runtime.py +++ b/src/lightcone/engine/compute/local_runtime.py @@ -4,6 +4,7 @@ import atexit import os +import select import signal import sys import threading @@ -25,9 +26,10 @@ def main() -> None: os.umask(0o077) directory = private_directory(Path(sys.argv[1])) launch = read_private_json(directory / "launch.json") - if "lock_fd" in launch: - # Keep the inherited lock until process exit, but exclude worker execs. - os.set_inheritable(int(launch["lock_fd"]), False) + # Keep the inherited lock until process exit, but exclude worker execs. + os.set_inheritable(int(launch["lock_fd"]), False) + startup_fd = int(launch["startup_fd"]) + os.set_inheritable(startup_fd, False) if os.getsid(0) != os.getpid() or os.getpgrp() != os.getpid(): raise RuntimeError("the local allocation owner must lead its own process session") stopped = threading.Event() @@ -49,10 +51,17 @@ def expire(_signum: int, _frame: FrameType | None) -> None: remaining = float(launch["deadline"]) - time.monotonic() signal.setitimer(signal.ITIMER_REAL, max(0.001, remaining)) try: - while not launch["identity"]: - if stopped.wait(0.01): + try: + startup = select.poll() + startup.register(startup_fd, select.POLLIN) + while not startup.poll(100): + if stopped.is_set(): + return + if os.read(startup_fd, 1) != b"1": return - launch = read_private_json(directory / "launch.json") + finally: + os.close(startup_fd) + launch = read_private_json(directory / "launch.json") import dask from distributed import LocalCluster diff --git a/tests/conftest.py b/tests/conftest.py index 260e2cc8..e8ade6fd 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -28,6 +28,23 @@ def runner() -> CliRunner: return CliRunner() +@pytest.fixture(scope="session") +def local_allocation_lock_root(tmp_path_factory: pytest.TempPathFactory) -> Path: + """Give each test session a lock root shared with its independent launchers.""" + return tmp_path_factory.mktemp("local-allocation-locks") + + +@pytest.fixture +def local_allocation_lock( + local_allocation_lock_root: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + """Exercise the lifetime lock without contending with other test sessions.""" + from lightcone.engine.compute import local + + monkeypatch.setattr(local, "_LOCK_ROOT", local_allocation_lock_root) + monkeypatch.delenv("NERSC_HOST", raising=False) + + @pytest.fixture(autouse=True) def ambient_uv(monkeypatch: pytest.MonkeyPatch) -> None: """Strip scrubbable ``UV_*`` out of the suite's environment. diff --git a/tests/test_compute.py b/tests/test_compute.py index 31be98ce..16868ed7 100644 --- a/tests/test_compute.py +++ b/tests/test_compute.py @@ -53,6 +53,7 @@ def expand(path: Path) -> Path: monkeypatch.setattr(Path, "expanduser", expand) monkeypatch.delenv("LC_COMPUTE_CONFIG", raising=False) + monkeypatch.delenv("NERSC_HOST", raising=False) return tmp_path @@ -493,6 +494,87 @@ def test_local_shortcut_uses_detected_capacity_without_writing_files( assert not list(default_home.iterdir()) +def test_nersc_login_nodes_block_first_launch_without_writing_a_catalog( + default_home: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.delenv("NERSC_HOST", raising=False) + monkeypatch.setattr("dask.system.CPU_COUNT", 1) + monkeypatch.setattr("distributed.system.MEMORY_LIMIT", GIB) + before = compute.Compute() + plan = before.plan_local() + monkeypatch.setenv("NERSC_HOST", "perlmutter") + monkeypatch.setenv("SLURM_JOB_ID", "123") + monkeypatch.setattr("socket.gethostname", lambda: "login07.nersc.gov") + service = compute.Compute() + assert not service.catalog.local.enabled + assert service.catalog.connections["local"] == plan.connection + assert service.resources()["offers"] == [] + for flags in ([], ["--dry-run"]): + result = CliRunner().invoke(main, ["compute", "launch", *flags, "--json"]) + assert result.exit_code == 1, result.output + assert "disabled on NERSC login nodes" in json.loads(result.stdout)["error"] + with pytest.raises(ComputeError, match="disabled on NERSC login nodes"): + before.launch(plan) + assert not list(default_home.iterdir()) + + +@pytest.mark.parametrize("site, hostname", [ + ("perlmutter", "nid005678"), ("perlmutter", "workstation"), ("", "login07"), +]) +def test_local_compute_remains_available_outside_identified_nersc_login_nodes( + default_home: Path, monkeypatch: pytest.MonkeyPatch, site: str, hostname: str, +) -> None: + monkeypatch.setenv("NERSC_HOST", site) + monkeypatch.setenv("SLURM_JOB_ID", "123") + monkeypatch.setattr("socket.gethostname", lambda: hostname) + monkeypatch.setattr("dask.system.CPU_COUNT", 1) + monkeypatch.setattr("distributed.system.MEMORY_LIMIT", GIB) + result = CliRunner().invoke(main, ["compute", "launch", "--dry-run", "--json"]) + assert result.exit_code == 0, result.output + assert json.loads(result.stdout)["plan"]["offer"] == "local" + assert not list(default_home.iterdir()) + + +def test_nersc_login_guard_keeps_remote_compute_and_local_inspection_available( + catalog: Path, provider: MagicMock, monkeypatch: pytest.MonkeyPatch, +) -> None: + connection = compute.Compute().catalog.connections["local"] + identity = IDENTITY.replace(namespace=connection.namespace, native_id="5678") + data = yaml.safe_load(catalog.read_text()) + data["local"] = {"enabled": True} + data["connections"]["workstation"] = connection.model_dump() + data["offers"].insert(0, { + **data["offers"][0], "name": "workstation", "connection": "workstation", + }) + catalog.write_text(yaml.safe_dump(data)) + local = MagicMock() + local.discover.return_value = [] + local.inspect.return_value = Snapshot(identity=identity, phase="active", num_nodes=1) + local.connect.return_value.__enter__.return_value.scheduler_info.return_value = { + "workers": {"one": {}}, + } + monkeypatch.setitem(compute.PROVIDERS, "local", lambda connection: local) + monkeypatch.setenv("NERSC_HOST", "perlmutter") + monkeypatch.setattr("socket.gethostname", lambda: "login07") + service = compute.Compute() + assert not service.catalog.local.enabled + assert [offer.name for offer in service.catalog.offers] == ["quick", "large"] + plan = service.plan(Request.parse("4", "8")) + assert service.launch(plan) == IDENTITY + assert service.status(IDENTITY.encode()).ready + with compute.connect(IDENTITY.encode()) as client: + assert client is provider.connect.return_value.__enter__.return_value + with pytest.raises(ComputeError, match="disabled on NERSC login nodes"): + with compute.connect(identity.encode()): + pytest.fail("borrowed local compute on a login node") + local.connect.assert_not_called() + assert service.status(identity.encode()).ready + assert service.down(identity.encode()) == identity + local.terminate.assert_called_once_with(identity) + local.plan.assert_not_called() + local.launch.assert_not_called() + + def test_local_config_overrides_default_and_survives_disabling( default_home: Path, monkeypatch: pytest.MonkeyPatch, ) -> None: @@ -535,6 +617,37 @@ def test_catalog_adds_local_after_remote_offers_and_shortcut_never_selects_remot provider.launch.assert_not_called() +@pytest.mark.parametrize("kind, enabled", [ + ("connection", True), ("connection", False), ("offer", True), +]) +def test_builtin_name_conflicts_identify_the_catalog_and_remedy( + catalog: Path, kind: str, enabled: bool, +) -> None: + data = yaml.safe_load(catalog.read_text()) + data["local"] = {"enabled": enabled} + if kind == "connection": + data["connections"]["local"] = data["connections"].pop("test") + for offer in data["offers"]: + offer["connection"] = "local" + else: + data["offers"][0]["name"] = "local" + catalog.write_text(yaml.safe_dump(data)) + with pytest.raises(ComputeError) as error: + Catalog.load() + assert f"invalid compute catalog {catalog}" in str(error.value) + assert f"{kind} name 'local' is reserved for the built-in local backend" in str(error.value) + assert f"rename the configured {kind}" in str(error.value) + + +def test_disabled_builtin_does_not_reserve_remote_offer_names(catalog: Path) -> None: + data = yaml.safe_load(catalog.read_text()) + data["offers"][0]["name"] = "local" + catalog.write_text(yaml.safe_dump(data)) + loaded = Catalog.load() + assert [offer.name for offer in loaded.offers] == ["local", "large"] + assert loaded.connections[loaded.offers[0].connection].provider == "fake" + + def test_explicit_local_offers_keep_their_sizes_and_replace_the_implicit_offer( catalog: Path, monkeypatch: pytest.MonkeyPatch, ) -> None: @@ -581,7 +694,7 @@ def test_local_policy_validation(catalog: Path, settings: object, message: str) Catalog.load() -def test_disabled_policy_blocks_explicit_local_offers_and_execution_but_allows_down( +def test_disabled_policy_blocks_local_execution_but_allows_status_and_down( catalog: Path, provider: MagicMock, monkeypatch: pytest.MonkeyPatch, ) -> None: data = yaml.safe_load(catalog.read_text()) @@ -594,6 +707,7 @@ def test_disabled_policy_blocks_explicit_local_offers_and_execution_but_allows_d with compute.connect(IDENTITY.encode()): pytest.fail("borrowed disabled local compute") provider.connect.assert_not_called() + assert service.status(IDENTITY.encode()).ready service.down(IDENTITY.encode()) provider.terminate.assert_called_once_with(IDENTITY) diff --git a/tests/test_compute_local.py b/tests/test_compute_local.py index 70fcdbdb..6d0e1ce5 100644 --- a/tests/test_compute_local.py +++ b/tests/test_compute_local.py @@ -34,6 +34,7 @@ Request, Resources, TimeLimits, + UnavailableOfferError, ) from lightcone.engine.compute.runtime import ( configured_directory, @@ -43,6 +44,8 @@ write_private_json, ) +pytestmark = pytest.mark.usefixtures("local_allocation_lock") + @pytest.fixture def provider(tmp_path: Path) -> LocalProvider: @@ -74,10 +77,13 @@ def test_singleton_survives_cli_exit_and_spans_names_and_connection_roots( # Independent launchers compete for the same host/user lock, with distinct catalogs. script = """ import json, sys +from pathlib import Path +from lightcone.engine.compute import local from lightcone.engine.compute.local import LocalProvider from lightcone.engine.compute.model import ( Connection, ComputeError, Offer, Request, Resources, TimeLimits, ) +local._LOCK_ROOT = Path(sys.argv[3]) p = LocalProvider(Connection.model_validate_json(sys.argv[1])) offer = Offer(name='small', connection='workstation', resources=Resources(cpus=1, memory_gib=0.5), max_nodes=1, time=TimeLimits(default='1m', max='1m')) @@ -93,7 +99,8 @@ def test_singleton_survives_cli_exit_and_spans_names_and_connection_roots( )) providers = [provider, other] processes = [subprocess.Popen( - [sys.executable, "-c", script, item.connection.model_dump_json(), f"local-{index}"], + [sys.executable, "-c", script, item.connection.model_dump_json(), + f"local-{index}", str(local._LOCK_ROOT)], stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, ) for index, item in enumerate(providers)] identities: list[tuple[LocalProvider, Identity]] = [] @@ -114,8 +121,11 @@ def test_singleton_survives_cli_exit_and_spans_names_and_connection_roots( assert "already running or starting" in next(r["error"] for r in results if "error" in r) owner, identity = identities[0] _ready(owner, identity) - with pytest.raises(ComputeError, match="already running"): + with pytest.raises(ComputeError, match="already running") as conflict: _launch(other) + assert identity.encode() in str(conflict.value) + assert conflict.value.cluster_id is None + assert str(owner.root) in str(conflict.value) owner.terminate(identity) _ended(owner, identity) replacement = _launch(other) @@ -130,6 +140,51 @@ def test_singleton_survives_cli_exit_and_spans_names_and_connection_roots( item.terminate(identity) +def test_allocation_lock_preserves_errors_from_the_launch_body() -> None: + error = OSError("allocation storage is full") + with pytest.raises(OSError) as raised, local._allocation_lock(): + raise error + assert raised.value is error + with local._allocation_lock(): + pass + + +@pytest.mark.parametrize("closed", [(0,), (1,), (2,), (0, 1, 2)]) +def test_singleton_lock_survives_launch_with_closed_standard_descriptors( + provider: LocalProvider, tmp_path: Path, closed: tuple[int, ...], +) -> None: + result_path = tmp_path / "launched.json" + script = """ +import json, os, sys +from pathlib import Path +from lightcone.engine.compute import local +from lightcone.engine.compute.model import Connection, Offer, Request, Resources, TimeLimits +local._LOCK_ROOT = Path(sys.argv[2]) +p = local.LocalProvider(Connection.model_validate_json(sys.argv[1])) +offer = Offer(name='small', connection='workstation', resources=Resources(cpus=1, memory_gib=0.5), + max_nodes=1, time=TimeLimits(default='1m', max='1m')) +plan = p.plan(offer, Request(cpus=1, memory_bytes=512 * 1024**2)) +for descriptor in json.loads(sys.argv[4]): + os.close(descriptor) +identity = p.launch(plan) +Path(sys.argv[3]).write_text(identity.encode()) +""" + try: + result = subprocess.run( + [sys.executable, "-c", script, provider.connection.model_dump_json(), + str(local._LOCK_ROOT), str(result_path), json.dumps(closed)], + capture_output=True, text=True, timeout=15, + ) + assert result.returncode == 0, result.stderr + identity = Identity.decode(result_path.read_text()) + _ready(provider, identity) + with pytest.raises(ComputeError, match="already running"): + _launch(provider) + finally: + for snapshot in provider.discover(): + provider.terminate(snapshot.identity) + + def _ready(provider: LocalProvider, identity: Identity) -> dict[str, object]: deadline = time.monotonic() + 20 while time.monotonic() < deadline: @@ -176,8 +231,11 @@ def spawn() -> int: def test_allocation_survives_launcher_and_borrowed_client_exit(provider: LocalProvider) -> None: script = """ import json, sys +from pathlib import Path +from lightcone.engine.compute import local from lightcone.engine.compute.local import LocalProvider from lightcone.engine.compute.model import Connection, Offer, Request, Resources, TimeLimits +local._LOCK_ROOT = Path(sys.argv[2]) p = LocalProvider(Connection(**json.loads(sys.argv[1]))) offer = Offer( name='small', connection='workstation', resources=Resources(cpus=1, memory_gib=0.5), @@ -186,7 +244,8 @@ def test_allocation_survives_launcher_and_borrowed_client_exit(provider: LocalPr print(p.launch(p.plan(offer, Request(cpus=1, memory_bytes=512 * 1024**2))).encode()) """ launched = subprocess.run( - [sys.executable, "-c", script, json.dumps(provider.connection.model_dump())], + [sys.executable, "-c", script, json.dumps(provider.connection.model_dump()), + str(local._LOCK_ROOT)], check=True, capture_output=True, text=True, @@ -248,51 +307,44 @@ def test_builtin_allocation_can_be_reopened_in_another_process_without_a_catalog from lightcone.engine.compute import Compute from lightcone.engine.compute.model import GIB - expanduser = Path.expanduser - - def expand(path: Path) -> Path: - if str(path) == "~": - return tmp_path - if str(path).startswith("~/"): - return tmp_path / str(path)[2:] - return expanduser(path) - - monkeypatch.setattr(Path, "expanduser", expand) + monkeypatch.setenv("HOME", str(tmp_path)) monkeypatch.delenv("LC_COMPUTE_CONFIG", raising=False) monkeypatch.setattr("dask.system.CPU_COUNT", 1) monkeypatch.setattr("distributed.system.MEMORY_LIMIT", GIB) - result = CliRunner().invoke(main, [ - "compute", "launch", "--wait", "--time", "1m", "--timeout", "20", "--json", - ]) - data = json.loads(result.stdout) - identity = Identity.decode(data["id"]) try: + result = CliRunner().invoke(main, [ + "compute", "launch", "--wait", "--time", "1m", "--timeout", "20", "--json", + ]) assert result.exit_code == 0, result.output + data = json.loads(result.stdout) + identity = Identity.decode(data["id"]) assert data["ready"] and data["name"] == "local" script = """ import sys -from pathlib import Path -expanduser = Path.expanduser -def expand(path): - if str(path) == '~': - return Path(sys.argv[1]) - if str(path).startswith('~/'): - return Path(sys.argv[1]) / str(path)[2:] - return expanduser(path) -Path.expanduser = expand from lightcone.engine.compute import Compute, connect -assert Compute().status(sys.argv[2], wait=True, timeout=20).ready -with connect(sys.argv[2]) as client: +assert Compute().status(sys.argv[1], wait=True, timeout=20).ready +with connect(sys.argv[1]) as client: print(client.submit(sum, [4, 5]).result(timeout=5)) """ result = subprocess.run( - [sys.executable, "-c", script, str(tmp_path), identity.encode()], + [sys.executable, "-c", script, identity.encode()], capture_output=True, text=True, timeout=30, check=True, ) assert result.stdout.strip() == "9" assert not (tmp_path / ".lightcone" / "compute.yaml").exists() + builtin = LocalProvider(Compute().catalog.connections["local"]) + with monkeypatch.context() as other_catalog: + other_catalog.setenv("LC_COMPUTE_CONFIG", str(tmp_path / "other.yaml")) + with pytest.raises(ComputeError, match="already running") as conflict: + _launch(builtin) + command = ["env", "-u", "LC_COMPUTE_CONFIG", "lc", "compute", "down", identity.encode()] + assert " ".join(command) in str(conflict.value) + stopped = subprocess.run(command, capture_output=True, text=True, timeout=15) + assert stopped.returncode == 0, stopped.stderr finally: - Compute().down(identity.encode()) + # Discovery also covers a successful launch whose CLI output is malformed. + for snapshot in Compute().discover()[0]: + Compute().down(snapshot.identity.encode()) assert Compute().status(identity.encode()).phase == "ended" assert Compute().discover() == ([], {}) @@ -377,6 +429,13 @@ def test_named_local_allocation_is_discovered_and_name_can_be_reused_after_down( # A new adapter reconstructs the name from the existing native locator. assert [item.identity for item in LocalProvider(provider.connection).discover()] == [first] assert Compute().status("analysis", wait=True, timeout=20).identity == first + with monkeypatch.context() as other_catalog: + other_catalog.setenv("LC_COMPUTE_CONFIG", str(tmp_path / "other.yaml")) + with pytest.raises(ComputeError, match="already running") as conflict: + _launch(provider) + assert str(catalog) in str(conflict.value) + assert first.encode() in str(conflict.value) + assert conflict.value.cluster_id is None with connect("analysis") as client: assert client.submit(sum, [7, 8]).result(timeout=5) == 15 Compute().down("analysis") @@ -394,10 +453,14 @@ def test_named_local_allocation_is_discovered_and_name_can_be_reused_after_down( def test_walltime_expires_without_a_connected_client(provider: LocalProvider) -> None: identity = _launch(provider, seconds=2) + identities = [identity] try: _ended(provider, identity, timeout=6) + # Native exit and release of the owner's last kernel descriptor can differ. + identities.append(_launch(provider)) finally: - provider.terminate(identity) + for allocation in identities: + provider.terminate(allocation) @pytest.mark.parametrize("ending", ["down", "walltime"]) @@ -581,6 +644,8 @@ def fail(argv: list[str], **kwargs: Any) -> subprocess.Popen[bytes]: patch.setattr("lightcone.engine.compute.local.subprocess.Popen", fail) with pytest.raises(ComputeError, match="cannot execute"): _launch(provider) + assert list(provider.root.iterdir()) == [] + assert list(Path(provider.connection.launch["scratch_root"]).iterdir()) == [] # A failed spawn must release the singleton lock as well as its private files. identity = _launch(provider) try: @@ -594,6 +659,120 @@ def fail(argv: list[str], **kwargs: Any) -> subprocess.Popen[bytes]: provider.terminate(identity) +def test_failed_initial_launch_write_removes_unpublished_files_and_releases_lock( + provider: LocalProvider, monkeypatch: pytest.MonkeyPatch, +) -> None: + def fail(path: Path, value: dict[str, Any]) -> None: + if path.name == "launch.json": + raise OSError("allocation storage is full") + write_private_json(path, value) + + monkeypatch.setattr(local, "write_private_json", fail) + with pytest.raises(ComputeError, match="cannot start.*storage is full"): + _launch(provider) + assert list(provider.root.iterdir()) == [] + assert list(Path(provider.connection.launch["scratch_root"]).iterdir()) == [] + with local._allocation_lock(): + pass + + +def test_interrupted_launch_kills_the_unreturned_owner_and_releases_lock( + provider: LocalProvider, monkeypatch: pytest.MonkeyPatch, +) -> None: + children: list[subprocess.Popen[bytes]] = [] + popen = subprocess.Popen + + def capture(argv: list[str], **kwargs: Any) -> subprocess.Popen[bytes]: + process = popen(argv, **kwargs) + if "lightcone.engine.compute.local_runtime" in argv: + children.append(process) + return process + + def interrupt(path: Path, value: dict[str, Any]) -> None: + if path.name == "launch.json" and value["identity"]: + raise KeyboardInterrupt + write_private_json(path, value) + + monkeypatch.setattr(local.subprocess, "Popen", capture) + monkeypatch.setattr(local, "write_private_json", interrupt) + try: + with pytest.raises(KeyboardInterrupt): + _launch(provider) + assert len(children) == 1 + assert children[0].poll() is not None + assert provider.discover() == [] + with local._allocation_lock(): + pass + finally: + for child in children: + if child.poll() is None: + child.kill() + child.wait(timeout=5) + + +@pytest.mark.parametrize("publication", ["identity.json", "launch.json"]) +def test_owner_exits_when_launcher_is_killed_before_startup_commit( + provider: LocalProvider, tmp_path: Path, publication: str, +) -> None: + marker = tmp_path / "paused.json" + script = """ +import json, sys, time +from pathlib import Path +from lightcone.engine.compute import local +from lightcone.engine.compute.model import Connection, Offer, Request, Resources, TimeLimits +local._LOCK_ROOT = Path(sys.argv[2]) +write = local.write_private_json +def pause(path, value): + if path.name == sys.argv[4] and (path.name != 'launch.json' or value['identity']): + record = value if path.name == 'identity.json' else local.read_private_json( + path.parent / 'identity.json') + checkpoint = {'pid': record['pid'], 'directory': str(path.parent)} + marker = Path(sys.argv[3]) + pending = marker.with_suffix('.pending') + pending.write_text(json.dumps(checkpoint)) + pending.replace(marker) + time.sleep(60) + write(path, value) +local.write_private_json = pause +p = local.LocalProvider(Connection.model_validate_json(sys.argv[1])) +offer = Offer(name='small', connection='workstation', resources=Resources(cpus=1, memory_gib=0.5), + max_nodes=1, time=TimeLimits(default='1m', max='1m')) +p.launch(p.plan(offer, Request(cpus=1, memory_bytes=512 * 1024**2))) +""" + launcher = subprocess.Popen( + [sys.executable, "-c", script, provider.connection.model_dump_json(), + str(local._LOCK_ROOT), str(marker), publication], + stdout=subprocess.PIPE, stderr=subprocess.PIPE, + ) + owner: psutil.Process | None = None + try: + deadline = time.monotonic() + 15 + while not marker.exists(): + if launcher.poll() is not None: + pytest.fail(launcher.communicate()[1].decode(errors="replace")) + assert time.monotonic() < deadline, "launcher did not reach the publication checkpoint" + time.sleep(0.05) + checkpoint = json.loads(marker.read_text()) + owner = psutil.Process(checkpoint["pid"]) + launcher.kill() + launcher.communicate(timeout=5) + deadline = time.monotonic() + 5 + while owner.is_running() and owner.status() != psutil.STATUS_ZOMBIE: + assert time.monotonic() < deadline, "unpublished owner retained the singleton lock" + time.sleep(0.05) + assert not (Path(checkpoint["directory"]) / "connection.json").exists() + with local._allocation_lock(): + pass + finally: + if launcher.poll() is None: + launcher.kill() + launcher.communicate(timeout=5) + if owner is not None and owner.is_running(): + owner.kill() + for snapshot in provider.discover(): + provider.terminate(snapshot.identity) + + def test_reused_pid_and_boot_identity_are_never_signalled( provider: LocalProvider, monkeypatch: pytest.MonkeyPatch ) -> None: @@ -911,18 +1090,23 @@ def test_configured_roots_reject_relative_parent_and_unresolvable_paths(tmp_path private_directory(configured_directory(loop), create=True) -def test_local_plan_does_not_infer_policy_from_login_hostname_or_slurm_environment( +def test_local_provider_allows_interactive_nodes_but_refuses_login_nodes_before_allocation( provider: LocalProvider, monkeypatch: pytest.MonkeyPatch ) -> None: monkeypatch.setenv("NERSC_HOST", "perlmutter") monkeypatch.setenv("SLURM_JOB_ID", "123") - monkeypatch.setattr(socket, "gethostname", lambda: "login01") + monkeypatch.setattr(socket, "gethostname", lambda: "nid200021") offer = Offer( name="small", connection="workstation", resources=Resources(cpus=1, memory_gib=0.5), max_nodes=1, time=TimeLimits(default="1m", max="1m"), ) plan = provider.plan(offer, Request(cpus=1, memory_bytes=512 * 1024**2)) assert plan.resources == offer.resources + monkeypatch.setattr(socket, "gethostname", lambda: "login01") + with pytest.raises(UnavailableOfferError, match="NERSC login nodes"): + provider.plan(offer, plan.request) + with pytest.raises(ComputeError, match="NERSC login nodes"): + provider.launch(plan) assert not provider.root.exists() @@ -953,13 +1137,26 @@ def test_local_gpu_plan_freezes_native_mask_and_publishes_the_configured_envelop assert plan.details["cuda_visible_devices"] == (mask if gpus else "") assert plan.details["cuda_device_order"] == order assert not provider.root.exists() - monkeypatch.setattr(local, "_boot_identity", lambda: str(uuid4())) - popen = MagicMock(return_value=SimpleNamespace(pid=12345)) + boot = str(uuid4()) + monkeypatch.setattr(local, "_boot_identity", lambda: boot) + startup_readers: list[int] = [] + + def spawned(argv: list[str], **kwargs: Any) -> SimpleNamespace: + launch = read_private_json(Path(argv[-1]) / "launch.json") + startup_readers.append(os.dup(int(launch["startup_fd"]))) + return SimpleNamespace(pid=12345) + + popen = MagicMock(side_effect=spawned) monkeypatch.setattr(local.subprocess, "Popen", popen) monkeypatch.setenv("CUDA_VISIBLE_DEVICES", "an-ambient-mask") monkeypatch.setenv("CUDA_DEVICE_ORDER", "changed-after-planning") - identity = provider.launch(plan) + try: + identity = provider.launch(plan) + assert os.read(startup_readers[0], 1) == b"1" + finally: + for descriptor in startup_readers: + os.close(descriptor) assert popen.call_args.kwargs["env"]["CUDA_VISIBLE_DEVICES"] == (mask if gpus else "") assert popen.call_args.kwargs["env"].get("CUDA_DEVICE_ORDER") == order @@ -1042,10 +1239,6 @@ def test_local_runtime_advertises_configured_resources_and_preserves_native_gpu_ directory = private_directory(tmp_path / "allocation", create=True) scratch = private_directory(tmp_path / "scratch", create=True) - write_private_json(directory / "launch.json", { - "identity": "allocation", "deadline": time.monotonic() + 60, - "task_slots": 1, "scratch": str(scratch), - }) record = {"cpus": 1, "memory": 1024**3} if gpus is not None: record["gpus"] = gpus @@ -1066,6 +1259,21 @@ def test_local_runtime_advertises_configured_resources_and_preserves_native_gpu_ cluster.return_value.__enter__.return_value.scheduler.id = "Scheduler-gpu" monkeypatch.setattr(distributed, "LocalCluster", cluster) - local_runtime.main() + with local._allocation_lock() as lock_fd: + startup_read, startup_write = os.pipe() + try: + os.write(startup_write, b"1") + write_private_json(directory / "launch.json", { + "identity": "allocation", "deadline": time.monotonic() + 60, + "task_slots": 1, "scratch": str(scratch), + "lock_fd": lock_fd, "startup_fd": startup_read, + }) + local_runtime.main() + finally: + os.close(startup_write) + try: + os.close(startup_read) + except OSError: + pass # The runtime owns and closes its inherited startup reader. assert cluster.call_args.kwargs["resources"] == {"CPU": 1, "MEMORY": 1024**3, "GPU": gpus or 0} assert os.environ["CUDA_VISIBLE_DEVICES"] == ("3,1" if gpus else "") diff --git a/tests/test_compute_output.py b/tests/test_compute_output.py index 7a9ebf5e..320bec73 100644 --- a/tests/test_compute_output.py +++ b/tests/test_compute_output.py @@ -19,6 +19,8 @@ from lightcone.engine.compute.model import GIB, Request from lightcone.engine.sandbox.boundary import _STDERR_TAIL_BYTES, _Tail, write_output +pytestmark = pytest.mark.usefixtures("local_allocation_lock") + @pytest.fixture def detached_cluster(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Iterator[str]: From a9258eeb9aa6509d1f8d1e31c3ce5bbf757671c0 Mon Sep 17 00:00:00 2001 From: Francois Lanusse Date: Tue, 29 Sep 2026 10:28:38 -0700 Subject: [PATCH 4/4] Find a running local cluster in the process table instead of flock The one-local-cluster-per-user rule was an flock held in the account home for the owner's lifetime, and NERSC home filesystems do not support flock. A launch now scans the process table for a session leader of this user running the owner command, which spans every catalog and connection root with no file lock. The owner command is one shared constant for launch, the scan and the identity check. A refusal still names the running cluster and the catalog to stop it with, now recorded in its identity file. A record not yet written means the owner is starting ("retry shortly"); a missing or unreadable one never recovers, so that refusal names the owner's PID. An unreadable process table is a clean ComputeError. Accepted: overlapping launches can both start, and the scan covers one PID namespace. The suite scopes the scan to its own temporary tree, so a developer's running cluster does not refuse test launches. Co-Authored-By: Claude Opus 5.5 --- CLAUDE.md | 12 +- docs/api/compute.md | 27 ++- docs/cli/compute.md | 5 +- docs/user/cluster.md | 24 +- src/lightcone/engine/compute/local.py | 182 +++++++-------- src/lightcone/engine/compute/local_runtime.py | 2 - tests/conftest.py | 25 +- tests/test_compute_local.py | 219 ++++++++++-------- tests/test_compute_output.py | 2 +- 9 files changed, 267 insertions(+), 231 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 482b2768..7e7b2d10 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -2141,7 +2141,17 @@ unlinks before writing; a new tampering test should too. own offers. Login-node catalogs disable local launch and execution through `local.enabled: false`. Inspection and termination stay available. One local allocation per user per machine is enforced across catalogs and connection - roots by an OS lock held throughout the detached owner's lifetime. + roots by scanning the process table for a live owner before launch + (`local._running_owners`: a session leader of this user running + `_OWNER_ARGS`, the command launch and `_process` share). This replaced an + `flock` held for the owner's lifetime, because NERSC home filesystems do + not support `flock`. Accepted residue: overlapping launches can both start, + and the scan covers one PID namespace, so a container sharing the home + does not see the host's owner. A missing identity record is transient + ("retry shortly"); an unreadable one is not, so that refusal names the PID. + The suite scopes the scan to its own temporary tree + (`local_allocation_scope`), so a developer's running cluster does not refuse + test launches. - **The engine is the host's uv tool, never a project dependency** (2026-08, reversing spec ยง2's engine-in-lock rule and deleting layer 3). diff --git a/docs/api/compute.md b/docs/api/compute.md index 83a2d3e8..88d65532 100644 --- a/docs/api/compute.md +++ b/docs/api/compute.md @@ -221,21 +221,24 @@ command. A driver that exits before every task reports says so with `UNSTOPPED`: closing a client cannot prove that a remote subprocess stopped. Probes preserve both streams; materialization sends recipe output to stderr to leave stdout for its report. -A host-local OS file lock limits local allocations to one per user, independent of -connection roots and namespaces. The launcher acquires it before spawning and passes -the descriptor to the detached owner, which retains it for its lifetime. The lock -resides under the account's home directory at -`.lightcone/local-locks//allocation.lock`, independently of `HOME`, -`TMPDIR`, and catalog paths. Its boot UUID keeps hosts with a shared home separate. -The file is never unlinked; OS process exit releases the lock without stale-lock -cleanup. Advisory metadata identifies the owning cluster, catalog, and connection -root when a competing launch is refused. -The account-home filesystem must support `flock`; the login-node guard alone -does not establish filesystem support on an interactive compute node. +Local allocations are limited to one per user on each machine, independent of +connection roots and namespaces. Before spawning, the launcher scans the process +table for a live owner of the same user: a session leader running `-P -m +lightcone.engine.compute.local_runtime `, which excludes workers forked +from it. The process table spans every catalog and connection root and needs no +file lock, which some shared home filesystems, NERSC's included, do not support. +The owner's directory argument locates its identity record, so a refused launch +names the running cluster, its connection root, and the catalog it was launched +with. A record not yet written means the owner is still starting; one that is +missing or unreadable after that never recovers, so the refusal names the +owner's PID instead. Two limits are accepted rather than closed with a lock: +launches that overlap can both pass the scan, and the scan covers one PID +namespace, so a container sharing the home directory does not see the host's +owner. A startup pipe lets the owner proceed only after the launcher publishes its identity and launch records. If the launcher dies before completing publication, -the pipe closes and the owner exits, releasing the lock. Failures before identity +the pipe closes and the owner exits. Failures before identity publication remove the launcher's private files; a published identity remains inspectable after a startup failure. diff --git a/docs/cli/compute.md b/docs/cli/compute.md index 1a67187e..8b0cd6b4 100644 --- a/docs/cli/compute.md +++ b/docs/cli/compute.md @@ -105,8 +105,9 @@ name only once ready; JSON adds `ready: true`. A timeout or startup failure exit or terminates the accepted allocation. Only one local cluster may run per user on a machine, across names, namespaces, -and configured roots. Concurrent launches are serialized by an OS lifetime lock; -a second launch fails until the existing cluster ends. +and configured roots. A launch that finds one of your local clusters running in +the process table fails until that cluster ends; launches that overlap can both +succeed. `--json` emits versioned (`schema_version: 1`), allowlisted data without scheduler credentials: diff --git a/docs/user/cluster.md b/docs/user/cluster.md index 4c1c3636..be7c4f14 100644 --- a/docs/user/cluster.md +++ b/docs/user/cluster.md @@ -41,13 +41,18 @@ Use `lc compute status NAME` for resource details and Dask readiness. ## Local allocations Local resources are cooperative limits, not an exclusive CPU/RAM reservation. -Only one local cluster can run per user on each machine. An OS lock held by the -detached owner prevents concurrent launches, including through different names, -catalogs, namespaces, or connection roots. End the existing cluster before -launching another. The lock releases when its owner exits, including on failure -or walltime expiry. A refusal identifies the owning cluster and its original -catalog and connection root. Use that catalog to inspect or stop the cluster if -the current catalog no longer includes its connection. +Only one local cluster can run per user on each machine. A launch checks the +process table for a running local cluster of yours and refuses if it finds one, +including one launched through a different name, catalog, namespace, or +connection root. End the existing cluster before launching another; once its +owner process exits, including on failure or walltime expiry, a new launch +proceeds. A refusal identifies the running cluster and its original catalog and +connection root. Use that catalog to inspect or stop the cluster if the current +catalog no longer includes its connection. If the cluster's record is missing or +damaged, the refusal names its process ID instead, to stop with `kill`. +Launches that overlap can both succeed, and the check sees only the processes +visible where `lc` runs: a launch inside a container does not see a cluster +started outside it. An allocation owns a detached process session and standard `LocalCluster`: one worker process with `task_slots_per_node` threads, and a scheduler that listens on `127.0.0.1` over TLS. Its own logs are discarded; a startup failure is kept @@ -71,11 +76,6 @@ sessions. A `SLURM_JOB_ID` variable does not exempt a login node. See NERSC's [environment conventions](https://docs.nersc.gov/environment/) and [interactive sessions](https://docs.nersc.gov/connect/vscode/). -Local startup also requires `flock` support in the account home, where the -singleton lock lives. NERSC documents that its compute-node home mounts do not -support `flock`, so this separate requirement can prevent local startup there; -the Slurm provider does not use that lock. - A local connection's optional `launch` settings are `connection_root` (default `~/.lightcone/compute`), `scratch_root` (default: the temporary directory), `python` (default: the interpreter running `lc`), and `task_slots_per_node` diff --git a/src/lightcone/engine/compute/local.py b/src/lightcone/engine/compute/local.py index a5902757..0f626311 100644 --- a/src/lightcone/engine/compute/local.py +++ b/src/lightcone/engine/compute/local.py @@ -3,15 +3,12 @@ from __future__ import annotations import fcntl -import json import os -import pwd import re import shlex import shutil import signal import socket -import stat import subprocess import sys import tempfile @@ -48,12 +45,10 @@ write_private_json, ) -_OWNER_MODULE = "lightcone.engine.compute.local_runtime" +# The owner's command after its interpreter; launch and both identity checks share it. +_OWNER_ARGS = ("-P", "-m", "lightcone.engine.compute.local_runtime") _STOP_GRACE = 3.0 _RETIRED = "ended.json" -# The account home is stable even when a command overrides HOME. A boot UUID -# separates machines sharing that home without depending on a mutable hostname. -_LOCK_ROOT = Path(pwd.getpwuid(os.getuid()).pw_dir) / ".lightcone" / "local-locks" def _above_stdio(descriptor: int) -> int: @@ -65,75 +60,85 @@ def _above_stdio(descriptor: int) -> int: return descriptor -@contextmanager -def _allocation_lock() -> Iterator[int]: - """Hold a per-user, per-boot lock across catalogs and detached owner lifetime.""" - root = private_directory(configured_directory(_LOCK_ROOT) / _boot_identity(), create=True) - path = root / "allocation.lock" - descriptor: int | None = None +def _running_owners() -> list[tuple[int, Path]]: + """Find this user's live allocation owners in this host's process table. + + An owner leads its own session and runs ``python -P -m + ``. The process table spans every catalog and connection root, + and asking it needs no file lock, which shared home filesystems such as + NERSC's refuse. + + Returns: + Each owner's PID and allocation directory. + + Raises: + ComputeError: If the process table cannot be read at all. + """ try: + processes = list(psutil.process_iter()) + except (psutil.Error, OSError) as exc: + raise ComputeError(f"cannot read this host's process table: {exc}") from exc + owners = [] + for process in processes: try: - descriptor = os.open(path, os.O_CREAT | os.O_RDWR | os.O_NOFOLLOW, 0o600) - descriptor = _above_stdio(descriptor) - info = os.fstat(descriptor) - if (not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid() - or stat.S_IMODE(info.st_mode) != 0o600 or info.st_nlink != 1): - raise ComputeError("local allocation lock must be a private file owned by you") - deadline = time.monotonic() + 0.1 - while True: - try: - fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB) - break - except BlockingIOError as exc: - # A zombie owner can still have exiting threads whose - # descriptor references the kernel has not released yet. - remaining = deadline - time.monotonic() - if remaining > 0: - time.sleep(min(0.01, remaining)) - continue - message = ( - "a local cluster is already running or starting " - "for this user on this machine" - ) - try: - owner = json.loads(os.pread(descriptor, 65536, 0)) - except (OSError, ValueError): - owner = {} - cluster_id = owner.get("identity") if isinstance(owner, dict) else None - if isinstance(cluster_id, str) and cluster_id: - command = "lc compute down " + shlex.quote(cluster_id) - if owner.get("explicit_catalog") is True and isinstance( - owner.get("catalog"), str, - ): - command = ( - "LC_COMPUTE_CONFIG=" + shlex.quote(owner["catalog"]) + " " + command - ) - else: - command = "env -u LC_COMPUTE_CONFIG " + command - message += ( - f"; allocation {cluster_id} uses local namespace " - f"{owner.get('namespace')!r} with connection_root " - f"{owner.get('connection_root')!r} (locator {owner.get('locator')!r}); " - "using that launch catalog, " - f"stop it with `{command}`" - ) - else: - message += ( - "; its launcher is publishing the allocation identity; retry shortly" - ) - raise ComputeError(message) from exc - current = path.stat(follow_symlinks=False) - if (current.st_dev, current.st_ino) != (info.st_dev, info.st_ino): - raise ComputeError("the local allocation lock changed while it was being acquired") - # Stale diagnostics have no authority once the kernel grants a new lock. - os.ftruncate(descriptor, 0) - except OSError as exc: - raise ComputeError(f"cannot acquire the local allocation lock: {exc}") from exc - yield descriptor - finally: - if descriptor is not None: - # Do not unlock: the detached owner inherits this same open file description. - os.close(descriptor) + if process.uids().real != os.getuid(): + continue + argv = process.cmdline() + # An exiting owner has no command line left, and a forked worker + # keeps the owner's command but does not lead its session. + if ( + len(argv) == len(_OWNER_ARGS) + 2 + and tuple(argv[1:-1]) == _OWNER_ARGS + and os.getsid(process.pid) == process.pid + ): + owners.append((process.pid, Path(argv[-1]))) + except (psutil.Error, OSError): + continue + return owners + + +def _refuse_a_second_allocation() -> None: + """Refuse a launch while this user already runs a local allocation here. + + Launches that overlap can both pass; that race is accepted rather than + closed with a lock. + + Raises: + ComputeError: If an owner is running, naming how to stop it. + """ + owners = _running_owners() + if not owners: + return + pid, directory = owners[0] + message = "a local cluster is already running or starting for this user on this machine" + path = directory / "identity.json" + if directory.is_dir() and not path.exists(): + raise ComputeError( + f"{message}; its launcher (owner process {pid}) is publishing the allocation " + "identity; retry shortly" + ) + try: + record = read_private_json(path) + cluster_id = record.get("identity") + if not isinstance(cluster_id, str) or not cluster_id: + raise ComputeError(f"compute connection file has no identity: {path}") + except ComputeError as exc: + # A lost or damaged record never recovers by waiting; the process can still be stopped. + raise ComputeError( + f"{message}; owner process {pid} has no usable allocation record ({exc}); " + f"stop it with `kill {pid}`" + ) from exc + catalog = record.get("catalog") + command = ( + f"LC_COMPUTE_CONFIG={shlex.quote(catalog)}" if isinstance(catalog, str) + else "env -u LC_COMPUTE_CONFIG" + ) + " lc compute down " + shlex.quote(cluster_id) + raise ComputeError( + f"{message}; allocation {cluster_id} uses local namespace {directory.parent.name!r} " + f"with connection_root {str(directory.parent.parent)!r} " + f"(locator {str(directory.parent)!r}); using that launch catalog, stop it with " + f"`{command}`" + ) def _boot_identity() -> str: @@ -238,11 +243,7 @@ def launch(self, plan: LaunchPlan) -> Identity: raise ComputeError("local launch plan belongs to a different connection or node count") if plan.name is not None: validate_name(plan.name) - with _allocation_lock() as lock_fd: - return self._launch_locked(plan, lock_fd) - - def _launch_locked(self, plan: LaunchPlan, lock_fd: int) -> Identity: - """Transfer the singleton lock to the allocation owner before releasing our copy.""" + _refuse_a_second_allocation() boot = _boot_identity() token = uuid4().hex directory: Path | None = None @@ -271,20 +272,19 @@ def _launch_locked(self, plan: LaunchPlan, lock_fd: int) -> Identity: "task_slots": plan.details["task_slots_per_node"], "scratch": str(scratch), "identity": "", - "lock_fd": lock_fd, "startup_fd": startup_read, } write_private_json(directory / "launch.json", launch) # This allocation outlives a command; the ordinary run-to-completion # subprocess seam cannot own it. Logs are discarded rather than grow. process = subprocess.Popen( - [plan.details["python"], "-P", "-m", _OWNER_MODULE, str(directory)], + [plan.details["python"], *_OWNER_ARGS, str(directory)], stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, start_new_session=True, close_fds=True, - pass_fds=(lock_fd, startup_read), + pass_fds=(startup_read,), env=environment, ) os.close(startup_read) @@ -294,6 +294,7 @@ def _launch_locked(self, plan: LaunchPlan, lock_fd: int) -> Identity: host=socket.gethostname(), name=plan.name or "", ) + catalog = os.environ.get("LC_COMPUTE_CONFIG") record = { "identity": identity.encode(), "pid": process.pid, @@ -304,23 +305,13 @@ def _launch_locked(self, plan: LaunchPlan, lock_fd: int) -> Identity: "memory": plan.resources.memory_bytes, "gpus": plan.resources.gpus, "accelerator_name": plan.resources.accelerator_name or "GPU", + # A refused launch names the catalog that reaches this allocation. + "catalog": None if catalog is None else str(Path(catalog).expanduser().absolute()), } write_private_json(directory / "identity.json", record) published = True launch["identity"] = identity.encode() write_private_json(directory / "launch.json", launch) - owner = json.dumps({ - "identity": identity.encode(), - "namespace": self.connection.namespace, - "connection_root": str(self.root.parent), - "locator": str(self.root), - "explicit_catalog": "LC_COMPUTE_CONFIG" in os.environ, - "catalog": str(Path(os.environ.get( - "LC_COMPUTE_CONFIG", "~/.lightcone/compute.yaml", - )).expanduser().absolute()), - }).encode() - if os.pwrite(lock_fd, owner, 0) != len(owner): - raise OSError("cannot publish the local allocation lock owner") # EOF without this byte means the launcher died before publication, # including SIGKILL, which no Python exception handler can clean up. os.write(startup_write, b"1") @@ -407,8 +398,7 @@ def _process( except psutil.TimeoutExpired: raise ComputeError("the local process identity is unavailable") from None if ( - len(argv) != 5 - or argv[1:] != ["-P", "-m", _OWNER_MODULE, str(directory)] + argv[1:] != [*_OWNER_ARGS, str(directory)] or os.getpgid(process.pid) != process.pid or os.getsid(process.pid) != process.pid ): diff --git a/src/lightcone/engine/compute/local_runtime.py b/src/lightcone/engine/compute/local_runtime.py index cf643f06..1feea938 100644 --- a/src/lightcone/engine/compute/local_runtime.py +++ b/src/lightcone/engine/compute/local_runtime.py @@ -26,8 +26,6 @@ def main() -> None: os.umask(0o077) directory = private_directory(Path(sys.argv[1])) launch = read_private_json(directory / "launch.json") - # Keep the inherited lock until process exit, but exclude worker execs. - os.set_inheritable(int(launch["lock_fd"]), False) startup_fd = int(launch["startup_fd"]) os.set_inheritable(startup_fd, False) if os.getsid(0) != os.getpid() or os.getpgrp() != os.getpid(): diff --git a/tests/conftest.py b/tests/conftest.py index e8ade6fd..c6541b4f 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -28,21 +28,24 @@ def runner() -> CliRunner: return CliRunner() -@pytest.fixture(scope="session") -def local_allocation_lock_root(tmp_path_factory: pytest.TempPathFactory) -> Path: - """Give each test session a lock root shared with its independent launchers.""" - return tmp_path_factory.mktemp("local-allocation-locks") - - @pytest.fixture -def local_allocation_lock( - local_allocation_lock_root: Path, monkeypatch: pytest.MonkeyPatch, -) -> None: - """Exercise the lifetime lock without contending with other test sessions.""" +def local_allocation_scope( + tmp_path_factory: pytest.TempPathFactory, monkeypatch: pytest.MonkeyPatch, +) -> Path: + """Count only this session's local allocations, never the developer's or another suite's. + + Returns: + The session's base temporary directory, for independent launchers to scope by. + """ from lightcone.engine.compute import local - monkeypatch.setattr(local, "_LOCK_ROOT", local_allocation_lock_root) + root = tmp_path_factory.getbasetemp() + owners = local._running_owners + monkeypatch.setattr( + local, "_running_owners", lambda: [o for o in owners() if o[1].is_relative_to(root)], + ) monkeypatch.delenv("NERSC_HOST", raising=False) + return root @pytest.fixture(autouse=True) diff --git a/tests/test_compute_local.py b/tests/test_compute_local.py index 6d0e1ce5..9c6fb90d 100644 --- a/tests/test_compute_local.py +++ b/tests/test_compute_local.py @@ -44,7 +44,7 @@ write_private_json, ) -pytestmark = pytest.mark.usefixtures("local_allocation_lock") +pytestmark = pytest.mark.usefixtures("local_allocation_scope") @pytest.fixture @@ -72,86 +72,124 @@ def _launch(provider: LocalProvider, *, seconds: int = 60) -> Identity: def test_singleton_survives_cli_exit_and_spans_names_and_connection_roots( - provider: LocalProvider, tmp_path: Path, + provider: LocalProvider, tmp_path: Path, local_allocation_scope: Path, ) -> None: - # Independent launchers compete for the same host/user lock, with distinct catalogs. + # An independent launcher counts only this session's owners, like the fixture. script = """ -import json, sys -from pathlib import Path +import sys from lightcone.engine.compute import local -from lightcone.engine.compute.local import LocalProvider -from lightcone.engine.compute.model import ( - Connection, ComputeError, Offer, Request, Resources, TimeLimits, -) -local._LOCK_ROOT = Path(sys.argv[3]) -p = LocalProvider(Connection.model_validate_json(sys.argv[1])) +from lightcone.engine.compute.model import Connection, Offer, Request, Resources, TimeLimits +owners = local._running_owners +local._running_owners = lambda: [o for o in owners() if o[1].is_relative_to(sys.argv[2])] +p = local.LocalProvider(Connection.model_validate_json(sys.argv[1])) offer = Offer(name='small', connection='workstation', resources=Resources(cpus=1, memory_gib=0.5), max_nodes=1, time=TimeLimits(default='1m', max='1m')) -plan = p.plan(offer, Request(cpus=1, memory_bytes=512 * 1024**2)).replace(name=sys.argv[2]) -input() -try: - print(json.dumps({'id': p.launch(plan).encode()})) -except ComputeError as exc: - print(json.dumps({'error': str(exc)})) +plan = p.plan(offer, Request(cpus=1, memory_bytes=512 * 1024**2)).replace(name='first') +print(p.launch(plan).encode()) """ + launched = subprocess.run( + [sys.executable, "-c", script, provider.connection.model_dump_json(), + str(local_allocation_scope)], + capture_output=True, text=True, timeout=20, check=True, + ) + identity = Identity.decode(launched.stdout.strip()) other = LocalProvider(provider.connection.replace( namespace=str(uuid4()), launch={"connection_root": str(tmp_path / "other")}, )) - providers = [provider, other] - processes = [subprocess.Popen( - [sys.executable, "-c", script, item.connection.model_dump_json(), - f"local-{index}", str(local._LOCK_ROOT)], - stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, - ) for index, item in enumerate(providers)] - identities: list[tuple[LocalProvider, Identity]] = [] + identities = [(provider, identity)] try: - for process in processes: - assert process.stdin is not None - process.stdin.write("go\n") - process.stdin.flush() - results = [] - for item, process in zip(providers, processes): - stdout, stderr = process.communicate(timeout=20) - assert process.returncode == 0, stderr - result = json.loads(stdout) - results.append(result) - if "id" in result: - identities.append((item, Identity.decode(result["id"]))) - assert len(identities) == 1, results - assert "already running or starting" in next(r["error"] for r in results if "error" in r) - owner, identity = identities[0] - _ready(owner, identity) + _ready(provider, identity) + # The owner is found by its command, never by its Dask worker processes. + assert local._running_owners() == [ + (int(identity.native_id), provider.root / identity.token), + ] with pytest.raises(ComputeError, match="already running") as conflict: _launch(other) assert identity.encode() in str(conflict.value) assert conflict.value.cluster_id is None - assert str(owner.root) in str(conflict.value) - owner.terminate(identity) - _ended(owner, identity) + assert str(provider.root) in str(conflict.value) + provider.terminate(identity) + _ended(provider, identity) replacement = _launch(other) identities.append((other, replacement)) _ready(other, replacement) finally: - for process in processes: - if process.poll() is None: - process.kill() - process.wait(timeout=5) - for item, identity in identities: - item.terminate(identity) + for item, allocation in identities: + item.terminate(allocation) + + +def test_a_starting_owner_refuses_a_launch_until_its_identity_is_published( + provider: LocalProvider, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + directory = private_directory(tmp_path / "allocation", create=True) + monkeypatch.setattr(local, "_running_owners", lambda: [(4242, directory)]) + with pytest.raises(ComputeError, match="publishing the allocation identity; retry"): + _launch(provider) + assert not provider.root.exists() -def test_allocation_lock_preserves_errors_from_the_launch_body() -> None: - error = OSError("allocation storage is full") - with pytest.raises(OSError) as raised, local._allocation_lock(): - raise error - assert raised.value is error - with local._allocation_lock(): - pass +@pytest.mark.parametrize("record", ["missing", "corrupt", "anonymous"]) +def test_an_owner_without_a_usable_record_is_named_by_its_process( + provider: LocalProvider, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, record: str, +) -> None: + directory = tmp_path / "allocation" + if record != "missing": + path = private_directory(directory, create=True) / "identity.json" + if record == "corrupt": + path.write_text("{") + path.chmod(0o600) + else: + write_private_json(path, {"pid": 4242}) + monkeypatch.setattr(local, "_running_owners", lambda: [(4242, directory)]) + # Waiting never repairs a lost record, so the refusal names the process instead. + with pytest.raises(ComputeError, match="no usable allocation record") as refused: + _launch(provider) + assert "retry" not in str(refused.value) + assert "`kill 4242`" in str(refused.value) + + +def test_the_owner_scan_matches_only_session_leaders_of_this_user( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + owner = [sys.executable, "-P", "-m", "lightcone.engine.compute.local_runtime", str(tmp_path)] + table: list[tuple[int, int, list[str] | Exception]] = [ + (101, os.getuid(), owner), + (102, os.getuid(), owner), # a worker forked from the owner: not a session leader + (103, os.getuid() + 1, owner), # another user's owner + (104, os.getuid(), [sys.executable, "-X", "dev", *owner[1:]]), + (105, os.getuid(), []), # an exiting owner + (106, os.getuid(), psutil.AccessDenied(106)), + (107, os.getuid(), psutil.NoSuchProcess(107)), + ] + + def process(pid: int, uid: int, argv: list[str] | Exception) -> SimpleNamespace: + def cmdline() -> list[str]: + if isinstance(argv, Exception): + raise argv + return argv + + return SimpleNamespace(pid=pid, uids=lambda: SimpleNamespace(real=uid), cmdline=cmdline) + + monkeypatch.setattr(local.psutil, "process_iter", lambda: [process(*row) for row in table]) + monkeypatch.setattr(local.os, "getsid", lambda pid: pid if pid != 102 else 101) + assert local._running_owners() == [(101, tmp_path)] + + +def test_an_unreadable_process_table_refuses_the_launch_cleanly( + provider: LocalProvider, monkeypatch: pytest.MonkeyPatch, +) -> None: + def unreadable() -> None: + raise FileNotFoundError("/proc is not mounted") + + monkeypatch.setattr(local.psutil, "process_iter", unreadable) + with pytest.raises(ComputeError, match="cannot read this host's process table"): + _launch(provider) @pytest.mark.parametrize("closed", [(0,), (1,), (2,), (0, 1, 2)]) -def test_singleton_lock_survives_launch_with_closed_standard_descriptors( - provider: LocalProvider, tmp_path: Path, closed: tuple[int, ...], +def test_launch_survives_closed_standard_descriptors( + provider: LocalProvider, tmp_path: Path, local_allocation_scope: Path, + closed: tuple[int, ...], ) -> None: result_path = tmp_path / "launched.json" script = """ @@ -159,7 +197,8 @@ def test_singleton_lock_survives_launch_with_closed_standard_descriptors( from pathlib import Path from lightcone.engine.compute import local from lightcone.engine.compute.model import Connection, Offer, Request, Resources, TimeLimits -local._LOCK_ROOT = Path(sys.argv[2]) +owners = local._running_owners +local._running_owners = lambda: [o for o in owners() if o[1].is_relative_to(sys.argv[2])] p = local.LocalProvider(Connection.model_validate_json(sys.argv[1])) offer = Offer(name='small', connection='workstation', resources=Resources(cpus=1, memory_gib=0.5), max_nodes=1, time=TimeLimits(default='1m', max='1m')) @@ -172,7 +211,7 @@ def test_singleton_lock_survives_launch_with_closed_standard_descriptors( try: result = subprocess.run( [sys.executable, "-c", script, provider.connection.model_dump_json(), - str(local._LOCK_ROOT), str(result_path), json.dumps(closed)], + str(local_allocation_scope), str(result_path), json.dumps(closed)], capture_output=True, text=True, timeout=15, ) assert result.returncode == 0, result.stderr @@ -228,14 +267,16 @@ def spawn() -> int: return psutil.Process(client.submit(spawn).result(timeout=5)) -def test_allocation_survives_launcher_and_borrowed_client_exit(provider: LocalProvider) -> None: +def test_allocation_survives_launcher_and_borrowed_client_exit( + provider: LocalProvider, local_allocation_scope: Path, +) -> None: script = """ import json, sys -from pathlib import Path from lightcone.engine.compute import local from lightcone.engine.compute.local import LocalProvider from lightcone.engine.compute.model import Connection, Offer, Request, Resources, TimeLimits -local._LOCK_ROOT = Path(sys.argv[2]) +owners = local._running_owners +local._running_owners = lambda: [o for o in owners() if o[1].is_relative_to(sys.argv[2])] p = LocalProvider(Connection(**json.loads(sys.argv[1]))) offer = Offer( name='small', connection='workstation', resources=Resources(cpus=1, memory_gib=0.5), @@ -245,7 +286,7 @@ def test_allocation_survives_launcher_and_borrowed_client_exit(provider: LocalPr """ launched = subprocess.run( [sys.executable, "-c", script, json.dumps(provider.connection.model_dump()), - str(local._LOCK_ROOT)], + str(local_allocation_scope)], check=True, capture_output=True, text=True, @@ -456,7 +497,7 @@ def test_walltime_expires_without_a_connected_client(provider: LocalProvider) -> identities = [identity] try: _ended(provider, identity, timeout=6) - # Native exit and release of the owner's last kernel descriptor can differ. + # An owner ended by its walltime no longer refuses the next launch. identities.append(_launch(provider)) finally: for allocation in identities: @@ -646,7 +687,6 @@ def fail(argv: list[str], **kwargs: Any) -> subprocess.Popen[bytes]: _launch(provider) assert list(provider.root.iterdir()) == [] assert list(Path(provider.connection.launch["scratch_root"]).iterdir()) == [] - # A failed spawn must release the singleton lock as well as its private files. identity = _launch(provider) try: _ready(provider, identity) @@ -659,7 +699,7 @@ def fail(argv: list[str], **kwargs: Any) -> subprocess.Popen[bytes]: provider.terminate(identity) -def test_failed_initial_launch_write_removes_unpublished_files_and_releases_lock( +def test_failed_initial_launch_write_removes_unpublished_files( provider: LocalProvider, monkeypatch: pytest.MonkeyPatch, ) -> None: def fail(path: Path, value: dict[str, Any]) -> None: @@ -672,11 +712,9 @@ def fail(path: Path, value: dict[str, Any]) -> None: _launch(provider) assert list(provider.root.iterdir()) == [] assert list(Path(provider.connection.launch["scratch_root"]).iterdir()) == [] - with local._allocation_lock(): - pass -def test_interrupted_launch_kills_the_unreturned_owner_and_releases_lock( +def test_interrupted_launch_kills_the_unreturned_owner( provider: LocalProvider, monkeypatch: pytest.MonkeyPatch, ) -> None: children: list[subprocess.Popen[bytes]] = [] @@ -701,8 +739,7 @@ def interrupt(path: Path, value: dict[str, Any]) -> None: assert len(children) == 1 assert children[0].poll() is not None assert provider.discover() == [] - with local._allocation_lock(): - pass + assert local._running_owners() == [] finally: for child in children: if child.poll() is None: @@ -712,7 +749,7 @@ def interrupt(path: Path, value: dict[str, Any]) -> None: @pytest.mark.parametrize("publication", ["identity.json", "launch.json"]) def test_owner_exits_when_launcher_is_killed_before_startup_commit( - provider: LocalProvider, tmp_path: Path, publication: str, + provider: LocalProvider, tmp_path: Path, local_allocation_scope: Path, publication: str, ) -> None: marker = tmp_path / "paused.json" script = """ @@ -720,7 +757,8 @@ def test_owner_exits_when_launcher_is_killed_before_startup_commit( from pathlib import Path from lightcone.engine.compute import local from lightcone.engine.compute.model import Connection, Offer, Request, Resources, TimeLimits -local._LOCK_ROOT = Path(sys.argv[2]) +owners = local._running_owners +local._running_owners = lambda: [o for o in owners() if o[1].is_relative_to(sys.argv[2])] write = local.write_private_json def pause(path, value): if path.name == sys.argv[4] and (path.name != 'launch.json' or value['identity']): @@ -741,7 +779,7 @@ def pause(path, value): """ launcher = subprocess.Popen( [sys.executable, "-c", script, provider.connection.model_dump_json(), - str(local._LOCK_ROOT), str(marker), publication], + str(local_allocation_scope), str(marker), publication], stdout=subprocess.PIPE, stderr=subprocess.PIPE, ) owner: psutil.Process | None = None @@ -758,11 +796,10 @@ def pause(path, value): launcher.communicate(timeout=5) deadline = time.monotonic() + 5 while owner.is_running() and owner.status() != psutil.STATUS_ZOMBIE: - assert time.monotonic() < deadline, "unpublished owner retained the singleton lock" + assert time.monotonic() < deadline, "unpublished owner kept running" time.sleep(0.05) assert not (Path(checkpoint["directory"]) / "connection.json").exists() - with local._allocation_lock(): - pass + assert local._running_owners() == [] finally: if launcher.poll() is None: launcher.kill() @@ -1259,21 +1296,15 @@ def test_local_runtime_advertises_configured_resources_and_preserves_native_gpu_ cluster.return_value.__enter__.return_value.scheduler.id = "Scheduler-gpu" monkeypatch.setattr(distributed, "LocalCluster", cluster) - with local._allocation_lock() as lock_fd: - startup_read, startup_write = os.pipe() - try: - os.write(startup_write, b"1") - write_private_json(directory / "launch.json", { - "identity": "allocation", "deadline": time.monotonic() + 60, - "task_slots": 1, "scratch": str(scratch), - "lock_fd": lock_fd, "startup_fd": startup_read, - }) - local_runtime.main() - finally: - os.close(startup_write) - try: - os.close(startup_read) - except OSError: - pass # The runtime owns and closes its inherited startup reader. + startup_read, startup_write = os.pipe() + os.write(startup_write, b"1") + os.close(startup_write) + write_private_json(directory / "launch.json", { + "identity": "allocation", "deadline": time.monotonic() + 60, + "task_slots": 1, "scratch": str(scratch), "startup_fd": startup_read, + }) + # The runtime owns and closes its startup reader; closing it here again + # could close whatever descriptor reused the number since. + local_runtime.main() assert cluster.call_args.kwargs["resources"] == {"CPU": 1, "MEMORY": 1024**3, "GPU": gpus or 0} assert os.environ["CUDA_VISIBLE_DEVICES"] == ("3,1" if gpus else "") diff --git a/tests/test_compute_output.py b/tests/test_compute_output.py index 320bec73..407929b1 100644 --- a/tests/test_compute_output.py +++ b/tests/test_compute_output.py @@ -19,7 +19,7 @@ from lightcone.engine.compute.model import GIB, Request from lightcone.engine.sandbox.boundary import _STDERR_TAIL_BYTES, _Tail, write_output -pytestmark = pytest.mark.usefixtures("local_allocation_lock") +pytestmark = pytest.mark.usefixtures("local_allocation_scope") @pytest.fixture