aboutsummaryrefslogtreecommitdiffziptar.gz
diff options
context:
space:
mode:
authorgodosa <godosa@godosa.eu>2026-10-07 20:09:02 +0200
committergodosa <godosa@godosa.eu>2026-10-07 20:09:02 +0200
commitf2135607b20931c65e88fede53014a1346fa54ee (patch)
treea32c77e07d317f9e662321934717f75f244a1bd9
parentaedabd0983c53edea462d41a7dd7ce511f59465d (diff)
downloadworkflow-f2135607b20931c65e88fede53014a1346fa54ee.tar.gz
workflow-f2135607b20931c65e88fede53014a1346fa54ee.zip
res: queue coalesces same-lock gates on descendant commitsHEADmaster
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JEAjUkQRrCYX5MZhWxdtj2
-rw-r--r--CHANGES.md1
-rw-r--r--docs/resource-ledger.md16
-rw-r--r--tests/test_res.py68
-rw-r--r--tests/test_res_io.py69
-rwxr-xr-xwf.py4
-rw-r--r--wf_res.py59
-rw-r--r--wflib/res.py53
7 files changed, 252 insertions, 18 deletions
diff --git a/CHANGES.md b/CHANGES.md
index a25df8b..f55a7e4 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -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()
diff --git a/wf.py b/wf.py
index 052ab63..2f75a20 100755
--- a/wf.py
+++ b/wf.py
@@ -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:
diff --git a/wf_res.py b/wf_res.py
index 9e90661..54987e9 100644
--- a/wf_res.py
+++ b/wf_res.py
@@ -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 '([^']+)'")