From 81d4e80fd5aabe4e80f58e960affa795cf7d34ec Mon Sep 17 00:00:00 2001 From: godosa Date: Wed, 7 Oct 2026 07:27:17 +0200 Subject: workflow: initial public history --- wflib/__init__.py | 0 wflib/areas.py | 199 +++++++++++ wflib/check.py | 378 +++++++++++++++++++++ wflib/cloud.py | 442 ++++++++++++++++++++++++ wflib/config.py | 225 +++++++++++++ wflib/lanes.py | 297 ++++++++++++++++ wflib/ledgers.py | 35 ++ wflib/migrate.py | 237 +++++++++++++ wflib/refs.py | 107 ++++++ wflib/res.py | 993 ++++++++++++++++++++++++++++++++++++++++++++++++++++++ wflib/search.py | 82 +++++ wflib/tasks.py | 728 +++++++++++++++++++++++++++++++++++++++ wflib/usage.py | 375 +++++++++++++++++++++ 13 files changed, 4098 insertions(+) create mode 100644 wflib/__init__.py create mode 100644 wflib/areas.py create mode 100644 wflib/check.py create mode 100644 wflib/cloud.py create mode 100644 wflib/config.py create mode 100644 wflib/lanes.py create mode 100644 wflib/ledgers.py create mode 100644 wflib/migrate.py create mode 100644 wflib/refs.py create mode 100644 wflib/res.py create mode 100644 wflib/search.py create mode 100644 wflib/tasks.py create mode 100644 wflib/usage.py (limited to 'wflib') diff --git a/wflib/__init__.py b/wflib/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/wflib/areas.py b/wflib/areas.py new file mode 100644 index 0000000..ce23d12 --- /dev/null +++ b/wflib/areas.py @@ -0,0 +1,199 @@ +"""Area notes (`## Areas` in the project's CLAUDE.md): code-map anchors, staleness, Checked marker.""" +from __future__ import annotations + +import fnmatch +import re +import subprocess +from dataclasses import dataclass, replace +from pathlib import Path + + +@dataclass(frozen=True) +class Area: + name: str + slug: str + anchors: list[str] + paths: list[str] + checked: str | None + line: int + + +def slug(name: str) -> str: + return re.sub(r"[^a-z0-9]+", "-", name.lower()).strip("-") + + +def parse(text: str) -> list[Area]: + out: list[Area] = [] + in_areas = False + cur: dict | None = None + + def flush(): + if cur: + out.append(Area(cur["name"], slug(cur["name"]), cur["anchors"], cur["paths"], cur["checked"], cur["line"])) + + for n, line in enumerate(text.replace("\r\n", "\n").split("\n"), 1): + if line.startswith("## "): + flush() + cur = None + in_areas = line[3:].strip().lower() == "areas" + elif in_areas and line.startswith("### "): + flush() + cur = {"name": line[4:].strip(), "anchors": [], "paths": [], "checked": None, "line": n} + elif cur is not None: + if line.startswith("- Code map:"): + cur["anchors"] = re.findall(r"`([^`]+)`", line) + elif line.startswith("- Paths:"): + cur["paths"] = line[len("- Paths:"):].split() + elif line.startswith("- Checked:"): + words = line[len("- Checked:"):].split() + cur["checked"] = words[0] if words and re.fullmatch(r"[0-9a-fA-F]{4,40}", words[0]) else None + flush() + return out + + +def matches(area: Area, text: str, skip: set[str] = frozenset()) -> bool: + """text (a task) names the area, one of its anchors or Paths (whole word / path, case-insensitive name); + Paths in `skip` (shared by several areas) never match.""" + def word(tok: str, flags=0) -> bool: + return re.search(r"(? set[str]: + """Paths (normalised) listed by more than one area.""" + seen: dict[str, int] = {} + for a in found: + for p in {p.removeprefix("./").rstrip("/") for p in a.paths}: + seen[p] = seen.get(p, 0) + 1 + return {p for p, n in seen.items() if n > 1} + + +def block(text: str, name: str) -> str: + """The area's notes: its ### heading up to the next heading, trailing blank lines dropped; unknown → "".""" + lines = text.replace("\r\n", "\n").split("\n") + in_areas, out = False, None + for line in lines: + if line.startswith("#"): + if out is not None: + break + if line.startswith("## "): + in_areas = line[3:].strip().lower() == "areas" + elif in_areas and line.startswith("### ") and line[4:].strip() == name: + out = [line] + elif out is not None: + out.append(line) + while out and not out[-1].strip(): + out.pop() + return "\n".join(out or []) + + +def covers(paths: list[str], file: str) -> bool: + """file (project-relative) matches one of an area's Paths: dir prefix, exact file or glob.""" + for p in paths: + p = p.removeprefix("./").rstrip("/") + if file == p or file.startswith(p + "/") or fnmatch.fnmatchcase(file, p): + return True + return False + + +def uncovered(files: list[str], found: list[Area], ignore: list[str]) -> list[tuple[str, int]]: + """Folders (≤ 2 levels) of files no area's Paths covers → [(folder, file count)], most files first. + Root files, dot folders and `ignore` prefixes skipped; no area with Paths → [].""" + if not any(a.paths for a in found): + return [] + count: dict[str, int] = {} + for f in files: + parts = f.split("/")[:-1] + if not parts or any(x.startswith(".") for x in parts) or covers(ignore, f): + continue + if any(covers(a.paths, f) for a in found): + continue + d = "/".join(parts[:2]) + count[d] = count.get(d, 0) + 1 + return sorted(count.items(), key=lambda kv: (-kv[1], kv[0])) + + +def changed_files(here: Path, base: str | None) -> list[str]: + """Files (relative to here) changed since merge-base(HEAD, base) incl. uncommitted + untracked; + base None → HEAD (uncommitted only).""" + ref = "HEAD" + if base: + r = _git(here, "merge-base", "HEAD", base) + if r.returncode == 0 and r.stdout.strip(): + ref = r.stdout.strip() + out = [] + for args in (("diff", "--name-only", "--relative", ref), ("ls-files", "--others", "--exclude-standard")): + r = _git(here, *args) + if r.returncode == 0: + out += [l for l in r.stdout.splitlines() if l] + return sorted(set(out)) + + +def with_anchor_paths(root: Path, area: Area) -> Area: + """No Paths → the folders of the files holding its anchors (git grep) stand in for them.""" + if area.paths or not area.anchors: + return area + args = [x for tok in area.anchors for x in ("-e", tok)] + r = _git(root, "grep", "-l", "-F", *args) + files = r.stdout.splitlines() if r.returncode == 0 else [] + return replace(area, paths=sorted({f.rsplit("/", 1)[0] if "/" in f else f for f in files})) + + +def _git(root: Path, *args: str) -> subprocess.CompletedProcess: + return subprocess.run(["git", "-C", str(root), *args], capture_output=True, text=True) + + +def missing(root: Path, area: Area) -> list[str]: + out = [] + for tok in area.anchors: + r = _git(root, "grep", "-F", "-q", "-e", tok, *(["--", *area.paths] if area.paths else [])) + if r.returncode != 0: + out.append(tok) + return out + + +def commits_since(root: Path, area: Area) -> int | None: + if not area.checked or _git(root, "cat-file", "-e", f"{area.checked}^{{commit}}").returncode != 0: + return None + r = _git(root, "rev-list", "--count", f"{area.checked}..HEAD", *(["--", *area.paths] if area.paths else [])) + return int(r.stdout.strip()) if r.returncode == 0 and r.stdout.strip().isdigit() else None + + +def status(root: Path, area: Area, limit: int) -> tuple[bool, list[str], int | None]: + miss = missing(root, area) + commits = commits_since(root, area) + return (bool(miss) or (commits is not None and commits >= limit)), miss, commits + + +def mark(text: str, name: str, sha: str) -> str: + lines = text.split("\n") + in_areas = False + start = None + for i, line in enumerate(lines): + if line.startswith("## "): + if start is not None: + break + in_areas = line[3:].strip().lower() == "areas" + elif in_areas and line.startswith("### "): + if start is not None: + break + if line[4:].strip() == name: + start = i + if start is None: + return text + end = start + 1 + while end < len(lines) and not lines[end].startswith("#"): + end += 1 + block_end = end + while block_end > start + 1 and not lines[block_end - 1].strip(): + block_end -= 1 + new = f"- Checked: {sha}" + for i in range(start + 1, block_end): + if lines[i].startswith("- Checked:"): + lines[i] = new + return "\n".join(lines) + lines.insert(block_end, new) + return "\n".join(lines) diff --git a/wflib/check.py b/wflib/check.py new file mode 100644 index 0000000..f86b2fe --- /dev/null +++ b/wflib/check.py @@ -0,0 +1,378 @@ +"""Validation: everything `wf check` reports.""" +from __future__ import annotations + +import re +import subprocess +import time +from dataclasses import dataclass +from pathlib import Path + +from . import areas, refs, tasks +from .config import NAME, Config + +NUMBERED_RE = re.compile(r"^\d+\. ") +NUMBER_REF_RE = re.compile(r"(? tuple[str, str, str]: + return (self.path, self.id, self.message) + + def __str__(self) -> str: + where = self.path + (f":{self.line}" if self.line else "") + return f"{where}: " + (f"{self.id}: " if self.id else "") + self.message + + +def _read(path: Path) -> str: + return path.read_text(encoding="utf-8", errors="replace").replace("\r\n", "\n") + + +def _prose_lines(text: str): + """(1-based line number, line) outside code fences.""" + fenced = False + for n, line in enumerate(text.split("\n"), 1): + if refs.FENCE_RE.match(line): + fenced = not fenced + elif not fenced: + yield n, line + + +def _link_problems(path: str, text: str, known: set[str]) -> list[Problem]: + return [Problem(path, n, f"[[{id}]]", "no such id in TASKS or archive") + for n, line in _prose_lines(text) + for id in dict.fromkeys(refs.links(line)) + if TASK_LINK_RE.match(id) and id not in known] + + +def _cycle(start: str, edges: dict[str, list[str]]) -> list[str] | None: + def walk(node: str, path: list[str]) -> list[str] | None: + for nxt in edges.get(node, []): + if nxt == start: + return path + [nxt] + if nxt not in path and (found := walk(nxt, path + [nxt])): + return found + return None + return walk(start, [start]) + + +def _index_sections(cfg: Config) -> list[refs.DocSection] | None: + """Headings under the anchors index section; None when no [anchors].""" + if not cfg.anchors_index or not cfg.anchors_index.is_file(): + return None + found = refs.sections(_read(cfg.anchors_index)) + if not cfg.anchors_section: + return [s for s in found if s.level] + top = next((s for s in found if s.level and s.heading == cfg.anchors_section), None) + if top is None: + return [] + return [s for s in found if s.level > top.level and top.start < s.start < top.end] + + +def check_tasks(text: str, archive_text: str, cfg: Config) -> tuple[list[Problem], list[Problem]]: + path = cfg.rel(cfg.tasks) + errors: list[Problem] = [] + warnings: list[Problem] = [] + doc = tasks.parse(text) + archived = tasks.archive_ids(archive_text) + + def err(line, id, message): + errors.append(Problem(path, line, id, message)) + + for key, heading in tasks.SECTIONS.items(): + if not any(s.key == key for s in doc.sections): + err(None, "", f"no '## {heading}' section") + + errors += _link_problems(path, text, doc.ids() | archived) + + index = _index_sections(cfg) + index_anchors = {a for s in index for a in s.anchors} if index is not None else None + open_awaiting = {i.id for s in doc.sections if s.key == "awaiting" for i in s.items} + deferred = {i.id for s in doc.sections if s.key == "deferred" for i in s.items} + edges = {i.id: [a for a in i.after if a in doc.ids()] for i in doc.all_items()} + seen: dict[str, int] = {} + in_cycle: set[str] = set() + + for section in doc.sections: + if section.key is None: + continue + for n, line in enumerate(section.prefix, section.line + 1): + if NUMBERED_RE.match(line): + err(n, "", "old numbered item (wf migrate)") + for n, line in enumerate(section.suffix, section.suffix_line or 0): + if NUMBERED_RE.match(line): + err(n, "", "old numbered item (wf migrate)") + elif line.startswith(tasks.ITEM_START): + err(n, tasks.parse_header(line).id, + f"item after a flush-left prose line (line {section.suffix_line}): indent or move the prose") + for at, item in enumerate(section.items): + if item.error: + err(item.line, item.id, item.error) + continue + if not tasks.ID_RE.match(item.id): + err(item.line, item.id, "bad id (want t-… or a-…, lowercase a-z 0-9 -)") + continue + if item.id in seen: + err(item.line, item.id, f"duplicate id (also line {seen[item.id]})") + seen.setdefault(item.id, item.line) + if item.id in archived: + err(item.line, item.id, "id already used in the archive (ids are never reused)") + if item.id.startswith("a-") != (section.key == "awaiting"): + err(item.line, item.id, f"{item.id[:2]} item in {section.heading} ({KINDS})") + continue + if section.key == "awaiting": + continue + if item.prio is None: + err(item.line, item.id, "no priority [P0]-[P3]") + elif item.prio > 3: + err(item.line, item.id, f"priority P{item.prio} (want P0-P3)") + if item.effort is None: + err(item.line, item.id, f"no effort (want {EFFORTS})") + elif item.effort not in tasks.EFFORTS: + err(item.line, item.id, f"effort '{item.effort}' (want {EFFORTS})") + if item.status and item.status.startswith("blocked:") and item.blocked_on not in open_awaiting: + err(item.line, item.id, f"blocked on '{item.blocked_on}', which is not an open Awaiting item") + if item.status and re.fullmatch(r"in progress:\s*", item.status): + warnings.append(Problem(path, item.line, item.id, "in progress without a branch or note")) + if item.id not in in_cycle and (cycle := _cycle(item.id, edges)): + in_cycle.update(cycle) + err(item.line, item.id, "After: cycle " + " → ".join(cycle)) + later = {o.id for o in section.items[at + 1:]} + for dep in item.after: + if dep in later: + err(item.line, item.id, f"placed before '{dep}', which it is After:") + for label, pattern in (("After", tasks.AFTER_RE), ("Slices", tasks.SLICES_RE)): + if (i := item._line(pattern)) is None: + continue + rest = tasks.LINK_RE.sub(" ", pattern.match(item.body[i]).group(1)) + for word in re.split(r"[\s,;]+", rest): + if tasks.ID_RE.match(word): + err(item.line + 1 + i, item.id, f"{label}: '{word}' is not a link (want [[{word}]])") + if (i := item._line(tasks.AFTER_RE)) is not None: + rest = tasks.LINK_RE.sub(" ", tasks.AFTER_RE.match(item.body[i]).group(1)) + if re.search(r"[A-Za-z0-9]", rest): + warnings.append(Problem(path, item.line + 1 + i, item.id, + "After: line has prose; every [[id]] in it is a dependency")) + for dep in item.after: + if dep in deferred and section.key != "deferred": + warnings.append(Problem(path, item.line + 1 + i, item.id, + f"After: '{dep}' is deferred (never runs; blocks this task)")) + models = [(n, m.group(1).split()) for n, l in enumerate(item.body) if (m := tasks.MODEL_RE.match(l))] + for n, words in models[:1]: + if not words or words[0] not in tasks.MODELS: + err(item.line + 1 + n, item.id, + f"Model '{' '.join(words)}' (want {', '.join(tasks.MODELS)})") + elif len(words) > 1: + warnings.append(Problem(path, item.line + 1 + n, item.id, f"Model line '{' '.join(words)}': " + f"write 'Model: {words[0]}' (wf set --model)")) + if len(models) > 1: + err(item.line + 1 + models[1][0], item.id, "two Model lines") + sess = [(n, m.group(1).split()) for n, l in enumerate(item.body) if (m := tasks.SESSIONS_RE.match(l))] + for n, words in sess[:1]: + if not words or words[0] not in tasks.SESSIONS: + err(item.line + 1 + n, item.id, + f"Sessions '{' '.join(words)}' (want {', '.join(tasks.SESSIONS)})") + if len(sess) > 1: + err(item.line + 1 + sess[1][0], item.id, "two Sessions lines") + clouds = [(n, m.group(1).split()) for n, l in enumerate(item.body) if (m := tasks.CLOUD_RE.match(l))] + for n, words in clouds[:1]: + if not words or words[0] not in tasks.CLOUDS: + err(item.line + 1 + n, item.id, f"Cloud '{' '.join(words)}' (want {', '.join(tasks.CLOUDS)})") + if len(clouds) > 1: + err(item.line + 1 + clouds[1][0], item.id, "two Cloud lines") + if item.interactive: + warnings.append(Problem(path, item.line, item.id, "'interactive' flag: write 'Sessions: owner' " + f"(wf set {item.id} --sessions owner)")) + if section.key == "pending" and item.human_done_match and item.sessions != "owner": + warnings.append(Problem(path, item.line, item.id, f"not runner-ready: Done reads as human action " + f"('{item.human_done_match}'); reword or set Sessions: owner")) + ref_at = item._line(tasks.REF_RE) + ref_line = item.line + 1 + ref_at if ref_at is not None else item.line + for target, anchor in item.refs: + shown = target + (f"#{anchor}" if anchor else "") + file = cfg.root / target + if not file.exists(): + err(ref_line, item.id, f"Ref '{shown}' does not exist") + elif anchor and file.is_file(): + if anchor not in refs.anchors(_read(file)): + err(ref_line, item.id, f"Ref '{shown}': no such anchor") + elif index_anchors is not None and file.resolve() == cfg.anchors_index.resolve() \ + and anchor not in index_anchors: + err(ref_line, item.id, f"Ref '{shown}' is not a heading under '{cfg.anchors_section}'") + + for n, line in _prose_lines(text): + for m in NUMBER_REF_RE.finditer(refs.strip_code(line)): + warnings.append(Problem(path, n, "", f"'{m.group(0)}': number ref (tasks have ids: [[t-…]])")) + return errors, warnings + + +def _anchor_problems(cfg: Config) -> list[Problem]: + index = _index_sections(cfg) + if index is None: + return [] + out: list[Problem] = [] + path = cfg.rel(cfg.anchors_index) + lines = _read(cfg.anchors_index).split("\n") + where = f"'{cfg.anchors_section}'" if cfg.anchors_section else "the index" + seen: dict[str, int] = {} + for s in index: + slug = refs.slugify(s.heading) + if slug in seen: + out.append(Problem(path, s.start + 1, "", f"two {where} headings give anchor '{slug}' (also line {seen[slug]})")) + seen.setdefault(slug, s.start + 1) + for n in range(s.start, s.end): + for target in MD_LINK_RE.findall(refs.strip_code(lines[n])): + if re.match(r"^[a-z][a-z0-9+.-]*:", target): + continue + file_part, _, frag = target.partition("#") + file = (cfg.anchors_index.parent / file_part) if file_part else cfg.anchors_index + if not file.exists(): + out.append(Problem(path, n + 1, "", f"link '{target}': file does not exist")) + elif frag and file.suffix == ".md" and frag not in refs.anchors(_read(file)): + out.append(Problem(path, n + 1, "", f"link '{target}': no such anchor")) + if cfg.anchors_specs and cfg.anchors_specs.is_dir(): + known = {a for s in index for a in s.anchors} + for spec in sorted(cfg.anchors_specs.rglob("*.md")): + for n, line in _prose_lines(_read(spec)): + for anchor in refs.ANCHOR_RE.findall(line): + if anchor not in known: + out.append(Problem(cfg.rel(spec), n, "", + f"explicit anchor '{anchor}' has no heading under {where} in {path}")) + return out + + +def doc_files(cfg: Config) -> list[Path]: + out: list[Path] = [] + for d in cfg.docs: + if d.is_file(): + out.append(d) + elif d.is_dir(): + out += sorted(p for p in d.rglob("*.md") if p.is_file()) + skip = {cfg.tasks.resolve(), cfg.archive.resolve()} + return [p for p in dict.fromkeys(out) if p.resolve() not in skip] + + +def _age_days(root: Path, path: Path, line: int) -> float | None: + try: + out = subprocess.run(["git", "-C", str(root), "blame", "-L", f"{line},{line}", "--porcelain", "--", str(path)], + capture_output=True, text=True, timeout=10) + except (OSError, subprocess.SubprocessError): + return None + m = re.search(r"^author-time (\d+)$", out.stdout, re.M) if out.returncode == 0 else None + if not m or re.match(r"^0{40}", out.stdout): + return None + return (time.time() - int(m.group(1))) / 86400 + + +def _stale_awaiting(cfg: Config, text: str) -> list[Problem]: + doc = tasks.parse(text) + waited = {i.blocked_on for i in doc.all_items()} | {l for i in doc.all_items() if i.id.startswith("t-") + for l in refs.links("\n".join(i.lines()))} + out = [] + for s in doc.sections: + if s.key != "awaiting": + continue + for item in s.items: + if item.id in waited or item.error: + continue + age = _age_days(cfg.root, cfg.tasks, item.line) + if age is not None and age > STALE_DAYS: + out.append(Problem(cfg.rel(cfg.tasks), item.line, item.id, + f"waiting {int(age)} days, no task references it")) + return out + + +def _git(root: Path, *args: str) -> str | None: + try: + out = subprocess.run(["git", "-C", str(root), *args], capture_output=True, text=True, timeout=10) + except (OSError, subprocess.SubprocessError): + return None + return out.stdout.strip() if out.returncode == 0 else None + + +def merged_tool_worktrees(tool_root: Path) -> list[Problem]: + """Worktrees of the wf tool repo whose branch is behind master (merged, leftover after release). + A branch equal to master, or a worktree created < 2h ago, is not flagged.""" + listing = _git(tool_root, "worktree", "list", "--porcelain") + master = _git(tool_root, "rev-parse", "master") + if not listing or not master: + return [] + out = [] + for block in listing.split("\n\n"): + path = branch = head = None + for l in block.splitlines(): + if l.startswith("worktree "): + path = l[9:] + elif l.startswith("HEAD "): + head = l[5:] + elif l.startswith("branch refs/heads/"): + branch = l[18:] + if not path or not branch or branch == "master" or head == master: + continue + try: # created < 2h ago: a running worker's fresh worktree, not a leftover + if time.time() - (Path(path) / ".git").stat().st_mtime < 7200: + continue + except OSError: + pass + if _git(tool_root, "merge-base", "--is-ancestor", branch, "master") is not None: + out.append(Problem("wf-tool", None, "", f"tool worktree '{path}' (branch {branch}) is merged into " + f"master: git worktree remove it, git branch -d {branch}")) + return out + + +def _area_problems(cfg: Config) -> list[Problem]: + file = cfg.areas_file + if not file.is_file(): + return [] + out = [] + for a in areas.parse(_read(file)): + for tok in areas.missing(cfg.code, a): + out.append(Problem(cfg.rel(file), a.line, "", f"area {a.name}: anchor {tok} not found")) + if a.checked and areas.commits_since(cfg.code, a) is None: + out.append(Problem(cfg.rel(file), a.line, "", f"area {a.name}: Checked {a.checked} unknown")) + return out + + +def check(cfg: Config, tasks_text: str | None = None, archive_text: str | None = None, + slow: bool = True) -> tuple[list[Problem], list[Problem]]: + """All errors and warnings of the project. `tasks_text` / `archive_text` + replace what is on disk (to judge a change before it is written); + `slow=False` leaves out the git-based warnings.""" + errors: list[Problem] = [] + for name, path in (("tasks", cfg.tasks), ("archive", cfg.archive)): + given = tasks_text if name == "tasks" else archive_text + if given is None and not path.is_file(): + errors.append(Problem(NAME, None, "", f"{name} '{cfg.rel(path)}' does not exist")) + for name, path in [("docs", d) for d in cfg.docs] + [("anchors.index", cfg.anchors_index), + ("anchors.specs", cfg.anchors_specs)]: + if path is not None and not path.exists(): + errors.append(Problem(NAME, None, "", f"{name} '{cfg.rel(path)}' does not exist")) + if tasks_text is None and not cfg.tasks.is_file(): + return errors, [] + text = (_read(cfg.tasks) if tasks_text is None else tasks_text).replace("\r\n", "\n") + archive = archive_text if archive_text is not None else (_read(cfg.archive) if cfg.archive.is_file() else "") + found, warnings = check_tasks(text, archive, cfg) + errors += found + known = tasks.parse(text).ids() | tasks.archive_ids(archive) + for file in doc_files(cfg): + errors += _link_problems(cfg.rel(file), _read(file), known) + errors += _anchor_problems(cfg) + if slow: + warnings += _stale_awaiting(cfg, text) + warnings += _area_problems(cfg) + order: dict[str, int] = {} + for p in errors + warnings: + order.setdefault(p.path, len(order)) + by_place = lambda p: (order[p.path], p.line or 0) + return sorted(errors, key=by_place), sorted(warnings, key=by_place) diff --git a/wflib/cloud.py b/wflib/cloud.py new file mode 100644 index 0000000..5878b87 --- /dev/null +++ b/wflib/cloud.py @@ -0,0 +1,442 @@ +"""Cloud-lane ledger: pure logic (spec docs/specs cloud-lane-design §4.4). IO: wf_cloud.py.""" +from __future__ import annotations + +import base64 +import binascii +import dataclasses +import datetime as dt +import gzip +import hashlib +import json +import re + +from . import tasks as TK, usage as U + +OVERHEAD = 1.15 # title generation, setup: what the session can't see +LOST_AFTER = dt.timedelta(hours=24) +ENDED = ("done", "awaiting", "handback", "lost") + + +class CloudError(Exception): + pass + + +def new() -> dict: + return {"budget": 240.0, "spent": 0.0, "reserve_per_task": 4.0, "max_parallel": 3, "entries": []} + + +def loads(text: str) -> dict: + if not text.strip(): + return new() + try: + led = json.loads(text) + led["budget"], led["spent"], led["entries"] + except (ValueError, KeyError, TypeError) as e: + raise CloudError(f"cloud ledger unreadable: {e}") + return {**new(), **led} + + +def dumps(led: dict) -> str: + return json.dumps(led, indent=1) + "\n" + + +def running(led: dict) -> int: + return sum(e["state"] == "running" for e in led["entries"]) + + +def balance(led: dict) -> float: + return led["budget"] - led["spent"] - led["reserve_per_task"] * running(led) + + +def refusal(led: dict) -> str | None: + """One line why `send` must refuse (exit 3), or None.""" + bal, res = balance(led), led["reserve_per_task"] + if bal < res: + return f"cloud busy: balance ${bal:.2f} < reserve ${res:.2f}" + if running(led) >= led["max_parallel"]: + return f"cloud busy: {running(led)} running >= max_parallel {led['max_parallel']}" + return None + + +def add(led: dict, id: str, project: str, sid: str, model: str, now: dt.datetime) -> dict: + e = {"id": id, "project": project, "sid": sid, "model": model, "sent": now.isoformat(), + "state": "running", "usd": None, "usd_source": None} + led["entries"].append(e) + return e + + +def find(led: dict, sid: str) -> dict: + for e in led["entries"]: + if e["sid"] == sid: + return e + raise CloudError(f"no cloud entry with sid {sid}") + + +def charge_usd(model: str, u: U.Usage | None, reserve: float) -> tuple[float, str]: + """(usd, source): WF-USAGE x PRICES + 15%; missing/unknown model -> reserve, 'est'.""" + c = U.cost(model, u) if u is not None else None + return (c * OVERHEAD, "self") if c is not None else (reserve, "est") + + +def end(led: dict, sid: str, state: str, u: U.Usage | None = None) -> dict: + """Running entry ends: state, charge added to spent.""" + if state not in ENDED: + raise CloudError(f"bad end state {state}") + e = find(led, sid) + if e["state"] != "running": + raise CloudError(f"entry {sid} is {e['state']}, not running") + if state == "lost": + usd, src = led["reserve_per_task"], "est" + else: + usd, src = charge_usd(e["model"], u, led["reserve_per_task"]) + e.update(state=state, usd=usd, usd_source=src) + led["spent"] += usd + return e + + +def expire(led: dict, now: dt.datetime) -> list[dict]: + """Running > 24h (no WF-RESULT seen) -> lost, charged reserve. Returns the lost entries.""" + out = [] + for e in led["entries"]: + if e["state"] == "running" and now - dt.datetime.fromisoformat(e["sent"]) > LOST_AFTER: + out.append(end(led, e["sid"], "lost")) + return out + + +def set_balance(led: dict, bal: float, now: dt.datetime) -> dict: + """Owner reconcile (balance read from claude.ai): spent = budget - balance.""" + old = led["spent"] + led["spent"] = led["budget"] - bal + row = {"id": "-", "project": "-", "sid": "-", "model": "-", "sent": now.isoformat(), "state": "done", + "usd": led["spent"] - old, "usd_source": "owner"} + led["entries"].append(row) + return row + + +def summary(led: dict) -> str: + return (f"budget ${led['budget']:.2f} spent ${led['spent']:.2f} balance ${balance(led):.2f} " + f"running {running(led)}/{led['max_parallel']} (reserve ${led['reserve_per_task']:.2f}/task)") + + + +# --- archive ended sessions (spec §4.3 Archive): the CLI's internal archiveRemoteSession, CLI 2.1.291 --- +API = "https://api.anthropic.com" +_ARCH_SID = re.compile(r"session_[A-Za-z0-9]+") + + +def archive_request(sid: str, creds: str, version: str) -> tuple[str, dict]: + """(url, headers) for POST .../archive; creds = the CLI's .credentials.json text. Never logged.""" + if not _ARCH_SID.fullmatch(sid): + raise CloudError(f"{sid!r} is not a cloud session id") + try: + d = json.loads(creds) + tok = d["claudeAiOauth"]["accessToken"] + except (ValueError, KeyError, TypeError): + tok = None + if not tok: + raise CloudError("no claude.ai login token (run claude once)") + h = {"Authorization": f"Bearer {tok}", "Content-Type": "application/json", + "anthropic-version": "2023-06-01", "User-Agent": f"claude-code/{version}"} + if d.get("trustedDeviceToken"): + h["X-Trusted-Device-Token"] = d["trustedDeviceToken"] + return f"{API}/v1/code/sessions/{sid}/archive", h + + +def archive_problem(status: int, body: str) -> str | None: + """None = archived (200, or 409 = already); else one short why.""" + if status in (200, 409): + return None + if status == 401: + return "HTTP 401 (login expired? run claude once)" + return f"HTTP {status}: {' '.join(body.split())}"[:120] + + +def unarchived(led: dict) -> list[str]: + """sids of ended sessions not archived yet (owner reconcile rows and pending sids skipped).""" + return [e["sid"] for e in led["entries"] + if e["state"] in ENDED and _ARCH_SID.fullmatch(e["sid"]) and not e.get("archived")] + +# --- pty driver helpers (spec §3 F1-F3, F7, F8) --------------------------------------------------- +import re # noqa: E402 + +_ANSI = re.compile(r"\x1b(\[[0-?]*[ -/]*[@-~]|\][^\x07\x1b]*(\x07|\x1b\\)|[PX^_][^\x1b]*\x1b\\|[@-Z\\-_])") +_SID = re.compile(r"(session_[A-Za-z0-9]{8,})") + + +def strip_ansi(text: str) -> str: + """Terminal output -> plain text: cursor-forward (CSI n C) -> n spaces, cursor-to-column (CSI n G) + -> one space (the TUI draws word gaps that way), other escape sequences dropped, CR -> LF.""" + text = re.sub(r"\x1b\[(\d*)C", lambda m: " " * int(m.group(1) or 1), text) + text = re.sub(r"\x1b\[\d*G", " ", text) + return _ANSI.sub("", text).replace("\r\n", "\n").replace("\r", "\n") + + +def parse_sid(text: str) -> str | None: + """Session id from `claude --cloud` output: 'Resume with: claude --teleport ', else the View URL.""" + plain = strip_ansi(text) + for pat in (r"Resume with:\s*claude --teleport\s+" + _SID.pattern, r"View:\s*\S*?/code/" + _SID.pattern): + m = re.findall(pat, plain) + if m: + return m[-1] + return None + + +def last_lines(text: str, n: int = 5) -> list[str]: + return [ln.strip() for ln in strip_ansi(text).splitlines() if ln.strip()][-n:] + + +def trust(claude_json: str, path: str) -> str: + """~/.claude.json text with projects[path].hasTrustDialogAccepted = true (F3). Other keys kept.""" + try: + data = json.loads(claude_json) if claude_json.strip() else {} + projects = data.setdefault("projects", {}) + if not isinstance(projects, dict): + raise TypeError("projects is not an object") + except (ValueError, TypeError, AttributeError) as e: + raise CloudError(f"~/.claude.json unreadable: {e}") + projects.setdefault(path, {})["hasTrustDialogAccepted"] = True + return json.dumps(data, indent=2) + "\n" + + +TRUST_PROMPT = re.compile(r"trust\s*(the\s*files\s*in\s*)?this\s*folder|Do\s*you\s*trust", re.I) +RESUMED = re.compile(r"Session\s*resumed", re.I) # teleport reached the prompt (F7) + + +# --- send: snapshot + prompt (spec §4.1, §4.2) ----------------------------------------------------- + +MODEL = "claude-opus-5-5" # what cloud sessions run (F4): ledger rows price with it +CAP = 90_000_000 # bytes: packed snapshot limit (upload ≤ 100 MB, F1) +NEVER = (".wf", ".worktrees", "out") # + TASKS / archive: never in a snapshot (cloud_include re-adds) + + +def include_problem(path: str) -> str | None: + """Why a cloud_include entry is unusable (absolute, leaves the project, empty), None = fine.""" + parts = path.replace("\\", "/").strip("/").split("/") + if not path.strip() or path.startswith("/") or ".." in parts or parts in (["."], [".git"]) or parts[0] == ".git": + return f"cloud_include '{path}': must be a path inside the project" + return None + + +def size_refusal(size: int, cap: int = CAP) -> str | None: + if size > cap: + return f"snapshot {size / 1e6:.0f} MB > {cap / 1e6:.0f} MB; trim cloud_include" + return None + + +def fill(template: str, id: str, base: str, task: str, recipe: str, note: str | None) -> str: + """templates/cloud-prompt.md with its {{…}} fields filled.""" + values = {"id": id, "base": base, "task": task.strip(), "recipe": recipe.strip() or "(none: see CLAUDE.md)", + "note": f"\n{note.strip()}\n" if note and note.strip() else ""} + unknown = sorted(set(re.findall(r"\{\{(\w+)\}\}", template)) - set(values)) + if unknown: + raise CloudError(f"cloud prompt template: unknown field {{{{{unknown[0]}}}}}") + out = re.sub(r"\{\{(\w+)\}\}", lambda m: values[m.group(1)], template) # one pass: task text kept as is + return out + + +# The export renders the prompt as '❯ first line' + ' ' continuations and every assistant message as +# '● first line' + ' ' continuations (F12): only an assistant message can start with '● WF-RESULT'. +RESULT = re.compile(r"^● WF-RESULT\b") + + +def final_message(export: str) -> list[str] | None: + """Lines of the last assistant message that starts with WF-RESULT (bullet / indent and trailing blanks + removed), None = no result yet. The prompt never matches, whatever it contains.""" + lines = export.replace("\r\n", "\n").split("\n") + starts = [i for i, ln in enumerate(lines) if RESULT.match(ln)] + if not starts: + return None + out = [] + for ln in lines[starts[-1]:]: + if out and ln[:2] not in (" ", ""): + break # next message / prompt + out.append(ln[2:].rstrip()) + while out and not out[-1]: + out.pop() + return out + + +def keyless_message(export: str) -> list[str] | None: + """No WF-RESULT anywhere, but the last assistant message ends in a bare WF-PATCH-END line (the session + dropped the key words): its lines (as final_message), else None. A non-empty prompt after it (a redo + sent) -> None: still running.""" + lines = export.replace("\r\n", "\n").split("\n") + blocks: list[tuple[str, list[str]]] = [] + for ln in lines: + if ln[:2] in (" ", "") and blocks: + blocks[-1][1].append(ln[2:].rstrip()) + elif ln.strip(): + blocks.append((ln[:1], [ln[2:].rstrip()])) + found = None + for kind, body in blocks: + if kind == "❯" and any(b.strip() for b in body): + found = None # a later prompt: the message before it is answered + elif kind == "●": + text = [b for b in body if b.strip()] + if text and text[-1].strip() == "WF-PATCH-END": + found = body + if found is None: + return None + out = list(found) + while out and not out[-1]: + out.pop() + return out + + +def keyless_result(lines: list[str]) -> Result: + """keyless_message() lines -> a handback Result: usage + patch (sha256=HEX bytes=N header anywhere before + WF-PATCH-END, base64 after it) when found, so the patch can be kept.""" + text = " ".join(ln.strip() for ln in lines[:-1]) + r = Result("handback", "(no WF-RESULT key)") + if m := _USAGE.search(text): + r.usage = U.Usage(turns=1, inp=int(m[1]), cw5=int(m[2]), cr=int(m[3]), out=int(m[4])) + r.model = m[5] + blob = "".join(text.split()) + if m := _HEADER.search(blob): + r.sha, r.nbytes, r.b64, r.has_patch = m[1].lower(), int(m[2]), blob[m.end():], True + return r + + +# --- pull: result + patch (spec §4.3) --------------------------------------------------------------- + +KEYS = ("WF-RESULT", "WF-REPORT", "WF-USAGE", "WF-PATCH-BEGIN", "WF-PATCH-END") +STATES = ("done", "awaiting", "handback") +_USAGE = re.compile(r"in=(\d+)\s+cw=(\d+)\s+cr=(\d+)\s+out=(\d+)(?:\s+model=(\S+))?") +# whitespace removed: the export may wrap the header anywhere; gzip base64 starts 'H4sI', never a digit +_HEADER = re.compile(r"sha256=([0-9a-fA-F]{64})bytes=(\d+)") + + +@dataclasses.dataclass +class Result: + state: str + report: str + usage: U.Usage | None = None + model: str | None = None + sha: str | None = None + nbytes: int | None = None + b64: str = "" + has_patch: bool = False + + +def parse_result(lines: list[str]) -> Result: + """final_message() lines -> Result. A line starting with a key opens its section, other lines continue + the open one (the export wraps long lines). CloudError: bad state / missing patch markers.""" + sec: dict[str, list[str]] = {} + cur = None + for ln in lines: + s = ln.strip() + key = next((k for k in KEYS if s == k or s.startswith(k + " ")), None) + if key: + cur = key + sec.setdefault(key, []).append(s[len(key):].strip()) + elif cur: + sec[cur].append(s) + state = " ".join(sec.get("WF-RESULT", [""])).strip().lower() + if state not in STATES: + raise CloudError(f"bad WF-RESULT '{state}' (want {'/'.join(STATES)})") + r = Result(state, " ".join(" ".join(sec.get("WF-REPORT", [])).split()) or "(no report)") + if m := _USAGE.search(" ".join(sec.get("WF-USAGE", []))): + r.usage = U.Usage(turns=1, inp=int(m[1]), cw5=int(m[2]), cr=int(m[3]), out=int(m[4])) + r.model = m[5] + if "WF-PATCH-BEGIN" in sec: + if "WF-PATCH-END" not in sec: + raise CloudError("WF-PATCH-BEGIN without WF-PATCH-END (truncated message)") + blob = "".join("".join(sec["WF-PATCH-BEGIN"]).split()) + m = _HEADER.match(blob) + if not m: + raise CloudError("WF-PATCH-BEGIN needs sha256=HEX bytes=N") + r.sha, r.nbytes, r.b64, r.has_patch = m[1].lower(), int(m[2]), blob[m.end():], True + return r + + +def decode_patch(r: Result) -> bytes: + """base64 -> check bytes + sha256 -> gunzip. CloudError on any mismatch.""" + if not r.has_patch: + raise CloudError("no WF-PATCH-BEGIN/END in the result") + try: + gz = base64.b64decode(r.b64, validate=True) + except (binascii.Error, ValueError) as e: + raise CloudError(f"patch base64 broken: {e}") + if len(gz) != r.nbytes: + raise CloudError(f"patch bytes {len(gz)} != {r.nbytes} announced") + if hashlib.sha256(gz).hexdigest() != r.sha: + raise CloudError("patch sha256 mismatch") + try: + return gzip.decompress(gz) + except (OSError, EOFError) as e: + raise CloudError(f"patch gzip broken: {e}") + + +_DIFF = re.compile(r'^diff --git "?a/(.+?)"? "?b/(.+?)"?$') +_MOVE = re.compile(r"^(?:rename|copy) (?:from|to) (.+)$") + + +def patch_files(patch: str) -> list[str]: + """Repo-relative paths a format-patch touches (both sides of renames/copies), in order, unique.""" + out = [] + for ln in patch.split("\n"): + if m := _DIFF.match(ln): + out += [m[1], m[2]] + elif m := _MOVE.match(ln): + out.append(m[1].strip('"')) + return list(dict.fromkeys(out)) + + +def patch_problem(files: list[str], project: str, never: list[str]) -> str | None: + """Why a patch may not be applied, None = fine. files: repo-relative; project: the project folder + relative to the repo ('' = repo root); never: project-relative paths the cloud may not touch + (TASKS, archive, .wf, .worktrees, out, cloud_include).""" + pre = project.strip("/") + for f in files: + parts = f.split("/") + if not f or f.startswith("/") or ".." in parts or parts[0] == ".git": + return f"patch touches {f}: outside the tree" + if pre and not (f + "/").startswith(pre + "/"): + return f"patch touches {f}: outside the project folder {pre}" + rel = f[len(pre) + 1:] if pre else f + for n in never: + n = n.strip("/") + if rel == n or rel.startswith(n + "/"): + return f"patch touches {f}: {n} is never sent (wf bookkeeping / out / cloud_include)" + return None + + +# ------------------------------------------------------------------ fit (spec §4.5) + +# (reason, regex) over title + body; `Cloud: yes` skips them (Model opus still required). wf res before the generic wf-step rule. +UNFIT = ( + ("names `wf res` (local resource ledger)", re.compile(r"\bwf res\b")), + ("names `wf` commands as steps", re.compile(r"\bwf (?!res\b)[a-z][a-z-]*")), + ("needs a GUI", re.compile(r"\bGUI\b|\bscreenshot", re.I)), + ("needs a LAN host", re.compile(r"\bLAN\b|\b192\.168\.\d+\.\d+|\b10\.\d+\.\d+\.\d+|\b[\w-]+\.local\b|\bssh [\w@.-]+")), + ("needs a live server", re.compile(r"live server|running server|\blocalhost\b|127\.0\.0\.1", re.I)), +) + + +def fit(item: TK.Item, cfg) -> str | None: + """None when the task fits the cloud lane, else the reason it does not (spec §4.5).""" + if not cfg.cloud: + return "project not opted in (cloud = true)" + if item.cloud == "no": + return "Cloud: no" + if not item.runner_ready: + return "not runner-ready (needs Done, not owner-bound)" + if item.sessions in ("owner", "solo"): + return f"Sessions: {item.sessions}" + if item.effort not in TK.EFFORTS or TK.EFFORTS.index(item.effort) > TK.EFFORTS.index(cfg.slice_above): + return "slice job (effort above slice_above)" + if item.model != "opus": + return f"Model {item.model} stays local" + if item.cloud == "yes": + return None + text = "\n".join([item.text, *item.body]) + for reason, rx in UNFIT: + if rx.search(text): + return reason + return None + + +def pick_key(item: TK.Item) -> tuple: + """wf orch pick cloud order (spec §4.5): Cloud: yes first, then prio, then larger effort (all fits are opus).""" + eff = TK.EFFORTS.index(item.effort) if item.effort in TK.EFFORTS else -1 + return (item.cloud != "yes", item.prio if item.prio is not None else 9, -eff) diff --git a/wflib/config.py b/wflib/config.py new file mode 100644 index 0000000..8f75c89 --- /dev/null +++ b/wflib/config.py @@ -0,0 +1,225 @@ +"""Find the project (nearest ancestor with workflow.toml) and load its config.""" +from __future__ import annotations + +import tomllib +from dataclasses import dataclass, field +from pathlib import Path + +from . import lanes as lanes_mod + +NAME = "workflow.toml" +FORMAT = 1 +KEYS = {"format", "tasks", "archive", "docs", "verify", "done", "ledgers", "anchors", "ctx_hint", "worktree_setup", "quick_gate", + "slice_above", "lanes", "areas", "code_root", "area_stale_commits", "area_ignore", + "cloud", "cloud_include", "cloud_note"} +AREA_IGNORE = ("tests", "test", "docs", "doc") +ANCHOR_KEYS = {"index", "index_section", "specs"} + + +class ConfigError(Exception): + pass + + +@dataclass +class Config: + root: Path + format: int + tasks: Path + archive: Path + docs: list[Path] = field(default_factory=list) + verify: list[str] = field(default_factory=list) + done: list[str] = field(default_factory=list) + ledgers: Path | None = None + anchors_index: Path | None = None + anchors_section: str | None = None + anchors_specs: Path | None = None + ctx_hint: int = 100_000 # tokens; wf done/next suggest /clear above it (0 = off) + worktree_setup: list[str] = field(default_factory=list) # `wf setup`: cwd = worktree, env WF_MAIN + quick_gate: list[str] = field(default_factory=list) # `wf gate`: worker runs it before `wf done` + slice_above: str = "1h" + lanes: tuple = lanes_mod.DEFAULT_LANES + areas: Path | None = None # area notes file; None = CLAUDE.md + code_root: Path | None = None # repo area anchors / Checked / staleness refer to; None = root + area_stale_commits: int = 20 + area_ignore: list[str] = field(default_factory=lambda: list(AREA_IGNORE)) # never an uncovered area + local: Path | None = None # lane worktree's copy of root whose LOCAL_KEYS were used + cloud: bool = False # cloud lane opt-in (wf cloud send) + cloud_include: list[str] = field(default_factory=list) # git-ignored paths force-added to the snapshot + cloud_note: str | None = None # appended to the cloud prompt + + @property + def areas_file(self) -> Path: + return self.areas or (self.local or self.root) / "CLAUDE.md" + + @property + def code(self) -> Path: + return self.code_root or self.root + + def rel(self, path: Path) -> str: + try: + return str(path.relative_to(self.root)) + except ValueError: + return str(path) + + @property + def areas_rel(self) -> str: + """areas_file relative to its project folder (the lane worktree's copy or root).""" + try: + return str(self.areas_file.relative_to(self.local or self.root)) + except ValueError: + return str(self.areas_file) + + +def git_top(folder: Path) -> Path | None: + """The work tree top holding folder (.git dir or file), or None.""" + for f in (folder, *folder.parents): + if (f / ".git").exists(): + return f + return None + + +def git_branch(gitdir: Path) -> str | None: + try: + head = (gitdir / "HEAD").read_text().strip() + except OSError: + return None + return head.removeprefix("ref: refs/heads/") if head.startswith("ref: refs/heads/") else None + + +def linked_worktree(folder: Path) -> tuple[Path, Path, Path] | None: + """(worktree top, its gitdir, main tree top) when folder is in a linked git worktree.""" + top = git_top(folder) + if not top or not (top / ".git").is_file(): + return None + try: + line = (top / ".git").read_text().strip() + gitdir = (top / line.removeprefix("gitdir:").strip()).resolve() + common = (gitdir / (gitdir / "commondir").read_text().strip()).resolve() + except OSError: + return None + if common.name != ".git" or not line.startswith("gitdir:"): + return None + return top, gitdir, common.parent + + +def find_root(start: Path) -> Path: + """Folder with workflow.toml at or above start; in a linked git worktree the main tree's + same folder, when it is a project too (one TASKS.md for all worktrees).""" + start = start.resolve() + for folder in (start, *start.parents): + if (folder / NAME).is_file(): + wt = linked_worktree(folder) + if wt: + main = wt[2] / folder.relative_to(wt[0]) + if (main / NAME).is_file(): + return main + return folder + raise ConfigError(f"no {NAME} in {start} or above: not a wf project (wf init makes one)") + + +def find_local(start: Path) -> Path | None: + """In a linked worktree whose project root is the main tree's (find_root): the worktree's own copy + of that folder when it has a workflow.toml, else None.""" + start = start.resolve() + wt = linked_worktree(start) + if not wt: + return None + root = find_root(start) + try: + local = wt[0] / root.relative_to(wt[2]) + except ValueError: + return None + return local if local != root and (local / NAME).is_file() else None + + +def load_at(start: Path) -> "Config": + """Config for a command run in start: the main tree's books, branch-testable keys (LOCAL_KEYS: + worktree_setup, quick_gate, areas file, area settings) from a lane worktree's own workflow.toml.""" + return load(find_root(start), find_local(start)) + + +def _strings(data: dict, key: str) -> list[str]: + value = data.get(key, []) + if not isinstance(value, list) or not all(isinstance(v, str) for v in value): + raise ConfigError(f"{NAME}: '{key}' must be a list of strings") + return value + + +def _string(data: dict, key: str, shown: str | None = None) -> str | None: + value = data.get(key) + if value is not None and not isinstance(value, str): + raise ConfigError(f"{NAME}: '{shown or key}' must be a string") + return value + + +LOCAL_KEYS = ("worktree_setup", "quick_gate", "areas", "area_stale_commits", "area_ignore") + + +def _read(path: Path) -> dict: + try: + return tomllib.loads(path.read_text(encoding="utf-8")) + except (tomllib.TOMLDecodeError, OSError) as e: + raise ConfigError(f"{path.name}: {e}") from None + + +def load(root: Path, local: Path | None = None) -> Config: + """root's workflow.toml; with local (a lane worktree's copy of root), LOCAL_KEYS come from local's + workflow.toml and the areas file resolves there (a branch can test and --mark them before merge).""" + data = _read(root / NAME) + if local: + mine = _read(local / NAME) + for key in LOCAL_KEYS: + data.pop(key, None) + if key in mine: + data[key] = mine[key] + for key in data: + if key not in KEYS: + raise ConfigError(f"{NAME}: unknown key '{key}'") + anchors = data.get("anchors", {}) + if not isinstance(anchors, dict): + raise ConfigError(f"{NAME}: 'anchors' must be a table") + for key in anchors: + if key not in ANCHOR_KEYS: + raise ConfigError(f"{NAME}: unknown key 'anchors.{key}'") + for key in ("format", "tasks", "archive"): + if key not in data: + raise ConfigError(f"{NAME}: missing '{key}'") + if not isinstance(data["format"], int) or isinstance(data["format"], bool): + raise ConfigError(f"{NAME}: 'format' must be a number") + ctx_hint = data.get("ctx_hint", 100_000) + if not isinstance(ctx_hint, int) or isinstance(ctx_hint, bool) or ctx_hint < 0: + raise ConfigError(f"{NAME}: 'ctx_hint' must be a number of tokens (0 = off)") + index = _string(anchors, "index", "anchors.index") + specs = _string(anchors, "specs", "anchors.specs") + ledgers = _string(data, "ledgers") + stale = data.get("area_stale_commits", 20) + if not isinstance(stale, int) or isinstance(stale, bool) or stale < 1: + raise ConfigError(f"{NAME}: 'area_stale_commits' must be a number ≥ 1") + slice_above = _string(data, "slice_above") or "1h" + try: + lane_list = lanes_mod.parse_lanes(data.get("lanes"), slice_above) + except ValueError as e: + raise ConfigError(f"{NAME}: {e}") from None + areas = _string(data, "areas") + code_root = _string(data, "code_root") + if not isinstance(data.get("cloud", False), bool): + raise ConfigError(f"{NAME}: 'cloud' must be true or false") + return Config( + root=root, format=data["format"], + tasks=root / _string(data, "tasks"), archive=root / _string(data, "archive"), + docs=[root / d for d in _strings(data, "docs")], + verify=_strings(data, "verify"), done=_strings(data, "done"), + ledgers=root / ledgers if ledgers else None, + anchors_index=root / index if index else None, + anchors_section=_string(anchors, "index_section", "anchors.index_section"), + anchors_specs=root / specs if specs else None, + ctx_hint=ctx_hint, + worktree_setup=_strings(data, "worktree_setup"), quick_gate=_strings(data, "quick_gate"), + slice_above=slice_above, lanes=lane_list, areas=(local or root) / areas if areas else None, + code_root=(root / code_root).resolve() if code_root else None, + area_stale_commits=stale, + area_ignore=_strings(data, "area_ignore") if "area_ignore" in data else list(AREA_IGNORE), + local=local, + cloud=data.get("cloud", False), cloud_include=_strings(data, "cloud_include"), + cloud_note=_string(data, "cloud_note"), + ) diff --git a/wflib/lanes.py b/wflib/lanes.py new file mode 100644 index 0000000..d9b6fe6 --- /dev/null +++ b/wflib/lanes.py @@ -0,0 +1,297 @@ +"""Lanes (by task size): config records, membership, pick order, output blocks. Sessions are dicts +{"lane", "model", "socket", "alive", …} keyed by lane name or "all" (registry I/O lives in wf.py).""" +import re +from dataclasses import dataclass + +from . import tasks + +LANE_KEYS = {"efforts", "order", "fallback", "slices"} +ORDERS = ("unblock", "priority") +LANE_NAME_RE = re.compile(r"^[a-z][a-z0-9-]*$") + + +@dataclass(frozen=True) +class Lane: + name: str + efforts: tuple[str, ...] + order: str = "priority" + fallback: str | None = None + slices: bool = False + + +DEFAULT_LANES = (Lane("fast", ("<1h",), "unblock", None, True), Lane("slow", ("1h",), "priority", "fast", False)) + + +def parse_lanes(table: dict | None, slice_above: str) -> tuple[Lane, ...]: + """[lanes.*] of workflow.toml → lanes in file order; None → DEFAULT_LANES. Raises ValueError.""" + if slice_above not in tasks.EFFORTS: + raise ValueError(f"slice_above '{slice_above}' (want {', '.join(tasks.EFFORTS)})") + allowed = tasks.EFFORTS[:tasks.EFFORTS.index(slice_above) + 1] + if table is None: + out = DEFAULT_LANES + else: + if not isinstance(table, dict): + raise ValueError("'lanes' must be a table of tables") + out = [] + for name, t in table.items(): + if not LANE_NAME_RE.match(name): + raise ValueError(f"bad lane name '{name}' (want a-z 0-9 -)") + if name == "all": + raise ValueError("lane name 'all' is reserved") + if not isinstance(t, dict): + raise ValueError(f"'lanes.{name}' must be a table") + for k in t: + if k not in LANE_KEYS: + raise ValueError(f"unknown key 'lanes.{name}.{k}'") + efforts = t.get("efforts") + if not isinstance(efforts, list) or not efforts or not all(isinstance(e, str) for e in efforts): + raise ValueError(f"lane {name}: 'efforts' must be a non-empty list of strings") + order = t.get("order", "priority") + if order not in ORDERS: + raise ValueError(f"lane {name}: order '{order}' (want {', '.join(ORDERS)})") + fallback, slices = t.get("fallback"), t.get("slices", False) + if fallback is not None and not isinstance(fallback, str) or not isinstance(slices, bool): + raise ValueError(f"lane {name}: 'fallback' must be a string, 'slices' true/false") + out.append(Lane(name, tuple(efforts), order, fallback, slices)) + names = [l.name for l in out] + owner: dict[str, str] = {} + for l in out: + for e in l.efforts: + if e not in tasks.EFFORTS: + raise ValueError(f"lane {l.name}: effort '{e}' (want {', '.join(tasks.EFFORTS)})") + if e not in allowed: + raise ValueError(f"lane {l.name}: effort '{e}' above slice_above '{slice_above}'") + if e in owner: + raise ValueError(f"effort '{e}' in lanes {owner[e]} and {l.name}") + owner[e] = l.name + if l.fallback is not None and (l.fallback not in names or l.fallback == l.name): + raise ValueError(f"lane {l.name}: fallback '{l.fallback}' is not another lane") + for e in allowed: + if e not in owner: + raise ValueError(f"effort '{e}' (≤ slice_above) in no lane") + by = {l.name: l for l in out} + for l in out: + if l.fallback and by[l.fallback].fallback == l.name: + a, b = sorted((l.name, l.fallback)) + raise ValueError(f"fallback cycle: {a} ↔ {b}") + flagged = [l for l in out if l.slices] + if len(flagged) > 1: + raise ValueError("slices = true on more than one lane") + if not flagged: + out = [Lane(l.name, l.efforts, l.order, l.fallback, owner["<1h"] == l.name) for l in out] + return tuple(out) + + + +def _pending(doc: tasks.Doc) -> list[tasks.Item]: + return [i for s in doc.sections if s.key == "pending" for i in s.items if not i.error] + + +def counts(doc: tasks.Doc, archived: set[str], lanes, slice_above: str, + held: dict[str, str] | None = None, runner: bool = False) -> dict[str, tuple[int, int, int]]: + """lane → (pickable, waiting, pickable slice jobs) over Pending, in lane order. held: claimed ids count waiting. + runner: pickable = what `list --runner` shows (has a usable Done); the rest is not_ready().""" + out = {l.name: [0, 0, 0] for l in lanes} + for item in _pending(doc): + c = out[lane_of(item, lanes, slice_above)] + if _why(doc, item, archived, slice_above, held, 0, True) is None: + if runner and not item.runner_ready: + continue + c[0] += 1 + c[2] += is_slice_job(item, slice_above) + else: + c[1] += 1 + return {name: tuple(c) for name, c in out.items()} + + +def not_ready(doc: tasks.Doc, archived: set[str], lanes, slice_above: str, + held: dict[str, str] | None = None) -> dict[str, int]: + """lane → pickable items without a Done that a prep worker can fill (owner-bound ones not counted).""" + out = {l.name: 0 for l in lanes} + for item in _pending(doc): + if _why(doc, item, archived, slice_above, held, 0, True) is None \ + and not (item.done_text and item.done_text.strip()) and item.sessions != "owner": + out[lane_of(item, lanes, slice_above)] += 1 + return out + + +def prep_targets(doc: tasks.Doc, archived: set[str], lanes, slice_above: str, + only: set[str] | None = None) -> list[tasks.Item]: + """Pending items a prep worker can make runner-ready: no Done, not owner-bound, no status, no open + slices; only: lane names. Order: pickable first, then priority, then file order.""" + out = [] + for n, item in enumerate(_pending(doc)): + if item.done_text and item.done_text.strip() or item.sessions == "owner" or item.status \ + or tasks.open_slices(doc, item.id): + continue + if only is not None and lane_of(item, lanes, slice_above) not in only: + continue + waits = _why(doc, item, archived, slice_above, None, 0, True) is not None + out.append(((waits, 9 if item.prio is None else item.prio, n), item)) + return [item for _, item in sorted(out, key=lambda x: x[0])] + + +def cross_waits(doc: tasks.Doc, archived: set[str], lane: str, lanes, + slice_above: str) -> list[tuple[tasks.Item, tasks.Item]]: + """(my waiting item, pickable item of another lane it waits on, transitively).""" + items = _pending(doc) + wait = tasks.waiters(doc) + name = {i.id: lane_of(i, lanes, slice_above) for i in items} + out = [] + for mine in items: + if name[mine.id] != lane or tasks.why_not(doc, mine, archived) is None: + continue + for other in items: + if name[other.id] != lane and mine.id in wait.get(other.id, ()) \ + and tasks.why_not(doc, other, archived) is None: + out.append((mine, other)) + return out + + +def lanes_block(counts: dict[str, tuple[int, int, int]], sessions: dict[str, dict], me: str | None, + nready: dict[str, int] | None = None) -> list[str]: + """One line per lane (always all: lanes are few). nready: lane → not runner-ready count.""" + out = [] + for name, (p, w, sl) in counts.items(): + line = (f"{name}{' (you)' if name == me else ''}: {p} pickable" + + (f" ({sl} slice{'s' if sl > 1 else ''})" if sl else "") + f" · {w} waiting") + if (nready or {}).get(name): + line += f" · {nready[name]} not runner-ready (no Done)" + if name != me: + s = sessions.get(name) + if s and s.get("alive"): + line += f" · session uds:{s['socket']} (alive, {s.get('model') or '?'})" + else: + line += " · no session" + (f" → orchestrator or owner: start one (wf next --lane {name})" if p else "") + out.append(line) + return out + + +def waiting_block(waits: list[tuple[tasks.Item, tasks.Item]], sessions: dict[str, dict], lanes, + slice_above: str) -> list[str]: + out = [] + for mine, other in waits: + name = lane_of(other, lanes, slice_above) + s = sessions.get(name) + head = f"- {mine.id} waits on {other.id} ({name} lane) → " + if s and s.get("alive"): + out.append(head + f'message uds:{s["socket"]}: "{other.id} blocks my {mine.id}, please take it"') + else: + out.append(head + f"no {name} session: tell the owner") + return out + + +def pickable_ids(doc: tasks.Doc, archived: set[str]) -> set[str]: + return {i.id for i in _pending(doc) if tasks.why_not(doc, i, archived) is None} + + +def notify_block(freed: list[tasks.Item], sessions: dict[str, dict], lanes, slice_above: str) -> list[str]: + """One line per lane with newly pickable tasks, in lane order.""" + out = [] + for lane in lanes: + ids = ", ".join(i.id for i in freed if lane_of(i, lanes, slice_above) == lane.name) + if not ids: + continue + s = sessions.get(lane.name) + if s and s.get("alive"): + out.append(f"notify {lane.name} uds:{s['socket']}: now pickable {ids}") + else: + out.append(f"{lane.name} work now pickable: {ids} (no session: tell the owner)") + return out + + +def solo_done_block(solo: list[str], sessions: dict[str, dict], me_pid: str | None) -> list[str]: + """After solo task(s) are done: one line per other live session, so it can pick again.""" + if not solo: + return [] + ids = ", ".join(solo) + return [f"notify {name} uds:{s['socket']}: solo {ids} done, run wf next" + for name, s in sessions.items() if s.get("alive") and str(s.get("pid")) != me_pid] + + +MODEL_RANK = {m: n for n, m in enumerate(tasks.MODELS)} + + +def over(effort: str | None, slice_above: str) -> bool: + return effort in tasks.EFFORTS and tasks.EFFORTS.index(effort) > tasks.EFFORTS.index(slice_above) + + +def is_slice_job(item: tasks.Item, slice_above: str) -> bool: + """Too big to implement and not sliced yet: the slices lane splits it.""" + return over(item.effort, slice_above) and not item.slices + + +def sliced_out(item: tasks.Item, slice_above: str) -> bool: + """Too big and already sliced once: never picked again (wf done it or add slices by hand).""" + return over(item.effort, slice_above) and bool(item.slices) + + +def lane_of(item: tasks.Item, lanes: tuple[Lane, ...], slice_above: str) -> str: + if over(item.effort, slice_above): + return next(l.name for l in lanes if l.slices) + effort = item.effort if item.effort in tasks.EFFORTS else "<1h" + return next(l.name for l in lanes if effort in l.efforts) + + +def _why(doc, item, archived, slice_above, held, others, owner) -> str | None: + why = tasks.why_not(doc, item, archived, held, others, owner) + if why is None and not item.error and sliced_out(item, slice_above): + return f"sliced: all slices done → wf done {item.id} or add slices" + return why + + +def rank(doc: tasks.Doc, items: list[tasks.Item], order: str, lane: dict[str, str]) -> dict[str, tuple]: + """Sort key per item. lane: item id → lane name (for 'waited on by another lane').""" + wait = tasks.waiters(doc) + by_id = {i.id: i for s in doc.sections if s.key in ("pending", "human") for i in s.items if not i.error} + out = {} + for n, item in enumerate(items): + ws = [by_id[w] for w in wait.get(item.id, ())] + prio = min([p for p in [item.prio, *(w.prio for w in ws)] if p is not None], default=9) + cross = any(lane.get(w.id) != lane.get(item.id) for w in ws) + out[item.id] = ((not ws, -len(ws), prio, n) if order == "unblock" + else (prio, not cross, -len(ws), n)) + return out + + +def _candidates(doc, lanes, slice_above, lane: str | None, model: str | None) -> list[tasks.Item]: + items = doc.section("pending").items + return [i for i in items if i.error or ( + (lane is None or lane_of(i, lanes, slice_above) == lane) + and (model is None or MODEL_RANK[i.model] <= MODEL_RANK[model]))] + + +def _ordered(doc, lanes, slice_above, lane: str | None, model: str | None) -> list[tasks.Item]: + items = _candidates(doc, lanes, slice_above, lane, model) + order = next((l.order for l in lanes if l.name == lane), "priority") + names = {i.id: lane_of(i, lanes, slice_above) for i in doc.section("pending").items if not i.error} + keys = rank(doc, items, order, names) + return sorted(items, key=lambda i: keys.get(i.id, (9, True, 0, items.index(i)))) + + +def pick(doc, archived, lanes, slice_above, lane=None, model=None, held=None, others=0, owner=False): + """Best pickable Pending item of the lane (None = all lanes, priority order), Model ≤ model; + the lane empty → its fallback lane (one hop). Plus unpickable items ranked before it, with the reason.""" + skipped = [] + hops = [lane] + fb = next((l.fallback for l in lanes if l.name == lane), None) + if fb: + hops.append(fb) + for name in hops: + for item in _ordered(doc, lanes, slice_above, name, model): + why = _why(doc, item, archived, slice_above, held, others, owner) + if why is None: + return item, skipped + if (item, why) not in skipped: + skipped.append((item, why)) + return None, skipped + + +def ranked(doc, archived, lanes, slice_above, lane: str | None) -> list[tasks.Item]: + """Every pickable item of the lane in pick order, then its fallback lane's (runner lists).""" + hops = [lane] + [l.fallback for l in lanes if l.name == lane and l.fallback] + out = [] + for name in hops: + out += [i for i in _ordered(doc, lanes, slice_above, name, None) + if not i.error and _why(doc, i, archived, slice_above, None, 0, False) is None and i not in out] + return out diff --git a/wflib/ledgers.py b/wflib/ledgers.py new file mode 100644 index 0000000..e22caa5 --- /dev/null +++ b/wflib/ledgers.py @@ -0,0 +1,35 @@ +"""In-flight plan ledgers (superpowers executing-plans / subagent-driven-development): +//progress.md, first line `# SDD ledger — plan: `.""" +from __future__ import annotations + +import re +from pathlib import Path + +TASK_RE = re.compile(r"^### Task (\S+?):", re.M) +DONE_RE = re.compile(r"^Task (\S+?): complete", re.M) + + +def summary(plan_text: str, ledger_text: str) -> str: + """'done 1, 2; resume at Task 3 (…)' / 'done 1, 2, 3; all tasks done; Final review: …'.""" + planned = TASK_RE.findall(plan_text) + done = DONE_RE.findall(ledger_text) + left = [t for t in planned if t not in done] + finals = [l for l in ledger_text.split("\n") if l.startswith("Final review")] + head = f"done {', '.join(done) if done else 'none'}" + if left: + return f"{head}; resume at Task {left[0]} (task-start PLAN {left[0]})" + return f"{head}; all tasks done; {finals[-1] if finals else 'final review not dispatched yet'}" + + +def in_flight(root: Path, ledgers: Path) -> list[str]: + out = [] + for ledger in sorted(ledgers.glob("*/progress.md")): + text = ledger.read_text(encoding="utf-8", errors="replace") + first = text.split("\n", 1)[0] + plan = first.split("plan:", 1)[1].strip() if "plan:" in first else "" + plan_path = root / plan + plan_text = plan_path.read_text(encoding="utf-8", errors="replace") if plan and plan_path.is_file() else "" + rulings = sum(1 for l in text.split("\n") if "Ruling:" in l) + out.append(f"- {plan or ledger.parent.name}: {summary(plan_text, text)} " + f"({rulings} rulings; ledger {ledger.relative_to(root)})") + return out diff --git a/wflib/migrate.py b/wflib/migrate.py new file mode 100644 index 0000000..04f606c --- /dev/null +++ b/wflib/migrate.py @@ -0,0 +1,237 @@ +"""One-time conversion of the numbered TASKS.md format to the id format. + +Old: `N. **[P1] Title** (Effort: 1h) (in progress: x) — goal.` + indented body, + `After: `, `Reference: …`, Awaiting bullets without ids. +New: docs/design.md, section 3. Input already in the new format comes back unchanged. +""" +from __future__ import annotations + +import re +from dataclasses import dataclass, field + +from . import tasks + +NUMBERED_RE = re.compile(r"^\d+\. ") +OLD_HEADER_RE = re.compile(r"^\d+\. (?:N\. )?\*\*\[P(\d)\] (.+?)\*\*(.*)$") +SLICE_RE = re.compile(r"^(.+?) (\d+)/\d+(?::|$)") +MD_LINK_RE = re.compile(r"\[[^\]]*\]\(([^)\s]+)\)") +AFTER_RE = re.compile(r"^(\s*)(?:-\s+)?After:\s*(.*)$") +REFERENCE_RE = re.compile(r"^\s*(?:-\s+)?(?:Reference|Ref):\s*(.*)$") +BULLET_RE = re.compile(r"^- (?:\*\*(.+?)\*\*:?\s*)?(.*)$") + + +@dataclass +class Report: + ids: dict[str, str] = field(default_factory=dict) # old title -> new id + notes: list[str] = field(default_factory=list) + + +@dataclass +class _Old: + title: str + item: tasks.Item + after: list[str] + refs: str | None + + +def _groups(rest: str) -> tuple[list[str], str]: + """Leading balanced '(…)' groups of `rest`, and what follows them.""" + out, i = [], 0 + while True: + j = i + while j < len(rest) and rest[j] == " ": + j += 1 + if j >= len(rest) or rest[j] != "(": + return out, rest[i:] + depth, k = 0, j + while k < len(rest): + depth += rest[k] == "(" + depth -= rest[k] == ")" + k += 1 + if depth == 0: + break + if depth: + return out, rest[i:] + out.append(rest[j + 1:k - 1]) + i = k + + +def _id_title(title: str) -> str: + return re.sub(r"\s*\([^)]*\)", "", title).strip() or title + + +def _body(lines: list[str]) -> list[str]: + """Re-indented body. A line one column deeper than the top level is a + child when the line above it opens a list (ends with ':'), else a typo.""" + out: list[str] = [] + child = False + for line in tasks.indent_body(lines): + if re.match(r"^ \S", line): + line = (2 * tasks.INDENT if child else tasks.INDENT) + line[3:] + elif re.match(r"^ \S", line): + child = line.rstrip().endswith(":") + out.append(line) + return out + + +def _is_path(part: str) -> bool: + first = part.split(" ", 1)[0] + return "/" in first or bool(re.search(r"\.\w+(#|$)", first)) + + +def _refs(text: str) -> tuple[str | None, str | None]: + """(Ref: value, leftover words that name no file).""" + parts: list[str] = [] + for part in tasks.split_refs(MD_LINK_RE.sub(r"\1", text).replace(" · ", ", ").strip().rstrip(".")): + if _is_path(part): + parts.append(part) + elif parts: + last = parts[-1] + parts[-1] = f"{last[:-1]}; {part})" if last.endswith(")") else f"{last} ({part})" + else: + return None, text.strip() + return (", ".join(parts) or None), None + + +def _sentence(title: str, goal: str) -> str: + title = title.strip().rstrip(".").replace(". ", ", ") + return f"{title}. {goal.strip()}" if goal.strip() else f"{title}." + + +def _split_body(lines: list[str]) -> tuple[list[str], list[str], str | None]: + """(body re-indented, After: titles, Ref: value).""" + body, after, ref, notes = [], [], None, [] + for line in lines: + if (m := AFTER_RE.match(line)) and not tasks.LINK_RE.search(line): + after.append(m.group(2).strip()) + elif (m := REFERENCE_RE.match(line)): + ref, words = _refs(m.group(1)) + if words: + notes.append(f"{tasks.INDENT}- Reference: {words}") + else: + body.append(line) + return _body(body) + notes, after, ref + + +def _task(lines: list[str], taken: set[str]) -> _Old | None: + m = OLD_HEADER_RE.match(lines[0]) + if not m: + return None + title = m.group(2).strip() + groups, rest = _groups(m.group(3)) + effort, interactive, status, extra = None, False, None, [] + for g in groups: + if g.startswith("Effort:"): + parts = [p.strip() for p in g[len("Effort:"):].split(",")] + effort, interactive = parts[0], "interactive" in parts[1:] + elif g.startswith(("in progress:", "blocked:")): + status = g.replace("(", "[").replace(")", "]") + else: + extra.append(f"({g})") + goal = re.sub(r"^\s*[—–-]\s*", "", rest).strip() + goal = " ".join(extra + ([goal] if goal else [])) + slice_ = SLICE_RE.match(title) + if slice_: + base = tasks.make_id(_id_title(slice_.group(1)), set()) + id = f"{base}-{int(slice_.group(2))}" + n = 2 + while id in taken: + id, n = f"{base}-{int(slice_.group(2))}-{n}", n + 1 + else: + id = tasks.make_id(_id_title(title), taken) + body, after, ref = _split_body(lines[1:]) + item = tasks.Item(id=id, prio=int(m.group(1)), effort=effort, interactive=interactive, + status=status, text=_sentence(title, goal), body=body) + return _Old(title, item, after, ref) + + +def _awaiting(lines: list[str], taken: set[str]) -> _Old: + m = BULLET_RE.match(lines[0]) + bold, rest = m.group(1), m.group(2).strip() + if bold: + title, text = bold.strip(), _sentence(bold, rest) + else: + title = rest.split(". ", 1)[0].rstrip(".") + text = rest + id = tasks.make_id(_id_title(title), taken, "a-") + body, _, _ = _split_body(lines[1:]) + return _Old(title, tasks.Item(id=id, text=text, body=body), [], None) + + +def _is_new_item(line: str) -> bool: + m = re.match(r"^- \*\*([^*\s]+)\*\*", line) + return bool(m and tasks.ID_RE.match(m.group(1))) + + +def migrate(text: str, taken: set[str] | None = None) -> tuple[str, Report]: + taken = set(taken or ()) + report = Report() + newline = "\r\n" if "\r\n" in text else "\n" + lines = text.replace("\r\n", "\n").split("\n") + while lines and not lines[-1].strip(): + lines.pop() + taken |= {m.group(1) for l in lines if (m := re.match(r"^- \*\*([^*\s]+)\*\*", l)) and _is_new_item(l)} + + out: list[object] = [] # str lines and _Old items, in order + key = None + i = 0 + while i < len(lines): + line = lines[i] + if line.startswith("## "): + heading = line[3:].strip() + key = next((k for k, h in tasks.SECTIONS.items() if heading == h or heading.startswith(h + " (")), None) + if key and heading != tasks.SECTIONS[key]: + out += [f"## {tasks.SECTIONS[key]}", "", heading[len(tasks.SECTIONS[key]):].strip()] + else: + out.append(line) + i += 1 + continue + old_task = key in tasks.TASK_KEYS and NUMBERED_RE.match(line) + old_bullet = key == "awaiting" and line.startswith("- ") and not _is_new_item(line) + if not (old_task or old_bullet): + out.append(line) + i += 1 + continue + j = i + 1 + while j < len(lines) and (not lines[j].strip() or lines[j][0].isspace()): + j += 1 + block = lines[i:j] + while block and not block[-1].strip(): + block.pop() + old = _task(block, taken) if old_task else _awaiting(block, taken) + if old is None: + report.notes.append(f"line {i + 1}: numbered item not understood, left as it is: {line}") + out += lines[i:j] + else: + taken.add(old.item.id) + report.ids[old.title] = old.item.id + out += [old, ""] + i = j + + by_title = {t.lower().rstrip("."): id for t, id in report.ids.items()} + result: list[str] = [] + for entry in out: + if isinstance(entry, str): + result.append(entry) + continue + item, ids = entry.item, [] + for wanted in entry.after: + whole = by_title.get(wanted.lower().rstrip(".")) + parts = [whole] if whole else [by_title.get(p.strip().lower().rstrip(".")) for p in wanted.split("; ")] + if not any(parts): + report.notes.append(f"{item.id}: dropped 'After: {wanted}' (no open task with that title: done)") + ids += [p for p in parts if p] + if ids: + item.body.append(f"{tasks.INDENT}- After: " + ", ".join(f"[[{a}]]" for a in ids)) + if entry.refs: + item.body.append(f"{tasks.INDENT}Ref: {entry.refs}") + result += item.lines() + + doc = tasks.parse("\n".join(result) + "\n") + for key, heading in tasks.SECTIONS.items(): + if not any(s.key == key for s in doc.sections): + if doc.sections and doc.sections[-1].lines()[-1].strip(): + doc.sections[-1].suffix.append("") if doc.sections[-1].key else doc.sections[-1].prefix.append("") + doc.sections.append(tasks.Section(heading, key, prefix=[""])) + doc.newline = newline + return tasks.render(doc), report diff --git a/wflib/refs.py b/wflib/refs.py new file mode 100644 index 0000000..d1e853d --- /dev/null +++ b/wflib/refs.py @@ -0,0 +1,107 @@ +"""Links, anchors and doc sections: what `[[id]]` and `Ref: path#anchor` point at.""" +from __future__ import annotations + +import re +from dataclasses import dataclass +from pathlib import Path + +LINK_RE = re.compile(r"\[\[([^\]|#]+)(?:#[^\]|]+)?(?:\|[^\]]+)?\]\]") +HEADING_RE = re.compile(r"^(#{1,6})\s+(.*?)\s*$") +ANCHOR_RE = re.compile(r'<a\s+id="([^"]+)"\s*>\s*</a>', re.I) +FENCE_RE = re.compile(r"^\s*(```|~~~)") + + +def strip_code(text: str) -> str: + text = re.sub(r"```.*?```", "", text, flags=re.S) + return re.sub(r"`[^`\n]*`", "", text) + + +def links(text: str) -> list[str]: + return [m.group(1).strip() for m in LINK_RE.finditer(strip_code(text))] + + +def slugify(heading: str) -> str: + text = ANCHOR_RE.sub("", heading).strip().lower() + text = re.sub(r"[^\w\s-]", "", text) + return re.sub(r"[\s_]+", "-", text).strip("-").replace("--", "-") + + +@dataclass +class DocSection: + level: int # 0 = an anchor standing alone, not a heading + heading: str + anchors: set[str] + start: int # 0-based line of the heading (or of the lone anchor) + end: int # exclusive + + +def _lines(text: str) -> list[str]: + lines = text.replace("\r\n", "\n").split("\n") + while lines and not lines[-1].strip(): + lines.pop() + return lines + + +def _is_lone_anchor(line: str) -> bool: + return bool(ANCHOR_RE.search(line)) and not ANCHOR_RE.sub("", line).strip() + + +def sections(text: str) -> list[DocSection]: + lines = _lines(text) + out: list[DocSection] = [] + fenced = False + for n, line in enumerate(lines): + if FENCE_RE.match(line): + fenced = not fenced + if fenced or FENCE_RE.match(line): + continue + m = HEADING_RE.match(line) + tags = set(ANCHOR_RE.findall(line)) + if m: + heading = ANCHOR_RE.sub("", m.group(2)).strip() + above = set(ANCHOR_RE.findall(lines[n - 1])) if n and _is_lone_anchor(lines[n - 1]) else set() + out.append(DocSection(len(m.group(1)), heading, {slugify(heading)} | tags | above, n, len(lines))) + elif tags and not (_is_lone_anchor(line) and n + 1 < len(lines) and HEADING_RE.match(lines[n + 1])): + out.append(DocSection(0, "", tags, n, len(lines))) + headings = [s for s in out if s.level] + for s in out: + nxt = next((h for h in headings if h.start > s.start and (s.level == 0 or h.level <= s.level)), None) + if nxt: + above = nxt.start - 1 + s.end = above if above > s.start and _is_lone_anchor(lines[above]) else nxt.start + return out + + +def anchors(text: str) -> set[str]: + return {a for s in sections(text) for a in s.anchors} + + +def find(text: str, anchor: str) -> DocSection | None: + found = sections(text) + return next((s for s in found if anchor in s.anchors and s.level), None) or \ + next((s for s in found if anchor in s.anchors), None) + + +def resolve(root: Path, path: str, anchor: str | None, limit: int = 80) -> str: + target = root / path + shown = path.rstrip("/") + if target.is_dir(): + return f"{shown}/ (directory)" + if not target.is_file(): + return f"(missing: {shown})" + text = target.read_text(encoding="utf-8", errors="replace") + lines = _lines(text) + if anchor is None: + heading = next((m.group(2) for l in lines if (m := HEADING_RE.match(l))), "") + goal = next((l[len("**Goal:**"):].strip() for l in lines if l.startswith("**Goal:**")), "") + return shown + (f": {heading}" if heading else "") + (f" — {goal}" if goal else "") + s = find(text, anchor) + if s is None: + return f"(no anchor '{anchor}' in {shown})" + body = lines[s.start:s.end] + while body and not body[-1].strip(): + body.pop() + out = [f"===== {shown}#{anchor} (line {s.start + 1}) =====", *body[:limit]] + if len(body) > limit: + out.append(f"… ({len(body) - limit} more lines: {shown}:{s.start + limit + 1})") + return "\n".join(out) diff --git a/wflib/res.py b/wflib/res.py new file mode 100644 index 0000000..fd1fba9 --- /dev/null +++ b/wflib/res.py @@ -0,0 +1,993 @@ +"""wf res — pure part: parsing, ledger, liveness, capacity, queue, game, scratch, output lines. +Text and data in, data and lines out; no files, no processes. +Spec: docs/resource-ledger.md.""" +from __future__ import annotations + +import datetime as dt +import json +import math +import re +import shlex +import tomllib +from dataclasses import dataclass, field, fields +from typing import Callable + +GB = 1024 ** 3 +SLICE = "agents.slice" +JOBS_SLICE = "agents-jobs.slice" +DONE_KEEP = dt.timedelta(hours=24) +CLEAN_EVERY = dt.timedelta(minutes=10) +EPS = 1e-9 + + +class ResError(Exception): + """Expected failure: `wf: <msg>`, exit 1.""" + + +class Busy(Exception): + """Request does not fit now: the busy line, exit 3.""" + + +# ---------------------------------------------------------------- parsing + +SIZE_RE = re.compile(r"^(\d+(?:\.\d+)?)([MG])B?$", re.I) +DUR_RE = re.compile(r"^(?:(\d+)h)?(?:(\d+)m)?$") + + +def parse_size(text: str) -> float: + """'10G' → 10.0, '512M' → 0.5 (GiB).""" + m = SIZE_RE.match(text.strip()) + if not m: + raise ResError(f"size '{text}' (want e.g. 10G or 512M)") + n = float(m.group(1)) + return n if m.group(2).upper() == "G" else n / 1024 + + +def parse_duration(text: str) -> int: + """'40m' → 40, '4h' → 240, '1h30m' → 90, '90' → 90 (minutes, > 0).""" + t = text.strip() + if t.isdigit(): + n = int(t) + else: + m = DUR_RE.match(t) + if not t or not m or not any(m.groups()): + raise ResError(f"duration '{text}' (want e.g. 40m, 4h, 1h30m)") + n = int(m.group(1) or 0) * 60 + int(m.group(2) or 0) + if n <= 0: + raise ResError(f"duration '{text}' must be > 0") + return n + + +def meminfo(text: str) -> tuple[float, float]: + """(MemTotal, MemAvailable) in GB from /proc/meminfo.""" + vals = {} + for line in text.splitlines(): + key, _, rest = line.partition(":") + if key in ("MemTotal", "MemAvailable"): + vals[key] = int(rest.split()[0]) * 1024 / GB + if len(vals) != 2: + raise ResError("no MemTotal/MemAvailable in /proc/meminfo") + return vals["MemTotal"], vals["MemAvailable"] + + +def show_units(text: str) -> dict[str, dict[str, str]]: + """`systemctl show -p Id,…` output (blank-line separated blocks) → {Id: {prop: value}}.""" + out = {} + for block in text.strip().split("\n\n"): + props = dict(l.split("=", 1) for l in block.splitlines() if "=" in l) + if "Id" in props: + out[props["Id"]] = props + return out + + +def psi_some_avg60(text: str) -> float | None: + """cgroup memory.pressure → `some avg60` (percent).""" + for line in text.splitlines(): + if line.startswith("some "): + vals = dict(kv.split("=", 1) for kv in line.split()[1:] if "=" in kv) + try: + return float(vals["avg60"]) + except (KeyError, ValueError): + return None + return None + + +def cache_gb(stat: str) -> float: + """cgroup memory.stat → reclaimable file cache GB: `file` − `shmem` (tmpfs pages count as file but stay).""" + vals = dict(l.split(" ", 1) for l in stat.splitlines() if " " in l) + def get(k): + return int(vals[k]) if vals.get(k, "").strip().isdigit() else 0 + return max(0, get("file") - get("shmem")) / GB + + +def gb_or_none(text: str) -> float | None: + return int(text) / GB if text.strip().isdigit() else None + + +# ---------------------------------------------------------------- formatting + +def fmt_gb(x: float) -> str: + return f"{x:.1f} GB" + + +def fmt_bytes(n: int) -> str: + return fmt_gb(n / GB) if n >= GB else f"{n / 1024 ** 2:.0f} MB" + + +def fmt_dur(minutes: int) -> str: + h, m = divmod(int(minutes), 60) + return (f"{h}h" if h else "") + (f"{m}m" if m or not h else "") + + +def hhmm(t: dt.datetime) -> str: + return t.strftime("%H:%M") + + +def rc_text(rc: int | None) -> str: + return "?" if rc is None else str(rc) + + +# ---------------------------------------------------------------- config + +@dataclass +class Config: + user_reserve_gb: float = 6.0 + user_reserve_cpus: int = 4 + game_reserve_gb: float = 12.0 + game_reserve_cpus: int = 8 + game_hours: float = 4.0 + small_headroom_gb: float = 2.0 + scratch_hours: float = 2.0 + session_mem_gb: float = 6.0 # MemoryHigh per Claude session scope (throttle, no kill); 0 = off + + +def load_config(text: str) -> Config: + try: + data = tomllib.loads(text) + except tomllib.TOMLDecodeError as e: + raise ResError(f"resources.toml: {e}") from None + cfg = Config() + unknown = sorted(set(data) - {f.name for f in fields(Config)}) + if unknown: + raise ResError("resources.toml: unknown key(s) " + ", ".join(unknown)) + for key, value in data.items(): + if isinstance(value, bool) or not isinstance(value, (int, float)) or value < 0: + raise ResError(f"resources.toml: {key} must be a number ≥ 0") + setattr(cfg, key, type(getattr(cfg, key))(value)) + return cfg + + +# ---------------------------------------------------------------- ledger + +TOP_KEYS = ("next", "game_until", "last_clean", "entries") +TIME_FIELDS = ("queued", "started", "expires", "ended") + + +@dataclass +class Entry: + id: str + project: str + owner: int + title: str + mem_gb: float + cpus: int + est_min: int + state: str # queued | running | note | done + cwd: str = "" + cmd: list[str] = field(default_factory=list) + unit: str = "" + log: str = "" + queued: dt.datetime | None = None + started: dt.datetime | None = None + expires: dt.datetime | None = None + ended: dt.datetime | None = None + rc: int | None = None + peak_gb: float | None = None + why: str = "" + env: dict = field(default_factory=dict) # caller variables the user manager lacks (env_diff) + by: dict = field(default_factory=dict) # who started it (owner_by): name, task, batch, address + lock: str = "" # '<key>@<main tree>': one queued/running job per lock (lock_for) + extra: dict = field(default_factory=dict, repr=False, compare=False) # unknown keys of a newer version: kept on rewrite + + def eta(self) -> dt.datetime | None: + return self.started + dt.timedelta(minutes=self.est_min) if self.started else None + + +@dataclass +class Ledger: + next: int = 1 + game_until: dt.datetime | None = None + last_clean: dt.datetime | None = None + entries: list[Entry] = field(default_factory=list) + extra: dict = field(default_factory=dict, repr=False, compare=False) # unknown top-level keys: kept on rewrite + + def get(self, id: str) -> Entry: + for e in self.entries: + if e.id == id: + return e + raise ResError(f"no entry '{id}'") + + def new_id(self) -> str: + id = f"r-{self.next}" + self.next += 1 + return id + + +def _time(v): + return dt.datetime.fromisoformat(v) if v else None + + +def _iso(t): + return t.isoformat(timespec="seconds") if t else None + + +def entry_dict(e: Entry) -> dict: + d = {f.name: getattr(e, f.name) for f in fields(Entry) if f.name != "extra"} + d.update({k: v for k, v in e.extra.items() if k not in d}) + for k in TIME_FIELDS: + d[k] = _iso(d[k]) + return d + + +def loads(text: str) -> Ledger: + if not text.strip(): + return Ledger() + try: + d = json.loads(text) + entries, known = [], {f.name for f in fields(Entry)} - {"extra"} + for raw in d.get("entries", []): + extra = {k: v for k, v in raw.items() if k not in known} # newer version's keys: kept, not fatal + raw = {k: v for k, v in raw.items() if k in known} + raw["extra"] = extra + for k in TIME_FIELDS: + raw[k] = _time(raw.get(k)) + entries.append(Entry(**raw)) + return Ledger(next=int(d.get("next", 1)), game_until=_time(d.get("game_until")), + last_clean=_time(d.get("last_clean")), entries=entries, + extra={k: v for k, v in d.items() if k not in TOP_KEYS}) + except (ValueError, TypeError, AttributeError) as e: + raise ResError(f"corrupt ledger: {e}") from None + + +def salvage_next(text: str) -> int: + """Id counter of an unreadable ledger, so a fresh one never reuses ids (logs/<id>.rc, wf-<id>.service).""" + nums = [int(n) for n in re.findall(r'"next":\s*(\d+)', text)] + \ + [int(n) + 1 for n in re.findall(r'"r-(\d+)"', text)] + return max(nums, default=1) + + +def dumps(led: Ledger) -> str: + return json.dumps({**{k: v for k, v in led.extra.items() if k not in TOP_KEYS}, "next": led.next, "game_until": _iso(led.game_until), "last_clean": _iso(led.last_clean), + "entries": [entry_dict(e) for e in led.entries]}, indent=1) + "\n" + + +# ---------------------------------------------------------------- facts, liveness + +@dataclass +class Unit: + active: bool + current_gb: float | None = None + stall: float | None = None # memory.pressure `some avg60`: % of the last minute stalled on memory + + +@dataclass +class Facts: + """Snapshot of the machine, gathered by wf_res under the lock.""" + now: dt.datetime + total_gb: float + available_gb: float + nproc: int + units: dict[str, Unit] = field(default_factory=dict) # unit name → state, for running entries + live: set[int] = field(default_factory=set) # note owner pids still alive + results: dict[str, tuple] = field(default_factory=dict) # id → (rc, peak_gb, why) from logs/, journal + slice_gb: float = 0.0 # agents.slice MemoryCurrent + slice_cache_gb: float = 0.0 # of it reclaimable file cache (cache_gb) + + +def _finish(e: Entry, now: dt.datetime, why: str) -> None: + e.state, e.ended, e.why = "done", now, why + + +def prune(led: Ledger, facts: Facts) -> list[str]: + """Liveness and expiry. Running jobs past their ETA are kept (never killed).""" + out = [] + now = facts.now + for e in led.entries: + if e.state == "running": + unit = facts.units.get(e.unit) + if unit is None or not unit.active: + e.rc, e.peak_gb, why = facts.results.get(e.id, (None, None, "")) + _finish(e, now, why or "exited") + out.append(f"{e.id} {why}" if why else f"{e.id} exited rc={rc_text(e.rc)}") + elif e.state == "note": + if e.owner not in facts.live: + _finish(e, now, "owner gone") + out.append(f"{e.id} freed (owner gone)") + elif e.expires and now >= e.expires: + _finish(e, now, "expired") + out.append(f"{e.id} freed (expired)") + led.entries = [e for e in led.entries if not (e.state == "done" and e.ended and now - e.ended > DONE_KEEP)] + if led.game_until and now >= led.game_until: + led.game_until = None + out.append("game off (expired)") + return out + + +# ---------------------------------------------------------------- capacity + +def gaming(led: Ledger, now: dt.datetime) -> bool: + return bool(led.game_until and now < led.game_until) + + +def reserve(cfg: Config, led: Ledger, now: dt.datetime) -> tuple[float, int]: + if gaming(led, now): + return cfg.game_reserve_gb, cfg.game_reserve_cpus + return cfg.user_reserve_gb, cfg.user_reserve_cpus + + +def _used(e: Entry, facts: Facts) -> float: + unit = facts.units.get(e.unit) + return unit.current_gb if unit and unit.current_gb is not None else 0.0 + + +def held_gb(e: Entry, facts: Facts) -> float: + """What an entry still claims beyond what MemAvailable already shows as used.""" + if e.state == "note": + return e.mem_gb + if e.state == "running": + return max(0.0, e.mem_gb - _used(e, facts)) + return 0.0 + + +def frees_gb(e: Entry, facts: Facts) -> float: + """Budget gained when the entry ends.""" + return max(e.mem_gb, _used(e, facts)) if e.state == "running" else e.mem_gb + + +def _claims(led: Ledger) -> list[Entry]: + return [e for e in led.entries if e.state in ("running", "note")] + + +def budget(cfg: Config, led: Ledger, facts: Facts) -> tuple[float, int]: + res_gb, res_cpus = reserve(cfg, led, facts.now) + claims = _claims(led) + mem = facts.available_gb - res_gb - cfg.small_headroom_gb - sum(held_gb(e, facts) for e in claims) + return round(mem, 6), facts.nproc - res_cpus - sum(e.cpus for e in claims) + + +def force_room(cfg: Config, led: Ledger, facts: Facts) -> tuple[float, int]: + """Budget for --force: what the machine really has free beyond the user reserve and headroom; ledger + claims (unused reservations, notes) and the queue are ignored.""" + res_gb, res_cpus = reserve(cfg, led, facts.now) + return round(facts.available_gb - res_gb - cfg.small_headroom_gb, 6), facts.nproc - res_cpus + + +def fits(mem_gb: float, cpus: int, b: tuple[float, int]) -> bool: + return mem_gb <= b[0] + EPS and cpus <= b[1] + + +def queued_total(led: Ledger) -> tuple[float, int]: + q = [e for e in led.entries if e.state == "queued"] + return sum(e.mem_gb for e in q), sum(e.cpus for e in q) + + +def never_fits(cfg: Config, facts: Facts, mem_gb: float, cpus: int) -> str | None: + """Larger than the whole agent budget on an empty machine (normal reserve).""" + max_gb = facts.total_gb - cfg.user_reserve_gb - cfg.small_headroom_gb + max_cpus = facts.nproc - cfg.user_reserve_cpus + if mem_gb > max_gb + EPS: + return f"{fmt_gb(mem_gb)} can never fit (max {fmt_gb(max_gb)} for agents)" + if cpus > max_cpus: + return f"{cpus} cpus can never fit (max {max_cpus} for agents)" + return None + + +def end_time(e: Entry) -> dt.datetime: + return (e.expires if e.state == "note" else e.eta()) or e.started + + +def until(e: Entry, now: dt.datetime) -> str: + t = end_time(e) + return f"~{hhmm(t)}" if t > now else "overdue" + + +def needed(led: Ledger, facts: Facts, need_gb: float, need_cpus: int) -> tuple[list[Entry], bool]: + """Claims that must end (earliest first) to free need_gb and need_cpus; and whether that is enough.""" + got_gb, got_cpus, out = 0.0, 0, [] + for e in sorted(_claims(led), key=end_time): + if got_gb >= need_gb - EPS and got_cpus >= need_cpus: + break + out.append(e) + got_gb += frees_gb(e, facts) + got_cpus += e.cpus + return out, got_gb >= need_gb - EPS and got_cpus >= need_cpus + + +def force_hint(cfg: Config, led: Ledger, facts: Facts, mem_gb: float, cpus: int) -> str: + """Busy-line suffix naming --force when the request fits the memory that is really free.""" + room = force_room(cfg, led, facts) + if not fits(mem_gb, cpus, room): + return "" + return (f"; {fmt_gb(room[0])} really free beyond the reserve: --force starts it past the ledger " + "(only if the holders will not use what they reserved)") + + +def busy_line(cfg: Config, led: Ledger, facts: Facts, mem_gb: float, cpus: int, hint: str = "") -> str: + b_gb, b_cpus = budget(cfg, led, facts) + need_gb, need_cpus = mem_gb - b_gb, cpus - b_cpus + holders, _ = needed(led, facts, need_gb, need_cpus) + if not holders: + return (f"busy: only {fmt_gb(max(0.0, b_gb))} free for agents and no agent job holds any; " + "other programs use the rest; retry later or work on something else" + hint) + by_mem = need_gb > EPS + + def held(e): + what = fmt_gb(frees_gb(e, facts)) if by_mem else f"{e.cpus} cpus" + return f'{what} held by {e.project} "{e.title}" ({e.id}) until {until(e, facts.now)}' + + free = fmt_gb(max(0.0, b_gb)) if by_mem else f"{max(0, b_cpus)} cpus" + last = max(end_time(e) for e in holders) + retry = f"after ~{hhmm(last)}" if last > facts.now else "later" + return (f"busy: {', '.join(held(e) for e in holders)}; {free} free for agents; " + f"retry {retry} or work on something else" + hint) + + +# ---------------------------------------------------------------- locks + +LOCK_TITLE_RE = re.compile(r"gate\b", re.I) # 'gate …' titles lock 'gate' unasked: one shared gate checkout + + +def lock_for(title: str, key: str, main: str) -> str: + """Ledger lock of a run: --lock KEY, else 'gate' for a title starting with the word gate; '' = none. + Scoped to the main tree (lane worktrees share their project's gate dir).""" + key = key or ("gate" if LOCK_TITLE_RE.match(title.strip()) else "") + return f"{key}@{main}" if key else "" + + +def lock_holder(led: Ledger, lock: str) -> Entry | None: + """The running (else first queued) entry holding this lock.""" + if not lock: + return None + mine = [e for e in led.entries if e.lock == lock and e.state in ("running", "queued")] + return next((e for e in mine if e.state == "running"), None) or next(iter(_queue_of(mine)), None) + + +def _queue_of(entries: list[Entry]) -> list[Entry]: + return sorted((e for e in entries if e.state == "queued"), key=lambda e: (e.queued, int(e.id[2:]))) + + +def lock_busy_line(h: Entry, facts: Facts) -> str: + key = h.lock.split("@", 1)[0] + when = f"ETA {until(h, facts.now)}" if h.state == "running" else "queued" + return (f"busy: lock '{key}' held by {h.id} \"{h.title}\" ({when}): one at a time (shared checkout); " + f"--queue waits for it (--force does not override a lock); never a hand-written copy without the lock") + + +# ---------------------------------------------------------------- run, status lines + +RES_ID_VAR = "WF_RES_ID" # set in every job's unit: wf res calls from inside it name it as their batch +ENV_SKIP = {"PWD", "OLDPWD", "SHLVL", "_", RES_ID_VAR} +SECRET_RE = re.compile(r"TOKEN|SECRET|PASSWORD|PASSWD|CREDENTIAL", re.I) # never written to the ledger + + +def env_diff(caller: dict[str, str], manager: dict[str, str]) -> dict[str, str]: + """Caller variables a systemd --user unit would not get as is (it starts from the manager's environment).""" + return {k: v for k, v in caller.items() + if manager.get(k) != v and k not in ENV_SKIP and not SECRET_RE.search(k)} + + +def parse_show_environment(text: str) -> dict[str, str]: + return dict(l.split("=", 1) for l in text.splitlines() if "=" in l) + + +def run_argv(e: Entry, logdir: str) -> list[str]: + """systemd-run argv; the sh wrapper records exit code and peak bytes (unit is collected on exit).""" + rc, peak = shlex.quote(f"{logdir}/{e.id}.rc"), shlex.quote(f"{logdir}/{e.id}.peak") + script = (f'"$@"; rc=$?; cat /sys/fs/cgroup$(cut -d: -f3 /proc/self/cgroup)/memory.peak > {peak} ' + f'2>/dev/null; echo $rc > {rc}') + return ["systemd-run", "--user", "--quiet", "--collect", f"--slice={JOBS_SLICE}", f"--unit={e.unit}", + f"--working-directory={e.cwd}", *(f"--setenv={k}={v}" for k, v in sorted(e.env.items())), + f"--setenv={RES_ID_VAR}={e.id}", + "-p", f"MemoryMax={int(e.mem_gb * GB)}", "-p", f"MemoryHigh={int(e.mem_gb * 0.9 * GB)}", + "-p", "MemorySwapMax=0", "-p", "Nice=10", + "-p", f"StandardOutput=append:{e.log}", "-p", f"StandardError=append:{e.log}", + "/bin/sh", "-c", script, "sh", *e.cmd] + + +TASK_BRANCH_RE = re.compile(r"^[\w.-]+/([\w.-]+)$") # <lane>/<task id> (wf worktree branches) + + +def owner_by(environ: dict[str, str], branch: str, records: list[dict], name: str = "") -> dict[str, str]: + """Who starts an entry, empty keys dropped. name: --by, else WF_SESSION_NAME, else '<lane> session' from the + wf lane session record (.wf/sessions) with this CLAUDE_PID; task: WF_TASK, else the <lane>/<id> branch; + batch: WF_RES_ID (the wf res job, e.g. wf batch, this runs in); address: uds:<CLAUDE_CODE_MESSAGING_SOCKET> + (a batch worker's = its orchestrator's: SendMessage there reaches the batch).""" + pid = environ.get("CLAUDE_PID", "") + name = name or environ.get("WF_SESSION_NAME", "") + if not name and pid: + name = next((f"{r['lane']} session" for r in records + if isinstance(r, dict) and r.get("lane") and str(r.get("pid")) == pid), "") + m = TASK_BRANCH_RE.match(branch) + sock = environ.get("CLAUDE_CODE_MESSAGING_SOCKET", "") + d = {"name": name, "task": environ.get("WF_TASK") or (m.group(1) if m else ""), + "batch": environ.get(RES_ID_VAR, ""), "address": f"uds:{sock}" if sock else ""} + return {k: v for k, v in d.items() if v} + + +def by_text(e: Entry) -> str: + """'slow session t-x batch r-3, message uds:/s' — empty when nothing is known.""" + b = e.by or {} + who = " ".join(x for x in (b.get("name"), b.get("task"), f"batch {b['batch']}" if b.get("batch") else "") if x) + msg = f"message {b['address']}" if b.get("address") else "" + return ", ".join(x for x in (who, msg) if x) + + +def position(led: Ledger, e: Entry) -> int: + return _queue(led).index(e) + 1 + + +def entry_line(e: Entry, led: Ledger, facts: Facts) -> str: + by = by_text(e) if e.state != "done" else "" + return _entry_line(e, led, facts) + (f" [by {by}]" if by else "") + + +def _entry_line(e: Entry, led: Ledger, facts: Facts) -> str: + head = f'{e.id} {e.project} "{e.title}"' + if e.state == "running": + unit = facts.units.get(e.unit) + used = fmt_gb(unit.current_gb) if unit and unit.current_gb is not None else "?" + return (f"{head} running {fmt_gb(e.mem_gb)} used {used} {e.cpus} cpu since {hhmm(e.started)} " + f"ETA {until(e, facts.now)}" + (" (still running, not killed)" if end_time(e) <= facts.now else "")) + if e.state == "queued": + return f"{head} queued #{position(led, e)} {fmt_gb(e.mem_gb)} {e.cpus} cpu" + if e.state == "note": + return f"{head} note {fmt_gb(e.mem_gb)} {e.cpus} cpu until {until(e, facts.now)}" + peak = fmt_gb(e.peak_gb) if e.peak_gb is not None else "?" + return f"{head} done ({e.why}) rc={rc_text(e.rc)} peak {peak} at {hhmm(e.ended)}" + + +def status_lines(cfg: Config, led: Ledger, facts: Facts) -> list[str]: + active = [e for e in led.entries if e.state != "done"] + done = [e for e in led.entries if e.state == "done"] + out = [entry_line(e, led, facts) for e in active + done] or ["no reservations"] + b_gb, b_cpus = budget(cfg, led, facts) + r_gb, r_cpus = reserve(cfg, led, facts.now) + out.append(f"agents may use {fmt_gb(max(0.0, b_gb))}, {max(0, b_cpus)} cpus now; " + f"reserve {fmt_gb(r_gb)}/{r_cpus} cpus") + if facts.slice_gb > 0: + jobs = sum(_used(e, facts) for e in led.entries if e.state == "running") + free = max(0.0, facts.slice_gb - jobs) + cache = min(free, facts.slice_cache_gb) + out.append(f"unreserved agent memory {fmt_gb(free)}" + + (f" ({fmt_gb(cache)} of it file cache, reclaimable)" if cache >= 0.05 else "")) + if gaming(led, facts.now): + left = int((led.game_until - facts.now).total_seconds() // 60) + out.append(f"game on until {hhmm(led.game_until)} ({fmt_dur(left)} left)") + out.extend(shortfall_lines(cfg, led, facts)) + return out + + +def status_json(cfg: Config, led: Ledger, facts: Facts) -> str: + b_gb, b_cpus = budget(cfg, led, facts) + return json.dumps({"budget_gb": b_gb, "budget_cpus": b_cpus, "gaming": gaming(led, facts.now), + "game_until": _iso(led.game_until), "entries": [entry_dict(e) for e in led.entries]}, + indent=1) + + +def detail_lines(e: Entry) -> list[str]: + d = entry_dict(e) + d["cmd"] = shlex.join(e.cmd) + return [f"{k}: {v}" for k, v in d.items() if v not in (None, "", [])] + + +def slice_unit_text() -> str: + return ("[Unit]\nDescription=wf agents: Claude sessions and wf res jobs\n\n" + "[Slice]\nCPUWeight=20\nIOWeight=20\n") + + +def done_line(e: Entry) -> str: + minutes = round((e.ended - e.started).total_seconds() / 60) if e.ended and e.started else 0 + peak = fmt_gb(e.peak_gb) if e.peak_gb is not None else "?" + why = f" ({e.why})" if e.why and e.why != "exited" else "" + return f"{e.id} done rc={rc_text(e.rc)} peak {peak} in {minutes} min{why}" + + +JOURNAL_RESULT = re.compile(r"Failed with result '([^']+)'") +JOURNAL_MAIN = re.compile(r"Main process exited, code=\w+, status=(\S+)") +JOURNAL_PEAK = re.compile(r"(\d+(?:\.\d+)?)([BKMGT]) memory peak") +UNIT_SCALE = {"B": 1 / GB, "K": 1 / 1024 ** 2, "M": 1 / 1024, "G": 1.0, "T": 1024.0} + + +def journal_reason(text: str, mem_gb: float) -> tuple[str, float | None]: + """Why a unit died, from `journalctl -u UNIT -o cat` (used when the job wrote no .rc): (why, peak_gb).""" + m = JOURNAL_PEAK.search(text) + peak = round(float(m.group(1)) * UNIT_SCALE[m.group(2)], 3) if m else None + m = JOURNAL_RESULT.search(text) + result = m.group(1) if m else "" + m = JOURNAL_MAIN.search(text) + status = m.group(1) if m else "?" + if result == "oom-kill": + by = "by systemd-oomd" if "systemd-oomd killed" in text else "at MemoryMax" + return f"killed: oom-kill {by}, limit {fmt_gb(mem_gb)}; raise --mem", peak + if result in ("signal", "core-dump"): + return f"killed: signal {status}", peak + if result == "exit-code": + return f"failed: exit {status}", peak + return (f"killed: {result}" if result else ""), peak + + +def _queue(led: Ledger) -> list[Entry]: + return _queue_of(led.entries) + + +def to_start(cfg: Config, led: Ledger, facts: Facts) -> list[Entry]: + """Strict FIFO: queued entries to start now, stopping at the first that does not fit; an entry whose lock a + running job (or one started in this pass) holds is skipped, not blocking the rest.""" + b_gb, b_cpus = budget(cfg, led, facts) + out, held = [], {e.lock for e in led.entries if e.state == "running" and e.lock} + for e in _queue(led): + if e.lock and e.lock in held: + continue + if not fits(e.mem_gb, e.cpus, (b_gb, b_cpus)): + break + out.append(e) + held.add(e.lock) + b_gb, b_cpus = b_gb - e.mem_gb, b_cpus - e.cpus + return out + + +def queue_estimate(cfg: Config, led: Ledger, facts: Facts, e: Entry) -> dt.datetime | None: + """When enough claims end for e and everything ahead of it; None if they never free enough.""" + q = _queue(led) + ahead = q[:q.index(e) + 1] + b_gb, b_cpus = budget(cfg, led, facts) + holders, enough = needed(led, facts, sum(x.mem_gb for x in ahead) - b_gb, sum(x.cpus for x in ahead) - b_cpus) + if not enough: + return None + lock = [end_time(x) for x in led.entries if x is not e and x.lock and x.lock == e.lock and x.state == "running"] + return max([*(end_time(h) for h in holders), *lock], default=facts.now) + + +def slice_props(cfg: Config, led: Ledger, facts: Facts) -> dict[str, str]: + """agents.slice values: normal, or gaming (MemoryHigh never below current use: slow, never squeeze).""" + if gaming(led, facts.now): + high, weight = max(facts.total_gb - cfg.game_reserve_gb, facts.slice_gb), "5" + else: + high, weight = facts.total_gb - cfg.user_reserve_gb, "20" + return {"CPUWeight": weight, "IOWeight": weight, "MemoryHigh": str(int(round(high * GB)))} + + +def shortfall_lines(cfg: Config, led: Ledger, facts: Facts) -> list[str]: + """Empty when the game reserve is free now; else who holds it and when it frees.""" + free = facts.available_gb - sum(held_gb(e, facts) for e in _claims(led)) + short = round(cfg.game_reserve_gb - free, 6) + if short <= EPS: + return [] + holders, enough = needed(led, facts, short, 0) + head = f"short {fmt_gb(short)} of {fmt_gb(cfg.game_reserve_gb)}: " + if not holders: + return [head + "other programs, not agent jobs", "agent jobs alone cannot free it; close other programs"] + parts = ", ".join(f'{e.id} {e.project} "{e.title}" {fmt_gb(frees_gb(e, facts))} {until(e, facts.now)}' + for e in holders) + if not enough: + return [head + parts, "agent jobs alone cannot free it; close other programs"] + last = max(end_time(e) for e in holders) + big = max(holders, key=lambda e: frees_gb(e, facts)) + when = f"~{hhmm(last)}" if last > facts.now else "soon" + return [head + parts, f"full reserve free {when} (est.); free now: wf res release {big.id} --stop"] + + +def game_on_lines(cfg: Config, led: Ledger, facts: Facts, minutes: int) -> list[str]: + return [f"game on until {hhmm(led.game_until)} ({fmt_dur(minutes)}); CPU/IO now yours", + *shortfall_lines(cfg, led, facts)] + + +def scratch_victims(tree: dict, now_ts: float, hours: float, keep: set[str] = frozenset()) -> list[str]: + """tree = {top path: (newest mtime, {child path: newest mtime})}. A stale top goes whole; under a fresh + top, its stale children go. Children named in keep (live Claude session ids) never go; a stale top + holding one is pruned per child. Sorted.""" + limit = now_ts - hours * 3600 + out = [] + for top, (newest, children) in tree.items(): + held = any(c.rsplit("/", 1)[-1] in keep for c in children) + if newest <= limit and not held: + out.append(top) + else: + out.extend(c for c, m in children.items() if m <= limit and c.rsplit("/", 1)[-1] not in keep) + return sorted(out) + + +DOTNET_PIPE_RE = re.compile(r"^(?:clr-debug-pipe|dotnet-diagnostic)-(\d+)-") + + +def litter_victims(entries: list[tuple[str, str, float]], now_ts: float, hours: float, + alive: Callable[[int], bool]) -> list[str]: + """Top-level temp-dir litter (own entries only): (path, kind emptydir|dir|file|other, mtime) → sorted paths. + Empty dirs idle for `hours`; .NET debug pipes/sockets of dead pids. Content is never swept.""" + limit = now_ts - hours * 3600 + out = [] + for path, kind, mtime in entries: + m = DOTNET_PIPE_RE.match(path.rsplit("/", 1)[-1]) + if (kind == "emptydir" and mtime <= limit) or (m and kind != "dir" and not alive(int(m.group(1)))): + out.append(path) + return sorted(out) + + +def live_session_id(text: str, proc_start: str) -> str | None: + """~/.claude/sessions/<pid>.json of a live pid → its sessionId (= its scratch dir name), unless the + recorded procStart shows the pid now belongs to another process.""" + try: + d = json.loads(text) + except ValueError: + return None + if not isinstance(d, dict) or not isinstance(d.get("sessionId"), str): + return None + if d.get("procStart") is not None and str(d["procStart"]) != proc_start: + return None + return d["sessionId"] + + +AGE_RE = re.compile(r"^(.*):(\d+)d$") + + +def cleanup_rule(text: str) -> tuple[str, int | None]: + """'out/logs/*.log:30d' → ('out/logs/*.log', 30); relative, no '..'.""" + m = AGE_RE.match(text) + pattern, days = (m.group(1), int(m.group(2))) if m else (text, None) + if not pattern or pattern.startswith("/") or ".." in pattern.split("/"): + raise ResError(f"cleanup pattern '{text}' (want a path relative to the project, no '..')") + return pattern, days + + +CLAUDE_BIN_RE = re.compile(r"/claude/versions/[^/]+$") + + +def proc_name(comm: str, argv0: str) -> str: + """`claude` for a Claude Code process, else comm. Background sessions run the versioned binary, so comm + is the version ('2.1.283'); argv0 is …/claude/versions/X, or a rewritten title 'claude bg-pty-host'. Tools it + re-execs (ugrep) keep their own argv0.""" + if argv0.split(" ", 1)[0].rsplit("/", 1)[-1] == "claude" or CLAUDE_BIN_RE.search(argv0): + return "claude" + return comm + + +def adopt_groups(procs: list[tuple[int, int, str, str]]) -> dict[int, list[int]]: + """procs = (pid, ppid, comm, cgroup). Claude session roots outside agents.slice (named `claude`, parent not + `claude`) → sorted pids of the root and its descendants that are outside the slice.""" + inside = f"/{SLICE}/" + by_pid = {p[0]: p for p in procs} + kids: dict[int, list[int]] = {} + for pid, ppid, _, _ in procs: + kids.setdefault(ppid, []).append(pid) + out = {} + for pid, ppid, comm, cgroup in procs: + parent = by_pid.get(ppid) + if comm != "claude" or inside in cgroup or (parent and parent[2] == "claude"): + continue + tree, stack = [], [pid] + while stack: + p = stack.pop() + if inside not in by_pid[p][3]: + tree.append(p) + stack.extend(kids.get(p, [])) + out[pid] = sorted(tree) + return out + + +THROTTLE_FILL = 0.85 # used ≥ this × --mem: at MemoryHigh (0.9 × --mem), where the kernel throttles +THROTTLE_STALL = 20.0 # % of the last minute stalled on memory: page cache alone at the limit stays near 0 + + +def throttle_lines(led: Ledger, facts: Facts) -> list[str]: + """Running jobs stuck at their own memory limit (they crawl, then systemd-oomd kills them).""" + out = [] + for e in led.entries: + u = facts.units.get(e.unit) + if (e.state == "running" and u and u.active and u.current_gb is not None and u.stall is not None + and u.current_gb >= THROTTLE_FILL * e.mem_gb and u.stall >= THROTTLE_STALL): + out.append(f"{e.id} throttled at its memory limit ({u.current_gb:.1f} of {fmt_gb(e.mem_gb)}, " + f"stalled {u.stall:.0f}% of the last minute): likely too small; " + f"`wf res release {e.id} --stop` and re-run with a bigger --mem" + + (f"; started by {by_text(e)}" if by_text(e) else "")) + return out + + +DESKTOP_UNIT = ("app-", "dbus-") # transient units the desktop starts (XDG app launch, dbus activation) + + +def bare_units(blocks: dict[str, dict[str, str]]) -> list[tuple[str, float | None]]: + """Running transient services outside agents.slice not from the desktop = jobs from a bare `systemd-run`: + sorted (unit, used GB).""" + return sorted((n, gb_or_none(b.get("MemoryCurrent", ""))) for n, b in blocks.items() + if b.get("Transient") == "yes" and not n.startswith(DESKTOP_UNIT) + and not b.get("Slice", "").startswith("agents")) + + +def session_caps(blocks: dict[str, dict[str, str]], cap_gb: float) -> list[tuple[str, str]]: + """Session scopes in agents.slice whose MemoryHigh differs from the cap → sorted (unit, value) to set. + MemoryHigh only: over it the session is throttled; a MemoryMax kill could pick `claude` itself.""" + want = str(int(cap_gb * GB)) if cap_gb > 0 else "infinity" + return sorted((n, want) for n, b in blocks.items() + if n.endswith(".scope") and b.get("Slice") == SLICE and b.get("MemoryHigh") != want) + + +def session_cap_lines(blocks: dict[str, dict[str, str]], stall: dict[str, float], cap_gb: float) -> list[str]: + """Sessions stuck at their cap: used ≥ 0.9 × cap and stalled ≥ THROTTLE_STALL %.""" + if cap_gb <= 0: + return [] + out = [] + for n, b in sorted(blocks.items()): + cur = gb_or_none(b.get("MemoryCurrent", "")) + if cur is not None and cur >= 0.9 * cap_gb and stall.get(n, 0.0) >= THROTTLE_STALL: + who = n.removeprefix("wf-claude-").removesuffix(".scope") + out.append(f"warning: claude session {who} at its memory cap ({cur:.1f} of {fmt_gb(cap_gb)}, " + f"stalled {stall[n]:.0f}% of the last minute): run big work with wf res run") + return out + + +def warning_lines(total_gb: float, tmp_gb: float, outside: int, + bare: list[tuple[str, float | None]] = ()) -> list[str]: + out = [] + if tmp_gb > 0.25 * total_gb: + out.append(f"warning: /tmp (RAM) holds {fmt_gb(tmp_gb)}; wf res clean") + if outside: + out.append(f"warning: {outside} claude sessions outside {SLICE} (wf res adopt)") + if bare: + units = ", ".join(f"{n} {'?' if gb is None else f'{gb:.1f}'} GB" for n, gb in bare) + out.append(f"warning: {len(bare)} jobs outside wf res (bare systemd-run): {units}; " + "start jobs with wf res run, also from project scripts") + return out + + +def adopt_argv(root: int, pids: list[int]) -> list[str]: + """Move live pids into a new scope in agents.slice (systemd StartTransientUnit with PIDs).""" + return ["busctl", "--user", "call", "org.freedesktop.systemd1", "/org/freedesktop/systemd1", + "org.freedesktop.systemd1.Manager", "StartTransientUnit", "ssa(sv)a(sa(sv))", + f"wf-claude-{root}.scope", "fail", "2", "PIDs", "au", str(len(pids)), *map(str, pids), + "Slice", "s", SLICE, "0"] + + +def service_text(python: str, wf: str) -> str: + return ("[Unit]\nDescription=wf res tick\n\n[Service]\nType=oneshot\n" + f"ExecStart={python} {wf} res tick\n") + + +def timer_text() -> str: + return ("[Unit]\nDescription=wf res tick every minute\n\n[Timer]\nOnBootSec=1min\nOnUnitActiveSec=1min\n\n" + "[Install]\nWantedBy=timers.target\n") + + +def hook_output(payload: dict, game_until: dt.datetime | None, now: dt.datetime) -> dict | None: + """Claude Code PreToolUse hook: while game mode is on, Bash commands run without a display (no windows pop + up over the game). None = no output, the command runs unchanged.""" + cmd = (payload.get("tool_input") or {}).get("command") + if payload.get("tool_name") != "Bash" or not cmd or not (game_until and now < game_until): + return None + return {"hookSpecificOutput": { + "hookEventName": "PreToolUse", + "updatedInput": {**payload["tool_input"], "command": f"unset DISPLAY WAYLAND_DISPLAY; {cmd}"}, + "additionalContext": f"wf res game mode until {hhmm(game_until)}: this command has no display, so GUI " + "windows fail. Run headless/offscreen, or do other work until game off."}} + + +def shell_init_line() -> str: + return f"alias claude='systemd-run --user --scope --quiet --slice={SLICE} claude'" + + +# ---------------------------------------------------------------- history (reservation sizes) + +HIST_MIN_RUNS = 3 +HIST_MEM_PCT, HIST_MEM_X, HIST_MEM_FLOOR = 0.95, 1.15, 0.2 # suggest mem = p95 peak ×1.15, ≥ 0.2 GB +HIST_DUR_PCT, HIST_DUR_X, HIST_DUR_FLOOR = 0.90, 1.5, 5 # suggest for = p90 duration ×1.5, ≥ 5 min +HIST_OVER = 2.0 # request > 2× suggestion → hint +HIST_KEEP = 2000 # history file: newest runs kept +_PAREN_TAIL = re.compile(r"\s*\([^()]*\)\s*$") +_HEX_TAIL = re.compile(r"\s+(?=[0-9a-f]*\d)[0-9a-f]{7,40}$") + + +def title_kind(title: str) -> str: + """A title without its trailing (…) and sha/hex words: 'gate 2c91cac7 (fix x)' → 'gate'.""" + t = title.strip() + while True: + s = _HEX_TAIL.sub("", _PAREN_TAIL.sub("", t)).strip() + if s == t or not s: + return t + t = s + + +def pct(xs: list[float], p: float) -> float: + """Nearest-rank percentile.""" + s = sorted(xs) + return s[max(0, math.ceil(round(p * len(s), 9)) - 1)] + + +def history_record(e: Entry) -> dict | None: + """A finished, successful run with a measured peak → history record; else None.""" + if e.state != "done" or e.rc != 0 or e.peak_gb is None or not e.started or not e.ended: + return None + return {"id": e.id, "project": e.project, "title": e.title, "mem_gb": e.mem_gb, "est_min": e.est_min, + "rc": e.rc, "peak_gb": e.peak_gb, "min": round((e.ended - e.started).total_seconds() / 60, 2)} + + +def all_runs(records: list[dict], led: Ledger) -> list[dict]: + """History file records + the ledger's finished runs not yet in it (ids never repeat).""" + seen = {r.get("id") for r in records} + return records + [r for r in map(history_record, led.entries) if r and r["id"] not in seen] + + +def _kind_runs(runs: list[dict], project: str, kind: str) -> list[dict]: + return [r for r in runs if r.get("project") == project and title_kind(r.get("title", "")) == kind + and r.get("rc") == 0 and r.get("peak_gb") is not None and r.get("min") is not None] + + +def _ceil1(x: float) -> float: + return math.ceil(round(x * 10, 6)) / 10 + + +def suggest(runs: list[dict], project: str, kind: str) -> tuple[float, int, int] | None: + """(mem GB, minutes, n) from ≥ 3 runs of this kind in this project, else None.""" + rs = _kind_runs(runs, project, kind) + if len(rs) < HIST_MIN_RUNS: + return None + mem = max(HIST_MEM_FLOOR, _ceil1(pct([r["peak_gb"] for r in rs], HIST_MEM_PCT) * HIST_MEM_X)) + minutes = max(HIST_DUR_FLOOR, math.ceil(round(pct([r["min"] for r in rs], HIST_DUR_PCT) * HIST_DUR_X, 6))) + return mem, minutes, len(rs) + + +def hint_line(mem_gb: float, est_min: int, sug: tuple[float, int, int] | None) -> str | None: + """Request far over the history (no auto-resize): one line, else None.""" + if not sug or (mem_gb <= HIST_OVER * sug[0] + EPS and est_min <= HIST_OVER * sug[1]): + return None + return f"hint: history says ~{sug[0]:g} GB / {sug[1]} min ({sug[2]} runs)" + + +def hist_lines(runs: list[dict], project: str | None) -> list[str]: + """wf res hist: per project and title kind: n, request, peak p50/p95, estimate, duration p50/p90, suggestion.""" + keys = sorted({(r.get("project", ""), title_kind(r.get("title", ""))) for r in runs + if project is None or r.get("project") == project}) + out = [] + for p, k in keys: + rs = _kind_runs(runs, p, k) + if not rs: + continue + peaks, durs = [r["peak_gb"] for r in rs], [r["min"] for r in rs] + sug = suggest(runs, p, k) + out.append(f"{p} · {k} · n {len(rs)} · req {fmt_gb(pct([r['mem_gb'] for r in rs], 0.5))} · " + f"peak {pct(peaks, 0.5):.1f}/{pct(peaks, HIST_MEM_PCT):.1f} GB · " + f"est {fmt_dur(pct([r['est_min'] for r in rs], 0.5))} · " + f"dur {fmt_dur(round(pct(durs, 0.5)))}/{fmt_dur(round(pct(durs, HIST_DUR_PCT)))} · " + + (f"suggest {sug[0]:g} GB {fmt_dur(sug[1])}" if sug else f"suggest - (< {HIST_MIN_RUNS} runs)")) + return out or ["no finished runs with a peak yet"] + + +# ---------------------------------------------------------------- batch fit (wf batch --left / --time-left) + +TASK_P90_DEFAULT, TASK_P90_MIN_RUNS = 30, 3 # minutes per task when out/wf-orch.log has < 3 timed rows +_ORCH_ROW = re.compile(r"^\d{4}-\d\d-\d\dT\S+ (\S+) \S+ \S+ (\S+) \S+ (?:(\d+)h(\d+)m|(\d+)m(\d+)s)(?:\s|$)") + + +def orch_durations(text: str) -> list[float]: + """Minutes of done*/handback rows of local lanes in out/wf-orch.log with a real (> 0) duration.""" + out = [] + for line in text.splitlines(): + m = _ORCH_ROW.match(line) + if not m or m.group(1) == "cloud" or not (m.group(2).startswith("done") or m.group(2) == "handback"): + continue + h, hm, mm, s = (int(x or 0) for x in m.groups()[2:]) + minutes = h * 60 + hm + mm + s / 60 + if minutes > 0: + out.append(round(minutes, 4)) + return out + + +def task_p90(text: str) -> tuple[int, int]: + """(p90 task minutes rounded up, runs) from out/wf-orch.log text; < 3 runs → (30, runs).""" + durs = orch_durations(text) + if len(durs) < TASK_P90_MIN_RUNS: + return TASK_P90_DEFAULT, len(durs) + return math.ceil(round(pct(durs, 0.90), 6)), len(durs) + + +def batch_fit(k: int, left_min: int, p90_min: int) -> int: + """Tasks that fit: min(K, floor(left / p90)).""" + return max(0, min(k, left_min // max(1, p90_min))) diff --git a/wflib/search.py b/wflib/search.py new file mode 100644 index 0000000..600902d --- /dev/null +++ b/wflib/search.py @@ -0,0 +1,82 @@ +"""Ranked word search over open tasks, doc sections and the archive. + +A word matches at the start of a word, any case ("guard" finds "Guards"). +Rank: more of the words found first, then the weight of where they were +found, then open tasks before docs before archive. No index: read per call. +""" +from __future__ import annotations + +import re +from dataclasses import dataclass + +from . import refs, tasks + +W_TITLE, W_HEADING, W_GOAL, W_ARCHIVE, W_TEXT = 5, 4, 3, 2, 1 +KINDS = ("task", "doc", "archive") + + +@dataclass +class Hit: + found: int # how many of the words + score: int + kind: str + where: str # id, or path:line + label: str + line: str # the best matching line + + +def _patterns(words: list[str]) -> list[re.Pattern]: + return [re.compile(r"(?<![A-Za-z0-9])" + re.escape(w), re.I) for w in words if w.strip()] + + +def _rate(patterns: list[re.Pattern], parts: list[tuple[int, list[str]]]) -> tuple[int, int, str]: + """parts = (weight, lines). Returns (words found, score, best line).""" + found = score = 0 + best, best_weight = "", 0 + for p in patterns: + hit = max(((w, l) for w, lines in parts for l in lines if p.search(l)), key=lambda x: x[0], default=None) + if hit: + found += 1 + score += hit[0] + if hit[0] > best_weight: + best_weight, best = hit + return found, score, " ".join(best.split())[:120] + + +def search(words: list[str], doc: tasks.Doc, archive_text: str, docs: dict[str, str], + kinds: set[str] | None = None, limit: int = 15, archive_name: str = "archive") -> list[Hit]: + patterns = _patterns(words) + if not patterns: + return [] + kinds = kinds or set(KINDS) + hits: list[Hit] = [] + + def add(kind, where, label, parts, header=None): + found, score, line = _rate(patterns, parts) + if found: + in_header = header and any(line == " ".join(l.split())[:120] for w, ls in parts[:2] for l in ls) + hits.append(Hit(found, score, kind, where, label, header if in_header else line)) + + if "task" in kinds: + for item in doc.all_items(): + if item.error: + add("task", item.id, "", [(W_TEXT, [item.raw, *item.body])]) + else: + add("task", item.id, item.title, + [(W_TITLE, [item.id, item.title]), (W_GOAL, [item.goal]), (W_TEXT, item.body)], + header=item.text) + if "doc" in kinds: + for path, text in docs.items(): + lines = text.replace("\r\n", "\n").split("\n") + starts = [s for s in refs.sections(text) if s.level] + for n, s in enumerate(starts): + end = starts[n + 1].start if n + 1 < len(starts) else len(lines) + add("doc", f"{path}:{s.start + 1}", s.heading, + [(W_HEADING, [s.heading]), (W_TEXT, lines[s.start + 1:end])]) + if "archive" in kinds: + for n, line in enumerate(archive_text.replace("\r\n", "\n").split("\n"), 1): + if line.startswith("- "): + add("archive", f"{archive_name}:{n}", "", [(W_ARCHIVE, [line[2:]])]) + hits.sort(key=lambda h: (-h.found, -h.score, KINDS.index(h.kind))) + return hits[:limit] + diff --git a/wflib/tasks.py b/wflib/tasks.py new file mode 100644 index 0000000..80264c7 --- /dev/null +++ b/wflib/tasks.py @@ -0,0 +1,728 @@ +"""TASKS.md as data: parse into sections and items, change, render back. + +Pure text in, text out; no file access. Format: docs/design.md (section 3). +""" +from __future__ import annotations + +import difflib +import re +from dataclasses import dataclass, field + +SECTIONS = { + "awaiting": "Awaiting your decision", + "pending": "Pending", + "human": "Needs human", + "deferred": "Deferred", +} +TASK_KEYS = ("pending", "human", "deferred") +EFFORTS = ("<1h", "1h", "5h", "10h", "100h") +INDENT = " " + +ID_RE = re.compile(r"^[ta]-[a-z0-9]+(?:-[a-z0-9]+)*$") +LINK_RE = re.compile(r"\[\[([^\]|#]+)\]\]") +HEADER_RE = re.compile( + r"^- \*\*(?P<id>[^*\s]+)\*\*" + r"(?: \[P(?P<prio>\d)\])?" + r"(?: \((?!in progress:|blocked:)(?P<effort>[^)]*)\))?" + r"(?: \((?P<status>(?:in progress|blocked):[^)]*)\))?" + r": (?P<text>.*)$" +) +ITEM_START = "- **" +AFTER_RE = re.compile(r"^\s*(?:-\s+)?After:\s*(.*)$") +REF_RE = re.compile(r"^\s*(?:-\s+)?Ref:\s*(.*)$") +SLICES_RE = re.compile(r"^\s*-\s+Slices:\s*(.*)$") +MODEL_RE = re.compile(r"^\s*(?:-\s+)?Model:\s*(.*)$") +MODELS = ("haiku", "sonnet", "opus") # low → high; no Model line = opus +DONE_RE = re.compile(r"^\s*(?:-\s+)?Done:\s*(.*)$") +HUMAN_DONE_RE = re.compile( + r"\bowner(?:'s)?\s+(?:approves?|approval|tests?|checks?|runs?|says?|confirms?|verif(?:y|ies)|decides?|reviews?|signs?|sees?|ok)\b" + r"|\breport to\b|\bconfirm with\b|\bplease confirm\b|^\s*confirm\b", re.I) +SESSIONS_RE = re.compile(r"^\s*(?:-\s+)?Sessions:\s*(.*)$") +SESSIONS = ("parallel", "solo", "owner") # no Sessions line = parallel +CLOUD_RE = re.compile(r"^\s*(?:-\s+)?Cloud:\s*(.*)$") +CLOUDS = ("yes", "no") # no Cloud line = decided by the fit rules (wflib/cloud.py) +LEADING_ID_RE = re.compile(r"^([ta]-[a-z0-9]+(?:-[a-z0-9]+)*):\s+(.*)$") + + +class TaskError(ValueError): + pass + + +@dataclass +class Item: + id: str + prio: int | None = None + effort: str | None = None + interactive: bool = False + status: str | None = None + text: str = "" + body: list[str] = field(default_factory=list) + error: str | None = None + raw: str | None = None + line: int | None = None # 1-based line of the header in the parsed text + + def header(self) -> str: + if self.error: + return self.raw or "" + out = f"- **{self.id}**" + if self.prio is not None: + out += f" [P{self.prio}]" + if self.effort is not None: + out += f" ({self.effort}{', interactive' if self.interactive else ''})" + if self.status: + out += f" ({self.status})" + return f"{out}: {self.text}" + + def lines(self) -> list[str]: + return [self.header(), *self.body] + + @property + def title(self) -> str: + return self.text.split(". ", 1)[0].rstrip(".").strip() + + @property + def goal(self) -> str: + return self.text.split(". ", 1)[1].strip() if ". " in self.text else "" + + @property + def blocked_on(self) -> str | None: + if self.status and self.status.startswith("blocked:"): + found = LINK_RE.findall(self.status) + return found[0] if found else None + return None + + def _line(self, pattern: re.Pattern) -> int | None: + return next((i for i, l in enumerate(self.body) if pattern.match(l)), None) + + @property + def after(self) -> list[str]: + i = self._line(AFTER_RE) + return LINK_RE.findall(self.body[i]) if i is not None else [] + + @property + def slices(self) -> list[str]: + i = self._line(SLICES_RE) + return LINK_RE.findall(self.body[i]) if i is not None else [] + + @property + def model(self) -> str: + i = self._line(MODEL_RE) + words = MODEL_RE.match(self.body[i]).group(1).split() if i is not None else [] + return words[0] if words else "opus" + + @property + def done_text(self) -> str | None: + """Text of the Done line, None without one.""" + i = self._line(DONE_RE) + return DONE_RE.match(self.body[i]).group(1) if i is not None else None + + @property + def runner_ready(self) -> bool: + """Headless worker can finish it: has a Done line, not owner-bound, Done not about the owner/confirming.""" + done = self.done_text + return bool(done and done.strip()) and self.sessions != "owner" and not HUMAN_DONE_RE.search(done) + + @property + def human_done_match(self) -> str | None: + """The Done phrase that makes it not runner-ready (human action), None if none.""" + m = HUMAN_DONE_RE.search(self.done_text or "") + return m.group(0) if m else None + + @property + def sessions(self) -> str: + """parallel | solo (no other live session) | owner (owner present); header ', interactive' = owner.""" + i = self._line(SESSIONS_RE) + words = SESSIONS_RE.match(self.body[i]).group(1).split() if i is not None else [] + if words and words[0] in SESSIONS: + return words[0] + return "owner" if self.interactive else "parallel" + + @property + def cloud(self) -> str | None: + """yes | no | None (no Cloud line, or an unknown word).""" + i = self._line(CLOUD_RE) + words = CLOUD_RE.match(self.body[i]).group(1).split() if i is not None else [] + return words[0] if words and words[0] in CLOUDS else None + + @property + def refs(self) -> list[tuple[str, str | None]]: + i = self._line(REF_RE) + if i is None: + return [] + out = [] + for part in split_refs(REF_RE.match(self.body[i]).group(1)): + target = re.sub(r"\s*\(.*$", "", part).strip() + if target: + path, _, anchor = target.partition("#") + out.append((path, anchor or None)) + return out + + +def split_refs(text: str) -> list[str]: + """Split on commas that are not inside parentheses.""" + parts, depth, cur = [], 0, "" + for ch in text: + if ch == "(": + depth += 1 + elif ch == ")": + depth = max(0, depth - 1) + if ch == "," and depth == 0: + parts.append(cur) + cur = "" + else: + cur += ch + return [p.strip() for p in parts + [cur] if p.strip()] + + +@dataclass +class Section: + heading: str + key: str | None + prefix: list[str] = field(default_factory=list) + items: list[Item] = field(default_factory=list) + suffix: list[str] = field(default_factory=list) + line: int | None = None # 1-based line of the heading + suffix_line: int | None = None + + def lines(self) -> list[str]: + out = [f"## {self.heading}", *self.prefix] + for item in self.items: + out += [*item.lines(), ""] + return out + self.suffix + + +@dataclass +class Doc: + preamble: list[str] + sections: list[Section] + newline: str = "\n" + + def section(self, key: str) -> Section: + for s in self.sections: + if s.key == key: + return s + raise TaskError(f"no '## {SECTIONS[key]}' section") + + def all_items(self) -> list[Item]: + return [i for s in self.sections for i in s.items] + + def ids(self) -> set[str]: + return {i.id for i in self.all_items() if i.id} + + def find(self, id: str) -> tuple[Section, int]: + for s in self.sections: + for n, item in enumerate(s.items): + if item.id == id: + return s, n + raise TaskError(unknown_id(id, self.ids())) + + def item(self, id: str) -> Item: + s, n = self.find(id) + return s.items[n] + + +def unknown_id(id: str, known: set[str]) -> str: + near = difflib.get_close_matches(id, sorted(known), n=3, cutoff=0.0) + return f"unknown id '{id}'" + (f" (nearest: {', '.join(near)})" if near else "") + + +def parse_header(line: str) -> Item: + m = HEADER_RE.match(line) + if not m: + found = re.match(r"^- \*\*([^*]+)\*\*", line) + return Item(id=found.group(1).strip() if found else "", raw=line, + error="bad header: want '- **id** [Pn] (effort) [(status)]: Title. Goal.'") + effort, interactive = m.group("effort"), False + if effort is not None: + parts = [p.strip() for p in effort.split(",")] + interactive = "interactive" in parts[1:] + effort = parts[0] + status = m.group("status") + return Item(id=m.group("id"), prio=int(m.group("prio")) if m.group("prio") else None, + effort=effort, interactive=interactive, + status=status.strip() if status else None, text=m.group("text").strip()) + + +def _strip_blank_tail(lines: list[str]) -> list[str]: + while lines and not lines[-1].strip(): + lines.pop() + return lines + + +def _parse_section(heading: str, lines: list[str], at: int = 0) -> Section: + """`at` = 1-based line number of the heading.""" + key = next((k for k, h in SECTIONS.items() if h == heading), None) + section = Section(heading, key, line=at) + if key is None: + section.prefix = lines + return section + for n, line in enumerate(lines, at + 1): + if section.suffix: + section.suffix.append(line) + elif line.startswith(ITEM_START): + section.items.append(parse_header(line)) + section.items[-1].line = n + elif not section.items: + section.prefix.append(line) + elif line.strip() and not line[0].isspace(): + section.suffix.append(line) + section.suffix_line = n + else: + section.items[-1].body.append(line) + for item in section.items: + _strip_blank_tail(item.body) + return section + + +def parse(text: str) -> Doc: + newline = "\r\n" if "\r\n" in text else "\n" + lines = _strip_blank_tail(text.replace("\r\n", "\n").split("\n")) + starts = [i for i, l in enumerate(lines) if l.startswith("## ")] + doc = Doc(lines[: starts[0]] if starts else lines, [], newline) + for n, start in enumerate(starts): + end = starts[n + 1] if n + 1 < len(starts) else len(lines) + doc.sections.append(_parse_section(lines[start][3:].strip(), lines[start + 1:end], start + 1)) + return doc + + +def render(doc: Doc) -> str: + lines = list(doc.preamble) + for s in doc.sections: + lines += s.lines() + _strip_blank_tail(lines) + return doc.newline.join(lines) + doc.newline + + +# ---------------------------------------------------------------- ids + +def make_id(title: str, taken: set[str], prefix: str = "t-") -> str: + words = [w for w in re.sub(r"[^a-z0-9]+", " ", title.lower().encode("ascii", "ignore").decode()).split() if w] + slug = "" + for w in words: + longer = f"{slug}-{w}" if slug else w + if len(longer) > 40: + break + slug = longer + if not slug and words: + slug = words[0][:40] + if not slug: + raise TaskError(f"no id can be made from title '{title}': give one with --id") + base, n, id = prefix + slug, 2, prefix + slug + while id in taken: + id = f"{base}-{n}" + n += 1 + return id + + +def slice_id(parent: str, taken: set[str]) -> str: + n = 1 + while f"{parent}-{n}" in taken: + n += 1 + return f"{parent}-{n}" + + +def open_slices(doc: Doc, id: str) -> list[str]: + """Open items listed in the Slices: line of `id`, then open `<id>-N` items.""" + pattern = re.compile(re.escape(id) + r"-\d+$") + open_ = doc.ids() + named = [s for s in doc.item(id).slices if s in open_] + return named + [i.id for i in doc.all_items() if pattern.match(i.id) and i.id not in named] + + +def parent_of(doc: Doc, id: str) -> str | None: + for item in doc.all_items(): + if id in item.slices: + return item.id + m = re.match(r"(.+)-\d+$", id) + return m.group(1) if m and m.group(1) in doc.ids() else None + + +def rename(doc: Doc, old: str, new: str, taken: set[str]) -> int: + """Change the id of an open item and every [[old]] link; returns the link count. + `taken` = ids used in the archive.""" + item = doc.item(old) + if not ID_RE.match(new): + raise TaskError(f"bad id '{new}' (want t-… or a-…, lowercase a-z 0-9 -)") + if new[:2] != old[:2]: + raise TaskError(f"a rename keeps the kind ({old[:2]}…)") + if new in doc.ids(): + raise TaskError(f"id '{new}' already exists") + if new in taken: + raise TaskError(f"id '{new}' already used in the archive (ids are never reused)") + link, n = re.compile(r"\[\[" + re.escape(old) + r"\]\]"), 0 + for other in doc.all_items(): + other.text, k = link.subn(f"[[{new}]]", other.text) + n += k + if other.status: + other.status, k = link.subn(f"[[{new}]]", other.status) + n += k + for i, line in enumerate(other.body): + other.body[i], k = link.subn(f"[[{new}]]", line) + n += k + item.id = new + return n + + +# ---------------------------------------------------------------- placing + +def _prio(item: Item) -> int: + return 9 if item.prio is None else item.prio + + +def place(section: Section, item: Item) -> int: + """Insert rule: after the last item of the same or higher priority, and + after every item named in the new item's After:.""" + at = max((n + 1 for n, other in enumerate(section.items) if _prio(other) <= _prio(item)), default=0) + deps = set(item.after) + return max([at, *(n + 1 for n, other in enumerate(section.items) if other.id in deps)]) + + +def _check_kind(item: Item, key: str) -> None: + if key not in SECTIONS: + raise TaskError(f"unknown section '{key}' (want {', '.join(SECTIONS)})") + if item.id.startswith("a-") != (key == "awaiting"): + raise TaskError("a- items belong in Awaiting, t- items in Pending / Needs human / Deferred") + if key != "awaiting" and item.prio is None: + raise TaskError(f"task '{item.id}' needs a priority [P0]-[P3]") + + +def insert(doc: Doc, item: Item, key: str) -> None: + if item.error: + raise TaskError(item.error) + if item.id in doc.ids(): + raise TaskError(f"id '{item.id}' already exists") + _check_kind(item, key) + section = doc.section(key) + section.items.insert(place(section, item), item) + + +def remove(doc: Doc, id: str) -> Item: + section, n = doc.find(id) + return section.items.pop(n) + + +def _section_key(doc: Doc, id: str) -> str: + return doc.find(id)[0].key + + +def set_prio(doc: Doc, id: str, prio: int) -> None: + if prio not in (0, 1, 2, 3): + raise TaskError("priority 0-3") + key = _section_key(doc, id) + item = doc.item(id) + if item.id.startswith("a-"): + raise TaskError("Awaiting items have no priority") + remove(doc, id) + item.prio = prio + insert(doc, item, key) + + +def move_to(doc: Doc, id: str, key: str) -> None: + item = doc.item(id) + _check_kind(item, key) + doc.section(key) + remove(doc, id) + insert(doc, item, key) + + +def move_rel(doc: Doc, id: str, other: str, before: bool, force: bool = False) -> None: + if id == other: + raise TaskError("cannot move an item relative to itself") + item = doc.item(id) + target, _ = doc.find(other) + _check_kind(item, target.key) + source, n = doc.find(id) + source.items.pop(n) + at = next(i for i, o in enumerate(target.items) if o.id == other) + (0 if before else 1) + target.items.insert(at, item) + + def undo(): + target.items.pop(at) + source.items.insert(n, item) + + later = {o.id for o in target.items[at + 1:]} + for dep in item.after: + if dep in later: + undo() + raise TaskError(f"'{id}' is After: [[{dep}]], it cannot go before it") + for o in target.items[:at]: + if id in o.after: + undo() + raise TaskError(f"'{o.id}' is After: [[{id}]], '{id}' cannot go after it") + prev = target.items[at - 1] if at else None + nxt = target.items[at + 1] if at + 1 < len(target.items) else None + if not force and ((prev and _prio(prev) > _prio(item)) or (nxt and _prio(nxt) < _prio(item))): + undo() + raise TaskError(f"moving '{id}' there breaks priority order (wf prio, or --force)") + + +# ---------------------------------------------------------------- fields + +def set_status(doc: Doc, id: str, status: str | None) -> None: + item = doc.item(id) + if item.id.startswith("a-"): + raise TaskError("Awaiting items have no status") + if status: + status = status.strip() + if status.startswith("blocked:"): + on = LINK_RE.findall(status) + open_awaiting = {i.id for s in doc.sections if s.key == "awaiting" for i in s.items} + if len(on) != 1 or on[0] not in open_awaiting: + raise TaskError(f"'{on[0] if on else status}' is not an open Awaiting item") + elif not status.startswith("in progress:"): + raise TaskError("status is 'in progress: <note>' or 'blocked: [[a-id]]'") + if ")" in status or "(" in status: + raise TaskError("status note cannot contain parentheses (the status line ends in one); use - or :") + item.status = status or None + + +def _tail_start(item: Item) -> int: + """Index where the trailing Sessions:/Model:/After:/Ref: lines of the body begin.""" + n = len(item.body) + while n and any(p.match(item.body[n - 1]) for p in (AFTER_RE, REF_RE, MODEL_RE, SESSIONS_RE, CLOUD_RE)): + n -= 1 + return n + + +def _set_line(item: Item, pattern: re.Pattern, line: str | None, at_end: bool) -> None: + i = item._line(pattern) + if i is not None: + if line is None: + item.body.pop(i) + else: + item.body[i] = line + elif line is not None: + if at_end: + item.body.append(line) + else: + ref = item._line(REF_RE) + item.body.insert(len(item.body) if ref is None else ref, line) + + +def set_fields(doc: Doc, id: str, *, title: str | None = None, effort: str | None = None, + interactive: bool | None = None, after: list[str] | None = None, + refs: list[str] | None = None, model: str | None = None, + sessions: str | None = None, done: str | None = None, + cloud: str | None = None) -> None: + item = doc.item(id) + if done is not None: + done = done.strip() + if (i := item._line(DONE_RE)) is not None: + if done: + item.body[i] = INDENT + "Done: " + done + else: + item.body.pop(i) + elif done: + item.body.insert(_tail_start(item), INDENT + "Done: " + done) + if sessions is not None: + if sessions not in ("", *SESSIONS): + raise TaskError(f"sessions '{sessions}' (want {', '.join(SESSIONS)})") + item.interactive = False + line = INDENT + "Sessions: " + sessions if sessions else None + if (i := item._line(SESSIONS_RE)) is not None: + if line: + item.body[i] = line + else: + item.body.pop(i) + elif line: + m = item._line(MODEL_RE) + item.body.insert(m if m is not None and m >= _tail_start(item) else _tail_start(item), line) + if cloud is not None: + if cloud not in ("", *CLOUDS): + raise TaskError(f"cloud '{cloud}' (want {', '.join(CLOUDS)})") + line = INDENT + "Cloud: " + cloud if cloud else None + if (i := item._line(CLOUD_RE)) is not None: + if line: + item.body[i] = line + else: + item.body.pop(i) + elif line: + item.body.insert(_tail_start(item), line) + if model: + if model not in MODELS: + raise TaskError(f"model '{model}' (want {', '.join(MODELS)})") + line = INDENT + "Model: " + model + if (i := item._line(MODEL_RE)) is not None: + item.body[i] = line + else: + item.body.insert(_tail_start(item), line) + elif model is not None and (i := item._line(MODEL_RE)) is not None: + item.body.pop(i) + if title is not None: + title = title.strip().rstrip(".") + if not title: + raise TaskError("empty title") + if ". " in title or title.endswith(("?", "!")): + item.text = title if title.endswith(("?", "!")) else f"{title}." + else: + item.text = f"{title}. {item.goal}" if item.goal else f"{title}." + if effort is not None: + if effort not in EFFORTS: + raise TaskError(f"effort '{effort}' (want {', '.join(EFFORTS)})") + item.effort = effort + if interactive is not None: + item.interactive = interactive + if after is not None: + line = INDENT + "- After: " + ", ".join(f"[[{a}]]" for a in after) if after else None + _set_line(item, AFTER_RE, line, at_end=False) + if refs is not None: + _set_line(item, REF_RE, INDENT + "Ref: " + ", ".join(refs) if refs else None, at_end=True) + + +def add_note(doc: Doc, id: str, line: str) -> None: + item = doc.item(id) + line = line.strip() + if not line: + raise TaskError("empty note") + item.body.insert(_tail_start(item), f"{INDENT}- {line}") + + +def _link_after(line: str) -> str: + """'After: t-a, t-b' -> 'After: [[t-a]], [[t-b]]' (bare ids accepted, written as links).""" + m = re.match(r"^(\s*(?:-\s+)?After:\s*)(.*)$", line) + if not m: + return line + parts = [w for w in re.split(r"[\s,;]+", m.group(2)) if w] + if not parts or not all(LINK_RE.fullmatch(w) or ID_RE.match(w) for w in parts): + return line + return m.group(1) + ", ".join(w if w.startswith("[[") else f"[[{w}]]" for w in parts) + + +def indent_body(lines: list[str]) -> list[str]: + import textwrap + text = textwrap.dedent("\n".join(_link_after(l.rstrip()) for l in lines)) + return _strip_blank_tail([INDENT + l if l.strip() else "" for l in text.split("\n")]) if text.strip() else [] + + +def set_body(doc: Doc, id: str, lines: list[str]) -> None: + item = doc.item(id) + new = indent_body(lines) + kept = [l for l in item.body[_tail_start(item):] + if not any(p.match(l) and any(p.match(n) for n in new) for p in (AFTER_RE, REF_RE))] + item.body = new + kept + + +BOX_RE = re.compile(r"^(\s*- \[)([ xX])(\] )(.*)$") + + +def tick(doc: Doc, id: str, which: str) -> str: + item = doc.item(id) + boxes = [(i, BOX_RE.match(l)) for i, l in enumerate(item.body) if BOX_RE.match(l)] + if which.isdigit(): + if not 1 <= int(which) <= len(boxes): + raise TaskError(f"no box {which} in '{id}' ({len(boxes)} boxes)") + hits = [boxes[int(which) - 1]] + else: + hits = [(i, m) for i, m in boxes if which.lower() in m.group(4).lower()] + if not hits: + raise TaskError(f"no box matches '{which}' in '{id}'") + if len(hits) > 1: + raise TaskError(f"'{which}' matches {len(hits)} boxes in '{id}'") + i, m = hits[0] + if m.group(2) != " ": + raise TaskError(f"box already ticked: {m.group(4)}") + item.body[i] = f"{m.group(1)}x{m.group(3)}{m.group(4)}" + return m.group(4) + + +def unblock(doc: Doc, a_id: str) -> list[str]: + out = [] + for item in doc.all_items(): + if item.blocked_on == a_id: + item.status = None + out.append(item.id) + return out + + +# ---------------------------------------------------------------- pick + +def why_not(doc: Doc, item: Item, archived: set[str], held: dict[str, str] | None = None, + others: int = 0, owner: bool = True) -> str | None: + """Why a Pending item can't be picked now; None = pickable. held: id → holder of a live + claim by another session; others: other live sessions in the project; owner: owner present.""" + if item.error: + return "bad header" + if held and item.id in held: + return f"in progress by {held[item.id]}" + if item.sessions == "owner" and not owner: + return "owner: needs the owner (wf next --owner)" + if item.sessions == "solo" and others: + return f"solo: {others} other live session{'s' if others > 1 else ''}" + if item.status and item.status.startswith("blocked:"): + return f"blocked: {item.blocked_on}" + if open_slices(doc, item.id): + return "open slices: " + ", ".join(open_slices(doc, item.id)) + if [a for a in item.after if a not in archived]: + return "after: " + ", ".join(a for a in item.after if a not in archived) + return None + + +def waiters(doc: Doc) -> dict[str, set[str]]: + """Open task id → every open task (Pending, Needs human) that waits on it, transitively: + via After:, and a parent waits on its open slices.""" + items = {i.id: i for s in doc.sections if s.key in ("pending", "human") for i in s.items if not i.error} + direct: dict[str, set[str]] = {id: set() for id in items} + for item in items.values(): + for dep in item.after: + if dep in direct: + direct[dep].add(item.id) + for s in item.slices: + if s in direct: + direct[s].add(item.id) + out: dict[str, set[str]] = {} + for id in direct: + seen, todo = set(), list(direct[id]) + while todo: + w = todo.pop() + if w not in seen and w != id: + seen.add(w) + todo += direct[w] + out[id] = seen + return out + + +def solo_running(doc: Doc, held: dict[str, str]) -> tuple[str, str] | None: + """(id, holder) of a solo task in progress by another live session; held: see why_not.""" + return next(((i.id, held[i.id]) for i in doc.section("pending").items + if i.id in held and not i.error and i.sessions == "solo"), None) + + +# ---------------------------------------------------------------- archive + +ARCHIVE_ID_RE = re.compile(r"^- .*?\*\*([ta]-[a-z0-9-]+)\*\*") + + +def archive_ids(text: str) -> set[str]: + return {m.group(1) for l in text.replace("\r\n", "\n").split("\n") if (m := ARCHIVE_ID_RE.match(l))} + + +def archive_line(date: str, item: Item, entry: str) -> str: + entry = " ".join(entry.split()) + return f"- {date} **{item.id}** {item.title}" + (f" — {entry}" if entry else "") + + +def archive_prepend(text: str, line: str) -> str: + newline = "\r\n" if "\r\n" in text else "\n" + lines = _strip_blank_tail(text.replace("\r\n", "\n").split("\n")) + at = next((i for i, l in enumerate(lines) if l.startswith("- ")), None) + if at is None: + lines += ["", line] + else: + lines.insert(at, line) + return newline.join(lines) + newline + + +# ---------------------------------------------------------------- blocks + +def parse_block(text: str) -> Item: + lines = _strip_blank_tail(text.replace("\r\n", "\n").strip("\n").split("\n")) + if not lines or not lines[0].strip(): + raise TaskError("empty item") + first = lines[0].strip() + if re.match(r"^\d+\. ", first): + raise TaskError("old numbered format: write '- **id** [Pn] (effort): Title. Goal.'") + item = parse_header(first) + if item.error: + raise TaskError(item.error) + item.body = indent_body(lines[1:]) + return item diff --git a/wflib/usage.py b/wflib/usage.py new file mode 100644 index 0000000..5b7c1ee --- /dev/null +++ b/wflib/usage.py @@ -0,0 +1,375 @@ +"""Token usage and API-price cost from Claude Code transcripts (jsonl), per model. + +One API request is streamed as several transcript entries (one per content block) that repeat its +usage: count each (message.id, requestId) once, taking its last entry (final output_tokens). +Subagent transcripts never get a final entry (stop_reason null, output_tokens = the message_start +placeholder): there output is estimated from the content (out_estimate), counted in Usage.est. +""" +from __future__ import annotations + +import json +import re +from dataclasses import dataclass, fields +from statistics import median + +# $ per million tokens: input, output, cache read. Cache writes: 1.25x input (5 min), 2x input (1 h). +# Source: Anthropic pricing as of 2026-09-25. Longest matching prefix of the model id wins. +PRICES = { + "claude-fable-5-1": (10.0, 50.0, 0.25), + "claude-fable-5": (10.0, 50.0, 1.0), + "claude-opus-5-5": (4.0, 20.0, 0.20), + "claude-opus-5": (5.0, 25.0, 0.50), + "claude-opus-4-8": (5.0, 25.0, 0.50), + "claude-sonnet-5-5": (2.0, 10.0, 0.20), + "claude-sonnet-5": (2.0, 10.0, 0.20), + "claude-sonnet-4-6": (3.0, 15.0, 0.30), + "claude-haiku-4-5": (1.0, 5.0, 0.10), +} + + +@dataclass +class Usage: + turns: int = 0 + inp: int = 0 + cw5: int = 0 + cw1h: int = 0 + cr: int = 0 + out: int = 0 + est: int = 0 # requests whose out is estimated from content + + @property + def cw(self) -> int: + return self.cw5 + self.cw1h + + def add(self, other: "Usage") -> None: + for f in fields(self): + setattr(self, f.name, getattr(self, f.name) + getattr(other, f.name)) + + +def price(model: str) -> tuple[float, float, float] | None: + keys = [k for k in PRICES if model == k or model.startswith(k + "-")] + return PRICES[max(keys, key=len)] if keys else None + + +def cost(model: str, u: Usage) -> float | None: + """API-price $ (subscription sessions: what it would cost on the API); None = unknown model.""" + p = price(model) + if p is None: + return None + inp, out, read = p + return (u.inp * inp + u.cw5 * inp * 1.25 + u.cw1h * inp * 2 + u.cr * read + u.out * out) / 1e6 + + +def _usage(raw: dict) -> Usage: + cw = raw.get("cache_creation_input_tokens") or 0 + cw1h = (raw.get("cache_creation") or {}).get("ephemeral_1h_input_tokens") or 0 + return Usage(1, raw.get("input_tokens") or 0, cw - cw1h, cw1h, + raw.get("cache_read_input_tokens") or 0, raw.get("output_tokens") or 0) + + +# Output-token estimate per content block, fitted on ~100k main-session requests with final usage +# (2026-10): per-session totals within ±5%. Thinking text is usually empty; its signature grows ~3.3 +# chars per thinking token over a ~800-char base. +EST_TEXT, EST_TOOL_CHAR, EST_TOOL, EST_SIG, EST_SIG_BASE = 0.3, 0.44, 30, 0.3, 800 + + +def out_estimate(blocks: list) -> int: + t = 0.0 + for b in blocks: + if not isinstance(b, dict): + continue + kind = b.get("type") + if kind == "text": + t += len(b.get("text") or "") * EST_TEXT + elif kind == "tool_use": + t += EST_TOOL + len(json.dumps(b.get("input") or {}, ensure_ascii=False)) * EST_TOOL_CHAR + elif kind == "thinking": + t += (len(b.get("thinking") or "") * EST_TEXT + + max(0, len(b.get("signature") or "") - EST_SIG_BASE) * EST_SIG) + return round(t) + + +def parse(text: str, since: str | None = None, until: str | None = None) -> dict[str, Usage]: + """Model id → summed Usage. since/until: ISO prefixes compared with the UTC timestamps.""" + last: dict[tuple, tuple[str, dict]] = {} + first_ts: dict[tuple, str] = {} + final: set[tuple] = set() + blocks: dict[tuple, list] = {} + seen: set[str] = set() + for line in text.split("\n"): + try: + d = json.loads(line) + except ValueError: + continue + if not isinstance(d, dict) or d.get("type") != "assistant": + continue + m = d.get("message") or {} + model, raw = m.get("model") or "?", m.get("usage") + if not raw or model == "<synthetic>": + continue + key = (m.get("id"), d.get("requestId")) + first_ts.setdefault(key, d.get("timestamp") or "") + last[key] = (model, raw) + if m.get("stop_reason"): + final.add(key) + uid = d.get("uuid") + if uid not in seen and isinstance(m.get("content"), list): + blocks.setdefault(key, []).extend(m["content"]) + if uid: + seen.add(uid) + out: dict[str, Usage] = {} + for key, (model, raw) in last.items(): + ts = first_ts[key] + if (since and ts < since) or (until and ts >= until): + continue + u = _usage(raw) + if key not in final: + guess = out_estimate(blocks.get(key, [])) + if guess > u.out: + u.out, u.est = guess, 1 + out.setdefault(model, Usage()).add(u) + return out + + +def context_tokens(text: str) -> int | None: + """Prompt size of the last main-thread request (input + cache write + cache read); None = no request.""" + size = None + for line in text.split("\n"): + try: + d = json.loads(line) + except ValueError: + continue + if not isinstance(d, dict) or d.get("type") != "assistant" or d.get("isSidechain"): + continue + m = d.get("message") or {} + raw = m.get("usage") + if raw and m.get("model") != "<synthetic>": + u = _usage(raw) + size = u.inp + u.cw + u.cr + return size + + +def ctx_hint(tokens: int | None, limit: int) -> str | None: + if not limit or tokens is None or tokens <= limit: + return None + return (f"context ~{tokens // 1000}k tokens (> {limit // 1000}k): ask the owner to /clear, then continue " + "(subagent: ignore, this is the main session)") + + +def short(model: str) -> str: + return re.sub(r"-\d{8}$", "", model.removeprefix("claude-")) + + +def fmt(n: int) -> str: + if n < 1000: + return str(n) + return f"{n / 1e3:.1f}k" if n < 1e6 else f"{n / 1e6:.2f}M" + + +# ------------------------------------------------------------------ cost log (out/wf-cost.log) + +def log_line(time: str, project: str, task: str, effort: str, outcome: str, agent: str, + by_model: dict[str, Usage], lane: str | None = None, dur: int | None = None) -> str: + """One key=value line per worker; model = the model with most turns, lane = given or model; usd = priced models only.""" + total = Usage() + for u in by_model.values(): + total.add(u) + model = short(max(by_model, key=lambda m: by_model[m].turns)).split("-")[0] if by_model else "?" + usd = sum(cost(m, u) or 0 for m, u in by_model.items()) + return (f"{time} project={project} task={task} lane={lane or model} model={model} effort={effort} outcome={outcome} " + f"turns={total.turns} in={total.inp} cw={total.cw} cr={total.cr} out={total.out} " + + (f"est={total.est} " if total.est else "") + f"usd={usd:.4f} " + (f"dur={dur} " if dur is not None else "") + f"agent={agent}") + + +def parse_log(text: str, since: str | None = None) -> list[dict]: + out = [] + for line in text.split("\n"): + time, _, rest = line.partition(" ") + fields_ = dict(f.split("=", 1) for f in rest.split() if "=" in f) + if not re.match(r"\d{4}-\d\d-\d\dT", time) or "lane" not in fields_ or (since and time < since): + continue + try: + out.append({**fields_, "time": time, "usd": float(fields_.get("usd", 0)), + "turns": int(fields_.get("turns", 0)), + **({"dur": int(fields_["dur"])} if "dur" in fields_ else {})}) + except ValueError: + continue + return out + + +def _num(x: float) -> float | int: + x = round(x, 4) + return int(x) if x == int(x) and not isinstance(x, bool) else x + + +def report(entries: list[dict], efforts: tuple[str, ...]) -> list[tuple]: + """Per lane, then per lane+effort: (lane, effort, n, done, other, median $, total $, $/done, median turns, median dur s or None).""" + def row(lane, effort, es): + done = sum(e.get("outcome") in ("done", "done+gate-red") for e in es) + total = sum(e["usd"] for e in es) + return (lane, effort, len(es), done, len(es) - done, round(median(e["usd"] for e in es), 4), + round(total, 4), round(total / done, 4) if done else None, _num(median(e["turns"] for e in es)), + _num(median(durs)) if (durs := [e["dur"] for e in es if "dur" in e]) else None) + + order = {e: i for i, e in enumerate(efforts)} + rows = [] + key = lambda e: f"{e['lane']}/{e['model']}" if "model" in e and e["model"] != e["lane"] else e["lane"] + for lane in sorted({key(e) for e in entries}): + es = [e for e in entries if key(e) == lane] + rows.append(row(lane, "all", es)) + for effort in sorted({e.get("effort", "-") for e in es}, key=lambda x: (order.get(x, len(order)), x)): + rows.append(row(lane, effort, [e for e in es if e.get("effort", "-") == effort])) + return rows + + +# --- explore: where an agent's billed input goes (exploring vs editing/running/bookkeeping) --------------------- +_MUT = re.compile(r"sed -i|python3 - <<|cat > |tee |expect-init|>> [a-zA-Z]|apply_patch|patch ") +_RUN = re.compile(r"dotnet (test|build|run)|own-fx\.py (check|measure|dump)|npx |npm |gate-bg|check\.sh|wf res run|pytest|unittest") +_BOOK = re.compile(r"wf\.py (done|finish|add|status|note|set|merge|check|tick|body)" + r"|git (commit|add|switch|rebase|merge|push)|git worktree (add|remove)|git branch -[dDmM]") +_FILE = re.compile(r"(?<![\w.-])((?:\./|/)?[\w.-]+(?:/[\w.-]+)+\.\w+)(?![\w/-])") +_RESULT_KINDS = ( # first match wins; Bash result kinds + ("git-hist", re.compile(r"git (show|log|diff)")), + ("wf", re.compile(r"wf\.py")), + ("build/test", re.compile(r"dotnet|npx|npm")), + ("grep", re.compile(r"\bgrep|\brg\b|find ")), + ("sed-cat", re.compile(r"sed -n|cat |head|tail")), +) + + +def call_kind(tool: dict) -> str: + """mut (edit/write) | book (wf, git writes) | run (build/test) | exp (everything else: exploring).""" + if tool.get("name") in ("Edit", "Write", "NotebookEdit"): + return "mut" + if tool.get("name") != "Bash": + return "exp" + c = (tool.get("input") or {}).get("command") or "" + return "book" if _BOOK.search(c) else "run" if _RUN.search(c) else "mut" if _MUT.search(c) else "exp" + + +def result_kind(tool: dict) -> str: + if tool.get("name") == "Bash": + c = (tool.get("input") or {}).get("command") or "" + return next((k for k, rx in _RESULT_KINDS if rx.search(c)), "other") + return tool.get("name") or "other" + + +def tool_files(tool: dict, root: str = "") -> list[str]: + """Files a Read / Bash call touches, relative to the project (.worktrees/<x>/ and root stripped).""" + inp = tool.get("input") or {} + if tool.get("name") == "Read": + paths = [inp.get("file_path") or ""] + elif tool.get("name") == "Bash": + paths = [p for p in _FILE.findall(inp.get("command") or "") if not p.startswith("/") or p.startswith(root + "/")] + else: + return [] + out = [] + for p in paths: + p = re.sub(r"^.*?/\.worktrees/[^/]+/", "", p) + if root and p.startswith(root.rstrip("/") + "/"): + p = p[len(root.rstrip("/")) + 1:] + if p and p not in out: + out.append(p) + return out + + +def _ctx(u: dict) -> int: + return (u.get("input_tokens") or 0) + (u.get("cache_read_input_tokens") or 0) + (u.get("cache_creation_input_tokens") or 0) + + +@dataclass +class Explore: + calls: int = 0 + first_edit: int | None = None # calls before the first edit; None = never edited + cost: dict = None # kind → billed input tokens + pre: int = 0 # billed input before the first edit + ctx_growth: int = 0 # context growth up to the first edit + results: dict = None # result kind → tokens + files: dict = None # file → result tokens (read by this agent) + + @property + def total(self) -> int: + return sum(self.cost.values()) + + @property + def edited(self) -> bool: + return self.first_edit is not None + + +def explore(text: str, since: str | None = None, root: str = "") -> Explore: + """One transcript → billed-input split by call kind, result tokens per tool kind, files read.""" + calls: dict[str, dict] = {} + for line in text.split("\n"): + try: + d = json.loads(line) + except ValueError: + continue + if not isinstance(d, dict) or d.get("type") != "assistant": + continue + m = d.get("message") or {} + key = m.get("id") or d.get("uuid") + c = calls.setdefault(key, {"u": {}, "t": [], "ts": d.get("timestamp") or ""}) + c["u"] = m.get("usage") or c["u"] + c["t"] += [t for t in m.get("content") or [] if isinstance(t, dict) and t.get("type") == "tool_use"] + kept = [c for c in calls.values() if not since or c["ts"] >= since] + tools = {t["id"]: t for c in kept for t in c["t"] if "id" in t} + ex = Explore(calls=len(kept), cost={}, results={}, files={}) + for i, c in enumerate(kept): + ks = [call_kind(t) for t in c["t"]] or ["exp"] + k = "mut" if "mut" in ks else ks[0] + if k == "mut" and ex.first_edit is None: + ex.first_edit = i + ex.pre = sum(_ctx(x["u"]) for x in kept[:i]) + ex.ctx_growth = _ctx(c["u"]) - _ctx(kept[0]["u"]) + ex.cost[k] = ex.cost.get(k, 0) + _ctx(c["u"]) + for line in text.split("\n"): + try: + d = json.loads(line) + except ValueError: + continue + content = (d.get("message") or {}).get("content") if isinstance(d, dict) and d.get("type") == "user" else None + for c in content if isinstance(content, list) else []: + if not isinstance(c, dict) or c.get("type") != "tool_result" or c.get("tool_use_id") not in tools: + continue + body = c.get("content") + size = len(body if isinstance(body, str) else json.dumps(body)) // 4 + t = tools[c["tool_use_id"]] + k = result_kind(t) + ex.results[k] = ex.results.get(k, 0) + size + for f in tool_files(t, root): + ex.files[f] = ex.files.get(f, 0) + size + return ex + + +def explore_report(agents: list[tuple[str, Explore]], min_calls: int = 4, min_pre: int = 10, top: int = 20) -> dict: + """agents: (label, Explore). Agents under min_calls are skipped; the pre-edit median uses agents with an + edit and >= min_pre calls. Returns rows, totals, results, files (≥2 agents: (file, agents, tokens)).""" + use = [(l, e) for l, e in agents if e.calls >= min_calls] + rows = [] + for label, e in use: + t = e.total or 1 + rows.append((label, e.calls, e.first_edit if e.edited else e.calls, e.cost.get("exp", 0) / t, e.pre / t)) + cost: dict[str, int] = {} + results: dict[str, int] = {} + for _, e in use: + for k, v in e.cost.items(): + cost[k] = cost.get(k, 0) + v + for k, v in e.results.items(): + results[k] = results.get(k, 0) + v + pre = [e for _, e in use if e.edited and e.calls >= min_pre] + seen: dict[str, list[int]] = {} + for _, e in use: + for f, n in e.files.items(): + s = seen.setdefault(f, [0, 0]) + s[0] += 1 + s[1] += n + files = sorted(((f, a, n) for f, (a, n) in seen.items() if a >= 2), key=lambda x: (-x[1], -x[2], x[0]))[:top] + total = sum(cost.values()) + return { + "rows": rows, "agents": len(use), "total": total, + "cost": cost, "results": results, "files": files, + "pre_n": len(pre), + "pre_share": median(e.pre / (e.total or 1) for e in pre) if pre else None, + "pre_calls": median(e.first_edit for e in pre) if pre else None, + "pre_ctx": median(e.ctx_growth for e in pre) if pre else None, + } -- cgit