diff options
Diffstat (limited to 'tests/test_cloud.py')
| -rw-r--r-- | tests/test_cloud.py | 1174 |
1 files changed, 1174 insertions, 0 deletions
diff --git a/tests/test_cloud.py b/tests/test_cloud.py new file mode 100644 index 0000000..e922f9c --- /dev/null +++ b/tests/test_cloud.py @@ -0,0 +1,1174 @@ +import datetime as dt +import io +import json +import os +import shutil +import sys +import tempfile +import unittest +import subprocess +import textwrap +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_cloud # noqa: E402 +from wflib import cloud as C, usage as U # noqa: E402 + +T0 = dt.datetime(2026, 10, 6, 12, 0, tzinfo=dt.timezone.utc) + + +def led_with(n_running, **kw): + led = {**C.new(), **kw} + for i in range(n_running): + C.add(led, f"t{i}", "p", f"s{i}", "claude-sonnet-5-5", T0) + return led + + +class Pure(unittest.TestCase): + def test_default_budget(self): + self.assertEqual(C.new()["budget"], 240.0) + + def test_balance(self): + # 240 - 10 - 4*2 = 222 + self.assertEqual(C.balance(led_with(2, spent=10.0)), 222.0) + self.assertEqual(C.balance(C.new()), 240.0) + + def test_refuse_low_balance(self): + # 20 - 17 - 0 = 3 < 4 + self.assertIn("balance", C.refusal(led_with(0, budget=20.0, spent=17.0))) + # exactly 4 left: ok + self.assertIsNone(C.refusal(led_with(0, budget=20.0, spent=16.0))) + + def test_refuse_max_parallel(self): + self.assertIn("max_parallel", C.refusal(led_with(3))) + self.assertIsNone(C.refusal(led_with(2))) + + def test_balance_counts_reserve_toward_refusal(self): + # 10 - 0 - 4*2 = 2 < 4 + self.assertIn("balance", C.refusal(led_with(2, budget=10.0))) + + def test_set_balance(self): + led = led_with(0, spent=50.0) + row = C.set_balance(led, 200.0, T0) + self.assertEqual(led["spent"], 40.0) + self.assertEqual((row["usd_source"], row["usd"]), ("owner", -10.0)) + + def test_charge_from_usage(self): + led = led_with(1) + u = U.Usage(turns=1, inp=1_000_000, out=1_000_000) # sonnet-5-5: 2 + 10 = 12 ; x1.15 = 13.8 + e = C.end(led, "s0", "done", u) + self.assertAlmostEqual(e["usd"], 13.8) + self.assertEqual(e["usd_source"], "self") + self.assertAlmostEqual(led["spent"], 13.8) + self.assertEqual(C.running(led), 0) + + def test_charge_missing_usage_is_reserve(self): + led = led_with(1) + e = C.end(led, "s0", "handback", None) + self.assertEqual((e["usd"], e["usd_source"]), (4.0, "est")) + + def test_end_twice_refused(self): + led = led_with(1) + C.end(led, "s0", "done", None) + with self.assertRaises(C.CloudError): + C.end(led, "s0", "done", None) + + def test_expire_lost_after_24h(self): + led = led_with(2) + led["entries"][1]["sent"] = (T0 + dt.timedelta(hours=20)).isoformat() + lost = C.expire(led, T0 + dt.timedelta(hours=25)) + self.assertEqual([e["sid"] for e in lost], ["s0"]) + self.assertEqual(led["entries"][0]["state"], "lost") + self.assertEqual(led["spent"], 4.0) + self.assertEqual(C.running(led), 1) + + def test_roundtrip_and_bad(self): + led = led_with(1) + self.assertEqual(C.loads(C.dumps(led)), led) + with self.assertRaises(C.CloudError): + C.loads("{nope") + + +class ArchivePure(unittest.TestCase): + CREDS = '{"claudeAiOauth": {"accessToken": "tok-1", "refreshToken": "r"}, "other": 1}' + + def test_request_url_and_headers(self): + url, h = C.archive_request("session_01AbCdEf", self.CREDS, "2.1.291") + self.assertEqual(url, "https://api.anthropic.com/v1/code/sessions/session_01AbCdEf/archive") + self.assertEqual(h, {"Authorization": "Bearer tok-1", "Content-Type": "application/json", + "anthropic-version": "2023-06-01", "User-Agent": "claude-code/2.1.291"}) + + def test_trusted_device_token_sent(self): + creds = '{"claudeAiOauth": {"accessToken": "tok-1"}, "trustedDeviceToken": "dev-9"}' + _, h = C.archive_request("session_01AbCdEf", creds, "2.1.291") + self.assertEqual(h["X-Trusted-Device-Token"], "dev-9") + + def test_no_token_or_bad_sid(self): + for creds in ("", "{}", '{"claudeAiOauth": {}}', "not json"): + with self.assertRaisesRegex(C.CloudError, "no claude.ai login token"): + C.archive_request("session_01AbCdEf", creds, "1") + for sid in ("pending:x", "-", "session_01/../x"): + with self.assertRaisesRegex(C.CloudError, "not a cloud session id"): + C.archive_request(sid, self.CREDS, "1") + + def test_problem(self): + self.assertIsNone(C.archive_problem(200, "")) + self.assertIsNone(C.archive_problem(409, "already")) + self.assertEqual(C.archive_problem(401, "x"), "HTTP 401 (login expired? run claude once)") + self.assertEqual(C.archive_problem(404, "no\nsuch " + "y" * 200), "HTTP 404: no such " + "y" * 102) + + def test_unarchived(self): + led = led_with(3) + C.end(led, "s0", "done") + C.end(led, "s1", "lost") + led["entries"][1]["archived"] = True + for e in led["entries"]: + e["sid"] = "session_" + e["sid"] + C.set_balance(led, 200, T0) + self.assertEqual(C.unarchived(led), ["session_s0"]) + + +class Cli(unittest.TestCase): + def run_wf(self, st, *argv): + out = io.StringIO() + with redirect_stdout(out): + code = wf_cloud.main(list(argv), state=st, now=T0) + return code, out.getvalue() + + def test_ledger_flow(self): + with tempfile.TemporaryDirectory() as d: + st = Path(d) + code, out = self.run_wf(st, "ledger") + self.assertEqual(code, 0) + self.assertIn("budget $240.00 spent $0.00 balance $240.00 running 0/3", out) + _, out = self.run_wf(st, "ledger", "--budget", "100", "--set-balance", "70") + self.assertIn("budget $100.00 spent $30.00 balance $70.00", out) + _, out = self.run_wf(st, "ledger") # persisted + self.assertIn("spent $30.00", out) + + +# Synthetic, same shape as a real `claude --cloud` run in a pty (escape codes, CRLF, title line). +SEND_OUT = ("\x1b[?2004l\x1b]3008;start=x;type=command\x1b\\\x1b7\x1b[r\x1b8\x1b[?25h\x1b[>4;2m\x1b[c" + "Created cloud session: Some title\r\n" + "View: https://claude.ai/code/session_01AbCdEfGhIjKlMnOpQrStUv?from=cli&m=0\r\n" + "Resume with: claude --teleport session_01AbCdEfGhIjKlMnOpQrStUv\r\n") + +FAKE = r"""#!/usr/bin/env python3 +import os, sys, time, tty +mode, st = os.environ.get("FAKE_MODE", ""), os.environ["FAKE_STATE"] +def say(s): sys.stdout.write(s); sys.stdout.flush() +def line(): + return sys.stdin.readline().strip() +n = int(open(st).read()) if os.path.exists(st) else 0 +open(st, "w").write(str(n + 1)) +if "trust" in mode and n == 0: + say("\x1b[2J Do\x1b[4Gyou trust the files in this folder?\r\n 1. No, exit\r\n 2. Yes, proceed\r\n") + if line() != "\x1b[B": + say("refused\r\n"); sys.exit(1) +if sys.argv[1] == "--cloud": + if os.environ.get("FAKE_ARGS"): + open(os.environ["FAKE_ARGS"], "w").write(os.getcwd() + "\n" + sys.argv[2]) + if "nosid" in mode: + say("Error: something\r\nline2\r\n"); sys.exit(1) + say(%r); sys.exit(0) +assert sys.argv[1] == "--teleport", sys.argv +say("\x1b[3G◯\x1b[5GChecking\x1b[14Gout\x1b[18Gbranch\r\n") +if "noresume" in mode or ("flaky" in mode and n == 0): + sys.exit(0) +say("●\x1b[3GSession\x1b[11Gresumed\r\n❯ ") +while True: + cmd = line() + if cmd.startswith("/export "): + if line() == "": # 2nd Enter (autocomplete) needed + open(cmd.split(" ", 1)[1], "w").write("● WF-RESULT done\n") + say("Conversation exported\r\n") + elif cmd == "/exit": + sys.exit(0) +""" % SEND_OUT + + +class Parse(unittest.TestCase): + def test_sid_from_send_output(self): + self.assertEqual(C.parse_sid(SEND_OUT), "session_01AbCdEfGhIjKlMnOpQrStUv") + + def test_sid_from_view_only(self): + self.assertEqual(C.parse_sid("View: https://claude.ai/code/session_01Zz9Zz9Zz9Zz9?from=cli\n"), + "session_01Zz9Zz9Zz9Zz9") + + def test_sid_absent(self): + self.assertIsNone(C.parse_sid("Error: not logged in\nclaude --teleport\n")) + + def test_strip_ansi_gaps(self): + # CSI n G (column) and CSI n C (forward) are word gaps in the TUI + self.assertEqual(C.strip_ansi("\x1b[39m●\x1b[3GSession\x1b[11Gresumed\r\n"), "● Session resumed\n") + self.assertEqual(C.strip_ansi("a\x1b[2Cb"), "a b") + + def test_last_lines(self): + self.assertEqual(C.last_lines("a\r\n\r\nb\nc\n", 2), ["b", "c"]) + + def test_trust_keeps_other_keys(self): + out = json.loads(C.trust('{"x": 1, "projects": {"/a": {"k": 2}}}', "/b")) + self.assertEqual(out, {"x": 1, "projects": {"/a": {"k": 2}, "/b": {"hasTrustDialogAccepted": True}}}) + out = json.loads(C.trust('{"projects": {"/a": {"k": 2}}}', "/a")) + self.assertEqual(out["projects"]["/a"], {"k": 2, "hasTrustDialogAccepted": True}) + self.assertEqual(json.loads(C.trust("", "/c")), {"projects": {"/c": {"hasTrustDialogAccepted": True}}}) + with self.assertRaises(C.CloudError): + C.trust("[1]", "/c") + + +class PtyDriver(unittest.TestCase): + """send/export against a fake claude on a real pty.""" + + def setUp(self): + self.d = Path(tempfile.mkdtemp()) + fake = self.d / "claude" + fake.write_text(FAKE) + fake.chmod(0o755) + self.repo = self.d / "repo" + self.repo.mkdir() + self.cj = self.d / "claude.json" + self.cj.write_text('{"keep": true}') + self.env = {"WF_CLAUDE_JSON": str(self.cj), "FAKE_STATE": str(self.d / "n")} + self.old = (wf_cloud.CLAUDE, wf_cloud.SETTLE, dict(os.environ)) + wf_cloud.CLAUDE, wf_cloud.SETTLE = str(fake), 0.3 + os.environ.update(self.env) + + def tearDown(self): + wf_cloud.CLAUDE, wf_cloud.SETTLE, env = self.old + os.environ.clear() + os.environ.update(env) + shutil.rmtree(self.d) + + def test_send_parses_sid_and_pre_trusts(self): + self.assertEqual(wf_cloud.send(self.repo, "do it", timeout=10), "session_01AbCdEfGhIjKlMnOpQrStUv") + data = json.loads(self.cj.read_text()) + self.assertEqual(data["projects"][str(self.repo.resolve())], {"hasTrustDialogAccepted": True}) + self.assertTrue(data["keep"]) + + def test_send_answers_trust_dialog(self): + os.environ["FAKE_MODE"] = "trust" + self.assertEqual(wf_cloud.send(self.repo, "do it", timeout=10), "session_01AbCdEfGhIjKlMnOpQrStUv") + + def test_send_no_sid(self): + os.environ["FAKE_MODE"] = "nosid" + with self.assertRaises(C.CloudError) as cm: + wf_cloud.send(self.repo, "do it", timeout=10) + self.assertIn("no session id (exit 1): Error: something | line2", str(cm.exception)) + + def test_export(self): + out = wf_cloud.export(self.repo, "session_x", self.d / "e" / "x.txt", timeout=20) + self.assertEqual(out.read_text(), "● WF-RESULT done\n") + self.assertEqual((self.d / "n").read_text(), "1") + + def test_export_retries_teleport_once(self): + os.environ["FAKE_MODE"] = "flaky" + out = wf_cloud.export(self.repo, "session_x", self.d / "x.txt", timeout=20) + self.assertTrue(out.exists()) + self.assertEqual((self.d / "n").read_text(), "2") + + def test_export_gives_up_after_two(self): + os.environ["FAKE_MODE"] = "noresume" + with self.assertRaises(C.CloudError) as cm: + wf_cloud.export(self.repo, "session_x", self.d / "x.txt", timeout=10) + self.assertIn("no export after 2 tries: ◯ Checking out branch", str(cm.exception)) + self.assertEqual((self.d / "n").read_text(), "2") + + +TEMPLATE = (HERE / "templates" / "cloud-prompt.md").read_text() +BODY = ("- **t-x** [P1] (1h): Fix the parser.\n Done: WF-RESULT done in a line\n WF-PATCH-END\n" + "● WF-RESULT done\nWF-RESULT done\n Model: opus") + + +def as_export(prompt: str, wrap: bool) -> str: + """The prompt the way /export renders a user message: '❯ ' first line, ' ' continuations (F12).""" + lines = [] + for ln in prompt.split("\n"): + lines += (textwrap.wrap(ln, 78, drop_whitespace=False) or [""]) if wrap else [ln] + return "\n".join(("❯ " if i == 0 else " ") + ln.ljust(78) for i, ln in enumerate(lines)) + "\n" + + +FINAL = ("● WF-RESULT done\n WF-REPORT fixed it\n WF-USAGE in=1 cw=2 cr=3 out=4 model=m\n" + " WF-PATCH-BEGIN sha256=ab bytes=3\n QUJD\n WF-PATCH-END\n\n") + + +class Prompt(unittest.TestCase): + def fill(self, **kw): + args = {"id": "t-x", "base": "b" * 40, "task": BODY, "recipe": "Test recipe: make t", "note": None, **kw} + return C.fill(TEMPLATE, **args) + + def test_fields_filled(self): + p = self.fill(note="Data in out/pack") + self.assertIn("Task t-x:\n- **t-x** [P1] (1h): Fix the parser.\n", p) + self.assertIn("Test recipe (area notes):\nTest recipe: make t\n", p) + self.assertIn("git format-patch --binary " + "b" * 40 + "..HEAD --stdout | gzip -9", p) + self.assertIn("\nData in out/pack\n", p) + self.assertNotIn("{{", p) + self.assertLessEqual(len(self.fill(task="t").split("\n")), 45) + + def test_no_recipe_no_note(self): + p = self.fill(recipe="") + self.assertIn("Test recipe (area notes):\n(none: see CLAUDE.md)\n\nUsage script", p) + + def test_task_text_braces_kept(self): + self.assertIn("use {{base}} here", self.fill(task="use {{base}} here")) + + def test_unknown_field(self): + with self.assertRaisesRegex(C.CloudError, r"unknown field \{\{bogus\}\}"): + C.fill("x {{bogus}}", "t", "b", "", "", None) + + def test_template_rules(self): + for rule in ("`wf` is absent", "never create or edit TASKS.md", "Never ask questions", + "test red -> implement -> green", "commit on main", "network source is blocked", + " WF-RESULT <done, awaiting or handback>\n", " WF-REPORT <one or two lines", + 'print("WF-USAGE in=%d cw=%d cr=%d out=%d model=%s"', 'echo "WF-PATCH-BEGIN sha256=$(', + "echo WF-PATCH-END", "base64 -w 76", "~/.claude/projects", "never -A", "__pycache__"): + self.assertIn(rule, TEMPLATE) + # live run 2026-10-06: '[WF-RESULT R]' placeholders made the session drop every key word + self.assertNotIn("[WF-", TEMPLATE) + self.assertIn("key word", TEMPLATE) + + def test_parser_never_matches_prompt(self): + p = self.fill() + for wrap in (False, True): + self.assertIsNone(C.final_message(as_export(p, wrap)), wrap) + + def test_parser_takes_final_message_after_prompt(self): + exp = as_export(self.fill(), True) + "\n● Working on it.\n Ran 3 shell commands\n\n" + FINAL + self.assertEqual(C.final_message(exp), ["WF-RESULT done", "WF-REPORT fixed it", + "WF-USAGE in=1 cw=2 cr=3 out=4 model=m", + "WF-PATCH-BEGIN sha256=ab bytes=3", "QUJD", "WF-PATCH-END"]) + + def test_parser_last_result_wins_and_stops_at_next_message(self): + exp = "● WF-RESULT handback\n WF-REPORT old\n\n● WF-RESULT awaiting\n WF-REPORT q?\n\n❯ more\n" + self.assertEqual(C.final_message(exp), ["WF-RESULT awaiting", "WF-REPORT q?"]) + + def test_include_problem(self): + for bad in ("/abs", "../x", "a/../../b", "", ".git", ".git/config"): + self.assertIsNotNone(C.include_problem(bad), bad) + for ok in ("out/pack", "data.bin", "out/a/b/"): + self.assertIsNone(C.include_problem(ok), ok) + + def test_size_refusal(self): + self.assertIsNone(C.size_refusal(90_000_000)) + self.assertEqual(C.size_refusal(91_400_000), "snapshot 91 MB > 90 MB; trim cloud_include") + + +from test_claims import FOUR, git # noqa: E402 +from test_setup import GitCli # noqa: E402 + +AREAS = "# Demo\n\n## Areas\n\n### app\n- Test recipe: RECIPE-MARK python3 -m unittest\n- Paths: src/app.py\n" + + + +def stub_archive(tc, d: Path, status=200): + """No network: wf_cloud.post records (url, headers) and answers `status`; fake credentials file.""" + creds = d / "credentials.json" + creds.write_text('{"claudeAiOauth": {"accessToken": "tok-1"}}') + tc.posts, tc.answer = [], status + saved = (wf_cloud.post, wf_cloud.cli_version, os.environ.get("WF_CLAUDE_CREDENTIALS")) + + def post(url, headers, timeout=10): + tc.posts.append((url, headers)) + if isinstance(tc.answer, Exception): + raise tc.answer + return tc.answer, "body" + + def restore(): + wf_cloud.post, wf_cloud.cli_version, env = saved + os.environ.pop("WF_CLAUDE_CREDENTIALS", None) + if env is not None: + os.environ["WF_CLAUDE_CREDENTIALS"] = env + wf_cloud.post, wf_cloud.cli_version = post, lambda: "9.9.9" + os.environ["WF_CLAUDE_CREDENTIALS"] = str(creds) + tc.addCleanup(restore) + + +class CloudProject(GitCli): + """A cloud-opted git project + fake claude; wf cloud in-process.""" + tasks_text = FOUR.replace("Four.", "Four: fix src/app.py.") + toml = GitCli.toml + 'cloud = true\ncloud_include = ["out/pack"]\ncloud_note = "NOTE-MARK data in out/pack"\n' + + def setUp(self): + super().setUp() + (self.root / "src").mkdir() + (self.root / "src" / "app.py").write_text("v1\n") + (self.root / "CLAUDE.md").write_text(AREAS) + (self.root / "out").mkdir() + (self.root / "out" / "tracked.txt").write_text("tracked under out/\n") + git(self.root, "add", "-A") + git(self.root, "add", "-f", "out/tracked.txt") + git(self.root, "commit", "-qm", "app") + (self.root / "src" / "app.py").write_text("dirty\n") # working tree: not in the snapshot + (self.root / "out" / "pack").mkdir(parents=True) + (self.root / "out" / "pack" / "data.bin").write_text("pack\n") + (self.root / "out" / "other.txt").write_text("not included\n") + (self.root / ".wf").mkdir(exist_ok=True) + (self.root / ".wf" / "x").write_text("state\n") + self.d = Path(tempfile.mkdtemp()) + fake = self.d / "claude" + fake.write_text(FAKE) + fake.chmod(0o755) + self.st = self.d / "state" + self.old = (wf_cloud.CLAUDE, wf_cloud.CAP, dict(os.environ)) + wf_cloud.CLAUDE = str(fake) + os.environ.update({"WF_CLAUDE_JSON": str(self.d / "claude.json"), "FAKE_STATE": str(self.d / "n"), + "FAKE_ARGS": str(self.d / "args"), "CLAUDE_CODE_MESSAGING_SOCKET": "", + "WF_INBOX": str(self.d / "inbox.md")}) + self.snap = self.root / "out" / "cloud" / "t-four" + self.rec = self.root / ".wf" / "cloud" / "t-four.json" + stub_archive(self, self.d) + + def tearDown(self): + wf_cloud.CLAUDE, wf_cloud.CAP, env = self.old + os.environ.clear() + os.environ.update(env) + shutil.rmtree(self.d) + super().tearDown() + + def send(self, *extra): + out, err = io.StringIO(), io.StringIO() + with redirect_stdout(out), redirect_stderr(err): + code = wf_cloud.main(["send", "t-four", "--project", str(self.root), *extra], state=self.st, now=T0) + return code, out.getvalue(), err.getvalue() + + def ledger(self): + return json.loads((self.st / "cloud.json").read_text()) if (self.st / "cloud.json").exists() else None + + def sh(self, cwd, *args): + return subprocess.run(["git", "-C", str(cwd), *args], capture_output=True, text=True).stdout.strip() + +class Send(CloudProject): + """wf cloud send against the fake claude.""" + + def test_send_snapshot_prompt_record_claim(self): + master = self.sh(self.root, "rev-parse", "master") + code, out, err = self.send() + self.assertEqual((code, err), (0, ""), out) + sid = "session_01AbCdEfGhIjKlMnOpQrStUv" + files = sorted(str(p.relative_to(self.snap)) for p in self.snap.rglob("*") + if p.is_file() and ".git" not in p.relative_to(self.snap).parts) + self.assertEqual(files, [".gitignore", "CLAUDE.md", "DESIGN.md", "docs/plan.md", "out/pack/data.bin", + "src/app.py", "workflow.toml"]) + self.assertEqual((self.snap / "src" / "app.py").read_text(), "v1\n") + self.assertEqual(self.sh(self.snap, "log", "--format=%s"), f"base {master}") + self.assertEqual(self.sh(self.snap, "ls-files").split("\n"), files) # all committed (-f) + self.assertEqual(self.sh(self.snap, "status", "--porcelain"), "") + self.assertEqual(self.sh(self.snap, "branch", "--show-current"), "main") + base = self.sh(self.snap, "rev-parse", "HEAD") + cwd, prompt = (self.d / "args").read_text().split("\n", 1) + self.assertEqual(Path(cwd), self.snap.resolve()) + self.assertIn("Task t-four:\n- **t-four** [P3] (1h): Four: fix src/app.py.\n", prompt) + self.assertIn("RECIPE-MARK", prompt) + self.assertIn("NOTE-MARK data in out/pack", prompt) + self.assertIn(f"format-patch --binary {base}..HEAD", prompt) + rec = json.loads(self.rec.read_text()) + self.assertEqual({k: rec[k] for k in ("id", "sid", "master", "base", "lane", "sent", "model")}, + {"id": "t-four", "sid": sid, "master": master, "base": base, "lane": "slow", + "sent": "2026-10-06T12:00:00+00:00", "model": "claude-opus-5-5"}) + self.assertIn(f"**t-four** [P3] (1h) (in progress: cloud:{sid}): Four", self.tasks()) + led = self.ledger() + self.assertEqual([(e["id"], e["sid"], e["state"]) for e in led["entries"]], [("t-four", sid, "running")]) + self.assertIn(f"sent t-four: {sid}", out) + code, _, err = self.send() # twice: refused + self.assertEqual(code, 1) + self.assertIn(f"t-four already sent ({sid})", err) + + def test_dry_run_sends_nothing(self): + code, out, err = self.send("--dry-run") + self.assertEqual((code, err), (0, "")) + self.assertRegex(out, r"^snapshot 0\.0 MB \(master [0-9a-f]{12}, base [0-9a-f]{12}\)\n\nYou are a remote") + self.assertIn("Task t-four:", out) + self.assertFalse(self.snap.exists()) + self.assertFalse(self.snap.parent.exists()) # out/cloud gone too, out/ kept + self.assertFalse(self.rec.exists()) + self.assertIsNone(self.ledger()) + self.assertFalse((self.d / "n").exists()) # claude never ran + self.assertNotIn("in progress", self.tasks()) + + def test_dry_run_keeps_live_snapshot(self): + self.assertEqual(self.send()[0], 0) + head = self.sh(self.snap, "rev-parse", "HEAD") + code, out, err = self.send("--dry-run") + self.assertEqual((code, err), (0, "")) + self.assertEqual(self.sh(self.snap, "rev-parse", "HEAD"), head) # live session's folder intact + self.assertEqual(sorted(p.name for p in self.snap.parent.iterdir()), ["t-four"]) # dry-run folder dropped + + def test_size_refused(self): + wf_cloud.CAP = 10 + code, _, err = self.send() + self.assertEqual(code, 1) + self.assertIn("MB > 0 MB; trim cloud_include", err) + self.assertFalse(self.snap.exists()) + self.assertIsNone(self.ledger()) + self.assertFalse((self.d / "n").exists()) + + def test_ledger_refused_exit_3(self): + with wf_cloud.locked(self.st) as led: + for i in range(3): + C.add(led, f"t{i}", "p", f"s{i}", "claude-opus-5-5", T0) + code, out, err = self.send() + self.assertEqual(code, 3) + self.assertEqual(err.count("\n"), 1) + self.assertIn("cloud refused: ", err) + self.assertIn("max_parallel", err) + self.assertEqual(len(self.ledger()["entries"]), 3) + self.assertFalse(self.snap.exists()) + self.assertFalse(self.rec.exists()) + self.assertFalse((self.d / "n").exists()) + + def test_send_failure_releases_reserve(self): + os.environ["FAKE_MODE"] = "nosid" + code, _, err = self.send() + self.assertEqual(code, 1) + self.assertIn("no session id", err) + self.assertEqual(self.ledger()["entries"], []) + self.assertFalse(self.snap.exists()) + self.assertFalse(self.rec.exists()) + + def test_not_opted_in(self): + (self.root / "workflow.toml").write_text(GitCli.toml) + code, _, err = self.send() + self.assertEqual(code, 1) + self.assertIn("not opted in", err) + + def test_include_outside_refused(self): + (self.root / "workflow.toml").write_text(GitCli.toml + 'cloud = true\ncloud_include = ["../x"]\n') + code, _, err = self.send("--dry-run") + self.assertEqual(code, 1) + self.assertIn("cloud_include '../x': must be a path inside the project", err) + + +if __name__ == "__main__": + unittest.main() + + +class PullPure(unittest.TestCase): + def result(self, raw: bytes, sha=None, n=None): + import base64, gzip, hashlib + gz = gzip.compress(raw) + b64 = base64.b64encode(gz).decode() + return ["WF-RESULT done", "WF-REPORT fixed the parser,", "tests green", + "WF-USAGE in=1000 cw=2000 cr=3000 out=4000 model=claude-opus-5-5", + f"WF-PATCH-BEGIN sha256={sha or hashlib.sha256(gz).hexdigest()} bytes={n or len(gz)}", + *[b64[i:i + 76] for i in range(0, len(b64), 76)], "WF-PATCH-END"] + + def test_parse_and_decode(self): + r = C.parse_result(self.result(b"diff --git a/x b/x\n")) + self.assertEqual((r.state, r.report, r.model), ("done", "fixed the parser, tests green", "claude-opus-5-5")) + self.assertEqual((r.usage.inp, r.usage.cw, r.usage.cr, r.usage.out), (1000, 2000, 3000, 4000)) + self.assertEqual(C.decode_patch(r), b"diff --git a/x b/x\n") + + def test_header_wrapped_anywhere(self): + lines = self.result(b"abc") + head = lines[4] + lines[4:5] = [head[:30], head[30:61], head[61:]] # sha split over three lines + self.assertEqual(C.decode_patch(C.parse_result(lines)), b"abc") + + def test_mismatches(self): + with self.assertRaisesRegex(C.CloudError, "sha256 mismatch"): + C.decode_patch(C.parse_result(self.result(b"abc", sha="0" * 64))) + with self.assertRaisesRegex(C.CloudError, "bytes .* != 7 announced"): + C.decode_patch(C.parse_result(self.result(b"abc", n=7))) + with self.assertRaisesRegex(C.CloudError, "bad WF-RESULT 'maybe'"): + C.parse_result(["WF-RESULT maybe"]) + with self.assertRaisesRegex(C.CloudError, "without WF-PATCH-END"): + C.parse_result(self.result(b"abc")[:-1]) + self.assertIsNone(C.parse_result(["WF-RESULT handback", "WF-REPORT no"]).usage) + + def keyless(self, raw: bytes) -> str: + """The keys-dropped shape seen live: '● done', report, usage, header split, base64, bare WF-PATCH-END.""" + import base64, gzip, hashlib + gz = gzip.compress(raw) + b64 = base64.b64encode(gz).decode() + msg = ["done", "Added sub and a test.", "usage in=4 cw=5 cr=6 out=7 model=claude-opus-5-5", + f"sha256={hashlib.sha256(gz).hexdigest()}", f"bytes={len(gz)}", + *[b64[i:i + 76] for i in range(0, len(b64), 76)], "WF-PATCH-END"] + return "\n".join(("● " if i == 0 else " ") + ln for i, ln in enumerate(msg)) + "\n\n" + + def test_keyless_final_message(self): + prompt = as_export("do it\nend with\n [WF-PATCH-END]\nWF-PATCH-END", wrap=False) + text = prompt + "● working\n\n" + self.keyless(b"diff --git a/x b/x\n") + "● Session resumed\n" + self.assertIsNone(C.final_message(text)) + lines = C.keyless_message(text) + self.assertEqual((lines[0], lines[-1]), ("done", "WF-PATCH-END")) + r = C.keyless_result(lines) + self.assertEqual((r.state, r.model, r.usage.inp, r.usage.out), ("handback", "claude-opus-5-5", 4, 7)) + self.assertEqual(C.decode_patch(r), b"diff --git a/x b/x\n") + self.assertIsNone(C.keyless_message(prompt)) # prompt only: running + self.assertIsNone(C.keyless_message(prompt + "● still working\n")) + self.assertIsNone(C.keyless_message(text + as_export("redo with the keys", wrap=False))) # redo sent + self.assertIsNotNone(C.keyless_message(text + "❯ \n")) # empty input line + r = C.keyless_result(["done", "no patch here", "WF-PATCH-END"]) + self.assertEqual((r.has_patch, r.usage), (False, None)) + + def test_patch_files_and_problem(self): + patch = ("diff --git a/src/a.py b/src/a.py\n--- a/src/a.py\n+++ b/src/a.py\n" + "diff --git a/old.txt b/new.txt\nrename from old.txt\nrename to new.txt\n") + self.assertEqual(C.patch_files(patch), ["src/a.py", "old.txt", "new.txt"]) + never = ["TASKS.md", "tasks/archive.md", ".wf", "out", "data/pack"] + self.assertIsNone(C.patch_problem(["src/a.py", "outish.txt"], "", never)) + self.assertIn("TASKS.md", C.patch_problem(["TASKS.md"], "", never)) + self.assertIn("data/pack", C.patch_problem(["data/pack/x"], "", never)) + self.assertIn("outside the tree", C.patch_problem(["../x"], "", never)) + self.assertIn("outside the tree", C.patch_problem([".git/hooks/x"], "", never)) + self.assertIn("outside the project folder", C.patch_problem(["other/x"], "proj", never)) + self.assertIsNone(C.patch_problem(["proj/src/a.py"], "proj", never)) + self.assertIn("out", C.patch_problem(["proj/out/x"], "proj", never)) + + +class Pull(CloudProject): + """wf cloud pull --export FIXTURE on a git project with a recorded send.""" + toml = None # set in setUp: no worktree_setup (it writes log.txt), a quick_gate + + def setUp(self): + from test_cli import TOML + self.toml = (TOML + 'cloud = true\ncloud_include = ["out/pack"]\n' + 'quick_gate = ["grep -q fixed src/app.py"]\n') + super().setUp() + os.environ.update({"GIT_AUTHOR_NAME": "t", "GIT_AUTHOR_EMAIL": "t@t", "GIT_COMMITTER_NAME": "t", + "GIT_COMMITTER_EMAIL": "t@t"}) + (self.root / "src" / "app.py").write_text("v1\n") + import wf + from wflib import config + cfg = config.load(self.root) + self.master, self.base, _ = wf_cloud.snapshot(cfg, self.snap) + self.sid = "session_01PullPullPullPull" + with wf_cloud.locked(self.st) as led: + C.add(led, "t-four", str(self.root), self.sid, C.MODEL, T0) + self.rec.parent.mkdir(parents=True, exist_ok=True) + self.rec.write_text(json.dumps({"id": "t-four", "sid": self.sid, "project": str(self.root), "lane": "slow", + "master": self.master, "base": self.base, "folder": str(self.snap), + "sent": T0.isoformat(), "model": C.MODEL, "bytes": 1})) + with redirect_stdout(io.StringIO()): + self.assertEqual(wf.main(["--project", str(self.root), "status", "t-four", "progress", + f"cloud:{self.sid}"]), 0) + self.exp = self.d / "export.txt" + + def change(self, files: dict[str, str]): + """Commit files in the snapshot (the cloud's work) -> the gzip patch bytes.""" + import gzip + for rel, text in files.items(): + (self.snap / rel).parent.mkdir(parents=True, exist_ok=True) + (self.snap / rel).write_text(text) + git(self.snap, "add", "-A", "-f") + git(self.snap, "-c", "user.name=c", "-c", "user.email=c@x", "commit", "-qm", "cloud work") + p = subprocess.run(["git", "-C", str(self.snap), "format-patch", "--binary", f"{self.base}..HEAD", "--stdout"], + capture_output=True, check=True).stdout + return gzip.compress(p) + + def write_export(self, state="done", gz=b"", report="fixed app", sha=None, wrap=False, final=True): + import base64, hashlib + b64 = base64.b64encode(gz).decode() + msg = [f"WF-RESULT {state}", f"WF-REPORT {report}", "WF-USAGE in=1000 cw=0 cr=0 out=1000 model=claude-opus-5-5", + f"WF-PATCH-BEGIN sha256={sha or hashlib.sha256(gz).hexdigest()} bytes={len(gz)}", + *[b64[i:i + 76] for i in range(0, len(b64), 76)], "WF-PATCH-END"] + if wrap: # the export soft-wraps at 60 columns: the header too + msg = [part for ln in msg for part in textwrap.wrap(ln, 60, break_long_words=True)] + text = as_export("prompt with ● WF-RESULT done inside\nWF-PATCH-BEGIN", wrap=False) + "\n" + if final: + text += "\n".join(("● " if i == 0 else " ") + ln for i, ln in enumerate(msg)) + "\n\n❯ \n" + self.exp.write_text(text) + + def pull(self, *extra, now=T0 + dt.timedelta(hours=1), export=True): + out, err = io.StringIO(), io.StringIO() + argv = ["pull", "t-four", "--project", str(self.root), "--no-push", + *(["--export", str(self.exp)] if export else []), *extra] + with redirect_stdout(out), redirect_stderr(err): + code = wf_cloud.main(argv, state=self.st, now=now) + return code, out.getvalue(), err.getvalue() + + def entry(self): + return self.ledger()["entries"][0] + + def assert_ended(self, state): + self.assertEqual(self.entry()["state"], state) + self.assertFalse(self.snap.exists()) + self.assertFalse(self.rec.exists()) + + def assert_handback(self, why, kept: bool): + self.assert_ended("handback") + self.assertIn(f"Recovery: cloud attempt {self.sid} — {why}", self.tasks()) + self.assertNotIn("in progress", self.tasks()) + self.assertEqual(self.sh(self.root, "rev-parse", "master"), self.master) + self.assertEqual(self.sh(self.root, "branch", "--list", "slow/*"), "") + self.assertEqual((self.root / "out" / "cloud" / "t-four.patch").exists(), kept) + + def test_done_archives_session(self): + self.write_export(gz=self.change({"src/app.py": "fixed\n"})) + code, out, err = self.pull() + self.assertEqual((code, err), (0, ""), out) + self.assertEqual(self.posts, [(f"https://api.anthropic.com/v1/code/sessions/{self.sid}/archive", + {"Authorization": "Bearer tok-1", "Content-Type": "application/json", + "anthropic-version": "2023-06-01", "User-Agent": "claude-code/9.9.9"})]) + self.assertIs(self.entry()["archived"], True) + self.assertNotIn("archive", out.replace("archive.md", "")) + + def test_second_pull_same_id_skips(self): + """Two pulls of one id at once (batch sidecar + orchestrator): the one without the record lock skips + with rc 0, no 'branch … exists (local WIP?)'; the holder's pull still ends the task.""" + import fcntl + self.write_export(gz=self.change({"src/app.py": "fixed\n"})) + fd = os.open(self.rec, os.O_RDONLY) + try: + fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) # the other pull, mid-flight + code, out, err = self.pull() + finally: + os.close(fd) + self.assertEqual((code, err), (0, ""), out) + self.assertEqual(out, "t-four: pulled by another wf cloud pull, skipped\n") + self.assertTrue(self.rec.exists()) + self.assertEqual(self.entry()["state"], "running") + self.assertEqual(self.sh(self.root, "branch", "--list", "slow/*"), "") + code, out, err = self.pull() + self.assertEqual((code, err), (0, ""), out) + self.assert_ended("done") + + def test_pull_claim_record_gone(self): + with wf_cloud.pull_claim(self.d / "gone.json") as mine: + self.assertFalse(mine) + with wf_cloud.pull_claim(self.rec) as mine: + self.assertTrue(mine) + self.rec.unlink() # winner ended it while we waited: re-checked on the next claim + with wf_cloud.pull_claim(self.rec) as mine: + self.assertFalse(mine) + + def archive_fails(self, answer, why): + self.answer = answer + self.write_export(state="handback", report="cannot") + code, out, err = self.pull() + self.assertEqual((code, err), (0, ""), out) + self.assert_ended("handback") + self.assertIn(f"archive {self.sid} failed: {why}; archive by hand (wf cloud archive --ended)", out) + self.assertNotIn("archived", self.entry()) + + def test_archive_http_error_warns_pull_still_ends(self): + self.archive_fails(403, "HTTP 403: body") + + def test_archive_network_error_warns_pull_still_ends(self): + self.archive_fails(C.CloudError("timed out"), "timed out") + + def test_running_not_archived(self): + self.write_export(final=False) + code, out, err = self.pull() + self.assertEqual(code, 4, out + err) + self.assertEqual(self.posts, []) + + def test_archive_command_ended(self): + self.write_export(state="handback", report="cannot") + self.answer = 500 + self.pull() + self.answer = 200 + out = io.StringIO() + with redirect_stdout(out): + code = wf_cloud.main(["archive", "--ended"], state=self.st, now=T0) + self.assertEqual((code, out.getvalue()), (0, f"archived {self.sid}\n")) + self.assertIs(self.entry()["archived"], True) + with redirect_stdout(out := io.StringIO()): + code = wf_cloud.main(["archive", "--ended"], state=self.st, now=T0) + self.assertEqual((code, out.getvalue()), (0, "no ended cloud session left to archive\n")) + + def test_archive_command_one_sid_failure_exit_1(self): + self.answer = 401 + err = io.StringIO() + with redirect_stdout(io.StringIO()), redirect_stderr(err): + code = wf_cloud.main(["archive", "session_01Other"], state=self.st, now=T0) + self.assertEqual(code, 1) + self.assertEqual(err.getvalue(), "wf: archive session_01Other failed: HTTP 401 (login expired? run claude once)\n") + + def test_done_applies_gates_merges(self): + self.write_export(gz=self.change({"src/app.py": "fixed\n", "src/new.py": "new\n"})) + code, out, err = self.pull() + self.assertEqual((code, err), (0, ""), out) + self.assertEqual(self.sh(self.root, "show", "master:src/app.py"), "fixed") + self.assertEqual(self.sh(self.root, "show", "master:src/new.py"), "new") + self.assertIn("cloud work", self.sh(self.root, "log", "--format=%s", "master")) + self.assertNotIn("t-four", self.tasks()) + self.assertIn(f"fixed app (cloud {self.sid})", self.archive()) + self.assertEqual(self.sh(self.root, "status", "--porcelain", "TASKS.md", "tasks"), "") # bookkeeping committed + self.assert_ended("done") + e = self.entry() + self.assertEqual(e["usd_source"], "self") + self.assertAlmostEqual(e["usd"], U.cost(C.MODEL, U.Usage(inp=1000, out=1000)) * C.OVERHEAD) + self.assertFalse((self.root / "out" / "cloud" / "t-four.patch").exists()) + self.assertRegex(out, r"report: commit [0-9a-f]{7,}\n$") + self.assertIn("$ grep -q fixed src/app.py", out) + + def test_wrapped_indented_base64_rejoined(self): + self.write_export(gz=self.change({"src/app.py": "fixed\n"}), wrap=True) + code, out, err = self.pull() + self.assertEqual((code, err), (0, ""), out) + self.assertEqual(self.sh(self.root, "show", "master:src/app.py"), "fixed") + + def test_sha_mismatch_hands_back(self): + self.write_export(gz=self.change({"src/app.py": "fixed\n"}), sha="0" * 64) + code, out, err = self.pull() + self.assertEqual(code, 0, err) + self.assert_handback("patch sha256 mismatch", kept=False) + + def test_gate_red_hands_back_keeps_patch(self): + self.write_export(gz=self.change({"src/app.py": "broken\n"})) + code, out, err = self.pull() + self.assertEqual(code, 0, err) + self.assert_handback("quick_gate 'grep -q fixed src/app.py' red (exit 1)", kept=True) + wt = self.root / ".worktrees" / "slow" + self.assertEqual(self.sh(wt, "status", "--porcelain"), "") + self.assertEqual(self.sh(wt, "branch", "--show-current"), "") + + def test_patch_touching_cloud_include_refused(self): + self.write_export(gz=self.change({"src/app.py": "fixed\n", "out/pack/data.bin": "x\n"})) + code, out, err = self.pull() + self.assertEqual(code, 0, err) + self.assert_handback("patch touches out/pack/data.bin", kept=True) + + def test_patch_touching_tasks_refused(self): + self.write_export(gz=self.change({"src/app.py": "fixed\n", "TASKS.md": "x\n"})) + code, out, err = self.pull() + self.assertEqual(code, 0, err) + self.assert_handback("patch touches TASKS.md", kept=True) + + def test_awaiting_adds_question_and_blocks(self): + self.write_export(state="awaiting", report="Which format: csv or json?") + code, out, err = self.pull() + self.assertEqual(code, 0, err) + self.assert_ended("awaiting") + m = __import__("re").search(r"\*\*(a-[^*]+)\*\*: Which format: csv or json\? \[\[t-four\]\]", self.tasks()) + self.assertTrue(m, self.tasks()) + self.assertIn(f"(blocked: [[{m[1]}]])", self.tasks()) + self.assertEqual(self.sh(self.root, "rev-parse", "master"), self.master) + + def test_cloud_handback(self): + self.write_export(state="handback", report="needs the LAN host") + code, out, err = self.pull() + self.assertEqual(code, 0, err) + self.assert_handback("cloud handback: needs the LAN host", kept=False) + + def test_no_marker_running_exit_4(self): + self.write_export(final=False) + code, out, err = self.pull() + self.assertEqual((code, err), (4, "")) + self.assertIn(f"t-four: running ({self.sid}, sent 1h00m ago)", out) + self.assertTrue(self.rec.exists() and self.snap.exists()) + self.assertEqual(self.entry()["state"], "running") + self.assertIn(f"in progress: cloud:{self.sid}", self.tasks()) + + def test_keyless_final_message_hands_back_keeps_patch(self): + import gzip + self.write_export(final=False) + msg = PullPure().keyless(gzip.decompress(self.change({"src/app.py": "fixed\n"}))) + self.exp.write_text(self.exp.read_text() + msg + "● Session resumed\n") + code, out, err = self.pull() + self.assertEqual(code, 0, err) + self.assert_handback("bad result: no WF-RESULT key; patch out/cloud/t-four.patch", kept=True) + self.assertIn(b"+fixed", (self.root / "out" / "cloud" / "t-four.patch").read_bytes()) + self.assertEqual(self.entry()["usd_source"], "self") + + def test_teleport_export_used_and_removed(self): + self.write_export(final=False) + seen = [] + + def fake_export(folder, sid, out, timeout=300): + seen.append((Path(folder), sid)) + Path(out).parent.mkdir(parents=True, exist_ok=True) + Path(out).write_text(self.exp.read_text()) + return Path(out) + + old, wf_cloud.export = wf_cloud.export, fake_export + try: + code, out, err = self.pull(export=False) + finally: + wf_cloud.export = old + self.assertEqual((code, seen), (4, [(self.snap, self.sid)])) + self.assertFalse((self.root / "out" / "cloud" / "t-four.export.txt").exists()) + + def test_over_24h_lost(self): + self.write_export(final=False) + code, out, err = self.pull(now=T0 + dt.timedelta(hours=25)) + self.assertEqual(code, 0, err) + self.assert_ended("lost") + self.assertEqual((self.entry()["usd"], self.entry()["usd_source"]), (4.0, "est")) + self.assertIn(f"Recovery: cloud attempt {self.sid} — lost: no WF-RESULT after 25h00m", self.tasks()) + self.assertNotIn("in progress", self.tasks()) + + def test_all_and_not_sent(self): + os.environ["FAKE_MODE"] = "noresume" # teleport never resumes: export fails + out, err = io.StringIO(), io.StringIO() + with redirect_stdout(out), redirect_stderr(err): + code = wf_cloud.main(["pull", "--all", "--project", str(self.root)], state=self.st, + now=T0 + dt.timedelta(hours=25)) + self.assertEqual(code, 0, err.getvalue()) # no export, > 24h: lost + self.assert_ended("lost") + code, _, err = self.pull() + self.assertEqual(code, 1) + self.assertIn("t-four is not out in the cloud", err) + + +# ------------------------------------------------------------------ fit (spec §4.5) + +FIT_DOC = """\ +# Tasks — demo + +## Awaiting your decision + +## Pending + +- **t-ok** [P1] (1h): Plain code. Goal. + Done: tests pass +- **t-sonnet** [P1] (1h): Plain code. Goal. + Done: tests pass + Model: sonnet +- **t-slice** [P1] (5h): Big. Goal. + Done: tests pass +- **t-owner** [P1] (1h): Plain. Goal. + Done: tests pass + Sessions: owner +- **t-solo** [P1] (1h): Plain. Goal. + Done: tests pass + Sessions: solo +- **t-nodone** [P1] (1h): Plain. Goal. +- **t-no** [P1] (1h): Plain. Goal. + Done: tests pass + Cloud: no +- **t-wf** [P1] (1h): Plain. Goal. + Steps: run wf set x --model opus then check + Done: tests pass +- **t-res** [P1] (1h): Plain. Goal. + Steps: wf res run the build + Done: tests pass +- **t-gui** [P1] (1h): Plain. Goal. + Steps: click through the GUI + Done: tests pass +- **t-lan** [P1] (1h): Plain. Goal. + Steps: copy to 10.0.0.20 + Done: tests pass +- **t-srv** [P1] (1h): Plain. Goal. + Steps: hit the live server + Done: tests pass +- **t-yes-wf** [P1] (1h): Plain. Goal. + Steps: run wf set x --model opus then check + Done: tests pass + Cloud: yes +- **t-yes-sonnet** [P1] (1h): Plain. Goal. + Done: tests pass + Model: sonnet + Cloud: yes +- **t-yes-owner** [P1] (1h): Plain. Goal. + Done: tests pass + Sessions: owner + Cloud: yes +- **t-yes-haiku** [P1] (<1h): Plain. Goal. + Done: tests pass + Model: haiku + Cloud: yes +- **t-no-opus** [P1] (1h): Plain. Goal. + Done: tests pass + Model: opus + Cloud: no +- **t-yes-slice** [P1] (5h): Big. Goal. + Done: tests pass + Cloud: yes + +## Needs human + +## Deferred +""" + + +class FitTest(unittest.TestCase): + def setUp(self): + from wflib import tasks as TK + self.TK = TK + self.doc = TK.parse(FIT_DOC) + self.cfg = type("Cfg", (), {"cloud": True, "slice_above": "1h"})() + + def fit(self, id): + return C.fit(self.doc.item(id), self.cfg) + + def test_fits(self): + self.assertIsNone(self.fit("t-ok")) + + def test_exclusions(self): + for id, why in [("t-sonnet", "Model sonnet"), ("t-slice", "slice"), ("t-owner", "not runner-ready"), + ("t-solo", "Sessions: solo"), ("t-nodone", "not runner-ready"), ("t-no", "Cloud: no"), + ("t-wf", "`wf` commands"), ("t-res", "wf res"), ("t-gui", "GUI"), ("t-lan", "LAN"), + ("t-srv", "live server")]: + with self.subTest(id): + self.assertIn(why, self.fit(id) or "") + + def test_not_opted_in(self): + self.cfg.cloud = False + self.assertIn("not opted in", self.fit("t-ok")) + + def test_cloud_yes_skips_regex_only_opus_only(self): + self.assertIsNone(self.fit("t-yes-wf")) + self.assertEqual(self.fit("t-yes-sonnet"), "Model sonnet stays local") + self.assertEqual(self.fit("t-yes-haiku"), "Model haiku stays local") + self.assertEqual(self.fit("t-no-opus"), "Cloud: no") + self.assertEqual(self.fit("t-sonnet"), "Model sonnet stays local") # unset: old rule + self.assertIn("not runner-ready", self.fit("t-yes-owner")) + self.assertIn("slice", self.fit("t-yes-slice")) + + def test_cloud_property_and_set(self): + self.assertEqual(self.doc.item("t-no").cloud, "no") + self.assertIsNone(self.doc.item("t-ok").cloud) + self.TK.set_fields(self.doc, "t-ok", cloud="yes") + self.assertEqual(self.doc.item("t-ok").body, [" Done: tests pass", " Cloud: yes"]) + self.TK.set_fields(self.doc, "t-ok", cloud="no", model="sonnet") + self.assertEqual(self.doc.item("t-ok").body, [" Done: tests pass", " Model: sonnet", " Cloud: no"]) + self.TK.set_fields(self.doc, "t-ok", cloud="") + self.assertEqual(self.doc.item("t-ok").body, [" Done: tests pass", " Model: sonnet"]) + with self.assertRaises(self.TK.TaskError): + self.TK.set_fields(self.doc, "t-ok", cloud="maybe") + + +class PickKeyTest(unittest.TestCase): + def test_order(self): + from wflib import tasks as TK + doc = TK.parse("""# Tasks + +## Pending + +- **t-p1-small** [P1] (<1h): A. G. + Done: x +- **t-p1-big** [P1] (1h): A. G. + Done: x +- **t-p2-yes** [P2] (<1h): A. G. + Done: x + Cloud: yes +- **t-p0** [P0] (<1h): A. G. + Done: x + +## Needs human + +## Deferred +""") + items = [doc.item(i) for i in ("t-p1-small", "t-p1-big", "t-p2-yes", "t-p0")] + got = [i.id for i in sorted(items, key=C.pick_key)] + # yes first; then prio; then 1h before <1h + self.assertEqual(got, ["t-p2-yes", "t-p0", "t-p1-big", "t-p1-small"]) + + +class CloudLineCheck(unittest.TestCase): + def test_check(self): + sys.path.insert(0, str(HERE / "tests")) + from test_check import Base + + class B(Base): + def runTest(self): + self.pending("- **t-a** [P1] (1h): A. G.\n Done: x\n Cloud: yes\n") + self.assertEqual(self.run_check(), ([], [])) + self.pending("- **t-a** [P1] (1h): A. G.\n Done: x\n Cloud: maybe\n") + self.assertTrue(any("Cloud 'maybe'" in e for e in self.errors())) + self.pending("- **t-a** [P1] (1h): A. G.\n Done: x\n Cloud: yes\n Cloud: no\n") + self.assertTrue(any("two Cloud lines" in e for e in self.errors())) + r = unittest.TextTestRunner(stream=io.StringIO()).run(B()) + self.assertTrue(r.wasSuccessful(), r.failures + r.errors) + + +ORCH_TASKS = """\ +# Tasks — demo + +## Awaiting your decision + +## Pending + +- **t-fast** [P1] (<1h): Fast sonnet task. + - Done: fast works + - Model: sonnet + +- **t-gui** [P1] (<1h): Check the GUI screenshot. + - Done: gui works + +- **t-hai** [P1] (<1h): Haiku forced to the cloud. + - Done: hai works + - Model: haiku + - Cloud: yes + +- **t-slow** [P2] (1h): Slow opus task. + - Done: slow works + +- **t-yes** [P3] (<1h): Opus task marked for the cloud. + - Done: yes works + - Cloud: yes + +## Needs human + +## Deferred +""" + + +class OrchCloud(CloudProject): + """wf orch pick cloud (subprocess, fake claude via WF_CLAUDE, ledger via XDG_STATE_HOME).""" + tasks_text = ORCH_TASKS + + def setUp(self): + super().setUp() + from test_merge import IDENT + self.env = {**IDENT, "WF_CLAUDE": wf_cloud.CLAUDE, "XDG_STATE_HOME": str(self.d / "xdg")} + + def orch(self, *args): + out = self.ok("orch", *args, env=self.env) + return out + + def put_ledger(self, **kw): + n = kw.pop("running", 0) + led = {**C.new(), **kw} + for i in range(n): + C.add(led, f"t-x{i}", "/x", f"session_x{i}", "opus", T0) + (self.d / "xdg" / "wf").mkdir(parents=True, exist_ok=True) + (self.d / "xdg" / "wf" / "cloud.json").write_text(C.dumps(led)) + + def status(self, id): + return next(l for l in (self.root / "TASKS.md").read_text().splitlines() if f"**{id}**" in l) + + def test_pick_cloud_yes_first_then_opus_then_none(self): + out = self.orch("pick", "cloud") + self.assertIn("pick: t-yes (lane cloud, model opus", out) # Cloud: yes beats P2 1h; haiku never + out = self.orch("pick", "cloud") + self.assertIn("pick: t-slow (lane cloud, model opus", out) + self.assertIn("in progress: cloud:session_01AbCdEfGhIjKlMnOpQrStUv", self.status("t-slow")) + rec = json.loads((self.root / ".wf" / "orch" / "t-slow.json").read_text()) + self.assertEqual((rec["lane"], rec["task_lane"], rec["branch"]), ("cloud", "slow", "slow/t-slow")) + self.assertTrue((self.root / ".wf" / "cloud" / "t-slow.json").exists()) + out = self.orch("pick", "cloud") + self.assertEqual(out.strip().splitlines()[0][:26], "stop lane cloud: none fit ") + led = json.loads((self.d / "xdg" / "wf" / "cloud.json").read_text()) + self.assertEqual([e["id"] for e in led["entries"]], ["t-yes", "t-slow"]) + + def test_local_lane_skips_cloud_claim(self): + self.orch("pick", "cloud", "--id", "t-slow") + out = self.orch("pick", "slow") + self.assertNotIn("t-slow", out.splitlines()[0]) + self.assertIn("pick: t-fast", out) + + def test_stop_ledger(self): + self.put_ledger(budget=3.0) + out = self.orch("pick", "cloud") + self.assertIn("stop lane cloud: ledger (balance $3.00 < reserve $4.00)", out) + self.assertNotIn("in progress", self.status("t-slow")) + + def test_stop_max_parallel(self): + self.put_ledger(max_parallel=2, running=2) + out = self.orch("pick", "cloud") + self.assertIn("stop lane cloud: max parallel (2 running >= max_parallel 2)", out) + self.assertFalse((self.root / ".wf" / "cloud" / "t-slow.json").exists()) + + def test_id_unfit_refused(self): + out = self.orch("pick", "cloud", "--id", "t-gui") + self.assertIn("stop lane cloud: none fit (t-gui: needs a GUI)", out) + + def test_not_opted_in(self): + (self.root / "workflow.toml").write_text(GitCli.toml) + out = self.orch("pick", "cloud") + self.assertIn("stop lane cloud: none fit (project not opted in", out) + + def test_post_handback_commits_and_stops(self): + self.orch("pick", "cloud", "--id", "t-slow") + self.ok("note", "t-slow", "Recovery: cloud attempt x — red") # what wf cloud pull leaves behind + self.ok("status", "t-slow", "clear") + out = self.orch("post", "t-slow", "cloud", "--result", "handback") + self.assertIn("post: t-slow handback", out) + self.assertIn("committed leftover", out) + self.assertIn("stop lane cloud: handback", out) + self.assertIn(" cloud opus t-slow handback ", (self.root / "out" / "wf-orch.log").read_text()) |
