diff options
Diffstat (limited to 'wf_res.py')
| -rw-r--r-- | wf_res.py | 59 |
1 files changed, 48 insertions, 11 deletions
@@ -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") |
