aboutsummaryrefslogtreecommitdiffziptar.gz
path: root/wflib
diff options
context:
space:
mode:
authorgodosa <godosa@godosa.eu>2026-10-07 19:59:24 +0200
committergodosa <godosa@godosa.eu>2026-10-07 19:59:24 +0200
commitea8bcc80ccbb3b7556d226d391189312fdd07321 (patch)
tree723e92b5c5d8d39f96c51942ee7d93094d3e360f /wflib
parente6f7306c472a7589d08f4098e476d54bffd1f761 (diff)
downloadworkflow-ea8bcc80ccbb3b7556d226d391189312fdd07321.tar.gz
workflow-ea8bcc80ccbb3b7556d226d391189312fdd07321.zip
res: batch gate borrows batch claim; drop queued jobs of dead batch
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JEAjUkQRrCYX5MZhWxdtj2
Diffstat (limited to 'wflib')
-rw-r--r--wflib/res.py35
1 files changed, 33 insertions, 2 deletions
diff --git a/wflib/res.py b/wflib/res.py
index fd1fba9..e110daf 100644
--- a/wflib/res.py
+++ b/wflib/res.py
@@ -306,6 +306,10 @@ def prune(led: Ledger, facts: Facts) -> list[str]:
elif e.expires and now >= e.expires:
_finish(e, now, "expired")
out.append(f"{e.id} freed (expired)")
+ for e in led.entries: # a queued job of a batch that ended: nobody waits for it, never let it head the queue
+ if e.state == "queued" and parent_of(led, e) is None and (e.by or {}).get("batch"):
+ _finish(e, now, "owner gone")
+ out.append(f"{e.id} dropped (batch {e.by['batch']} gone)")
led.entries = [e for e in led.entries if not (e.state == "done" and e.ended and now - e.ended > DONE_KEEP)]
if led.game_until and now >= led.game_until:
led.game_until = None
@@ -355,6 +359,24 @@ def budget(cfg: Config, led: Ledger, facts: Facts) -> tuple[float, int]:
return round(mem, 6), facts.nproc - res_cpus - sum(e.cpus for e in claims)
+def parent_of(led: Ledger, e: Entry) -> Entry | None:
+ """The live job (e.g. wf batch) e was started from (by.batch = its WF_RES_ID), else None."""
+ pid = (e.by or {}).get("batch")
+ return next((p for p in led.entries if p.id == pid and p.state in ("running", "note")), None) if pid else None
+
+
+def loan(led: Ledger, e: Entry, lent: dict | None = None) -> tuple[float, int]:
+ """Room e borrows from its parent job's claim: a batch's gate runs inside the batch's reservation (the batch
+ waits for it, so a gate bigger than the free budget would otherwise deadlock). The claim is lent once: running
+ children of the parent and loans in `lent` (id → (gb, cpus), this pass) are subtracted."""
+ p = parent_of(led, e)
+ if p is None:
+ return 0.0, 0
+ kids = [k for k in led.entries if k is not e and k.state == "running" and (k.by or {}).get("batch") == p.id]
+ gb, cpus = (lent or {}).get(p.id, (0.0, 0))
+ return (max(0.0, p.mem_gb - gb - sum(k.mem_gb for k in kids)), max(0, p.cpus - cpus - sum(k.cpus for k in kids)))
+
+
def force_room(cfg: Config, led: Ledger, facts: Facts) -> tuple[float, int]:
"""Budget for --force: what the machine really has free beyond the user reserve and headroom; ledger
claims (unused reservations, notes) and the queue are ignored."""
@@ -625,14 +647,21 @@ def to_start(cfg: Config, led: Ledger, facts: Facts) -> list[Entry]:
running job (or one started in this pass) holds is skipped, not blocking the rest."""
b_gb, b_cpus = budget(cfg, led, facts)
out, held = [], {e.lock for e in led.entries if e.state == "running" and e.lock}
+ lent: dict[str, tuple[float, int]] = {}
for e in _queue(led):
if e.lock and e.lock in held:
continue
- if not fits(e.mem_gb, e.cpus, (b_gb, b_cpus)):
+ l_gb, l_cpus = loan(led, e, lent)
+ if not fits(e.mem_gb, e.cpus, (b_gb + l_gb, b_cpus + l_cpus)):
break
out.append(e)
held.add(e.lock)
- b_gb, b_cpus = b_gb - e.mem_gb, b_cpus - e.cpus
+ u_gb, u_cpus = min(e.mem_gb, l_gb), min(e.cpus, l_cpus) # the loan first, then the free budget
+ if u_gb or u_cpus:
+ p = (e.by or {})["batch"]
+ g0, c0 = lent.get(p, (0.0, 0))
+ lent[p] = (g0 + u_gb, c0 + u_cpus)
+ b_gb, b_cpus = b_gb - (e.mem_gb - u_gb), b_cpus - (e.cpus - u_cpus)
return out
@@ -641,6 +670,8 @@ def queue_estimate(cfg: Config, led: Ledger, facts: Facts, e: Entry) -> dt.datet
q = _queue(led)
ahead = q[:q.index(e) + 1]
b_gb, b_cpus = budget(cfg, led, facts)
+ l_gb, l_cpus = loan(led, e)
+ b_gb, b_cpus = b_gb + l_gb, b_cpus + l_cpus
holders, enough = needed(led, facts, sum(x.mem_gb for x in ahead) - b_gb, sum(x.cpus for x in ahead) - b_cpus)
if not enough:
return None