diff --git a/CLAUDE.md b/CLAUDE.md index 7e7b2d10..cf7933d9 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -1726,14 +1726,15 @@ 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.** 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. +logical CPUs and RAM, one node, fast startup, no walltime and a 30-minute idle +timeout. 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. +and `local.time` override the built-in CPU/RAM budget and time limits 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 @@ -2134,6 +2135,34 @@ unlinks before writing; a new tampering test should too. override flag is needed. Local compute in interactive compute-node sessions remains available, subject to the configured local policy. +- **Local compute ends when idle, not at a fixed age (2026-09, issue #233).** + An offer's `time` is `{default?, max?, idle?}` and needs a `default` or an + `idle`; the built-in local offer is `{idle: 30m}`, so a long recipe finishes + and the cluster stops 30 minutes after the last task. `idle` is handed to the + scheduler as Dask's own `idle_timeout` — its activity test (running, queued or + unrunnable tasks and any transition reset it; clients and `scheduler_info` + polls do not, measured) rather than a tracker of ours. `--time` stays a hard + walltime (SIGALRM, unconditional on activity); with both, the first to fire + ends the allocation, and the built-in offer has no `max` because a ceiling on + `--time` is meaningless when omitting it means unbounded. When the scheduler + closes without the owner asking (no SIGTERM yet), a `SchedulerPlugin.close` + hook records the reason in `error.json` and SIGKILLs the session itself, like + the walltime path — from the hook, so a scheduler that idles out before + `LocalCluster(...)` returns still ends the owner, and never through + `LocalCluster`'s own close, which waits ~34 s on the departed scheduler + (measured), holding the one-per-machine slot and the name. `SCHEDULER_CONFIG` + pins `idle-timeout: None`, so only an offer sets one — ambient Dask config + reaches neither local nor Slurm schedulers. `compute.connect` (execution only; + `status` connects through the provider) submits one no-op task, so a + driver's preparation — annex fetch, image build, sync — starts with a full + countdown. Slurm refuses `time.idle` and still needs `time.default`: its + allocations end at the native walltime, and an ignored idle timeout would be + a lie. Accepted residue: a preparation longer than the timeout still loses + the cluster; a driver pausing between tasks (a long annex commit) counts as + idle; and without `--time` nothing bounds an owner whose scheduler loop + wedges — the walltime's SIGALRM was that bound, and a second timer only for + it was judged not worth its code. + - **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. diff --git a/docs/api/compute.md b/docs/api/compute.md index 88d65532..76704d08 100644 --- a/docs/api/compute.md +++ b/docs/api/compute.md @@ -13,16 +13,17 @@ It owns no service, registry, or saved current-cluster selection. | `Compute.discover()` | Snapshots and per-connection errors, querying each authority once. | | `Compute.status(cluster_id, wait=False, timeout=300)` | Resolve a name or full ID; return native allocation state plus authenticated Dask readiness. Waiting backs off from one to 30 seconds between native queries. | | `Compute.down(cluster_id)` | Resolve a name or full ID, request native termination independent of scheduler health, and return the canonical `Identity`. | -| `connect(cluster_id, timeout=10)` | Resolve a name or full ID; borrow a standard Dask client, closing the client but never the allocation. | +| `connect(cluster_id, timeout=10)` | Resolve a name or full ID; borrow a standard Dask client, closing the client but never the allocation. Submits one no-op task, so a caller's preparation restarts the idle countdown. | | `Provider` | `plan`, `launch`, `discover`, `inspect`, `connect`, `terminate`. | `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; +offer uses detected usable CPUs and RAM, one node, fast startup, and no walltime: +it ends after 30 minutes without task activity. `local.resources` overrides its +CPU/RAM budget and `local.time` its time limits; `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. +cannot be combined with `local.resources` or `local.time`. 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]+`. @@ -48,8 +49,12 @@ duplicate and non-string mapping keys before model validation. Provider-specific Units are explicit. `Resources.memory_gib` stores exact decimal GiB (the YAML key is `memory`), and `memory_bytes` derives an exact integer. Native observations use `Resources.from_bytes(...)`; requests store `Request.memory_bytes`. `TimeLimits` -keeps the configured `default` and `max` duration strings and exposes -`default_seconds` and `max_seconds`. `Startup.class_` corresponds to YAML `class`. +keeps the configured `default`, `max`, and `idle` duration strings, each optional +but requiring a `default` or an `idle`, and exposes `default_seconds`, +`max_seconds`, and `idle_seconds` (`None` when unset). A `LaunchPlan` carries the +resolved hard walltime as `seconds` and derives `idle_seconds` from its offer; +either may be `None` for a local plan, while Slurm plans always have `seconds` and +never `idle_seconds`. `Startup.class_` corresponds to YAML `class`. Connection names exist only as catalog mapping keys, referenced by `Offer.connection`. Compute memory accepts bare GiB quantities and SkyPilot-style binary units: @@ -257,6 +262,7 @@ management; callers must respect the documented execution limits. Containers managed outside that process group can survive local teardown. Tests cover deterministic selection, malformed identities and catalogs, partial -native failures, acceptance ambiguity, PID reuse, detached local lifetime, standard +native failures, acceptance ambiguity, PID reuse, detached local walltime and idle +expiry, standard Dask bootstrap, and explicit execution through borrowed clients. Slurm command contracts are simulated; a real NERSC submission remains a deployment check. diff --git a/docs/cli/compute.md b/docs/cli/compute.md index 8b0cd6b4..5486dc4f 100644 --- a/docs/cli/compute.md +++ b/docs/cli/compute.md @@ -13,16 +13,17 @@ lc compute down CLUSTER [--json] 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. +The built-in offer has one node, fast startup, and no fixed lifetime: it ends +after 30 minutes without task activity. `--name` overrides the name, and `--time` +adds a hard lifetime that ends the cluster even while work is running. 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; +can override the built-in CPU/RAM budget or time limits, 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 @@ -36,7 +37,7 @@ See [local allocations](../user/cluster.md#local-allocations) for detection deta | Command | Behavior | |---|---| -| `resources` | Ordered available offers, per-node shape, node limit, default/maximum time, and startup class. Free capacity remains unknown. | +| `resources` | Ordered available offers, per-node shape, node limit, default/maximum walltime, idle timeout, 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. | @@ -83,11 +84,28 @@ maintain SkyPilot's accelerator alias registry: use the labels configured in Time accepts positive durations with day/hour/minute/second units, such as `30m`, `1h30m`, or `45s`. Without -`--time`, the chosen offer's default applies. `fast` is a service class, not a -queue-time promise. Limits apply to each allocation; aggregate quotas remain -with the native backend. +`--time`, the chosen offer's default walltime applies, if it has one. `fast` is a +service class, not a queue-time promise. Limits apply to each allocation; aggregate +quotas remain with the native backend. + +A local allocation ends at its walltime, after its idle timeout, or at whichever +comes first when it has both. The idle timeout is Dask's scheduler +`idle-timeout`: running or queued tasks keep the allocation alive, and new work +restarts the countdown; connected clients and `status` queries do not. `lc run` +and `lc materialize` restart it when they connect, so their preparation (fetching +inputs, building the image, syncing the environment) starts with the full timeout. +When it expires, the allocation ends, `status` gives that as the reason, and both +its name and this machine's one local allocation are free again. A walltime ends +the allocation even during active work. `down` still ends it at once. + +For Slurm, time is a finite native `--time` request, so a Slurm offer needs a +`time.default` and cannot declare `time.idle`: -For Slurm, time is a finite native `--time` request. Slurm's overtime and +```text +Error: Slurm allocations end at their native walltime: set the offer's time.default and remove time.idle +``` + +Slurm's overtime and termination-grace policy determines actual expiry and can allow unlimited overrun; Lightcone supplies no independent Slurm runtime deadline. A partition is passed only when explicitly set in the offer's configuration. diff --git a/docs/user/cluster.md b/docs/user/cluster.md index be7c4f14..32e8e7ea 100644 --- a/docs/user/cluster.md +++ b/docs/user/cluster.md @@ -9,8 +9,12 @@ present. `lc materialize --check` and `lc status` remain local project inspectio 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). +A local cluster has no fixed lifetime: it ends after 30 minutes without task +activity, so a two-hour recipe finishes normally and the cluster stops 30 minutes +later if no further work arrives. Running or queued tasks keep it alive; connecting +a client or checking its status does not. Use `--time` to add a hard lifetime, which +ends the cluster even while work is running, or `lc compute down` to stop it now. +GPUs require explicit offers; see [GPU allocations](#gpu-allocations). ```bash lc compute resources @@ -45,8 +49,8 @@ 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 +owner process exits, including on failure, idle expiry, 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`. @@ -55,9 +59,13 @@ 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 -and shown as the reason by `lc compute status`. At its time limit the whole -process session is killed with SIGKILL, so a recipe still running stops mid-write. +on `127.0.0.1` over TLS. Its own logs are discarded; a startup failure, or the +scheduler closing after its idle timeout, is kept and shown as the reason by +`lc compute status`. At its walltime the whole process session is killed with +SIGKILL, so a recipe still running stops mid-write. The idle timeout ends the +session the same way. Dask tracks tasks, not processes: after an interrupted +`lc run` or `lc materialize`, a recipe can keep running once its task is gone, +and the idle timeout stops it too. `down` sends SIGTERM, waits three seconds, then sends SIGKILL. Private process locators are checked against the native boot UUID, UID, process session, and exact command containing the allocation's random token before @@ -128,6 +136,15 @@ 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. +`local.time` replaces the default offer's time limits, for example a longer idle +timeout with a ceiling on `--time`: + +```yaml +version: 1 +local: + time: {idle: 1h, max: 8h} +``` + NERSC login nodes are guarded without setup. For other sites, or to disable local compute on every node using the catalog, set: @@ -151,7 +168,7 @@ 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 +and `local.time`, and size and time 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 @@ -175,7 +192,7 @@ offers: connection: workstation resources: {cpus: 4, memory: 8} max_nodes: 1 - time: {default: 30m, max: 2h} + time: {idle: 30m} startup: {class: fast} ``` @@ -196,8 +213,11 @@ and an ordered `offers` list. Connections and offers default to empty: optional `context`, and optional provider `launch` settings. Namespaces must be unique, and so must each provider/`context` pair. - An offer has a unique `name`, the `connection` it uses, per-node `resources` - (`cpus`, `memory`, and optional `accelerators`), `max_nodes`, and `time` with a - `default` no longer than its `max`. `startup` is optional (`fast`, `batch`, or the default + (`cpus`, `memory`, and optional `accelerators`), `max_nodes`, and `time`. `time` + holds an optional walltime `default`, no longer than an optional `max`, and an + optional `idle` timeout; it needs a `default` or an `idle`, so every allocation + can end. Local offers end at whichever limit comes first. Slurm offers need a + `default` and refuse `idle`. `startup` is optional (`fast`, `batch`, or the default `unknown`), written either as a bare class or as `{class: …, source: …}`. `config` holds provider-specific settings. @@ -420,7 +440,7 @@ For local GPUs, add an offer to the [workstation catalog above](#customize-resou connection: workstation resources: {cpus: 4, memory: 8GB, accelerators: 'GPU:1'} max_nodes: 1 - time: {default: 30m, max: 2h} + time: {idle: 30m} startup: fast ``` @@ -513,7 +533,7 @@ CUDA mask is empty. `lc run` reserves the worker's entire CPU, memory, and GPU budgets; direct and podman-hpc probes inherit that allocation mask. Recipe `time_limit` is not supported and is refused before preparation or -execution. Set the allocation lifetime with `lc compute launch --time` instead. +execution. Bound the allocation instead, with `lc compute launch --time`. Fractional CPU/GPU counts, GPU model requests inside a recipe, and disk requests are also rejected rather than ignored. diff --git a/docs/user/getting-started.md b/docs/user/getting-started.md index 7e2e5df6..a0e6e6f6 100644 --- a/docs/user/getting-started.md +++ b/docs/user/getting-started.md @@ -250,8 +250,8 @@ the record of what your results were computed with. ## 5. Materialize Launch the built-in local offer; no compute configuration is needed. It provides -all usable CPUs and RAM for 30 minutes. Keep the returned ID in `CLUSTER` for this -walkthrough. Configured remote offers coexist with that default. A catalog can +all usable CPUs and RAM, and stops once it has had no work 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. diff --git a/evals/prompt.md b/evals/prompt.md index d905996e..7675453d 100644 --- a/evals/prompt.md +++ b/evals/prompt.md @@ -53,9 +53,9 @@ inspect `lc compute status` and reuse it rather than launching another. 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. +A local cluster ends after 30 minutes without task activity; running work keeps +it alive, and `--time` adds a hard lifetime that ends it even mid-run. After it +ends, 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 diff --git a/src/lightcone/cli/compute.py b/src/lightcone/cli/compute.py index 01d221d4..729b5440 100644 --- a/src/lightcone/cli/compute.py +++ b/src/lightcone/cli/compute.py @@ -44,13 +44,15 @@ def _table(headers: list[str], rows: list[list[str]]) -> None: from rich.console import Console from rich.table import Table - table = Table(*headers, box=None, padding=(0, 2)) + table = Table(*headers, box=None, padding=(0, 1)) for row in rows: table.add_row(*row) Console(markup=False).print(table) -def _duration(seconds: int) -> str: +def _duration(seconds: int | None) -> str: + if seconds is None: + return "-" minutes, remainder = divmod(seconds, 60) return (f"{minutes}m" if minutes else "") + (f"{remainder}s" if remainder else "") @@ -61,8 +63,9 @@ def compute() -> None: 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. + explicit local connections supply their own offers; by default it ends + after 30 minutes without task activity rather than at a fixed age. Local + compute is automatically disabled on recognized NERSC login nodes. """ @@ -78,7 +81,10 @@ def resources(as_json: bool) -> None: click.echo(json.dumps(data)) return _table( - ["OFFER", "CPUS", "MEMORY", "GPUS", "MAX NODES", "DEFAULT", "MAX TIME", "STARTUP"], + [ + "OFFER", "CPUS", "MEMORY", "GPUS", "MAX NODES", + "DEFAULT", "MAX TIME", "IDLE", "STARTUP", + ], [ [ offer["name"], @@ -91,6 +97,7 @@ def resources(as_json: bool) -> None: str(offer["max_nodes"]), _duration(offer["time"]["default_seconds"]), _duration(offer["time"]["max_seconds"]), + _duration(offer["time"]["idle_seconds"]), offer["startup"], ] for offer in data["offers"] @@ -107,7 +114,8 @@ def resources(as_json: bool) -> None: help="Accelerator NAME[:COUNT] per node, e.g. A100:4 or GPU:1; 0 requests CPU only.") @click.option("--num-nodes", default=1, type=click.IntRange(min=1), show_default=True) @click.option( - "--time", "walltime", help="Requested walltime, e.g. 30m or 1h30m; defaults to the offer." + "--time", "walltime", + help="Hard walltime, e.g. 30m or 1h30m, even during active work; defaults to the offer's.", ) @click.option( "--startup", type=click.Choice(["fast"]), help="Require a fast startup service class." diff --git a/src/lightcone/engine/compute/__init__.py b/src/lightcone/engine/compute/__init__.py index 694b7998..353b8f93 100644 --- a/src/lightcone/engine/compute/__init__.py +++ b/src/lightcone/engine/compute/__init__.py @@ -115,6 +115,7 @@ def resources(self) -> dict[str, Any]: "time": { "default_seconds": offer.time.default_seconds, "max_seconds": offer.time.max_seconds, + "idle_seconds": offer.time.idle_seconds, }, "startup": offer.startup.class_, } @@ -148,7 +149,8 @@ def _plan_offer( 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: + limit = offer.time.max_seconds + if request.seconds is not None and limit is not None and request.seconds > limit: return None if ( offer.resources.cpus < request.cpus @@ -193,8 +195,9 @@ def plan_local( 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 "") + "no local offer matches this request; see lc compute resources for local shapes " + "and time limits, or supply --cpus and --memory for a remote allocation" + + ("; " + "; ".join(unavailable) if unavailable else "") ) def launch(self, plan: LaunchPlan) -> Identity: @@ -310,4 +313,7 @@ def connect(cluster_id: str, *, timeout: float = 10) -> Iterator[Any]: "cluster does not have its expected workers; inspect lc compute status", cluster_id=identity.encode(), ) + # A caller prepares (fetch, build, sync) before its first task, and Dask's + # idle test counts only tasks: a no-op task restarts the countdown for it. + client.submit(int, pure=False) yield client diff --git a/src/lightcone/engine/compute/catalog.py b/src/lightcone/engine/compute/catalog.py index efd23fb8..19b50995 100644 --- a/src/lightcone/engine/compute/catalog.py +++ b/src/lightcone/engine/compute/catalog.py @@ -63,6 +63,7 @@ class LocalSettings(ComputeModel): enabled: bool = True resources: Resources | None = None + time: TimeLimits | None = None @model_validator(mode="after") def cpu_only(self) -> Self: @@ -139,10 +140,10 @@ def _with_local(self) -> Catalog: 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: + if explicit and (local.resources is not None or local.time is not None): raise ComputeError( - "local.resources cannot be combined with explicit local connections; " - "set their offer resources instead" + "local.resources and local.time cannot be combined with explicit local " + "connections; set their offers' resources and time instead" ) if not explicit: if "local" in connections: @@ -164,9 +165,10 @@ def _with_local(self) -> Catalog: resources = local.resources or Resources.from_bytes( cpus=CPU_COUNT, memory_bytes=MEMORY_LIMIT, ) + # Interactive work comes and goes: end when idle, not at a fixed age. offers.append(Offer( name="local", connection="local", resources=resources, - max_nodes=1, time=TimeLimits(default="30m", max="2h"), + max_nodes=1, time=local.time or TimeLimits(idle="30m"), startup=Startup(class_="fast"), )) if not local.enabled: diff --git a/src/lightcone/engine/compute/local.py b/src/lightcone/engine/compute/local.py index 0f626311..32c5ad2c 100644 --- a/src/lightcone/engine/compute/local.py +++ b/src/lightcone/engine/compute/local.py @@ -210,8 +210,9 @@ def plan(self, offer: Offer, request: Request) -> LaunchPlan: "local GPU offers require an explicit nonempty CUDA_VISIBLE_DEVICES mask" ) seconds = request.seconds if request.seconds is not None else offer.time.default_seconds - if seconds <= 0 or seconds > offer.time.max_seconds: - raise ComputeError("local allocations require a finite time within the offer's limit") + limit = offer.time.max_seconds + if seconds is not None and limit is not None and seconds > limit: + raise ComputeError("the requested time exceeds the local offer's maximum") python = Path(self.connection.launch.get("python", sys.executable)).expanduser() scratch = configured_directory( Path(self.connection.launch.get("scratch_root", tempfile.gettempdir())) @@ -268,7 +269,8 @@ def launch(self, plan: LaunchPlan) -> Identity: startup_read = _above_stdio(startup_read) startup_write = _above_stdio(startup_write) launch = { - "deadline": started + plan.seconds, + "deadline": None if plan.seconds is None else started + plan.seconds, + "idle_timeout": plan.idle_seconds, "task_slots": plan.details["task_slots_per_node"], "scratch": str(scratch), "identity": "", diff --git a/src/lightcone/engine/compute/local_runtime.py b/src/lightcone/engine/compute/local_runtime.py index 1feea938..80043aaf 100644 --- a/src/lightcone/engine/compute/local_runtime.py +++ b/src/lightcone/engine/compute/local_runtime.py @@ -1,4 +1,4 @@ -"""Own one standard LocalCluster for a finite allocation lifetime.""" +"""Own one standard LocalCluster until it is stopped, idles out, or reaches its walltime.""" from __future__ import annotations @@ -22,7 +22,7 @@ def main() -> None: - """Run the detached allocation owner until shutdown or its walltime expires.""" + """Run the detached allocation owner until shutdown, idle expiry, or its walltime.""" os.umask(0o077) directory = private_directory(Path(sys.argv[1])) launch = read_private_json(directory / "launch.json") @@ -46,9 +46,10 @@ def expire(_signum: int, _frame: FrameType | None) -> None: signal.signal(signal.SIGTERM, stop) signal.signal(signal.SIGINT, stop) signal.signal(signal.SIGALRM, expire) - remaining = float(launch["deadline"]) - time.monotonic() - signal.setitimer(signal.ITIMER_REAL, max(0.001, remaining)) try: + if launch["deadline"] is not None: + remaining = float(launch["deadline"]) - time.monotonic() + signal.setitimer(signal.ITIMER_REAL, max(0.001, remaining)) try: startup = select.poll() startup.register(startup_fd, select.POLLIN) @@ -62,6 +63,23 @@ def expire(_signum: int, _frame: FrameType | None) -> None: launch = read_private_json(directory / "launch.json") import dask from distributed import LocalCluster + from distributed.diagnostics.plugin import SchedulerPlugin + + class Closed(SchedulerPlugin): + """End the session when the scheduler closes without the owner asking.""" + + async def close(self) -> None: + """Record why, then kill the session as the walltime does. + + ``LocalCluster``'s own close would wait on the departed + scheduler, and the owner may not have finished starting it. + """ + if not stopped.is_set(): + write_private_json(directory / "error.json", { + "error": "the scheduler closed itself, after its idle timeout " + "or at a client's request", + }) + os.killpg(os.getpgrp(), signal.SIGKILL) security = create_security(directory) allocation = read_private_json(directory / "identity.json") @@ -82,6 +100,10 @@ def expire(_signum: int, _frame: FrameType | None) -> None: "scheduler_file": str(directory / "scheduler.json"), "dashboard": False, "dashboard_address": "127.0.0.1:0", + # Dask's own activity test: tasks reset the timer; connected + # clients and status queries do not. + "idle_timeout": launch["idle_timeout"], + "plugins": [Closed()], }, local_directory=str(private_directory(Path(launch["scratch"]))), # Recipes use subprocesses: Dask's Python-process RSS cannot enforce diff --git a/src/lightcone/engine/compute/model.py b/src/lightcone/engine/compute/model.py index eb16380f..3c89a0cb 100644 --- a/src/lightcone/engine/compute/model.py +++ b/src/lightcone/engine/compute/model.py @@ -307,22 +307,35 @@ class Connection(ComputeModel): class TimeLimits(ComputeModel): - """Configured durations with derived seconds for native allocation requests.""" + """How an allocation ends: a hard walltime, an idle timeout, or both. - default: Duration - max: Duration + ``default`` and ``max`` bound the walltime; ``idle`` ends an allocation + once its scheduler has had no task activity for that long. Each is + optional, but an allocation must have some way to end. + """ + + default: Duration | None = None + max: Duration | None = None + idle: Duration | None = None + + @property + def default_seconds(self) -> int | None: + return None if self.default is None else duration(self.default) @property - def default_seconds(self) -> int: - return duration(self.default) + def max_seconds(self) -> int | None: + return None if self.max is None else duration(self.max) @property - def max_seconds(self) -> int: - return duration(self.max) + def idle_seconds(self) -> int | None: + return None if self.idle is None else duration(self.idle) @model_validator(mode="after") def ordered_limits(self) -> Self: - if self.default_seconds > self.max_seconds: + if self.default is None and self.idle is None: + raise ValueError("time needs a default walltime or an idle timeout") + default, limit = self.default_seconds, self.max_seconds + if default is not None and limit is not None and default > limit: raise ValueError("default time exceeds its maximum") return self @@ -412,7 +425,7 @@ class LaunchPlan(ComputeModel): connection: Connection offer: Offer request: Request - seconds: PositiveInt + seconds: PositiveInt | None details: dict[str, Any] = Field(default_factory=dict) name: str | None = None @@ -420,6 +433,10 @@ class LaunchPlan(ComputeModel): def resources(self) -> Resources: return self.offer.resources + @property + def idle_seconds(self) -> int | None: + return self.offer.time.idle_seconds + @property def num_nodes(self) -> int: return self.request.num_nodes @@ -435,6 +452,7 @@ def as_dict(self) -> dict[str, Any]: "num_nodes": self.num_nodes, "resources": self.resources.as_dict(), "time_seconds": self.seconds, + "idle_seconds": self.idle_seconds, "startup": self.offer.startup.class_, "launch": self.details, } diff --git a/src/lightcone/engine/compute/runtime.py b/src/lightcone/engine/compute/runtime.py index b27419b1..32c0a339 100644 --- a/src/lightcone/engine/compute/runtime.py +++ b/src/lightcone/engine/compute/runtime.py @@ -25,10 +25,12 @@ #: served. A departed client's events, which carry the recipe output it #: forwarded, are dropped at once rather than after Dask's default hour, so #: a long-lived allocation does not accumulate every invocation's output. -SCHEDULER_CONFIG = { +#: The idle timeout is the offer's alone, never the ambient Dask config's. +SCHEDULER_CONFIG: dict[str, Any] = { "distributed.scheduler.http.routes": [], "distributed.worker.http.routes": [], "distributed.scheduler.events-cleanup-delay": "0s", + "distributed.scheduler.idle-timeout": None, } diff --git a/src/lightcone/engine/compute/slurm.py b/src/lightcone/engine/compute/slurm.py index 8bdb7484..d5a13bd8 100644 --- a/src/lightcone/engine/compute/slurm.py +++ b/src/lightcone/engine/compute/slurm.py @@ -259,6 +259,11 @@ def plan(self, offer: Offer, request: Request) -> LaunchPlan: if interface is not None: interface = _value(interface, "interface") seconds = request.seconds or offer.time.default_seconds + if seconds is None or offer.time.idle is not None: + raise ComputeError( + "Slurm allocations end at their native walltime: " + "set the offer's time.default and remove time.idle" + ) hours, remainder = divmod(seconds, 3600) minutes, seconds_part = divmod(remainder, 60) args = self._scope() diff --git a/tests/test_compute.py b/tests/test_compute.py index 16868ed7..471c57fe 100644 --- a/tests/test_compute.py +++ b/tests/test_compute.py @@ -269,9 +269,11 @@ def test_missing_default_catalog_exposes_stable_local_resources_without_writing_ offer = first.offers[0] assert (offer.name, offer.connection) == ("local", "local") assert (offer.resources.cpus, offer.resources.memory_bytes, offer.max_nodes) == (1, GIB, 1) - assert (offer.time.default_seconds, offer.time.max_seconds, offer.startup.class_) == ( - 1800, 7200, "fast", + # No hard lifetime: the built-in offer ends after 30 minutes without task activity. + assert (offer.time.default_seconds, offer.time.max_seconds, offer.time.idle_seconds) == ( + None, None, 1800, ) + assert offer.startup.class_ == "fast" monkeypatch.setattr("dask.system.CPU_COUNT", 1) monkeypatch.setattr("distributed.system.MEMORY_LIMIT", GIB) runner = CliRunner() @@ -285,10 +287,13 @@ def test_missing_default_catalog_exposes_stable_local_resources_without_writing_ assert json.loads(planned.output)["plan"]["offer"] == "local" assert list(default_home.iterdir()) == [] service = compute.Compute() - assert service.plan(Request.parse("1", "1", time="2h", startup="fast")).seconds == 7200 + idle = service.plan(Request.parse("1", "1", startup="fast")) + assert (idle.seconds, idle.idle_seconds) == (None, 1800) + # An explicit hard lifetime has no built-in maximum, and the idle timeout still applies. + both = service.plan(Request.parse("1", "1", time="3h")) + assert (both.seconds, both.idle_seconds) == (10800, 1800) for request in ( - Request.parse("2", "1"), Request.parse("1", "2"), - Request.parse("1", "1", num_nodes=2), Request.parse("1", "1", time="3h"), + Request.parse("2", "1"), Request.parse("1", "2"), Request.parse("1", "1", num_nodes=2), ): with pytest.raises(ComputeError, match="no configured offer"): service.plan(request) @@ -639,6 +644,26 @@ def test_builtin_name_conflicts_identify_the_catalog_and_remedy( assert f"rename the configured {kind}" in str(error.value) +def test_local_time_replaces_the_builtin_idle_timeout( + 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 time: {idle: 1h, max: 8h}\n") + monkeypatch.setenv("LC_COMPUTE_CONFIG", str(path)) + service = compute.Compute() + plan = service.plan_local() + assert (plan.seconds, plan.idle_seconds) == (None, 3600) + assert plan.as_dict()["idle_seconds"] == 3600 + assert service.plan_local(time="8h").seconds == 8 * 3600 + with pytest.raises(ComputeError, match="local shapes and time limits"): + service.plan_local(time="9h") + listed = CliRunner().invoke(main, ["compute", "resources"]) + assert listed.exit_code == 0, listed.output + assert "IDLE" in listed.output and "60m" in listed.output + + 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" @@ -662,10 +687,12 @@ def test_explicit_local_offers_keep_their_sizes_and_replace_the_implicit_offer( 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() + for setting, value in ( + ("resources", {"cpus": 2, "memory": 2}), ("time", {"idle": "1h"}), + ): + catalog.write_text(yaml.safe_dump({**data, "local": {setting: value}})) + with pytest.raises(ComputeError, match="cannot be combined with explicit local"): + Catalog.load() def test_configured_local_budget_still_must_fit_host_capacity( @@ -809,7 +836,7 @@ def test_catalog_uses_the_public_models_and_roundtrips_without_an_adapter(catalo assert dumped["offers"][0]["resources"] == { "cpus": 4, "memory": Decimal(8), "accelerators": None, } - assert dumped["offers"][0]["time"] == {"default": "30m", "max": "2h"} + assert dumped["offers"][0]["time"] == {"default": "30m", "max": "2h", "idle": None} assert Catalog.model_validate(dumped) == loaded assert Catalog.model_validate_json(loaded.model_dump_json(by_alias=True)) == loaded @@ -950,16 +977,21 @@ def test_ended_allocation_cannot_be_borrowed(catalog: Path, provider: MagicMock) def test_borrowed_client_only_detaches(catalog: Path, provider: MagicMock) -> None: - with compute.connect(IDENTITY.encode()) as client: - assert client is provider.connect.return_value.__enter__.return_value + client = provider.connect.return_value.__enter__.return_value + with compute.connect(IDENTITY.encode()) as borrowed: + assert borrowed is client provider.connect.return_value.__exit__.assert_called_once() client.shutdown.assert_not_called() provider.terminate.assert_not_called() + # An execution connection restarts the idle countdown; status polling does not. + client.submit.assert_called_once_with(int, pure=False) + compute.Compute().status(IDENTITY.encode()) + client.submit.assert_called_once() @pytest.mark.parametrize("mutation", [ - "version_missing", "namespace", "context", "offer", "limits", "reference", "unknown", - "resources_extra", "time_extra", "startup_extra", "connection_extra", + "version_missing", "namespace", "context", "offer", "limits", "unbounded", "reference", + "unknown", "resources_extra", "time_extra", "startup_extra", "connection_extra", ]) def test_invalid_catalog_is_rejected(catalog: Path, mutation: str) -> None: data = yaml.safe_load(catalog.read_text()) @@ -976,6 +1008,8 @@ def test_invalid_catalog_is_rejected(catalog: Path, mutation: str) -> None: data["offers"][1]["name"] = "quick" elif mutation == "limits": data["offers"][0]["time"]["default"] = "3h" + elif mutation == "unbounded": + data["offers"][0]["time"] = {"max": "2h"} elif mutation == "reference": data["offers"][0]["connection"] = "missing" elif mutation == "connection_extra": diff --git a/tests/test_compute_local.py b/tests/test_compute_local.py index 9c6fb90d..138b51c0 100644 --- a/tests/test_compute_local.py +++ b/tests/test_compute_local.py @@ -2,6 +2,7 @@ from __future__ import annotations +import asyncio import json import os import signal @@ -15,12 +16,13 @@ from pathlib import Path from types import SimpleNamespace from typing import Any -from unittest.mock import MagicMock +from unittest.mock import DEFAULT, MagicMock from uuid import uuid4 import psutil import pytest from click.testing import CliRunner +from distributed import fire_and_forget from lightcone.cli.commands import main from lightcone.engine.compute import Compute, local, local_runtime @@ -61,10 +63,13 @@ def provider(tmp_path: Path) -> LocalProvider: ) -def _launch(provider: LocalProvider, *, seconds: int = 60) -> Identity: +def _launch( + provider: LocalProvider, *, seconds: int | None = 60, idle: str | None = None, +) -> Identity: offer = Offer( name="small", connection="workstation", resources=Resources(cpus=1, memory_gib=0.5), - max_nodes=1, time=TimeLimits(default="1m", max="1m"), + max_nodes=1, + time=TimeLimits(default="1m", max="1m") if idle is None else TimeLimits(idle=idle), ) return provider.launch( provider.plan(offer, Request(cpus=1, memory_bytes=512 * 1024**2, seconds=seconds)) @@ -504,6 +509,90 @@ def test_walltime_expires_without_a_connected_client(provider: LocalProvider) -> provider.terminate(allocation) +def test_an_unused_allocation_idles_out_despite_status_polling_and_frees_its_name( + provider: LocalProvider, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + catalog = tmp_path / "compute.yaml" + catalog.write_text(json.dumps({ + "version": 1, + "connections": { + "workstation": { + "provider": "local", + "namespace": provider.connection.namespace, + "launch": provider.connection.launch, + }, + }, + "offers": [{ + "name": "small", + "connection": "workstation", + "resources": {"cpus": 1, "memory": 0.5}, + "max_nodes": 1, + "time": {"idle": "8s"}, + }], + })) + monkeypatch.setenv("LC_COMPUTE_CONFIG", str(catalog)) + plan = Compute().plan(Request(cpus=1, memory_bytes=512 * 1024**2), name="analysis") + assert (plan.seconds, plan.idle_seconds) == (None, 8) + identities = [] + try: + first = Compute().launch(plan) + identities.append(first) + assert Compute().status("analysis", wait=True, timeout=20).ready + # Each status query connects a client; none of them is task activity. + deadline = time.monotonic() + 20 + while Compute().status(first.encode()).phase != "ended": + assert time.monotonic() < deadline, "an unused allocation outlived its idle timeout" + time.sleep(0.2) + assert "idle timeout" in Compute().status(first.encode()).reason + assert Compute().discover() == ([], {}) + # The ended owner no longer holds this machine's one local allocation, or its name. + second = Compute().launch(plan) + identities.append(second) + assert second.name == "analysis" + assert Compute().status("analysis", wait=True, timeout=20).identity == second + finally: + for identity in identities: + provider.terminate(identity) + + +def test_work_outlasts_the_idle_timeout_and_new_work_restarts_it( + provider: LocalProvider, +) -> None: + identity = _launch(provider, seconds=None, idle="8s") + try: + _ready(provider, identity) + with provider.connect(identity) as client: + # One task slot: the second task queues behind the first. + assert client.gather(client.map(time.sleep, [5, 5], pure=False)) == [None, None] + finished = time.monotonic() + time.sleep(5) + assert client.submit(sum, [1, 2]).result(timeout=5) == 3 + restarted = time.monotonic() + # Without the restart, the allocation would have ended about 8s after `finished`. + time.sleep(max(0.0, finished + 9.5 - time.monotonic())) + assert provider.inspect(identity).phase == "active" + _ended(provider, identity, timeout=max(1.0, restarted + 10 - time.monotonic())) + finally: + provider.terminate(identity) + + +def test_an_explicit_walltime_interrupts_active_work_despite_an_idle_timeout( + provider: LocalProvider, +) -> None: + identity = _launch(provider, seconds=10, idle="1m") + try: + _ready(provider, identity) + with provider.connect(identity) as client: + fire_and_forget(client.submit(time.sleep, 60, pure=False)) + deadline = time.monotonic() + 5 + while not any(client.processing().values()): + assert time.monotonic() < deadline, "the task never started" + time.sleep(0.05) + _ended(provider, identity, timeout=15) + finally: + provider.terminate(identity) + + @pytest.mark.parametrize("ending", ["down", "walltime"]) def test_an_ended_allocation_keeps_its_record_but_not_its_secrets_or_scratch( provider: LocalProvider, ending: str, @@ -1020,7 +1109,7 @@ def test_private_material_rejects_symlinks_broad_modes_and_hardlinks(tmp_path: P read_private_json(path) -def test_local_plan_is_one_node_finite_cooperative_and_does_not_allocate( +def test_local_plan_is_one_node_bounded_cooperative_and_does_not_allocate( provider: LocalProvider, ) -> None: offer = Offer( @@ -1032,7 +1121,7 @@ def test_local_plan_is_one_node_finite_cooperative_and_does_not_allocate( assert "cooperative" in plan.details["resource_enforcement"] with pytest.raises(ComputeError, match="one execution node"): provider.plan(offer, Request(cpus=1, memory_bytes=512 * 1024**2, num_nodes=2)) - with pytest.raises(ComputeError, match="finite time"): + with pytest.raises(ComputeError, match="exceeds the local offer's maximum"): provider.plan(offer, Request(cpus=1, memory_bytes=512 * 1024**2, seconds=61)) @@ -1272,6 +1361,7 @@ def test_local_gpu_offers_are_unavailable_outside_linux( def test_local_runtime_advertises_configured_resources_and_preserves_native_gpu_mask( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, gpus: int | None, ) -> None: + import dask import distributed directory = private_directory(tmp_path / "allocation", create=True) @@ -1294,17 +1384,34 @@ def test_local_runtime_advertises_configured_resources_and_preserves_native_gpu_ monkeypatch.setattr(local_runtime, "create_security", lambda _: None) cluster = MagicMock() cluster.return_value.__enter__.return_value.scheduler.id = "Scheduler-gpu" + ambient: dict[str, object] = {} + + def build(**_: object) -> object: + ambient["idle"] = dask.config.get("distributed.scheduler.idle-timeout") + return DEFAULT + + cluster.side_effect = build + killed: list[tuple[int, int]] = [] + monkeypatch.setattr(os, "killpg", lambda *args: killed.append(args)) 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, + "identity": "allocation", "deadline": time.monotonic() + 60, "idle_timeout": 30, "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() + with dask.config.set({"distributed.scheduler.idle-timeout": "10m"}): + local_runtime.main() + # The offer's idle timeout reaches the scheduler; an ambient one never does. + assert ambient["idle"] is None assert cluster.call_args.kwargs["resources"] == {"CPU": 1, "MEMORY": 1024**3, "GPU": gpus or 0} + scheduler_kwargs = cluster.call_args.kwargs["scheduler_kwargs"] + assert scheduler_kwargs["idle_timeout"] == 30 + # A close the owner asked for is Dask's graceful one, never the session kill. + asyncio.run(scheduler_kwargs["plugins"][0].close()) + assert killed == [] assert os.environ["CUDA_VISIBLE_DEVICES"] == ("3,1" if gpus else "") diff --git a/tests/test_compute_slurm.py b/tests/test_compute_slurm.py index cec8c871..b38e99e0 100644 --- a/tests/test_compute_slurm.py +++ b/tests/test_compute_slurm.py @@ -310,6 +310,17 @@ def test_plan_includes_only_an_explicit_partition_and_requested_walltime( assert calls == [] +@pytest.mark.parametrize("time", [ + TimeLimits(idle="30m"), TimeLimits(default="1h", max="4h", idle="30m"), +]) +def test_plan_refuses_an_offer_that_would_end_on_idle( + provider: slurm.SlurmProvider, offer: Offer, time: TimeLimits, +) -> None: + # Only the native walltime ends a Slurm allocation; an ignored idle timeout would lie. + with pytest.raises(ComputeError, match="native walltime"): + provider.plan(offer.replace(time=time), Request.parse("256", "480", time="30m")) + + def test_plan_refuses_multiple_partitions(provider: slurm.SlurmProvider, offer: Offer) -> None: with pytest.raises(ComputeError, match="one native partition"): provider.plan(