aboutsummaryrefslogtreecommitdiffziptar.gz
path: root/wflib
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 /wflib
parentaedabd0983c53edea462d41a7dd7ce511f59465d (diff)
downloadworkflow-master.tar.gz
workflow-master.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
Diffstat (limited to 'wflib')
-rw-r--r--wflib/res.py53
1 files changed, 52 insertions, 1 deletions
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 '([^']+)'")