aboutsummaryrefslogtreecommitdiffziptar.gz
path: root/wf_res.py
diff options
context:
space:
mode:
Diffstat (limited to 'wf_res.py')
-rw-r--r--wf_res.py59
1 files changed, 48 insertions, 11 deletions
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")