aboutsummaryrefslogtreecommitdiffziptar.gz
path: root/tests/test_cloud.py
diff options
context:
space:
mode:
Diffstat (limited to 'tests/test_cloud.py')
-rw-r--r--tests/test_cloud.py1174
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())