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 | |
| 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
| -rw-r--r-- | CHANGES.md | 1 | ||||
| -rw-r--r-- | docs/resource-ledger.md | 6 | ||||
| -rw-r--r-- | tests/test_res.py | 57 | ||||
| -rw-r--r-- | tests/test_res_io.py | 7 | ||||
| -rw-r--r-- | wf_res.py | 10 | ||||
| -rw-r--r-- | wflib/res.py | 35 |
6 files changed, 114 insertions, 2 deletions
@@ -1,4 +1,5 @@ # Changes (newest first) +- 2026-10-07 `wf res`: a job started from inside a live job (`by.batch`, e.g. a batch worker's gate) borrows the parent's claim (mem + cpus, lent once) when fitting, so a batch waiting on its own queued gate no longer deadlocks; a queued entry whose batch ended is dropped at prune (`r-N dropped (batch r-M gone)`), never heads the queue. Projects: nothing. - 2026-10-07 Lane cloud: one stuck task no longer stops it: `wf orch post <id> cloud` with awaiting / handback / lost parks only that task (awaiting: park as local lanes; handback / lost: `Cloud: no` + status cleared + note, local lanes may retry with pull's Recovery note), prints `alert: …; lane cloud keeps picking`, then the next cloud pick (never the same id). Projects: nothing. - 2026-10-07 Owner-absent default: workers never wait; owner-only step → `wf add -s human` + `wf set <id> --after <h-id>` + status clear, result `needs-owner <h-id>`. `wf orch post` parks it (After: the h-task, claim cleared, alert), lane keeps picking; owner's `wf done <h-id>` makes it pickable again. Worker def, shared CLAUDE.md, wf-orchestrate, wf-pilot, batch prompt, docs. Projects: nothing. - 2026-10-07 One stuck task no longer freezes its lane: `wf orch post` with awaiting / handback / post-check-red (local lanes; lane cloud unchanged: pull re-queues) parks only that task (blocked on the a-id named in the result, else `Sessions: owner` + an `orch …` note; handback also clears its claim), commits, prints `alert: <id> <outcome> → tell the owner (…); lane <lane> keeps picking`, then the lane's next pick. `stop lane` only for none / stop file / push-failed / wip / no report. Batch prompt, wf-orchestrate, wf-pilot, docs: keep picking, never wait on the owner while tasks are pickable. Projects: nothing. diff --git a/docs/resource-ledger.md b/docs/resource-ledger.md index 17a0443..4417517 100644 --- a/docs/resource-ledger.md +++ b/docs/resource-ledger.md @@ -165,6 +165,12 @@ Strict FIFO by `queued`: the head starts when it fits; entries behind a waiting starved); an entry whose lock a running job holds is skipped (it neither starts nor blocks the rest). Starting happens inside every `wf res` call (after prune) and every timer tick. A started queued job runs with the cwd and argv it was queued with; its ETA counts from the actual start. +Batch loan: a job started from inside a live job (`by.batch` = its `WF_RES_ID`, e.g. a `wf batch` worker's gate) +fits against budget + the parent's claim (mem and cpus) minus what the parent's other running children hold, so a +batch waiting for its own gate never deadlocks (direct run and queue alike; not with `--force`). A queued entry +whose parent job ended (done or gone) is dropped at prune (`r-N dropped (batch r-M gone)`, why `owner gone`): +nobody waits for it, and it must not head the queue. + ### 4.3 `wf res status [r-N] [--json]` Without id: one line per entry (id, project, title, state, reserved / used / peak GB, cpus, started, ETA or diff --git a/tests/test_res.py b/tests/test_res.py index 6c92549..e4dfb1c 100644 --- a/tests/test_res.py +++ b/tests/test_res.py @@ -385,6 +385,63 @@ class Queue(unittest.TestCase): led = R.Ledger(entries=[queued("r-1", 20)]) self.assertIsNone(R.queue_estimate(R.Config(), led, facts(), led.get("r-1"))) +class BatchLoan(unittest.TestCase): + """Real case: batch r-741 claims 6.3 GB, its gate r-757 (10 GB) queued, budget 7.6 → must not deadlock.""" + + def setUp(self): + self.batch = running("r-741", title="wf batch", mem=6.3, cpus=1) + # 21.9 − 6 reserve − 2 headroom − 6.3 unused claim = 7.6 GB budget + self.f = facts(available=21.9, units={"wf-r-741.service": R.Unit(True, 0.0)}) + + def gate(self, id="r-757", batch="r-741", mem=10.0, at=T(13, 0)): + e = queued(id, mem, at=at) + e.by = {"batch": batch} if batch else {} + return e + + def test_budget_is_real_case(self): + led = R.Ledger(entries=[self.batch]) + self.assertAlmostEqual(R.budget(R.Config(), led, self.f)[0], 7.6) + + def test_own_gate_borrows_batch_claim(self): + led = R.Ledger(entries=[self.batch, self.gate()]) + self.assertEqual([e.id for e in R.to_start(R.Config(), led, self.f)], ["r-757"]) + self.assertIsNotNone(R.queue_estimate(R.Config(), led, self.f, led.get("r-757"))) + + def test_foreign_job_does_not_borrow(self): + led = R.Ledger(entries=[self.batch, self.gate(batch="")]) + self.assertEqual(R.to_start(R.Config(), led, self.f), []) + led = R.Ledger(entries=[self.batch, self.gate(batch="r-1")]) # other batch, not live + self.assertEqual(R.to_start(R.Config(), led, self.f), []) + + def test_claim_lent_once(self): + # 7.6 + 6.3 = 13.9: first 10 GB gate fits, second does not (loan used up, 3.9 left) + led = R.Ledger(entries=[self.batch, self.gate("r-757"), self.gate("r-758", at=T(13, 1))]) + self.assertEqual([e.id for e in R.to_start(R.Config(), led, self.f)], ["r-757"]) + led.get("r-757").state = "running" + led.get("r-757").unit = "wf-r-757.service" + self.f.units["wf-r-757.service"] = R.Unit(True, 0.0) + self.assertEqual(R.loan(led, led.get("r-758")), (0.0, 0)) + self.assertEqual(R.to_start(R.Config(), led, self.f), []) + + def test_dead_batch_gate_dropped(self): + # r-740: gate of dead batch r-731 heads the queue; r-757 behind it must still start + dead = self.gate("r-740", batch="r-731", at=T(12, 0)) + led = R.Ledger(entries=[self.batch, dead, self.gate()]) + self.assertEqual(R.prune(led, self.f), ["r-740 dropped (batch r-731 gone)"]) + self.assertEqual((dead.state, dead.why), ("done", "owner gone")) + self.assertEqual([e.id for e in R.to_start(R.Config(), led, self.f)], ["r-757"]) + + def test_batch_ending_in_same_prune_drops_gate(self): + led = R.Ledger(entries=[self.batch, self.gate()]) + f = facts(units={"wf-r-741.service": R.Unit(False)}) + self.assertEqual(R.prune(led, f), ["r-741 exited rc=?", "r-757 dropped (batch r-741 gone)"]) + + def test_live_batch_gate_kept(self): + led = R.Ledger(entries=[self.batch, self.gate()]) + self.assertEqual(R.prune(led, self.f), []) + self.assertEqual(led.get("r-757").state, "queued") + + class Game(unittest.TestCase): def setUp(self): self.cfg = R.Config() diff --git a/tests/test_res_io.py b/tests/test_res_io.py index d91046a..16271e0 100644 --- a/tests/test_res_io.py +++ b/tests/test_res_io.py @@ -426,6 +426,13 @@ class QueueIO(IOBase): self.assertEqual((e.state, e.unit), ("running", "wf-r-2.service")) self.assertEqual(self.box.fake.ran("systemd-run")[-1][-2:], ["sh", "y"]) + def test_batch_gate_borrows_batch_claim(self): + self.box.wf("run", "--mem", "10G", "--for", "40m", "--title", "wf batch", "--", "x") + self.box.caller["WF_RES_ID"] = "r-1" # the gate is queued from inside batch r-1 + code, out, _ = self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "gate", "--queue", "--", "y") + self.assertEqual((code, out.split(";")[0]), (0, "r-2 started")) + self.assertEqual(self.box.ledger().get("r-2").by["batch"], "r-1") + def test_direct_run_does_not_jump_queue(self): self.box.wf("run", "--mem", "10G", "--for", "40m", "--title", "big", "--", "x") self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "two", "--queue", "--", "y") @@ -391,6 +391,12 @@ def _new_entry(env, led, facts, args, mem, est, state) -> res.Entry: by=owner_by(env, getattr(args, "by", "") or "")) +def _new_probe(env, args) -> res.Entry: + """Entry stand-in (no id) for capacity checks that depend on who starts it (res.loan).""" + return res.Entry(id="", project="", owner=0, title=args.title, mem_gb=0.0, cpus=0, est_min=0, state="queued", + by=owner_by(env, getattr(args, "by", "") or "")) + + def force_busy(cfg, led, facts, mem, cpus, room) -> str: """--force refused: not even the really free memory (beyond the user reserve) holds it.""" free = f"{res.fmt_gb(max(0.0, room[0]))}" if mem > room[0] + res.EPS else f"{max(0, room[1])} cpus" @@ -425,6 +431,10 @@ def cmd_run(env, cfg, args) -> int: room = res.force_room(cfg, led, facts) if args.force else res.budget(cfg, led, facts) if args.force: q_gb, q_cpus = 0.0, 0 + probe = _new_probe(env, args) + if not args.force: # a batch's own job borrows the batch's claim (res.loan) + l_gb, l_cpus = res.loan(led, probe) + room = (room[0] + l_gb, room[1] + l_cpus) if res.fits(mem + q_gb, args.cpus + q_cpus, room): e = _new_entry(env, led, facts, args, mem, est, "queued") e.cmd, e.env, e.lock = cmd, caller_env, lock 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 |
