import datetime as dt import io import json import os import subprocess import sys import tempfile import unittest os.environ.pop("FORCE_COLOR", None) # argparse colors help under FORCE_COLOR from contextlib import redirect_stderr, redirect_stdout from pathlib import Path HERE = Path(__file__).resolve().parent.parent sys.path.insert(0, str(HERE)) import wf_res # noqa: E402 from wflib import res as R # noqa: E402 UTC = dt.timezone.utc GIB_KB = 1024 * 1024 class Fake: """Records argv; answers systemctl show from self.units {name: (ActiveState, MemoryCurrent bytes|None)}.""" def __init__(self): self.calls, self.units, self.fail, self.journal, self.cgroups, self.bare = [], {}, {}, {}, {}, {} self.scopes = {} # name → {prop: value} for agent session scopes self.manager_env = "PATH=/usr/bin\nHOME=/h/u\n" self.git = {} # git subcommand (branch | rev-parse) → stdout self.real_git = False # True: git calls run for real (cwd = a synthetic repo) def __call__(self, argv): self.calls.append(list(argv)) key = argv[0] if argv[0] != "systemctl" else argv[2] if key in self.fail: return subprocess.CompletedProcess(argv, 1, "", self.fail[key] + "\n") out = "" if argv[0] == "git" and self.real_git: return subprocess.run(argv, capture_output=True, text=True) if argv[:3] == ["systemctl", "--user", "show"]: blocks = [] for name in argv[5:]: if name in self.scopes: blocks.append("\n".join([f"Id={name}"] + [f"{k}={v}" for k, v in self.scopes[name].items()])) continue if name in self.bare: blocks.append(f"Id={name}\nTransient=yes\nSlice=app.slice\nMemoryCurrent={self.bare[name]}") continue state, cur = self.units.get(name, ("inactive", None)) blocks.append(f"Id={name}\nActiveState={state}\nMemoryCurrent={'[not set]' if cur is None else cur}" f"\nControlGroup={self.cgroups.get(name, '')}") out = "\n\n".join(blocks) + "\n" elif argv[:3] == ["systemctl", "--user", "list-units"]: names = self.scopes if "--type=scope" in argv else self.bare out = "".join(f"{n} loaded active running x\n" for n in names) elif argv[:3] == ["systemctl", "--user", "set-property"]: if argv[4] in self.scopes: self.scopes[argv[4]].update(a.split("=", 1) for a in argv[5:]) elif argv[:3] == ["systemctl", "--user", "show-environment"]: out = self.manager_env elif argv[0] == "git": out = self.git.get(argv[3], "") elif argv[0] == "journalctl": out = self.journal.get(argv[argv.index("-u") + 1], "") elif argv[0] == "systemd-run": unit = next(a for a in argv if a.startswith("--unit="))[len("--unit="):] self.units[unit] = ("active", 0) elif argv[:3] == ["systemctl", "--user", "stop"]: self.units[argv[3]] = ("inactive", None) return subprocess.CompletedProcess(argv, 0, out, "") def ran(self, word): return [c for c in self.calls if word in c[:7]] class Box: """Temp dirs + fake machine; .env is a wf_res.Env.""" def __init__(self, tmp: Path, available_gb=20, total_gb=30, now=dt.datetime(2026, 10, 1, 14, 0, tzinfo=UTC)): self.tmp, self.fake, self.clock, self.alive = tmp, Fake(), [now], {4242} self.available_gb, self.total_gb = available_gb, total_gb (tmp / "proj").mkdir(exist_ok=True) self.env = wf_res.Env( state=tmp / "state", config=tmp / "cfg" / "resources.toml", units=tmp / "units", scratch=tmp / "scratch", run=self.fake, meminfo=lambda: f"MemTotal: {int(self.total_gb * GIB_KB)} kB\nMemAvailable: {int(self.available_gb * GIB_KB)} kB\n", nproc=16, now=lambda: self.clock[0], pid_alive=lambda p: p in self.alive, owner=lambda: 4242, root=lambda: tmp / "proj", cwd=lambda: str(tmp / "proj"), procs=lambda: [], sleep=self.tick, tmp_used_gb=lambda: 0.0, cgread=lambda cg, name: self.psi.get(cg, "") if name == "memory.pressure" else self.stat.get(cg, "")) self.psi, self.stat = {}, {} self.env.live_sessions = lambda: set() self.caller = {"PATH": "/usr/bin", "HOME": "/h/u"} self.env.environ = lambda: dict(self.caller) self.env.tmp_dirs = [] def tick(self, seconds): self.clock[0] += dt.timedelta(seconds=seconds) def wf(self, *argv): out, err = io.StringIO(), io.StringIO() with redirect_stdout(out), redirect_stderr(err): code = wf_res.main(list(argv), self.env) return code, out.getvalue(), err.getvalue() def ledger(self): return R.loads((self.env.state / "resources.json").read_text()) class IOBase(unittest.TestCase): def setUp(self): self._tmp = tempfile.TemporaryDirectory() self.box = Box(Path(self._tmp.name)) def tearDown(self): self._tmp.cleanup() class Run(IOBase): def test_run_starts_unit(self): code, out, err = self.box.wf("run", "--mem", "4G", "--for", "40m", "--title", "build", "--", "make", "-j8") self.assertEqual((code, err), (0, "")) log = self.box.env.state / "logs" / "r-1.log" self.assertEqual(out, f"r-1 started; log {log}; ETA ~14:40 (estimate: not killed when over)\n") argv = self.box.fake.ran("systemd-run")[0] self.assertEqual(argv[-3:], ["sh", "make", "-j8"]) self.assertIn("--unit=wf-r-1.service", argv) e = self.box.ledger().get("r-1") self.assertEqual((e.state, e.project, e.owner, e.cwd), ("running", "proj", 4242, str(self.box.tmp / "proj"))) self.assertTrue((self.box.env.units / "agents.slice").exists()) def test_run_keeps_caller_env(self): self.box.caller.update({"APP_DIR": "/data/arc", "PATH": "/venv/bin:/usr/bin", "PWD": "/x", "SHLVL": "2", "CLAUDE_CODE_MESSAGING_TOKEN": "t", "GH_TOKEN": "t", "DB_PASSWORD": "p"}) self.box.wf("run", "--mem", "1G", "--for", "1m", "--title", "x", "--", "x") argv = self.box.fake.ran("systemd-run")[0] self.assertEqual([a for a in argv if a.startswith("--setenv=")], ["--setenv=APP_DIR=/data/arc", "--setenv=PATH=/venv/bin:/usr/bin", "--setenv=WF_RES_ID=r-1"]) self.assertEqual(self.box.ledger().get("r-1").env, {"APP_DIR": "/data/arc", "PATH": "/venv/bin:/usr/bin"}) def test_run_records_owner(self): main = self.box.tmp / "main" (main / ".wf" / "sessions").mkdir(parents=True) (main / ".wf" / "sessions" / "slow.json").write_text('{"lane": "slow", "pid": 77, "socket": "/s/o"}\n') self.box.fake.git = {"branch": "slow/t-big\n", "rev-parse": f"{main}/.git\n"} self.box.caller.update({"CLAUDE_PID": "77", "CLAUDE_CODE_MESSAGING_SOCKET": "/s/o", "WF_RES_ID": "r-9"}) self.box.wf("run", "--mem", "1G", "--for", "1m", "--title", "x", "--", "x") self.assertEqual(self.box.ledger().get("r-1").by, {"name": "slow session", "task": "t-big", "batch": "r-9", "address": "uds:/s/o"}) argv = self.box.fake.ran("systemd-run")[0] self.assertEqual([a for a in argv if a.startswith("--setenv=WF_RES_ID")], ["--setenv=WF_RES_ID=r-1"]) self.assertIn('"x" running', self.box.wf("status")[1]) self.assertIn("[by slow session t-big batch r-9, message uds:/s/o]", self.box.wf("status")[1]) self.box.wf("note", "--mem", "1G", "--for", "5m", "--by", "pilot", "edit") self.assertEqual(self.box.ledger().get("r-2").by["name"], "pilot") def test_queued_run_starts_later_with_its_env(self): self.box.wf("run", "--mem", "10G", "--for", "40m", "--title", "big", "--", "x") self.box.caller["FOO"] = "1" self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "q", "--queue", "--", "y") self.box.caller.pop("FOO") self.box.wf("release", "r-1", "--stop") self.box.wf("status") starts = self.box.fake.ran("systemd-run") self.assertEqual([[a for a in s if a.startswith("--setenv=")] for s in starts], [["--setenv=WF_RES_ID=r-1"], ["--setenv=FOO=1", "--setenv=WF_RES_ID=r-2"]]) def test_busy_exit_3(self): self.box.wf("run", "--mem", "10G", "--for", "40m", "--title", "big", "--", "x") code, out, _ = self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "two", "--", "y") # 20 − 6 − 2 − 10 unused = 2.0 free self.assertEqual(code, 3) self.assertEqual(out, 'busy: 10.0 GB held by proj "big" (r-1) until ~14:40; 2.0 GB free for agents; ' 'retry after ~14:40 or work on something else; 12.0 GB really free beyond the reserve: ' '--force starts it past the ledger (only if the holders will not use what they reserved)\n') self.assertEqual([e.id for e in self.box.ledger().entries], ["r-1"]) def test_force_runs_past_unused_claim(self): self.box.wf("run", "--mem", "10G", "--for", "40m", "--title", "big", "--", "x") code, out, _ = self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "two", "--force", "--", "y") # really free: 20 − 6 reserve − 2 headroom = 12 ≥ 8 (r-1's unused 10 GB ignored) self.assertEqual(code, 0) self.assertTrue(out.startswith("r-2 started"), out) self.assertEqual([(e.id, e.state) for e in self.box.ledger().entries], [("r-1", "running"), ("r-2", "running")]) def test_force_refused_when_memory_really_used(self): self.box.available_gb = 12 # really free beyond reserve: 12 − 6 − 2 = 4 self.box.wf("run", "--mem", "3G", "--for", "40m", "--title", "a", "--", "x") code, out, _ = self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "two", "--", "y") self.assertEqual(code, 3) self.assertNotIn("--force", out) # hint only when --force would fit code, out, _ = self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "two", "--force", "--", "y") self.assertEqual((code, out), (3, "busy even with --force: only 4.0 GB really free beyond the reserve; " "retry later or work on something else\n")) self.assertEqual([e.id for e in self.box.ledger().entries], ["r-1"]) def test_never_fits(self): code, _, err = self.box.wf("run", "--mem", "25G", "--for", "1h", "--title", "x", "--", "x") self.assertEqual((code, err), (1, "wf: 25.0 GB can never fit (max 22.0 GB for agents)\n")) def test_systemd_run_failure(self): self.box.fake.fail["systemd-run"] = "Failed to start transient service unit: boom" code, _, err = self.box.wf("run", "--mem", "1G", "--for", "1m", "--title", "x", "--", "x") self.assertEqual((code, err), (1, "wf: systemd-run: Failed to start transient service unit: boom\n")) self.assertEqual(self.box.ledger().entries, []) def test_no_command(self): code, _, err = self.box.wf("run", "--mem", "1G", "--for", "1m", "--title", "x") self.assertEqual((code, err), (1, "wf: no command after --\n")) def test_usage_error_exit_2(self): with redirect_stderr(io.StringIO()): self.assertEqual(self.box.wf("run", "--for", "1m")[0], 2) class Status(IOBase): def test_status_and_exit_pruned(self): self.box.wf("run", "--mem", "4G", "--for", "40m", "--title", "build", "--", "make") self.box.fake.units["wf-r-1.service"] = ("active", 2 * 1024 ** 3) code, out, _ = self.box.wf("status") self.assertEqual(out.splitlines()[0], 'r-1 proj "build" running 4.0 GB used 2.0 GB 1 cpu since 14:00 ETA ~14:40') logs = self.box.env.state / "logs" (logs / "r-1.rc").write_text("0\n") (logs / "r-1.peak").write_text(str(3 * 1024 ** 3) + "\n") self.box.fake.units["wf-r-1.service"] = ("inactive", None) code, out, _ = self.box.wf("status") self.assertEqual(out.splitlines()[0], 'r-1 proj "build" done (exited) rc=0 peak 3.0 GB at 14:00') def test_status_warns_throttled(self): self.box.wf("run", "--mem", "4G", "--for", "40m", "--title", "build", "--", "make") self.box.fake.units["wf-r-1.service"] = ("active", int(3.8 * 1024 ** 3)) self.box.fake.cgroups["wf-r-1.service"] = "/x/wf-r-1.service" self.box.psi["/x/wf-r-1.service"] = "some avg10=50.00 avg60=40.00 avg300=10.00 total=1\n" out = self.box.wf("status")[1] self.assertIn("r-1 throttled at its memory limit (3.8 of 4.0 GB, stalled 40% of the last minute): " "likely too small; `wf res release r-1 --stop` and re-run with a bigger --mem\n", out) def test_status_warns_bare_units(self): self.box.fake.bare["wh-t27.service"] = 2 * 1024 ** 3 out = self.box.wf("status")[1] self.assertIn("warning: 1 jobs outside wf res (bare systemd-run): wh-t27.service 2.0 GB; " "start jobs with wf res run, also from project scripts\n", out) self.assertIn(["systemctl", "--user", "list-units", "--type=service", "--state=running", "--no-legend", "--plain"], self.box.fake.calls) def test_status_splits_slice_cache(self): self.box.fake.units["agents.slice"] = ("active", 5 * 1024 ** 3) self.box.fake.cgroups["agents.slice"] = "/a.slice" self.box.stat["/a.slice"] = f"anon {3 * 1024 ** 3}\nfile {2 * 1024 ** 3}\nshmem {1024 ** 3 // 2}\n" self.assertIn("unreserved agent memory 5.0 GB (1.5 GB of it file cache, reclaimable)\n", self.box.wf("status")[1]) def test_status_id_and_json(self): self.box.wf("run", "--mem", "4G", "--for", "40m", "--title", "build", "--", "make", "a b") out = self.box.wf("status", "r-1")[1] self.assertIn("cmd: make 'a b'", out) self.assertEqual(json.loads(self.box.wf("status", "--json")[1])["entries"][0]["id"], "r-1") def test_show_failure_keeps_ledger(self): self.box.wf("run", "--mem", "4G", "--for", "40m", "--title", "build", "--", "make") self.box.fake.fail["show"] = "Failed to connect to bus" code, _, err = self.box.wf("status") self.assertEqual((code, err), (1, "wf: systemctl show: Failed to connect to bus\n")) self.assertEqual(self.box.ledger().get("r-1").state, "running") def test_corrupt_ledger_moved_aside(self): self.box.env.state.mkdir(parents=True) (self.box.env.state / "resources.json").write_text("{oops") code, out, err = self.box.wf("status") self.assertEqual(code, 0) self.assertRegex(err, r"^wf: corrupt ledger: .*; moved to resources\.json\.bad-20261001-140000, starting empty\n$") self.assertTrue((self.box.env.state / "resources.json.bad-20261001-140000").exists()) self.assertEqual(out.splitlines()[0], "no reservations") class LedgerCompat(IOBase): """Readers of another code version (a long `wf res wait` started before a release) must keep entries.""" def _write(self, extra: dict, drop=()): self.box.wf("run", "--mem", "1G", "--for", "5m", "--title", "t", "--", "x") path = self.box.env.state / "resources.json" d = json.loads(path.read_text()) for e in d["entries"]: for k in drop: e.pop(k, None) e.update(extra) path.write_text(json.dumps(d)) return path def _kept(self, path): for argv in (("status",), ("tick",), ("wait", "r-1", "--timeout", "1m")): self.box.wf(*argv) self.assertEqual(self.box.ledger().get("r-1").state, "running", argv) self.assertEqual(sorted(p.name for p in path.parent.glob("resources.json*")), ["resources.json"]) self.assertIn("r-1", self.box.wf("status")[1]) def test_pre_owner_ledger_kept(self): # written before e84d7ca/9533211: no by field self._kept(self._write({}, drop=("by", "env"))) def test_future_fields_kept(self): # written by newer code: unknown keys self._kept(self._write({"by": {"name": "x"}, "added_later": 1})) def test_corrupt_ledger_keeps_id_counter(self): self.box.env.state.mkdir(parents=True) (self.box.env.state / "resources.json").write_text('{"next": 652, "entries": [{"id": "r-651"}]}') self.box.wf("run", "--mem", "1G", "--for", "5m", "--title", "t", "--", "x") self.assertEqual([e.id for e in self.box.ledger().entries], ["r-652"]) def test_start_clears_stale_results(self): logs = self.box.env.state / "logs" logs.mkdir(parents=True) (logs / "r-1.rc").write_text("1\n") (logs / "r-1.peak").write_text("5\n") self.box.wf("run", "--mem", "1G", "--for", "5m", "--title", "t", "--", "x") self.assertFalse((logs / "r-1.rc").exists() or (logs / "r-1.peak").exists()) class HookIO(IOBase): def hook(self, text): self.box.env.stdin = lambda: text return self.box.wf("hook") def test_game_on_rewrites_bash(self): self.box.wf("game", "on", "--for", "1h") code, out, err = self.hook('{"tool_name": "Bash", "tool_input": {"command": "./app"}}') self.assertEqual((code, err), (0, "")) self.assertEqual(json.loads(out)["hookSpecificOutput"]["updatedInput"], {"command": "unset DISPLAY WAYLAND_DISPLAY; ./app"}) calls = len(self.box.fake.calls) self.hook('{"tool_name": "Bash", "tool_input": {"command": "./app"}}') self.assertEqual(len(self.box.fake.calls), calls) # read-only: no systemctl, no prune def test_silent_otherwise(self): self.assertEqual(self.hook('{"tool_name": "Bash", "tool_input": {"command": "./app"}}'), (0, "", "")) self.box.wf("game", "on") self.assertEqual(self.hook("not json"), (0, "", "")) (self.box.env.state / "resources.json").write_text("{oops") self.assertEqual(self.hook('{"tool_name": "Bash", "tool_input": {"command": "x"}}'), (0, "", "")) self.assertTrue((self.box.env.state / "resources.json").exists()) # hook never moves a bad ledger class Release(IOBase): def test_release_running_needs_stop(self): self.box.wf("run", "--mem", "4G", "--for", "40m", "--title", "build", "--", "make") code, _, err = self.box.wf("release", "r-1") self.assertEqual((code, err), (1, "wf: r-1 is running; --stop to kill it\n")) code, out, _ = self.box.wf("release", "r-1", "--stop") self.assertEqual((code, out), (0, "r-1 stopped\n")) self.assertEqual(self.box.fake.ran("stop")[0], ["systemctl", "--user", "stop", "wf-r-1.service"]) e = self.box.ledger().get("r-1") self.assertEqual((e.state, e.why), ("done", "released")) def test_release_unknown(self): self.assertEqual(self.box.wf("release", "r-7")[2], "wf: no entry 'r-7'\n") class Dispatch(unittest.TestCase): def test_wf_forwards_res(self): r = subprocess.run([sys.executable, str(HERE / "wf.py"), "res", "-h"], capture_output=True, text=True) self.assertEqual(r.returncode, 0) self.assertIn("usage: wf res", r.stdout) class Wait(IOBase): def test_wait_until_done(self): self.box.wf("run", "--mem", "1G", "--for", "5m", "--title", "t", "--", "x") logs = self.box.env.state / "logs" polls = [] def sleep(seconds): polls.append(seconds) self.box.tick(seconds) if len(polls) == 2: (logs / "r-1.rc").write_text("0\n") (logs / "r-1.peak").write_text(str(1024 ** 3 // 2) + "\n") self.box.fake.units["wf-r-1.service"] = ("inactive", None) self.box.env.sleep = sleep code, out, _ = self.box.wf("wait", "r-1") self.assertEqual((code, out, polls), (0, "r-1 done rc=0 peak 0.5 GB in 0 min\n", [15, 15])) def test_wait_warns_throttled_once(self): self.box.wf("run", "--mem", "4G", "--for", "5m", "--title", "t", "--", "x") self.box.fake.units["wf-r-1.service"] = ("active", int(3.8 * 1024 ** 3)) self.box.fake.cgroups["wf-r-1.service"] = "/x/wf-r-1.service" self.box.psi["/x/wf-r-1.service"] = "some avg10=50.00 avg60=40.00 avg300=10.00 total=1\n" code, _, err = self.box.wf("wait", "r-1", "--timeout", "1m") self.assertEqual(err, "wf: r-1 throttled at its memory limit (3.8 of 4.0 GB, stalled 40% of the last minute): " "likely too small; `wf res release r-1 --stop` and re-run with a bigger --mem\n" "wf: r-1 not done after 1m (running)\n") def test_wait_timeout(self): self.box.wf("run", "--mem", "1G", "--for", "5m", "--title", "t", "--", "x") code, _, err = self.box.wf("wait", "r-1", "--timeout", "1m") self.assertEqual((code, err), (1, "wf: r-1 not done after 1m (running)\n")) def test_wait_killed_job(self): self.box.wf("run", "--mem", "1G", "--for", "5m", "--title", "t", "--", "x") self.box.fake.units["wf-r-1.service"] = ("failed", None) self.assertEqual(self.box.wf("wait", "r-1")[1], "r-1 done rc=? peak ? in 0 min (no exit code; see journalctl --user -u wf-r-1.service)\n") def test_wait_oom_killed_job(self): self.box.wf("run", "--mem", "4G", "--for", "5m", "--title", "t", "--", "x") self.box.fake.units["wf-r-1.service"] = ("failed", None) self.box.fake.journal["wf-r-1.service"] = ( "wf-r-1.service: systemd-oomd killed some process(es) in this unit.\n" "wf-r-1.service: Failed with result 'oom-kill'.\n" "wf-r-1.service: Consumed 11min 2.837s CPU time, 3.9G memory peak.\n") self.assertEqual(self.box.wf("wait", "r-1")[1], "r-1 done rc=? peak 3.9 GB in 0 min " "(killed: oom-kill by systemd-oomd, limit 4.0 GB; raise --mem)\n") j = self.box.fake.ran("journalctl")[0] self.assertEqual(j, ["journalctl", "--user", "-u", "wf-r-1.service", "-o", "cat", "--no-pager", "--since", "@" + str(int(dt.datetime(2026, 10, 1, 14, 0, tzinfo=UTC).timestamp()))]) out = self.box.wf("status")[1] self.assertEqual(out.splitlines()[0], 'r-1 proj "t" done (killed: oom-kill by systemd-oomd, limit 4.0 GB; ' 'raise --mem) rc=? peak 3.9 GB at 14:00') def test_rc_file_skips_journal(self): self.box.wf("run", "--mem", "1G", "--for", "5m", "--title", "t", "--", "x") (self.box.env.state / "logs" / "r-1.rc").write_text("2\n") self.box.fake.units["wf-r-1.service"] = ("inactive", None) self.assertEqual(self.box.wf("wait", "r-1")[1], "r-1 done rc=2 peak ? in 0 min\n") self.assertEqual(self.box.fake.ran("journalctl"), []) class QueueIO(IOBase): def test_queue_then_start_when_free(self): self.box.wf("run", "--mem", "10G", "--for", "40m", "--title", "big", "--", "x") code, out, _ = self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "two", "--queue", "--", "y") self.assertEqual((code, out), (0, "r-2 queued, position 1; est. start ~14:40; cancel: wf res release r-2\n")) self.box.fake.units["wf-r-1.service"] = ("inactive", None) self.box.wf("status") e = self.box.ledger().get("r-2") self.assertEqual((e.state, e.unit), ("running", "wf-r-2.service")) self.assertEqual(self.box.fake.ran("systemd-run")[-1][-2:], ["sh", "y"]) def test_batch_gate_borrows_batch_claim(self): self.box.wf("run", "--mem", "10G", "--for", "40m", "--title", "wf batch", "--", "x") self.box.caller["WF_RES_ID"] = "r-1" # the gate is queued from inside batch r-1 code, out, _ = self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "gate", "--queue", "--", "y") self.assertEqual((code, out.split(";")[0]), (0, "r-2 started")) self.assertEqual(self.box.ledger().get("r-2").by["batch"], "r-1") def test_direct_run_does_not_jump_queue(self): self.box.wf("run", "--mem", "10G", "--for", "40m", "--title", "big", "--", "x") self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "two", "--queue", "--", "y") code, out, _ = self.box.wf("run", "--mem", "1G", "--for", "10m", "--title", "small", "--", "z") self.assertEqual(code, 3) def test_queued_start_failure_moves_on(self): self.box.wf("run", "--mem", "10G", "--for", "40m", "--title", "big", "--", "x") self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "two", "--queue", "--", "y") self.box.fake.units["wf-r-1.service"] = ("inactive", None) self.box.fake.fail["systemd-run"] = "boom" code, out, err = self.box.wf("status") self.assertEqual(code, 0) e = self.box.ledger().get("r-2") self.assertEqual((e.state, e.why), ("done", "start failed: systemd-run: boom")) self.box.fake.fail.clear() self.assertEqual(self.box.wf("status")[0], 0) class Note(IOBase): def test_note_and_owner_exit(self): code, out, _ = self.box.wf("note", "--mem", "3G", "--for", "20m", "dotnet test") self.assertEqual((code, out), (0, "r-1 noted 3.0 GB until ~14:20 (then freed; nothing is killed)\n")) self.assertEqual(self.box.ledger().get("r-1").owner, 4242) self.box.alive.clear() self.box.wf("status") self.assertEqual(self.box.ledger().get("r-1").why, "owner gone") def test_note_busy(self): self.box.wf("note", "--mem", "10G", "--for", "20m", "a") code, out, _ = self.box.wf("note", "--mem", "3G", "--for", "20m", "b") self.assertEqual((code, out), (3, 'busy: 10.0 GB held by proj "a" (r-1) until ~14:20; 2.0 GB free for agents; ' 'retry after ~14:20 or work on something else; 12.0 GB really free ' 'beyond the reserve: --force starts it past the ledger (only if the ' 'holders will not use what they reserved)\n')) code, out, _ = self.box.wf("note", "--mem", "3G", "--for", "20m", "--force", "b") self.assertEqual(code, 0) self.assertTrue(out.startswith("r-2 noted 3.0 GB"), out) def _race_child(tmp, q): import time as _t box = Box(Path(tmp), available_gb=14) # budget 14 − 6 − 2 = 6 GB: one 5G note fits, not two slow = box.env.meminfo box.env.meminfo = lambda: (_t.sleep(0.3), slow())[1] with redirect_stdout(io.StringIO()), redirect_stderr(io.StringIO()): q.put(wf_res.main(["note", "--mem", "5G", "--for", "10m", "x"], box.env)) class LockRace(unittest.TestCase): def test_two_processes_never_overbook(self): import multiprocessing as mp ctx = mp.get_context("fork") with tempfile.TemporaryDirectory() as tmp: q = ctx.Queue() ps = [ctx.Process(target=_race_child, args=(tmp, q)) for _ in range(2)] for p in ps: p.start() for p in ps: p.join(10) self.assertEqual([p.exitcode for p in ps], [0, 0]) # a crashed child must fail, not hang q.get() self.assertEqual(sorted([q.get(timeout=5), q.get(timeout=5)]), [0, 3]) class GameIO(IOBase): def props(self): return [c[5:] for c in self.box.fake.ran("set-property")] def test_on_off(self): code, out, _ = self.box.wf("game", "on") self.assertEqual((code, out), (0, "game on until 18:00 (4h); CPU/IO now yours\n")) self.assertEqual(self.props()[-1], ["CPUWeight=5", "IOWeight=5", f"MemoryHigh={18 * 1024 ** 3}"]) self.assertEqual(self.box.fake.ran("set-property")[-1][:5], ["systemctl", "--user", "set-property", "--runtime", "agents.slice"]) self.assertEqual(self.box.wf("game", "off")[1], "game off\n") self.assertEqual(self.props()[-1], ["CPUWeight=20", "IOWeight=20", f"MemoryHigh={24 * 1024 ** 3}"]) def test_expiry_restores(self): self.box.wf("game", "on", "--for", "1h") self.box.tick(3601) self.box.wf("status") self.assertIsNone(self.box.ledger().game_until) self.assertEqual(self.props()[-1][0], "CPUWeight=20") def test_game_refuses_new_job(self): self.box.wf("game", "on") # 20 − 12 − 2 = 6 GB budget self.assertEqual(self.box.wf("run", "--mem", "7G", "--for", "5m", "--title", "x", "--", "x")[0], 3) def test_on_again_extends(self): self.box.wf("game", "on", "--for", "2h") self.box.wf("game", "on", "--for", "1h") self.assertEqual(self.box.ledger().game_until.hour, 16) class CleanIO(IOBase): def make(self, rel, age_h, size=10): p = self.box.env.scratch / rel p.parent.mkdir(parents=True, exist_ok=True) p.write_bytes(b"x" * size) t = self.box.clock[0].timestamp() - age_h * 3600 os.utime(p, (t, t)) for d in [p.parent, *p.parent.parents]: if d == self.box.env.scratch.parent: break os.utime(d, (t, t)) return p def test_auto_clean_sweeps_tmp_litter(self): t = self.box.tmp / "tmpdir" (t / "old-empty").mkdir(parents=True) (t / "full").mkdir() (t / "full" / "x").write_text("x") (t / "clr-debug-pipe-999999-1-in").write_text("") (t / "keep.txt").write_text("x") old = self.box.clock[0].timestamp() - 3 * 3600 for n in ("old-empty", "full", "keep.txt"): os.utime(t / n, (old, old)) self.box.env.tmp_dirs = [t] self.box.env.pid_alive = lambda p: p != 999999 and p in self.box.alive self.box.wf("status") self.assertEqual(sorted(os.listdir(t)), ["full", "keep.txt"]) def test_auto_clean_keeps_live_session(self): live = self.make("-projects-a/live-1/scratchpad/notes.txt", 9) dead = self.make("-projects-a/dead-2/scratchpad/x.txt", 9) self.box.env.live_sessions = lambda: {"live-1"} self.box.wf("status") self.assertEqual((live.exists(), dead.exists()), (True, False)) def test_auto_clean_on_status_and_throttle(self): old = self.make("-projects-a/s1/scratchpad/big.bin", 3) new = self.make("-projects-b/s2/scratchpad/x.txt", 0.5) self.box.wf("status") self.assertFalse(old.exists()) self.assertTrue(new.exists()) self.assertTrue(any(c[2] == "reset-failed" for c in self.box.fake.calls if c[0] == "systemctl")) again = self.make("-projects-c/s3/a.txt", 3) self.box.wf("status") self.assertTrue(again.exists()) # < 10 min since last clean self.box.tick(601) self.box.wf("status") self.assertFalse(again.exists()) def test_clean_prints_and_lists_project_patterns(self): self.make("-projects-a/s1/f.bin", 3, size=2048) proj = self.box.tmp / "proj" (proj / "workflow.toml").write_text('cleanup = ["out/prof", "out/logs/*.log:30d"]\n') (proj / "out" / "prof").mkdir(parents=True) (proj / "out" / "prof" / "p.dat").write_bytes(b"x" * 1024) logs = proj / "out" / "logs" logs.mkdir() (logs / "new.log").write_text("n") old = logs / "old.log" old.write_text("o") t = self.box.clock[0].timestamp() - 31 * 86400 os.utime(old, (t, t)) (logs / "link.log").symlink_to("/etc/hostname") code, out, _ = self.box.wf("clean") self.assertEqual(out.splitlines(), [ f"freed 0 MB: {self.box.env.scratch / '-projects-a'}", "would delete 0 MB: out/prof (wf res clean --yes)", "would delete 0 MB: out/logs/old.log (wf res clean --yes)"]) self.assertTrue((proj / "out" / "prof").exists()) self.box.wf("clean", "--yes") self.assertFalse((proj / "out" / "prof").exists()) self.assertFalse(old.exists()) self.assertTrue((logs / "new.log").exists()) self.assertTrue((logs / "link.log").is_symlink()) def test_status_warnings(self): self.box.env.tmp_used_gb = lambda: 8.0 self.box.env.procs = lambda: [(10, 1, "claude", "/u/app.slice/tab.scope"), (20, 1, "claude", "/u/agents.slice/wf-claude-20.scope")] out = self.box.wf("status")[1].splitlines() self.assertEqual(out[-2:], ["warning: /tmp (RAM) holds 8.0 GB; wf res clean", "warning: 1 claude sessions outside agents.slice (wf res adopt)"]) class TimerIO(IOBase): def test_timer_on_off(self): code, out, _ = self.box.wf("timer", "on") self.assertEqual((code, out), (0, "timer on: wf-res.timer every 1 min\n")) units = self.box.env.units self.assertIn("res tick", (units / "wf-res.service").read_text()) self.assertTrue((units / "wf-res.timer").exists()) self.assertIn(["systemctl", "--user", "enable", "--now", "wf-res.timer"], self.box.fake.calls) self.assertEqual(self.box.wf("timer", "off")[1], "timer off\n") self.assertFalse((units / "wf-res.timer").exists()) self.assertIn(["systemctl", "--user", "disable", "--now", "wf-res.timer"], self.box.fake.calls) def test_adopt(self): self.box.env.procs = lambda: [(101, 1, "claude", "/u/app.slice/t.scope"), (102, 101, "bash", "/u/app.slice/t.scope"), (201, 1, "claude", "/u/app.slice/u.scope")] code, out, _ = self.box.wf("adopt") self.assertEqual(out, "adopted claude 101 (2 processes)\nadopted claude 201 (1 processes)\n") self.assertEqual(self.box.fake.ran("StartTransientUnit")[0][8], "wf-claude-101.scope") def test_adopt_failure_reported(self): self.box.env.procs = lambda: [(101, 1, "claude", "/u/app.slice/t.scope")] self.box.fake.fail["busctl"] = "Call failed: No such process" self.assertEqual(self.box.wf("adopt")[1], "claude 101: Call failed: No such process\n") def test_adopt_none(self): self.assertEqual(self.box.wf("adopt")[1], "all claude sessions already in agents.slice\n") def test_tick_adopts_silently(self): self.box.env.procs = lambda: [(101, 1, "claude", "/u/app.slice/t.scope")] self.assertEqual(self.box.wf("tick"), (0, "", "")) self.assertEqual(len(self.box.fake.ran("StartTransientUnit")), 1) def test_tick_caps_sessions(self): self.box.fake.scopes = {"wf-claude-5.scope": {"Slice": "agents.slice", "MemoryHigh": "infinity", "MemoryCurrent": str(1024 ** 3), "ControlGroup": "/a/5"}, "init.scope": {"Slice": "-.slice", "MemoryHigh": "infinity", "MemoryCurrent": "1"}} self.assertEqual(self.box.wf("tick"), (0, "", "")) self.assertEqual(self.box.fake.ran("set-property"), [["systemctl", "--user", "set-property", "--runtime", "wf-claude-5.scope", f"MemoryHigh={6 * 1024 ** 3}"]]) self.box.wf("tick") self.assertEqual(len(self.box.fake.ran("set-property")), 1) # already capped → no call def test_status_warns_session_at_cap(self): self.box.fake.scopes = {"wf-claude-5.scope": {"Slice": "agents.slice", "MemoryHigh": str(6 * 1024 ** 3), "MemoryCurrent": str(int(5.7 * 1024 ** 3)), "ControlGroup": "/a/5"}} self.box.psi["/a/5"] = "some avg10=50.00 avg60=40.00 avg300=10.00 total=1\n" self.assertIn("warning: claude session 5 at its memory cap (5.7 of 6.0 GB, stalled 40% of the last minute): " "run big work with wf res run\n", self.box.wf("status")[1]) def test_tick_silent_and_starts_queue(self): self.box.wf("run", "--mem", "10G", "--for", "40m", "--title", "big", "--", "x") self.box.wf("run", "--mem", "8G", "--for", "10m", "--title", "two", "--queue", "--", "y") self.box.fake.units["wf-r-1.service"] = ("inactive", None) self.assertEqual(self.box.wf("tick"), (0, "", "")) self.assertEqual(self.box.ledger().get("r-2").state, "running") def test_shell_init(self): self.assertEqual(self.box.wf("shell-init")[1], "alias claude='systemd-run --user --scope --quiet --slice=agents.slice claude'\n") class Robust(IOBase): def test_vanished_paths_are_skipped(self): gone = self.box.tmp / "gone" self.assertEqual((wf_res._newest(gone), wf_res._size(gone)), (0.0, 0)) def test_os_error_is_one_line(self): def boom(): raise PermissionError(13, "Permission denied", "/proc/meminfo") self.box.env.meminfo = boom code, _, err = self.box.wf("status") self.assertEqual((code, err), (1, "wf: [Errno 13] Permission denied: '/proc/meminfo'\n")) class CoalesceIO(IOBase): """Batch of commits a ← b ← c, each queueing its gate (real repo): one queued gate per lock, on c.""" def setUp(self): super().setUp() repo = self.box.tmp / "proj" git = lambda *a: subprocess.run(["git", "-C", str(repo), *a], check=True, capture_output=True, text=True).stdout git("init", "-q", "-b", "master") self.sha = {} for c in ("a", "b", "c"): git("-c", "user.name=t", "-c", "user.email=t@t", "commit", "-q", "--allow-empty", "-m", c) self.sha[c] = git("rev-parse", "HEAD").strip() git("checkout", "-q", "-b", "side", self.sha["a"]) git("-c", "user.name=t", "-c", "user.email=t@t", "commit", "-q", "--allow-empty", "-m", "x") self.sha["x"] = git("rev-parse", "HEAD").strip() self.box.fake.real_git = True def gate(self, c, *extra): return self.box.wf("run", "--queue", "--mem", "2G", "--for", "30m", "--title", f"gate {self.sha[c][:7]}", *extra, "--", "x") def test_batch_yields_one_queued_gate(self): self.box.wf("run", "--mem", "2G", "--for", "30m", "--title", "gate running", "--", "x") # r-1 holds the lock self.assertEqual(self.gate("a")[0], 0) self.box.tick(60) self.gate("b") self.box.tick(60) code, out, _ = self.gate("c") self.assertEqual(code, 0) self.assertIn("r-4 queued, position 1 (lock held by r-1)", out) self.assertIn("supersedes r-3", out) led = self.box.ledger() self.assertEqual([(e.id, e.commit) for e in led.entries if e.state == "queued"], [("r-4", self.sha["c"])]) self.assertEqual(led.get("r-4").covers, [self.sha["a"], self.sha["b"]]) self.assertEqual(led.get("r-4").queued, led.get("r-2").queued) self.assertEqual([led.get(i).superseded_by for i in ("r-2", "r-3")], ["r-3", "r-4"]) # chain; wait follows it code, out, _ = self.gate("b") # late ancestor: covered, nothing new self.assertEqual(code, 0) self.assertTrue(out.startswith("r-4 queued"), out) self.assertIn(f"covers {self.sha['b'][:7]}", out) self.assertEqual(len([e for e in self.box.ledger().entries if e.state == "queued"]), 1) code, out, _ = self.gate("x") # side branch: own gate self.assertTrue(out.startswith("r-5 queued"), out) code, out, _ = self.box.wf("run", "--queue", "--mem", "2G", "--for", "30m", "--title", "gate", "--commit", "master", "--", "x") self.assertTrue(out.startswith("r-4 queued"), out) # --commit REV, no hex in title self.assertNotEqual(self.box.wf("run", "--queue", "--mem", "2G", "--for", "30m", "--title", "gate", "--commit", "nope", "--", "x")[0], 0) def test_wait_follows_superseded_and_red_names_range(self): self.box.wf("run", "--mem", "2G", "--for", "30m", "--title", "gate running", "--", "x") self.gate("a") self.gate("c") self.box.fake.units["wf-r-1.service"] = ("inactive", None) self.box.wf("status") # r-3 starts logs = self.box.env.state / "logs" logs.mkdir(parents=True, exist_ok=True) (logs / "r-3.rc").write_text("1\n") self.box.fake.units["wf-r-3.service"] = ("inactive", None) code, out, _ = self.box.wf("wait", "r-2") self.assertEqual(code, 0) self.assertIn("r-2 superseded by r-3", out) self.assertIn(f"r-3 done rc=1", out) self.assertIn(f"{self.sha['a'][:7]}^..{self.sha['c'][:7]}", out) if __name__ == "__main__": unittest.main() class LockIO(IOBase): """--lock / 'gate …' titles: one job per lock key and main tree at a time (shared gate checkout).""" def gate(self, title, *extra): return self.box.wf("run", "--mem", "2G", "--for", "30m", "--title", title, *extra, "--", "x") def test_second_gate_refused_even_with_force(self): self.assertEqual(self.gate("gate c61907f0")[0], 0) for extra in ((), ("--force",)): code, out, err = self.gate("gate 2c91cac7 (fix)", *extra) self.assertEqual(code, 3, err) self.assertIn("lock 'gate' held by r-1", out + err) self.assertIn("--queue", out + err) self.assertEqual(len(self.box.fake.ran("systemd-run")), 1) def test_second_gate_queues_and_starts_after_first(self): self.gate("gate c61907f0") code, out, _ = self.gate("gate 2c91cac7", "--queue") self.assertEqual(code, 0) self.assertTrue(out.startswith("r-2 queued, position 1 (lock held by r-1); est. start ~14:30"), out) self.box.wf("status") # memory is free, lock is not self.assertEqual(self.box.ledger().get("r-2").state, "queued") other = self.box.wf("run", "--mem", "1G", "--for", "5m", "--title", "build", "--", "y") self.assertEqual(other[0], 0) # unlocked work is not held by the locked queue self.box.fake.units["wf-r-1.service"] = ("inactive", None) self.box.wf("status") self.assertEqual(self.box.ledger().get("r-2").state, "running") def test_two_queued_gates_start_one_at_a_time(self): self.gate("gate a") self.gate("gate b", "--queue") self.gate("gate c", "--queue") self.box.fake.units["wf-r-1.service"] = ("inactive", None) self.box.wf("status") led = self.box.ledger() self.assertEqual([led.get(i).state for i in ("r-2", "r-3")], ["running", "queued"]) def test_explicit_lock_and_unlocked_titles(self): self.box.wf("run", "--mem", "2G", "--for", "30m", "--title", "e2e", "--lock", "e2e", "--", "x") self.assertEqual(self.box.wf("run", "--mem", "2G", "--for", "30m", "--title", "e2e 2", "--lock", "e2e", "--", "x")[0], 3) self.assertEqual(self.gate("gate a")[0], 0) # other key self.assertEqual(self.box.wf("run", "--mem", "2G", "--for", "30m", "--title", "gatekeeper", "--", "x")[0], 0) self.assertEqual(self.box.ledger().get("r-1").lock, f"e2e@{self.box.tmp / 'proj'}") def test_lock_scoped_to_main_tree(self): self.gate("gate a") self.box.fake.git["rev-parse"] = str(self.box.tmp / "other" / ".git") # another project's tree self.assertEqual(self.gate("gate b")[0], 0) self.box.fake.git["rev-parse"] = str(self.box.tmp / "proj" / ".git") # lane worktree of proj self.assertEqual(self.gate("gate c")[0], 3) def test_project_is_main_tree_name_from_lane_worktree(self): self.box.fake.git["rev-parse"] = str(self.box.tmp / "home" / ".git") # cwd = lane worktree 'proj' of home self.box.wf("run", "--mem", "1G", "--for", "1m", "--title", "build 1", "--", "x") self.box.fake.git["rev-parse"] = "" # no git → root name self.box.wf("run", "--mem", "1G", "--for", "1m", "--title", "build 2", "--", "x") self.box.fake.git["rev-parse"] = str(self.box.tmp / "home" / ".git" / "modules" / "m") # submodule → root self.box.wf("run", "--mem", "1G", "--for", "1m", "--title", "build 3", "--", "x") self.assertEqual([self.box.ledger().get(f"r-{i}").project for i in (1, 2, 3)], ["home", "proj", "proj"]) def test_percent_args_reach_systemd_verbatim(self): # r-671 'fatal: ambiguous argument %s"': caller quoting, not wf; argv passes through untouched self.box.wf("run", "--mem", "1G", "--for", "5m", "--title", "log", "--", "git", "log", "--format=%h %s", "HEAD") self.assertEqual(self.box.fake.ran("systemd-run")[0][-4:], ["git", "log", "--format=%h %s", "HEAD"]) def hist_rec(id, title="build 1234567", peak=1.0, minutes=10.0, project="proj", mem=4.0, est=40): return {"id": id, "project": project, "title": title, "mem_gb": mem, "est_min": est, "rc": 0, "peak_gb": peak, "min": minutes} def seed_history(box, recs): box.env.state.mkdir(parents=True, exist_ok=True) (box.env.state / "resources-history.jsonl").write_text("".join(json.dumps(r) + "\n" for r in recs)) class HistoryIO(IOBase): def finish(self, id, peak_gb): logs = self.box.env.state / "logs" (logs / f"{id}.rc").write_text("0\n") (logs / f"{id}.peak").write_text(str(int(peak_gb * 1024 ** 3)) + "\n") self.box.fake.units[f"wf-{id}.service"] = ("inactive", None) def test_pruned_run_kept_in_history(self): self.box.wf("run", "--mem", "4G", "--for", "40m", "--title", "build abc1234", "--", "make") self.box.tick(12 * 60) self.finish("r-1", 3.0) self.box.wf("status") hist = self.box.env.state / "resources-history.jsonl" self.assertFalse(hist.exists()) # still in the ledger: not copied yet self.box.tick(25 * 3600) self.box.wf("status") self.assertEqual(self.box.ledger().entries, []) self.assertEqual([json.loads(x) for x in hist.read_text().splitlines()], [ {"id": "r-1", "project": "proj", "title": "build abc1234", "mem_gb": 4.0, "est_min": 40, "rc": 0, "peak_gb": 3.0, "min": 12.0}]) def test_history_trimmed(self): seed_history(self.box, [hist_rec(f"r-{i}") for i in range(100, 100 + R.HIST_KEEP)]) self.box.wf("run", "--mem", "4G", "--for", "40m", "--title", "build", "--", "make") self.finish("r-1", 1.0) self.box.wf("status") self.box.tick(25 * 3600) self.box.wf("status") lines = (self.box.env.state / "resources-history.jsonl").read_text().splitlines() self.assertEqual(len(lines), R.HIST_KEEP) self.assertEqual((json.loads(lines[0])["id"], json.loads(lines[-1])["id"]), ("r-101", "r-1")) def test_run_hint_over_history(self): # peaks 1,2,3 → p95 3 ×1.15 = 3.45 → 3.5 GB; durations 10,20,30 → p90 30 ×1.5 = 45 min seed_history(self.box, [hist_rec("r-90", peak=1.0, minutes=10.0), hist_rec("r-91", peak=2.0, minutes=20.0), hist_rec("r-92", title="build 89abcde (retry)", peak=3.0, minutes=30.0)]) code, out, err = self.box.wf("run", "--mem", "8G", "--for", "40m", "--title", "build 7654321", "--", "make") self.assertEqual((code, err), (0, "hint: history says ~3.5 GB / 45 min (3 runs)\n")) code, out, err = self.box.wf("run", "--mem", "7G", "--for", "90m", "--title", "build", "--", "make") self.assertEqual(err, "") code, out, err = self.box.wf("run", "--mem", "1G", "--for", "91m", "--title", "build", "--", "make") self.assertEqual(err, "hint: history says ~3.5 GB / 45 min (3 runs)\n") code, out, err = self.box.wf("note", "--mem", "8G", "--for", "10m", "build") self.assertEqual(err, ("hint: history says ~3.5 GB / 45 min (3 runs)\n")) code, out, err = self.box.wf("run", "--mem", "8G", "--for", "40m", "--title", "other", "--", "make") self.assertEqual(err, "") def test_hist_command(self): seed_history(self.box, [hist_rec("r-90", peak=1.0, minutes=10.0), hist_rec("r-91", peak=2.0, minutes=20.0), hist_rec("r-92", peak=3.0, minutes=30.0), hist_rec("r-93", project="zzz")]) code, out, err = self.box.wf("hist") self.assertEqual((code, err), (0, "")) self.assertEqual(out, "proj · build · n 3 · req 4.0 GB · peak 2.0/3.0 GB · est 40m · dur 20m/30m · suggest 3.5 GB 45m\n" "zzz · build · n 1 · req 4.0 GB · peak 1.0/1.0 GB · est 40m · dur 10m/10m · suggest - (< 3 runs)\n") self.assertEqual(self.box.wf("hist", "--project", "zzz")[1].count("\n"), 1)