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.py | 2239 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 2239 insertions(+) create mode 100755 wf.py (limited to 'wf.py') diff --git a/wf.py b/wf.py new file mode 100755 index 0000000..4b901ab --- /dev/null +++ b/wf.py @@ -0,0 +1,2239 @@ +#!/usr/bin/env python3 +"""wf — shared task workflow tool. Run inside a project (folder with workflow.toml). + +Look: next [--brief] · list [filters] · show ID… · ctx ID|PATH#ANCHOR + search WORDS… · log [-n N] [WORDS] · projects · check +Change: add "Title. Goal." -p N -e EFFORT · done ID… -m "entry" · prio ID N + move ID SECTION | --before/--after ID · status ID progress NOTE|blocked A-ID|clear + set ID --title/--effort/--after/--ref/--model/--sessions/--cloud · note ID "line" + body ID None: + """One-line notice for an option kept one more release under its old name.""" + print(f"wf: {old} is now {new}", file=sys.stderr) + + +# ------------------------------------------------------------------ files + +def stamp(path: Path) -> tuple[int, int]: + s = path.stat() + return (s.st_mtime_ns, s.st_size) + + +def write_if_unchanged(path: Path, text: str, read_stamp: tuple[int, int]) -> None: + """Replace `path` by temp file + rename, unless it changed since `read_stamp`.""" + if stamp(path) != read_stamp: + raise Failure(f"{path.name} changed on disk since it was read: nothing written, run again") + mode = path.stat().st_mode & 0o777 + fd, tmp = tempfile.mkstemp(dir=path.parent, prefix=f".{path.name}.", suffix=".tmp") + try: + with os.fdopen(fd, "wb") as f: + f.write(text.encode("utf-8")) + os.chmod(tmp, mode) + if stamp(path) != read_stamp: + raise Failure(f"{path.name} changed on disk since it was read: nothing written, run again") + os.replace(tmp, path) + finally: + if os.path.exists(tmp): + os.unlink(tmp) + + +def read(path: Path) -> str: + return path.read_bytes().decode("utf-8", errors="replace") + + +class Project: + def __init__(self, cfg: config.Config): + self.cfg = cfg + for name, path in (("tasks", cfg.tasks), ("archive", cfg.archive)): + if not path.is_file(): + raise Failure(f"workflow.toml: {name} '{cfg.rel(path)}' does not exist") + self.tasks_stamp, self.archive_stamp = stamp(cfg.tasks), stamp(cfg.archive) + self.tasks_text, self.archive_text = read(cfg.tasks), read(cfg.archive) + self.doc = tasks.parse(self.tasks_text) + self.new_archive: str | None = None + + @property + def archived(self) -> set[str]: + return tasks.archive_ids(self.new_archive or self.archive_text) + + def taken(self) -> set[str]: + return self.doc.ids() | self.archived + + def resolve(self, given: str) -> str: + """Exact id, or a unique prefix of an open id (read commands).""" + known = self.doc.ids() + if given in known: + return given + hits = sorted(i for i in known if i.startswith(given)) + if len(hits) == 1: + return hits[0] + if hits: + raise Failure(f"'{given}' matches {', '.join(hits)}") + raise Failure(tasks.unknown_id(given, known)) + + def save(self, dry_run: bool) -> None: + cfg = self.cfg + new_tasks = tasks.render(self.doc) + before = {p.key for p in checks.check(cfg, slow=False)[0]} + after = checks.check(cfg, tasks_text=new_tasks, archive_text=self.new_archive, slow=False)[0] + added = [p for p in after if p.key not in before] + if added: + raise Failure("refused: nothing written, the change adds problems:\n" + "\n".join(f" {p}" for p in added)) + changes = [(cfg.tasks, self.tasks_text, new_tasks, self.tasks_stamp)] + if self.new_archive is not None: + changes.insert(0, (cfg.archive, self.archive_text, self.new_archive, self.archive_stamp)) + if dry_run: + for path, old, new, _ in changes: + name = cfg.rel(path) + sys.stdout.writelines(difflib.unified_diff( + old.replace("\r\n", "\n").splitlines(keepends=True), + new.replace("\r\n", "\n").splitlines(keepends=True), name, f"{name} (new)")) + return + written = [] + for path, old, new, read_stamp in changes: + if new == old: + continue + try: + write_if_unchanged(path, new, read_stamp) + except Failure as e: + if written: + raise Failure(f"{e}; {written[0]} was already written: remove its new top line(s)") from None + raise + written.append(cfg.rel(path)) + + +def load_project(args, write: bool = False) -> Project: + start = Path(args.project) if args.project else Path.cwd() + if not start.is_dir(): + raise Failure(f"--project {start}: no such folder") + cfg = config.load_at(start) + if cfg.format > config.FORMAT: + raise Failure(f"project format {cfg.format} is newer than this wf ({config.FORMAT}): update /projects/public/workflow") + if write and cfg.format < config.FORMAT: + raise Failure(f"TASKS format {cfg.format}, wf needs {config.FORMAT}: " + "run wf migrate --write (idle project, one commit)") + return Project(cfg) + + +# ------------------------------------------------------------------ output + +def show(item: tasks.Item) -> str: + return "\n".join(item.lines()) + + +def status_word(item: tasks.Item) -> str: + if not item.status: + return "-" + return "blkd" if item.status.startswith("blocked") else "prog" + + +def row(item: tasks.Item, width: int, cfg: config.Config) -> str: + prio = "--" if item.prio is None else f"P{item.prio}" + model = item.model if item.prio is not None else "-" + lane = lanes.lane_of(item, cfg.lanes, cfg.slice_above) if item.prio is not None else "-" + lane_w = max(len(l.name) for l in cfg.lanes) + mark = f"[{item.sessions}] " if item.prio is not None and item.sessions != "parallel" else "" + line = (f"{item.id:<{width}} {prio} {item.effort or '-':<4} {status_word(item):<4} {model:<6} " + f"{lane:<{lane_w}} {mark}{item.title}") + return line if len(line) <= WIDTH else line[:WIDTH - 1] + "…" + + +def counts(doc: tasks.Doc) -> str: + def n(key): + return sum(len(s.items) for s in doc.sections if s.key == key) + return " · ".join(f"{key} {n(key)}" for key in ("pending", "human", "awaiting", "deferred")) + + +def resolve_refs(cfg: config.Config, item: tasks.Item) -> list[str]: + out = [] + for path, anchor in item.refs: + out.append(refs.resolve(cfg.root, path, anchor)) + if anchor and cfg.anchors_index and cfg.anchors_specs and cfg.anchors_specs.is_dir() \ + and (cfg.root / path).resolve() == cfg.anchors_index.resolve(): + for spec in sorted(cfg.anchors_specs.rglob("*.md")): + text = read(spec) + if any(anchor in refs.ANCHOR_RE.findall(l) for l in text.split("\n")): + out.append(refs.resolve(cfg.root, cfg.rel(spec), anchor)) + return out + + +def inbox_path() -> Path: + return Path(os.environ.get("WF_INBOX") or HERE / "inbox.md") + + +def inbox_count() -> int: + path = inbox_path() + if not path.is_file(): + return 0 + return sum(1 for l in read(path).split("\n") if l.startswith("- ")) + + +# ------------------------------------------------------------------ sessions (lanes) + +def sessions_dir(cfg: config.Config) -> Path: + return cfg.root / ".wf" / "sessions" + + +def alive(pid, socket) -> bool: + try: + os.kill(int(pid), 0) + except (OSError, ValueError, TypeError): + return False + return bool(socket) and Path(socket).exists() + + +def lane_names(cfg: config.Config) -> list[str]: + return [l.name for l in cfg.lanes] + + +def check_lane(cfg: config.Config, lane: str | None) -> None: + if lane and lane not in lane_names(cfg): + raise Usage(f"unknown lane '{lane}' ({', '.join(lane_names(cfg))})") + + +def sessions(cfg: config.Config) -> dict[str, dict]: + """Registered sessions keyed by lane name or "all" (no lane: every lane).""" + out = {} + folder = sessions_dir(cfg) + for name in [*lane_names(cfg), "all"]: + try: + s = json.loads((folder / f"{name}.json").read_text()) + except (OSError, ValueError): + continue + if isinstance(s, dict): + s["alive"] = alive(s.get("pid"), s.get("socket")) + out[name] = s + return out + + +def my_session() -> tuple[str, str] | None: + """(socket, pid) of this agent session from Claude Code's env; None outside one.""" + socket, pid = os.environ.get("CLAUDE_CODE_MESSAGING_SOCKET"), os.environ.get("CLAUDE_PID") + return (socket, pid) if socket and pid else None + + +def wf_folder(cfg: config.Config, name: str) -> Path: + """.wf/ in the project root, created; .wf is git-ignored.""" + folder = cfg.root / ".wf" / name + folder.mkdir(parents=True, exist_ok=True) + ignore = folder.parent / ".gitignore" + if not ignore.exists(): + ignore.write_text("*\n") + return folder + + +@contextlib.contextmanager +def project_lock(root: Path): + """Exclusive project lock (.wf/lock): every wf write and merge-back runs under it. Not reentrant.""" + folder = root / ".wf" + folder.mkdir(parents=True, exist_ok=True) + ignore = folder / ".gitignore" + if not ignore.exists(): + ignore.write_text("*\n") + with open(folder / "lock", "w") as f: + fcntl.flock(f, fcntl.LOCK_EX) + yield + + +def register(cfg: config.Config, lane: str, model: str) -> str | None: + """Record this agent session as the lane's session (lane "all" = every lane; env from Claude Code; + skipped without it). Another live session already holding the lane keeps it; returns a warning line then.""" + me = my_session() + if not me: + return None + socket, pid = me + other = sessions(cfg).get(lane) + if other and other["alive"] and str(other.get("pid")) != pid: + return (f"another live {lane} session holds this lane: uds:{other['socket']} " + "(claims keep tasks apart; tell the owner if unintended)") + try: + record = {"lane": lane, "model": model, "socket": socket, "pid": int(pid) if pid.isdigit() else pid, + "session": os.environ.get("CLAUDE_CODE_SESSION_ID", ""), + "at": datetime.datetime.now().isoformat(timespec="seconds")} + (wf_folder(cfg, "sessions") / f"{lane}.json").write_text(json.dumps(record) + "\n") + except OSError as e: + print(f"wf: session not registered: {e}", file=sys.stderr) + return None + + +def unregister(cfg: config.Config) -> None: + """Drop session records whose pid is this session (orchestrator holds no lane).""" + me = my_session() + if not me: + return + for name in [*lane_names(cfg), "all"]: + f = sessions_dir(cfg) / f"{name}.json" + try: + if str(json.loads(f.read_text()).get("pid")) == me[1]: + f.unlink() + except (OSError, ValueError, AttributeError): + pass + + +# ------------------------------------------------------------------ claims (wf status progress) + +def claim(cfg: config.Config, id: str, lane: str = "") -> None: + """Mark task id as held by this session (skipped outside one). lane: the task's lane.""" + me = my_session() + if not me: + return + socket, pid = me + model = next((s.get("model", name) for name, s in sessions(cfg).items() + if str(s.get("pid")) == pid and s.get("socket") == socket), "") + try: + record = {"id": id, "lane": lane, "model": model, "socket": socket, "pid": int(pid) if pid.isdigit() else pid, + "session": os.environ.get("CLAUDE_CODE_SESSION_ID", ""), + "at": datetime.datetime.now().isoformat(timespec="seconds")} + (wf_folder(cfg, "claims") / f"{id}.json").write_text(json.dumps(record) + "\n") + except OSError as e: + print(f"wf: claim not recorded: {e}", file=sys.stderr) + + +def unclaim(cfg: config.Config, ids) -> None: + for id in ids: + try: + (cfg.root / ".wf" / "claims" / f"{id}.json").unlink(missing_ok=True) + except OSError: + pass + + +def held(p: "Project") -> dict[str, str]: + """id → holder, for tasks in progress claimed by another live session.""" + me = my_session() + pending = {i.id: i for i in p.doc.section("pending").items} + out = {} + for path in sorted((p.cfg.root / ".wf" / "claims").glob("*.json")): + try: + c = json.loads(path.read_text()) + except (OSError, ValueError): + continue + if not isinstance(c, dict) or (me and str(c.get("pid")) == me[1]) or not alive(c.get("pid"), c.get("socket")): + continue + item = pending.get(path.stem) + if item is None or not (item.status or "").startswith("in progress"): + continue + out[item.id] = f"{c.get('model') + ' ' if c.get('model') else ''}session uds:{c.get('socket')}" + return out + + +def stale(p: "Project") -> list: + """Items in progress whose claim is missing or held by a dead session.""" + out = [] + for item in p.doc.section("pending").items: + if not (item.status or "").startswith("in progress"): + continue + try: + c = json.loads((p.cfg.root / ".wf" / "claims" / f"{item.id}.json").read_text()) + except (OSError, ValueError): + c = None + if not isinstance(c, dict) or not alive(c.get("pid"), c.get("socket")): + out.append(item) + return out + + +def other_live(cfg: config.Config) -> int: + """Live sessions in the project other than this one.""" + me = my_session() + return sum(1 for s in sessions(cfg).values() if s["alive"] and not (me and str(s.get("pid")) == me[1])) + + +def multi_lines(p: Project, args, lane: str, extra: int) -> list[str]: + """Worktree mode advice when more than one session is live in the project (git only).""" + n = sum(1 for s in sessions(p.cfg).values() if s["alive"]) + extra + main_top = config.git_top(p.cfg.root) + if n < 2 or not main_top or not (main_top / ".git").is_dir(): + return [] + head = f"{n} live sessions here: " + wt = config.linked_worktree(Path(args.project or Path.cwd()).resolve()) + if wt: + return [head + f"you are in worktree {p.cfg.rel(wt[0])} (branch {config.git_branch(wt[1])}); " + "wf writes the main tree's TASKS.md"] + folder, master = p.cfg.root / ".worktrees" / lane, config.git_branch(main_top / ".git") or "master" + cmd = (f"cd {p.cfg.rel(folder)} && git switch -c {lane}/ {master}" if folder.is_dir() + else f"git worktree add {p.cfg.rel(folder)} -b {lane}/ {master}") + if p.cfg.worktree_setup: + cmd += " && wf setup" if folder.is_dir() else f" && cd {p.cfg.rel(folder)} && wf setup" + return [head + f"work in your lane's worktree, never on {master}:", f" {cmd} (wf there writes this TASKS.md)"] + + +def lane_lines(p: Project, me: str | None, runner: bool = False) -> list[str]: + a = (p.doc, p.archived, p.cfg.lanes, p.cfg.slice_above, held(p)) + return lanes.lanes_block(lanes.counts(*a, runner=runner), sessions(p.cfg), me, + lanes.not_ready(*a) if runner else None) + + +# ------------------------------------------------------------------ read commands + +def cmd_next(args) -> int: + p = load_project(args) + lane = args.lane + check_lane(p.cfg, lane) + model = args.as_ or "haiku" + out: list[str] = [] + warning = None + if not args.as_: + out.append("no --as: treated as haiku; pass --as " + "|".join(tasks.MODELS)) + elif warning := register(p.cfg, lane or "all", model): + out.append(warning) + claims = held(p) + if solo := tasks.solo_running(p.doc, claims): + if out: + print("\n".join(out)) + raise Failure(f"solo {solo[0]} in progress by {solo[1]}: wait (its wf done notifies you)") + item, skipped = lanes.pick(p.doc, p.archived, p.cfg.lanes, p.cfg.slice_above, lane, model, claims, + other_live(p.cfg), args.owner) + multi = multi_lines(p, args, lane or "all", 1 if warning else 0) + if not args.brief: + awaiting = p.doc.section("awaiting").items + if awaiting: + out += ["===== Awaiting your decision (mention, don't block) =====", + *(f"- {i.id}: {i.text}" for i in awaiting), ""] + human = p.doc.section("human").items + if human: + out += ["===== Needs human (not picked) =====", *(f"- {i.id}: {i.title}" for i in human), ""] + flight = ledgers.in_flight(p.cfg.root, p.cfg.ledgers) if p.cfg.ledgers and p.cfg.ledgers.is_dir() else [] + if flight: + out += ["===== In-flight plans =====", *flight, ""] + if skipped: + out += ["===== Skipped =====", *(f"- {i.id}: {why}" for i, why in skipped), ""] + if multi: + out += ["===== Multi-session =====", *multi, ""] + lines = lane_lines(p, lane) + if len(lines) > 1: + out += ["===== Lanes =====", *lines, ""] + waits = lanes.waiting_block(lanes.cross_waits(p.doc, p.archived, lane, p.cfg.lanes, p.cfg.slice_above), + sessions(p.cfg), p.cfg.lanes, p.cfg.slice_above) if lane else [] + if waits: + out += ["===== Waiting on other lanes =====", *waits, ""] + if item: + if not args.brief: + out.append("===== Next task =====") + out.append(show(item)) + if lanes.is_slice_job(item, p.cfg.slice_above): + out += ["", "===== Slice job (no code) =====", + f"effort {item.effort} > slice_above {p.cfg.slice_above}: split it, don't implement.", + f' wf add ". <goal>" -e <1h|1h> --parent {item.id} --model haiku|sonnet|opus' + " (then wf body: Steps/Done/Ref)", + f' wf note {item.id} "sliced into …" · wf status {item.id} clear · not wf done' + " (the parent waits on its slices)"] + if not args.brief: + for block in resolve_refs(p.cfg, item): + out += ["", block] + n = 0 if args.brief else inbox_count() + if n: + out += ["", f"workflow inbox: {n} reports (triage: workflow session)"] + if hint := ctx_hint_line(p.cfg): + out += ["", hint] + while out and not out[-1]: + out.pop() + if out: + print("\n".join(out)) + if not item: + raise Failure(f"nothing pickable for {lane or 'all lanes'} ({model}) in Pending") + return 0 + + +def cmd_areas(args) -> int: + p = load_project(args, write=bool(args.mark)) + file = p.cfg.areas_file + rel = p.cfg.areas_rel + text = file.read_text(encoding="utf-8", errors="replace") if file.is_file() else "" + found = areas_mod.parse(text) + if not found: + print(f"no areas in {rel}") + return 0 + names = ", ".join(a.name for a in found) + wanted = args.name or args.mark + if wanted and wanted not in [a.name for a in found]: + hits = [a.name for a in found if a.name.startswith(wanted)] + if len(hits) == 1: + if args.name: + args.name = hits[0] + else: + args.mark = hits[0] + wanted = hits[0] + elif hits: + raise Usage(f"ambiguous area '{wanted}' in {rel}: " + ", ".join(hits)) + else: + raise Usage(f"no area '{wanted}' in {rel} ({names})") + if args.mark: + r = git_run(p.cfg.code, "rev-parse", "--short", "HEAD") + if r.returncode != 0: + raise Failure("git rev-parse HEAD failed: " + r.stderr.strip()) + text = areas_mod.mark(text, args.mark, r.stdout.strip()) + file.write_text(text, encoding="utf-8") + found = areas_mod.parse(text) + for a in found: + if wanted and a.name != wanted: + continue + stale, miss, commits = areas_mod.status(p.cfg.code, a, p.cfg.area_stale_commits) + line = " · ".join([f"{a.name}: " + ("stale" if stale else "ok")] + + ([f"missing: {', '.join(miss)}"] if miss else []) + + ([f"{commits} commits since {a.checked}"] if commits is not None else ["no Checked"] if not a.checked else [])) + print(line) + if not wanted: + for folder, n in uncovered_areas(p, args): + print(f"uncovered: {folder} ({n} file{'s' * (n != 1)}): no area's Paths covers it") + return 0 + + +def cmd_lanes(args) -> int: + p = load_project(args) + check_lane(p.cfg, args.lane) + if args.unregister: + unregister(p.cfg) + elif args.lane or args.as_: + register(p.cfg, args.lane or "all", args.as_ or "haiku") + if args.wait is not None: + import time + poll = float(os.environ.get("WF_LANES_POLL", "30")) + end = time.monotonic() + args.wait + while True: + p = load_project(args) + if (p.cfg.root / "out" / "wf-batch.stop").exists(): + print("stop requested: out/wf-batch.stop") + return 2 + if any(c[0] for c in lanes.counts(p.doc, p.archived, p.cfg.lanes, p.cfg.slice_above, held(p), runner=True).values()): + break + left = end - time.monotonic() + if left <= 0: + print("\n".join(lane_lines(p, args.lane, True))) + return 1 + time.sleep(min(poll, left)) + print("\n".join(lane_lines(p, args.lane, True))) + return 0 + + +def runner_ids(p: Project) -> set[str]: + """Pending ids a headless worker can take: not blocked, After done, runner-ready, no open slices, not sliced out.""" + return {i.id for i in p.doc.section("pending").items + if not i.error and not (i.status or "").startswith("blocked") and all(a in p.archived for a in i.after) + and i.runner_ready and not tasks.open_slices(p.doc, i.id) and not lanes.sliced_out(i, p.cfg.slice_above)} + + +def cmd_list(args) -> int: + p = load_project(args) + check_lane(p.cfg, args.lane) + keys = list(tasks.SECTIONS) if args.section == "all" else [args.section] + ready = None + if args.ready: + ready = {i.id for i in p.doc.section("pending").items + if not i.error and not (i.status or "").startswith("blocked") + and all(a in p.archived for a in i.after)} + if args.runner: + ready = runner_ids(p) + stale_ids = {i.id for i in stale(p)} if args.stale else None + want_ref = tuple(args.ref.split("#", 1)) if args.ref else None + + def keep(i: tasks.Item) -> bool: + if args.prio is not None and (i.prio is None or i.prio > args.prio): + return False + if ready is not None and i.id not in ready: + return False + if args.blocked and status_word(i) != "blkd": + return False + if args.progress and status_word(i) != "prog": + return False + if stale_ids is not None and i.id not in stale_ids: + return False + if args.model and (i.prio is None or i.model != args.model): + return False + if args.lane and (i.prio is None or lanes.lane_of(i, p.cfg.lanes, p.cfg.slice_above) != args.lane): + return False + if want_ref and not any(path == want_ref[0] and (len(want_ref) == 1 or anchor == want_ref[1]) + for path, anchor in i.refs): + return False + return True + + groups = [(s, [i for i in s.items if keep(i)]) for k in keys for s in p.doc.sections if s.key == k] + if args.runner and args.lane: # pick order of the lane, then its fallback lane's + order = lanes.ranked(p.doc, p.archived, p.cfg.lanes, p.cfg.slice_above, args.lane) + keys, groups = ["pending"], [(None, [i for i in order if i.id in ready + and (args.prio is None or i.prio <= args.prio)])] + if args.n is not None: + left = args.n + for n, (s, items) in enumerate(groups): + groups[n] = (s, items[:left]) + left -= len(groups[n][1]) + width = max((len(i.id) for _, items in groups for i in items), default=0) + for s, items in groups: + if len(keys) > 1: + print(f"== {s.heading}") + for i in items: + print(row(i, width, p.cfg)) + print(counts(p.doc)) + return 0 + + +def cmd_show(args) -> int: + p = load_project(args) + print("\n\n".join(show(p.doc.item(p.resolve(i))) for i in args.ids)) + return 0 + + +def cmd_ctx(args) -> int: + p = load_project(args) + target = args.target + if "#" in target or "/" in target or (p.cfg.root / target).exists(): + path, _, anchor = target.partition("#") + users = [i.id for i in p.doc.all_items() + if any(r == path and (not anchor or a == anchor) for r, a in i.refs)] + print(refs.resolve(p.cfg.root, path, anchor or None)) + print("\nTasks: " + (", ".join(users) if users else "none")) + return 0 + if target not in p.doc.ids() and target in p.archived: + line = next(l for l in p.archive_text.replace("\r\n", "\n").split("\n") + if (m := tasks.ARCHIVE_ID_RE.match(l)) and m.group(1) == target) + print(f"done: {line}") + return 0 + id = p.resolve(target) + section, _ = p.doc.find(id) + item = p.doc.item(id) + where = {i.id: s.heading for s in p.doc.sections for i in s.items} + out = [show(item), "", f"Section: {section.heading}"] + if item.after: + out.append("After: " + ", ".join( + f"{a} ({'open, ' + where[a] if a in where else 'done' if a in p.archived else 'unknown'})" + for a in item.after)) + needed = [i.id for i in p.doc.all_items() if id in i.after] + blocks = [i.id for i in p.doc.all_items() if i.blocked_on == id] + linked = [i.id for i in p.doc.all_items() + if i.id != id and i.id not in needed and i.id not in blocks + and id in refs.links("\n".join(i.lines()))] + for label, ids in (("Needed by", needed), ("Blocks", blocks), ("Linked from", linked)): + if ids: + out.append(f"{label}: {', '.join(ids)}") + if p.cfg.verify and not id.startswith("a-"): + out.append("Verify (run before wf finish): " + " · ".join(p.cfg.verify)) + for block in resolve_refs(p.cfg, item): + out += ["", block] + out += area_blocks(p, "\n".join(item.lines())) + print("\n".join(out)) + return 0 + + +def area_blocks(p: Project, text: str) -> list[str]: + """Notes of the areas the task text names (area name, anchor or a path no other area lists) → ctx lines.""" + file = p.cfg.areas_file + if not file.is_file(): + return [] + notes = file.read_text(encoding="utf-8", errors="replace") + out = [] + found = areas_mod.parse(notes) + skip = areas_mod.shared_paths(found) + for a in found: + if areas_mod.matches(a, text, skip): + out += ["", f"Area {a.name} ({p.cfg.areas_rel}):", areas_mod.block(notes, a.name)] + return out + + +def cmd_search(args) -> int: + p = load_project(args) + kinds = {k for k, on in (("task", args.tasks), ("archive", args.archive), ("doc", args.docs)) if on} or None + docs = {} + if kinds is None or "doc" in kinds: + docs = {p.cfg.rel(f): read(f) for f in checks.doc_files(p.cfg)} + hits = search.search(args.words, p.doc, p.archive_text, docs, kinds=kinds, limit=args.n, + archive_name=p.cfg.rel(p.cfg.archive)) + if not hits: + raise Failure("no hits") + for h in hits: + if h.kind == "task": + header = p.doc.item(h.where).text if h.where in p.doc.ids() else "" + text = h.line if h.line == header else f"{h.label} — {h.line}" + elif h.kind == "archive": + text = re.sub(r"^\d{4}-\d\d-\d\d \*\*([^*]+)\*\*", r"\1", h.line) + else: + text = f"{h.label} — {h.line}" if h.line and h.line != h.label else h.label + line = f"{h.where} {text}" + print(line if len(line) <= 2 * WIDTH else line[:2 * WIDTH - 1] + "…") + return 0 + + +def cmd_log(args) -> int: + p = load_project(args) + lines = [l for l in p.archive_text.replace("\r\n", "\n").split("\n") if l.startswith("- ")] + words = [w.lower() for w in args.words] + lines = [l for l in lines if all(w in l.lower() for w in words)] + print("\n".join(lines[:args.n])) + return 0 + + +def cmd_check(args) -> int: + p = load_project(args) + errors, warnings = checks.check(p.cfg) + warnings += checks.merged_tool_worktrees(Path(os.environ.get("WF_TOOL_ROOT") or HERE)) + for w in warnings: + print(f"warn: {w}") + for e in errors: + print(f"ERROR: {e}") + print(("OK: " if not errors else "") + f"{len(errors)} errors · {len(warnings)} warnings") + return 1 if errors else 0 + + +def find_projects(root: Path, depth: int = 3) -> list[Path]: + out = [] + + def walk(folder: Path, left: int) -> None: + if (folder / "wf.py").is_file() and (folder / "wflib").is_dir(): + return # a wf checkout (its templates/ is no project) + if (folder / config.NAME).is_file(): + out.append(folder) + return + if left == 0: + return + try: + children = sorted(c for c in folder.iterdir() if c.is_dir() and not c.is_symlink()) + except OSError: + return + for c in children: + if not c.name.startswith(".") and c.name not in SKIP_DIRS: + walk(c, left - 1) + + walk(root, depth) + return out + + +def default_root(tool: Path) -> Path: + """Folder `wf projects` scans: the tool's parent, else its grandparent (tool under e.g. public/), if it holds projects.""" + return next((d for d in (tool.parent, tool.parent.parent) if find_projects(d)), tool.parent) + + +def projects_root() -> Path: + return Path(os.environ.get("WF_ROOT") or default_root(HERE)) + + +def cmd_projects(args) -> int: + root = projects_root() + rows = [] + for folder in find_projects(root): + name = str(folder.relative_to(root)) + try: + cfg = config.load(folder) + doc = tasks.parse(read(cfg.tasks)) + archived = tasks.archive_ids(read(cfg.archive)) if cfg.archive.is_file() else set() + errors = len(checks.check(cfg, slow=False)[0]) + item = lanes.pick(doc, archived, cfg.lanes, cfg.slice_above, owner=True)[0] + n = {k: sum(len(s.items) for s in doc.sections if s.key == k) for k in tasks.SECTIONS} + rows.append((name, f"pending {n['pending']} human {n['human']} awaiting {n['awaiting']} " + f"errors {errors} next: " + (f"{item.id} {item.title}" if item else "-"))) + except (config.ConfigError, tasks.TaskError, OSError) as e: + rows.append((name, f"broken: {e}")) + width = max((len(n) for n, _ in rows), default=0) + for name, text in rows: + line = f"{name:<{width}} {text}" + print(line if len(line) <= WIDTH else line[:WIDTH - 1] + "…") + if not rows: + raise Failure(f"no projects under {root}") + return 0 + + +def claude_projects() -> Path: + return Path(os.environ.get("CLAUDE_CONFIG_DIR") or Path.home() / ".claude") / "projects" + + +def since_time(value: str | None) -> str | None: + """ISO (UTC) as is; 90m / 2h / 1d = that long ago.""" + m = re.fullmatch(r"(\d+)([mhd])", value or "") + if not m: + return value + delta = datetime.timedelta(**{{"m": "minutes", "h": "hours", "d": "days"}[m[2]]: int(m[1])}) + return (datetime.datetime.now(datetime.timezone.utc) - delta).strftime("%Y-%m-%dT%H:%M:%S") + + +def agent_label(path: Path) -> str: + id = path.stem.removeprefix("agent-") + try: + desc = json.loads(path.with_name(path.stem + ".meta.json").read_text()).get("description") or "" + except (OSError, ValueError, AttributeError): + desc = "" + return f"{id} {desc}".strip() + + +def usage_transcripts(args) -> list[tuple[str, Path]]: + """(label, jsonl): one agent, or a session's main transcript + its subagents.""" + base = claude_projects() + if args.agent: + id = args.agent.removeprefix("agent-") + found = sorted(base.glob(f"*/*/subagents/agent-{id}.jsonl")) + if not found: + raise Failure(f"no transcript for agent {args.agent}") + return [(agent_label(found[0]), found[0])] + sid = args.session or os.environ.get("CLAUDE_CODE_SESSION_ID") + if sid: + found = sorted(base.glob(f"*/{sid}*.jsonl")) + if len(found) != 1: + raise Failure(f"session {sid}: {len(found)} transcripts match") + main = found[0] + else: + folder = base / re.sub(r"[^A-Za-z0-9]", "-", str(Path.cwd())) + found = sorted(folder.glob("*.jsonl"), key=lambda p: p.stat().st_mtime) + if not found: + raise Failure(f"no transcripts in {folder} (--session ID)") + main = found[-1] + agents = sorted((main.parent / main.stem / "subagents").glob("agent-*.jsonl")) + return [("main", main)] + [(agent_label(p), p) for p in agents] + + +def cmd_usage(args) -> int: + since = since_time(args.since) + if args.report: + return usage_report(since) + if args.explore: + return usage_explore(args, since) + if args.log and not args.agent: + raise Usage("--log needs --agent") + if args.log: + root = config.find_root(Path(args.project) if args.project else Path.cwd()) + path = usage_transcripts(args)[0][1] + now = datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + line = usage.log_line(now, root.name, args.log[0], args.effort or "-", args.log[1], + path.stem.removeprefix("agent-"), usage.parse(read(path), since=since), lane=args.lane, dur=args.duration) + (root / "out").mkdir(exist_ok=True) + with open(root / "out" / "wf-cost.log", "a") as f: + f.write(line + "\n") + print(line) + return 0 + rows, total, dollars, unknown = [], usage.Usage(), 0.0, False + for label, path in usage_transcripts(args): + for model, u in usage.parse(read(path), since=since).items(): + c = usage.cost(model, u) + total.add(u) + dollars += c or 0 + unknown |= c is None + rows.append((label, usage.short(model), u, "?" if c is None else f"{c:.2f}")) + width = min(max([len(r[0]) for r in rows] + [5]), 40) + + def nums(u): + out = ("~" if u.est else "") + usage.fmt(u.out) # ~ = estimated from content (subagent) + return " ".join(f"{usage.fmt(n):>7}" for n in (u.inp, u.cw, u.cr)) + f" {out:>7}" + + print(f"{'agent':<{width}} {'model':<12} {'turns':>6} {'in':>7} {'cw':>7} {'cr':>7} {'out':>7} {'$':>7}") + for label, model, u, c in rows: + label = label if len(label) <= width else label[:width - 1] + "…" + print(f"{label:<{width}} {model:<12} {u.turns:>6} {nums(u)} {c:>7}") + print(f"{'total':<{width}} {'':<12} {total.turns:>6} {nums(total)} {dollars:>7.2f}{'+' if unknown else ''}") + return 0 + + +def usage_explore(args, since: str | None) -> int: + root = str(Path.cwd()) + agents = [(label, usage.explore(read(path), since=since, root=root)) + for label, path in usage_transcripts(args) if label != "main"] + r = usage.explore_report(agents) + if not r["rows"]: + print("no subagent with >= 4 calls") + return 0 + width = min(max([len(x[0]) for x in r["rows"]] + [5]), 40) + print(f"{'agent':<{width}} {'calls':>5} {'1stEdit':>7} {'explore':>7} {'pre-edit':>8}") + for label, calls, fe, exp, pre in r["rows"]: + label = label if len(label) <= width else label[:width - 1] + "…" + print(f"{label:<{width}} {calls:>5} {fe:>7} {exp:>7.0%} {pre:>8.0%}") + tot = r["total"] or 1 + print(f"ALL {r['agents']} agents, {usage.fmt(r['total'])} billed input: " + + ", ".join(f"{k} {100 * v // tot}%" for k, v in sorted(r["cost"].items(), key=lambda kv: -kv[1]))) + if r["pre_n"]: + print(f"before 1st edit (median, {r['pre_n']} agents): {r['pre_share']:.0%} of input, " + f"{r['pre_calls']:g} calls, +{usage.fmt(int(r['pre_ctx']))} context") + rt = sum(r["results"].values()) or 1 + print("result tokens: " + ", ".join(f"{k} {usage.fmt(v)} {100 * v // rt}%" + for k, v in sorted(r["results"].items(), key=lambda kv: -kv[1]))) + if r["files"]: + print("files read by >= 2 agents (agents, tokens):") + for f, a, n in r["files"]: + print(f"{a:>4} {usage.fmt(n):>7} {f}") + return 0 + + +def fmt_dur(s) -> str: + return "-" if s is None else f"{int(s) // 3600}h{int(s) % 3600 // 60:02d}m" if s >= 3600 else f"{int(s) // 60}m{int(s) % 60:02d}s" + + +def usage_report(since: str | None) -> int: + root = projects_root() + entries = [] + for folder in find_projects(root): + log = folder / "out" / "wf-cost.log" + if log.is_file(): + entries += usage.parse_log(read(log), since=since) + if not entries: + raise Failure(f"no out/wf-cost.log lines in the projects under {root}") + print(f"{'lane':<8} {'effort':<6} {'n':>4} {'done':>4} {'other':>5} {'med$':>7} {'total$':>8} {'$/done':>7} med_turns med_dur") + for lane, effort, n, done, other, med, total, per, turns, dur in usage.report(entries, tasks.EFFORTS): + per = "-" if per is None else f"{per:.2f}" + print(f"{lane:<8} {effort:<6} {n:>4} {done:>4} {other:>5} {med:>7.2f} {total:>8.2f} {per:>7} {turns:>9} {fmt_dur(dur):>8}") + return 0 + + +# ------------------------------------------------------------------ write commands + +def split_list(value: str) -> list[str]: + return tasks.split_refs(value) + + +def stdin_lines() -> list[str]: + return sys.stdin.read().replace("\r\n", "\n").split("\n") + + +def set_slices_line(parent: tasks.Item, id: str) -> None: + at = parent._line(tasks.SLICES_RE) + if at is None: + parent.body.insert(tasks._tail_start(parent), f"{tasks.INDENT}- Slices: [[{id}]]") + else: + parent.body[at] = parent.body[at].rstrip() + f", [[{id}]]" + + +def cmd_add(args) -> int: + p = load_project(args, write=True) + doc = p.doc + if args.title == "-": + item = tasks.parse_block(sys.stdin.read()) + key = args.section or ("awaiting" if item.id.startswith("a-") else "pending") + else: + key = args.section or "pending" + prio, after = args.prio, split_list(args.after or "") + text = " ".join(args.title.split()) + leading = tasks.LEADING_ID_RE.match(text) + if leading and not args.id: + args.id, text = leading.groups() + if args.parent: + parent_section, _ = doc.find(args.parent) + parent = doc.item(args.parent) + key = args.section or parent_section.key + prio = parent.prio if prio is None else prio + id = args.id or tasks.slice_id(args.parent, p.taken()) + m = re.match(re.escape(args.parent) + r"-(\d+)$", id) + prev = parent.slices or ([f"{args.parent}-{int(m.group(1)) - 1}"] if m and int(m.group(1)) > 1 else []) + if not after: + # a deferred previous slice would block a live one forever: chain to the last live slice + live = [s for s in prev if s in p.archived or key == "deferred" + or (s in doc.ids() and doc.find(s)[0].key != "deferred")] + after = live[-1:] + task = key != "awaiting" + if task and prio is None: + raise Usage("add: a task needs -p 0-3") + if task and args.effort is None: + raise Usage("add: a task needs -e " + "|".join(tasks.EFFORTS)) + if task and args.effort not in tasks.EFFORTS: + raise Failure(f"effort '{args.effort}' (want {', '.join(tasks.EFFORTS)})") + if not text: + raise Usage("add: empty title") + if not text.endswith((".", "?", "!")): + text += "." + if not args.parent: + title = text.split(". ", 1)[0] + id = args.id or tasks.make_id(title, p.taken(), "t-" if task else "a-") + if not tasks.ID_RE.match(id): + raise Failure(f"bad id '{id}' (want t-… or a-…, lowercase a-z 0-9 -)") + if id in p.archived: + raise Failure(f"id '{id}' already used in the archive (ids are never reused)") + item = tasks.Item(id=id, prio=prio if task else None, effort=args.effort if task else None, text=text) + if args.body: + item.body = tasks.indent_body(stdin_lines()) + if args.interactive: + old_spelling("--interactive", "--sessions owner") + args.sessions = args.sessions or "owner" + if args.done and task: + item.body.append(f"{tasks.INDENT}Done: {args.done}") + if args.sessions and args.sessions != "parallel" and task: + item.body.append(f"{tasks.INDENT}Sessions: {args.sessions}") + if args.model and task: + item.body.append(f"{tasks.INDENT}Model: {args.model}") + if args.cloud and task: + item.body.append(f"{tasks.INDENT}Cloud: {args.cloud}") + if after: + item.body.append(f"{tasks.INDENT}- After: " + ", ".join(f"[[{a}]]" for a in after)) + if args.ref: + item.body.append(f"{tasks.INDENT}Ref: " + ", ".join(split_list(args.ref))) + if args.parent: + set_slices_line(parent, id) + tasks.insert(doc, item, key) + p.save(args.dry_run) + if not args.dry_run: + print(item.header()) + if item.prio is not None and item.done_text is None: + if item.prio == 0: + print("wf: warning: P0 without Done line is not runner-pickable; use --done \"<text>\"", file=sys.stderr) + else: + print("hint: no Done line; add one (--done) so runners can pick it", file=sys.stderr) + return 0 + + +RELEARN_LINE = "re-learned anything (>3 greps to find)? one anchor line → that area's code map" + + +def area_tasks(p: Project) -> list[str]: + """Add t-map-<area> refresh tasks for stale areas without an open one; returns output lines.""" + file = p.cfg.areas_file + if not file.is_file(): + return [] + out = [] + for a in areas_mod.parse(file.read_text(encoding="utf-8", errors="replace")): + if not areas_mod.status(p.cfg.code, a, p.cfg.area_stale_commits)[0]: + continue + prefix = f"t-map-{a.slug}" + if any(i == prefix or i.startswith(prefix + "-") for i in p.doc.ids()): + continue + id, n = prefix, 1 + while id in p.archived: + n += 1 + id = f"{prefix}-{n}" + item = tasks.Item(id, prio=1, effort="<1h", text=f"Refresh {a.name} area map.") + item.body = tasks.indent_body([ + f"Steps: wf areas {a.name}; fix missing anchors + test recipe from git log --stat " + f"<Checked>..HEAD -- <paths>; wf areas --mark {a.name}.", + f"Done: wf areas {a.name} shows ok.", + "Model: sonnet"]) + tasks.insert(p.doc, item, "pending") + out.append(f"added {id} (area map stale)") + return out + + +def uncovered_areas(p: Project, args) -> list[tuple[str, int]]: + """Folders this task's diff (branch since its base + uncommitted) touches outside every area's Paths.""" + file = p.cfg.areas_file + if not file.is_file() or p.cfg.code_root: # code in another repo: this diff is not its code + return [] + start = Path(getattr(args, "project", None) or Path.cwd()).resolve() + here = _here(start, p.cfg) + found = [areas_mod.with_anchor_paths(here, a) + for a in areas_mod.parse(file.read_text(encoding="utf-8", errors="replace"))] + if not any(a.paths for a in found): + return [] + wt = config.linked_worktree(start) + if wt: + base = config.git_branch(wt[2] / ".git") or "master" + else: + base = next((b for b in ("master", "main") + if git_run(here, "rev-parse", "-q", "--verify", f"refs/heads/{b}").returncode == 0), None) + books = {p.cfg.rel(f) for f in (p.cfg.tasks, p.cfg.archive, p.cfg.root / config.NAME)} | {p.cfg.areas_rel} + files = [f for f in areas_mod.changed_files(here, base) if f not in books] + return areas_mod.uncovered(files, found, p.cfg.area_ignore) + + +UNCOVERED_MAX = 3 # map tasks one wf done adds at most + + +def uncovered_tasks(p: Project, args) -> list[str]: + """Add t-map-<folder> tasks for folders the diff touches outside every area's Paths (once per id).""" + out = [] + for folder, n in uncovered_areas(p, args)[:UNCOVERED_MAX]: + id = f"t-map-{areas_mod.slug(folder)}" + if any(i == id or i.startswith(id + "-") for i in [*p.doc.ids(), *p.archived]): + continue + rel = p.cfg.areas_rel + item = tasks.Item(id, prio=1, effort="<1h", text=f"Map {folder} area. No area's Paths covers it.") + item.body = tasks.indent_body([ + f"Steps: git log --stat -- {folder}; new ### section under ## Areas in {rel} (Code map anchors, " + f"Test recipe, Paths: {folder}/) or add {folder}/ to an existing area's Paths; wf areas --mark <area>.", + f"Done: wf areas lists the area ok and no longer names {folder} uncovered.", + "Model: sonnet"]) + tasks.insert(p.doc, item, "pending") + out.append(f"added {id} (uncovered area: {folder}, {n} file{'s' * (n != 1)})") + return out + + +def cmd_done(args) -> int: + p = load_project(args, write=True) + doc = p.doc + task_ids = [i for i in args.ids if not i.startswith("a-")] + if len(task_ids) == 1 and not args.m: + raise Usage('done: -m "<entry>" is required for a single task') + if len(task_ids) > 1 and args.m: + raise Usage("done: -m goes with one task; several ids use each task's goal") + out = [] + today = datetime.date.today().isoformat() + for id in args.ids: + doc.find(id) + if not config.linked_worktree(Path(args.project or Path.cwd()).resolve()): + for id in task_ids: + if wt := claimed_worktree(p.cfg.root, doc.item(id)): + raise Failure(f"{id} is in progress in worktree {wt[0]} (branch {wt[1]}): run wf finish/done there " + f"(its bookkeeping goes in with wf merge), or wf status {id} clear first") + before = lanes.pickable_ids(doc, p.archived) + lane = lambda i: lanes.lane_of(i, p.cfg.lanes, p.cfg.slice_above) + done_lanes = {lane(doc.item(i)) for i in task_ids} + solo = [i for i in task_ids if doc.item(i).sessions == "solo"] + for id in args.ids: + if id.startswith("a-"): + tasks.remove(doc, id) + out.append(f"removed: {id}") + freed = tasks.unblock(doc, id) + if freed: + out.append("unblocked: " + ", ".join(freed)) + continue + open_ = [s for s in tasks.open_slices(doc, id) if s not in args.ids] + if open_: + raise Failure(f"'{id}' has open slices: {', '.join(open_)}") + item = tasks.remove(doc, id) + line = tasks.archive_line(today, item, args.m or item.goal) + p.new_archive = tasks.archive_prepend(p.new_archive or p.archive_text, line) + out.append(f"done: {id} → {p.cfg.rel(p.cfg.archive)}") + parent = tasks.parent_of(doc, id) + if parent and not tasks.open_slices(doc, parent): + out.append(f"last slice of {parent}: finish it with wf done {parent} -m \"…\"") + if task_ids and not args.dry_run: + out += area_tasks(p) + out += uncovered_tasks(p, args) + p.save(args.dry_run) + if args.dry_run: + return 0 + unclaim(p.cfg, task_ids) + after = lanes.pickable_ids(doc, p.archived | set(task_ids)) + freed = [i for s in doc.sections if s.key == "pending" for i in s.items + if i.id in after - before and lane(i) not in done_lanes] + if task_ids: + out += lanes.notify_block(freed, sessions(p.cfg), p.cfg.lanes, p.cfg.slice_above) + me = my_session() + out += lanes.solo_done_block(solo, sessions(p.cfg), me[1] if me else None) + if task_ids: + if p.cfg.verify: + out += ["verify:", *(f" {v}" for v in p.cfg.verify)] + out += ["checklist:", f" - {RELEARN_LINE}", *(f" - {d}" for d in p.cfg.done)] + if not getattr(args, "finishing", False): + out += worktree_done_lines(p, args) + if hint := ctx_hint_line(p.cfg): + out.append(hint) + print("\n".join(out)) + return 0 + + +CTX_TAIL = 1 << 20 # bytes of the transcript read for the last request + + +def ctx_hint_line(cfg) -> str | None: + """This session's prompt size over cfg.ctx_hint → a /clear hint; no session id / transcript → None.""" + sid = os.environ.get("CLAUDE_CODE_SESSION_ID") + if not cfg.ctx_hint or not sid: + return None + found = sorted(claude_projects().glob(f"*/{glob.escape(sid)}.jsonl")) + if len(found) != 1: + return None + try: + with open(found[0], "rb") as f: + f.seek(max(0, f.seek(0, 2) - CTX_TAIL)) + text = f.read().decode("utf-8", "replace") + except OSError: + return None + return usage.ctx_hint(usage.context_tokens(text), cfg.ctx_hint) + + +def claimed_worktree(root: Path, item: tasks.Item) -> tuple[str, str] | None: + """(worktree path, branch) of a linked worktree on the branch named by item's 'in progress: <branch>'.""" + status = item.status or "" + if not status.startswith("in progress: "): + return None + words = status.removeprefix("in progress: ").split() + branch = words[0] if words else "" + path = None + for line in git_run(root, "worktree", "list", "--porcelain").stdout.splitlines(): + if line.startswith("worktree "): + path = line.removeprefix("worktree ") + elif branch and line == f"branch refs/heads/{branch}" and path and Path(path).resolve() != root.resolve(): + if config.linked_worktree(Path(path)): + return path, branch + return None + + +def worktree_done_lines(p: Project, args) -> list[str]: + """Merge-back steps when wf done runs in a linked worktree (multi-session mode).""" + wt = config.linked_worktree(Path(args.project or Path.cwd()).resolve()) + if not wt: + return [] + branch, main = config.git_branch(wt[1]), wt[2] + master = config.git_branch(main / ".git") or "master" + files = " ".join(str(f.relative_to(main)) for f in (p.cfg.tasks, p.cfg.archive)) + warn = [] + if not branch and (n := unmerged_count(wt[0], master)): + warn = [f"detached HEAD has {n} commit{'s' * (n != 1)} not in {master}: wf merge merges " + f"{'them' if n != 1 else 'it'} (never leave them unmerged)"] + return warn + [f"worktree mode (branch {branch or 'none: detached HEAD'}), after verify: commit your code here (explicit paths), then:", + f" wf merge (rebase, ff-merge into {master}, commit {files}, push home; " + f"conflict → git rebase {master}, resolve, verify, wf merge again)"] + + +def git_run(cwd: Path, *args: str) -> subprocess.CompletedProcess: + return subprocess.run(["git", *args], cwd=cwd, capture_output=True, text=True) + + +def unmerged_count(top: Path, master: str) -> int: + """Commits on this worktree's HEAD not in master (0 when git fails).""" + r = git_run(top, "rev-list", "--count", f"{master}..HEAD") + return int(r.stdout.strip()) if r.returncode == 0 and r.stdout.strip().isdigit() else 0 + + +def cmd_merge(args) -> int: + """Merge-back of a lane worktree's branch, under the project lock (main() takes it).""" + start = Path(args.project or Path.cwd()).resolve() + wt = config.linked_worktree(start) + if not wt: + raise Failure("merge runs inside a linked git worktree (lane worktree)") + top, gitdir, main = wt + branch = config.git_branch(gitdir) + master = config.git_branch(main / ".git") or "master" + cfg = config.load_at(start) + files = [str(f.relative_to(main)) for f in (cfg.tasks, cfg.archive)] + if git_run(top, "status", "--porcelain").stdout.strip(): + raise Failure("worktree has uncommitted changes: commit them first") + if not branch and not unmerged_count(top, master): + # after a merge: notes / follow-ups written since → bookkeeping commit only + if not git_run(main, "status", "--porcelain", "--", *files).stdout.strip(): + raise Failure("worktree is on a detached HEAD: nothing to merge") + r = git_run(main, "commit", "-q", "-m", args.m or "bookkeeping", "--", *files) + if r.returncode: + raise Failure(f"bookkeeping commit failed: {(r.stderr.strip() or r.stdout.strip() or 'git error').splitlines()[-1]}") + print(f"committed {' '.join(files)}") + return 0 + out = [] # a branch, or a detached HEAD with commits not in master (never left behind) + if git_run(top, "rebase", master).returncode: + git_run(top, "rebase", "--abort") + raise Failure(f"rebase onto {master} conflicts: git rebase {master}, resolve, verify, then wf merge again") + out.append(f"rebased onto {master}") + args.merged_sha = git_run(top, "rev-parse", "--short", "HEAD").stdout.strip() + r = git_run(main, "merge", "--ff-only", branch or git_run(top, "rev-parse", "HEAD").stdout.strip()) + if r.returncode: + raise Failure(f"ff-merge into {master} failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") + out.append(f"fast-forwarded {master}") + if git_run(main, "status", "--porcelain", "--", *files).stdout.strip(): + msg = args.m or (f"{branch.rsplit('/', 1)[-1]} done" if branch else "bookkeeping") + r = git_run(main, "commit", "-q", "-m", msg, "--", *files) + if r.returncode: + raise Failure(f"bookkeeping commit failed: {(r.stderr.strip() or r.stdout.strip() or 'git error').splitlines()[-1]}") + out.append(f"committed {' '.join(files)}") + git_run(top, "switch", "-q", "--detach", master) + if branch: + git_run(top, "branch", "-q", "-d", branch) + if not args.no_push and _push_home(main): + out.append("pushed home") + out.append(f"merged {branch or 'detached HEAD'} into {master}") + print("\n".join(out)) + return 0 + + +def _push_home(main: Path) -> bool: + """git push home --all/--tags when remote home exists; True if pushed.""" + if "home" not in git_run(main, "remote").stdout.split(): + return False + for extra in ("--all", "--tags"): + r = git_run(main, "push", "-q", "home", extra) + if r.returncode: + raise Failure(f"push home failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") + return True + + +def _dirty_outside(top: Path, keep: list[Path]) -> list[str]: + """Uncommitted paths of the repo at `top` not under any of `keep` (absolute paths).""" + stray = [] + for line in git_run(top, "status", "--porcelain").stdout.splitlines(): + rel = line[3:].split(" -> ")[-1].strip('"').rstrip("/") + path = (top / rel).resolve() + if not any(path == k or k in path.parents for k in keep): + stray.append(rel) + return stray + + +GATE_RED_HELP = ( + "\ngate red procedure (also for a red gate run after wf finish, on the merged sha):" + "\n- culprit = your diff (git show --stat <sha>) explains the failure -> fix now (new commit, wf merge) or hand back" + "\n- not your diff (earlier task broke it; see git log of the failing area) -> wf add -p 0 --model <m> -e <effort>" + ' --done "gate ALL GREEN" -b "Fix gate red <test>. <goal>" with body line "Steps: <failing test, culprit sha/id>"' + "\n- then report result: done+gate-red <culprit sha or id> <fix id> (your task is done, the lane goes on with the fix)") + + +def cmd_finish(args) -> int: + """gate → done → commit the given paths → merge (lane worktree), in one call; checks first, so a refusal + leaves the task open.""" + if args.paths and not args.commit: + raise Usage("finish: paths need --commit MSG") + start = Path(args.project or Path.cwd()).resolve() + root = config.find_root(start) + cfg = config.load(root) + wt = config.linked_worktree(start) + top = Path(git_run(start, "rev-parse", "--show-toplevel").stdout.strip() or start) + paths = [(start / f).resolve() for f in args.paths] + books = [] if wt else [cfg.tasks.resolve(), cfg.archive.resolve()] + for f, path in zip(args.paths, paths): + if path != top and top not in path.parents: + raise Failure(f"path {f} is outside this repo ({top}): commit it in its own repo first, " + "then wf finish without it") + if wt and path in (top / cfg.tasks.relative_to(wt[2]), top / cfg.archive.relative_to(wt[2])): + raise Failure(f"path {f} = this worktree's copy of the books: wf finish writes the main tree's " + "and wf merge commits them; drop it") + if not path.exists() and git_run(top, "ls-files", "--error-unmatch", "--", str(path)).returncode: + raise Failure(f"path {f} does not exist and is not tracked: nothing to commit") + if wt and paths and not git_run(top, "status", "--porcelain", "--", *map(str, paths)).stdout.strip(): + raise Failure("nothing to commit in the --commit paths: drop them (wf finish without paths) or fix them") + if stray := _dirty_outside(top, paths + books): + raise Failure("uncommitted changes outside the --commit paths: " + " ".join(stray)) + if cfg.quick_gate: + _run_lines(cfg.quick_gate, _here(start, cfg), cfg, + "quick_gate '{line}' red (exit {rc}): fix it, then wf finish again; later commands skipped" + GATE_RED_HELP) + print(f"quick_gate: {len(cfg.quick_gate)} green", flush=True) + with project_lock(root): + p = load_project(args) + if args.id not in p.doc.ids() and args.id in p.archived: + print(f"{args.id} already done ({cfg.rel(cfg.archive)}): resuming with commit + merge", flush=True) + else: + done = argparse.Namespace(project=args.project, ids=[args.id], m=args.m, dry_run=False, finishing=True) + cmd_done(done) + sys.stdout.flush() + resume = f" ({args.id} is done already: fix it, then rerun the same wf finish, it resumes)" + commit = list(dict.fromkeys(map(str, [*paths, *(books if args.commit or not wt else [])]))) + if commit: + r = git_run(top, "add", "--", *map(str, paths)) if paths else None + if r and r.returncode: + raise Failure(f"git add failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}" + resume) + msg = args.commit or f"{args.id} done" + r = git_run(top, "commit", "-q", "-m", msg, "--", *commit) + if r.returncode: + raise Failure(f"commit failed: {(r.stderr.strip() or r.stdout.strip() or 'git error').splitlines()[-1]}" + + resume) + print("committed " + " ".join(os.path.relpath(c, top) for c in commit), flush=True) + if not wt: + sha = git_run(top, "rev-parse", "--short", "HEAD").stdout.strip() + tool = getattr(args, "tool_commit", None) + print(f"report: commit {sha}" + (f" tool {tool}" if tool else ""), flush=True) + if not wt and not args.no_push and _push_home(top): + print("pushed home", flush=True) + if wt: + ns = argparse.Namespace(project=args.project, m=None, no_push=args.no_push) + rc = cmd_merge(ns) + if not rc and getattr(ns, "merged_sha", ""): + tool = getattr(args, "tool_commit", None) + print(f"report: commit {ns.merged_sha}" + (f" tool {tool}" if tool else ""), flush=True) + return rc + return 0 + + +def cmd_wip(args) -> int: + """Wrap-up in one call (lane worktree): commit the given paths on the branch, note the state, status clear.""" + start = Path(args.project or Path.cwd()).resolve() + if not config.linked_worktree(start): + raise Failure("wip runs inside a lane worktree") + top = Path(git_run(start, "rev-parse", "--show-toplevel").stdout.strip() or start) + paths = [str((start / f).resolve()) for f in args.paths] + if paths: + r = git_run(top, "add", "--", *paths) + if r.returncode: + raise Failure(f"git add failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") + r = git_run(top, "commit", "-q", "-m", args.commit or f"{args.id} WIP", "--", *paths) + if r.returncode: + raise Failure(f"commit failed: {(r.stderr.strip() or r.stdout.strip() or 'git error').splitlines()[-1]}") + print("committed " + " ".join(os.path.relpath(c, top) for c in paths), flush=True) + p = load_project(args) + id = p.resolve(args.id) + with project_lock(p.cfg.root): + cmd_note(argparse.Namespace(project=args.project, id=id, line=args.m, dry_run=False)) + cmd_status(argparse.Namespace(project=args.project, id=id, kind="clear", value=[], dry_run=False, + clear_stale=False)) + return 0 + + +def _here(start: Path, cfg) -> Path: + """Project folder to run config commands in: the worktree's copy of cfg.root, else cfg.root.""" + wt = config.linked_worktree(start) + if not wt: + return cfg.root + top, _, main = wt + try: + return top / cfg.root.relative_to(main) + except ValueError: + return top + + +def _run_lines(lines: list[str], here: Path, cfg, fail: str) -> None: + """Run shell lines in `here` with env WF_MAIN = main tree project folder; first nonzero → Failure.""" + env = {**os.environ, "WF_MAIN": str(cfg.root)} + for line in lines: + print(f"$ {line}", flush=True) + rc = subprocess.run(line, shell=True, cwd=here, env=env).returncode + if rc: + raise Failure(fail.format(line=line, rc=rc)) + + +def cmd_setup(args) -> int: + """Run workflow.toml worktree_setup in this lane worktree (cwd = its project folder, env WF_MAIN).""" + start = Path(args.project or Path.cwd()).resolve() + wt = config.linked_worktree(start) + if not wt: + raise Failure("setup runs inside a linked git worktree (lane worktree)") + cfg = config.load_at(start) + if not cfg.worktree_setup: + print(f"no worktree_setup in {config.NAME}: nothing to do") + return 0 + _run_lines(cfg.worktree_setup, _here(start, cfg), cfg, + "worktree_setup '{line}' failed (exit {rc}): later commands skipped") + n = len(cfg.worktree_setup) + print(f"worktree_setup: {n} command{'s' * (n != 1)} ok in {os.path.relpath(wt[0], cfg.root)}") + return 0 + + +def cmd_start(args) -> int: + """Worker setup in one call (main tree): worktree + branch (clean check), worktree_setup, status progress, + then wf ctx and a ready verify && wf finish line.""" + start = Path(args.project or Path.cwd()).resolve() + if config.linked_worktree(start): + raise Failure("start runs in the main tree (it creates the worktree)") + p = load_project(args) + id = p.resolve(args.id) + main = Path(git_run(start, "rev-parse", "--show-toplevel").stdout.strip() or p.cfg.root) + master = config.git_branch(main / ".git") or "master" + wt = (Path.cwd() / args.worktree).resolve() + shown = p.cfg.rel(wt) + b = args.branch + + def git_ok(cwd, *a): + r = git_run(cwd, *a) + if r.returncode: + raise Failure(f"git {a[0]} failed: {(r.stderr.strip() or r.stdout.strip() or 'git error').splitlines()[-1]}") + return r.stdout + + has_branch = not git_run(main, "rev-parse", "--verify", "-q", f"refs/heads/{b}").returncode + extra = [] + if not wt.exists(): + if has_branch: + git_ok(main, "worktree", "add", "-q", str(wt), b) + how = f"new, existing branch {b}: earlier WIP, read the notes" + else: + git_ok(main, "worktree", "add", "-q", str(wt), "-b", b, master) + how = f"new, branch {b} from {master}" + else: + dirty = git_run(wt, "status", "--porcelain").stdout.rstrip("\n") + cur = config.git_branch(Path(git_run(wt, "rev-parse", "--absolute-git-dir").stdout.strip())) + if dirty and not args.recovery: + raise Failure(f"worktree {shown} has uncommitted changes: " + + " ".join(l.strip() for l in dirty.splitlines()) + + " (hand back, or --recovery when a dead worker left them)") + if dirty and cur != b: + raise Failure(f"worktree {shown} has uncommitted changes on another branch ({cur or 'detached'}), " + f"not {b}: hand back") + if cur == b: + how = f"on {b}" + elif has_branch: + git_ok(wt, "switch", "-q", b) + how = f"switched to existing branch {b}: earlier WIP, read the notes" + else: + git_ok(wt, "switch", "-q", "-c", b, master) + how = f"switched to new branch {b} from {master}" + if args.recovery: + log = git_run(wt, "log", "--oneline", f"{master}..HEAD").stdout.strip() + extra = [f"recovery: uncommitted:\n{dirty}" if dirty else "recovery: uncommitted: none", + f"recovery: commits {master}..HEAD:" + (f"\n{log}" if log else " none")] + print(f"worktree: {shown} ({how})", *extra, sep="\n", flush=True) + try: + here = wt / p.cfg.root.relative_to(main) + except ValueError: + here = wt + cmd_setup(argparse.Namespace(project=str(here))) + sys.stdout.flush() + with project_lock(p.cfg.root), contextlib.redirect_stdout(io.StringIO()): + cmd_status(argparse.Namespace(project=args.project, id=id, kind="progress", value=[b], dry_run=False, + clear_stale=False)) + print(flush=True) + cmd_ctx(argparse.Namespace(project=args.project, target=id)) + steps = [f"cd {here}", *(f"({v})" for v in p.cfg.verify), + f'python3 {HERE / "wf.py"} finish {id} -m "<entry>" --commit "<msg + footer>" <paths>'] + print("\nFinish (after the work; fill in the quoted parts and the paths):\n " + " && ".join(steps)) + return 0 + + +def cmd_gate(args) -> int: + """Run workflow.toml quick_gate (fast regression check) in this tree's project folder, env WF_MAIN.""" + start = Path(args.project or Path.cwd()).resolve() + cfg = config.load_at(start) + if not cfg.quick_gate: + print(f"no quick_gate in {config.NAME}: nothing to do") + return 0 + _run_lines(cfg.quick_gate, _here(start, cfg), cfg, + "quick_gate '{line}' red (exit {rc}): fix it before wf done; later commands skipped" + GATE_RED_HELP) + n = len(cfg.quick_gate) + print(f"quick_gate: {n} command{'s' * (n != 1)} green") + return 0 + + +# ------------------------------------------------------------------ orchestrator + +GO_ON = ("done", "done+gate-red", "sliced") # outcomes after which the lane picks its next task + + +def orch_main(args) -> tuple[Path, Project]: + """(main tree, project) for wf orch; refuses a linked worktree.""" + start = Path(args.project or Path.cwd()).resolve() + if config.linked_worktree(start): + raise Failure("orch runs in the main tree (the orchestrator's)") + p = load_project(args) + main = Path(git_run(start, "rev-parse", "--show-toplevel").stdout.strip() or p.cfg.root) + return main, p + + +def orch_record(cfg: config.Config, id: str) -> dict: + try: + r = json.loads((cfg.root / ".wf" / "orch" / f"{id}.json").read_text()) + except (OSError, ValueError): + return {} + return r if isinstance(r, dict) else {} + + +def live_session_cwds() -> list[Path]: + """cwd of every live Claude Code session (~/.claude/sessions/*.json).""" + out = [] + for f in (claude_projects().parent / "sessions").glob("*.json"): + try: + s = json.loads(f.read_text()) + os.kill(int(s["pid"]), 0) + out.append(Path(s["cwd"]).resolve()) + except (OSError, ValueError, KeyError, TypeError): + continue + return out + + +def worktree_free(path: Path, branch: str, busy: set[Path], cwds: list[Path]) -> bool: + """No live session in it, no other picked worker on it, and clean on a detached HEAD (or on branch = recovery).""" + if path in busy: + return False + if not path.exists(): + return True + if any(c == path or path in c.parents for c in cwds): + return False + cur = config.git_branch(Path(git_run(path, "rev-parse", "--absolute-git-dir").stdout.strip() or path / ".git")) + if cur == branch: + return True + return not cur and not git_run(path, "status", "--porcelain").stdout.strip() + + +def lane_worktree(main: Path, p: Project, lane: str, id: str) -> Path: + """First free main/.worktrees/<lane>[-n] for branch <lane>/<id> (worktree_free; other picked workers' trees busy).""" + cfg, busy = p.cfg, set() + for f in (cfg.root / ".wf" / "orch").glob("*.json"): + r = orch_record(cfg, f.stem) + other = p.doc.item(f.stem) if f.stem in p.doc.ids() else None + if f.stem != id and r.get("worktree") and other and (other.status or "").startswith("in progress"): + busy.add(Path(r["worktree"])) + cwds = live_session_cwds() + n = 1 + while not worktree_free(wt := main / ".worktrees" / (lane if n == 1 else f"{lane}-{n}"), f"{lane}/{id}", busy, cwds): + n += 1 + return wt + + +CLOUD = "cloud" # virtual lane: wf cloud send instead of a local worker (spec cloud-lane §4.6) + + +def cloud_pick(main: Path, p: Project, id: str | None) -> int: + """wf orch pick cloud: ledger check, first fitting task across the lanes (cloud.pick_key) -> wf cloud send.""" + import wf_cloud + from wflib import cloud + cfg = p.cfg + if not cfg.cloud: + print(f"stop lane {CLOUD}: none fit (project not opted in: workflow.toml cloud = true)") + return 0 + with wf_cloud.locked(wf_cloud.state_dir()) as led: + run, cap, bal, res = cloud.running(led), led["max_parallel"], cloud.balance(led), led["reserve_per_task"] + if run >= cap: + print(f"stop lane {CLOUD}: max parallel ({run} running >= max_parallel {cap})") + return 0 + if bal < res: + print(f"stop lane {CLOUD}: ledger (balance ${bal:.2f} < reserve ${res:.2f})") + return 0 + if id: + item = p.doc.item(p.resolve(id)) + if why := cloud.fit(item, cfg): + print(f"stop lane {CLOUD}: none fit ({item.id}: {why})") + return 0 + else: + ready, seen = runner_ids(p), [] + for l in cfg.lanes: + seen += [i for i in lanes.ranked(p.doc, p.archived, cfg.lanes, cfg.slice_above, l.name) + if i not in seen and i.id in ready and not (i.status or "").startswith("in progress")] + fits = sorted((i for i in seen if cloud.fit(i, cfg) is None), key=cloud.pick_key) + if not fits: + print(f"stop lane {CLOUD}: none fit (wf cloud fit rules: wf set ID --cloud yes skips the regex list, opus only)") + return 0 + item = fits[0] + buf = io.StringIO() + with contextlib.redirect_stdout(buf): + code = wf_cloud.main(["send", item.id, "--project", str(cfg.root)]) + sent = buf.getvalue().strip() + if code == 3: + print(f"stop lane {CLOUD}: ledger (send refused, see above)") + return 0 + if code: + if sent: + print(sent) + raise Failure(f"wf cloud send {item.id} failed (exit {code}): lane {CLOUD} stops, tell the owner") + rec = json.loads(wf_cloud.record_path(cfg, item.id).read_text()) + (wf_folder(cfg, "orch") / f"{item.id}.json").write_text(json.dumps( + {"id": item.id, "lane": CLOUD, "task_lane": rec["lane"], "model": item.model, "effort": item.effort, + "sid": rec["sid"], "branch": f"{rec['lane']}/{item.id}", "at": int(datetime.datetime.now().timestamp())}) + "\n") + print(f"pick: {item.id} (lane {CLOUD}, model {item.model}, effort {item.effort or '-'}) · {sent.splitlines()[-1]}" + f" · claimed (in progress: cloud:{rec['sid']})") + print(f"agent: none (cloud session). On each wake (>= 10 min apart): wf cloud pull --all; per ended task " + f"'<id>: <state> …' → wf orch post <id> {CLOUD} --result <state> [--commit <sha>]") + return 0 + + +def orch_pick(main: Path, p: Project, lane: str, id: str | None = None, recovery: str | None = None) -> int: + cfg = p.cfg + stop = cfg.root / "out" / "wf-batch.stop" + if stop.exists() and not recovery: + print(f"stop: {cfg.rel(stop)} exists: spawn nothing (let running workers finish)") + return 0 + if lane == CLOUD: + return cloud_pick(main, p, id) + if id: + id = p.resolve(id) + else: + ready = runner_ids(p) + order = [i for i in lanes.ranked(p.doc, p.archived, cfg.lanes, cfg.slice_above, lane) + if i.id in ready and not (i.status or "").startswith("in progress")] + if not order: + print(f"none: lane {lane} has no runner-ready task (wf list --runner --lane {lane}; wf lanes)") + return 0 + id = order[0].id + item = p.doc.item(id) + branch = f"{lane}/{id}" + wt = lane_worktree(main, p, lane, id) + with project_lock(cfg.root), contextlib.redirect_stdout(io.StringIO()): + cmd_status(argparse.Namespace(project=str(cfg.root), id=id, kind="progress", value=["worker"], + dry_run=False, clear_stale=False)) + model = item.model + (wf_folder(cfg, "orch") / f"{id}.json").write_text(json.dumps( + {"id": id, "lane": lane, "model": model, "effort": item.effort, "worktree": str(wt), "branch": branch, + "at": int(datetime.datetime.now().timestamp())}) + "\n") + kind = ", slice job: the worker only slices" if lanes.is_slice_job(item, cfg.slice_above) else "" + print(f"pick: {id} (lane {lane}, model {model}, effort {item.effort or '-'}{kind}) · claimed (in progress: worker)") + print(f"agent: subagent_type wf-worker, model {model}, no isolation; prompt:") + print(f"Task: {id} Lane: {lane} Model: {model}\nMain tree: {main} Worktree: {wt} Branch: {branch}" + + (f"\nRecovery: {recovery}" if recovery else "") + "\nFinal message: the 4 report lines only.") + return 0 + + +def cost_line(root: Path, agent: str, task: str, outcome: str, effort: str | None, lane: str, dur: int | None) -> str: + """Append one usage line for a worker to <root>/out/wf-cost.log; returns it.""" + path = usage_transcripts(argparse.Namespace(agent=agent, session=None))[0][1] + now = datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + line = usage.log_line(now, root.name, task, effort or "-", outcome, path.stem.removeprefix("agent-"), + usage.parse(read(path)), lane=lane, dur=dur) + (root / "out").mkdir(exist_ok=True) + with open(root / "out" / "wf-cost.log", "a") as f: + f.write(line + "\n") + return line + + +def archived_meta(cfg: config.Config, id: str) -> tuple[str | None, str | None]: + """(model, effort) of a task already archived: its block in TASKS.md just before the commit that removed it.""" + rel = str(cfg.tasks.relative_to(cfg.root)) + r = git_run(cfg.root, "log", "-n1", "--format=%H", f"-S**{id}**", "--", rel) + if r.returncode or not r.stdout.strip(): + return None, None + old = git_run(cfg.root, "show", f"{r.stdout.strip()}^:{rel}") + try: + it = tasks.parse(old.stdout).item(id) if old.returncode == 0 and id in tasks.parse(old.stdout).ids() else None + except Exception: + it = None + return (it.model, it.effort) if it else (None, None) + + +def orch_post(main: Path, p: Project, args) -> int: + cfg, id, lane = p.cfg, args.id, args.lane + rec = orch_record(cfg, id) + wt = Path(rec.get("worktree") or main / ".worktrees" / lane) + item = p.doc.item(id) if id in p.doc.ids() else None + if item is None and id not in p.archived: + raise Failure(tasks.unknown_id(id, p.doc.ids())) + words = (args.result or "").split() + outcome = words[0] if words else ("done" if id in p.archived else "no-report") + if outcome == "done" and item is not None and id not in p.archived and tasks.open_slices(p.doc, id): + outcome = "sliced" # slice job: worker said done but the task is open with slices + master = config.git_branch(main / ".git") or "master" + files = [str(f.relative_to(main)) for f in (cfg.tasks, cfg.archive)] + problems, out = [], [] + with project_lock(cfg.root): + if outcome.startswith("done"): + if id not in p.archived: + problems.append(f"no archive line for {id}") + if (wt / ".git").exists(): + dirty = git_run(wt, "status", "--porcelain").stdout.strip() + if dirty: + problems.append(f"worktree {wt} dirty: " + " ".join(l.strip() for l in dirty.splitlines())) + elif git_run(wt, "merge-base", "--is-ancestor", "HEAD", master).returncode: + buf = io.StringIO() + try: + with contextlib.redirect_stdout(buf): + cmd_merge(argparse.Namespace(project=str(wt), m=None, no_push=args.no_push)) + out.append(f"merged worktree HEAD: {buf.getvalue().strip().splitlines()[-1]}") + except Failure as e: + problems.append(f"wf merge in {wt}: {e}") + branch = rec.get("branch") or f"{lane}/{id}" + if not git_run(main, "rev-parse", "--verify", "-q", f"refs/heads/{branch}").returncode: + problems.append(f"branch {branch} still there") + errors = checks.check(config.load(cfg.root))[0] + if errors: + problems.append(f"wf check: {len(errors)} errors: {errors[0]}") + elif git_run(main, "status", "--porcelain", "--", *files).stdout.strip(): + r = git_run(main, "commit", "-q", "-m", f"{id} {outcome} (orchestrator)", "--", *files) + out.append("committed leftover " + " ".join(files) if not r.returncode + else f"leftover commit failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") + raised = bool(item and rec.get("model") and not outcome.startswith("done") + and tasks.MODELS.index(item.model) > tasks.MODELS.index(rec["model"])) + final = "post-check-red" if problems else ("model-raised" if raised else outcome) + am, ae = (None, None) if item or (rec.get("model") and rec.get("effort")) else archived_meta(cfg, id) + model = rec.get("model") or (item.model if item else am) or "?" + dur = args.duration if args.duration else ( # None/0 (orchestrator lacked duration_ms) → since the pick + int(datetime.datetime.now().timestamp()) - rec["at"] if isinstance(rec.get("at"), int) else None) + commit = args.commit or (git_run(main, "rev-parse", "--short", master).stdout.strip() + if outcome.startswith("done") else "-") + now = datetime.datetime.now().isoformat(timespec="seconds") + (cfg.root / "out").mkdir(exist_ok=True) + with open(cfg.root / "out" / "wf-orch.log", "a") as f: + f.write(f"{now} {lane} {model} {id} {final} {commit} {fmt_dur(dur) if dur is not None else '-'}" + + (f" ({'; '.join(problems)})" if problems else "") + "\n") + if args.agent: + try: + out.append("cost: " + cost_line(cfg.root, args.agent, id, final, rec.get("effort") or (item and item.effort) or ae, + lane, dur)) + except Failure as e: + out.append(f"cost: not logged ({e})") + if final != "post-check-red": # kept: the re-post after the fix still knows the pick time + (cfg.root / ".wf" / "orch" / f"{id}.json").unlink(missing_ok=True) + print(f"post: {id} {final}" + "".join(f"\n {l}" for l in problems + out)) + if outcome == "done+gate-red" and len(words) >= 3: + q = load_project(args) + fix = q.doc.item(words[2]) if words[2] in q.doc.ids() else None + if fix and not fix.runner_ready: + print(f"fix {fix.id} not runner-ready: add its Done/Model (wf set {fix.id} --done … --model …), " + f"then wf orch pick {lane} --id {fix.id}") + return 0 + if final not in GO_ON and final != "model-raised": + print(f"stop lane {lane}: {final} → tell the owner" + + (f" (crash/no report: one fresh worker: wf orch pick {lane} --id {id} --recovery \"<why>\", then stop)" + if final == "no-report" else "")) + return 0 + if args.no_pick: + return 0 + print() + return orch_pick(main, load_project(args), lane) + + +def cmd_orch(args) -> int: + main, p = orch_main(args) + if args.lane != CLOUD or CLOUD in lane_names(p.cfg): + check_lane(p.cfg, args.lane) + if args.action == "pick": + return orch_pick(main, p, args.lane, args.id, args.recovery) + return orch_post(main, p, args) + + +def change(args, action) -> int: + """Run one edit on one item, save, print its new header.""" + p = load_project(args, write=True) + note = action(p) + p.save(args.dry_run) + if not args.dry_run: + print(note if note else p.doc.item(args.id).header()) + return 0 + + +def cmd_prio(args) -> int: + return change(args, lambda p: tasks.set_prio(p.doc, args.id, args.prio)) + + +def cmd_move(args) -> int: + targets = [t for t in (args.section, args.before, args.after) if t] + if len(targets) != 1: + raise Usage("move: give a section, or --before ID, or --after ID") + if args.section: + return change(args, lambda p: tasks.move_to(p.doc, args.id, args.section)) + return change(args, lambda p: tasks.move_rel(p.doc, args.id, args.before or args.after, + before=bool(args.before), force=args.force)) + + +def clear_stale(args) -> int: + p = load_project(args, write=True) + items = stale(p) + for i in items: + tasks.set_status(p.doc, i.id, None) + p.save(args.dry_run) + if not args.dry_run: + unclaim(p.cfg, [i.id for i in items]) + for i in items: + print(f"cleared {i.id}") + return 0 + + +def cmd_status(args) -> int: + if args.clear_stale: + if args.id or args.kind or args.value: + raise Usage("status: --clear-stale takes no id/kind/value") + return clear_stale(args) + if not args.id or not args.kind: + raise Usage("status: ID and progress NOTE|blocked A-ID|clear (or --clear-stale)") + if args.kind == "clear": + if args.value: + raise Usage("status: clear takes no value") + status = None + elif not args.value: + raise Usage("status: progress needs a note, blocked needs an a-id") + elif args.kind == "progress": + status = "in progress: " + " ".join(args.value) + else: + status = f"blocked: [[{args.value[0]}]]" + code = change(args, lambda p: tasks.set_status(p.doc, args.id, status)) + if not args.dry_run: + proj = load_project(args) + cfg = proj.cfg + if args.kind == "progress": + claim(cfg, args.id, lanes.lane_of(proj.doc.item(args.id), cfg.lanes, cfg.slice_above)) + else: + unclaim(cfg, [args.id]) + return code + + +def cmd_set(args) -> int: + if args.interactive is not None: + old_spelling("--interactive", "--sessions owner") + if args.sessions is None: + args.sessions = "owner" if args.interactive == "yes" else "" + fields = dict( + title=args.title, effort=args.effort, + after=None if args.after is None else split_list(args.after), + refs=None if args.ref is None else split_list(args.ref), model=args.model, sessions=args.sessions, + done=args.done, cloud=args.cloud) + if args.model not in (None, "", *tasks.MODELS): + raise Usage(f"set: --model {args.model} (want {', '.join(tasks.MODELS)}, or \"\" to remove)") + if args.sessions not in (None, "", *tasks.SESSIONS): + raise Usage(f"set: --sessions {args.sessions} (want {', '.join(tasks.SESSIONS)}, or \"\" to remove)") + if args.cloud not in (None, "", *tasks.CLOUDS): + raise Usage(f"set: --cloud {args.cloud} (want {', '.join(tasks.CLOUDS)}, or \"\" to remove)") + if all(v is None for v in fields.values()): + raise Usage("set: nothing to set (--title --effort --after --ref --model --sessions --cloud --done)") + return change(args, lambda p: tasks.set_fields(p.doc, args.id, **fields)) + + +def cmd_rename(args) -> int: + def action(p): + n = tasks.rename(p.doc, args.id, args.new, p.archived) + return f"renamed: {args.id} → {args.new} ({n} link{'' if n == 1 else 's'})" + return change(args, action) + + +def cmd_note(args) -> int: + return change(args, lambda p: tasks.add_note(p.doc, args.id, " ".join(args.line.split()))) + + +def cmd_body(args) -> int: + if args.text: + raise Failure("body text goes on stdin (wf body ID <<'EOF' … EOF), not as an argument") + lines = stdin_lines() + return change(args, lambda p: tasks.set_body(p.doc, args.id, lines)) + + +def cmd_tick(args) -> int: + return change(args, lambda p: "ticked: " + tasks.tick(p.doc, args.id, args.which)) + + +# ------------------------------------------------------------------ other commands + +def wf_version() -> str: + try: + out = subprocess.run(["git", "-C", str(HERE), "rev-parse", "--short=7", "HEAD"], + capture_output=True, text=True, timeout=10) + return out.stdout.strip() if out.returncode == 0 and out.stdout.strip() else "0000000" + except (OSError, subprocess.SubprocessError): + return "0000000" + + +def cmd_report(args) -> int: + text = " ".join(args.text.split()) + if not text: + raise Usage("report: say what happened") + start = Path(args.project) if args.project else Path.cwd() + try: + project = config.find_root(start).name + except config.ConfigError: + project = start.resolve().name + cmd = f" (cmd: {' '.join(args.cmd.split())})" if args.cmd else "" + line = f"- {datetime.date.today().isoformat()} {project} {args.kind}: {text}{cmd} @{wf_version()}\n" + fd = os.open(inbox_path(), os.O_WRONLY | os.O_APPEND | os.O_CREAT, 0o644) + try: + os.write(fd, line.encode("utf-8")) + finally: + os.close(fd) + print("reported (workflow inbox); carry on") + return 0 + + +def cmd_init(args) -> int: + root = (Path(args.project) if args.project else Path.cwd()).resolve() + if (root / config.NAME).exists(): + raise Failure(f"{root}/{config.NAME} exists already") + templates = HERE / "templates" + for target, source in ((config.NAME, "workflow.toml"), ("TASKS.md", "TASKS.md"), + ("tasks/archive.md", "archive.md"), ("CLAUDE.md", "CLAUDE.md")): + path = root / target + if path.exists(): + print(f"kept {target} (exists" + ("; old format → wf migrate)" if target == "TASKS.md" else ")")) + continue + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(read(templates / source).replace("{name}", root.name), encoding="utf-8") + print(f"wrote {target}") + return 0 + + +def cmd_migrate(args) -> int: + from wflib import migrate + p = load_project(args) + cfg = p.cfg + new, report = migrate.migrate(p.tasks_text, taken=p.archived) + bump = cfg.format < config.FORMAT + if new == p.tasks_text and not bump: + print("nothing to migrate") + return 0 + name = cfg.rel(cfg.tasks) + sys.stdout.writelines(difflib.unified_diff( + p.tasks_text.replace("\r\n", "\n").splitlines(keepends=True), + new.replace("\r\n", "\n").splitlines(keepends=True), name, f"{name} (new)")) + if report.ids: + print("\nids:") + for title, id in report.ids.items(): + print(f" {id} {title}") + if report.notes: + print("\nnotes:") + for note in report.notes: + print(f" {note}") + errors = checks.check(cfg, tasks_text=new, slow=False)[0] + for e in errors: + print(f"ERROR: {e}") + print(f"check after migrate: {len(errors)} errors") + if not args.write: + print("dry run: nothing written (wf migrate --write)") + return 0 + wrote = [] + if new != p.tasks_text: + write_if_unchanged(cfg.tasks, new, p.tasks_stamp) + wrote.append(name) + if bump: + toml = cfg.root / config.NAME + text = read(toml) + changed = re.sub(r"(?m)^format(\s*)=(\s*)\d+", rf"format\g<1>=\g<2>{config.FORMAT}", text, count=1) + write_if_unchanged(toml, changed, stamp(toml)) + wrote.append(f"{config.NAME} (format = {config.FORMAT})") + print("wrote " + ", ".join(wrote)) + return 0 + + +# ------------------------------------------------------------------ main + +def parser() -> argparse.ArgumentParser: + ap = argparse.ArgumentParser(prog="wf", description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("--project", metavar="DIR", help="project folder (default: found from the working directory)") + sub = ap.add_subparsers(dest="command", metavar="COMMAND") + sections = list(tasks.SECTIONS) + + def cmd(name, func, help, write=False): + sp = sub.add_parser(name, help=help, description=help) + sp.set_defaults(func=func, locks=write) + if write: + sp.add_argument("--dry-run", action="store_true", help="print the diff, write nothing") + return sp + + sp = cmd("next", cmd_next, "cold start: awaiting, needs human, plans in flight, next task with its refs") + sp.add_argument("--brief", "-b", action="store_true", help="the task only") + sp.add_argument("--lane", help="lane name (wf lanes); none = all lanes") + sp.add_argument("--as", dest="as_", choices=tasks.MODELS, + help="your model: takes tasks with Model ≤ it (default haiku)") + sp.add_argument("--owner", action="store_true", help="the owner is present: Sessions: owner tasks pickable") + + sp = cmd("areas", cmd_areas, "area notes: anchors missing, commits since Checked (stale = missing or " + "≥ area_stale_commits); no NAME → also uncovered: folders the diff since master touches outside " + "every Paths (area_ignore skipped); --mark NAME stamps HEAD", write=True) + sp.add_argument("name", nargs="?", help="one area only") + sp.add_argument("--mark", metavar="NAME", help="set the area's Checked to HEAD (after a refresh)") + sp = cmd("lanes", cmd_lanes, "lanes: pickable/waiting per lane, the session of each " + "(--lane/--as registers yours; neither lane = all)") + sp.add_argument("--lane", help="your lane (wf lanes lists them)") + sp.add_argument("--as", dest="as_", choices=tasks.MODELS, help="your model (default haiku)") + sp.add_argument("--unregister", action="store_true", help="drop this session's lane record (orchestrators hold none)") + sp.add_argument("--wait", type=float, metavar="SECS", + help="block until any lane has a pickable task (exit 0, prints the lanes line) or SECS pass " + "(exit 1, prints the lanes line); out/wf-batch.stop exists (wf batch --stop) → exit 2 at once; " + "polls every 30 s (env WF_LANES_POLL)") + sp = cmd("list", cmd_list, "one line per item: id, priority, effort, status, model, lane, title") + sp.add_argument("-s", "--section", choices=[*sections, "all"], default="pending") + sp.add_argument("-p", "--prio", type=int, choices=range(4), help="this priority or higher") + sp.add_argument("--ready", action="store_true", help="pickable now") + sp.add_argument("--runner", action="store_true", + help="--ready + Done line, no Sessions: owner, Done not about owner/confirming, no open slices") + sp.add_argument("--blocked", action="store_true") + sp.add_argument("--progress", action="store_true") + sp.add_argument("--stale", action="store_true", help="in progress with no live claim (dead or no holder)") + sp.add_argument("--model", choices=tasks.MODELS, help="tasks with this Model line (none = opus)") + sp.add_argument("--lane", help="tasks of this lane (by effort); with --runner: in pick order") + sp.add_argument("--ref", metavar="PATH[#ANCHOR]") + sp.add_argument("-n", type=int, metavar="N") + + sp = cmd("show", cmd_show, "the raw item(s)") + sp.add_argument("ids", nargs="+", metavar="ID") + + sp = cmd("ctx", cmd_ctx, "item with refs resolved, its relations and the areas it names; or a doc section and its tasks") + sp.add_argument("target", metavar="ID|PATH#ANCHOR") + + sp = cmd("search", cmd_search, "ranked search over tasks, archive and docs") + sp.add_argument("words", nargs="+") + sp.add_argument("--tasks", action="store_true") + sp.add_argument("--archive", action="store_true") + sp.add_argument("--docs", action="store_true") + sp.add_argument("-n", type=int, default=15, metavar="N") + + sp = cmd("log", cmd_log, "newest archive lines") + sp.add_argument("words", nargs="*") + sp.add_argument("-n", type=int, default=10, metavar="N") + + sp = cmd("usage", cmd_usage, "tokens and API-price $ per agent of a session, from Claude Code transcripts") + sp.add_argument("--session", metavar="ID", help="session id or prefix (default: this session, else newest here)") + sp.add_argument("--agent", metavar="ID", help="one subagent only (agent id from its completion notice)") + sp.add_argument("--since", metavar="TIME", help="ISO UTC time, or 90m / 2h / 1d ago") + sp.add_argument("--log", nargs=2, metavar=("TASK", "OUTCOME"), + help="with --agent: append one line to <project>/out/wf-cost.log (orchestrator, per worker)") + sp.add_argument("--effort", choices=tasks.EFFORTS, metavar="|".join(tasks.EFFORTS), + help="the task's estimate (not reasoning effort), for --log") + sp.add_argument("--duration", type=int, metavar="S", help="agent wall time in seconds (duration_ms / 1000), for --log") + sp.add_argument("--lane", metavar="L", help="lane of the logged worker, for --log (default: its model)") + sp.add_argument("--explore", action="store_true", + help="exploring cost per subagent: pct billed input exploring / before the first edit, " + "result tokens per tool, files read by >= 2 agents (same --session / --since)") + sp.add_argument("--report", action="store_true", + help="per lane and effort: n, done, $ (median, total, per done), turns, duration, from every project's log") + + cmd("projects", cmd_projects, "every wf project: counts, errors, next task") + cmd("check", cmd_check, "validate ids, links, refs, order") + + sp = cmd("add", cmd_add, "new item; `add -` reads a whole item block from stdin", write=True) + sp.add_argument("title", metavar='"Title. Goal."|-') + sp.add_argument("-p", "--prio", type=int, choices=range(4)) + sp.add_argument("-e", "--effort", metavar="|".join(tasks.EFFORTS)) + sp.add_argument("-s", "--section", choices=sections) + sp.add_argument("--id", help="the id (default: from the title, or <parent>-N with --parent); " + "a leading 'ID: ' in the text works too") + sp.add_argument("--after", metavar="ID,…") + sp.add_argument("--ref", metavar="REF,…") + sp.add_argument("--parent", metavar="ID", help="add as the next slice of ID, After: the previous one") + sp.add_argument("--interactive", action="store_true", help=argparse.SUPPRESS) + sp.add_argument("--model", choices=tasks.MODELS, help="Model line (default: none = opus)") + sp.add_argument("--done", metavar="TEXT", help="Done line (makes the task runner-pickable)") + sp.add_argument("--cloud", choices=tasks.CLOUDS, help="Cloud line: yes = cloud lane may take it (opus only), no = never") + sp.add_argument("--sessions", choices=tasks.SESSIONS, + help="Sessions line: solo = no other live session, owner = owner present (default parallel)") + sp.add_argument("-b", "--body", action="store_true", help="read body lines from stdin") + + sp = cmd("done", cmd_done, "finish task(s): remove, archive line, print verify + checklist", write=True) + sp.add_argument("ids", nargs="+", metavar="ID") + sp.add_argument("-m", metavar="ENTRY", help="archive entry, ≤2 lines") + + sp = cmd("prio", cmd_prio, "set priority, reposition", write=True) + sp.add_argument("id") + sp.add_argument("prio", type=int, choices=range(4)) + + sp = cmd("move", cmd_move, "move to a section, or before/after another item", write=True) + sp.add_argument("id") + sp.add_argument("section", nargs="?", choices=sections) + sp.add_argument("--before", metavar="ID") + sp.add_argument("--after", metavar="ID") + sp.add_argument("--force", action="store_true", help="allow breaking priority order") + + sp = cmd("status", cmd_status, "progress NOTE | blocked A-ID | clear", write=True) + sp.add_argument("id", nargs="?") + sp.add_argument("kind", nargs="?", choices=["progress", "blocked", "clear"]) + sp.add_argument("value", nargs="*") + sp.add_argument("--clear-stale", action="store_true", help="clear in progress of tasks with no live claim") + + sp = cmd("set", cmd_set, "change header fields, Done:, Model:, After:, Ref:", write=True) + sp.add_argument("id") + sp.add_argument("--title", help="new title, goal kept; 'Title. Goal.' or a question replaces the text") + sp.add_argument("--effort") + sp.add_argument("--after", metavar="ID,…", help='"" removes the line') + sp.add_argument("--ref", metavar="REF,…", help='"" removes the line') + sp.add_argument("--interactive", choices=["yes", "no"], help=argparse.SUPPRESS) + sp.add_argument("--sessions", metavar="|".join(tasks.SESSIONS), help='Sessions line; "" removes it (= parallel)') + sp.add_argument("--model", metavar="|".join(tasks.MODELS), help='Model line; "" removes it (= opus)') + sp.add_argument("--cloud", metavar="|".join(tasks.CLOUDS), help='Cloud line (yes: skip the fit regex list, opus only; no: never cloud); "" removes it') + sp.add_argument("--done", metavar="TEXT", help='Done line (makes it runner-pickable); "" removes it') + + sp = cmd("rename", cmd_rename, "change an open item's id and every [[link]] to it", write=True) + sp.add_argument("id") + sp.add_argument("new", metavar="NEW") + + sp = cmd("note", cmd_note, "append one line to the body", write=True) + sp.add_argument("id") + sp.add_argument("line") + + sp = cmd("body", cmd_body, "replace the body with stdin; Model:, Sessions:, After:, Ref: lines kept", write=True) + sp.formatter_class = argparse.RawDescriptionHelpFormatter + sp.epilog = "example (body text goes on stdin, not as an argument):\n wf body ID <<'EOF'\n - Steps: …\n - Done: …\n EOF" + sp.add_argument("id") + sp.add_argument("text", nargs="*", help=argparse.SUPPRESS) + + sp = cmd("tick", cmd_tick, "check a `- [ ]` box by number or text", write=True) + sp.add_argument("id") + sp.add_argument("which", metavar="N|TEXT") + sp = cmd("merge", cmd_merge, "in a lane worktree: rebase, ff-merge into master, commit TASKS/archive, " + "detach, push home (under the project lock)") + sp.set_defaults(locks=True) + sp.add_argument("-m", help="bookkeeping commit message (default '<id> done', id from the branch)") + sp.add_argument("--no-push", action="store_true", help="skip git push home") + + sp = cmd("finish", cmd_finish, "worker's last step in one call: quick_gate → done → commit PATHS (explicit) → " + "wf merge (lane worktree; main tree: TASKS/archive go into the commit). Run verify first. " + "Refuses before done if the gate is red, files outside PATHS are uncommitted, or a PATH is outside " + "this repo / missing / unchanged; an id already archived resumes (commit + merge). " + "Prints the done output (notify / checklist) and the merge lines, then 'report: commit <sha> [tool <sha>]' " + "(lane worktree) = the sha for your report") + sp.add_argument("id") + sp.add_argument("-m", required=True, metavar="ENTRY", help="archive entry, ≤2 lines") + sp.add_argument("--commit", metavar="MSG", help="commit message for PATHS (incl. footer lines)") + sp.add_argument("paths", nargs="*", metavar="PATH", help="files/folders to commit (relative to cwd)") + sp.add_argument("--no-push", action="store_true", help="skip git push home") + sp.add_argument("--tool-commit", metavar="SHA", help="tool-repo sha merged by hand; added to the report line") + + sp = cmd("wip", cmd_wip, "wrap-up in one call (lane worktree): commit PATHS on the branch, note MSG " + "(state + next step), status clear") + sp.add_argument("id") + sp.add_argument("-m", required=True, metavar="NOTE", help="state + next step, one line") + sp.add_argument("--commit", metavar="MSG", help="commit message (default '<id> WIP')") + sp.add_argument("paths", nargs="*", metavar="PATH", help="files/folders to commit (relative to cwd)") + + cmd("setup", cmd_setup, "in a lane worktree: run workflow.toml worktree_setup (cwd = worktree, env WF_MAIN = " + "main tree project folder); provides git-ignored inputs, commands must be idempotent") + + sp = cmd("start", cmd_start, "worker setup in one call, in the main tree: worktree at PATH on BRANCH " + "(new: from master or the existing branch; existing: must be clean, switched to BRANCH), " + "worktree_setup, status progress BRANCH, then wf ctx and a ready verify && wf finish line") + sp.add_argument("id") + sp.add_argument("--worktree", required=True, metavar="PATH", help="worktree folder (relative to cwd)") + sp.add_argument("--branch", required=True, metavar="BRANCH") + sp.add_argument("--recovery", action="store_true", + help="a dead worker's WIP: uncommitted changes on BRANCH kept, printed with its commits") + + cmd("gate", cmd_gate, "run workflow.toml quick_gate (fast regression check a worker runs before wf done; " + "cwd = this tree's project folder, env WF_MAIN = main tree project folder); red = exit 1, nothing = exit 0") + + sp = cmd("report", cmd_report, "report a workflow problem or idea to the workflow inbox") + sp.add_argument("text") + sp.add_argument("--kind", choices=["bug", "idea", "friction"], default="bug") + sp.add_argument("--cmd", metavar="COMMAND") + + cmd("init", cmd_init, "make this folder a wf project") + + sp = cmd("migrate", cmd_migrate, "convert a numbered TASKS.md to the id format (dry run unless --write)") + sp.add_argument("--write", action="store_true") + + sp = cmd("orch", cmd_orch, "orchestrator, main tree: pick LANE (pick + claim + free worktree + agent prompt) · " + "post ID LANE (post-check, merge if needed, orch log + cost line, next pick)") + osub = sp.add_subparsers(dest="action", required=True) + op = osub.add_parser("pick", help="first runner-ready task of LANE: claim it, choose a free worktree " + "(.worktrees/LANE, -2, …), print the wf-worker prompt; 'none:'/'stop:' line = spawn nothing. " + "LANE cloud: ledger check, first cloud-fitting task of any lane (opus first) -> wf cloud send; " + "'stop lane cloud: ledger|none fit|max parallel'") + op.add_argument("lane") + op.add_argument("--id", help="this task instead of the pick (a P0 fix, a recovery)") + op.add_argument("--recovery", metavar="WHY", help="prompt gets 'Recovery: WHY' (crashed worker, one retry)") + op = osub.add_parser("post", help="after the worker's report: post-check (done: archive line, branch gone, " + "worktree clean and in master, else wf merge; wf check), leftover commit otherwise, " + "out/wf-orch.log line, out/wf-cost.log line (--agent), then the next pick or 'stop lane'") + op.add_argument("id") + op.add_argument("lane") + op.add_argument("--result", metavar="LINE", help="the report's result line (default: done if archived, else no-report)") + op.add_argument("--commit", metavar="SHA", help="the report's commit (default: master's sha when done)") + op.add_argument("--agent", metavar="ID", help="agent id from the completion notice: cost line (wf usage --log)") + op.add_argument("--duration", type=int, metavar="S", help="agent duration_ms / 1000 (default or 0: since the pick)") + op.add_argument("--no-pick", "--no-next", dest="no_pick", action="store_true", + help="no next pick, claims nothing (batch end, stop file, cross-project rollout)") + op.add_argument("--no-push", action="store_true", help="the wf merge it may run skips git push home") + + cmd("res", None, "shared memory/CPU ledger for agent jobs (wf res -h)") + cmd("cloud", None, "cloud-lane ledger (wf cloud -h)") + cmd("prep", None, "prep tasks without Done: wf prep [N|all] = wf batch --prep (wf batch -h)") + cmd("batch", None, "start an unattended batch orchestrator (claude -p) under wf res (wf batch -h)") + return ap + + +def main(argv: list[str]) -> int: + if argv[:1] == ["cloud"]: + import wf_cloud + return wf_cloud.main(argv[1:]) + if argv[:1] == ["res"]: + import wf_res + return wf_res.main(argv[1:]) + if argv[:1] == ["prep"]: + import wf_res + return wf_res.prep_main(argv[1:]) + if argv[:1] == ["batch"]: + import wf_res + return wf_res.batch_main(argv[1:]) + ap = parser() + args = ap.parse_args(argv) + if not args.command: + ap.print_usage(sys.stderr) + return 2 + try: + if getattr(args, "locks", False): + root = config.find_root(Path(args.project) if args.project else Path.cwd()) + with project_lock(root): + return args.func(args) + return args.func(args) + except Failure as e: + print(f"wf: {e}", file=sys.stderr) + return e.code + except (tasks.TaskError, config.ConfigError) as e: + print(f"wf: {e}", file=sys.stderr) + return 1 + except BrokenPipeError: + return 0 + + +if __name__ == "__main__": + sys.exit(main(sys.argv[1:])) -- cgit