diff options
Diffstat (limited to 'wf_res.py')
| -rw-r--r-- | wf_res.py | 1208 |
1 files changed, 1208 insertions, 0 deletions
diff --git a/wf_res.py b/wf_res.py new file mode 100644 index 0000000..cc9a415 --- /dev/null +++ b/wf_res.py @@ -0,0 +1,1208 @@ +"""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) + 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 +""" +from __future__ import annotations + +import re +import argparse +import contextlib +import datetime as dt +import fcntl +import json +import os +import shlex +import shutil +import stat +import subprocess +import sys +import time +import tomllib +from dataclasses import dataclass +from pathlib import Path +from typing import Callable + +HERE = Path(__file__).resolve().parent +sys.path.insert(0, str(HERE)) + +from wflib import res # noqa: E402 + +LEDGER, LOCK = "resources.json", "resources.lock" +HISTORY = "resources-history.jsonl" # finished runs pruned from the ledger: reservation sizes (wf res hist) +ACTIVE = {"active", "activating", "deactivating", "reloading", "refreshing"} + + +@dataclass +class Env: + """Everything wf res touches outside its own state; tests swap the callables.""" + state: Path + config: Path + units: Path + scratch: Path + run: Callable[[list[str]], subprocess.CompletedProcess] + meminfo: Callable[[], str] + nproc: int + now: Callable[[], dt.datetime] + pid_alive: Callable[[int], bool] + owner: Callable[[], int] + root: Callable[[], Path] + cwd: Callable[[], str] + procs: Callable[[], list[tuple[int, int, str, str]]] # (pid, ppid, comm, cgroup) + sleep: Callable[[float], None] + tmp_used_gb: Callable[[], float] + stdin: Callable[[], str] = lambda: sys.stdin.read() + live_sessions: Callable[[], set[str]] = lambda: _live_sessions() + tmp_dirs: tuple = (Path("/tmp"), Path("/var/tmp")) # swept for litter (res.litter_victims) + cgread: Callable[[str, str], str] = lambda cg, name: _read(f"/sys/fs/cgroup{cg}/{name}") # ControlGroup, file + environ: Callable[[], dict] = lambda: dict(os.environ) + + +# ---------------------------------------------------------------- real machine + +def _read(path: str) -> str: + try: + return Path(path).read_text() + except OSError: + return "" + + +def _live_sessions() -> set[str]: + """Session ids of running Claude processes (Claude Code's ~/.claude/sessions/<pid>.json); unreadable → empty.""" + out = set() + for f in (Path.home() / ".claude" / "sessions").glob("*.json"): + try: + stat = Path(f"/proc/{f.stem}/stat").read_text() + sid = res.live_session_id(f.read_text(), stat.rsplit(")", 1)[1].split()[19]) + except (OSError, IndexError): + continue + if sid: + out.add(sid) + return out + + +def _run(argv): + return subprocess.run(argv, capture_output=True, text=True) + + +def _xdg(var: str, default: str) -> Path: + return Path(os.environ.get(var) or Path.home() / default) + + +def _argv0(d: Path) -> str: + try: + return (d / "cmdline").read_bytes().split(b"\0", 1)[0].decode(errors="replace") + except OSError: + return "" + + +def _comm(pid: int) -> str: + """Process name; `claude` for any Claude Code process (res.proc_name).""" + d = Path(f"/proc/{pid}") + try: + comm = (d / "comm").read_text().strip() + except OSError: + return "" + return res.proc_name(comm, _argv0(d)) + + +def _ppid(pid: int) -> int: + return int(Path(f"/proc/{pid}/stat").read_text().rsplit(")", 1)[1].split()[1]) + + +def _owner() -> int: + """Nearest ancestor named `claude` (the session outlives its shells), else wf's parent.""" + parent = p = os.getppid() + while p > 1: + if _comm(p) == "claude": + return p + try: + p = _ppid(p) + except (OSError, ValueError, IndexError): + break + return parent + + +def _root() -> Path: + r = _run(["git", "rev-parse", "--show-toplevel"]) + return Path(r.stdout.strip()) if r.returncode == 0 and r.stdout.strip() else Path.cwd() + + +def _procs() -> list[tuple[int, int, str, str]]: + out = [] + for d in Path("/proc").iterdir(): + if not d.name.isdigit(): + continue + try: + stat = (d / "stat").read_text() + comm = stat[stat.index("(") + 1:stat.rindex(")")] + ppid = int(stat.rsplit(")", 1)[1].split()[1]) + cgroup = (d / "cgroup").read_text().strip().split(":", 2)[2] + except (OSError, ValueError, IndexError): + continue + out.append((int(d.name), ppid, res.proc_name(comm, _argv0(d)), cgroup)) + return out + + +def _tmp_used_gb() -> float: + st = os.statvfs("/tmp") + return (st.f_blocks - st.f_bfree) * st.f_frsize / res.GB + + +def real_env() -> Env: + runtime = Path(os.environ.get("XDG_RUNTIME_DIR") or f"/run/user/{os.getuid()}") + if not (runtime / "systemd").is_dir(): + raise res.ResError("wf res needs a systemd user session") + return Env(state=_xdg("XDG_STATE_HOME", ".local/state") / "wf", + config=_xdg("XDG_CONFIG_HOME", ".config") / "wf" / "resources.toml", + units=_xdg("XDG_CONFIG_HOME", ".config") / "systemd" / "user", + scratch=Path(f"/tmp/claude-{os.getuid()}"), run=_run, + meminfo=lambda: Path("/proc/meminfo").read_text(), nproc=os.cpu_count() or 1, + now=lambda: dt.datetime.now().astimezone(), pid_alive=lambda p: Path(f"/proc/{p}").exists(), + owner=_owner, root=_root, cwd=os.getcwd, procs=_procs, sleep=time.sleep, + tmp_used_gb=_tmp_used_gb) + + +# ---------------------------------------------------------------- ledger file + +def warn(msg: str) -> None: + print(f"wf: {msg}", file=sys.stderr) + + +def write_atomic(path: Path, text: str) -> None: + tmp = path.with_name(path.name + ".tmp") + tmp.write_text(text) + os.replace(tmp, path) + + +@contextlib.contextmanager +def locked(env: Env): + """Exclusive flock; yields the ledger; writes it back atomically when the block ends without error.""" + env.state.mkdir(parents=True, exist_ok=True) + with open(env.state / LOCK, "w") as lock: + fcntl.flock(lock, fcntl.LOCK_EX) + path = env.state / LEDGER + text = path.read_text() if path.exists() else "" + try: + led = res.loads(text) + except res.ResError as e: + bad = path.with_name(f"{LEDGER}.bad-{env.now():%Y%m%d-%H%M%S}") + path.rename(bad) + warn(f"{e}; moved to {bad.name}, starting empty") + led = res.Ledger(next=res.salvage_next(text)) + yield led + write_atomic(path, res.dumps(led)) + + +# ---------------------------------------------------------------- systemd + +def systemctl(env: Env, args: list[str]) -> str: + r = env.run(["systemctl", "--user", *args]) + if r.returncode != 0: + first = (r.stderr or r.stdout).strip().splitlines() + raise res.ResError(f"systemctl {args[0]}: " + (first[0] if first else f"exit {r.returncode}")) + return r.stdout + + +def ensure_slice(env: Env) -> None: + f = env.units / res.SLICE + if not f.exists(): + env.units.mkdir(parents=True, exist_ok=True) + f.write_text(res.slice_unit_text()) + systemctl(env, ["daemon-reload"]) + + +def _result(env: Env, e: res.Entry) -> tuple: + """(rc, peak_gb, why); no .rc (wrapper killed with the job) → why from the journal.""" + def read(name): + try: + return (env.state / "logs" / name).read_text().strip() + except OSError: + return "" + rc, peak = read(f"{e.id}.rc"), read(f"{e.id}.peak") + rc, peak = int(rc) if rc.lstrip("-").isdigit() else None, int(peak) / res.GB if peak.isdigit() else None + if rc is not None: + return rc, peak, "" + argv = ["journalctl", "--user", "-u", e.unit, "-o", "cat", "--no-pager"] + if e.started: + argv += ["--since", f"@{int(e.started.timestamp())}"] + why, jpeak = res.journal_reason(env.run(argv).stdout or "", e.mem_gb) + return rc, peak if peak is not None else jpeak, why or f"no exit code; see journalctl --user -u {e.unit}" + + +def bare_units(env: Env) -> list: + names = [l.split()[0] for l in systemctl(env, ["list-units", "--type=service", "--state=running", "--no-legend", + "--plain"]).splitlines() if l.strip()] + if not names: + return [] + return res.bare_units(res.show_units(systemctl(env, ["show", "-p", "Id,Transient,Slice,MemoryCurrent", *names]))) + + +def session_scopes(env: Env) -> dict[str, dict[str, str]]: + """Running scopes in agents.slice (Claude sessions) → show props.""" + names = [l.split()[0] for l in systemctl(env, ["list-units", "--type=scope", "--state=running", "--no-legend", + "--plain"]).splitlines() if l.strip()] + if not names: + return {} + blocks = res.show_units(systemctl(env, ["show", "-p", "Id,Slice,MemoryHigh,MemoryCurrent,ControlGroup", *names])) + return {n: b for n, b in blocks.items() if b.get("Slice") == res.SLICE} + + +def cap_sessions(env: Env, cfg: res.Config) -> None: + for unit, value in res.session_caps(session_scopes(env), cfg.session_mem_gb): + systemctl(env, ["set-property", "--runtime", unit, f"MemoryHigh={value}"]) + + +def gather(env: Env, led: res.Ledger) -> res.Facts: + total, available = res.meminfo(env.meminfo()) + names = sorted({e.unit for e in led.entries if e.state == "running"}) + blocks = res.show_units(systemctl(env, ["show", "-p", "Id,ActiveState,MemoryCurrent,ControlGroup", + res.SLICE, *names])) + units = {} + for n in names: + b = blocks.get(n, {}) + active, cg = b.get("ActiveState") in ACTIVE, b.get("ControlGroup", "") + units[n] = res.Unit(active, res.gb_or_none(b.get("MemoryCurrent", "")), + res.psi_some_avg60(env.cgread(cg, "memory.pressure")) if active and cg else None) + scg = blocks.get(res.SLICE, {}).get("ControlGroup", "") + return res.Facts( + now=env.now(), total_gb=total, available_gb=available, nproc=env.nproc, units=units, + live={e.owner for e in led.entries if e.state == "note" and env.pid_alive(e.owner)}, + results={e.id: _result(env, e) for e in led.entries if e.state == "running" and not units[e.unit].active}, + slice_gb=res.gb_or_none(blocks.get(res.SLICE, {}).get("MemoryCurrent", "")) or 0.0, + slice_cache_gb=res.cache_gb(env.cgread(scg, "memory.stat")) if scg else 0.0) + + +def start(env: Env, e: res.Entry, now: dt.datetime) -> str | None: + """Launch e as a unit; on success e is running. Returns an error line or None.""" + logs = env.state / "logs" + logs.mkdir(parents=True, exist_ok=True) + e.unit, e.log = f"wf-{e.id}.service", str(logs / f"{e.id}.log") + for stale in (f"{e.id}.rc", f"{e.id}.peak"): # a reused id must not inherit an old exit code + (logs / stale).unlink(missing_ok=True) + ensure_slice(env) + r = env.run(res.run_argv(e, str(logs))) + if r.returncode != 0: + first = (r.stderr or r.stdout).strip().splitlines() + return "systemd-run: " + (first[0] if first else f"exit {r.returncode}") + e.state, e.started = "running", now + e.expires = now + dt.timedelta(minutes=2 * e.est_min) + return None + + +def housekeep(env: Env, cfg: res.Config, led: res.Ledger, facts: res.Facts) -> list[str]: + """Every command, under the lock: prune, start queued work.""" + done = [e for e in led.entries if e.state == "done"] + lines = res.prune(led, facts) + kept = {id(e) for e in led.entries} + save_history(env, [r for r in (res.history_record(e) for e in done if id(e) not in kept) if r]) + for e in res.to_start(cfg, led, facts): + err = start(env, e, facts.now) + if err: + e.state, e.ended, e.why = "done", facts.now, f"start failed: {err}" + lines.append(f"{e.id} {e.why}") + else: + lines.append(f"{e.id} started (queued)") + if res.gaming(led, facts.now) or "game off (expired)" in lines: + apply_slice(env, cfg, led, facts) + if not led.last_clean or facts.now - led.last_clean >= res.CLEAN_EVERY: + lines.extend(clean_auto(env, cfg, led, facts.now)) + return lines + + +def read_history(env: Env) -> list[dict]: + path = env.state / HISTORY + out = [] + for line in (path.read_text().splitlines() if path.exists() else []): + try: + r = json.loads(line) + except ValueError: + continue + if isinstance(r, dict): + out.append(r) + return out + + +def save_history(env: Env, recs: list[dict]) -> None: + """Append runs leaving the ledger; keep the newest res.HIST_KEEP (call under the lock).""" + if not recs: + return + path = env.state / HISTORY + old = path.read_text().splitlines() if path.exists() else [] + lines = old + [json.dumps(r) for r in recs] + if len(lines) > res.HIST_KEEP: + write_atomic(path, "".join(x + "\n" for x in lines[-res.HIST_KEEP:])) + else: + with open(path, "a") as f: + f.write("".join(json.dumps(r) + "\n" for r in recs)) + + +def suggestion(env: Env, led: res.Ledger, title: str) -> tuple[float, int, int] | None: + """History's (mem GB, minutes, n) for this project's kind of title (res.suggest), else None.""" + return res.suggest(res.all_runs(read_history(env), led), main_tree(env).name, res.title_kind(title)) + + +def hint(env: Env, led: res.Ledger, title: str, mem: float, est: int) -> None: + line = res.hint_line(mem, est, suggestion(env, led, title)) + if line: + print(line, file=sys.stderr) + + +def session(env: Env, cfg: res.Config, led: res.Ledger) -> res.Facts: + facts = gather(env, led) + housekeep(env, cfg, led, facts) + return facts + + +# ---------------------------------------------------------------- commands + +def _git(env: Env, *args) -> str: + r = env.run(["git", "-C", env.cwd(), *args]) + return r.stdout.strip() if r.returncode == 0 else "" + + +def main_tree(env: Env) -> Path: + """Project's main tree: git common dir's parent (same for its lane worktrees), else env.root().""" + common = _git(env, "rev-parse", "--path-format=absolute", "--git-common-dir") + return Path(common).parent if common and Path(common).name == ".git" else env.root() + + +def _session_records(env: Env) -> list[dict]: + """wf lane session records (.wf/sessions/*.json) of this project: main tree + root.""" + out = [] + for top in dict.fromkeys([main_tree(env), env.root()]): + for f in sorted((top / ".wf" / "sessions").glob("*.json")): + try: + out.append(json.loads(f.read_text())) + except (OSError, ValueError): + pass + return out + + +def owner_by(env: Env, name: str = "") -> dict: + return res.owner_by(env.environ(), _git(env, "branch", "--show-current"), _session_records(env), name) + + +def _new_entry(env, led, facts, args, mem, est, state) -> res.Entry: + return res.Entry(id=led.new_id(), project=main_tree(env).name, owner=env.owner(), title=args.title, mem_gb=mem, + cpus=args.cpus, est_min=est, state=state, cwd=env.cwd(), queued=facts.now, + by=owner_by(env, getattr(args, "by", "") or "")) + + +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" + return f"busy even with --force: only {free} really free beyond the reserve; retry later or work on something else" + + +def cmd_run(env, cfg, args) -> int: + mem, est = res.parse_size(args.mem), res.parse_duration(args.for_) + cmd = args.cmd[1:] if args.cmd[:1] == ["--"] else args.cmd + if not cmd: + raise res.ResError("no command after --") + fail = busy = out = None + 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))) + with locked(env) as led: + facts = session(env, cfg, led) + hint(env, led, args.title, mem, est) + fail = res.never_fits(cfg, facts, mem, args.cpus) + holder = res.lock_holder(led, lock) + 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) + 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}") + 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) + if args.force: + q_gb, q_cpus = 0.0, 0 + 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 + 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) + 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}") + elif args.force: + busy = force_busy(cfg, led, facts, mem, args.cpus, room) + else: + busy = res.busy_line(cfg, led, facts, mem + q_gb, args.cpus + q_cpus, + res.force_hint(cfg, led, facts, mem, args.cpus)) + if fail: + raise res.ResError(fail) + if busy: + raise res.Busy(busy) + print(out) + return 0 + + +def cmd_status(env, cfg, args) -> int: + with locked(env) as led: + facts = session(env, cfg, led) + if args.json: + text = res.status_json(cfg, led, facts) + elif args.id: + text = "\n".join(res.detail_lines(led.get(args.id))) + else: + outside = len(res.adopt_groups(env.procs())) + scopes = session_scopes(env) + stall = {n: res.psi_some_avg60(env.cgread(b.get("ControlGroup", ""), "memory.pressure")) or 0.0 + for n, b in scopes.items() if b.get("ControlGroup")} + text = "\n".join(res.status_lines(cfg, led, facts) + res.throttle_lines(led, facts) + + res.session_cap_lines(scopes, stall, cfg.session_mem_gb) + + res.warning_lines(facts.total_gb, env.tmp_used_gb(), outside, bare_units(env))) + print(text) + return 0 + + +def cmd_hist(env, cfg, args) -> int: + with locked(env) as led: + runs = res.all_runs(read_history(env), led) + print("\n".join(res.hist_lines(runs, args.project))) + return 0 + + +def cmd_release(env, cfg, args) -> int: + fail = out = None + with locked(env) as led: + facts = session(env, cfg, led) + e = led.get(args.id) + if e.state == "running" and not args.stop: + fail = f"{e.id} is running; --stop to kill it" + elif e.state == "running": + systemctl(env, ["stop", e.unit]) + e.rc, e.peak_gb, _ = _result(env, e) + e.state, e.ended, e.why = "done", facts.now, "released" + out = f"{e.id} stopped" + else: + led.entries.remove(e) + out = f"{e.id} released" + if fail: + raise res.ResError(fail) + print(out) + return 0 + + +POLL_SECONDS = 15 + + +def cmd_wait(env, cfg, args) -> int: + limit = res.parse_duration(args.timeout) + deadline = env.now() + dt.timedelta(minutes=limit) + warned = False + while True: + with locked(env) as led: + facts = session(env, cfg, led) + e = led.get(args.id) + 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 + warn(hot[0]) + warned = True + if line: + print(line) + return 0 + if env.now() >= deadline: + raise res.ResError(f"{args.id} not done after {res.fmt_dur(limit)} ({state})") + env.sleep(POLL_SECONDS) +def cmd_note(env, cfg, args) -> int: + mem, est = res.parse_size(args.mem), res.parse_duration(args.for_) + fail = busy = out = None + 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 "")) + with locked(env) as led: + facts = session(env, cfg, led) + hint(env, led, args.title, mem, est) + fail = res.never_fits(cfg, facts, mem, args.cpus) + q_gb, q_cpus = (0.0, 0) if args.force else res.queued_total(led) + room = res.force_room(cfg, led, facts) if args.force else res.budget(cfg, led, facts) + if not fail and res.fits(mem + q_gb, args.cpus + q_cpus, room): + e = _new_entry(env, led, facts, args, mem, est, "note") + e.started, e.expires = facts.now, facts.now + dt.timedelta(minutes=est) + led.entries.append(e) + out = f"{e.id} noted {res.fmt_gb(mem)} until ~{res.hhmm(e.expires)} (then freed; nothing is killed)" + elif not fail and args.force: + busy = force_busy(cfg, led, facts, mem, args.cpus, room) + elif not fail: + busy = res.busy_line(cfg, led, facts, mem + q_gb, args.cpus + q_cpus, + res.force_hint(cfg, led, facts, mem, args.cpus)) + if fail: + raise res.ResError(fail) + if busy: + raise res.Busy(busy) + print(out) + return 0 +def apply_slice(env: Env, cfg: res.Config, led: res.Ledger, facts: res.Facts) -> None: + ensure_slice(env) + props = res.slice_props(cfg, led, facts) + systemctl(env, ["set-property", "--runtime", res.SLICE, *(f"{k}={v}" for k, v in props.items())]) + + +def cmd_game(env, cfg, args) -> int: + with locked(env) as led: + facts = session(env, cfg, led) + if args.state == "on": + minutes = res.parse_duration(args.for_) if args.for_ else int(cfg.game_hours * 60) + led.game_until = max(led.game_until or facts.now, facts.now + dt.timedelta(minutes=minutes)) + lines = res.game_on_lines(cfg, led, facts, minutes) + else: + led.game_until = None + lines = ["game off"] + apply_slice(env, cfg, led, facts) + print("\n".join(lines)) + return 0 +def _lstat(path) -> os.stat_result | None: + """None when the path vanished meanwhile (scratch is shared with running sessions).""" + try: + return os.lstat(path) + except OSError: + return None + + +def _newest(p: Path) -> float: + st = _lstat(p) + newest = st.st_mtime if st else 0.0 + if p.is_dir() and not p.is_symlink(): + for dirpath, dirnames, filenames in os.walk(p, followlinks=False): + for name in dirnames + filenames: + if st := _lstat(os.path.join(dirpath, name)): + newest = max(newest, st.st_mtime) + return newest + + +def _size(p: Path) -> int: + if p.is_symlink() or not p.is_dir(): + st = _lstat(p) + return st.st_size if st else 0 + return sum(st.st_size for d, _, fs in os.walk(p, followlinks=False) for f in fs + if (st := _lstat(os.path.join(d, f)))) + + +def _remove(p: Path) -> None: + if p.is_dir() and not p.is_symlink(): + shutil.rmtree(p, ignore_errors=True) + else: + p.unlink(missing_ok=True) + + +def scan_scratch(env: Env) -> dict: + tree = {} + if not env.scratch.is_dir(): + return tree + for top in env.scratch.iterdir(): + kids = {} + if top.is_dir() and not top.is_symlink(): + kids = {str(c): _newest(c) for c in top.iterdir()} + tree[str(top)] = (_newest(top), kids) + return tree + + +def scan_litter(env: Env) -> list[tuple[str, str, float]]: + uid, out = os.getuid(), [] + for d in env.tmp_dirs: + try: + names = os.listdir(d) + except OSError: + continue + for n in names: + p = Path(d) / n + st = _lstat(p) + if st is None or st.st_uid != uid: + continue + if stat.S_ISDIR(st.st_mode): + try: + kind = "dir" if os.listdir(p) else "emptydir" + except OSError: + continue + else: + kind = "file" if stat.S_ISREG(st.st_mode) and st.st_size else "other" + out.append((str(p), kind, st.st_mtime)) + return out + + +def sweep_litter(env: Env, cfg: res.Config, now: dt.datetime) -> int: + n = 0 + for path in res.litter_victims(scan_litter(env), now.timestamp(), cfg.scratch_hours, env.pid_alive): + try: + (os.rmdir if os.path.isdir(path) and not os.path.islink(path) else os.unlink)(path) + n += 1 + except OSError: + pass # refilled or gone meanwhile + return n + + +def clean_auto(env: Env, cfg: res.Config, led: res.Ledger, now: dt.datetime) -> list[str]: + """Scratch untouched for scratch_hours (D4) and failed wf units; no confirmation needed.""" + lines = [] + for path in res.scratch_victims(scan_scratch(env), now.timestamp(), cfg.scratch_hours, + env.live_sessions()): + p = Path(path) + size = _size(p) + _remove(p) + lines.append(f"freed {res.fmt_bytes(size)}: {p}") + swept = sweep_litter(env, cfg, now) + if swept: + lines.append(f"swept {swept} temp-dir litter entries (empty dirs, dead .NET pipes)") + env.run(["systemctl", "--user", "reset-failed", "wf-r-*"]) + led.last_clean = now + return lines + + +def project_rules(env: Env) -> tuple[Path, list[tuple[str, int | None]]]: + root = env.root() + toml = root / "workflow.toml" + if not toml.exists(): + return root, [] + try: + patterns = tomllib.loads(toml.read_text()).get("cleanup", []) + except tomllib.TOMLDecodeError as e: + raise res.ResError(f"workflow.toml: {e}") from None + return root, [res.cleanup_rule(p) for p in patterns] + + +def project_matches(env: Env, now: dt.datetime) -> list[Path]: + root, rules = project_rules(env) + out = [] + for pattern, days in rules: + for p in sorted(root.glob(pattern)): + if p.is_symlink() or not p.resolve().is_relative_to(root.resolve()): + continue + if days is not None and now.timestamp() - _newest(p) < days * 86400: + continue + out.append(p) + return out + + +def cmd_clean(env, cfg, args) -> int: + with locked(env) as led: + facts = gather(env, led) + lines = [l for l in housekeep(env, cfg, led, facts) if l.startswith("freed ")] + lines += clean_auto(env, cfg, led, facts.now) + root = env.root() + for p in project_matches(env, facts.now): + size, rel = res.fmt_bytes(_size(p)), p.relative_to(root) + if args.yes: + _remove(p) + lines.append(f"deleted {size}: {rel}") + else: + lines.append(f"would delete {size}: {rel} (wf res clean --yes)") + print("\n".join(lines) if lines else "nothing to clean") + return 0 + + +def cmd_timer(env, cfg, args) -> int: + service, timer = env.units / "wf-res.service", env.units / "wf-res.timer" + if args.state == "on": + env.units.mkdir(parents=True, exist_ok=True) + service.write_text(res.service_text(sys.executable, str(HERE / "wf.py"))) + timer.write_text(res.timer_text()) + systemctl(env, ["daemon-reload"]) + systemctl(env, ["enable", "--now", "wf-res.timer"]) + print("timer on: wf-res.timer every 1 min") + else: + if timer.exists(): + systemctl(env, ["disable", "--now", "wf-res.timer"]) + service.unlink(missing_ok=True) + timer.unlink(missing_ok=True) + systemctl(env, ["daemon-reload"]) + print("timer off") + return 0 + + +def adopt(env: Env) -> list[str]: + """Move running Claude sessions (and their descendants) into agents.slice; one line per session.""" + lines = [] + groups = res.adopt_groups(env.procs()) + if groups: + ensure_slice(env) + for root, pids in groups.items(): + r = env.run(res.adopt_argv(root, pids)) + if r.returncode == 0: + lines.append(f"adopted claude {root} ({len(pids)} processes)") + else: + first = (r.stderr or r.stdout).strip().splitlines() + lines.append(f"claude {root}: " + (first[0] if first else f"exit {r.returncode}")) + return lines + + +def cmd_adopt(env, cfg, args) -> int: + print("\n".join(adopt(env)) or "all claude sessions already in agents.slice") + return 0 + + +def cmd_tick(env, cfg, args) -> int: + with locked(env) as led: + session(env, cfg, led) + adopt(env) # silent; a session that vanished mid-move is retried next minute + try: + cap_sessions(env, cfg) + except res.ResError: + pass # a session that ended between list and set-property; retried next minute + return 0 + + +def cmd_hook(env, cfg, args) -> int: + """PreToolUse hook; read-only (no lock, no systemctl) and silent on any error: it runs before every command.""" + try: + payload = json.loads(env.stdin()) + led = res.loads((env.state / LEDGER).read_text()) + out = res.hook_output(payload, led.game_until, env.now()) + except (OSError, ValueError, AttributeError, res.ResError): + return 0 + if out: + print(json.dumps(out)) + return 0 + + +def cmd_shell_init(env, cfg, args) -> int: + print(res.shell_init_line()) + return 0 +# ---------------------------------------------------------------- main + +def parser() -> argparse.ArgumentParser: + ap = argparse.ArgumentParser(prog="wf res", description=__doc__, + formatter_class=argparse.RawDescriptionHelpFormatter) + sub = ap.add_subparsers(dest="command", metavar="COMMAND", required=True) + + def cmd(name, func, help): + sp = sub.add_parser(name, help=help, description=help) + sp.set_defaults(func=func) + return sp + + sp = cmd("run", cmd_run, "reserve and start a job as a systemd user unit; busy → exit 3") + sp.add_argument("--mem", required=True, metavar="SIZE", help="e.g. 10G, 512M") + sp.add_argument("--cpus", type=int, default=1, choices=range(1, 1025), metavar="N") + sp.add_argument("--for", dest="for_", required=True, metavar="DURATION", help="estimate, e.g. 40m, 2h (ETA and queue; the job is never killed at it)") + sp.add_argument("--title", required=True) + sp.add_argument("--queue", action="store_true", help="wait in the FIFO instead of exit 3") + sp.add_argument("--force", action="store_true", help=FORCE_HELP) + sp.add_argument("--by", default="", metavar="NAME", help="your session name (ListAgents), shown as the owner; default: WF_SESSION_NAME, else the wf lane session record") + 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("cmd", nargs=argparse.REMAINDER, metavar="-- CMD …") + + sp = cmd("status", cmd_status, "reservations, budget, game mode, warnings") + sp.add_argument("id", nargs="?") + sp.add_argument("--json", action="store_true") + + sp = cmd("wait", cmd_wait, "poll every 15 s until the entry is done; prints rc and peak") + sp.add_argument("id") + sp.add_argument("--timeout", default="2h", metavar="DURATION") + sp = cmd("hist", cmd_hist, "reservation sizes from history, per project and title kind (title minus trailing " + "sha/hex words and (…)): runs n, request, peak p50/p95, estimate, duration p50/p90, suggestion " + "(≥ 3 runs: mem = p95 peak ×1.15 ≥ 0.2 GB, for = p90 duration ×1.5 ≥ 5 min); run/note print a hint " + "when asking > 2× it") + sp.add_argument("--project", metavar="P", help="only this project (dir name of its main tree)") + sp = cmd("note", cmd_note, "reserve for foreground work; freed when this Claude session exits or at --for") + sp.add_argument("--mem", required=True, metavar="SIZE") + sp.add_argument("--cpus", type=int, default=1, choices=range(1, 1025), metavar="N") + sp.add_argument("--for", dest="for_", required=True, metavar="DURATION", help="reservation freed after this (nothing is killed)") + sp.add_argument("--force", action="store_true", help=FORCE_HELP) + sp.add_argument("--by", default="", metavar="NAME", help="your session name (ListAgents), shown as the owner; default: WF_SESSION_NAME, else the wf lane session record") + sp.add_argument("title") + sp = cmd("game", cmd_game, "gaming reserve on (auto-off, default 4h) or off; jobs slow, never stop") + sp.add_argument("state", choices=["on", "off"]) + sp.add_argument("--for", dest="for_", metavar="DURATION") + sp = cmd("clean", cmd_clean, "delete stale /tmp/claude-* scratch (auto) and list project cleanup patterns") + sp.add_argument("--yes", action="store_true", help="also delete the project's workflow.toml cleanup matches") + + cmd("adopt", cmd_adopt, "move running Claude sessions into agents.slice (no restart); the timer does this every minute") + sp = cmd("timer", cmd_timer, "install/remove the 1-minute wf-res.timer (queue, game expiry, clean)") + sp.add_argument("state", choices=["on", "off"]) + cmd("tick", cmd_tick, "housekeeping once (run by the timer); silent") + cmd("hook", cmd_hook, "Claude Code PreToolUse hook (Bash): no display for agent commands while game is on") + cmd("shell-init", cmd_shell_init, "print the alias that starts Claude inside agents.slice (add to ~/.bashrc)") + sp = cmd("release", cmd_release, "drop an entry (cancels a queued job, frees a note); --stop kills a running job") + sp.add_argument("id") + sp.add_argument("--stop", action="store_true") + return ap + + +# ---------------------------------------------------------------- wf batch + +FORCE_HELP = ("memory is really free though the ledger says busy (unused reservations, stale notes): fit against " + "MemAvailable minus the user reserve and headroom only, ignoring claims and the queue; still recorded; " + "the busy line names --force when it would fit") +CEILING_VAR = "CLAUDE_CODE_PRINT_BG_WAIT_CEILING_MS" +BATCH_TITLE = "wf-batch" +BATCH_FOR = "6h" # default --for and ceiling (no history) +PREP_MEM = "1G" +STOP_FILE = "out/wf-batch.stop" + + +def cmd_batch(env, cfg, args) -> int: + if args.stop: + f = env.root() / STOP_FILE + f.parent.mkdir(exist_ok=True) + f.touch() + print(f"stop requested: {STOP_FILE} (a batch or loop stops before its next spawn; running workers finish)") + return 0 + if args.status: + return batch_status(env, cfg) + if args.time_left: + return batch_time_left(env, args.time_left) + if args.n is None: + batch_parser().print_usage(sys.stderr) + print("wf batch: error: N is required (or --status)", file=sys.stderr) + return 2 + if args.left: + if args.prep: + raise res.ResError("--left is for work batches, not --prep (use --for)") + left = res.parse_duration(args.left) + p90, runs = _task_p90(env) + k = res.batch_fit(args.n, left, p90) + line = f"fit: {k} of {args.n} tasks (left {res.fmt_dur(left)}, task p90 {_p90_text(p90, runs)})" + if not k: + print(f"{line}; nothing started") + return 0 + print(line) + args.n = k + args.for_ = args.for_ or args.left + sug = None + if not args.prep and (args.mem is None or args.for_ is None): # prep jobs share the title, not the size + with locked(env) as led: + sug = suggestion(env, led, BATCH_TITLE) + minutes = res.parse_duration(args.for_ or BATCH_FOR) # the ceiling: never shortened by history + if args.for_ is None: + args.for_ = f"{sug[1]}m" if sug else BATCH_FOR + if args.mem is None: # prep: one claude -p, subagents in-process, no builds → small + args.mem = PREP_MEM if args.prep else f"{sug[0]:g}G" if sug else f"{min(8.0, max(4.0, 1.5 * args.n)):g}G" + claude = shutil.which("claude", path=env.environ().get("PATH")) + if not claude: + raise res.ResError("claude not found on PATH") + out = f"out/wf-batch-{env.now().strftime('%Y-%m-%d-%H%M')}.md" + names = [x.strip() for x in (args.lanes or "").split(",") if x.strip()] + stamp = env.now().strftime('%Y-%m-%d %H:%M') + if args.prep: + ids, tasks_rel = prep_ids(env.root(), set(names) if names else None, args.n) + if not ids: + print(f"prep: no pending task without Done (lanes: {', '.join(names) or 'all'}); nothing started") + return 0 + prompt = (HERE / "templates" / "prep-prompt.md").read_text().strip().format( + here=HERE, root=env.root(), ids=" ".join(ids), out=out, tasks=tasks_rel) + header = f"prep K={args.n}: {' '.join(ids)}" + else: + lanes = ", ".join(names) if names else \ + "every lane with ready tasks, see wf lanes" + prompt = (HERE / "templates" / "batch-prompt.md").read_text().strip().format( + here=HERE, n=args.n, lanes=lanes, out=out, stop=env.root() / STOP_FILE, + deadline=(env.now() + dt.timedelta(minutes=minutes)).isoformat(timespec="minutes")) + header = f"N={args.n}, lanes: {lanes}" + summary = env.root() / out + if not args.dry_run: # dry-run touches nothing; an existing summary (same minute) is appended to, never clobbered + summary.parent.mkdir(exist_ok=True) + with summary.open("a") as f: + f.write(f"# wf-batch {env.root().name} {stamp} ({header})\n") + cmd = [claude, "-p", "--model", args.model, "--permission-mode", "auto", "--permission-prompts", "none", prompt] + if not args.prep and _cloud_project(env.root()): + cmd = [sys.executable, str(HERE / "wf_res.py"), SIDECAR, "--root", str(env.root()), "--summary", out, + "--every", str(SIDECAR_EVERY), "--", *cmd] + ceiling = str(minutes * 60_000) + if args.dry_run: + print(f"{CEILING_VAR}={ceiling} wf res run --mem {args.mem} --for {args.for_} --title {BATCH_TITLE} -- " + + " ".join(shlex.quote(a) for a in cmd)) + return 0 + stop = env.root() / STOP_FILE + if stop.exists(): + stop.unlink() + print(f"warning: stale {STOP_FILE} cleared; starting batch", file=sys.stderr) + base = env.environ + env.environ = lambda: {**base(), CEILING_VAR: ceiling} + try: + code = cmd_run(env, cfg, argparse.Namespace(mem=args.mem, cpus=1, for_=args.for_, title=BATCH_TITLE, + queue=False, force=args.force, cmd=cmd)) + finally: + env.environ = base + if stop.exists(): + stop.unlink() + print(f"summary: {out} · status: wf batch --status") + return code + + +SIDECAR = "batch-sidecar" +SIDECAR_EVERY = 600 # s between the sidecar's wf cloud pull --all runs +PULLED = re.compile(r"^(\S+): (\w+) \(") +TAIL_CAP = 24 * 3600 # s: the tail after the orchestrator exits pulls at most this long +TAIL_MEM_GB, TAIL_CPUS = 0.2, 1 # the job's reservation during the tail + + +def _cloud_project(root: Path) -> bool: + from wflib import config + try: + return config.load(config.find_root(root)).cloud + except (config.ConfigError, OSError): + return False + + +def sidecar(child: list[str], root: Path, summary: Path, every: float, wf_cmd: list[str], + cap: float = TAIL_CAP, shrink: Callable[[], str] | None = None) -> int: + """Run the batch orchestrator `child`; while it lives, every `every` s (until the stop file appears): + wf cloud pull --all, each ended '<id>: <state> (…' line → wf orch post <id> cloud --result <state> + [--commit <sha>] --no-pick + one summary line. After it exits with .wf/cloud records left: the tail — + shrink() the job's reservation, keep pulling every `every` s (no picks, no sends) until no record is left, + the stop file appears or `cap` s passed; one summary line says why it ended. Returns the child's exit code.""" + proc = subprocess.Popen(child) + while True: + try: + rc = proc.wait(timeout=every) + break + except subprocess.TimeoutExpired: + pass + if (root / STOP_FILE).exists(): + print("sidecar: stop file, no more pulls", flush=True) + return proc.wait() + _pull_post(root, summary, wf_cmd) + if not _cloud_records(root): + return rc + print(f"sidecar: orchestrator exited ({rc}), cloud records left: {' '.join(_cloud_records(root))}; " + f"pulling every {every:g}s (no picks) until none, stop file or {cap / 3600:g}h", flush=True) + if shrink: + print(shrink(), flush=True) + t0, pulls = time.monotonic(), 0 + while True: + time.sleep(every) + if (root / STOP_FILE).exists(): + why = "stop file" + elif time.monotonic() - t0 >= cap: + why = f"{cap / 3600:g}h cap" + else: + _pull_post(root, summary, wf_cmd) + pulls += 1 + if _cloud_records(root): + continue + why = "no cloud records left" + break + left = " ".join(_cloud_records(root)) or "none" + line = f"cloud tail ended: {why}; {pulls} pulls; left: {left}" + print(f"sidecar: {line}", flush=True) + with summary.open("a") as f: + f.write(f"{time.strftime('%H:%M')} {line} (sidecar)\n") + return rc + + +def _cloud_records(root: Path) -> list[str]: + return sorted(p.stem for p in (root / ".wf" / "cloud").glob("*.json")) + + +def _pull_post(root: Path, summary: Path, wf_cmd: list[str]) -> None: + pull = subprocess.run([*wf_cmd, "cloud", "pull", "--all", "--project", str(root)], cwd=root, + capture_output=True, text=True) + print(pull.stdout + pull.stderr, end="", flush=True) + for line in pull.stdout.splitlines(): + m = PULLED.match(line) + if not m or m[2] == "running" or line.startswith("archive "): + continue + sha = re.search(r"merged ([0-9a-f]{7,40})\b", line) + post = [*wf_cmd, "orch", "post", m[1], "cloud", "--result", m[2], "--no-pick"] + post += ["--commit", sha[1]] if sha else [] + r = subprocess.run(post, cwd=root, capture_output=True, text=True) + print(r.stdout + r.stderr, end="", flush=True) + with summary.open("a") as f: + f.write(f"{time.strftime('%H:%M')} cloud {m[1]} {m[2]} {sha[1] if sha else '-'} (sidecar pull)\n") + + +def shrink_reservation(env: Env, res_id: str, mem_gb: float = TAIL_MEM_GB, cpus: int = TAIL_CPUS) -> str: + """The batch job's ledger entry down to the tail's size (ledger only: the unit's MemoryMax stays, a pull + must never be OOM-killed).""" + with locked(env) as led: + e = next((x for x in led.entries if x.id == res_id), None) + if not e or e.state != "running": + return f"sidecar: {res_id} not running in the ledger; reservation kept" + old = e.mem_gb + e.mem_gb, e.cpus = min(e.mem_gb, mem_gb), min(e.cpus, cpus) + return (f"sidecar: {res_id} reservation {res.fmt_gb(old)} -> {res.fmt_gb(e.mem_gb)}, {e.cpus} cpu " + "for the cloud tail") + + +def sidecar_main(argv: list[str]) -> int: + ap = argparse.ArgumentParser(prog=f"wf_res.py {SIDECAR}", description="wf batch job wrapper of cloud projects " + "(started by wf batch, not by hand): runs the orchestrator and pulls cloud tasks") + ap.add_argument("--root", type=Path, required=True) + ap.add_argument("--summary", required=True, help="summary file, relative to --root") + ap.add_argument("--every", type=float, default=SIDECAR_EVERY) + ap.add_argument("--wf", default=f"{sys.executable} {HERE / 'wf.py'}", help="wf command (tests: a fake)") + ap.add_argument("--cap", type=float, default=TAIL_CAP, help="s the tail (cloud records left after the " + "orchestrator exits) pulls at most") + ap.add_argument("cmd", nargs=argparse.REMAINDER) + a = ap.parse_args(argv) + cmd = a.cmd[1:] if a.cmd[:1] == ["--"] else a.cmd + if not cmd: + ap.error("no command after --") + rid = os.environ.get(res.RES_ID_VAR, "") + + def shrink() -> str: + try: + return shrink_reservation(real_env(), rid) + except (res.ResError, OSError) as e: + return f"sidecar: reservation of {rid} kept: {e}" + return sidecar(cmd, a.root, a.root / a.summary, a.every, shlex.split(a.wf), a.cap, shrink if rid else None) + + +def prep_ids(root: Path, only: set[str] | None, k: int) -> tuple[list[str], str]: + """Up to k ids a prep batch works (wflib.lanes.prep_targets) + the task file path relative to root.""" + from wflib import config, lanes, tasks + try: + cfg = config.load(config.find_root(root)) + doc = tasks.parse(cfg.tasks.read_text(encoding="utf-8")) + archived = tasks.archive_ids(cfg.archive.read_text(encoding="utf-8")) + except (config.ConfigError, OSError) as e: + raise res.ResError(f"batch --prep: {e}") from None + items = lanes.prep_targets(doc, archived, cfg.lanes, cfg.slice_above, only) + return [i.id for i in items[:k]], cfg.rel(cfg.tasks) + + +def _orch_progress(root: Path, stem: str) -> list[str]: + """out/wf-orch.log lines (start 'YYYY-MM-DD HH:MM' or 'T') at/after the batch start (from the summary name).""" + m = re.fullmatch(r"wf-batch-(\d{4}-\d\d-\d\d)-(\d\d)(\d\d)", stem) + log = root / "out" / "wf-orch.log" + if not m or not log.is_file(): + return [] + start = f"{m[1]} {m[2]}:{m[3]}" + done = [ln for ln in log.read_text().splitlines() + if re.match(r"\d{4}-\d\d-\d\d[ T]\d\d:\d\d", ln) and ln[:16].replace("T", " ") >= start] + return ["finished since batch start (out/wf-orch.log):"] + done[-30:] if done else [] + + +def _task_p90(env) -> tuple[int, int]: + log = env.root() / "out" / "wf-orch.log" + return res.task_p90(log.read_text(errors="replace") if log.exists() else "") + + +def _p90_text(p90: int, runs: int) -> str: + return f"{res.fmt_dur(p90)} from {runs} runs" if runs >= res.TASK_P90_MIN_RUNS else \ + f"{res.fmt_dur(p90)} default ({runs} runs)" + + +def batch_time_left(env, deadline: str) -> int: + """Per-round check of a running batch: spawn while the time to the deadline is ≥ the task p90.""" + try: + end = dt.datetime.fromisoformat(deadline) + except ValueError: + raise res.ResError(f"--time-left '{deadline}' (want an ISO time, e.g. 2026-10-01T20:00+02:00)") + now = env.now() + if end.tzinfo is None: + end = end.replace(tzinfo=now.tzinfo) + left = max(0, int((end - now).total_seconds() // 60)) + p90, runs = _task_p90(env) + word = "spawn" if left >= p90 else "stop" + print(f"time left {res.fmt_dur(left)}, task p90 {res.fmt_dur(p90)} ({runs} runs): {word}") + return 0 + + +def batch_status(env, cfg) -> int: + root = env.root() + with locked(env) as led: + facts = session(env, cfg, led) + mine = [e for e in led.entries if e.title == BATCH_TITLE and e.project == main_tree(env).name] + lines = [res.entry_line(mine[-1], led, facts), f" log {mine[-1].log}"] if mine else [] + summaries = sorted((root / "out").glob("wf-batch-*.md")) + if summaries: + rel = summaries[-1].relative_to(root) + lines += [f"{rel}:"] + summaries[-1].read_text().rstrip("\n").splitlines()[-30:] + lines += _orch_progress(root, summaries[-1].stem) + if (root / STOP_FILE).exists(): + lines.append(f"stop requested: {STOP_FILE}") + print("\n".join(lines) or "no batch in this project") + return 0 + + +def _count(text: str) -> int: + if not text.isdigit() or int(text) < 1: + raise argparse.ArgumentTypeError("N must be ≥ 1") + return int(text) + + +def batch_parser() -> argparse.ArgumentParser: + ap = argparse.ArgumentParser( + prog="wf batch", description="Start an unattended batch orchestrator: claude -p (auto mode, no prompts) " + "as a wf res job, prompt templates/batch-prompt.md; it works up to N tasks and writes out/wf-batch-*.md. " + "Owner allow rule (once): Bash(python3 /projects/public/workflow/wf.py batch:*).") + ap.add_argument("n", nargs="?", type=_count, metavar="N", help="at most N tasks") + ap.add_argument("--lanes", metavar="L,…", help="lanes to work (default: every lane with ready tasks)") + ap.add_argument("--model", default="opus", help="the batch orchestrator's model (default opus)") + ap.add_argument("--for", dest="for_", default=None, metavar="DURATION", + help=f"estimate and {CEILING_VAR} (default: estimate from wf res hist for wf-batch in this " + f"project when ≥ 3 runs (not --prep), else {BATCH_FOR}; the ceiling stays {BATCH_FOR} unless given)") + ap.add_argument("--mem", default=None, metavar="SIZE", + help="reservation (default: wf res hist suggestion for wf-batch in this project when ≥ 3 runs, " + "else max(4G, 1.5G x N), cap 8G; --prep: 1G; explicit wins; a request > 2x history prints a hint)") + ap.add_argument("--dry-run", action="store_true", help="print the command, start nothing") + ap.add_argument("--force", action="store_true", help=FORCE_HELP) + ap.add_argument("--prep", action="store_true", + help="prep batch: for up to N pending tasks without Done (pickable first, then priority) sonnet " + "workers write Done (+Model) from the task text, unclear ones → awaiting; prompt " + "templates/prep-prompt.md; summary lines 'id → Done: …' | 'id → awaiting a-…'. " + "None to prep → prints so, starts nothing") + ap.add_argument("--left", metavar="DURATION", + help="time left to the pilot's deadline: N becomes min(N, floor(left / task p90)) (p90 of done/" + f"handback durations in out/wf-orch.log, < {res.TASK_P90_MIN_RUNS} runs → {res.TASK_P90_DEFAULT}m), " + "prints 'fit: K of N tasks (…)'; K = 0 → nothing started (rc 0); --for defaults to DURATION") + ap.add_argument("--time-left", metavar="DEADLINE", + help="(run by the batch orchestrator before each round) print 'time left X, task p90 Y (n runs): " + "spawn|stop' for the ISO DEADLINE its prompt names (= start + ceiling); stop = left < p90") + ap.add_argument("--status", action="store_true", help="newest batch job of this project + its summary file + tasks finished since start (out/wf-orch.log)") + ap.add_argument("--stop", action="store_true", + help=f"graceful stop: create {STOP_FILE}; batch and loop check it before every spawn (latency ≤ the running " + "tasks), stop with 'stopped: stop file' and removes it") + ap.set_defaults(func=cmd_batch) + return ap + + +def batch_main(argv: list[str], env: Env | None = None) -> int: + return main(argv, env, batch_parser()) + + +PREP_ALL = 1000 # "wf prep all": more than any task list + + +def prep_main(argv: list[str], env: Env | None = None) -> int: + """wf prep [N|all] [batch opts] = wf batch --prep N (all = every not-ready task).""" + if argv and argv[0] == "all": + argv = [str(PREP_ALL), *argv[1:]] + return batch_main([*argv, "--prep"], env) + + +def load_cfg(env: Env) -> res.Config: + return res.load_config(env.config.read_text() if env.config.exists() else "") + + +def main(argv: list[str], env: Env | None = None, ap: argparse.ArgumentParser | None = None) -> int: + try: + args = (ap or parser()).parse_args(argv) + except SystemExit as e: # usage error → 2, -h → 0; callers (and tests) get a code, not an exception + return e.code if isinstance(e.code, int) else 2 + try: + env = env or real_env() + return args.func(env, load_cfg(env), args) + except res.Busy as e: + print(e) + return 3 + except res.ResError as e: + warn(str(e)) + return 1 + except BrokenPipeError: + return 0 + except OSError as e: + warn(str(e)) + return 1 + + +if __name__ == "__main__": + sys.exit(sidecar_main(sys.argv[2:]) if sys.argv[1:2] == [SIDECAR] else main(sys.argv[1:])) |
