From 81d4e80fd5aabe4e80f58e960affa795cf7d34ec Mon Sep 17 00:00:00 2001 From: godosa Date: Wed, 7 Oct 2026 07:27:17 +0200 Subject: workflow: initial public history --- wf_cloud.py | 720 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 720 insertions(+) create mode 100644 wf_cloud.py (limited to 'wf_cloud.py') 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:`, 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 / from the recorded master in a free lane worktree, + quick_gate, wf done -m " (cloud )" + wf merge; awaiting -> wf add -s awaiting + + status blocked; handback / red gate / bad patch -> note "Recovery: cloud attempt — ", + 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 ` 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 ' 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 / 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 → ': 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:])) -- cgit