#!/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"] if code := config.code_main(p.cfg): master = config.git_branch(code / ".git") or "master" cmd = f"wf start --worktree {code / '.worktrees' / lane} --branch {lane}/" return [head + f"work in your lane's code worktree, never on {master}:", f" {cmd} (run in {p.cfg.root}; wf there writes this 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 (marker := p.cfg.root / ".wf" / "push-failed").is_file(): out += [f"push failed: {marker.read_text().splitlines()[0]} → wf push (in {p.cfg.root})"] 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-"): item = tasks.remove(doc, id) if args.m: line = tasks.archive_line(today, item, args.m) p.new_archive = tasks.archive_prepend(p.new_archive or p.archive_text, line) out.append(f"removed: {id}" + (f" → {p.cfg.rel(p.cfg.archive)}" if args.m else "")) 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" bmain = config.git_top(p.cfg.root) or main # books repo (private repo of a split project) files = " ".join(str(f.relative_to(bmain)) 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) bmain = config.git_top(cfg.root) or main # books repo: the private repo of a split project, else main private = [str(Path(x).resolve().relative_to(bmain.resolve())) for x in getattr(args, "private", None) or []] files = [str(f.relative_to(bmain)) for f in (cfg.tasks, cfg.archive)] + private if _dirty_outside(top, []): raise Failure("worktree has uncommitted changes: commit them first") had = unmerged_count(top, master) if not branch and not had: # after a merge: notes / follow-ups written since → bookkeeping commit only if not git_run(bmain, "status", "--porcelain", "--", *files).stdout.strip(): raise Failure("worktree is on a detached HEAD: nothing to merge") if private and (r := git_run(bmain, "add", "--", *private)).returncode: raise Failure(f"git add failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") r = git_run(bmain, "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() if had or bmain == main else "-") 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 private and (r := git_run(bmain, "add", "--", *private)).returncode: raise Failure(f"git add failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") if git_run(bmain, "status", "--porcelain", "--", *files).stdout.strip(): msg = (getattr(args, "private_msg", None) if private else None) or args.m \ or (f"{branch.rsplit('/', 1)[-1]} done" if branch else "bookkeeping") if bmain != main and args.merged_sha != "-": msg += f" (code {args.merged_sha})" r = git_run(bmain, "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)}") args.books_sha = git_run(bmain, "rev-parse", "--short", "HEAD").stdout.strip() git_run(top, "switch", "-q", "--detach", master) if branch: git_run(top, "branch", "-q", "-d", branch) args.push_rc = 0 if not args.no_push: args.push_rc, lines = run_push(cfg, [bmain, main]) out += lines out.append(f"merged {branch or 'detached HEAD'} into {master}") print("\n".join(out)) return 0 def run_push(cfg: config.Config, repos: list[Path]) -> tuple[int, list[str]]: """Push after a merge: workflow.toml push lines (cwd = root, env WF_MAIN; first nonzero stops and writes .wf/push-failed), else _push_home on each repo. (exit code, lines to print).""" marker = cfg.root / ".wf" / "push-failed" if not cfg.push: pushed = [r for r in dict.fromkeys(repos) if _push_home(r)] return 0, ["pushed home"] if pushed else [] env = {**os.environ, "WF_MAIN": str(cfg.root)} for line in cfg.push: r = subprocess.run(line, shell=True, cwd=cfg.root, env=env, capture_output=True, text=True) if r.returncode: tail = (r.stdout + r.stderr).strip().splitlines()[-5:] wf_folder(cfg, "orch") # creates .wf with its .gitignore marker.write_text(f"{line}\nexit {r.returncode}\n" + "\n".join(tail) + "\n") return r.returncode, [f"not pushed (exit {r.returncode}): {line}", *(f" {t}" for t in tail), "rerun: wf push"] marker.unlink(missing_ok=True) return 0, [f"pushed ({len(cfg.push)} push command{'s' * (len(cfg.push) != 1)})"] def cmd_push(args) -> int: """Run the push of a merge again (workflow.toml push, else push home), from the project or a code worktree.""" cfg = config.load_at(Path(args.project or Path.cwd()).resolve()) repos = [r for r in (config.git_top(cfg.root), config.code_main(cfg)) if r] rc, lines = run_push(cfg, repos) print("\n".join(lines) or "nothing to push (no push key, no remote home)") return rc 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", "-uall").stdout.splitlines(): rel = line[3:].split(" -> ")[-1].strip('"').rstrip("/") path = (top / rel).resolve() if rel == config.WF_HOME: # wf start's pointer in a code worktree, never committed continue 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)" "\n- coalesced gate (wf res wait prints 'covers …; red: culprit is any commit in A^..B'): bisect that range for the" " culprit, else the P0 fix task's Steps name the range") 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()] btop = (config.git_top(cfg.root) or cfg.root).resolve() split = bool(wt) and btop != wt[2].resolve() # code worktree of a split project: books in another repo proot, priv, code = cfg.root.resolve(), [], [] for f, path in zip(args.paths, paths): if split and (path == proot or proot in path.parents): if not path.exists() and git_run(btop, "ls-files", "--error-unmatch", "--", str(path)).returncode: raise Failure(f"path {f} does not exist and is not tracked: nothing to commit") priv.append(path) continue code.append(path) if path != top and top not in path.parents: if split: raise Failure(f"path {f} is in neither repo ({top}, {cfg.root}): commit it in its own repo " "first, then wf finish without it") 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 not split 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") paths = code 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) push_rc = 0 if not wt and not args.no_push: push_rc, lines = run_push(cfg, [top]) if lines: print("\n".join(lines), flush=True) if wt: ns = argparse.Namespace(project=args.project, m=None, no_push=args.no_push, private=priv, private_msg=args.commit) rc = cmd_merge(ns) if not rc and getattr(ns, "merged_sha", ""): tool = getattr(args, "tool_commit", None) books_sha = f" books {ns.books_sha}" if split else "" failed = f" push-failed {ns.push_rc}" if getattr(ns, "push_rc", 0) else "" print(f"report: commit {ns.merged_sha}{books_sha}" + (f" tool {tool}" if tool else "") + failed, 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() wt = config.linked_worktree(start) if not wt: raise Failure("wip runs inside a lane worktree") top = Path(git_run(start, "rev-parse", "--show-toplevel").stdout.strip() or start) p = load_project(args) id = p.resolve(args.id) btop = (config.git_top(p.cfg.root) or p.cfg.root).resolve() proot, priv, paths = p.cfg.root.resolve(), [], [] for f in args.paths: # split project: paths under the private root are committed there path = (start / f).resolve() (priv if btop != wt[2].resolve() and (path == proot or proot in path.parents) else paths).append(str(path)) if priv: r = git_run(btop, "add", "--", *priv) if r.returncode: raise Failure(f"git add failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") r = git_run(btop, "commit", "-q", "-m", f"{id} wip: {args.m}", "--", *priv) 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, btop) for c in priv), flush=True) 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) 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) repo = config.code_main(p.cfg) or main # split project: worktree + branch in the code repo master = config.git_branch(repo / ".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(repo, "rev-parse", "--verify", "-q", f"refs/heads/{b}").returncode extra = [] if not wt.exists(): if has_branch: git_ok(repo, "worktree", "add", "-q", str(wt), b) how = f"new, existing branch {b}: earlier WIP, read the notes" else: git_ok(repo, "worktree", "add", "-q", str(wt), "-b", b, master) how = f"new, branch {b} from {master}" else: own = git_run(wt, "rev-parse", "--path-format=absolute", "--git-common-dir").stdout.strip() want = git_run(repo, "rev-parse", "--path-format=absolute", "--git-common-dir").stdout.strip() if not own or Path(own).resolve() != Path(want).resolve(): raise Failure(f"{shown} is not a worktree of {repo}" + (f" (it belongs to {Path(own).resolve().parent})" if own else "") + f": pass a worktree path of {repo} (e.g. {repo}/.worktrees/<lane>)") 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) if repo != main: (wt / config.WF_HOME).write_text(f"{p.cfg.root}\n") 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 PARK = ("handback", "awaiting", "needs-owner", "post-check-red") # block only that task (park_task), lane keeps picking 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() base = config.code_main(cfg) or main n = 1 while not worktree_free(wt := base / ".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) ARCHIVED = "archived, nothing to mark" def park_task(root: Path, id: str, final: str, words: list[str]) -> str: """Block only this task after a handback/awaiting/needs-owner/post-check-red: needs-owner → After: its open Needs-human item (status cleared); awaiting → blocked on its open a-id; else Sessions: owner (+ status cleared on a handback) and a note. Returns what it did (alert text).""" p = load_project(argparse.Namespace(project=str(root)), write=True) if id not in p.doc.ids(): return ARCHIVED item = p.doc.item(id) why = " ".join(words[1:]) or "-" open_a = {i.id for s in p.doc.sections if s.key == "awaiting" for i in s.items} a = next((w for w in words[1:] if w.strip("[]") in open_a), None) if final == "awaiting" else None open_h = {i.id for i in p.doc.section("human").items} h = next((w.strip("[]") for w in words[1:] if w.strip("[]") in open_h), None) \ if final == "needs-owner" else None if h: # owner-only step: After: it, picked again once the owner's wf done archives it if h not in item.after: tasks.set_fields(p.doc, id, after=item.after + [h]) if item.status: tasks.set_status(p.doc, id, None) done = f"after {h}" elif (item.status or "").startswith("blocked:"): done = f"already {item.status}" elif a: tasks.set_status(p.doc, id, f"blocked: [[{a.strip('[]')}]]") done = f"blocked on {a.strip('[]')}" else: tasks.set_fields(p.doc, id, sessions="owner") done = "Sessions: owner" if final == "handback" and item.status: tasks.set_status(p.doc, id, None) done += ", status cleared" tasks.add_note(p.doc, id, f"orch {datetime.date.today().isoformat()}: {final} ({why}) → {done}; " "lane kept picking, owner decides") p.save(False) if done.endswith("status cleared") or h: unclaim(p.cfg, [id]) return done def cloud_park(root: Path, id: str, final: str, words: list[str]) -> str: """Lane cloud after a handback/lost/awaiting/needs-owner (wf cloud pull already noted/blocked it): awaiting/needs-owner → park_task; else Cloud: no (the cloud lane never re-sends it; local lanes may retry with pull's Recovery note) + status cleared + note. Returns what it did (alert text).""" if final in ("awaiting", "needs-owner"): return park_task(root, id, final, words) p = load_project(argparse.Namespace(project=str(root)), write=True) if id not in p.doc.ids(): return ARCHIVED item, done = p.doc.item(id), "Cloud: no" if item.cloud != "no": tasks.set_fields(p.doc, id, cloud="no") if (item.status or "").startswith("in progress"): tasks.set_status(p.doc, id, None) done += ", status cleared" tasks.add_note(p.doc, id, f"orch {datetime.date.today().isoformat()}: cloud {final} " f"({' '.join(words[1:]) or '-'}) → {done}; lane cloud kept picking, local lanes may retry") p.save(False) unclaim(p.cfg, [id]) return done + " (local lanes may retry)" def orch_post(main: Path, p: Project, args) -> int: cfg, id, lane = p.cfg, args.id, args.lane rec = orch_record(cfg, id) repo = config.code_main(cfg) or main wt = Path(rec.get("worktree") or repo / ".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())) line = (args.result or "").strip() if line[:7].lower() == "result:": line = line[7:].strip() # literal report line 'result: done' words = line.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(repo / ".git") or "master" files = [str(f.relative_to(main)) for f in (cfg.tasks, cfg.archive)] problems, out = [], [] split = repo.resolve() != main.resolve() def report_commit() -> str: # after any merge below; split: a private sha → the code repo's master c = args.commit or (git_run(repo, "rev-parse", "--short", master).stdout.strip() if outcome.startswith("done") else "-") if split and c != "-" and git_run(repo, "cat-file", "-e", f"{c}^{{commit}}").returncode: fixed = git_run(repo, "rev-parse", "--short", master).stdout.strip() out.append(f"commit {c} not in code repo {repo} (a private sha?): using its {master} {fixed}") c = fixed return c commit, parked = None, None raised = bool(item and rec.get("model") and not outcome.startswith("done") and tasks.MODELS.index(item.model) > tasks.MODELS.index(rec["model"])) 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(repo, "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]}") commit = report_commit() if not problems and git_run(main, "status", "--porcelain", "--", *files).stdout.strip(): msg = f"{id} done" + (f" (code {commit})" if split and commit != "-" else " (orchestrator)") r = git_run(main, "commit", "-q", "-m", msg, "--", *files) out.append("committed leftover " + " ".join(files) if not r.returncode else f"leftover commit failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") if problems or not outcome.startswith("done"): if lane != CLOUD and (problems or outcome in PARK and not raised): parked = park_task(cfg.root, id, "post-check-red" if problems else outcome, words) elif lane == CLOUD and not problems and outcome in PARK + ("lost",) and not raised: # no re-send loop parked = cloud_park(cfg.root, id, outcome, words) if (not problems or parked and parked != ARCHIVED) \ and git_run(main, "status", "--porcelain", "--", *files).stdout.strip(): r = git_run(main, "commit", "-q", "-m", f"{id} {'post-check-red' if problems else 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]}") final = "post-check-red" if problems else ("model-raised" if raised else outcome) if outcome.startswith("done") and not problems and (cfg.root / ".wf" / "push-failed").is_file(): final = "push-failed" 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 = commit or report_commit() 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 parked: print(f"alert: {id} {final} → tell the owner ({parked}); lane {lane} keeps picking") elif final not in GO_ON and final != "model-raised": print(f"stop lane {lane}: {final} → tell the owner" + (f" (wf push in {cfg.root})" if final == "push-failed" else "") + (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("push", cmd_push, "rerun a merge's push: workflow.toml push lines, else git push home") sp.set_defaults(locks=True) 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); awaiting/needs-owner/handback/post-check-red park only the task (blocked / After: h-task / Sessions: owner (lane cloud: Cloud: no) + note, 'alert:' line); 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:]))