From 81d4e80fd5aabe4e80f58e960affa795cf7d34ec Mon Sep 17 00:00:00 2001 From: godosa Date: Wed, 7 Oct 2026 07:27:17 +0200 Subject: workflow: initial public history --- wflib/res.py | 993 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 993 insertions(+) create mode 100644 wflib/res.py (limited to 'wflib/res.py') diff --git a/wflib/res.py b/wflib/res.py new file mode 100644 index 0000000..fd1fba9 --- /dev/null +++ b/wflib/res.py @@ -0,0 +1,993 @@ +"""wf res — pure part: parsing, ledger, liveness, capacity, queue, game, scratch, output lines. +Text and data in, data and lines out; no files, no processes. +Spec: docs/resource-ledger.md.""" +from __future__ import annotations + +import datetime as dt +import json +import math +import re +import shlex +import tomllib +from dataclasses import dataclass, field, fields +from typing import Callable + +GB = 1024 ** 3 +SLICE = "agents.slice" +JOBS_SLICE = "agents-jobs.slice" +DONE_KEEP = dt.timedelta(hours=24) +CLEAN_EVERY = dt.timedelta(minutes=10) +EPS = 1e-9 + + +class ResError(Exception): + """Expected failure: `wf: `, exit 1.""" + + +class Busy(Exception): + """Request does not fit now: the busy line, exit 3.""" + + +# ---------------------------------------------------------------- parsing + +SIZE_RE = re.compile(r"^(\d+(?:\.\d+)?)([MG])B?$", re.I) +DUR_RE = re.compile(r"^(?:(\d+)h)?(?:(\d+)m)?$") + + +def parse_size(text: str) -> float: + """'10G' → 10.0, '512M' → 0.5 (GiB).""" + m = SIZE_RE.match(text.strip()) + if not m: + raise ResError(f"size '{text}' (want e.g. 10G or 512M)") + n = float(m.group(1)) + return n if m.group(2).upper() == "G" else n / 1024 + + +def parse_duration(text: str) -> int: + """'40m' → 40, '4h' → 240, '1h30m' → 90, '90' → 90 (minutes, > 0).""" + t = text.strip() + if t.isdigit(): + n = int(t) + else: + m = DUR_RE.match(t) + if not t or not m or not any(m.groups()): + raise ResError(f"duration '{text}' (want e.g. 40m, 4h, 1h30m)") + n = int(m.group(1) or 0) * 60 + int(m.group(2) or 0) + if n <= 0: + raise ResError(f"duration '{text}' must be > 0") + return n + + +def meminfo(text: str) -> tuple[float, float]: + """(MemTotal, MemAvailable) in GB from /proc/meminfo.""" + vals = {} + for line in text.splitlines(): + key, _, rest = line.partition(":") + if key in ("MemTotal", "MemAvailable"): + vals[key] = int(rest.split()[0]) * 1024 / GB + if len(vals) != 2: + raise ResError("no MemTotal/MemAvailable in /proc/meminfo") + return vals["MemTotal"], vals["MemAvailable"] + + +def show_units(text: str) -> dict[str, dict[str, str]]: + """`systemctl show -p Id,…` output (blank-line separated blocks) → {Id: {prop: value}}.""" + out = {} + for block in text.strip().split("\n\n"): + props = dict(l.split("=", 1) for l in block.splitlines() if "=" in l) + if "Id" in props: + out[props["Id"]] = props + return out + + +def psi_some_avg60(text: str) -> float | None: + """cgroup memory.pressure → `some avg60` (percent).""" + for line in text.splitlines(): + if line.startswith("some "): + vals = dict(kv.split("=", 1) for kv in line.split()[1:] if "=" in kv) + try: + return float(vals["avg60"]) + except (KeyError, ValueError): + return None + return None + + +def cache_gb(stat: str) -> float: + """cgroup memory.stat → reclaimable file cache GB: `file` − `shmem` (tmpfs pages count as file but stay).""" + vals = dict(l.split(" ", 1) for l in stat.splitlines() if " " in l) + def get(k): + return int(vals[k]) if vals.get(k, "").strip().isdigit() else 0 + return max(0, get("file") - get("shmem")) / GB + + +def gb_or_none(text: str) -> float | None: + return int(text) / GB if text.strip().isdigit() else None + + +# ---------------------------------------------------------------- formatting + +def fmt_gb(x: float) -> str: + return f"{x:.1f} GB" + + +def fmt_bytes(n: int) -> str: + return fmt_gb(n / GB) if n >= GB else f"{n / 1024 ** 2:.0f} MB" + + +def fmt_dur(minutes: int) -> str: + h, m = divmod(int(minutes), 60) + return (f"{h}h" if h else "") + (f"{m}m" if m or not h else "") + + +def hhmm(t: dt.datetime) -> str: + return t.strftime("%H:%M") + + +def rc_text(rc: int | None) -> str: + return "?" if rc is None else str(rc) + + +# ---------------------------------------------------------------- config + +@dataclass +class Config: + user_reserve_gb: float = 6.0 + user_reserve_cpus: int = 4 + game_reserve_gb: float = 12.0 + game_reserve_cpus: int = 8 + game_hours: float = 4.0 + small_headroom_gb: float = 2.0 + scratch_hours: float = 2.0 + session_mem_gb: float = 6.0 # MemoryHigh per Claude session scope (throttle, no kill); 0 = off + + +def load_config(text: str) -> Config: + try: + data = tomllib.loads(text) + except tomllib.TOMLDecodeError as e: + raise ResError(f"resources.toml: {e}") from None + cfg = Config() + unknown = sorted(set(data) - {f.name for f in fields(Config)}) + if unknown: + raise ResError("resources.toml: unknown key(s) " + ", ".join(unknown)) + for key, value in data.items(): + if isinstance(value, bool) or not isinstance(value, (int, float)) or value < 0: + raise ResError(f"resources.toml: {key} must be a number ≥ 0") + setattr(cfg, key, type(getattr(cfg, key))(value)) + return cfg + + +# ---------------------------------------------------------------- ledger + +TOP_KEYS = ("next", "game_until", "last_clean", "entries") +TIME_FIELDS = ("queued", "started", "expires", "ended") + + +@dataclass +class Entry: + id: str + project: str + owner: int + title: str + mem_gb: float + cpus: int + est_min: int + state: str # queued | running | note | done + cwd: str = "" + cmd: list[str] = field(default_factory=list) + unit: str = "" + log: str = "" + queued: dt.datetime | None = None + started: dt.datetime | None = None + expires: dt.datetime | None = None + ended: dt.datetime | None = None + rc: int | None = None + peak_gb: float | None = None + why: str = "" + env: dict = field(default_factory=dict) # caller variables the user manager lacks (env_diff) + by: dict = field(default_factory=dict) # who started it (owner_by): name, task, batch, address + lock: str = "" # '@
': one queued/running job per lock (lock_for) + extra: dict = field(default_factory=dict, repr=False, compare=False) # unknown keys of a newer version: kept on rewrite + + def eta(self) -> dt.datetime | None: + return self.started + dt.timedelta(minutes=self.est_min) if self.started else None + + +@dataclass +class Ledger: + next: int = 1 + game_until: dt.datetime | None = None + last_clean: dt.datetime | None = None + entries: list[Entry] = field(default_factory=list) + extra: dict = field(default_factory=dict, repr=False, compare=False) # unknown top-level keys: kept on rewrite + + def get(self, id: str) -> Entry: + for e in self.entries: + if e.id == id: + return e + raise ResError(f"no entry '{id}'") + + def new_id(self) -> str: + id = f"r-{self.next}" + self.next += 1 + return id + + +def _time(v): + return dt.datetime.fromisoformat(v) if v else None + + +def _iso(t): + return t.isoformat(timespec="seconds") if t else None + + +def entry_dict(e: Entry) -> dict: + d = {f.name: getattr(e, f.name) for f in fields(Entry) if f.name != "extra"} + d.update({k: v for k, v in e.extra.items() if k not in d}) + for k in TIME_FIELDS: + d[k] = _iso(d[k]) + return d + + +def loads(text: str) -> Ledger: + if not text.strip(): + return Ledger() + try: + d = json.loads(text) + entries, known = [], {f.name for f in fields(Entry)} - {"extra"} + for raw in d.get("entries", []): + extra = {k: v for k, v in raw.items() if k not in known} # newer version's keys: kept, not fatal + raw = {k: v for k, v in raw.items() if k in known} + raw["extra"] = extra + for k in TIME_FIELDS: + raw[k] = _time(raw.get(k)) + entries.append(Entry(**raw)) + return Ledger(next=int(d.get("next", 1)), game_until=_time(d.get("game_until")), + last_clean=_time(d.get("last_clean")), entries=entries, + extra={k: v for k, v in d.items() if k not in TOP_KEYS}) + except (ValueError, TypeError, AttributeError) as e: + raise ResError(f"corrupt ledger: {e}") from None + + +def salvage_next(text: str) -> int: + """Id counter of an unreadable ledger, so a fresh one never reuses ids (logs/.rc, wf-.service).""" + nums = [int(n) for n in re.findall(r'"next":\s*(\d+)', text)] + \ + [int(n) + 1 for n in re.findall(r'"r-(\d+)"', text)] + return max(nums, default=1) + + +def dumps(led: Ledger) -> str: + return json.dumps({**{k: v for k, v in led.extra.items() if k not in TOP_KEYS}, "next": led.next, "game_until": _iso(led.game_until), "last_clean": _iso(led.last_clean), + "entries": [entry_dict(e) for e in led.entries]}, indent=1) + "\n" + + +# ---------------------------------------------------------------- facts, liveness + +@dataclass +class Unit: + active: bool + current_gb: float | None = None + stall: float | None = None # memory.pressure `some avg60`: % of the last minute stalled on memory + + +@dataclass +class Facts: + """Snapshot of the machine, gathered by wf_res under the lock.""" + now: dt.datetime + total_gb: float + available_gb: float + nproc: int + units: dict[str, Unit] = field(default_factory=dict) # unit name → state, for running entries + live: set[int] = field(default_factory=set) # note owner pids still alive + results: dict[str, tuple] = field(default_factory=dict) # id → (rc, peak_gb, why) from logs/, journal + slice_gb: float = 0.0 # agents.slice MemoryCurrent + slice_cache_gb: float = 0.0 # of it reclaimable file cache (cache_gb) + + +def _finish(e: Entry, now: dt.datetime, why: str) -> None: + e.state, e.ended, e.why = "done", now, why + + +def prune(led: Ledger, facts: Facts) -> list[str]: + """Liveness and expiry. Running jobs past their ETA are kept (never killed).""" + out = [] + now = facts.now + for e in led.entries: + if e.state == "running": + unit = facts.units.get(e.unit) + if unit is None or not unit.active: + e.rc, e.peak_gb, why = facts.results.get(e.id, (None, None, "")) + _finish(e, now, why or "exited") + out.append(f"{e.id} {why}" if why else f"{e.id} exited rc={rc_text(e.rc)}") + elif e.state == "note": + if e.owner not in facts.live: + _finish(e, now, "owner gone") + out.append(f"{e.id} freed (owner gone)") + elif e.expires and now >= e.expires: + _finish(e, now, "expired") + out.append(f"{e.id} freed (expired)") + led.entries = [e for e in led.entries if not (e.state == "done" and e.ended and now - e.ended > DONE_KEEP)] + if led.game_until and now >= led.game_until: + led.game_until = None + out.append("game off (expired)") + return out + + +# ---------------------------------------------------------------- capacity + +def gaming(led: Ledger, now: dt.datetime) -> bool: + return bool(led.game_until and now < led.game_until) + + +def reserve(cfg: Config, led: Ledger, now: dt.datetime) -> tuple[float, int]: + if gaming(led, now): + return cfg.game_reserve_gb, cfg.game_reserve_cpus + return cfg.user_reserve_gb, cfg.user_reserve_cpus + + +def _used(e: Entry, facts: Facts) -> float: + unit = facts.units.get(e.unit) + return unit.current_gb if unit and unit.current_gb is not None else 0.0 + + +def held_gb(e: Entry, facts: Facts) -> float: + """What an entry still claims beyond what MemAvailable already shows as used.""" + if e.state == "note": + return e.mem_gb + if e.state == "running": + return max(0.0, e.mem_gb - _used(e, facts)) + return 0.0 + + +def frees_gb(e: Entry, facts: Facts) -> float: + """Budget gained when the entry ends.""" + return max(e.mem_gb, _used(e, facts)) if e.state == "running" else e.mem_gb + + +def _claims(led: Ledger) -> list[Entry]: + return [e for e in led.entries if e.state in ("running", "note")] + + +def budget(cfg: Config, led: Ledger, facts: Facts) -> tuple[float, int]: + res_gb, res_cpus = reserve(cfg, led, facts.now) + claims = _claims(led) + mem = facts.available_gb - res_gb - cfg.small_headroom_gb - sum(held_gb(e, facts) for e in claims) + return round(mem, 6), facts.nproc - res_cpus - sum(e.cpus for e in claims) + + +def force_room(cfg: Config, led: Ledger, facts: Facts) -> tuple[float, int]: + """Budget for --force: what the machine really has free beyond the user reserve and headroom; ledger + claims (unused reservations, notes) and the queue are ignored.""" + res_gb, res_cpus = reserve(cfg, led, facts.now) + return round(facts.available_gb - res_gb - cfg.small_headroom_gb, 6), facts.nproc - res_cpus + + +def fits(mem_gb: float, cpus: int, b: tuple[float, int]) -> bool: + return mem_gb <= b[0] + EPS and cpus <= b[1] + + +def queued_total(led: Ledger) -> tuple[float, int]: + q = [e for e in led.entries if e.state == "queued"] + return sum(e.mem_gb for e in q), sum(e.cpus for e in q) + + +def never_fits(cfg: Config, facts: Facts, mem_gb: float, cpus: int) -> str | None: + """Larger than the whole agent budget on an empty machine (normal reserve).""" + max_gb = facts.total_gb - cfg.user_reserve_gb - cfg.small_headroom_gb + max_cpus = facts.nproc - cfg.user_reserve_cpus + if mem_gb > max_gb + EPS: + return f"{fmt_gb(mem_gb)} can never fit (max {fmt_gb(max_gb)} for agents)" + if cpus > max_cpus: + return f"{cpus} cpus can never fit (max {max_cpus} for agents)" + return None + + +def end_time(e: Entry) -> dt.datetime: + return (e.expires if e.state == "note" else e.eta()) or e.started + + +def until(e: Entry, now: dt.datetime) -> str: + t = end_time(e) + return f"~{hhmm(t)}" if t > now else "overdue" + + +def needed(led: Ledger, facts: Facts, need_gb: float, need_cpus: int) -> tuple[list[Entry], bool]: + """Claims that must end (earliest first) to free need_gb and need_cpus; and whether that is enough.""" + got_gb, got_cpus, out = 0.0, 0, [] + for e in sorted(_claims(led), key=end_time): + if got_gb >= need_gb - EPS and got_cpus >= need_cpus: + break + out.append(e) + got_gb += frees_gb(e, facts) + got_cpus += e.cpus + return out, got_gb >= need_gb - EPS and got_cpus >= need_cpus + + +def force_hint(cfg: Config, led: Ledger, facts: Facts, mem_gb: float, cpus: int) -> str: + """Busy-line suffix naming --force when the request fits the memory that is really free.""" + room = force_room(cfg, led, facts) + if not fits(mem_gb, cpus, room): + return "" + return (f"; {fmt_gb(room[0])} really free beyond the reserve: --force starts it past the ledger " + "(only if the holders will not use what they reserved)") + + +def busy_line(cfg: Config, led: Ledger, facts: Facts, mem_gb: float, cpus: int, hint: str = "") -> str: + b_gb, b_cpus = budget(cfg, led, facts) + need_gb, need_cpus = mem_gb - b_gb, cpus - b_cpus + holders, _ = needed(led, facts, need_gb, need_cpus) + if not holders: + return (f"busy: only {fmt_gb(max(0.0, b_gb))} free for agents and no agent job holds any; " + "other programs use the rest; retry later or work on something else" + hint) + by_mem = need_gb > EPS + + def held(e): + what = fmt_gb(frees_gb(e, facts)) if by_mem else f"{e.cpus} cpus" + return f'{what} held by {e.project} "{e.title}" ({e.id}) until {until(e, facts.now)}' + + free = fmt_gb(max(0.0, b_gb)) if by_mem else f"{max(0, b_cpus)} cpus" + last = max(end_time(e) for e in holders) + retry = f"after ~{hhmm(last)}" if last > facts.now else "later" + return (f"busy: {', '.join(held(e) for e in holders)}; {free} free for agents; " + f"retry {retry} or work on something else" + hint) + + +# ---------------------------------------------------------------- locks + +LOCK_TITLE_RE = re.compile(r"gate\b", re.I) # 'gate …' titles lock 'gate' unasked: one shared gate checkout + + +def lock_for(title: str, key: str, main: str) -> str: + """Ledger lock of a run: --lock KEY, else 'gate' for a title starting with the word gate; '' = none. + Scoped to the main tree (lane worktrees share their project's gate dir).""" + key = key or ("gate" if LOCK_TITLE_RE.match(title.strip()) else "") + return f"{key}@{main}" if key else "" + + +def lock_holder(led: Ledger, lock: str) -> Entry | None: + """The running (else first queued) entry holding this lock.""" + if not lock: + return None + mine = [e for e in led.entries if e.lock == lock and e.state in ("running", "queued")] + return next((e for e in mine if e.state == "running"), None) or next(iter(_queue_of(mine)), None) + + +def _queue_of(entries: list[Entry]) -> list[Entry]: + return sorted((e for e in entries if e.state == "queued"), key=lambda e: (e.queued, int(e.id[2:]))) + + +def lock_busy_line(h: Entry, facts: Facts) -> str: + key = h.lock.split("@", 1)[0] + when = f"ETA {until(h, facts.now)}" if h.state == "running" else "queued" + return (f"busy: lock '{key}' held by {h.id} \"{h.title}\" ({when}): one at a time (shared checkout); " + f"--queue waits for it (--force does not override a lock); never a hand-written copy without the lock") + + +# ---------------------------------------------------------------- run, status lines + +RES_ID_VAR = "WF_RES_ID" # set in every job's unit: wf res calls from inside it name it as their batch +ENV_SKIP = {"PWD", "OLDPWD", "SHLVL", "_", RES_ID_VAR} +SECRET_RE = re.compile(r"TOKEN|SECRET|PASSWORD|PASSWD|CREDENTIAL", re.I) # never written to the ledger + + +def env_diff(caller: dict[str, str], manager: dict[str, str]) -> dict[str, str]: + """Caller variables a systemd --user unit would not get as is (it starts from the manager's environment).""" + return {k: v for k, v in caller.items() + if manager.get(k) != v and k not in ENV_SKIP and not SECRET_RE.search(k)} + + +def parse_show_environment(text: str) -> dict[str, str]: + return dict(l.split("=", 1) for l in text.splitlines() if "=" in l) + + +def run_argv(e: Entry, logdir: str) -> list[str]: + """systemd-run argv; the sh wrapper records exit code and peak bytes (unit is collected on exit).""" + rc, peak = shlex.quote(f"{logdir}/{e.id}.rc"), shlex.quote(f"{logdir}/{e.id}.peak") + script = (f'"$@"; rc=$?; cat /sys/fs/cgroup$(cut -d: -f3 /proc/self/cgroup)/memory.peak > {peak} ' + f'2>/dev/null; echo $rc > {rc}') + return ["systemd-run", "--user", "--quiet", "--collect", f"--slice={JOBS_SLICE}", f"--unit={e.unit}", + f"--working-directory={e.cwd}", *(f"--setenv={k}={v}" for k, v in sorted(e.env.items())), + f"--setenv={RES_ID_VAR}={e.id}", + "-p", f"MemoryMax={int(e.mem_gb * GB)}", "-p", f"MemoryHigh={int(e.mem_gb * 0.9 * GB)}", + "-p", "MemorySwapMax=0", "-p", "Nice=10", + "-p", f"StandardOutput=append:{e.log}", "-p", f"StandardError=append:{e.log}", + "/bin/sh", "-c", script, "sh", *e.cmd] + + +TASK_BRANCH_RE = re.compile(r"^[\w.-]+/([\w.-]+)$") # / (wf worktree branches) + + +def owner_by(environ: dict[str, str], branch: str, records: list[dict], name: str = "") -> dict[str, str]: + """Who starts an entry, empty keys dropped. name: --by, else WF_SESSION_NAME, else ' session' from the + wf lane session record (.wf/sessions) with this CLAUDE_PID; task: WF_TASK, else the / branch; + batch: WF_RES_ID (the wf res job, e.g. wf batch, this runs in); address: uds: + (a batch worker's = its orchestrator's: SendMessage there reaches the batch).""" + pid = environ.get("CLAUDE_PID", "") + name = name or environ.get("WF_SESSION_NAME", "") + if not name and pid: + name = next((f"{r['lane']} session" for r in records + if isinstance(r, dict) and r.get("lane") and str(r.get("pid")) == pid), "") + m = TASK_BRANCH_RE.match(branch) + sock = environ.get("CLAUDE_CODE_MESSAGING_SOCKET", "") + d = {"name": name, "task": environ.get("WF_TASK") or (m.group(1) if m else ""), + "batch": environ.get(RES_ID_VAR, ""), "address": f"uds:{sock}" if sock else ""} + return {k: v for k, v in d.items() if v} + + +def by_text(e: Entry) -> str: + """'slow session t-x batch r-3, message uds:/s' — empty when nothing is known.""" + b = e.by or {} + who = " ".join(x for x in (b.get("name"), b.get("task"), f"batch {b['batch']}" if b.get("batch") else "") if x) + msg = f"message {b['address']}" if b.get("address") else "" + return ", ".join(x for x in (who, msg) if x) + + +def position(led: Ledger, e: Entry) -> int: + return _queue(led).index(e) + 1 + + +def entry_line(e: Entry, led: Ledger, facts: Facts) -> str: + by = by_text(e) if e.state != "done" else "" + return _entry_line(e, led, facts) + (f" [by {by}]" if by else "") + + +def _entry_line(e: Entry, led: Ledger, facts: Facts) -> str: + head = f'{e.id} {e.project} "{e.title}"' + if e.state == "running": + unit = facts.units.get(e.unit) + used = fmt_gb(unit.current_gb) if unit and unit.current_gb is not None else "?" + return (f"{head} running {fmt_gb(e.mem_gb)} used {used} {e.cpus} cpu since {hhmm(e.started)} " + f"ETA {until(e, facts.now)}" + (" (still running, not killed)" if end_time(e) <= facts.now else "")) + if e.state == "queued": + return f"{head} queued #{position(led, e)} {fmt_gb(e.mem_gb)} {e.cpus} cpu" + if e.state == "note": + return f"{head} note {fmt_gb(e.mem_gb)} {e.cpus} cpu until {until(e, facts.now)}" + peak = fmt_gb(e.peak_gb) if e.peak_gb is not None else "?" + return f"{head} done ({e.why}) rc={rc_text(e.rc)} peak {peak} at {hhmm(e.ended)}" + + +def status_lines(cfg: Config, led: Ledger, facts: Facts) -> list[str]: + active = [e for e in led.entries if e.state != "done"] + done = [e for e in led.entries if e.state == "done"] + out = [entry_line(e, led, facts) for e in active + done] or ["no reservations"] + b_gb, b_cpus = budget(cfg, led, facts) + r_gb, r_cpus = reserve(cfg, led, facts.now) + out.append(f"agents may use {fmt_gb(max(0.0, b_gb))}, {max(0, b_cpus)} cpus now; " + f"reserve {fmt_gb(r_gb)}/{r_cpus} cpus") + if facts.slice_gb > 0: + jobs = sum(_used(e, facts) for e in led.entries if e.state == "running") + free = max(0.0, facts.slice_gb - jobs) + cache = min(free, facts.slice_cache_gb) + out.append(f"unreserved agent memory {fmt_gb(free)}" + + (f" ({fmt_gb(cache)} of it file cache, reclaimable)" if cache >= 0.05 else "")) + if gaming(led, facts.now): + left = int((led.game_until - facts.now).total_seconds() // 60) + out.append(f"game on until {hhmm(led.game_until)} ({fmt_dur(left)} left)") + out.extend(shortfall_lines(cfg, led, facts)) + return out + + +def status_json(cfg: Config, led: Ledger, facts: Facts) -> str: + b_gb, b_cpus = budget(cfg, led, facts) + return json.dumps({"budget_gb": b_gb, "budget_cpus": b_cpus, "gaming": gaming(led, facts.now), + "game_until": _iso(led.game_until), "entries": [entry_dict(e) for e in led.entries]}, + indent=1) + + +def detail_lines(e: Entry) -> list[str]: + d = entry_dict(e) + d["cmd"] = shlex.join(e.cmd) + return [f"{k}: {v}" for k, v in d.items() if v not in (None, "", [])] + + +def slice_unit_text() -> str: + return ("[Unit]\nDescription=wf agents: Claude sessions and wf res jobs\n\n" + "[Slice]\nCPUWeight=20\nIOWeight=20\n") + + +def done_line(e: Entry) -> str: + minutes = round((e.ended - e.started).total_seconds() / 60) if e.ended and e.started else 0 + peak = fmt_gb(e.peak_gb) if e.peak_gb is not None else "?" + why = f" ({e.why})" if e.why and e.why != "exited" else "" + return f"{e.id} done rc={rc_text(e.rc)} peak {peak} in {minutes} min{why}" + + +JOURNAL_RESULT = re.compile(r"Failed with result '([^']+)'") +JOURNAL_MAIN = re.compile(r"Main process exited, code=\w+, status=(\S+)") +JOURNAL_PEAK = re.compile(r"(\d+(?:\.\d+)?)([BKMGT]) memory peak") +UNIT_SCALE = {"B": 1 / GB, "K": 1 / 1024 ** 2, "M": 1 / 1024, "G": 1.0, "T": 1024.0} + + +def journal_reason(text: str, mem_gb: float) -> tuple[str, float | None]: + """Why a unit died, from `journalctl -u UNIT -o cat` (used when the job wrote no .rc): (why, peak_gb).""" + m = JOURNAL_PEAK.search(text) + peak = round(float(m.group(1)) * UNIT_SCALE[m.group(2)], 3) if m else None + m = JOURNAL_RESULT.search(text) + result = m.group(1) if m else "" + m = JOURNAL_MAIN.search(text) + status = m.group(1) if m else "?" + if result == "oom-kill": + by = "by systemd-oomd" if "systemd-oomd killed" in text else "at MemoryMax" + return f"killed: oom-kill {by}, limit {fmt_gb(mem_gb)}; raise --mem", peak + if result in ("signal", "core-dump"): + return f"killed: signal {status}", peak + if result == "exit-code": + return f"failed: exit {status}", peak + return (f"killed: {result}" if result else ""), peak + + +def _queue(led: Ledger) -> list[Entry]: + return _queue_of(led.entries) + + +def to_start(cfg: Config, led: Ledger, facts: Facts) -> list[Entry]: + """Strict FIFO: queued entries to start now, stopping at the first that does not fit; an entry whose lock a + running job (or one started in this pass) holds is skipped, not blocking the rest.""" + b_gb, b_cpus = budget(cfg, led, facts) + out, held = [], {e.lock for e in led.entries if e.state == "running" and e.lock} + for e in _queue(led): + if e.lock and e.lock in held: + continue + if not fits(e.mem_gb, e.cpus, (b_gb, b_cpus)): + break + out.append(e) + held.add(e.lock) + b_gb, b_cpus = b_gb - e.mem_gb, b_cpus - e.cpus + return out + + +def queue_estimate(cfg: Config, led: Ledger, facts: Facts, e: Entry) -> dt.datetime | None: + """When enough claims end for e and everything ahead of it; None if they never free enough.""" + q = _queue(led) + ahead = q[:q.index(e) + 1] + b_gb, b_cpus = budget(cfg, led, facts) + holders, enough = needed(led, facts, sum(x.mem_gb for x in ahead) - b_gb, sum(x.cpus for x in ahead) - b_cpus) + if not enough: + return None + lock = [end_time(x) for x in led.entries if x is not e and x.lock and x.lock == e.lock and x.state == "running"] + return max([*(end_time(h) for h in holders), *lock], default=facts.now) + + +def slice_props(cfg: Config, led: Ledger, facts: Facts) -> dict[str, str]: + """agents.slice values: normal, or gaming (MemoryHigh never below current use: slow, never squeeze).""" + if gaming(led, facts.now): + high, weight = max(facts.total_gb - cfg.game_reserve_gb, facts.slice_gb), "5" + else: + high, weight = facts.total_gb - cfg.user_reserve_gb, "20" + return {"CPUWeight": weight, "IOWeight": weight, "MemoryHigh": str(int(round(high * GB)))} + + +def shortfall_lines(cfg: Config, led: Ledger, facts: Facts) -> list[str]: + """Empty when the game reserve is free now; else who holds it and when it frees.""" + free = facts.available_gb - sum(held_gb(e, facts) for e in _claims(led)) + short = round(cfg.game_reserve_gb - free, 6) + if short <= EPS: + return [] + holders, enough = needed(led, facts, short, 0) + head = f"short {fmt_gb(short)} of {fmt_gb(cfg.game_reserve_gb)}: " + if not holders: + return [head + "other programs, not agent jobs", "agent jobs alone cannot free it; close other programs"] + parts = ", ".join(f'{e.id} {e.project} "{e.title}" {fmt_gb(frees_gb(e, facts))} {until(e, facts.now)}' + for e in holders) + if not enough: + return [head + parts, "agent jobs alone cannot free it; close other programs"] + last = max(end_time(e) for e in holders) + big = max(holders, key=lambda e: frees_gb(e, facts)) + when = f"~{hhmm(last)}" if last > facts.now else "soon" + return [head + parts, f"full reserve free {when} (est.); free now: wf res release {big.id} --stop"] + + +def game_on_lines(cfg: Config, led: Ledger, facts: Facts, minutes: int) -> list[str]: + return [f"game on until {hhmm(led.game_until)} ({fmt_dur(minutes)}); CPU/IO now yours", + *shortfall_lines(cfg, led, facts)] + + +def scratch_victims(tree: dict, now_ts: float, hours: float, keep: set[str] = frozenset()) -> list[str]: + """tree = {top path: (newest mtime, {child path: newest mtime})}. A stale top goes whole; under a fresh + top, its stale children go. Children named in keep (live Claude session ids) never go; a stale top + holding one is pruned per child. Sorted.""" + limit = now_ts - hours * 3600 + out = [] + for top, (newest, children) in tree.items(): + held = any(c.rsplit("/", 1)[-1] in keep for c in children) + if newest <= limit and not held: + out.append(top) + else: + out.extend(c for c, m in children.items() if m <= limit and c.rsplit("/", 1)[-1] not in keep) + return sorted(out) + + +DOTNET_PIPE_RE = re.compile(r"^(?:clr-debug-pipe|dotnet-diagnostic)-(\d+)-") + + +def litter_victims(entries: list[tuple[str, str, float]], now_ts: float, hours: float, + alive: Callable[[int], bool]) -> list[str]: + """Top-level temp-dir litter (own entries only): (path, kind emptydir|dir|file|other, mtime) → sorted paths. + Empty dirs idle for `hours`; .NET debug pipes/sockets of dead pids. Content is never swept.""" + limit = now_ts - hours * 3600 + out = [] + for path, kind, mtime in entries: + m = DOTNET_PIPE_RE.match(path.rsplit("/", 1)[-1]) + if (kind == "emptydir" and mtime <= limit) or (m and kind != "dir" and not alive(int(m.group(1)))): + out.append(path) + return sorted(out) + + +def live_session_id(text: str, proc_start: str) -> str | None: + """~/.claude/sessions/.json of a live pid → its sessionId (= its scratch dir name), unless the + recorded procStart shows the pid now belongs to another process.""" + try: + d = json.loads(text) + except ValueError: + return None + if not isinstance(d, dict) or not isinstance(d.get("sessionId"), str): + return None + if d.get("procStart") is not None and str(d["procStart"]) != proc_start: + return None + return d["sessionId"] + + +AGE_RE = re.compile(r"^(.*):(\d+)d$") + + +def cleanup_rule(text: str) -> tuple[str, int | None]: + """'out/logs/*.log:30d' → ('out/logs/*.log', 30); relative, no '..'.""" + m = AGE_RE.match(text) + pattern, days = (m.group(1), int(m.group(2))) if m else (text, None) + if not pattern or pattern.startswith("/") or ".." in pattern.split("/"): + raise ResError(f"cleanup pattern '{text}' (want a path relative to the project, no '..')") + return pattern, days + + +CLAUDE_BIN_RE = re.compile(r"/claude/versions/[^/]+$") + + +def proc_name(comm: str, argv0: str) -> str: + """`claude` for a Claude Code process, else comm. Background sessions run the versioned binary, so comm + is the version ('2.1.283'); argv0 is …/claude/versions/X, or a rewritten title 'claude bg-pty-host'. Tools it + re-execs (ugrep) keep their own argv0.""" + if argv0.split(" ", 1)[0].rsplit("/", 1)[-1] == "claude" or CLAUDE_BIN_RE.search(argv0): + return "claude" + return comm + + +def adopt_groups(procs: list[tuple[int, int, str, str]]) -> dict[int, list[int]]: + """procs = (pid, ppid, comm, cgroup). Claude session roots outside agents.slice (named `claude`, parent not + `claude`) → sorted pids of the root and its descendants that are outside the slice.""" + inside = f"/{SLICE}/" + by_pid = {p[0]: p for p in procs} + kids: dict[int, list[int]] = {} + for pid, ppid, _, _ in procs: + kids.setdefault(ppid, []).append(pid) + out = {} + for pid, ppid, comm, cgroup in procs: + parent = by_pid.get(ppid) + if comm != "claude" or inside in cgroup or (parent and parent[2] == "claude"): + continue + tree, stack = [], [pid] + while stack: + p = stack.pop() + if inside not in by_pid[p][3]: + tree.append(p) + stack.extend(kids.get(p, [])) + out[pid] = sorted(tree) + return out + + +THROTTLE_FILL = 0.85 # used ≥ this × --mem: at MemoryHigh (0.9 × --mem), where the kernel throttles +THROTTLE_STALL = 20.0 # % of the last minute stalled on memory: page cache alone at the limit stays near 0 + + +def throttle_lines(led: Ledger, facts: Facts) -> list[str]: + """Running jobs stuck at their own memory limit (they crawl, then systemd-oomd kills them).""" + out = [] + for e in led.entries: + u = facts.units.get(e.unit) + if (e.state == "running" and u and u.active and u.current_gb is not None and u.stall is not None + and u.current_gb >= THROTTLE_FILL * e.mem_gb and u.stall >= THROTTLE_STALL): + out.append(f"{e.id} throttled at its memory limit ({u.current_gb:.1f} of {fmt_gb(e.mem_gb)}, " + f"stalled {u.stall:.0f}% of the last minute): likely too small; " + f"`wf res release {e.id} --stop` and re-run with a bigger --mem" + + (f"; started by {by_text(e)}" if by_text(e) else "")) + return out + + +DESKTOP_UNIT = ("app-", "dbus-") # transient units the desktop starts (XDG app launch, dbus activation) + + +def bare_units(blocks: dict[str, dict[str, str]]) -> list[tuple[str, float | None]]: + """Running transient services outside agents.slice not from the desktop = jobs from a bare `systemd-run`: + sorted (unit, used GB).""" + return sorted((n, gb_or_none(b.get("MemoryCurrent", ""))) for n, b in blocks.items() + if b.get("Transient") == "yes" and not n.startswith(DESKTOP_UNIT) + and not b.get("Slice", "").startswith("agents")) + + +def session_caps(blocks: dict[str, dict[str, str]], cap_gb: float) -> list[tuple[str, str]]: + """Session scopes in agents.slice whose MemoryHigh differs from the cap → sorted (unit, value) to set. + MemoryHigh only: over it the session is throttled; a MemoryMax kill could pick `claude` itself.""" + want = str(int(cap_gb * GB)) if cap_gb > 0 else "infinity" + return sorted((n, want) for n, b in blocks.items() + if n.endswith(".scope") and b.get("Slice") == SLICE and b.get("MemoryHigh") != want) + + +def session_cap_lines(blocks: dict[str, dict[str, str]], stall: dict[str, float], cap_gb: float) -> list[str]: + """Sessions stuck at their cap: used ≥ 0.9 × cap and stalled ≥ THROTTLE_STALL %.""" + if cap_gb <= 0: + return [] + out = [] + for n, b in sorted(blocks.items()): + cur = gb_or_none(b.get("MemoryCurrent", "")) + if cur is not None and cur >= 0.9 * cap_gb and stall.get(n, 0.0) >= THROTTLE_STALL: + who = n.removeprefix("wf-claude-").removesuffix(".scope") + out.append(f"warning: claude session {who} at its memory cap ({cur:.1f} of {fmt_gb(cap_gb)}, " + f"stalled {stall[n]:.0f}% of the last minute): run big work with wf res run") + return out + + +def warning_lines(total_gb: float, tmp_gb: float, outside: int, + bare: list[tuple[str, float | None]] = ()) -> list[str]: + out = [] + if tmp_gb > 0.25 * total_gb: + out.append(f"warning: /tmp (RAM) holds {fmt_gb(tmp_gb)}; wf res clean") + if outside: + out.append(f"warning: {outside} claude sessions outside {SLICE} (wf res adopt)") + if bare: + units = ", ".join(f"{n} {'?' if gb is None else f'{gb:.1f}'} GB" for n, gb in bare) + out.append(f"warning: {len(bare)} jobs outside wf res (bare systemd-run): {units}; " + "start jobs with wf res run, also from project scripts") + return out + + +def adopt_argv(root: int, pids: list[int]) -> list[str]: + """Move live pids into a new scope in agents.slice (systemd StartTransientUnit with PIDs).""" + return ["busctl", "--user", "call", "org.freedesktop.systemd1", "/org/freedesktop/systemd1", + "org.freedesktop.systemd1.Manager", "StartTransientUnit", "ssa(sv)a(sa(sv))", + f"wf-claude-{root}.scope", "fail", "2", "PIDs", "au", str(len(pids)), *map(str, pids), + "Slice", "s", SLICE, "0"] + + +def service_text(python: str, wf: str) -> str: + return ("[Unit]\nDescription=wf res tick\n\n[Service]\nType=oneshot\n" + f"ExecStart={python} {wf} res tick\n") + + +def timer_text() -> str: + return ("[Unit]\nDescription=wf res tick every minute\n\n[Timer]\nOnBootSec=1min\nOnUnitActiveSec=1min\n\n" + "[Install]\nWantedBy=timers.target\n") + + +def hook_output(payload: dict, game_until: dt.datetime | None, now: dt.datetime) -> dict | None: + """Claude Code PreToolUse hook: while game mode is on, Bash commands run without a display (no windows pop + up over the game). None = no output, the command runs unchanged.""" + cmd = (payload.get("tool_input") or {}).get("command") + if payload.get("tool_name") != "Bash" or not cmd or not (game_until and now < game_until): + return None + return {"hookSpecificOutput": { + "hookEventName": "PreToolUse", + "updatedInput": {**payload["tool_input"], "command": f"unset DISPLAY WAYLAND_DISPLAY; {cmd}"}, + "additionalContext": f"wf res game mode until {hhmm(game_until)}: this command has no display, so GUI " + "windows fail. Run headless/offscreen, or do other work until game off."}} + + +def shell_init_line() -> str: + return f"alias claude='systemd-run --user --scope --quiet --slice={SLICE} claude'" + + +# ---------------------------------------------------------------- history (reservation sizes) + +HIST_MIN_RUNS = 3 +HIST_MEM_PCT, HIST_MEM_X, HIST_MEM_FLOOR = 0.95, 1.15, 0.2 # suggest mem = p95 peak ×1.15, ≥ 0.2 GB +HIST_DUR_PCT, HIST_DUR_X, HIST_DUR_FLOOR = 0.90, 1.5, 5 # suggest for = p90 duration ×1.5, ≥ 5 min +HIST_OVER = 2.0 # request > 2× suggestion → hint +HIST_KEEP = 2000 # history file: newest runs kept +_PAREN_TAIL = re.compile(r"\s*\([^()]*\)\s*$") +_HEX_TAIL = re.compile(r"\s+(?=[0-9a-f]*\d)[0-9a-f]{7,40}$") + + +def title_kind(title: str) -> str: + """A title without its trailing (…) and sha/hex words: 'gate 2c91cac7 (fix x)' → 'gate'.""" + t = title.strip() + while True: + s = _HEX_TAIL.sub("", _PAREN_TAIL.sub("", t)).strip() + if s == t or not s: + return t + t = s + + +def pct(xs: list[float], p: float) -> float: + """Nearest-rank percentile.""" + s = sorted(xs) + return s[max(0, math.ceil(round(p * len(s), 9)) - 1)] + + +def history_record(e: Entry) -> dict | None: + """A finished, successful run with a measured peak → history record; else None.""" + if e.state != "done" or e.rc != 0 or e.peak_gb is None or not e.started or not e.ended: + return None + return {"id": e.id, "project": e.project, "title": e.title, "mem_gb": e.mem_gb, "est_min": e.est_min, + "rc": e.rc, "peak_gb": e.peak_gb, "min": round((e.ended - e.started).total_seconds() / 60, 2)} + + +def all_runs(records: list[dict], led: Ledger) -> list[dict]: + """History file records + the ledger's finished runs not yet in it (ids never repeat).""" + seen = {r.get("id") for r in records} + return records + [r for r in map(history_record, led.entries) if r and r["id"] not in seen] + + +def _kind_runs(runs: list[dict], project: str, kind: str) -> list[dict]: + return [r for r in runs if r.get("project") == project and title_kind(r.get("title", "")) == kind + and r.get("rc") == 0 and r.get("peak_gb") is not None and r.get("min") is not None] + + +def _ceil1(x: float) -> float: + return math.ceil(round(x * 10, 6)) / 10 + + +def suggest(runs: list[dict], project: str, kind: str) -> tuple[float, int, int] | None: + """(mem GB, minutes, n) from ≥ 3 runs of this kind in this project, else None.""" + rs = _kind_runs(runs, project, kind) + if len(rs) < HIST_MIN_RUNS: + return None + mem = max(HIST_MEM_FLOOR, _ceil1(pct([r["peak_gb"] for r in rs], HIST_MEM_PCT) * HIST_MEM_X)) + minutes = max(HIST_DUR_FLOOR, math.ceil(round(pct([r["min"] for r in rs], HIST_DUR_PCT) * HIST_DUR_X, 6))) + return mem, minutes, len(rs) + + +def hint_line(mem_gb: float, est_min: int, sug: tuple[float, int, int] | None) -> str | None: + """Request far over the history (no auto-resize): one line, else None.""" + if not sug or (mem_gb <= HIST_OVER * sug[0] + EPS and est_min <= HIST_OVER * sug[1]): + return None + return f"hint: history says ~{sug[0]:g} GB / {sug[1]} min ({sug[2]} runs)" + + +def hist_lines(runs: list[dict], project: str | None) -> list[str]: + """wf res hist: per project and title kind: n, request, peak p50/p95, estimate, duration p50/p90, suggestion.""" + keys = sorted({(r.get("project", ""), title_kind(r.get("title", ""))) for r in runs + if project is None or r.get("project") == project}) + out = [] + for p, k in keys: + rs = _kind_runs(runs, p, k) + if not rs: + continue + peaks, durs = [r["peak_gb"] for r in rs], [r["min"] for r in rs] + sug = suggest(runs, p, k) + out.append(f"{p} · {k} · n {len(rs)} · req {fmt_gb(pct([r['mem_gb'] for r in rs], 0.5))} · " + f"peak {pct(peaks, 0.5):.1f}/{pct(peaks, HIST_MEM_PCT):.1f} GB · " + f"est {fmt_dur(pct([r['est_min'] for r in rs], 0.5))} · " + f"dur {fmt_dur(round(pct(durs, 0.5)))}/{fmt_dur(round(pct(durs, HIST_DUR_PCT)))} · " + + (f"suggest {sug[0]:g} GB {fmt_dur(sug[1])}" if sug else f"suggest - (< {HIST_MIN_RUNS} runs)")) + return out or ["no finished runs with a peak yet"] + + +# ---------------------------------------------------------------- batch fit (wf batch --left / --time-left) + +TASK_P90_DEFAULT, TASK_P90_MIN_RUNS = 30, 3 # minutes per task when out/wf-orch.log has < 3 timed rows +_ORCH_ROW = re.compile(r"^\d{4}-\d\d-\d\dT\S+ (\S+) \S+ \S+ (\S+) \S+ (?:(\d+)h(\d+)m|(\d+)m(\d+)s)(?:\s|$)") + + +def orch_durations(text: str) -> list[float]: + """Minutes of done*/handback rows of local lanes in out/wf-orch.log with a real (> 0) duration.""" + out = [] + for line in text.splitlines(): + m = _ORCH_ROW.match(line) + if not m or m.group(1) == "cloud" or not (m.group(2).startswith("done") or m.group(2) == "handback"): + continue + h, hm, mm, s = (int(x or 0) for x in m.groups()[2:]) + minutes = h * 60 + hm + mm + s / 60 + if minutes > 0: + out.append(round(minutes, 4)) + return out + + +def task_p90(text: str) -> tuple[int, int]: + """(p90 task minutes rounded up, runs) from out/wf-orch.log text; < 3 runs → (30, runs).""" + durs = orch_durations(text) + if len(durs) < TASK_P90_MIN_RUNS: + return TASK_P90_DEFAULT, len(durs) + return math.ceil(round(pct(durs, 0.90), 6)), len(durs) + + +def batch_fit(k: int, left_min: int, p90_min: int) -> int: + """Tasks that fit: min(K, floor(left / p90)).""" + return max(0, min(k, left_min // max(1, p90_min))) -- cgit