aboutsummaryrefslogtreecommitdiffziptar.gz
path: root/wf_res.py
diff options
context:
space:
mode:
Diffstat (limited to 'wf_res.py')
-rw-r--r--wf_res.py1208
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:]))