workflow

git clone https://git.godosa.eu/workflow

master

raw · 30712 bytes

  1
  2
  3
  4
  5
  6
  7
  8
  9
 10
 11
 12
 13
 14
 15
 16
 17
 18
 19
 20
 21
 22
 23
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
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:]))