Skip to content
Closed
10 changes: 8 additions & 2 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -1767,9 +1767,15 @@ does not independently guarantee a finite termination deadline for Slurm jobs.

**Execution borrows a client and leaves the allocation alive.** Validate native
identity and scheduler readiness. The driver keeps git and convergence. Use unique
invocation task keys. Interrupted unreported outputs remain in place because a
invocation task keys. `engine.execution` requires an atomic scheduler claim before
each recipe or probe runs. A repeated claim or missing invocation state refuses
execution, including recomputation after a completed result is lost. There are no
completion receipts, leases, heartbeats, or automatic recovery. `retries=0` alone
does not prevent Dask recomputation after worker loss. This is an invocation-local
replay guard, not a project lock or proof that a disconnected command stopped.
Interrupted unreported outputs remain in place because a
client disconnect does not prove remote subprocess termination. Comprehensive
cancellation/fencing and simultaneous writers are deferred by explicit user decision.
cancellation and simultaneous writers are deferred by explicit user decision.
Local containerized processes can outlive process-group shutdown; do not claim
that `down` or walltime proves an external runtime's containers have stopped.
Read-only project validation precedes cluster connection, and a run with no
Expand Down
14 changes: 14 additions & 0 deletions docs/api/compute.md
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,20 @@ Dask chooses the workers and handles dependencies; invocation-specific keys prev
unintended reuse across commands. There is no worker-selection layer, per-worker
preflight orchestration, source fingerprinting, or login-node guard. Driver-side
preparation and the existing task runtime/sandbox checks remain in their owners.

`engine.execution` prevents automatic replay of side effects. The driver registers
an invocation in the existing Dask scheduler; before running a recipe or probe,
its worker must receive an atomic claim for that task. A repeated claim, missing
invocation state, or absent driver client fails before the command runs. This also
refuses recomputation when a completed result was lost: Lightcone does not store
or recover results.
`retries=0` alone cannot prevent [Dask's recomputation after worker loss](https://distributed.dask.org/en/stable/resilience.html).
There are no leases, heartbeats, command supervisors, or separate execution service.
This guard prevents automatic replay within an invocation; it does not lock the
project against another invocation or prove that a disconnected worker stopped.
The context removes its record on exit; later invocations prune records left by
disconnected clients. Removing a record cannot authorize a later claim.

`output.py` transports byte chunks through standard Dask events so detached
workers' output reaches the invoking CLI. It uses the borrowed client's event
topic, which the schedulers lc launches drop as soon as the client disconnects
Expand Down
4 changes: 4 additions & 0 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,10 @@ client on exit. Both execution commands require a cluster ID. The materializatio
scheduler keeps its `submit`/`completed` seam. Driver preparation and existing
task runtime/sandbox checks remain unchanged. Tasks use ordinary Dask scheduling;
there is no separate worker-selection or preflight layer, or site-marker guard. No execution command implicitly allocates compute.
Before a recipe or probe runs, `engine.execution` claims its task once in the
existing scheduler. Repeated claims and missing invocation state refuse execution,
so Dask cannot silently replay side effects after a worker or result is lost.
This adds no heartbeat or execution service and does not guarantee cancellation.
See [compute internals](api/compute.md) and [deployment limits](user/cluster.md).

## The publication view
Expand Down
10 changes: 8 additions & 2 deletions docs/user/cluster.md
Original file line number Diff line number Diff line change
Expand Up @@ -332,9 +332,15 @@ directories and credential files still reject symlinks, retain ownership and
ancestor-permission checks, and require modes `0700` and `0600`, respectively.
The catalog's location is independent of the private connection files.

Lightcone prevents Dask from automatically rerunning a recipe or probe within the
same invocation: each task must claim permission once in the existing scheduler
before it runs. If a worker or its result is lost, an attempted replay fails
instead. Missing scheduler state also refuses execution. There is no automatic
recovery; inspect the allocation and outputs before starting another invocation.

Use one execution invocation per project at a time. Concurrent writers,
comprehensive cancellation, task fencing, and recovery after client/worker loss
are not guaranteed. A lost client does not prove its subprocesses stopped.
comprehensive cancellation, and recovery after client/worker loss are not
guaranteed. A lost client does not prove its subprocesses stopped.
Unreported partial outputs are retained after interruption rather than restored
while a task may still write them. End the allocation and establish that work has
stopped before inspecting or repairing that project's outputs.
Expand Down
91 changes: 91 additions & 0 deletions src/lightcone/engine/execution.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,91 @@
"""Refuse repeated execution of side-effecting tasks in the existing Dask scheduler."""

from __future__ import annotations

from collections.abc import Callable, Iterator
from contextlib import contextmanager
from dataclasses import dataclass
from typing import Any
from uuid import uuid4

from lightcone.engine.project import ProjectError


def _register(invocation: str, owner: str, *, dask_scheduler: Any) -> None:
records = dask_scheduler.extensions.setdefault("lightcone-executions", {})
if invocation in records:
raise ProjectError("execution is already registered")
# Collect abandoned invocations using Dask's own client membership, without
# a background service or a second liveness protocol.
for key, record in list(records.items()):
if record["client"] not in dask_scheduler.clients:
del records[key]
records[invocation] = {"client": owner, "tasks": set()}


def _claim(invocation: str, task: str, *, dask_scheduler: Any) -> None:
records = dask_scheduler.extensions.get("lightcone-executions", {})
record = records.get(invocation)
if record is None or record["client"] not in dask_scheduler.clients:
raise ProjectError("execution is no longer registered or its client disconnected")
# This synchronous callback runs atomically on the scheduler's event loop.
# Claims survive task failure, worker loss, and Dask forgetting task results.
if task in record["tasks"]:
raise ProjectError(f"{task}: already claimed; refusing duplicate execution")
record["tasks"].add(task)


def _forget(invocation: str, *, dask_scheduler: Any) -> None:
dask_scheduler.extensions.get("lightcone-executions", {}).pop(invocation, None)


def _rpc(client: Any, function: Callable[..., None], *args: Any) -> None:
try:
client.sync(client.run_on_scheduler, function, *args, callback_timeout=5)
except ProjectError:
raise
except Exception as exc:
raise ProjectError(f"cannot contact the Dask execution guard: {exc}") from exc


def _call(invocation: str, task: str, function: Callable[..., Any], *args: Any) -> Any:
from distributed import get_client

# A lost claim reply is ambiguous. Execute only after acknowledgment, and
# never recreate missing state or retry the claim on a worker's behalf.
_rpc(get_client(), _claim, invocation, task)
return function(*args)


@dataclass(frozen=True)
class Invocation:
"""Submit tasks that may begin at most once within this invocation."""

client: Any
id: str

def submit(self, function: Callable[..., Any], *args: Any, key: str) -> Any:
"""Guard each task before its first side effect, including Dask recomputation."""
return self.client.submit(
_call, self.id, key, function, *args,
key=f"lc-{self.id}-{key}", pure=False, retries=0,
)


@contextmanager
def invocation(client: Any) -> Iterator[Invocation]:
"""Register claims for one borrowed client; never stop its running commands.

Missing records refuse admission, so forgetting an invocation also prevents
late tasks from starting. Cleanup failures cannot discard received results;
abandoned records are collected when another invocation registers.
"""
run = Invocation(client, uuid4().hex)
_rpc(client, _register, run.id, client.id)
try:
yield run
finally:
try:
_rpc(client, _forget, run.id)
except ProjectError:
pass
15 changes: 8 additions & 7 deletions src/lightcone/engine/materialize.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,9 +39,8 @@
from dataclasses import asdict, dataclass, field, replace
from pathlib import Path
from typing import TYPE_CHECKING, Any, Protocol
from uuid import uuid4

from lightcone.engine import assets, container, dataset, identity, plan, project, worker
from lightcone.engine import assets, container, dataset, execution, identity, plan, project, worker
from lightcone.engine.plan import Graph, Key, Task
from lightcone.engine.project import ProjectError

Expand Down Expand Up @@ -676,16 +675,16 @@ class _Dask:
"""A borrowed Dask client, narrowed to what the graph driver needs."""

client: Any
invocation: str
invocation: execution.Invocation
output: Forwarder

def submit(self, fn: Any, *args: Any, key: str) -> Any:
"""Submit an ordinary Dask task with a unique key and forwarded output."""
from lightcone.engine.compute.output import call

return self.client.submit(
return self.invocation.submit(
call, fn, self.output.topic, key, *args,
key=f"lc-{self.invocation}-{key}", pure=False,
key=key,
)

def completed(self, handles: list[Any]) -> Iterator[worker.TaskResult]:
Expand Down Expand Up @@ -725,8 +724,10 @@ def cluster_for_run(cluster_id: str) -> Iterator[Scheduler]:
from lightcone.engine.compute.output import forwarding

with compute.connect(cluster_id) as client:
invocation = uuid4().hex
with forwarding(client, stdout="stderr") as output:
with (
forwarding(client, stdout="stderr") as output,
execution.invocation(client) as invocation,
):
yield _Dask(client, invocation, output)


Expand Down
10 changes: 4 additions & 6 deletions src/lightcone/engine/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,8 @@
from dataclasses import replace
from pathlib import Path
from typing import Any
from uuid import uuid4

from lightcone.engine import container, sandbox
from lightcone.engine import container, execution, sandbox
from lightcone.engine.project import (
SPEC_FILENAME,
ProjectError,
Expand Down Expand Up @@ -53,11 +52,10 @@ def probe(project: Path, command: Sequence[str], *, cluster_id: str) -> sandbox.
with compute.connect(cluster_id) as client:
runtime = container.runtime_for_run(project, build=False)
notes = [f"uv: {warning}" for warning in container.converge(runtime)]
invocation = uuid4().hex
with forwarding(client) as output:
future = client.submit(
with forwarding(client) as output, execution.invocation(client) as invocation:
future = invocation.submit(
call, _probe, output.topic, "probe", runtime, paths, tuple(command),
key=f"lc-{invocation}-probe", pure=False,
key="probe",
)
try:
outcome: sandbox.Outcome = future.result()
Expand Down
Loading
Loading