diff options
| author | godosa <godosa@godosa.eu> | 2026-10-07 07:27:17 +0200 |
|---|---|---|
| committer | godosa <godosa@godosa.eu> | 2026-10-07 07:27:17 +0200 |
| commit | 81d4e80fd5aabe4e80f58e960affa795cf7d34ec (patch) | |
| tree | e98eeac2af6af63aa4287bba1f6d4a3af26b5727 /tests/test_res.py | |
| download | workflow-81d4e80fd5aabe4e80f58e960affa795cf7d34ec.tar.gz workflow-81d4e80fd5aabe4e80f58e960affa795cf7d34ec.zip | |
workflow: initial public history
Diffstat (limited to 'tests/test_res.py')
| -rw-r--r-- | tests/test_res.py | 751 |
1 files changed, 751 insertions, 0 deletions
diff --git a/tests/test_res.py b/tests/test_res.py new file mode 100644 index 0000000..6c92549 --- /dev/null +++ b/tests/test_res.py @@ -0,0 +1,751 @@ +import datetime as dt +import json +import sys +import unittest +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) +from wflib import res as R + +UTC = dt.timezone.utc + + +def T(h, m=0, day=1): + return dt.datetime(2026, 10, day, h, m, tzinfo=UTC) + + +NOW = T(14, 10) +GIB = 1024 ** 3 + + +def running(id="r-4", project="proj-a", title="dotnet e2e", mem=10.0, cpus=2, est=40, started=T(14, 0)): + return R.Entry(id=id, project=project, owner=1, title=title, mem_gb=mem, cpus=cpus, est_min=est, + state="running", unit=f"wf-{id}.service", started=started, + expires=started + dt.timedelta(minutes=2 * est)) + + +def note(id="r-8", owner=77, mem=3.0, cpus=1, expires=T(14, 30)): + return R.Entry(id=id, project="proj-b", owner=owner, title="dotnet test", mem_gb=mem, cpus=cpus, + est_min=20, state="note", started=T(14, 0), expires=expires) + + +def facts(available=20.0, units=None, live=(), results=None, slice_gb=0.0, now=NOW): + return R.Facts(now=now, total_gb=30.0, available_gb=available, nproc=16, units=units or {}, + live=set(live), results=results or {}, slice_gb=slice_gb) + + +class Parse(unittest.TestCase): + def test_size(self): + self.assertEqual(R.parse_size("10G"), 10.0) + self.assertEqual(R.parse_size("512M"), 0.5) + self.assertEqual(R.parse_size("1.5g"), 1.5) + with self.assertRaisesRegex(R.ResError, "size '10'"): + R.parse_size("10") + + def test_duration(self): + self.assertEqual(R.parse_duration("40m"), 40) + self.assertEqual(R.parse_duration("4h"), 240) + self.assertEqual(R.parse_duration("1h30m"), 90) + self.assertEqual(R.parse_duration("90"), 90) + for bad in ("0m", "x", "", "h"): + with self.assertRaises(R.ResError): + R.parse_duration(bad) + + def test_meminfo(self): + text = "MemTotal: 31457280 kB\nMemFree: 1 kB\nMemAvailable: 12582912 kB\n" + self.assertEqual(R.meminfo(text), (30.0, 12.0)) + with self.assertRaises(R.ResError): + R.meminfo("MemFree: 1 kB\n") + + def test_show_units(self): + text = ("Id=agents.slice\nActiveState=active\nMemoryCurrent=0\n\n" + "Id=wf-r-1.service\nActiveState=inactive\nMemoryCurrent=[not set]\n") + self.assertEqual(R.show_units(text), { + "agents.slice": {"Id": "agents.slice", "ActiveState": "active", "MemoryCurrent": "0"}, + "wf-r-1.service": {"Id": "wf-r-1.service", "ActiveState": "inactive", "MemoryCurrent": "[not set]"}}) + + def test_gb_or_none(self): + self.assertEqual(R.gb_or_none("1073741824"), 1.0) + self.assertIsNone(R.gb_or_none("[not set]")) + self.assertIsNone(R.gb_or_none("infinity")) + + def test_formats(self): + self.assertEqual(R.fmt_gb(2.25), "2.2 GB") + self.assertEqual(R.fmt_bytes(1288490188), "1.2 GB") + self.assertEqual(R.fmt_bytes(300 * 1024 ** 2), "300 MB") + self.assertEqual(R.fmt_dur(240), "4h") + self.assertEqual(R.fmt_dur(90), "1h30m") + self.assertEqual(R.fmt_dur(40), "40m") + self.assertEqual(R.hhmm(T(9, 5)), "09:05") + + +class ConfigTest(unittest.TestCase): + def test_defaults_and_override(self): + cfg = R.load_config("user_reserve_gb = 8\n") + self.assertEqual((cfg.user_reserve_gb, cfg.user_reserve_cpus, cfg.game_reserve_gb, cfg.game_reserve_cpus, + cfg.game_hours, cfg.small_headroom_gb, cfg.scratch_hours), (8.0, 4, 12.0, 8, 4.0, 2.0, 2.0)) + + def test_unknown_key(self): + self.assertEqual(R.load_config("session_mem_gb = 12").session_mem_gb, 12.0) + self.assertEqual(R.Config().session_mem_gb, 6.0) + with self.assertRaisesRegex(R.ResError, "unknown key\\(s\\) foo"): + R.load_config("foo = 1\n") + + def test_negative(self): + with self.assertRaisesRegex(R.ResError, "game_hours must be a number ≥ 0"): + R.load_config("game_hours = -1\n") + + +class LedgerJson(unittest.TestCase): + def test_roundtrip(self): + led = R.Ledger(next=9, game_until=T(18), entries=[running(), note()]) + text = R.dumps(led) + self.assertIn('"started": "2026-10-01T14:00:00+00:00"', text) + self.assertEqual(R.loads(text), led) + + def test_empty(self): + self.assertEqual(R.loads(""), R.Ledger()) + + def test_corrupt(self): + for bad in ("{", '{"entries": [{"id": "r-1"}]}', "[]"): + with self.assertRaisesRegex(R.ResError, "corrupt ledger"): + R.loads(bad) + + def test_unknown_entry_keys_ignored(self): + d = json.loads(R.dumps(R.Ledger(next=9, entries=[running()]))) + d["entries"][0]["from_a_newer_version"] = 1 + self.assertEqual(R.loads(json.dumps(d)), R.Ledger(next=9, entries=[running()])) + + def test_unknown_keys_survive_rewrite(self): + d = json.loads(R.dumps(R.Ledger(next=9, entries=[running()]))) + d["entries"][0]["field_x"] = {"a": 1} + d["top_x"] = 5 + out = json.loads(R.dumps(R.loads(json.dumps(d)))) + self.assertEqual(out["entries"][0]["field_x"], {"a": 1}) + self.assertEqual(out["top_x"], 5) + + def test_new_id_and_get(self): + led = R.Ledger(next=3) + self.assertEqual(led.new_id(), "r-3") + self.assertEqual(led.next, 4) + with self.assertRaisesRegex(R.ResError, "no entry 'r-9'"): + led.get("r-9") + + +class Prune(unittest.TestCase): + def test_prune(self): + r1, r2 = running("r-1"), running("r-2") + n3, n5 = note("r-3", owner=77), note("r-5", owner=78, expires=T(14, 5)) + d6 = R.Entry(id="r-6", project="p", owner=1, title="old", mem_gb=1, cpus=1, est_min=1, state="done", + ended=NOW - dt.timedelta(hours=25)) + d7 = R.Entry(id="r-7", project="p", owner=1, title="new", mem_gb=1, cpus=1, est_min=1, state="done", + ended=NOW - dt.timedelta(hours=23)) + led = R.Ledger(game_until=T(14, 9), entries=[r1, r2, n3, n5, d6, d7]) + f = facts(units={"wf-r-1.service": R.Unit(False)}, live={78}, results={"r-1": (0, 9.4, "")}) + lines = R.prune(led, f) + self.assertEqual(lines, ["r-1 exited rc=0", "r-2 exited rc=?", "r-3 freed (owner gone)", + "r-5 freed (expired)", "game off (expired)"]) + self.assertEqual([e.id for e in led.entries], ["r-1", "r-2", "r-3", "r-5", "r-7"]) + self.assertEqual((r1.state, r1.rc, r1.peak_gb, r1.why, r1.ended), ("done", 0, 9.4, "exited", NOW)) + self.assertEqual((n3.why, n5.why), ("owner gone", "expired")) + self.assertIsNone(led.game_until) + + def test_prune_missing_results(self): + r = running("r-1") + R.prune(R.Ledger(entries=[r]), facts(units={"wf-r-1.service": R.Unit(False)})) + self.assertEqual((r.state, r.rc, r.peak_gb), ("done", None, None)) + + def test_overdue_running_is_kept(self): + r = running(started=T(10)) + R.prune(R.Ledger(entries=[r]), facts(units={"wf-r-4.service": R.Unit(True, 4.0)})) + self.assertEqual(r.state, "running") + + +class Capacity(unittest.TestCase): + def setUp(self): + self.cfg = R.Config() + + def test_budget(self): + led = R.Ledger(entries=[running(), note()]) + f = facts(units={"wf-r-4.service": R.Unit(True, 4.0)}, live={77}) + # 20 − 6 reserve − 2 headroom − (10−4) unused − 3 note = 3; 16 − 4 − (2+1) = 9 + self.assertEqual(R.budget(self.cfg, led, f), (3.0, 9)) + + def test_budget_gaming(self): + led = R.Ledger(game_until=T(15), entries=[running(), note()]) + f = facts(units={"wf-r-4.service": R.Unit(True, 4.0)}, live={77}) + # 20 − 12 − 2 − 6 − 3 = −3; 16 − 8 − 3 = 5 + self.assertEqual(R.budget(self.cfg, led, f), (-3.0, 5)) + + def test_held_over_reservation_is_zero(self): + self.assertEqual(R.held_gb(running(), facts(units={"wf-r-4.service": R.Unit(True, 12.0)})), 0.0) + self.assertEqual(R.frees_gb(running(), facts(units={"wf-r-4.service": R.Unit(True, 12.0)})), 12.0) + + def test_fits(self): + self.assertTrue(R.fits(3.0, 9, (3.0, 9))) + self.assertFalse(R.fits(3.1, 1, (3.0, 9))) + self.assertFalse(R.fits(1.0, 10, (3.0, 9))) + + def test_never_fits(self): + f = facts() + self.assertEqual(R.never_fits(self.cfg, f, 25.0, 1), "25.0 GB can never fit (max 22.0 GB for agents)") + self.assertEqual(R.never_fits(self.cfg, f, 1.0, 13), "13 cpus can never fit (max 12 for agents)") + self.assertIsNone(R.never_fits(self.cfg, f, 22.0, 12)) + + def test_queued_total(self): + q = running("r-9") + q.state = "queued" + self.assertEqual(R.queued_total(R.Ledger(entries=[q, running()])), (10.0, 2)) + + +class BusyLine(unittest.TestCase): + def setUp(self): + self.cfg = R.Config() + self.led = R.Ledger(entries=[running()]) + self.f = facts(units={"wf-r-4.service": R.Unit(True, 4.0)}) # budget 6.0 GB, 10 cpus + + def test_memory(self): + self.assertEqual(R.busy_line(self.cfg, self.led, self.f, 8.0, 1), + 'busy: 10.0 GB held by proj-a "dotnet e2e" (r-4) until ~14:40; 6.0 GB free for agents; ' + 'retry after ~14:40 or work on something else') + + def test_cpus(self): + self.assertEqual(R.busy_line(self.cfg, self.led, self.f, 1.0, 12), + 'busy: 2 cpus held by proj-a "dotnet e2e" (r-4) until ~14:40; 10 cpus free for agents; ' + 'retry after ~14:40 or work on something else') + + def test_no_holder(self): + self.assertEqual(R.busy_line(self.cfg, R.Ledger(), facts(available=7.0), 1.0, 1), + "busy: only 0.0 GB free for agents and no agent job holds any; other programs use the rest; " + "retry later or work on something else") + + def test_busy_overdue_holder(self): + led = R.Ledger(entries=[running(started=T(10))]) + self.assertEqual(R.busy_line(self.cfg, led, self.f, 8.0, 1), + 'busy: 10.0 GB held by proj-a "dotnet e2e" (r-4) until overdue; 6.0 GB free for agents; ' + 'retry later or work on something else') + + def test_two_holders_earliest_first(self): + a = running("r-1", project="a", title="A", mem=4.0, cpus=1, est=60) # ends 15:00 + b = running("r-2", project="b", title="B", mem=4.0, cpus=1, est=30) # ends 14:30 + f = facts(available=14.0, units={}) # 14 − 6 − 2 − 4 − 4 = −2.0 budget + self.assertEqual(R.busy_line(self.cfg, R.Ledger(entries=[a, b]), f, 5.0, 1), + 'busy: 4.0 GB held by b "B" (r-2) until ~14:30, 4.0 GB held by a "A" (r-1) until ~15:00; ' + '0.0 GB free for agents; retry after ~15:00 or work on something else') + + +class Lines(unittest.TestCase): + def test_run_argv(self): + e = running("r-3", mem=10.0) + e.cwd, e.log, e.cmd = "/projects/x", "/s/logs/r-3.log", ["make", "it big", "--fast"] + self.assertEqual(R.run_argv(e, "/s/logs"), [ + "systemd-run", "--user", "--quiet", "--collect", "--slice=agents-jobs.slice", "--unit=wf-r-3.service", + "--working-directory=/projects/x", "--setenv=WF_RES_ID=r-3", + "-p", "MemoryMax=10737418240", "-p", "MemoryHigh=9663676416", "-p", "MemorySwapMax=0", "-p", "Nice=10", + "-p", "StandardOutput=append:/s/logs/r-3.log", "-p", "StandardError=append:/s/logs/r-3.log", + "/bin/sh", "-c", + '"$@"; rc=$?; cat /sys/fs/cgroup$(cut -d: -f3 /proc/self/cgroup)/memory.peak > /s/logs/r-3.peak ' + '2>/dev/null; echo $rc > /s/logs/r-3.rc', + "sh", "make", "it big", "--fast"]) + + def test_run_argv_env(self): + e = running("r-3", mem=1.0) + e.cwd, e.log, e.cmd, e.env = "/p", "/s/r-3.log", ["x"], {"B": "2", "A": "1 2"} + argv = R.run_argv(e, "/s") + self.assertEqual(argv[argv.index("--working-directory=/p") + 1:][:2], ["--setenv=A=1 2", "--setenv=B=2"]) + + def test_env_diff(self): + self.assertEqual(R.env_diff({"A": "1", "B": "2", "C": "3", "PWD": "/", "OLDPWD": "/", "SHLVL": "1", "_": "/x", + "MY_TOKEN": "t", "API_SECRET": "s", "PGPASSWORD": "p", "SSH_AUTH_SOCK": "/s"}, + {"A": "1", "B": "9", "SSH_AUTH_SOCK": "/s"}), + {"B": "2", "C": "3"}) + + def test_manager_env_parse(self): + self.assertEqual(R.parse_show_environment("A=1\nB=x=y\n\nbad\n"), {"A": "1", "B": "x=y"}) + + def test_entry_lines(self): + f = facts(units={"wf-r-4.service": R.Unit(True, 4.0)}, live={77}) + q = running("r-9", project="p", title="Q", mem=1.0, cpus=1) + q.state = "queued" + d = R.Entry(id="r-2", project="p", owner=1, title="D", mem_gb=1, cpus=1, est_min=1, state="done", + ended=T(13, 50), rc=None, peak_gb=None, why="exited") + led = R.Ledger(entries=[running(), note(), q, d]) + self.assertEqual([R.entry_line(e, led, f) for e in led.entries], [ + 'r-4 proj-a "dotnet e2e" running 10.0 GB used 4.0 GB 2 cpu since 14:00 ETA ~14:40', + 'r-8 proj-b "dotnet test" note 3.0 GB 1 cpu until ~14:30', + 'r-9 p "Q" queued #1 1.0 GB 1 cpu', + 'r-2 p "D" done (exited) rc=? peak ? at 13:50']) + + def test_status_lines(self): + f = facts(units={"wf-r-4.service": R.Unit(True, 4.0)}, slice_gb=5.5) + self.assertEqual(R.status_lines(R.Config(), R.Ledger(entries=[running()]), f), [ + 'r-4 proj-a "dotnet e2e" running 10.0 GB used 4.0 GB 2 cpu since 14:00 ETA ~14:40', + "agents may use 6.0 GB, 10 cpus now; reserve 6.0 GB/4 cpus", + "unreserved agent memory 1.5 GB"]) + f.slice_cache_gb = 0.6 + self.assertEqual(R.status_lines(R.Config(), R.Ledger(entries=[running()]), f)[-1], + "unreserved agent memory 1.5 GB (0.6 GB of it file cache, reclaimable)") + f.slice_cache_gb = 3.0 # cache partly inside the jobs' 4.0 GB → capped at the unreserved figure + self.assertEqual(R.status_lines(R.Config(), R.Ledger(entries=[running()]), f)[-1], + "unreserved agent memory 1.5 GB (1.5 GB of it file cache, reclaimable)") + + def test_cache_gb(self): + # agents.slice memory.stat excerpt from this machine: file 5602209792, shmem (tmpfs, not reclaimable) 3861041152 + self.assertEqual(round(R.cache_gb("anon 11335516160\nfile 5602209792\nkernel 1\nshmem 3861041152\n"), 3), 1.622) + self.assertEqual(R.cache_gb(""), 0.0) + + def test_status_lines_empty_overdue(self): + self.assertEqual(R.status_lines(R.Config(), R.Ledger(), facts(available=7.0)), [ + "no reservations", "agents may use 0.0 GB, 12 cpus now; reserve 6.0 GB/4 cpus"]) + r = running(started=T(10)) + self.assertIn("ETA overdue (still running, not killed)", R.entry_line(r, R.Ledger(entries=[r]), facts())) + + def test_status_json(self): + import json + d = json.loads(R.status_json(R.Config(), R.Ledger(entries=[running()]), + facts(units={"wf-r-4.service": R.Unit(True, 4.0)}))) + self.assertEqual((d["budget_gb"], d["budget_cpus"], d["gaming"], d["entries"][0]["id"]), + (6.0, 10, False, "r-4")) + + def test_done_line(self): + e = running("r-3", started=T(14, 0)) + e.state, e.ended, e.rc, e.peak_gb = "done", T(14, 37), 0, 9.4 + self.assertEqual(R.done_line(e), "r-3 done rc=0 peak 9.4 GB in 37 min") + e.rc, e.peak_gb = None, None + self.assertEqual(R.done_line(e), "r-3 done rc=? peak ? in 37 min") + e.why = "killed: timeout" + self.assertEqual(R.done_line(e), "r-3 done rc=? peak ? in 37 min (killed: timeout)") + + def test_prune_killed(self): + r = running("r-1") + f = facts(units={"wf-r-1.service": R.Unit(False)}, results={"r-1": (None, 3.9, "killed: timeout")}) + self.assertEqual(R.prune(R.Ledger(entries=[r]), f), ["r-1 killed: timeout"]) + self.assertEqual((r.state, r.rc, r.peak_gb, r.why), ("done", None, 3.9, "killed: timeout")) + + +# journalctl --user -u wf-r-12.service -o cat, verbatim from a real systemd-oomd kill (systemd 258) +OOMD = """Started wf-r-12.service - [systemd-run] /bin/sh -c "x" sh dotnet test. +wf-r-12.service: systemd-oomd killed some process(es) in this unit. +wf-r-12.service: Main process exited, code=killed, status=9/KILL +wf-r-12.service: Failed with result 'oom-kill'. +wf-r-12.service: Consumed 11min 2.837s CPU time, 3.9G memory peak. +""" + + +class JournalReason(unittest.TestCase): + def test_oomd(self): + self.assertEqual(R.journal_reason(OOMD, 4.0), + ("killed: oom-kill by systemd-oomd, limit 4.0 GB; raise --mem", 3.9)) + + def test_kernel_oom(self): + text = ("u: A process of this unit has been killed by the OOM killer.\n" + "u: Main process exited, code=killed, status=9/KILL\n" + "u: Failed with result 'oom-kill'.\nu: Consumed 2s CPU time, 512M memory peak.\n") + self.assertEqual(R.journal_reason(text, 0.5), + ("killed: oom-kill at MemoryMax, limit 0.5 GB; raise --mem", 0.5)) + + def test_signal(self): + text = "u: Main process exited, code=killed, status=15/TERM\nu: Failed with result 'signal'.\n" + self.assertEqual(R.journal_reason(text, 1.0), ("killed: signal 15/TERM", None)) + + def test_timeout_and_exit(self): + self.assertEqual(R.journal_reason("u: Failed with result 'timeout'.\n", 1.0), ("killed: timeout", None)) + text = "u: Main process exited, code=exited, status=3/NOTIMPLEMENTED\nu: Failed with result 'exit-code'.\n" + self.assertEqual(R.journal_reason(text, 1.0), ("failed: exit 3/NOTIMPLEMENTED", None)) + + def test_nothing(self): + self.assertEqual(R.journal_reason("", 1.0), ("", None)) + self.assertEqual(R.journal_reason("u: Deactivated successfully.\nu: Consumed 1s CPU time, 1.5G memory peak.\n", 1.0), + ("", 1.5)) + +def queued(id, mem, cpus=1, at=T(13, 0)): + e = R.Entry(id=id, project="p", owner=1, title=id, mem_gb=mem, cpus=cpus, est_min=10, state="queued", + cmd=["x"], queued=at) + return e + + +class Queue(unittest.TestCase): + def test_fifo_blocked_head(self): + # budget 20 − 6 − 2 = 12 GB + led = R.Ledger(entries=[queued("r-1", 5, at=T(13, 0)), queued("r-2", 9, at=T(13, 1)), + queued("r-3", 1, at=T(13, 2))]) + self.assertEqual([e.id for e in R.to_start(R.Config(), led, facts())], ["r-1"]) + + def test_all_fit(self): + led = R.Ledger(entries=[queued("r-2", 4, at=T(13, 1)), queued("r-1", 4, at=T(13, 0))]) + self.assertEqual([e.id for e in R.to_start(R.Config(), led, facts())], ["r-1", "r-2"]) + + def test_estimate(self): + # running r-4 holds 10 (used 4), ends 14:40; budget 6; queue r-1 (5) ahead of r-2 (4): need 9 − 6 = 3 + led = R.Ledger(entries=[running(), queued("r-1", 5), queued("r-2", 4, at=T(13, 5))]) + f = facts(units={"wf-r-4.service": R.Unit(True, 4.0)}) + self.assertEqual(R.queue_estimate(R.Config(), led, f, led.get("r-2")), T(14, 40)) + + def test_estimate_unknown(self): + led = R.Ledger(entries=[queued("r-1", 20)]) + self.assertIsNone(R.queue_estimate(R.Config(), led, facts(), led.get("r-1"))) + +class Game(unittest.TestCase): + def setUp(self): + self.cfg = R.Config() + + def test_slice_props(self): + self.assertEqual(R.slice_props(self.cfg, R.Ledger(), facts()), + {"CPUWeight": "20", "IOWeight": "20", "MemoryHigh": str(24 * GIB)}) + on = R.Ledger(game_until=T(18)) + self.assertEqual(R.slice_props(self.cfg, on, facts(slice_gb=10.0))["MemoryHigh"], str(18 * GIB)) + self.assertEqual(R.slice_props(self.cfg, on, facts(slice_gb=20.0)), + {"CPUWeight": "5", "IOWeight": "5", "MemoryHigh": str(20 * GIB)}) + + def test_shortfall(self): + a = running("r-4", project="proj-a", title="dotnet e2e", mem=2.5, est=55, started=T(14, 10)) # ~15:05 + b = running("r-7", project="proj-b", title="r5 rebuild", mem=9.0, est=80, started=T(14, 10)) # ~15:30 + f = facts(available=11.5, units={"wf-r-4.service": R.Unit(True, 2.5), "wf-r-7.service": R.Unit(True, 9.0)}) + led = R.Ledger(game_until=T(18), entries=[a, b]) + # free for you 11.5 − 0 unused = 11.5 → short 0.5: r-4 alone frees 2.5 + self.assertEqual(R.shortfall_lines(self.cfg, led, f), [ + 'short 0.5 GB of 12.0 GB: r-4 proj-a "dotnet e2e" 2.5 GB ~15:05', + "full reserve free ~15:05 (est.); free now: wf res release r-4 --stop"]) + f.available_gb = 8.9 # short 3.1: needs both + self.assertEqual(R.shortfall_lines(self.cfg, led, f), [ + 'short 3.1 GB of 12.0 GB: r-4 proj-a "dotnet e2e" 2.5 GB ~15:05, r-7 proj-b "r5 rebuild" 9.0 GB ~15:30', + "full reserve free ~15:30 (est.); free now: wf res release r-7 --stop"]) + f.available_gb = 13.0 + self.assertEqual(R.shortfall_lines(self.cfg, led, f), []) + + def test_shortfall_not_agents(self): + self.assertEqual(R.shortfall_lines(self.cfg, R.Ledger(game_until=T(18)), facts(available=5.0)), [ + "short 7.0 GB of 12.0 GB: other programs, not agent jobs", + "agent jobs alone cannot free it; close other programs"]) + + def test_game_on_lines(self): + led = R.Ledger(game_until=T(18, 10)) + self.assertEqual(R.game_on_lines(self.cfg, led, facts(available=20.0), 240), + ["game on until 18:10 (4h); CPU/IO now yours"]) + + def test_status_shows_game(self): + lines = R.status_lines(self.cfg, R.Ledger(game_until=T(18, 10)), facts(available=5.0)) + self.assertEqual(lines[1:], ["agents may use 0.0 GB, 8 cpus now; reserve 12.0 GB/8 cpus", + "game on until 18:10 (4h left)", + "short 7.0 GB of 12.0 GB: other programs, not agent jobs", + "agent jobs alone cannot free it; close other programs"]) + +class Clean(unittest.TestCase): + def test_scratch_victims(self): + now = 10_000_000.0 + h = 3600 + tree = { + "/s/old-file.png": (now - 3 * h, {}), + "/s/-projects-a": (now - 3 * h, {"/s/-projects-a/s1": now - 3 * h}), + "/s/-projects-b": (now - 60, {"/s/-projects-b/s1": now - 5 * h, "/s/-projects-b/s2": now - 60}), + "/s/new.txt": (now - 60, {}), + } + self.assertEqual(R.scratch_victims(tree, now, 2.0), + ["/s/-projects-a", "/s/-projects-b/s1", "/s/old-file.png"]) + # live sessions keep their dir however idle; a stale top holding one is pruned per child + self.assertEqual(R.scratch_victims(tree, now, 2.0, keep={"s1"}), ["/s/old-file.png"]) + tree["/s/-projects-a"] = (now - 3 * h, {"/s/-projects-a/s1": now - 3 * h, "/s/-projects-a/s9": now - 3 * h}) + self.assertEqual(R.scratch_victims(tree, now, 2.0, keep={"s1"}), ["/s/-projects-a/s9", "/s/old-file.png"]) + + def test_live_session_id(self): + text = '{"pid": 817151, "sessionId": "4ad140ab-69ca", "procStart": "31379851", "kind": "interactive"}' + self.assertEqual(R.live_session_id(text, "31379851"), "4ad140ab-69ca") + self.assertIsNone(R.live_session_id(text, "999")) # pid reused by another process + self.assertEqual(R.live_session_id('{"sessionId": "x"}', "5"), "x") # no procStart recorded → trust pid + self.assertIsNone(R.live_session_id("{oops", "5")) + self.assertIsNone(R.live_session_id('{"pid": 1}', "5")) + + def test_cleanup_rule(self): + self.assertEqual(R.cleanup_rule("out/prof"), ("out/prof", None)) + self.assertEqual(R.cleanup_rule("out/history-logs/*.log:30d"), ("out/history-logs/*.log", 30)) + for bad in ("/etc", "../x", "out/../../x"): + with self.assertRaisesRegex(R.ResError, "cleanup pattern"): + R.cleanup_rule(bad) + + def test_proc_name(self): + # background sessions run the versioned binary: comm is the version, argv0 the path or `claude` + self.assertEqual(R.proc_name("2.1.283", "claude"), "claude") + self.assertEqual(R.proc_name("2.1.283", "claude bg-pty-host --bg-pty-host /tmp/x.sock"), "claude") # rewritten title + self.assertEqual(R.proc_name("2.1.283", "/h/u/.local/share/claude/versions/2.1.283"), "claude") + self.assertEqual(R.proc_name("claude", "/h/u/.local/bin/claude"), "claude") + self.assertEqual(R.proc_name("2.1.284", "ugrep"), "2.1.284") # tool re-exec of the binary: not a session + self.assertEqual(R.proc_name("bash", "/bin/bash"), "bash") + self.assertEqual(R.proc_name("x", ""), "x") # kernel thread / unreadable cmdline + + def test_adopt_groups(self): + procs = [ + (100, 1, "bash", "/u/app.slice/tab1.scope"), + (101, 100, "claude", "/u/app.slice/tab1.scope"), # root, outside + (102, 101, "bash", "/u/app.slice/tab1.scope"), # its shell + (103, 102, "make", "/u/app.slice/tab1.scope"), + (104, 101, "claude", "/u/app.slice/tab1.scope"), # child claude: not a root, but a descendant + (105, 101, "job", "/u/agents.slice/agents-jobs.slice/wf-r-1.service"), # already inside: skipped + (200, 1, "claude", "/u/agents.slice/run-9.scope"), # root already inside + (300, 1, "vim", "/u/app.slice/tab2.scope"), + ] + self.assertEqual(R.adopt_groups(procs), {101: [101, 102, 103, 104]}) + + def test_warning_lines(self): + self.assertEqual(R.warning_lines(30.0, 7.6, 2), [ + "warning: /tmp (RAM) holds 7.6 GB; wf res clean", + "warning: 2 claude sessions outside agents.slice (wf res adopt)"]) + self.assertEqual(R.warning_lines(30.0, 7.4, 0), []) + +class Units(unittest.TestCase): + def test_adopt_argv(self): + self.assertEqual(R.adopt_argv(101, [101, 102]), [ + "busctl", "--user", "call", "org.freedesktop.systemd1", "/org/freedesktop/systemd1", + "org.freedesktop.systemd1.Manager", "StartTransientUnit", "ssa(sv)a(sa(sv))", + "wf-claude-101.scope", "fail", "2", "PIDs", "au", "2", "101", "102", "Slice", "s", "agents.slice", "0"]) + + def test_texts(self): + self.assertEqual(R.service_text("/usr/bin/python3", "/projects/public/workflow/wf.py"), + "[Unit]\nDescription=wf res tick\n\n[Service]\nType=oneshot\n" + "ExecStart=/usr/bin/python3 /projects/public/workflow/wf.py res tick\n") + self.assertEqual(R.timer_text(), + "[Unit]\nDescription=wf res tick every minute\n\n[Timer]\nOnBootSec=1min\n" + "OnUnitActiveSec=1min\n\n[Install]\nWantedBy=timers.target\n") + self.assertEqual(R.shell_init_line(), + "alias claude='systemd-run --user --scope --quiet --slice=agents.slice claude'") + + +if __name__ == "__main__": + unittest.main() + + +class Throttle(unittest.TestCase): + def test_psi(self): + text = "some avg10=41.50 avg60=33.20 avg300=12.00 total=99\nfull avg10=30.00 avg60=25.00 avg300=9.00 total=88\n" + self.assertEqual(R.psi_some_avg60(text), 33.2) + self.assertIsNone(R.psi_some_avg60("")) + + def test_lines(self): + hot = running("r-1", mem=4.0) # 3.5 of 4.0 GB (≥ 0.85×), stall 33% → warn + cache = running("r-2", mem=4.0) # full of page cache, no stall → quiet + small = running("r-3", mem=10.0) # stalled but far below its limit (machine pressure) → quiet + f = facts(units={"wf-r-1.service": R.Unit(True, 3.5, 33.2), "wf-r-2.service": R.Unit(True, 3.9, 0.0), + "wf-r-3.service": R.Unit(True, 2.0, 50.0)}) + self.assertEqual(R.throttle_lines(R.Ledger(entries=[hot, cache, small]), f), [ + "r-1 throttled at its memory limit (3.5 of 4.0 GB, stalled 33% of the last minute): " + "likely too small; `wf res release r-1 --stop` and re-run with a bigger --mem"]) + + +class Owner(unittest.TestCase): + REC = [{"lane": "fast", "pid": 11, "socket": "/s/a"}, {"lane": "slow", "pid": 22, "socket": "/s/b"}] + + def test_from_env_branch_and_session_record(self): + env = {"CLAUDE_PID": "22", "CLAUDE_CODE_MESSAGING_SOCKET": "/s/b", "WF_RES_ID": "r-7"} + self.assertEqual(R.owner_by(env, "slow/t-res-owner", self.REC), + {"name": "slow session", "task": "t-res-owner", "batch": "r-7", "address": "uds:/s/b"}) + + def test_explicit_name_wins_then_env(self): + env = {"CLAUDE_PID": "22", "WF_SESSION_NAME": "pilot", "WF_TASK": "t-x"} + self.assertEqual(R.owner_by(env, "slow/t-y", self.REC, "worker-3"), {"name": "worker-3", "task": "t-x"}) + self.assertEqual(R.owner_by(env, "", self.REC), {"name": "pilot", "task": "t-x"}) + + def test_nothing_known(self): + self.assertEqual(R.owner_by({"CLAUDE_PID": "99"}, "master", self.REC), {}) + self.assertEqual(R.owner_by({}, "feature-x", [{"pid": 1}]), {}) + + def test_status_and_throttle_name_owner(self): + e = running("r-1", mem=4.0) + e.by = {"name": "slow session", "task": "t-a", "batch": "r-7", "address": "uds:/s/b"} + f = facts(units={"wf-r-1.service": R.Unit(True, 3.5, 33.2)}) + led = R.Ledger(entries=[e]) + self.assertTrue(R.status_lines(R.Config(), led, f)[0].endswith( + " [by slow session t-a batch r-7, message uds:/s/b]")) + self.assertTrue(R.throttle_lines(led, f)[0].endswith( + "bigger --mem; started by slow session t-a batch r-7, message uds:/s/b")) + + def test_ledger_roundtrip_and_old_entries(self): + e = running("r-1") + e.by = {"task": "t-a"} + self.assertEqual(R.loads(R.dumps(R.Ledger(entries=[e]))).get("r-1").by, {"task": "t-a"}) + d = json.loads(R.dumps(R.Ledger(entries=[running("r-2")]))) + for k in ("by", "lock"): # entries written before these fields + del d["entries"][0][k] + old = R.loads(json.dumps(d)).get("r-2") + self.assertEqual((old.by, old.lock), ({}, "")) + + +class BareUnits(unittest.TestCase): + def test_bare_units(self): + blocks = R.show_units( + "Id=wh-t27.service\nTransient=yes\nSlice=app.slice\nMemoryCurrent=2147483648\n\n" + "Id=run-u42.service\nTransient=yes\nSlice=app.slice\nMemoryCurrent=[not set]\n\n" + "Id=app-org.kde.konsole@65c0.service\nTransient=yes\nSlice=app.slice\nMemoryCurrent=1\n\n" + "Id=dbus-:1.1-org.kde.kwalletd6@0.service\nTransient=yes\nSlice=app.slice\nMemoryCurrent=1\n\n" + "Id=wf-r-3.service\nTransient=yes\nSlice=agents-jobs.slice\nMemoryCurrent=1\n\n" + "Id=pipewire.service\nTransient=no\nSlice=session.slice\nMemoryCurrent=1\n") + self.assertEqual(R.bare_units(blocks), [("run-u42.service", None), ("wh-t27.service", 2.0)]) + + def test_warning(self): + self.assertEqual(R.warning_lines(30.0, 0.0, 0, [("wh-t27.service", 2.0), ("run-u42.service", None)]), [ + "warning: 2 jobs outside wf res (bare systemd-run): wh-t27.service 2.0 GB, run-u42.service ? GB; " + "start jobs with wf res run, also from project scripts"]) + + +class Hook(unittest.TestCase): + BASH = {"tool_name": "Bash", "tool_input": {"command": "dotnet run --project Desktop", "timeout": 5000}} + + def test_gaming_hides_display(self): + out = R.hook_output(self.BASH, T(16, 0), NOW) + self.assertEqual(out, {"hookSpecificOutput": { + "hookEventName": "PreToolUse", + "updatedInput": {"command": "unset DISPLAY WAYLAND_DISPLAY; dotnet run --project Desktop", "timeout": 5000}, + "additionalContext": "wf res game mode until 16:00: this command has no display, so GUI windows fail. " + "Run headless/offscreen, or do other work until game off."}}) + + def test_quiet_when_not_gaming_or_not_bash(self): + self.assertIsNone(R.hook_output(self.BASH, None, NOW)) + self.assertIsNone(R.hook_output(self.BASH, T(14, 0), NOW)) # expired + self.assertIsNone(R.hook_output({"tool_name": "Read", "tool_input": {"file_path": "x"}}, T(16, 0), NOW)) + self.assertIsNone(R.hook_output({"tool_name": "Bash", "tool_input": {}}, T(16, 0), NOW)) + + +class Litter(unittest.TestCase): + def test_victims(self): + now, h = 10_000_000.0, 3600 + entries = [("/t/aB3xYz", "emptydir", now - 3 * h), # old empty dir → go + ("/t/MSBuildTemp1000", "emptydir", now - 5 * h), # → go + ("/t/fresh", "emptydir", now - 60), # too new + ("/t/full", "dir", now - 9 * h), # has content → never + ("/t/clr-debug-pipe-111-222-in", "other", now - 60), # pid 111 dead → go + ("/t/dotnet-diagnostic-333-444-socket", "other", now), # pid 333 alive + ("/t/notes.txt", "file", now - 9 * h)] # file → never + alive = {333}.__contains__ + self.assertEqual(R.litter_victims(entries, now, 2.0, alive), + ["/t/MSBuildTemp1000", "/t/aB3xYz", "/t/clr-debug-pipe-111-222-in"]) +class SessionCap(unittest.TestCase): + BLOCKS = { + "wf-claude-1.scope": {"Slice": "agents.slice", "MemoryHigh": "infinity", "MemoryCurrent": str(GIB)}, + "wf-claude-2.scope": {"Slice": "agents.slice", "MemoryHigh": str(10 * GIB), "MemoryCurrent": str(GIB)}, + "run-u7.scope": {"Slice": "agents.slice", "MemoryHigh": str(4 * GIB), "MemoryCurrent": str(GIB)}, + "podman-pause.scope": {"Slice": "user.slice", "MemoryHigh": "infinity", "MemoryCurrent": str(GIB)}, + } + + def test_caps_to_set(self): + self.assertEqual(R.session_caps(self.BLOCKS, 10.0), + [("run-u7.scope", str(10 * GIB)), ("wf-claude-1.scope", str(10 * GIB))]) + self.assertEqual(R.session_caps(self.BLOCKS, 0.0), # 0 = off → lift caps + [("run-u7.scope", "infinity"), ("wf-claude-2.scope", "infinity")]) + + def test_warn_at_cap(self): + blocks = {"wf-claude-5.scope": {"Slice": "agents.slice", "MemoryCurrent": str(int(9.6 * GIB))}, + "wf-claude-6.scope": {"Slice": "agents.slice", "MemoryCurrent": str(int(9.8 * GIB))}, + "wf-claude-7.scope": {"Slice": "agents.slice", "MemoryCurrent": str(int(2 * GIB))}} + stall = {"wf-claude-5.scope": 45.0, "wf-claude-6.scope": 1.0, "wf-claude-7.scope": 90.0} + self.assertEqual(R.session_cap_lines(blocks, stall, 10.0), [ + "warning: claude session 5 at its memory cap (9.6 of 10.0 GB, stalled 45% of the last minute): " + "run big work with wf res run"]) + self.assertEqual(R.session_cap_lines(blocks, stall, 0.0), []) + + +def run(title="gate 30735e0", peak=1.0, minutes=10.0, project="proj", mem=12.0, est=60, rc=0, id="r-1"): + return {"id": id, "project": project, "title": title, "mem_gb": mem, "est_min": est, "rc": rc, + "peak_gb": peak, "min": minutes} + + +TEN = [run(id=f"r-{i}", title=f"gate {i:07x}a", peak=float(i), minutes=6.0 * i) for i in range(1, 11)] + + +class History(unittest.TestCase): + def test_title_kind(self): + for title, kind in (("gate 30735e0", "gate"), ("gate 2c91cac7 (fix optimizer test)", "gate"), + ("mem-profile 16x5M (t-optimizer-x)", "mem-profile 16x5M"), ("wf-batch", "wf-batch"), + ("gate 97e6044b direct", "gate 97e6044b direct"), ("build deadbeef", "build deadbeef"), + ("gate a3cbde1 0cf9a79", "gate"), ("abc1234", "abc1234"), (" x ", "x")): + self.assertEqual(R.title_kind(title), kind, title) + + def test_pct_nearest_rank(self): + xs = [5.0, 1.0, 3.0, 2.0, 4.0] + self.assertEqual([R.pct(xs, p) for p in (0.5, 0.9, 0.95, 0.2)], [3.0, 5.0, 5.0, 1.0]) + + def test_history_record(self): + e = running(project="p", title="gate x", mem=10.0, est=40) + e.state, e.ended, e.rc, e.peak_gb = "done", T(14, 13, ), 0, 9.1 + e.ended = e.started + dt.timedelta(minutes=13, seconds=30) + self.assertEqual(R.history_record(e), {"id": "r-4", "project": "p", "title": "gate x", "mem_gb": 10.0, + "est_min": 40, "rc": 0, "peak_gb": 9.1, "min": 13.5}) + for change in ({"rc": 1}, {"peak_gb": None}, {"state": "running"}, {"ended": None}): + bad = R.Entry(**{**{f.name: getattr(e, f.name) for f in R.fields(R.Entry)}, **change}) + self.assertIsNone(R.history_record(bad), change) + + def test_suggest_p95_p90(self): + # peaks 1..10 → p95 = 10 → ×1.15 = 11.5; durations 6..60 → p90 = 54 → ×1.5 = 81 + self.assertEqual(R.suggest(TEN, "proj", "gate"), (11.5, 81, 10)) + + def test_suggest_floors_and_rounding(self): + tiny = [run(id=f"r-{i}", peak=0.01, minutes=1.0) for i in range(3)] + self.assertEqual(R.suggest(tiny, "proj", "gate"), (0.2, 5, 3)) + odd = [run(id=f"r-{i}", peak=1.01, minutes=7.1) for i in range(3)] + self.assertEqual(R.suggest(odd, "proj", "gate"), (1.2, 11, 3)) # 1.1615 → 1.2 up; 10.65 → 11 up + + def test_suggest_needs_three_matching(self): + two = TEN[:2] + [run(id="r-a", project="other"), run(id="r-b", title="build"), + run(id="r-c", rc=1), run(id="r-d", peak=None)] + self.assertIsNone(R.suggest(two, "proj", "gate")) + self.assertIsNone(R.suggest([], "proj", "gate")) + + def test_hint_line(self): + sug = (11.5, 81, 10) + line = "hint: history says ~11.5 GB / 81 min (10 runs)" + self.assertEqual(R.hint_line(23.1, 60, sug), line) + self.assertEqual(R.hint_line(10.0, 163, sug), line) + self.assertIsNone(R.hint_line(23.0, 162, sug)) + self.assertIsNone(R.hint_line(99.0, 999, None)) + + def test_merge_runs(self): + led = R.Ledger(entries=[running(id="r-4", project="proj", title="gate 1234567"), running(id="r-5")]) + led.entries[0].state, led.entries[0].rc, led.entries[0].peak_gb = "done", 0, 2.0 + led.entries[0].ended = led.entries[0].started + dt.timedelta(minutes=3) + merged = R.all_runs([run(id="r-1"), run(id="r-4", peak=7.0)], led) + self.assertEqual([(r["id"], r["peak_gb"]) for r in merged], [("r-1", 1.0), ("r-4", 7.0)]) + merged = R.all_runs([run(id="r-1")], led) + self.assertEqual([(r["id"], r["peak_gb"], r["min"]) for r in merged], [("r-1", 1.0, 10.0), ("r-4", 2.0, 3.0)]) + + def test_hist_lines(self): + runs = TEN + [run(id="r-x", title="wf-batch", peak=3.0, minutes=20.0, mem=6.0, est=400), + run(id="r-y", project="other", title="build")] + self.assertEqual(R.hist_lines(runs, "proj"), [ + "proj · gate · n 10 · req 12.0 GB · peak 5.0/10.0 GB · est 1h · dur 30m/54m · suggest 11.5 GB 1h21m", + "proj · wf-batch · n 1 · req 6.0 GB · peak 3.0/3.0 GB · est 6h40m · dur 20m/20m · suggest - (< 3 runs)"]) + self.assertEqual(len(R.hist_lines(runs, None)), 3) + self.assertEqual(R.hist_lines([], None), ["no finished runs with a peak yet"]) + + +ORCH_LOG = """\ +2026-10-06T10:00:00 slow opus t-a done abc1234 10m00s +2026-10-06T10:10:00 fast haiku t-b handback - 20m30s +2026-10-06T10:20:00 slow opus t-c done+gate-red abc1234 1h05m +2026-10-06T10:30:00 fast haiku t-d post-check-red abc1234 50m00s (branch x still there) +2026-10-06T10:40:00 fast haiku t-e done abc1234 0m00s +2026-10-06T10:50:00 slow opus t-f done abc1234 - +2026-10-06T11:00:00 cloud opus t-g done abc1234 3h00m +2026-10-06T11:10:00 slow sonnet t-h awaiting - 40m00s +21:30 pilot batch 1 rc=0 done=t-a stopped= added= +ALERT wf-pilot proj: awaiting a-1 +2026-10-06T11:20:00 slow opus t-i done abc1234 30m00s +""" + + +class TaskFit(unittest.TestCase): + def test_durations_from_log(self): + # done/handback rows of local lanes with a real duration only (minutes) + self.assertEqual(R.orch_durations(ORCH_LOG), [10.0, 20.5, 65.0, 30.0]) + + def test_p90(self): + # 4 rows: nearest rank ceil(0.9*4)=4 → 65 min + self.assertEqual(R.task_p90(ORCH_LOG), (65, 4)) + self.assertEqual(R.task_p90(""), (30, 0)) + two = "2026-10-06T10:00:00 slow opus t-a done abc 5m00s\n" * 2 + self.assertEqual(R.task_p90(two), (30, 2)) # < 3 runs → default 30 min + odd = "2026-10-06T10:00:00 slow opus t-a done abc 4m10s\n" * 3 + self.assertEqual(R.task_p90(odd), (5, 3)) # rounded up to whole minutes + + def test_batch_fit(self): + self.assertEqual(R.batch_fit(4, 150, 30), 4) + self.assertEqual(R.batch_fit(4, 90, 30), 3) + self.assertEqual(R.batch_fit(4, 59, 30), 1) + self.assertEqual(R.batch_fit(4, 29, 30), 0) + self.assertEqual(R.batch_fit(4, 0, 30), 0) |
