diff options
| author | godosa <godosa@godosa.eu> | 2026-10-07 20:09:02 +0200 |
|---|---|---|
| committer | godosa <godosa@godosa.eu> | 2026-10-07 20:09:02 +0200 |
| commit | f2135607b20931c65e88fede53014a1346fa54ee (patch) | |
| tree | a32c77e07d317f9e662321934717f75f244a1bd9 | |
| parent | aedabd0983c53edea462d41a7dd7ce511f59465d (diff) | |
| download | workflow-f2135607b20931c65e88fede53014a1346fa54ee.tar.gz workflow-f2135607b20931c65e88fede53014a1346fa54ee.zip | |
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 | 16 | ||||
| -rw-r--r-- | tests/test_res.py | 68 | ||||
| -rw-r--r-- | tests/test_res_io.py | 69 | ||||
| -rwxr-xr-x | wf.py | 4 | ||||
| -rw-r--r-- | wf_res.py | 59 | ||||
| -rw-r--r-- | wflib/res.py | 53 |
7 files changed, 252 insertions, 18 deletions
@@ -1,4 +1,5 @@ # Changes (newest first) +- 2026-10-07 `wf res` queue coalesces gates: a queued locked job on commit C (`--commit REV`, default a hex word of a locked job's title, `gate 1a2b3c4`) supersedes same-lock queued jobs on ancestors of C (they end `superseded by r-N`, it keeps the earliest turn, ledger `covers`); one on an ancestor of a queued job is covered, not queued. A batch of N commits → ≤ 1 queued gate per lock. `wf res wait` follows superseded ids; a red coalesced job names the range and bisect. Gate red procedure mentions it. Projects: a gate script that greps its own log by sha should `wf res wait` the printed id (covered/superseded logs are never written). - 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. diff --git a/docs/resource-ledger.md b/docs/resource-ledger.md index 4417517..a3ba3f8 100644 --- a/docs/resource-ledger.md +++ b/docs/resource-ledger.md @@ -157,6 +157,8 @@ id, systemd failure) · 2 usage. `gate` locks `gate` unasked (2026-10-06: two hand-started gates overlapped in one gate checkout). Held → exit 3 `busy: lock 'KEY' held by r-N "title" (ETA ~HH:MM|queued): …; --queue waits for it (--force does not override a lock)`; `--queue` → queued `(lock held by r-N)`. Ledger field `lock` = `KEY@<main tree>`. +- `--commit REV`: the commit the job tests (ledger `commit`, full sha); default for a locked job: a 7–40 hex word + of the title that resolves to a commit in the cwd (`gate 1a2b3c4`). Coalescing, §4.2. - The command reaches the job verbatim (no systemd `%` specifier expansion: `--format='%h %s'` is safe). ### 4.2 Queue @@ -171,6 +173,18 @@ batch waiting for its own gate never deadlocks (direct run and queue alike; not 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. +Coalescing (ruling 2026-10-07, t-coalesce-batch-gates; chosen over "batch queues one gate at its end": works for +every caller, batch or not, needs no batch-end hook): when a locked job with a `commit` is queued, same-lock queued +entries are compared by ancestry (`git merge-base --is-ancestor` in the cwd). Queued entries on ancestor commits end +`superseded by r-N` (ledger `superseded_by`); the new entry takes the earliest queued time of those (keeps its turn) +and lists their commits in `covers` (oldest queued first). A queued entry on the same or a descendant commit +already exists → nothing is queued: `r-M queued, position … ; covers <sha7> (queued on descendant <sha7>, nothing +new queued): wf res wait r-M` (r-M gains the commit in `covers`). Running entries and unrelated commits (side +branches) are never touched. So a batch of N commits yields ≤ 1 queued gate per lock. `wf res wait` on a +superseded id prints `r-N superseded by r-M …; waiting for it` and follows the chain. Done line of a coalesced +entry: `…; covers a,b`; red (rc ≠ 0): `; red: culprit is any commit in <oldest>^..<sha> -> bisect (git bisect +start <sha> <oldest>^, rerun the job per step) or name that range in the P0 fix task`. + ### 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 @@ -191,7 +205,7 @@ With id: that entry in full, including `rc`, `peak_gb`, `why` for done entries. ### 4.4 `wf res wait r-N [--timeout 2h]` Polls every 15 s (each poll is a normal locked call, so it also prunes and starts queued work) until the entry -is `done`; prints `r-N done rc=0 peak 9.4 GB in 37 min`. The throttle warning (§4.3) goes to stderr once. Exit 0 when done (whatever `rc`), 1 on timeout. If +is `done`; prints `r-N done rc=0 peak 9.4 GB in 37 min` (coalesced: superseded ids followed, §4.2). The throttle warning (§4.3) goes to stderr once. Exit 0 when done (whatever `rc`), 1 on timeout. If reaped, `status r-N` answers from the ledger. ### 4.5 `wf res release r-N [--stop]` diff --git a/tests/test_res.py b/tests/test_res.py index e4dfb1c..5f9ac92 100644 --- a/tests/test_res.py +++ b/tests/test_res.py @@ -566,10 +566,6 @@ class Units(unittest.TestCase): "alias claude='systemd-run --user --scope --quiet --slice=agents.slice claude'") -if __name__ == "__main__": - unittest.main() - - class Throttle(unittest.TestCase): def test_psi(self): text = "some avg10=41.50 avg60=33.20 avg300=12.00 total=99\nfull avg10=30.00 avg60=25.00 avg300=9.00 total=88\n" @@ -806,3 +802,67 @@ class TaskFit(unittest.TestCase): self.assertEqual(R.batch_fit(4, 59, 30), 1) self.assertEqual(R.batch_fit(4, 29, 30), 0) self.assertEqual(R.batch_fit(4, 0, 30), 0) + + +class Coalesce(unittest.TestCase): + """Batch r-740/744/757: one queued gate per lock on the newest commit (ruling 2026-10-07: the queue coalesces).""" + # history a ← b ← c (c newest); x on a side branch of a + ANC = {("a", "b"), ("a", "c"), ("b", "c"), ("a", "x")} + + def anc(self, x, y): + return (x, y) in self.ANC + + def gate(self, id, commit, at, lock="gate@/p"): + e = queued(id, 10.0, at=at) + e.lock, e.commit = lock, commit + return e + + def test_descendant_supersedes_queued_ancestors(self): + g1, g2 = self.gate("r-740", "a", T(12)), self.gate("r-744", "b", T(13)) + run = running("r-730", title="gate z") + run.lock, run.commit = "gate@/p", "z" + led = R.Ledger(entries=[run, g1, g2]) + new = self.gate("r-757", "c", T(14)) + cov, gone = R.coalesce(led, new, self.anc) + self.assertEqual((cov, [e.id for e in gone]), (None, ["r-740", "r-744"])) # running r-730 untouched + led.entries.append(new) + R.supersede(new, gone, NOW) + self.assertEqual([(e.state, e.superseded_by, e.why) for e in (g1, g2)], + [("done", "r-757", "superseded by r-757")] * 2) + self.assertEqual((new.covers, new.queued), (["a", "b"], T(12))) # keeps the oldest turn + self.assertEqual([e.id for e in R._queue(led)], ["r-757"]) + self.assertEqual(run.state, "running") + + def test_ancestor_is_covered_not_queued(self): + g = self.gate("r-757", "c", T(14)) + led = R.Ledger(entries=[g]) + for c in ("b", "c"): + cov, gone = R.coalesce(led, self.gate("r-760", c, T(15)), self.anc) + self.assertIs(cov, g) + self.assertEqual(gone, []) + R.cover(g, "b") + R.cover(g, "c") # own commit not listed + self.assertEqual(g.covers, ["b"]) + + def test_unrelated_other_lock_or_no_commit_kept(self): + led = R.Ledger(entries=[self.gate("r-1", "x", T(12)), self.gate("r-2", "a", T(12), lock="gate@/q"), + self.gate("r-3", "", T(12))]) + self.assertEqual(R.coalesce(led, self.gate("r-9", "c", T(14)), self.anc), (None, [])) + self.assertEqual(R.coalesce(led, self.gate("r-9", "", T(14)), self.anc), (None, [])) + + def test_red_coalesced_gate_names_range(self): + e = running("r-757", title="gate c") + e.commit, e.covers = "c" * 40, ["a" * 40, "b" * 40] + e.state, e.rc, e.ended = "done", 1, T(14, 30) + line = R.done_line(e) + self.assertIn("covers aaaaaaa,bbbbbbb", line) + self.assertIn("aaaaaaa^..ccccccc", line) + self.assertIn("bisect", line) + e.rc = 0 + self.assertNotIn("bisect", R.done_line(e)) + e.covers = [] + self.assertNotIn("covers", R.done_line(e)) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_res_io.py b/tests/test_res_io.py index c30b21f..3d5b789 100644 --- a/tests/test_res_io.py +++ b/tests/test_res_io.py @@ -27,6 +27,7 @@ class Fake: self.scopes = {} # name → {prop: value} for agent session scopes self.manager_env = "PATH=/usr/bin\nHOME=/h/u\n" self.git = {} # git subcommand (branch | rev-parse) → stdout + self.real_git = False # True: git calls run for real (cwd = a synthetic repo) def __call__(self, argv): self.calls.append(list(argv)) @@ -34,6 +35,8 @@ class Fake: if key in self.fail: return subprocess.CompletedProcess(argv, 1, "", self.fail[key] + "\n") out = "" + if argv[0] == "git" and self.real_git: + return subprocess.run(argv, capture_output=True, text=True) if argv[:3] == ["systemctl", "--user", "show"]: blocks = [] for name in argv[5:]: @@ -685,6 +688,72 @@ class Robust(IOBase): self.assertEqual((code, err), (1, "wf: [Errno 13] Permission denied: '/proc/meminfo'\n")) +class CoalesceIO(IOBase): + """Batch of commits a ← b ← c, each queueing its gate (real repo): one queued gate per lock, on c.""" + + def setUp(self): + super().setUp() + repo = self.box.tmp / "proj" + git = lambda *a: subprocess.run(["git", "-C", str(repo), *a], check=True, capture_output=True, text=True).stdout + git("init", "-q", "-b", "master") + self.sha = {} + for c in ("a", "b", "c"): + git("-c", "user.name=t", "-c", "user.email=t@t", "commit", "-q", "--allow-empty", "-m", c) + self.sha[c] = git("rev-parse", "HEAD").strip() + git("checkout", "-q", "-b", "side", self.sha["a"]) + git("-c", "user.name=t", "-c", "user.email=t@t", "commit", "-q", "--allow-empty", "-m", "x") + self.sha["x"] = git("rev-parse", "HEAD").strip() + self.box.fake.real_git = True + + def gate(self, c, *extra): + return self.box.wf("run", "--queue", "--mem", "2G", "--for", "30m", "--title", f"gate {self.sha[c][:7]}", + *extra, "--", "x") + + def test_batch_yields_one_queued_gate(self): + self.box.wf("run", "--mem", "2G", "--for", "30m", "--title", "gate running", "--", "x") # r-1 holds the lock + self.assertEqual(self.gate("a")[0], 0) + self.box.tick(60) + self.gate("b") + self.box.tick(60) + code, out, _ = self.gate("c") + self.assertEqual(code, 0) + self.assertIn("r-4 queued, position 1 (lock held by r-1)", out) + self.assertIn("supersedes r-3", out) + led = self.box.ledger() + self.assertEqual([(e.id, e.commit) for e in led.entries if e.state == "queued"], [("r-4", self.sha["c"])]) + self.assertEqual(led.get("r-4").covers, [self.sha["a"], self.sha["b"]]) + self.assertEqual(led.get("r-4").queued, led.get("r-2").queued) + self.assertEqual([led.get(i).superseded_by for i in ("r-2", "r-3")], ["r-3", "r-4"]) # chain; wait follows it + code, out, _ = self.gate("b") # late ancestor: covered, nothing new + self.assertEqual(code, 0) + self.assertTrue(out.startswith("r-4 queued"), out) + self.assertIn(f"covers {self.sha['b'][:7]}", out) + self.assertEqual(len([e for e in self.box.ledger().entries if e.state == "queued"]), 1) + code, out, _ = self.gate("x") # side branch: own gate + self.assertTrue(out.startswith("r-5 queued"), out) + code, out, _ = self.box.wf("run", "--queue", "--mem", "2G", "--for", "30m", "--title", "gate", + "--commit", "master", "--", "x") + self.assertTrue(out.startswith("r-4 queued"), out) # --commit REV, no hex in title + self.assertNotEqual(self.box.wf("run", "--queue", "--mem", "2G", "--for", "30m", "--title", "gate", + "--commit", "nope", "--", "x")[0], 0) + + def test_wait_follows_superseded_and_red_names_range(self): + self.box.wf("run", "--mem", "2G", "--for", "30m", "--title", "gate running", "--", "x") + self.gate("a") + self.gate("c") + self.box.fake.units["wf-r-1.service"] = ("inactive", None) + self.box.wf("status") # r-3 starts + logs = self.box.env.state / "logs" + logs.mkdir(parents=True, exist_ok=True) + (logs / "r-3.rc").write_text("1\n") + self.box.fake.units["wf-r-3.service"] = ("inactive", None) + code, out, _ = self.box.wf("wait", "r-2") + self.assertEqual(code, 0) + self.assertIn("r-2 superseded by r-3", out) + self.assertIn(f"r-3 done rc=1", out) + self.assertIn(f"{self.sha['a'][:7]}^..{self.sha['c'][:7]}", out) + + if __name__ == "__main__": unittest.main() @@ -1338,7 +1338,9 @@ GATE_RED_HELP = ( "\n- culprit = your diff (git show --stat <sha>) explains the failure -> fix now (new commit, wf merge) or hand back" "\n- not your diff (earlier task broke it; see git log of the failing area) -> wf add -p 0 --model <m> -e <effort>" ' --done "gate ALL GREEN" -b "Fix gate red <test>. <goal>" with body line "Steps: <failing test, culprit sha/id>"' - "\n- then report result: done+gate-red <culprit sha or id> <fix id> (your task is done, the lane goes on with the fix)") + "\n- then report result: done+gate-red <culprit sha or id> <fix id> (your task is done, the lane goes on with the fix)" + "\n- coalesced gate (wf res wait prints 'covers …; red: culprit is any commit in A^..B'): bisect that range for the" + " culprit, else the P0 fix task's Steps name the range") def cmd_finish(args) -> int: @@ -1,7 +1,7 @@ """wf res — shared memory/CPU ledger for agent jobs (files, /proc, systemd, printing). Pure logic: wflib/res.py. Spec: docs/resource-ledger.md. - run --mem 10G [--cpus N] --for 40m --title T [--queue|--force] [--by NAME] -- CMD… start a job or get `busy` (exit 3) + run --mem 10G [--cpus N] --for 40m --title T [--queue|--force] [--by NAME] [--commit REV] -- CMD… start a job or get `busy` (exit 3) status [ID] [--json] · wait ID · release ID [--stop] · note --mem 3G --for 20m [--force] [--by NAME] TITLE game on [--for 4h] | off · clean [--yes] · timer on|off · tick · shell-init """ @@ -397,6 +397,37 @@ def _new_probe(env, args) -> res.Entry: by=owner_by(env, getattr(args, "by", "") or "")) +def job_commit(env: Env, args, lock: str) -> str: + """Full sha the job tests: --commit REV, else (locked job) a hex word of the title that names a commit; else ''.""" + def sha(rev: str) -> str: + out = _git(env, "rev-parse", "--verify", "--quiet", f"{rev}^{{commit}}") + return out if re.fullmatch(r"[0-9a-f]{40}", out) else "" + rev = getattr(args, "commit", "") or "" + if rev: + return sha(rev) or _raise(res.ResError(f"--commit {rev}: not a commit here")) + m = res.COMMIT_TITLE_RE.search(args.title) if lock else None + return sha(m.group()) if m else "" + + +def _raise(err: Exception): + raise err + + +def enqueue(env, led, facts, args, mem, est, cmd, caller_env, lock, commit) -> tuple[res.Entry, str]: + """Queue a job; same-lock queued entries coalesce by commit (res.coalesce): (entry waited for, note).""" + e = _new_entry(env, led, facts, args, mem, est, "queued") + e.cmd, e.env, e.lock, e.commit = cmd, caller_env, lock, commit + anc = lambda a, b: env.run(["git", "-C", env.cwd(), "merge-base", "--is-ancestor", a, b]).returncode == 0 + cov, gone = res.coalesce(led, e, anc) + if cov: + led.next -= 1 # never used + res.cover(cov, commit) + return cov, f"; covers {commit[:7]} (queued on descendant {cov.commit[:7]}, nothing new queued): wf res wait {cov.id}" + led.entries.append(e) + res.supersede(e, gone, facts.now) + return e, "".join(f"; supersedes {g.id} ({g.commit[:7]})" for g in gone) + + 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" @@ -412,6 +443,7 @@ def cmd_run(env, cfg, args) -> int: shown = env.run(["systemctl", "--user", "show-environment"]) caller_env = res.env_diff(env.environ(), res.parse_show_environment(shown.stdout if shown.returncode == 0 else "")) lock = res.lock_for(args.title, getattr(args, "lock", ""), str(main_tree(env))) + commit = job_commit(env, args, lock) with locked(env) as led: facts = session(env, cfg, led) hint(env, led, args.title, mem, est) @@ -420,12 +452,12 @@ def cmd_run(env, cfg, args) -> int: if not fail and holder and not args.queue: busy = res.lock_busy_line(holder, facts) elif not fail and holder: - e = _new_entry(env, led, facts, args, mem, est, "queued") - e.cmd, e.env, e.lock = cmd, caller_env, lock - led.entries.append(e) + e, note = enqueue(env, led, facts, args, mem, est, cmd, caller_env, lock, commit) when = res.queue_estimate(cfg, led, facts, e) - out = (f"{e.id} queued, position {res.position(led, e)} (lock held by {holder.id}); est. start " - + (f"~{res.hhmm(when)}" if when else "unknown") + f"; cancel: wf res release {e.id}") + held = res.lock_holder(led, lock) + held = f" (lock held by {held.id})" if held and held is not e else "" + out = (f"{e.id} queued, position {res.position(led, e)}{held}; est. start " + + (f"~{res.hhmm(when)}" if when else "unknown") + f"; cancel: wf res release {e.id}{note}") elif not fail: q_gb, q_cpus = res.queued_total(led) room = res.force_room(cfg, led, facts) if args.force else res.budget(cfg, led, facts) @@ -437,18 +469,16 @@ def cmd_run(env, cfg, args) -> int: 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 + e.cmd, e.env, e.lock, e.commit = cmd, caller_env, lock, commit fail = start(env, e, facts.now) if not fail: led.entries.append(e) out = f"{e.id} started; log {e.log}; ETA ~{res.hhmm(e.eta())} (estimate: not killed when over)" elif args.queue: - e = _new_entry(env, led, facts, args, mem, est, "queued") - e.cmd, e.env, e.lock = cmd, caller_env, lock - led.entries.append(e) + e, note = enqueue(env, led, facts, args, mem, est, cmd, caller_env, lock, commit) when = res.queue_estimate(cfg, led, facts, e) out = (f"{e.id} queued, position {res.position(led, e)}; est. start " - + (f"~{res.hhmm(when)}" if when else "unknown") + f"; cancel: wf res release {e.id}") + + (f"~{res.hhmm(when)}" if when else "unknown") + f"; cancel: wf res release {e.id}{note}") elif args.force: busy = force_busy(cfg, led, facts, mem, args.cpus, room) else: @@ -520,6 +550,10 @@ def cmd_wait(env, cfg, args) -> int: with locked(env) as led: facts = session(env, cfg, led) e = led.get(args.id) + if e.superseded_by: # coalesced into a gate on a descendant commit: wait for that one + print(f"{e.id} superseded by {e.superseded_by} (same lock, descendant commit); waiting for it") + args.id = e.superseded_by + continue line, state = (res.done_line(e) if e.state == "done" else None), e.state hot = res.throttle_lines(res.Ledger(entries=[e]), facts) if hot and not warned: # once per wait @@ -806,6 +840,9 @@ def parser() -> argparse.ArgumentParser: sp.add_argument("--lock", default="", metavar="KEY", help="one queued/running job per KEY in this project (main tree " "and its lane worktrees) at a time, e.g. a shared checkout; a title starting with 'gate' locks " "'gate' unasked; held → exit 3, --queue waits for it, --force does not override it") + sp.add_argument("--commit", default="", metavar="REV", help="commit the job tests (default for a locked job: a " + "hex word of the title, e.g. 'gate 1a2b3c4'); a queued job on a descendant commit replaces queued " + "same-lock jobs on its ancestors (wf res wait follows), one on an ancestor is covered, not queued") sp.add_argument("cmd", nargs=argparse.REMAINDER, metavar="-- CMD …") sp = cmd("status", cmd_status, "reservations, budget, game mode, warnings") diff --git a/wflib/res.py b/wflib/res.py index e110daf..9add386 100644 --- a/wflib/res.py +++ b/wflib/res.py @@ -187,6 +187,9 @@ class Entry: env: dict = field(default_factory=dict) # caller variables the user manager lacks (env_diff) by: dict = field(default_factory=dict) # who started it (owner_by): name, task, batch, address lock: str = "" # '<key>@<main tree>': one queued/running job per lock (lock_for) + commit: str = "" # full sha the job tests (--commit / gate title): coalesce key + covers: list = field(default_factory=list) # commits of superseded/covered same-lock entries, oldest queued first + superseded_by: str = "" # id of the queued entry on a descendant commit that replaced this one extra: dict = field(default_factory=dict, repr=False, compare=False) # unknown keys of a newer version: kept on rewrite def eta(self) -> dt.datetime | None: @@ -474,6 +477,54 @@ def lock_holder(led: Ledger, lock: str) -> Entry | None: return next((e for e in mine if e.state == "running"), None) or next(iter(_queue_of(mine)), None) +COMMIT_TITLE_RE = re.compile(r"\b[0-9a-f]{7,40}\b") # 'gate 1a2b3c4': the commit a locked job tests + + +def coalesce(led: Ledger, new: Entry, is_ancestor) -> tuple[Entry | None, list[Entry]]: + """Same-lock queued entries vs a new queued entry on commit new.commit (ruling 2026-10-07: the queue coalesces, + one queued gate per lock and line of history). Returns (cover, gone): cover = a queued entry on the same or a + descendant commit (the new one is redundant, never queued); else gone = queued entries on ancestor commits + (the new one replaces them). is_ancestor(a, b) = a is an ancestor of b. Running entries are never touched.""" + if not (new.lock and new.commit): + return None, [] + same = [e for e in _queue_of(led.entries) if e is not new and e.lock == new.lock and e.commit] + for e in same: + if e.commit == new.commit or is_ancestor(new.commit, e.commit): + return e, [] + return None, [e for e in same if is_ancestor(e.commit, new.commit)] + + +def supersede(new: Entry, gone: list[Entry], now: dt.datetime) -> None: + """new replaces gone (queue order): they end 'superseded by new', new inherits their commits and the earliest + queue place (a coalesced gate never loses its turn).""" + for e in gone: + _finish(e, now, f"superseded by {new.id}") + e.superseded_by = new.id + new.covers = _uniq(new.covers + e.covers + [e.commit], new.commit) + if e.queued and (new.queued is None or e.queued < new.queued): + new.queued = e.queued + + +def cover(e: Entry, commit: str, covers: list[str] = ()) -> None: + """e (queued, descendant commit) also stands for commit.""" + e.covers = _uniq(e.covers + list(covers) + [commit], e.commit) + + +def _uniq(commits: list[str], own: str) -> list[str]: + return [c for c in dict.fromkeys(commits) if c and c != own] + + +def covers_text(e: Entry) -> str: + """' covers a1,b2 (red: culprit is one of them or <sha>; bisect …)' for a coalesced entry, else ''.""" + if not e.covers: + return "" + out = f"; covers {','.join(c[:7] for c in e.covers)}" + if e.state == "done" and e.rc not in (0, None) and not e.superseded_by: + out += (f"; red: culprit is any commit in {e.covers[0][:7]}^..{e.commit[:7]} -> bisect (git bisect start " + f"{e.commit[:7]} {e.covers[0][:7]}^, rerun the job per step) or name that range in the P0 fix task") + return out + + def _queue_of(entries: list[Entry]) -> list[Entry]: return sorted((e for e in entries if e.state == "queued"), key=lambda e: (e.queued, int(e.id[2:]))) @@ -611,7 +662,7 @@ def done_line(e: Entry) -> str: minutes = round((e.ended - e.started).total_seconds() / 60) if e.ended and e.started else 0 peak = fmt_gb(e.peak_gb) if e.peak_gb is not None else "?" why = f" ({e.why})" if e.why and e.why != "exited" else "" - return f"{e.id} done rc={rc_text(e.rc)} peak {peak} in {minutes} min{why}" + return f"{e.id} done rc={rc_text(e.rc)} peak {peak} in {minutes} min{why}{covers_text(e)}" JOURNAL_RESULT = re.compile(r"Failed with result '([^']+)'") |
