diff --git a/CLAUDE.md b/CLAUDE.md index 23386042..7e7b2d10 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,33 @@ 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 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). `lc init` scaffolds no `lightcone-cli` dependency, and there is no 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..88d65532 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,27 @@ 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. +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. +`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. @@ -136,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 @@ -208,6 +221,27 @@ 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. +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. 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 session and the exact command containing a random allocation token establish diff --git a/docs/architecture.md b/docs/architecture.md index 4c47acf4..f5fdd193 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -169,8 +169,11 @@ 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. 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 b082434b..8b0cd6b4 100644 --- a/docs/cli/compute.md +++ b/docs/cli/compute.md @@ -5,45 +5,54 @@ 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. + +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. | | `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 +65,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 +98,23 @@ 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. 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: | 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..be7c4f14 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,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. 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 @@ -57,9 +67,14 @@ 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; -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/). A local connection's optional `launch` settings are `connection_root` (default `~/.lightcone/compute`), `scratch_root` (default: the temporary directory), @@ -69,12 +84,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 ``` @@ -90,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 @@ -97,12 +115,54 @@ 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: + +```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. + +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 +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. +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 +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`. + +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. -For example, this catalog exposes a larger local allocation: +For example, this catalog explicitly defines a local allocation: ```yaml version: 1 @@ -121,16 +181,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 +238,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..7e2e5df6 100644 --- a/docs/user/getting-started.md +++ b/docs/user/getting-started.md @@ -250,13 +250,14 @@ 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 -walkthrough. If you already have a compute catalog, its offers replace that default; -see [Running on a Cluster](cluster.md). +all usable CPUs and RAM for 30 minutes. Keep the returned ID in `CLUSTER` for this +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 --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/evals/prompt.md b/evals/prompt.md index 6b0c9e3c..d905996e 100644 --- a/evals/prompt.md +++ b/evals/prompt.md @@ -38,16 +38,36 @@ 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, 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. +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. +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 aa139375..01d221d4 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 @@ -58,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. """ @@ -96,9 +99,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 +113,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 +164,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..694b7998 100644 --- a/src/lightcone/engine/compute/__init__.py +++ b/src/lightcone/engine/compute/__init__.py @@ -9,12 +9,13 @@ from typing import Any from uuid import uuid4 -from .catalog import Catalog +from .catalog import Catalog, local_disabled_reason from .model import ( ComputeError, 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 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 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] = [] + 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 ( + 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() @@ -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 ( + 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 a9b92aa8..efd23fb8 100644 --- a/src/lightcone/engine/compute/catalog.py +++ b/src/lightcone/engine/compute/catalog.py @@ -3,7 +3,8 @@ from __future__ import annotations import os -from decimal import Decimal +import re +import socket from pathlib import Path from typing import Annotated, Any, Self @@ -26,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.""" @@ -43,12 +58,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 +104,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 +119,62 @@ 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) - except ValidationError as exc: + catalog = cls.model_validate(raw) + return catalog._with_local() + 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 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; " + "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 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 = 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 local.enabled: + offers = [ + offer for offer in offers if connections[offer.connection].provider != "local" + ] + 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.""" for connection in self.connections.values(): diff --git a/src/lightcone/engine/compute/local.py b/src/lightcone/engine/compute/local.py index 85b05272..0f626311 100644 --- a/src/lightcone/engine/compute/local.py +++ b/src/lightcone/engine/compute/local.py @@ -2,8 +2,10 @@ from __future__ import annotations +import fcntl import os import re +import shlex import shutil import signal import socket @@ -19,6 +21,7 @@ import psutil +from lightcone.engine.compute.catalog import local_disabled_reason from lightcone.engine.compute.model import ( ComputeError, Connection, @@ -42,11 +45,102 @@ 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" +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 + + +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: + 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: try: if sys.platform == "linux": @@ -75,6 +169,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: @@ -141,48 +237,64 @@ 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) + _refuse_a_second_allocation() 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": "", - }, - ) 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": "", + "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=(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(), name=plan.name or "", ) + catalog = os.environ.get("LC_COMPUTE_CONFIG") record = { "identity": identity.encode(), "pid": process.pid, @@ -193,20 +305,18 @@ def launch(self, plan: LaunchPlan) -> 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) - # 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(), - }, - ) + published = True + launch["identity"] = identity.encode() + write_private_json(directory / "launch.json", launch) + # 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. @@ -215,21 +325,24 @@ def launch(self, plan: LaunchPlan) -> 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 ( @@ -285,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 1ca6d552..1feea938 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,6 +26,8 @@ def main() -> None: os.umask(0o077) directory = private_directory(Path(sys.argv[1])) launch = read_private_json(directory / "launch.json") + 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() @@ -46,10 +49,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..c6541b4f 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -28,6 +28,26 @@ def runner() -> CliRunner: return CliRunner() +@pytest.fixture +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 + + 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) 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 29dfead0..16868ed7 100644 --- a/tests/test_compute.py +++ b/tests/test_compute.py @@ -53,16 +53,18 @@ 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 @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 +253,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 +295,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 +303,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 +478,240 @@ 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_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: + 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() + + +@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: + 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_local_execution_but_allows_status_and_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() + assert service.status(IDENTITY.encode()).ready + 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 +1088,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 +1177,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..9c6fb90d 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 @@ -32,6 +34,7 @@ Request, Resources, TimeLimits, + UnavailableOfferError, ) from lightcone.engine.compute.runtime import ( configured_directory, @@ -41,6 +44,8 @@ write_private_json, ) +pytestmark = pytest.mark.usefixtures("local_allocation_scope") + @pytest.fixture def provider(tmp_path: Path) -> LocalProvider: @@ -66,6 +71,159 @@ 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, local_allocation_scope: Path, +) -> None: + # An independent launcher counts only this session's owners, like the fixture. + script = """ +import sys +from lightcone.engine.compute import local +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='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")}, + )) + identities = [(provider, identity)] + try: + _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(provider.root) in str(conflict.value) + provider.terminate(identity) + _ended(provider, identity) + replacement = _launch(other) + identities.append((other, replacement)) + _ready(other, replacement) + finally: + 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() + + +@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_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 = """ +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 +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)) +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_allocation_scope), 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: @@ -109,11 +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 lightcone.engine.compute import local from lightcone.engine.compute.local import LocalProvider 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 = LocalProvider(Connection(**json.loads(sys.argv[1]))) offer = Offer( name='small', connection='workstation', resources=Resources(cpus=1, memory_gib=0.5), @@ -122,7 +285,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_allocation_scope)], check=True, capture_output=True, text=True, @@ -184,44 +348,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) - 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) 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() == ([], {}) @@ -306,6 +470,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") @@ -323,10 +494,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) + # An owner ended by its walltime no longer refuses the next launch. + identities.append(_launch(provider)) finally: - provider.terminate(identity) + for allocation in identities: + provider.terminate(allocation) @pytest.mark.parametrize("ending", ["down", "walltime"]) @@ -500,22 +675,21 @@ 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) + assert list(provider.root.iterdir()) == [] + assert list(Path(provider.connection.launch["scratch_root"]).iterdir()) == [] 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) @@ -525,6 +699,117 @@ def fail(argv: list[str], **kwargs: Any) -> subprocess.Popen[bytes]: provider.terminate(identity) +def test_failed_initial_launch_write_removes_unpublished_files( + 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()) == [] + + +def test_interrupted_launch_kills_the_unreturned_owner( + 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() == [] + assert local._running_owners() == [] + 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, local_allocation_scope: 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 +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']): + 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_allocation_scope), 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 kept running" + time.sleep(0.05) + assert not (Path(checkpoint["directory"]) / "connection.json").exists() + assert local._running_owners() == [] + 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: @@ -842,18 +1127,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() @@ -884,13 +1174,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 @@ -973,10 +1276,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 @@ -997,6 +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) + 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 7a9ebf5e..407929b1 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_scope") + @pytest.fixture def detached_cluster(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Iterator[str]: