"""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) commit: str = "" # full sha the job tests (--commit / gate title): coalesce key covers: list = field(default_factory=list) # commits of superseded/covered same-lock entries, oldest queued first superseded_by: str = "" # id of the queued entry on a descendant commit that replaced this one 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)") for e in led.entries: # a queued job of a batch that ended: nobody waits for it, never let it head the queue if e.state == "queued" and parent_of(led, e) is None and (e.by or {}).get("batch"): _finish(e, now, "owner gone") out.append(f"{e.id} dropped (batch {e.by['batch']} gone)") 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 parent_of(led: Ledger, e: Entry) -> Entry | None: """The live job (e.g. wf batch) e was started from (by.batch = its WF_RES_ID), else None.""" pid = (e.by or {}).get("batch") return next((p for p in led.entries if p.id == pid and p.state in ("running", "note")), None) if pid else None def loan(led: Ledger, e: Entry, lent: dict | None = None) -> tuple[float, int]: """Room e borrows from its parent job's claim: a batch's gate runs inside the batch's reservation (the batch waits for it, so a gate bigger than the free budget would otherwise deadlock). The claim is lent once: running children of the parent and loans in `lent` (id → (gb, cpus), this pass) are subtracted.""" p = parent_of(led, e) if p is None: return 0.0, 0 kids = [k for k in led.entries if k is not e and k.state == "running" and (k.by or {}).get("batch") == p.id] gb, cpus = (lent or {}).get(p.id, (0.0, 0)) return (max(0.0, p.mem_gb - gb - sum(k.mem_gb for k in kids)), max(0, p.cpus - cpus - sum(k.cpus for k in kids))) 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) COMMIT_TITLE_RE = re.compile(r"\b[0-9a-f]{7,40}\b") # 'gate 1a2b3c4': the commit a locked job tests def coalesce(led: Ledger, new: Entry, is_ancestor) -> tuple[Entry | None, list[Entry]]: """Same-lock queued entries vs a new queued entry on commit new.commit (ruling 2026-10-07: the queue coalesces, one queued gate per lock and line of history). Returns (cover, gone): cover = a queued entry on the same or a descendant commit (the new one is redundant, never queued); else gone = queued entries on ancestor commits (the new one replaces them). is_ancestor(a, b) = a is an ancestor of b. Running entries are never touched.""" if not (new.lock and new.commit): return None, [] same = [e for e in _queue_of(led.entries) if e is not new and e.lock == new.lock and e.commit] for e in same: if e.commit == new.commit or is_ancestor(new.commit, e.commit): return e, [] return None, [e for e in same if is_ancestor(e.commit, new.commit)] def supersede(new: Entry, gone: list[Entry], now: dt.datetime) -> None: """new replaces gone (queue order): they end 'superseded by new', new inherits their commits and the earliest queue place (a coalesced gate never loses its turn).""" for e in gone: _finish(e, now, f"superseded by {new.id}") e.superseded_by = new.id new.covers = _uniq(new.covers + e.covers + [e.commit], new.commit) if e.queued and (new.queued is None or e.queued < new.queued): new.queued = e.queued def cover(e: Entry, commit: str, covers: list[str] = ()) -> None: """e (queued, descendant commit) also stands for commit.""" e.covers = _uniq(e.covers + list(covers) + [commit], e.commit) def _uniq(commits: list[str], own: str) -> list[str]: return [c for c in dict.fromkeys(commits) if c and c != own] def covers_text(e: Entry) -> str: """' covers a1,b2 (red: culprit is one of them or ; bisect …)' for a coalesced entry, else ''.""" if not e.covers: return "" out = f"; covers {','.join(c[:7] for c in e.covers)}" if e.state == "done" and e.rc not in (0, None) and not e.superseded_by: out += (f"; red: culprit is any commit in {e.covers[0][:7]}^..{e.commit[:7]} -> bisect (git bisect start " f"{e.commit[:7]} {e.covers[0][:7]}^, rerun the job per step) or name that range in the P0 fix task") return out 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}{covers_text(e)}" 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} lent: dict[str, tuple[float, int]] = {} for e in _queue(led): if e.lock and e.lock in held: continue l_gb, l_cpus = loan(led, e, lent) if not fits(e.mem_gb, e.cpus, (b_gb + l_gb, b_cpus + l_cpus)): break out.append(e) held.add(e.lock) u_gb, u_cpus = min(e.mem_gb, l_gb), min(e.cpus, l_cpus) # the loan first, then the free budget if u_gb or u_cpus: p = (e.by or {})["batch"] g0, c0 = lent.get(p, (0.0, 0)) lent[p] = (g0 + u_gb, c0 + u_cpus) b_gb, b_cpus = b_gb - (e.mem_gb - u_gb), b_cpus - (e.cpus - u_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) l_gb, l_cpus = loan(led, e) b_gb, b_cpus = b_gb + l_gb, b_cpus + l_cpus 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)))