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 --- tests/test_res_io.py | 815 +++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 815 insertions(+) create mode 100644 tests/test_res_io.py (limited to 'tests/test_res_io.py') diff --git a/tests/test_res_io.py b/tests/test_res_io.py new file mode 100644 index 0000000..d91046a --- /dev/null +++ b/tests/test_res_io.py @@ -0,0 +1,815 @@ +import datetime as dt +import io +import json +import os +import subprocess +import sys +import tempfile +import unittest +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 + + 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[: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_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")) + + +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) -- cgit