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