aboutsummaryrefslogtreecommitdiffziptar.gz
path: root/wf_cloud.py
diff options
context:
space:
mode:
Diffstat (limited to 'wf_cloud.py')
-rw-r--r--wf_cloud.py720
1 files changed, 720 insertions, 0 deletions
diff --git a/wf_cloud.py b/wf_cloud.py
new file mode 100644
index 0000000..10f1b0a
--- /dev/null
+++ b/wf_cloud.py
@@ -0,0 +1,720 @@
+"""wf cloud — cloud-lane ledger IO. Pure logic: wflib/cloud.py.
+
+ ledger [--set-balance USD] [--budget USD] print budget/spent/balance/running; reconcile / change cap
+ send ID [--dry-run] [--project DIR] snapshot the project's master (+ cloud_include) into
+ out/cloud/ID, prompt from templates/cloud-prompt.md, reserve in the ledger (refused: exit 3),
+ `claude --cloud`, claim `in progress: cloud:<sid>`, record .wf/cloud/ID.json.
+ --dry-run: snapshot size + prompt, nothing sent / reserved / kept
+ archive SID | --ended archive ended sessions in the claude.ai app (POST
+ /v1/code/sessions/SID/archive with the CLI's login token, never printed); pull does it after each end
+ and prints `archive SID failed: …` on failure (pull's exit code unchanged); --ended retries the rest.
+ pull ID | --all [--export FILE] [--no-push] teleport + /export, parse the final WF-RESULT message:
+ done -> check sha256/bytes, gunzip, refuse paths outside the project / TASKS / archive / .wf / out /
+ cloud_include, `git am --3way` on <lane>/<id> from the recorded master in a free lane worktree,
+ quick_gate, wf done -m "<WF-REPORT> (cloud <sid>)" + wf merge; awaiting -> wf add -s awaiting +
+ status blocked; handback / red gate / bad patch -> note "Recovery: cloud attempt <sid> — <why>",
+ status clear, patch kept in out/cloud/ID.patch. No result: running (exit 4), > 24h: lost.
+ Every end charges the ledger, deletes the snapshot + record and archives the session.
+ Exit: 0 ended, 4 still running, 1 local error (worktree busy, branch exists: session kept, pull again).
+
+Library (for send/pull): send(folder, prompt) -> sid, export(folder, sid, out) -> path drive the real
+CLI under a pty (env WF_CLAUDE = claude binary, WF_CLAUDE_JSON = trust file, default ~/.claude.json).
+"""
+from __future__ import annotations
+
+import argparse
+import contextlib
+import datetime as dt
+import fcntl
+import os
+import pty
+import re
+import select
+import shutil
+import signal
+import struct
+import subprocess
+import sys
+import termios
+import time
+from pathlib import Path
+
+HERE = Path(__file__).resolve().parent
+sys.path.insert(0, str(HERE))
+
+from wflib import cloud # noqa: E402
+
+LEDGER, LOCK = "cloud.json", "cloud.lock"
+
+
+def state_dir() -> Path:
+ base = os.environ.get("XDG_STATE_HOME") or str(Path.home() / ".local/state")
+ return Path(base) / "wf"
+
+
+@contextlib.contextmanager
+def locked(state: Path):
+ """Exclusive flock; yields the ledger dict; atomic write-back when the block ends without error."""
+ state.mkdir(parents=True, exist_ok=True)
+ with open(state / LOCK, "w") as lock:
+ fcntl.flock(lock, fcntl.LOCK_EX)
+ path = state / LEDGER
+ led = cloud.loads(path.read_text() if path.exists() else "")
+ yield led
+ tmp = path.with_name(path.name + ".tmp")
+ tmp.write_text(cloud.dumps(led))
+ os.replace(tmp, path)
+
+
+# --- pty driver (spec §3 F1-F3, F7, F8; §4.1, §4.3) ----------------------------------------------
+
+CLAUDE = os.environ.get("WF_CLAUDE", "claude")
+DOWN, ENTER = "\x1b[B", "\r"
+CAP = cloud.CAP # bytes; tests lower it
+SETTLE = 2.0 # s: TUI settle after resume / between Enter retries while /export runs
+
+
+class Pty:
+ """A child process on a pseudo-terminal (claude --cloud refuses pipes, F2). `out` = raw output so far."""
+
+ def __init__(self, argv: list[str], cwd: Path, cols: int = 200, rows: int = 50):
+ self.master, slave = pty.openpty()
+ fcntl.ioctl(slave, termios.TIOCSWINSZ, struct.pack("HHHH", rows, cols, 0, 0))
+ env = {**os.environ, "TERM": os.environ.get("TERM") or "xterm-256color"}
+ self.proc = subprocess.Popen(argv, cwd=cwd, stdin=slave, stdout=slave, stderr=slave, env=env,
+ start_new_session=True, close_fds=True)
+ os.close(slave)
+ self.out = ""
+
+ def _read(self, wait: float) -> bool:
+ """Read what is there within `wait` s. False = child output closed."""
+ r, _, _ = select.select([self.master], [], [], wait)
+ if not r:
+ return True
+ try:
+ data = os.read(self.master, 65536)
+ except OSError: # EIO: child side closed
+ return False
+ if not data:
+ return False
+ self.out += data.decode("utf-8", "replace")
+ return True
+
+ def expect(self, pred, timeout: float) -> bool:
+ """Read until pred(plain text so far) or timeout / EOF. Returns pred's final truth."""
+ end = time.monotonic() + timeout
+ while not pred(cloud.strip_ansi(self.out)):
+ left = end - time.monotonic()
+ if left <= 0 or not self._read(min(left, 0.5)):
+ return bool(pred(cloud.strip_ansi(self.out)))
+ return True
+
+ def alive(self) -> bool:
+ return self.proc.poll() is None
+
+ def write(self, s: str) -> None:
+ os.write(self.master, s.encode())
+
+ def type(self, s: str) -> None:
+ """Text then Enter, Enter as its own write (a TUI reads a pasted CR as text)."""
+ self.write(s)
+ time.sleep(0.3)
+ self.write(ENTER)
+
+ def close(self, timeout: float = 10) -> int:
+ """Wait for exit (draining output), then TERM / KILL the process group. Returns the exit code."""
+ end = time.monotonic() + timeout
+ while self.alive() and time.monotonic() < end:
+ self._read(0.2)
+ for sig in (signal.SIGTERM, signal.SIGKILL):
+ if not self.alive():
+ break
+ with contextlib.suppress(ProcessLookupError):
+ os.killpg(self.proc.pid, sig)
+ with contextlib.suppress(subprocess.TimeoutExpired):
+ self.proc.wait(3)
+ with contextlib.suppress(OSError):
+ while self._read(0):
+ pass
+ os.close(self.master)
+ return self.proc.wait()
+
+
+def claude_json() -> Path:
+ return Path(os.environ.get("WF_CLAUDE_JSON") or Path.home() / ".claude.json")
+
+
+def pre_trust(folder: Path) -> None:
+ """Accept the trust dialog for folder in ~/.claude.json, atomic rewrite (F3)."""
+ path = claude_json()
+ text = path.read_text() if path.exists() else ""
+ new = cloud.trust(text, str(folder.resolve()))
+ tmp = path.with_name(f"{path.name}.wf-{os.getpid()}.tmp")
+ tmp.write_text(new)
+ if path.exists():
+ os.chmod(tmp, path.stat().st_mode & 0o777)
+ os.replace(tmp, path)
+
+
+def _answer_trust(p: Pty, plain: str, done: set) -> None:
+ """Fallback when the pre-accept did not take: Yes = Down, Enter (F3). Once per run."""
+ if "trust" not in done and cloud.TRUST_PROMPT.search(plain):
+ done.add("trust")
+ p.write(DOWN)
+ time.sleep(0.3)
+ p.write(ENTER)
+
+
+def send(folder: Path, prompt: str, timeout: float = 300) -> str:
+ """`claude --cloud <prompt>` in folder under a pty -> session id (F1). CloudError with the last 5 lines."""
+ folder = Path(folder)
+ pre_trust(folder)
+ p, seen = Pty([CLAUDE, "--cloud", prompt], folder), set()
+
+ def ready(plain: str) -> bool:
+ _answer_trust(p, plain, seen)
+ return cloud.parse_sid(plain) is not None and "Resume with:" in plain or not p.alive()
+
+ p.expect(ready, timeout)
+ code = p.close(5)
+ sid = cloud.parse_sid(p.out)
+ if not sid:
+ tail = " | ".join(cloud.last_lines(p.out)) or "(no output)"
+ raise cloud.CloudError(f"claude --cloud gave no session id (exit {code}): {tail}")
+ return sid
+
+
+def _teleport_export(folder: Path, sid: str, out: Path, timeout: float) -> tuple[bool, str]:
+ """One teleport -> /export -> /exit round. (exported?, raw output)."""
+ p, seen = Pty([CLAUDE, "--teleport", sid], folder), set()
+ end = time.monotonic() + timeout
+ try:
+ def resumed(plain: str) -> bool:
+ _answer_trust(p, plain, seen)
+ return bool(cloud.RESUMED.search(plain)) or not p.alive()
+
+ if not p.expect(resumed, timeout) or not p.alive():
+ return False, p.out
+ p.expect(lambda _: False, SETTLE) # let the TUI settle
+ p.type(f"/export {out}")
+ while not (out.exists() and out.stat().st_size) and p.alive() and time.monotonic() < end:
+ p.expect(lambda _: False, SETTLE)
+ if not (out.exists() and out.stat().st_size):
+ p.write(ENTER) # autocomplete eats the first Enter (F8)
+ ok = out.exists() and out.stat().st_size > 0
+ if p.alive():
+ p.type("/exit")
+ return ok, p.out
+ finally:
+ p.close(10)
+
+
+def export(folder: Path, sid: str, out: Path, timeout: float = 300) -> Path:
+ """Teleport into sid, `/export out`, `/exit`; no model call (F7, F8). Teleport may exit at
+ 'Checking out branch' without resuming: retried once, then CloudError with the last 5 lines."""
+ folder, out = Path(folder), Path(out).resolve()
+ pre_trust(folder)
+ out.parent.mkdir(parents=True, exist_ok=True)
+ raw = ""
+ for _ in range(2):
+ with contextlib.suppress(FileNotFoundError):
+ out.unlink()
+ ok, raw = _teleport_export(folder, sid, out, timeout)
+ if ok:
+ return out
+ tail = " | ".join(cloud.last_lines(raw)) or "(no output)"
+ raise cloud.CloudError(f"teleport {sid}: no export after 2 tries: {tail}")
+
+
+# --- send (spec §4.1, §4.2) ----------------------------------------------------------------------
+
+def git(cwd: Path, *args: str, check: bool = True) -> str:
+ r = subprocess.run(["git", "-C", str(cwd), *args], capture_output=True, text=True)
+ if check and r.returncode:
+ raise cloud.CloudError(f"git {args[0]}: {(r.stderr.strip() or r.stdout.strip() or 'failed').splitlines()[-1]}")
+ return r.stdout.strip()
+
+
+def snapshot(cfg, folder: Path) -> tuple[str, str, int]:
+ """folder = git archive of the project's master minus TASKS/archive/.wf/.worktrees/out, plus
+ cloud_include, as a fresh repo with one commit 'base <sha>' on main. -> (master sha, base sha, packed bytes)."""
+ from wflib import config
+ top = Path(git(cfg.root, "rev-parse", "--show-toplevel"))
+ master = config.git_branch(top / ".git") or "master"
+ sha = git(top, "rev-parse", master)
+ for inc in cfg.cloud_include:
+ if why := cloud.include_problem(inc):
+ raise cloud.CloudError(why)
+ if not (cfg.root / inc).exists():
+ raise cloud.CloudError(f"cloud_include '{inc}': not found in {cfg.root}")
+ if folder.exists():
+ shutil.rmtree(folder)
+ folder.mkdir(parents=True)
+ archive = subprocess.run(["git", "-C", str(top), "archive", "--format=tar", sha], capture_output=True)
+ if archive.returncode:
+ raise cloud.CloudError(f"git archive: {archive.stderr.decode(errors='replace').strip()}")
+ subprocess.run(["tar", "-x", "-C", str(folder)], input=archive.stdout, check=True)
+ rel = cfg.root.resolve().relative_to(top.resolve())
+ drop = [cfg.tasks, cfg.archive] + [cfg.root / n for n in cloud.NEVER]
+ for path in drop:
+ target = folder / Path(path).resolve().relative_to(top.resolve())
+ if target.is_dir() and not target.is_symlink():
+ shutil.rmtree(target)
+ elif target.exists() or target.is_symlink():
+ target.unlink()
+ for inc in cfg.cloud_include:
+ src, dst = cfg.root / inc, folder / rel / inc
+ dst.parent.mkdir(parents=True, exist_ok=True)
+ if src.is_dir():
+ shutil.copytree(src, dst, symlinks=True, dirs_exist_ok=True)
+ else:
+ shutil.copy2(src, dst)
+ who = ["-c", "user.name=wf", "-c", "user.email=wf@localhost", "-c", "commit.gpgsign=false"]
+ git(folder, "init", "-q", "-b", "main")
+ git(folder, "add", "-A", "-f")
+ git(folder, *who, "commit", "-q", "--allow-empty", "--no-verify", "-m", f"base {sha}")
+ git(folder, "repack", "-adq")
+ size = sum(p.stat().st_size for p in (folder / ".git" / "objects" / "pack").glob("*.pack"))
+ return sha, git(folder, "rev-parse", "HEAD"), size
+
+
+def credentials_path() -> Path:
+ """The CLI's login file (env WF_CLAUDE_CREDENTIALS, else $CLAUDE_CONFIG_DIR or ~/.claude)."""
+ if os.environ.get("WF_CLAUDE_CREDENTIALS"):
+ return Path(os.environ["WF_CLAUDE_CREDENTIALS"])
+ return Path(os.environ.get("CLAUDE_CONFIG_DIR") or Path.home() / ".claude") / ".credentials.json"
+
+
+def cli_version() -> str:
+ try:
+ out = subprocess.run([CLAUDE, "--version"], capture_output=True, text=True, timeout=20).stdout
+ except (OSError, subprocess.SubprocessError):
+ out = ""
+ m = re.match(r"\s*(\d+\.\d+\.\d+)", out)
+ return m[1] if m else "unknown"
+
+
+def post(url: str, headers: dict, timeout: float = 10) -> tuple[int, str]:
+ """POST {} -> (status, body); network failure -> CloudError."""
+ import urllib.error
+ import urllib.request
+ req = urllib.request.Request(url, data=b"{}", headers=headers, method="POST")
+ try:
+ with urllib.request.urlopen(req, timeout=timeout) as r:
+ return r.status, r.read(500).decode("utf-8", "replace")
+ except urllib.error.HTTPError as e:
+ return e.code, e.read(500).decode("utf-8", "replace")
+ except (urllib.error.URLError, OSError) as e:
+ raise cloud.CloudError(str(getattr(e, "reason", e)))
+
+
+def archive(state: Path, sid: str) -> str | None:
+ """Archive one session in the app (undocumented endpoint, spec §4.3): None ok (ledger row marked), else why."""
+ try:
+ try:
+ creds = credentials_path().read_text()
+ except OSError:
+ creds = ""
+ url, headers = cloud.archive_request(sid, creds, cli_version())
+ why = cloud.archive_problem(*post(url, headers))
+ except cloud.CloudError as e:
+ why = str(e)
+ if why is None:
+ with locked(state) as led:
+ with contextlib.suppress(cloud.CloudError):
+ cloud.find(led, sid)["archived"] = True
+ return why
+
+
+def cmd_archive(a, state: Path) -> int:
+ if a.ended:
+ with locked(state) as led:
+ sids = cloud.unarchived(led)
+ if not sids:
+ print("no ended cloud session left to archive")
+ return 0
+ else:
+ sids = [a.sid]
+ code = 0
+ for sid in sids:
+ if why := archive(state, sid):
+ print(f"wf: archive {sid} failed: {why}", file=sys.stderr)
+ code = 1
+ else:
+ print(f"archived {sid}")
+ return code
+
+
+def drop(folder: Path) -> None:
+ """Delete a snapshot folder and its parents out/cloud, out when that leaves them empty."""
+ shutil.rmtree(folder, ignore_errors=True)
+ for parent in (folder.parent, folder.parent.parent):
+ with contextlib.suppress(OSError):
+ parent.rmdir()
+
+
+def record_path(cfg, id: str) -> Path:
+ return cfg.root / ".wf" / "cloud" / f"{id}.json"
+
+
+def cmd_send(a, state: Path, now: dt.datetime) -> int:
+ import json
+ import wf
+ from wflib import lanes
+ p = wf.load_project(a)
+ cfg, id = p.cfg, a.id
+ item = p.doc.item(p.resolve(id))
+ id = item.id
+ if not cfg.cloud:
+ raise cloud.CloudError("project not opted in to the cloud lane (workflow.toml: cloud = true)")
+ rec = record_path(cfg, id)
+ if rec.exists() and not a.dry_run:
+ raise cloud.CloudError(f"{id} already sent ({json.loads(rec.read_text()).get('sid')}): wf cloud pull it first")
+ folder = cfg.root / "out" / "cloud" / (id + ".dry-run" if a.dry_run else id) # never the live snapshot
+ try:
+ sha, base, size = snapshot(cfg, folder)
+ if why := cloud.size_refusal(size, CAP):
+ raise cloud.CloudError(why)
+ text = "\n".join(item.lines())
+ recipe = "\n".join(wf.area_blocks(p, text)).strip()
+ prompt = cloud.fill((HERE / "templates" / "cloud-prompt.md").read_text(), id, base, text, recipe,
+ cfg.cloud_note)
+ if a.dry_run:
+ print(f"snapshot {size / 1e6:.1f} MB (master {sha[:12]}, base {base[:12]})\n\n{prompt}")
+ return 0
+ pending = f"pending:{id}:{os.getpid()}"
+ with locked(state) as led:
+ if why := cloud.refusal(led):
+ print(f"wf: cloud refused: {why}", file=sys.stderr)
+ drop(folder)
+ return 3
+ cloud.add(led, id, str(cfg.root), pending, cloud.MODEL, now)
+ try:
+ sid = send(folder, prompt)
+ except BaseException:
+ with locked(state) as led:
+ led["entries"] = [e for e in led["entries"] if e["sid"] != pending]
+ raise
+ with locked(state) as led:
+ cloud.find(led, pending)["sid"] = sid
+ except BaseException:
+ drop(folder)
+ raise
+ finally:
+ if a.dry_run:
+ drop(folder)
+ lane = lanes.lane_of(item, cfg.lanes, cfg.slice_above)
+ rec.parent.mkdir(parents=True, exist_ok=True)
+ if not (rec.parent.parent / ".gitignore").exists():
+ (rec.parent.parent / ".gitignore").write_text("*\n")
+ rec.write_text(json.dumps({"id": id, "sid": sid, "project": str(cfg.root), "lane": lane, "master": sha,
+ "base": base, "folder": str(folder), "sent": now.isoformat(timespec="seconds"),
+ "model": cloud.MODEL, "bytes": size}, indent=2) + "\n")
+ code = wf.main(["--project", str(cfg.root), "status", id, "progress", f"cloud:{sid}"])
+ print(f"sent {id}: {sid} ({size / 1e6:.1f} MB){'' if code == 0 else '; claim failed, see above'}")
+ return 0 if code == 0 else 1
+
+
+# --- pull (spec §4.3) ----------------------------------------------------------------------------
+
+class Busy(cloud.CloudError):
+ """Local obstacle (lane worktree / branch / setup): the session stays out, pull again later."""
+
+
+def age_text(td: dt.timedelta) -> str:
+ m = int(td.total_seconds() // 60)
+ return f"{m // 60}h{m % 60:02d}m" if m >= 60 else f"{m}m"
+
+
+def apply_patch(p, rec: dict, lane: str, patch: Path) -> tuple[Path, Path, str]:
+ """Lane worktree on a new branch <lane>/<id> at the recorded master, `git am --3way` the patch, worktree_setup.
+ -> (worktree, its project folder, branch). Busy: worktree/branch unusable; CloudError: am failed (undone)."""
+ import argparse
+ import wf
+ cfg, id = p.cfg, rec["id"]
+ main = Path(git(cfg.root, "rev-parse", "--show-toplevel"))
+ branch = f"{lane}/{id}"
+ if not subprocess.run(["git", "-C", str(main), "rev-parse", "--verify", "-q", f"refs/heads/{branch}"],
+ capture_output=True).returncode:
+ raise Busy(f"branch {branch} exists (local WIP?): merge or delete it, then pull again")
+ wt = wf.lane_worktree(main, p, lane, id)
+ if not wt.exists():
+ git(main, "worktree", "add", "-q", "--detach", str(wt), rec["master"])
+ git(wt, "switch", "-q", "-c", branch, rec["master"])
+ here = wt / cfg.root.resolve().relative_to(main.resolve())
+
+ def undo():
+ git(wt, "am", "--abort", check=False)
+ git(wt, "reset", "-q", "--hard", check=False)
+ git(wt, "switch", "-q", "--detach", rec["master"], check=False)
+ git(main, "branch", "-q", "-D", branch, check=False)
+
+ r = subprocess.run(["git", "-C", str(wt), "am", "-q", "--3way", str(patch)], capture_output=True, text=True)
+ if r.returncode:
+ undo()
+ tail = (r.stderr.strip() or r.stdout.strip() or "failed").splitlines()[-1]
+ raise cloud.CloudError(f"git am --3way: {tail}")
+ try:
+ wf.cmd_setup(argparse.Namespace(project=str(here)))
+ except wf.Failure as e:
+ undo()
+ raise Busy(str(e))
+ return wt, here, branch
+
+
+def pull_one(a, p, rec_path: Path, state: Path, now: dt.datetime) -> int:
+ """One sent task: export, parse, apply/gate/done/merge or awaiting/handback/lost; ends the session
+ (ledger charge, snapshot + record deleted). 0 ended, 4 running, 1 local error (session kept)."""
+ import argparse
+ import contextlib as cl
+ import io
+ import json
+ import wf
+ from wflib import lanes
+ cfg = p.cfg
+ rec = json.loads(rec_path.read_text())
+ id, sid = rec["id"], rec["sid"]
+ age = now - dt.datetime.fromisoformat(rec["sent"])
+ lane = rec.get("lane") or lanes.lane_of(p.doc.item(id), cfg.lanes, cfg.slice_above)
+ out_dir = cfg.root / "out" / "cloud"
+ kept = out_dir / f"{id}.patch"
+
+ def quiet(*argv):
+ buf = io.StringIO()
+ with cl.redirect_stdout(buf):
+ code = wf.main(["--project", str(cfg.root), *argv])
+ if code:
+ raise cloud.CloudError(f"wf {argv[0]} {id} failed (exit {code})")
+ return buf.getvalue()
+
+ def finish(state_: str, res=None, line=""):
+ with locked(state) as led:
+ try:
+ e = cloud.end(led, sid, state_, res.usage if res else None)
+ charge = f"${e['usd']:.2f} {e['usd_source']}"
+ except cloud.CloudError as err:
+ charge = f"no ledger charge: {err}"
+ drop(Path(rec["folder"]))
+ rec_path.unlink(missing_ok=True)
+ print(f"{id}: {state_} ({sid}, {charge}){': ' + line if line else ''}")
+ if why := archive(state, sid):
+ print(f"archive {sid} failed: {why}; archive by hand (wf cloud archive --ended)")
+ return 0
+
+ def keep(raw: bytes | None) -> str:
+ if not raw or not raw.strip():
+ return ""
+ out_dir.mkdir(parents=True, exist_ok=True)
+ kept.write_bytes(raw)
+ return f"; patch {cfg.rel(kept)}"
+
+ def handback(why: str, res=None, raw=None, state_="handback"):
+ extra = keep(raw)
+ quiet("note", id, f"Recovery: cloud attempt {sid} — {why}{extra}")
+ quiet("status", id, "clear")
+ return finish(state_, res, why + extra)
+
+ # 1. export (or a given file), F12 parse
+ exp = out_dir / f"{id}.export.txt"
+ try:
+ if a.export:
+ text = Path(a.export).read_text()
+ else:
+ text = export(Path(rec["folder"]), sid, exp).read_text()
+ except (cloud.CloudError, OSError) as e:
+ if age > cloud.LOST_AFTER:
+ return handback(f"lost: no export after {age_text(age)} ({e})", state_="lost")
+ raise
+ finally:
+ exp.unlink(missing_ok=True)
+ lines = cloud.final_message(text)
+ if lines is None and (kl := cloud.keyless_message(text)):
+ res, raw = cloud.keyless_result(kl), None
+ with cl.suppress(cloud.CloudError):
+ raw = cloud.decode_patch(res)
+ return handback("bad result: no WF-RESULT key", res, raw)
+ if lines is None:
+ if age > cloud.LOST_AFTER:
+ return handback(f"lost: no WF-RESULT after {age_text(age)}", state_="lost")
+ print(f"{id}: running ({sid}, sent {age_text(age)} ago)")
+ return 4
+ try:
+ res = cloud.parse_result(lines)
+ except cloud.CloudError as e:
+ return handback(f"bad result: {e}")
+ if res.usage and res.model and cloud.U.price(res.model):
+ with locked(state) as led:
+ with cl.suppress(cloud.CloudError):
+ cloud.find(led, sid)["model"] = res.model
+ raw, why = None, None
+ with cl.suppress(cloud.CloudError):
+ raw = cloud.decode_patch(res)
+
+ # 2. awaiting / handback from the session
+ if res.state == "awaiting":
+ extra = keep(raw)
+ header = quiet("add", "-s", "awaiting", f"{res.report} [[{id}]]")
+ m = re.search(r"\*\*(a-[^*]+)\*\*", header)
+ if not m:
+ raise cloud.CloudError(f"wf add printed no a-id: {header.strip()}")
+ quiet("status", id, "blocked", m[1])
+ if extra:
+ quiet("note", id, f"cloud attempt {sid}: partial work{extra}")
+ return finish("awaiting", res, m[1] + extra)
+ if res.state == "handback":
+ return handback(f"cloud handback: {res.report}", res, raw)
+
+ # 3. done: patch checks, apply, gate, done + merge
+ try:
+ raw = cloud.decode_patch(res)
+ except cloud.CloudError as e:
+ return handback(str(e), res)
+ patch = raw.decode("utf-8", errors="replace")
+ files = cloud.patch_files(patch)
+ main = Path(git(cfg.root, "rev-parse", "--show-toplevel")).resolve()
+ proj = str(cfg.root.resolve().relative_to(main))
+ never = [cfg.rel(cfg.tasks), cfg.rel(cfg.archive), *cloud.NEVER, *cfg.cloud_include]
+ if not files:
+ return handback("done with an empty patch", res, raw)
+ if why := cloud.patch_problem(files, "" if proj == "." else proj, never):
+ return handback(why, res, raw)
+ out_dir.mkdir(parents=True, exist_ok=True)
+ kept.write_bytes(raw)
+ try:
+ wt, here, branch = apply_patch(p, rec, lane, kept)
+ except Busy:
+ kept.unlink(missing_ok=True)
+ raise
+ except cloud.CloudError as e:
+ return handback(str(e), res, raw)
+ try:
+ if cfg.quick_gate:
+ wf._run_lines(cfg.quick_gate, here, cfg, "quick_gate '{line}' red (exit {rc})")
+ except wf.Failure as e:
+ git(wt, "reset", "-q", "--hard", check=False)
+ git(wt, "switch", "-q", "--detach", rec["master"], check=False)
+ git(main, "branch", "-q", "-D", branch, check=False)
+ return handback(str(e), res, raw)
+ kept.unlink(missing_ok=True)
+ sys.stdout.flush()
+ merged, err = "", None
+ with wf.project_lock(cfg.root):
+ wf.cmd_done(argparse.Namespace(project=str(cfg.root), ids=[id], m=f"{res.report} (cloud {sid})",
+ dry_run=False, finishing=True))
+ sys.stdout.flush()
+ ns = argparse.Namespace(project=str(here), m=None, no_push=a.no_push)
+ try:
+ wf.cmd_merge(ns)
+ merged = getattr(ns, "merged_sha", "")
+ except wf.Failure as e:
+ err = e
+ finish("done", res, f"merged {merged}" if merged else f"merge failed: {err}")
+ if err:
+ print(f"wf: {id} done, not merged: {err} (in {wt})", file=sys.stderr)
+ return 1
+ print(f"report: commit {merged}")
+ return 0
+
+
+@contextlib.contextmanager
+def pull_claim(rec: Path):
+ """Exclusive non-blocking flock on the task's record: yields False when another pull holds it or already
+ ended it (record gone), so two pulls of one id never race into 'branch … exists (local WIP?)'."""
+ try:
+ fd = os.open(rec, os.O_RDONLY)
+ except FileNotFoundError:
+ yield False
+ return
+ try:
+ try:
+ fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
+ except BlockingIOError:
+ yield False
+ return
+ yield rec.exists()
+ finally:
+ os.close(fd)
+
+
+def cmd_pull(a, state: Path, now: dt.datetime) -> int:
+ import wf
+ p = wf.load_project(a)
+ folder = p.cfg.root / ".wf" / "cloud"
+ if a.all:
+ recs = sorted(folder.glob("*.json"))
+ if not recs:
+ print("no task out in the cloud")
+ return 0
+ else:
+ id = p.resolve(a.id)
+ recs = [folder / f"{id}.json"]
+ if not recs[0].exists():
+ raise cloud.CloudError(f"{id} is not out in the cloud (no {p.cfg.rel(recs[0])})")
+ codes = []
+ for rec in recs:
+ try:
+ with pull_claim(rec) as mine:
+ if not mine: # a concurrent pull (batch sidecar + orchestrator) has it: not WIP, not an error
+ print(f"{rec.stem}: pulled by another wf cloud pull, skipped")
+ codes.append(0)
+ continue
+ codes.append(pull_one(a, wf.load_project(a), rec, state, now))
+ except (cloud.CloudError, wf.Failure) as e:
+ print(f"wf: {rec.stem}: {e}", file=sys.stderr)
+ codes.append(1)
+ sys.stdout.flush()
+ return 1 if 1 in codes else 4 if 4 in codes else 0
+
+
+def main(argv: list[str], state: Path | None = None, now: dt.datetime | None = None) -> int:
+ ap = argparse.ArgumentParser(prog="wf cloud", description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)
+ sub = ap.add_subparsers(dest="sub", required=True)
+ lp = sub.add_parser("ledger", help="print budget/spent/balance/running")
+ lp.add_argument("--set-balance", type=float, metavar="USD", help="balance shown on claude.ai: spent = budget - balance")
+ lp.add_argument("--budget", type=float, metavar="USD", help="change the cap")
+ sp = sub.add_parser("send", help="snapshot + prompt -> claude --cloud; ledger reserve, claim, record")
+ sp.add_argument("id")
+ sp.add_argument("--dry-run", action="store_true", help="print snapshot size + prompt; send nothing")
+ sp.add_argument("--project", metavar="DIR", help="project folder (default: from the working directory)")
+ pp = sub.add_parser("pull", help="export result -> apply patch, gate, done/merge | awaiting | handback; ledger charge; an id another pull holds → '<id>: pulled by another …, skipped' (rc 0)")
+ g = pp.add_mutually_exclusive_group(required=True)
+ g.add_argument("id", nargs="?")
+ g.add_argument("--all", action="store_true", help="every task of the project out in the cloud")
+ pp.add_argument("--project", metavar="DIR", help="project folder (default: from the working directory)")
+ pp.add_argument("--export", metavar="FILE", help="parse this /export file instead of teleporting (one id)")
+ pp.add_argument("--no-push", action="store_true", help="merge without git push home")
+ ar = sub.add_parser("archive", help="archive ended sessions in the app (pull does it; this retries)")
+ g = ar.add_mutually_exclusive_group(required=True)
+ g.add_argument("sid", nargs="?")
+ g.add_argument("--ended", action="store_true", help="every ended ledger session not archived yet")
+ a = ap.parse_args(argv)
+ now = now or dt.datetime.now().astimezone()
+ if a.sub == "archive":
+ return cmd_archive(a, state or state_dir())
+ if a.sub in ("send", "pull"):
+ import wf
+ from wflib import config, tasks
+ try:
+ if a.sub == "pull":
+ if a.export and a.all:
+ raise cloud.CloudError("--export goes with one id, not --all")
+ return cmd_pull(a, state or state_dir(), now)
+ return cmd_send(a, state or state_dir(), now)
+ except (cloud.CloudError, wf.Failure, tasks.TaskError, config.ConfigError) as e:
+ print(f"wf: {e}", file=sys.stderr)
+ return 1
+ try:
+ with locked(state or state_dir()) as led:
+ if a.budget is not None:
+ led["budget"] = a.budget
+ if a.set_balance is not None:
+ cloud.set_balance(led, a.set_balance, now)
+ print(cloud.summary(led))
+ except cloud.CloudError as e:
+ print(f"wf: {e}", file=sys.stderr)
+ return 1
+ return 0
+
+
+if __name__ == "__main__":
+ sys.exit(main(sys.argv[1:]))