From d733329e406a08f4aa7d3a70680bbb256ce48fcd Mon Sep 17 00:00:00 2001 From: Jeremy Rapin Date: Mon, 14 Sep 2026 16:38:48 +0200 Subject: [PATCH 1/3] [Step] Fix nested deadlocks --- .../claim-inheritance-design-history.md | 55 +++++ .../proposals/claim-inheritance-design.md | 229 ++++++++++++++++++ exca/steps/backends.py | 48 +++- exca/steps/test_backends.py | 8 + 4 files changed, 331 insertions(+), 9 deletions(-) create mode 100644 docs/internal/proposals/claim-inheritance-design-history.md create mode 100644 docs/internal/proposals/claim-inheritance-design.md diff --git a/docs/internal/proposals/claim-inheritance-design-history.md b/docs/internal/proposals/claim-inheritance-design-history.md new file mode 100644 index 00000000..825db305 --- /dev/null +++ b/docs/internal/proposals/claim-inheritance-design-history.md @@ -0,0 +1,55 @@ +# Claim Inheritance — cut content + +Append-only. Verbatim content removed from `claim-inheritance-design.md`, with a +one-line note on why. + +## Open question: single-task array ids + +Cut: only matters if `Opt-self-job` is adopted, which `UC-nested-chain` now puts in +doubt. + +> Unconfirmed: whether a single-task `executor.batch()` yields an array-style +> (`"_"`) or plain job id. Confirm on-cluster before relying on +> `Opt-self-job`. + +## Opt-reject-config + +Cut: rejected — the config is legitimate, it should work. + +> Guard at validation or first dispatch when a chain and its last resolved step both +> carry a `_concurrent` backend on the same folder. Two resource requests for one +> cache cell is a contradiction, and the outer's worker requests a second allocation +> for work it was already allocated for. `Parallel._unify_infra` sets the precedent +> by refusing mismatched backends across itself and its steps. + +## Opt-prefix-split + +Cut: rejected — redefines what the chain's infra covers. + +> When the chain has infra and its last resolved step has its own, the chain's +> backend covers `steps[:-1]` (a distinct prefix `step_uid`) and the last step owns +> the shared cell. Honors both resource specs, no double allocation, no self-claim. +> Costs a cache cell for the intermediate. + +## Opt-outer-wins + +Cut: rejected — silently discards the tail's resource request, usually the expensive +one. + +> Chain's backend runs everything; the last step's infra degrades to caching only. + +## Opt-no-shared-identity + +Cut: rejected — drops a deliberate design property. + +> Give the chain its own uid segment so it no longer shares a cell with its last +> step. Kills the whole class of bugs, but duplicates the final result on disk. + +## Workaround: infra-free outer chain + +Cut: different topology, so it does not address `UC-tail`. + +> Nest instead of flatten, leaving the outer chain infra-free: +> `Chain(steps=[Chain(steps=[Preprocess()], infra=...), Train(infra=...)])`. +> Verified working — the inner chain takes the prefix uid, the tail takes the shared +> cell, and nobody self-claims. diff --git a/docs/internal/proposals/claim-inheritance-design.md b/docs/internal/proposals/claim-inheritance-design.md new file mode 100644 index 00000000..fe11f5d2 --- /dev/null +++ b/docs/internal/proposals/claim-inheritance-design.md @@ -0,0 +1,229 @@ +# Claim Inheritance + +## Problem + +A chain and its last step resolve to the same `step_uid` — `Chain._uid_steps()` +flattens to its children — so they share a `step_folder`, an `inflight.db`, and an +item uid. Call that shared entry the **cell**. + +The driver holds the cell's claim across the whole submit-and-wait (`Backend._run` +wraps `_execute` in `_claim`; `_SubmititBackend._execute` ends on `job.result()`). +Inside the job, the tail claims the same cell. Ownership is keyed on PID: + +```python +if info.pid == my_pid: # inflight.wait_for_inflight + remaining.discard(uid) + continue +``` + +Off-process, the PIDs differ: + +``` +driver claim(cell) -> row{pid=D} submit A wait on A +job A claim(cell) -> row{pid=D} wait on D <- deadlock +``` + +Reproduced with `LocalProcess`; the worker only moved once the driver was killed: + +``` +WARN Waiting for 1 in-flight items (of 1 requested) held by: pid=46687 [local] x1 +INFO Reclaimed 1 items from dead workers: pid=46687 [local] x1 +``` + +Trigger: + +- outer executes off-process — `Slurm`, `Auto`, `LocalProcess`, `ProcessPool` +- **and** the tail's backend is `_concurrent` + +Unaffected: + +- as outer: `Cached`, `SubmititDebug`, `ThreadPool` — same PID +- as tail: `Cached`, `SubmititDebug` — `_claim` only builds a registry when + `_concurrent` + +Slurm amplifies it: `WorkerInfo.wait()` calls `SlurmJob.wait()` on the job the worker +runs in, so it blocks to wall-clock timeout holding two allocations. + +## Use cases + +### UC-tail + +```python +chain = Chain( + steps=[Preprocess(), Train(infra={"backend": "Slurm", "gpus_per_node": 8})], + infra={"backend": "Slurm", "cpus_per_task": 4, "folder": "/cache"}, +) +chain.run() # today: deadlock +``` + +### UC-tail-force + +`UC-tail` with `mode="force"`, which fails earlier and louder: + +- `_prepare` clears before claiming +- `_clear_caches` cancels the submitit job in `inflight.db` — the outer's +- the worker `scancel`s its own array, by base id, so every task dies + +Slurm-only: `LocalProcess` rows have no `job_folder` and skip the cancel branch. + +### UC-nested-chain + +`Chain(infra) > Chain(infra) > Step(infra)` — identity flattens recursively, so three +claimants land on one cell. + +### UC-thread-tail + +`ThreadPool` outer, `LocalProcess` tail. Works today via the PID check, must keep +working, and never pickles. + +### UC-sibling-sessions + +Two independent sessions in one process must not treat each other's claims as theirs +— the deferred "PID is too broad" limitation in `inflight-registry.md`. + +### UC-third-party + +An unrelated driver on the same cell must still wait for the running job, and reclaim +it if the owner dies. + +## Properties + +Hard: + +- **P-no-self-wait**: a worker never waits on an ancestor's claim, at any depth. +- **P-no-self-cancel**: clearing never cancels a job in the caller's ancestry. +- **P-ancestor-row-intact**: a descendant never releases or repoints an ancestor's + row — `UC-third-party` reads liveness off it. +- **P-liveness-unchanged**: reclamation keeps working off the existing `pid` / + `job_id` / `job_folder` columns; new state brings no liveness story of its own. +- **P-no-livelock**: `claim()` agrees with the wait, or the deadlock becomes a + `"Claim race: got 0/1 items, re-waiting"` spin. +- **P-advisory**: failures degrade to duplicate work, never wrong results or a hang. + +Soft: + +- **P-own-narrow**: ownership identifies the session, not the process. +- **P-no-schema-change**: leave the `inflight` schema alone. +- **P-propagation-free**: correctness does not depend on state reaching the worker. + +`P-own-narrow` and `P-propagation-free` contradict — narrowing below PID needs an +identity only the session's descendants can see, which has to travel. + +## Options + +### Opt-inherited-claims + +Ancestors' claims travel with the work. `ComputeBatch.__getstate__` already strips +`info.claim`; keep a reduced form instead. + +```python +@dataclasses.dataclass +class CoordinationInfo: + ... + inherited: frozenset[tuple[str, str]] = frozenset() # (step_folder, uid) +``` + +`Backend._claim` subtracts `inherited` from the `inflight_session` request; the +batch's pending set is untouched, so the work still runs. + +``` +driver inherited={} claim(cell) submit A +job A inherited={cell} claim() -> {} submit B +job B inherited={cell} claim() -> {} runs +``` + +- Depth-agnostic: the set accumulates forward, never consulting who holds the row. +- Release is safe for free — never claimed, so the session's `finally` cannot release + it. +- Stamping is not free: both submit sites pass `uids=batch.items.uids`, which still + spans inherited cells, so they need narrowing to `claim.uids`. +- Costs `P-propagation-free` — a propagation gap silently restores the deadlock. +- Buys `P-own-narrow` only if it *replaces* the PID clause rather than joining it. + +### Opt-self-job + +Compare the row's `job_id` to the worker's own. + +```python +def _own_job_id() -> str | None: + env = submitit.JobEnvironment() + return env.job_id if env.activated() else None +``` + +The ids compare exactly: the driver stores `str(job.job_id)` = `"_"` +(`submitit/slurm/slurm.py:347`), in-job `SlurmJobEnvironment.job_id` is +`f"{SLURM_ARRAY_JOB_ID}_{SLURM_ARRAY_TASK_ID}"`, and `clean_env()` at submission +stops a sub-job inheriting the outer's SLURM vars. + +**Fails `UC-nested-chain`** — the row carries one job id, the driver's: + +``` +driver row{job=A} submit A +job A row.job=A == own A -> proceeds submit B +job B row.job=A != own B -> wait on A A waits on B <- deadlock +``` + +- Repointing the row per level would fix it and break `P-ancestor-row-intact`: once + the innermost job ends, the cell reads dead under still-running ancestors. +- Carrying the ancestor *chain* of ids is propagation, i.e. `Opt-inherited-claims`. + So depth-1 is the ceiling, and it cannot be lifted within `P-propagation-free`. +- Slurm-only. PID ancestry is no substitute: stdlib gives one level + (`os.getppid()`), and `LocalProcess` workers are grandchildren through submitit's + subprocess. +- Needs its clause in `claim()`, `pre_owned` and `WorkerInfo.wait`, not just the wait + loop — otherwise `P-no-livelock`. + +### Opt-owner-column + +The `owner_token` column sketched in `inflight-registry.md`, matched by prefix so +siblings stay mutually exclusive. + +- Breaks `P-ancestor-row-intact` exactly as `Opt-self-job` does: the descendant's + `record_worker_info` repoints the ancestor's row at the sub-job. +- Pays a schema change and still needs propagation, for guarantees + `Opt-inherited-claims` already gives. + +### Opt-combined + +`Opt-inherited-claims` as the mechanism, `Opt-self-job` as backstop, as one predicate +wherever ownership is decided: + +```python +def _is_own(info: WorkerInfo, cell: tuple[str, str]) -> bool: + return ( + info.pid == os.getpid() + or cell in _inherited.get() + or (info.job_id is not None and info.job_id == _own_job_id()) + ) +``` + +The backstop covers depth 1 on Slurm only, so it guards a narrower band than it first +appears: a propagation gap at depth ≥ 2 still deadlocks. + +## Touch points + +1. `CoordinationInfo.inherited`, surviving `ComputeBatch.__getstate__`. +2. `ComputeBatch.run_and_cache` binds a `ContextVar` to `inherited | own claim` for + the duration — a `ContextVar` not a module global, since `ThreadPool` workers + carry different sets in one process. +3. `Backend._prepare` seeds `inherited` from that `ContextVar`, so depth ≥ 2 + accumulates. +4. `Backend._claim` subtracts inherited cells from the `inflight_session` request. +5. `Backend._clear_caches` skips cancellation for inherited cells (`UC-tail-force`). +6. `_SubmititBackend._execute` (`backends.py:772`) and `_PoolBackend._submit_pool` + (`backends.py:974`) intersect their `uids=batch.items.uids` with `claim.uids`, so + an inherited cell keeps pointing at the ancestor's job. +7. `inflight.py` unchanged. + +`UC-thread-tail` never pickles, but needs nothing extra: threads keep `info.claim` +live, so step 2 reads it directly and the same `ContextVar` carries it. + +## Open questions + +- **Keep or drop the PID clause?** Dropping satisfies `P-own-narrow` and closes + `UC-sibling-sessions`, but puts every re-entrant case on propagation. Keeping is + deadlock-safer and leaves the duplicate-work looseness in place. +- **Is a depth-1 Slurm backstop worth `Opt-self-job`'s clauses?** It cannot cover + `UC-nested-chain`, so it guards only the shallow propagation gap. +- **Second allocation.** The tail still requests its own job while the outer holds + one. Accepted: the deadlock goes, the double-booking stays. diff --git a/exca/steps/backends.py b/exca/steps/backends.py index ea15f65a..62a92fa5 100644 --- a/exca/steps/backends.py +++ b/exca/steps/backends.py @@ -14,6 +14,7 @@ import collections import contextlib +import contextvars import dataclasses import datetime import logging @@ -282,15 +283,18 @@ def result(self) -> tp.Any: raise RuntimeError(f"No cached entry for {self._uid}") +# cache entries claimed by the running work or its ancestors +_HELD_ENTRIES: contextvars.ContextVar[frozenset[tuple[str, str]]] = ( + contextvars.ContextVar("exca_held_entries", default=frozenset()) +) + + @dataclasses.dataclass class CoordinationInfo: - """Driver-only per-run state for a ``ComputeBatch``; stripped from the - worker pickle (``ComputeBatch.__getstate__``). - """ - mode: identity.ModeType = "cached" upstream: tuple[Step, ...] = () # this step + everything before it claim: inflight.InflightClaim | None = None + held_entries: frozenset[tuple[str, str]] = frozenset() # {(folder, uid),...} @dataclasses.dataclass @@ -303,8 +307,20 @@ class ComputeBatch: items: items.StepItems info: CoordinationInfo = dataclasses.field(default_factory=CoordinationInfo) + def claimed_uids(self) -> list[str]: + """This batch's uids claimed by its own session (entries held by an ancestor + are excluded, so their rows keep pointing at the ancestor's job).""" + claim = self.info.claim + if claim is None: + raise RuntimeError(f"batch was never claimed: {self.paths.step_uid}") + owned = set(claim.uids) + return [uid for uid in self.items.uids if uid in owned] + def __getstate__(self) -> dict[str, tp.Any]: - return {**self.__dict__, "info": CoordinationInfo()} + return { + **self.__dict__, + "info": CoordinationInfo(held_entries=self.info.held_entries), + } def select(self, uids: tp.Sequence[str]) -> ComputeBatch: """Sub-batch over *uids*, sharing step/paths/cache; copies ``info`` @@ -335,9 +351,10 @@ def run_and_cache(self) -> None: folder = self.cache_dict.folder if folder is not None: folder.mkdir(parents=True, exist_ok=True) - result_items = self.step._run_items(self.items) written_uids: list[str] = [] + token = _HELD_ENTRIES.set(self.info.held_entries) try: + result_items = self.step._run_items(self.items) with self.cache_dict.write(): for i, result in enumerate(result_items): uid = self.items.uids[i] @@ -366,6 +383,8 @@ def run_and_cache(self) -> None: for uid in inflight: reg.record(uid, e, tb) raise + finally: + _HELD_ENTRIES.reset(token) def _multi_run_and_cache(batches: list[ComputeBatch]) -> None: @@ -554,12 +573,16 @@ def _clear_caches( # Other backends may have left inflight rows for this step folder. if paths.step_folder.exists(): try: + held = _HELD_ENTRIES.get() + folder_key = str(paths.step_folder) with inflight.InflightRegistry(paths.step_folder) as reg: info = reg.get(uids) jobs: dict[str, str] = {} for uid, worker in info.items(): if worker.job_id is None or worker.job_folder is None: continue # not submitit + if (folder_key, uid) in held: + continue # an ancestor's job — cancelling it kills us # Slurm array tasks share a scheduler job; avoid per-task cancels. job_id = worker.job_id.split("_", 1)[0] jobs[job_id] = worker.job_folder @@ -617,6 +640,7 @@ def _claim(self, cbatches: list[ComputeBatch]) -> _Claimed: if len(set(step_uids)) != len(step_uids): raise ValueError(f"one batch per step_uid required, got {step_uids}") claimed = _Claimed() + held = _HELD_ENTRIES.get() try: # sort by step_uid: concurrent dispatches claim in the same order for cb in sorted(cbatches, key=lambda cb: cb.paths.step_uid): @@ -628,10 +652,16 @@ def _claim(self, cbatches: list[ComputeBatch]) -> _Claimed: reg: inflight.InflightRegistry | None = None if self._concurrent: reg = inflight.InflightRegistry(cb.paths.step_folder) + # ancestors already hold their entries: claiming them self-deadlocks + folder_key = str(cb.paths.step_folder) + request = {u for u in pending if (folder_key, u) not in held} cb = cb.select(list(pending)) cb.info.claim = claimed.stack.enter_context( - inflight.inflight_session(reg, set(pending)) + inflight.inflight_session(reg, request) ) + cb.info.held_entries = held + if reg is not None: # registry-less claims hold nothing to inherit + cb.info.held_entries |= {(folder_key, u) for u in cb.info.claim.uids} claimed.batches.append(cb) claimed.ready = [ n @@ -769,7 +799,7 @@ def _execute(self, cbatches: list[ComputeBatch]) -> None: for task, job in zip(tasks, jobs): for batch in task: assert batch.info.claim is not None # inherited from its variant - batch.info.claim.record_worker_info(job, uids=batch.items.uids) + batch.info.claim.record_worker_info(job, uids=batch.claimed_uids()) folder = batch.paths.step_folder by_folder.setdefault(folder, {})[job.job_id] = batch.items.uids for folder, records in by_folder.items(): @@ -971,7 +1001,7 @@ def _submit_pool( for task in tasks: for batch in task: assert batch.info.claim is not None # inherited from its variant - batch.info.claim.record_worker_info(uids=batch.items.uids) + batch.info.claim.record_worker_info(uids=batch.claimed_uids()) pool = utils.make_pool_executor(self._POOL_TYPE, max_workers) logger.info("Sent %s items for %s steps into a %s", n_items, len(cbatches), pool) task_futs = {pool.submit(_multi_run_and_cache, task): task for task in tasks} diff --git a/exca/steps/test_backends.py b/exca/steps/test_backends.py index 72b3296f..29019929 100644 --- a/exca/steps/test_backends.py +++ b/exca/steps/test_backends.py @@ -365,3 +365,11 @@ def run(value: float, fail: bool, mode: str) -> float: # retry both; a uid-only _recomputed would make the 2nd re-raise the 1st's error assert run(1.0, fail=False, mode="retry") == 3.0 # 2 + 1 assert run(5.0, fail=False, mode="retry") == 7.0 # 2 + 5 + + +def test_nested_dispatch_on_shared_cell(tmp_path: Path) -> None: + infra: tp.Any = {"backend": "LocalProcess", "folder": tmp_path} + inner = Chain(steps=[conftest.Mult(coeff=3.0, infra=infra)], infra=infra) + chain = Chain(steps=[conftest.Add(value=1.0), inner], infra=infra) + # identity flattens recursively: the 3 infras share one cache entry + assert chain.run(1.0) == 6.0 From 0f82b162ecdd9c07540f6ad7c0c87a5bfc1c602a Mon Sep 17 00:00:00 2001 From: Jeremy Rapin Date: Mon, 14 Sep 2026 16:39:12 +0200 Subject: [PATCH 2/3] rm --- .../claim-inheritance-design-history.md | 55 ----- .../proposals/claim-inheritance-design.md | 229 ------------------ 2 files changed, 284 deletions(-) delete mode 100644 docs/internal/proposals/claim-inheritance-design-history.md delete mode 100644 docs/internal/proposals/claim-inheritance-design.md diff --git a/docs/internal/proposals/claim-inheritance-design-history.md b/docs/internal/proposals/claim-inheritance-design-history.md deleted file mode 100644 index 825db305..00000000 --- a/docs/internal/proposals/claim-inheritance-design-history.md +++ /dev/null @@ -1,55 +0,0 @@ -# Claim Inheritance — cut content - -Append-only. Verbatim content removed from `claim-inheritance-design.md`, with a -one-line note on why. - -## Open question: single-task array ids - -Cut: only matters if `Opt-self-job` is adopted, which `UC-nested-chain` now puts in -doubt. - -> Unconfirmed: whether a single-task `executor.batch()` yields an array-style -> (`"_"`) or plain job id. Confirm on-cluster before relying on -> `Opt-self-job`. - -## Opt-reject-config - -Cut: rejected — the config is legitimate, it should work. - -> Guard at validation or first dispatch when a chain and its last resolved step both -> carry a `_concurrent` backend on the same folder. Two resource requests for one -> cache cell is a contradiction, and the outer's worker requests a second allocation -> for work it was already allocated for. `Parallel._unify_infra` sets the precedent -> by refusing mismatched backends across itself and its steps. - -## Opt-prefix-split - -Cut: rejected — redefines what the chain's infra covers. - -> When the chain has infra and its last resolved step has its own, the chain's -> backend covers `steps[:-1]` (a distinct prefix `step_uid`) and the last step owns -> the shared cell. Honors both resource specs, no double allocation, no self-claim. -> Costs a cache cell for the intermediate. - -## Opt-outer-wins - -Cut: rejected — silently discards the tail's resource request, usually the expensive -one. - -> Chain's backend runs everything; the last step's infra degrades to caching only. - -## Opt-no-shared-identity - -Cut: rejected — drops a deliberate design property. - -> Give the chain its own uid segment so it no longer shares a cell with its last -> step. Kills the whole class of bugs, but duplicates the final result on disk. - -## Workaround: infra-free outer chain - -Cut: different topology, so it does not address `UC-tail`. - -> Nest instead of flatten, leaving the outer chain infra-free: -> `Chain(steps=[Chain(steps=[Preprocess()], infra=...), Train(infra=...)])`. -> Verified working — the inner chain takes the prefix uid, the tail takes the shared -> cell, and nobody self-claims. diff --git a/docs/internal/proposals/claim-inheritance-design.md b/docs/internal/proposals/claim-inheritance-design.md deleted file mode 100644 index fe11f5d2..00000000 --- a/docs/internal/proposals/claim-inheritance-design.md +++ /dev/null @@ -1,229 +0,0 @@ -# Claim Inheritance - -## Problem - -A chain and its last step resolve to the same `step_uid` — `Chain._uid_steps()` -flattens to its children — so they share a `step_folder`, an `inflight.db`, and an -item uid. Call that shared entry the **cell**. - -The driver holds the cell's claim across the whole submit-and-wait (`Backend._run` -wraps `_execute` in `_claim`; `_SubmititBackend._execute` ends on `job.result()`). -Inside the job, the tail claims the same cell. Ownership is keyed on PID: - -```python -if info.pid == my_pid: # inflight.wait_for_inflight - remaining.discard(uid) - continue -``` - -Off-process, the PIDs differ: - -``` -driver claim(cell) -> row{pid=D} submit A wait on A -job A claim(cell) -> row{pid=D} wait on D <- deadlock -``` - -Reproduced with `LocalProcess`; the worker only moved once the driver was killed: - -``` -WARN Waiting for 1 in-flight items (of 1 requested) held by: pid=46687 [local] x1 -INFO Reclaimed 1 items from dead workers: pid=46687 [local] x1 -``` - -Trigger: - -- outer executes off-process — `Slurm`, `Auto`, `LocalProcess`, `ProcessPool` -- **and** the tail's backend is `_concurrent` - -Unaffected: - -- as outer: `Cached`, `SubmititDebug`, `ThreadPool` — same PID -- as tail: `Cached`, `SubmititDebug` — `_claim` only builds a registry when - `_concurrent` - -Slurm amplifies it: `WorkerInfo.wait()` calls `SlurmJob.wait()` on the job the worker -runs in, so it blocks to wall-clock timeout holding two allocations. - -## Use cases - -### UC-tail - -```python -chain = Chain( - steps=[Preprocess(), Train(infra={"backend": "Slurm", "gpus_per_node": 8})], - infra={"backend": "Slurm", "cpus_per_task": 4, "folder": "/cache"}, -) -chain.run() # today: deadlock -``` - -### UC-tail-force - -`UC-tail` with `mode="force"`, which fails earlier and louder: - -- `_prepare` clears before claiming -- `_clear_caches` cancels the submitit job in `inflight.db` — the outer's -- the worker `scancel`s its own array, by base id, so every task dies - -Slurm-only: `LocalProcess` rows have no `job_folder` and skip the cancel branch. - -### UC-nested-chain - -`Chain(infra) > Chain(infra) > Step(infra)` — identity flattens recursively, so three -claimants land on one cell. - -### UC-thread-tail - -`ThreadPool` outer, `LocalProcess` tail. Works today via the PID check, must keep -working, and never pickles. - -### UC-sibling-sessions - -Two independent sessions in one process must not treat each other's claims as theirs -— the deferred "PID is too broad" limitation in `inflight-registry.md`. - -### UC-third-party - -An unrelated driver on the same cell must still wait for the running job, and reclaim -it if the owner dies. - -## Properties - -Hard: - -- **P-no-self-wait**: a worker never waits on an ancestor's claim, at any depth. -- **P-no-self-cancel**: clearing never cancels a job in the caller's ancestry. -- **P-ancestor-row-intact**: a descendant never releases or repoints an ancestor's - row — `UC-third-party` reads liveness off it. -- **P-liveness-unchanged**: reclamation keeps working off the existing `pid` / - `job_id` / `job_folder` columns; new state brings no liveness story of its own. -- **P-no-livelock**: `claim()` agrees with the wait, or the deadlock becomes a - `"Claim race: got 0/1 items, re-waiting"` spin. -- **P-advisory**: failures degrade to duplicate work, never wrong results or a hang. - -Soft: - -- **P-own-narrow**: ownership identifies the session, not the process. -- **P-no-schema-change**: leave the `inflight` schema alone. -- **P-propagation-free**: correctness does not depend on state reaching the worker. - -`P-own-narrow` and `P-propagation-free` contradict — narrowing below PID needs an -identity only the session's descendants can see, which has to travel. - -## Options - -### Opt-inherited-claims - -Ancestors' claims travel with the work. `ComputeBatch.__getstate__` already strips -`info.claim`; keep a reduced form instead. - -```python -@dataclasses.dataclass -class CoordinationInfo: - ... - inherited: frozenset[tuple[str, str]] = frozenset() # (step_folder, uid) -``` - -`Backend._claim` subtracts `inherited` from the `inflight_session` request; the -batch's pending set is untouched, so the work still runs. - -``` -driver inherited={} claim(cell) submit A -job A inherited={cell} claim() -> {} submit B -job B inherited={cell} claim() -> {} runs -``` - -- Depth-agnostic: the set accumulates forward, never consulting who holds the row. -- Release is safe for free — never claimed, so the session's `finally` cannot release - it. -- Stamping is not free: both submit sites pass `uids=batch.items.uids`, which still - spans inherited cells, so they need narrowing to `claim.uids`. -- Costs `P-propagation-free` — a propagation gap silently restores the deadlock. -- Buys `P-own-narrow` only if it *replaces* the PID clause rather than joining it. - -### Opt-self-job - -Compare the row's `job_id` to the worker's own. - -```python -def _own_job_id() -> str | None: - env = submitit.JobEnvironment() - return env.job_id if env.activated() else None -``` - -The ids compare exactly: the driver stores `str(job.job_id)` = `"_"` -(`submitit/slurm/slurm.py:347`), in-job `SlurmJobEnvironment.job_id` is -`f"{SLURM_ARRAY_JOB_ID}_{SLURM_ARRAY_TASK_ID}"`, and `clean_env()` at submission -stops a sub-job inheriting the outer's SLURM vars. - -**Fails `UC-nested-chain`** — the row carries one job id, the driver's: - -``` -driver row{job=A} submit A -job A row.job=A == own A -> proceeds submit B -job B row.job=A != own B -> wait on A A waits on B <- deadlock -``` - -- Repointing the row per level would fix it and break `P-ancestor-row-intact`: once - the innermost job ends, the cell reads dead under still-running ancestors. -- Carrying the ancestor *chain* of ids is propagation, i.e. `Opt-inherited-claims`. - So depth-1 is the ceiling, and it cannot be lifted within `P-propagation-free`. -- Slurm-only. PID ancestry is no substitute: stdlib gives one level - (`os.getppid()`), and `LocalProcess` workers are grandchildren through submitit's - subprocess. -- Needs its clause in `claim()`, `pre_owned` and `WorkerInfo.wait`, not just the wait - loop — otherwise `P-no-livelock`. - -### Opt-owner-column - -The `owner_token` column sketched in `inflight-registry.md`, matched by prefix so -siblings stay mutually exclusive. - -- Breaks `P-ancestor-row-intact` exactly as `Opt-self-job` does: the descendant's - `record_worker_info` repoints the ancestor's row at the sub-job. -- Pays a schema change and still needs propagation, for guarantees - `Opt-inherited-claims` already gives. - -### Opt-combined - -`Opt-inherited-claims` as the mechanism, `Opt-self-job` as backstop, as one predicate -wherever ownership is decided: - -```python -def _is_own(info: WorkerInfo, cell: tuple[str, str]) -> bool: - return ( - info.pid == os.getpid() - or cell in _inherited.get() - or (info.job_id is not None and info.job_id == _own_job_id()) - ) -``` - -The backstop covers depth 1 on Slurm only, so it guards a narrower band than it first -appears: a propagation gap at depth ≥ 2 still deadlocks. - -## Touch points - -1. `CoordinationInfo.inherited`, surviving `ComputeBatch.__getstate__`. -2. `ComputeBatch.run_and_cache` binds a `ContextVar` to `inherited | own claim` for - the duration — a `ContextVar` not a module global, since `ThreadPool` workers - carry different sets in one process. -3. `Backend._prepare` seeds `inherited` from that `ContextVar`, so depth ≥ 2 - accumulates. -4. `Backend._claim` subtracts inherited cells from the `inflight_session` request. -5. `Backend._clear_caches` skips cancellation for inherited cells (`UC-tail-force`). -6. `_SubmititBackend._execute` (`backends.py:772`) and `_PoolBackend._submit_pool` - (`backends.py:974`) intersect their `uids=batch.items.uids` with `claim.uids`, so - an inherited cell keeps pointing at the ancestor's job. -7. `inflight.py` unchanged. - -`UC-thread-tail` never pickles, but needs nothing extra: threads keep `info.claim` -live, so step 2 reads it directly and the same `ContextVar` carries it. - -## Open questions - -- **Keep or drop the PID clause?** Dropping satisfies `P-own-narrow` and closes - `UC-sibling-sessions`, but puts every re-entrant case on propagation. Keeping is - deadlock-safer and leaves the duplicate-work looseness in place. -- **Is a depth-1 Slurm backstop worth `Opt-self-job`'s clauses?** It cannot cover - `UC-nested-chain`, so it guards only the shallow propagation gap. -- **Second allocation.** The tail still requests its own job while the outer holds - one. Accepted: the deadlock goes, the double-booking stays. From f938d20977204c3698a7f9c6dc18815b6de66d8c Mon Sep 17 00:00:00 2001 From: Jeremy Rapin Date: Mon, 14 Sep 2026 16:43:55 +0200 Subject: [PATCH 3/3] changelog --- CHANGELOG.md | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 51095e13..eefdd229 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,7 +2,9 @@ ## [Unreleased] -- `DiscriminatedModel`: optimized look-up. +- `DiscriminatedModel`: optimized look-up. [#313] +- `steps`: fixed nested infra claim deadlock. [#323] + ## 0.5.29 - 26-07-28