diff options
| author | godosa <godosa@godosa.eu> | 2026-10-07 19:59:24 +0200 |
|---|---|---|
| committer | godosa <godosa@godosa.eu> | 2026-10-07 19:59:24 +0200 |
| commit | ea8bcc80ccbb3b7556d226d391189312fdd07321 (patch) | |
| tree | 723e92b5c5d8d39f96c51942ee7d93094d3e360f /wflib/res.py | |
| parent | e6f7306c472a7589d08f4098e476d54bffd1f761 (diff) | |
| download | workflow-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/res.py')
| -rw-r--r-- | wflib/res.py | 35 |
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 |
