raw · 119281 bytes
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 915 916 917 918 919 920 921 922 923 924 925 926 927 928 929 930 931 932 933 934 935 936 937 938 939 940 941 942 943 944 945 946 947 948 949 950 951 952 953 954 955 956 957 958 959 960 961 962 963 964 965 966 967 968 969 970 971 972 973 974 975 976 977 978 979 980 981 982 983 984 985 986 987 988 989 990 991 992 993 994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 1026 1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075 1076 1077 1078 1079 1080 1081 1082 1083 1084 1085 1086 1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146 1147 1148 1149 1150 1151 1152 1153 1154 1155 1156 1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188 1189 1190 1191 1192 1193 1194 1195 1196 1197 1198 1199 1200 1201 1202 1203 1204 1205 1206 1207 1208 1209 1210 1211 1212 1213 1214 1215 1216 1217 1218 1219 1220 1221 1222 1223 1224 1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237 1238 1239 1240 1241 1242 1243 1244 1245 1246 1247 1248 1249 1250 1251 1252 1253 1254 1255 1256 1257 1258 1259 1260 1261 1262 1263 1264 1265 1266 1267 1268 1269 1270 1271 1272 1273 1274 1275 1276 1277 1278 1279 1280 1281 1282 1283 1284 1285 1286 1287 1288 1289 1290 1291 1292 1293 1294 1295 1296 1297 1298 1299 1300 1301 1302 1303 1304 1305 1306 1307 1308 1309 1310 1311 1312 1313 1314 1315 1316 1317 1318 1319 1320 1321 1322 1323 1324 1325 1326 1327 1328 1329 1330 1331 1332 1333 1334 1335 1336 1337 1338 1339 1340 1341 1342 1343 1344 1345 1346 1347 1348 1349 1350 1351 1352 1353 1354 1355 1356 1357 1358 1359 1360 1361 1362 1363 1364 1365 1366 1367 1368 1369 1370 1371 1372 1373 1374 1375 1376 1377 1378 1379 1380 1381 1382 1383 1384 1385 1386 1387 1388 1389 1390 1391 1392 1393 1394 1395 1396 1397 1398 1399 1400 1401 1402 1403 1404 1405 1406 1407 1408 1409 1410 1411 1412 1413 1414 1415 1416 1417 1418 1419 1420 1421 1422 1423 1424 1425 1426 1427 1428 1429 1430 1431 1432 1433 1434 1435 1436 1437 1438 1439 1440 1441 1442 1443 1444 1445 1446 1447 1448 1449 1450 1451 1452 1453 1454 1455 1456 1457 1458 1459 1460 1461 1462 1463 1464 1465 1466 1467 1468 1469 1470 1471 1472 1473 1474 1475 1476 1477 1478 1479 1480 1481 1482 1483 1484 1485 1486 1487 1488 1489 1490 1491 1492 1493 1494 1495 1496 1497 1498 1499 1500 1501 1502 1503 1504 1505 1506 1507 1508 1509 1510 1511 1512 1513 1514 1515 1516 1517 1518 1519 1520 1521 1522 1523 1524 1525 1526 1527 1528 1529 1530 1531 1532 1533 1534 1535 1536 1537 1538 1539 1540 1541 1542 1543 1544 1545 1546 1547 1548 1549 1550 1551 1552 1553 1554 1555 1556 1557 1558 1559 1560 1561 1562 1563 1564 1565 1566 1567 1568 1569 1570 1571 1572 1573 1574 1575 1576 1577 1578 1579 1580 1581 1582 1583 1584 1585 1586 1587 1588 1589 1590 1591 1592 1593 1594 1595 1596 1597 1598 1599 1600 1601 1602 1603 1604 1605 1606 1607 1608 1609 1610 1611 1612 1613 1614 1615 1616 1617 1618 1619 1620 1621 1622 1623 1624 1625 1626 1627 1628 1629 1630 1631 1632 1633 1634 1635 1636 1637 1638 1639 1640 1641 1642 1643 1644 1645 1646 1647 1648 1649 1650 1651 1652 1653 1654 1655 1656 1657 1658 1659 1660 1661 1662 1663 1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678 1679 1680 1681 1682 1683 1684 1685 1686 1687 1688 1689 1690 1691 1692 1693 1694 1695 1696 1697 1698 1699 1700 1701 1702 1703 1704 1705 1706 1707 1708 1709 1710 1711 1712 1713 1714 1715 1716 1717 1718 1719 1720 1721 1722 1723 1724 1725 1726 1727 1728 1729 1730 1731 1732 1733 1734 1735 1736 1737 1738 1739 1740 1741 1742 1743 1744 1745 1746 1747 1748 1749 1750 1751 1752 1753 1754 1755 1756 1757 1758 1759 1760 1761 1762 1763 1764 1765 1766 1767 1768 1769 1770 1771 1772 1773 1774 1775 1776 1777 1778 1779 1780 1781 1782 1783 1784 1785 1786 1787 1788 1789 1790 1791 1792 1793 1794 1795 1796 1797 1798 1799 1800 1801 1802 1803 1804 1805 1806 1807 1808 1809 1810 1811 1812 1813 1814 1815 1816 1817 1818 1819 1820 1821 1822 1823 1824 1825 1826 1827 1828 1829 1830 1831 1832 1833 1834 1835 1836 1837 1838 1839 1840 1841 1842 1843 1844 1845 1846 1847 1848 1849 1850 1851 1852 1853 1854 1855 1856 1857 1858 1859 1860 1861 1862 1863 1864 1865 1866 1867 1868 1869 1870 1871 1872 1873 1874 1875 1876 1877 1878 1879 1880 1881 1882 1883 1884 1885 1886 1887 1888 1889 1890 1891 1892 1893 1894 1895 1896 1897 1898 1899 1900 1901 1902 1903 1904 1905 1906 1907 1908 1909 1910 1911 1912 1913 1914 1915 1916 1917 1918 1919 1920 1921 1922 1923 1924 1925 1926 1927 1928 1929 1930 1931 1932 1933 1934 1935 1936 1937 1938 1939 1940 1941 1942 1943 1944 1945 1946 1947 1948 1949 1950 1951 1952 1953 1954 1955 1956 1957 1958 1959 1960 1961 1962 1963 1964 1965 1966 1967 1968 1969 1970 1971 1972 1973 1974 1975 1976 1977 1978 1979 1980 1981 1982 1983 1984 1985 1986 1987 1988 1989 1990 1991 1992 1993 1994 1995 1996 1997 1998 1999 2000 2001 2002 2003 2004 2005 2006 2007 2008 2009 2010 2011 2012 2013 2014 2015 2016 2017 2018 2019 2020 2021 2022 2023 2024 2025 2026 2027 2028 2029 2030 2031 2032 2033 2034 2035 2036 2037 2038 2039 2040 2041 2042 2043 2044 2045 2046 2047 2048 2049 2050 2051 2052 2053 2054 2055 2056 2057 2058 2059 2060 2061 2062 2063 2064 2065 2066 2067 2068 2069 2070 2071 2072 2073 2074 2075 2076 2077 2078 2079 2080 2081 2082 2083 2084 2085 2086 2087 2088 2089 2090 2091 2092 2093 2094 2095 2096 2097 2098 2099 2100 2101 2102 2103 2104 2105 2106 2107 2108 2109 2110 2111 2112 2113 2114 2115 2116 2117 2118 2119 2120 2121 2122 2123 2124 2125 2126 2127 2128 2129 2130 2131 2132 2133 2134 2135 2136 2137 2138 2139 2140 2141 2142 2143 2144 2145 2146 2147 2148 2149 2150 2151 2152 2153 2154 2155 2156 2157 2158 2159 2160 2161 2162 2163 2164 2165 2166 2167 2168 2169 2170 2171 2172 2173 2174 2175 2176 2177 2178 2179 2180 2181 2182 2183 2184 2185 2186 2187 2188 2189 2190 2191 2192 2193 2194 2195 2196 2197 2198 2199 2200 2201 2202 2203 2204 2205 2206 2207 2208 2209 2210 2211 2212 2213 2214 2215 2216 2217 2218 2219 2220 2221 2222 2223 2224 2225 2226 2227 2228 2229 2230 2231 2232 2233 2234 2235 2236 2237 2238 2239 2240 2241 2242 2243 2244 2245 2246 2247 2248 2249 2250 2251 2252 2253 2254 2255 2256 2257 2258 2259 2260 2261 2262 2263 2264 2265 2266 2267 2268 2269 2270 2271 2272 2273 2274 2275 2276 2277 2278 2279 2280 2281 2282 2283 2284 2285 2286 2287 2288 2289 2290 2291 2292 2293 2294 2295 2296 2297 2298 2299 2300 2301 2302 2303 2304 2305 2306 2307 2308 2309 2310 2311 2312 2313 2314 2315 2316 2317 2318 2319 2320 2321 2322 2323 2324 2325 2326 2327 2328 2329 2330 2331 2332 2333 2334 2335 2336 2337 2338 2339 2340 2341 2342 2343 2344 2345 2346 2347 2348 2349 2350 2351 2352 2353 2354 2355 2356 2357 2358 2359 2360 2361 2362 2363 2364 2365 2366 2367 2368 2369 2370 2371 2372 2373 2374 2375 2376 2377 2378 2379 2380 2381 2382 2383 2384 2385 2386 2387 2388 2389 2390 2391 2392 2393 2394 2395 2396 2397 2398 2399 2400 2401 2402 2403 2404 2405 2406 2407 2408 2409 2410 2411 2412 2413 2414 2415 2416 2417 2418 2419 2420 2421 2422 2423 2424 2425 2426 2427 2428 2429 2430 2431 2432 2433 2434 2435 2436 2437 2438 2439 2440 2441 | #!/usr/bin/env python3 """wf — shared task workflow tool. Run inside a project (folder with workflow.toml). Look: next [--brief] · list [filters] · show ID… · ctx ID|PATH#ANCHOR search WORDS… · log [-n N] [WORDS] · projects · check Change: add "Title. Goal." -p N -e EFFORT · done ID… -m "entry" · prio ID N move ID SECTION | --before/--after ID · status ID progress NOTE|blocked A-ID|clear set ID --title/--effort/--after/--ref/--model/--sessions/--cloud · note ID "line" body ID <stdin (Model:/Sessions:/After:/Ref: kept) · tick ID N|TEXT · setup · merge (in a lane worktree) push (rerun a merge's push: workflow.toml push, else push home) Orch: orch pick LANE [--id ID] [--recovery WHY] · orch post ID LANE [--result LINE] [--agent A] [--duration S] [--no-next] Other: report "what happened" · init · migrate [--write] · res … (shared memory/CPU ledger; wf res -h) batch N [--lanes L,…] [--prep] | --status (unattended batch orchestrator; wf batch -h) prep [N|all] [--lanes L,…] (= batch --prep: sonnet workers write Done for tasks without one; all = every such task) Every change command takes --dry-run. `wf CMD -h` for details. Rules: /projects/CLAUDE.md. Design: docs/design.md. """ from __future__ import annotations import argparse import contextlib import datetime import difflib import glob import io import fcntl import os import re import subprocess import sys import tempfile from pathlib import Path HERE = Path(__file__).resolve().parent sys.path.insert(0, str(HERE)) from wflib import check as checks # noqa: E402 from wflib import areas as areas_mod # noqa: E402 from wflib import config, lanes, ledgers, refs, search, tasks, usage # noqa: E402 import json # noqa: E402 WIDTH = 100 SKIP_DIRS = {".worktrees", "node_modules", "lib", "__pycache__"} class Failure(Exception): code = 1 class Usage(Failure): code = 2 def old_spelling(old: str, new: str) -> None: """One-line notice for an option kept one more release under its old name.""" print(f"wf: {old} is now {new}", file=sys.stderr) # ------------------------------------------------------------------ files def stamp(path: Path) -> tuple[int, int]: s = path.stat() return (s.st_mtime_ns, s.st_size) def write_if_unchanged(path: Path, text: str, read_stamp: tuple[int, int]) -> None: """Replace `path` by temp file + rename, unless it changed since `read_stamp`.""" if stamp(path) != read_stamp: raise Failure(f"{path.name} changed on disk since it was read: nothing written, run again") mode = path.stat().st_mode & 0o777 fd, tmp = tempfile.mkstemp(dir=path.parent, prefix=f".{path.name}.", suffix=".tmp") try: with os.fdopen(fd, "wb") as f: f.write(text.encode("utf-8")) os.chmod(tmp, mode) if stamp(path) != read_stamp: raise Failure(f"{path.name} changed on disk since it was read: nothing written, run again") os.replace(tmp, path) finally: if os.path.exists(tmp): os.unlink(tmp) def read(path: Path) -> str: return path.read_bytes().decode("utf-8", errors="replace") class Project: def __init__(self, cfg: config.Config): self.cfg = cfg for name, path in (("tasks", cfg.tasks), ("archive", cfg.archive)): if not path.is_file(): raise Failure(f"workflow.toml: {name} '{cfg.rel(path)}' does not exist") self.tasks_stamp, self.archive_stamp = stamp(cfg.tasks), stamp(cfg.archive) self.tasks_text, self.archive_text = read(cfg.tasks), read(cfg.archive) self.doc = tasks.parse(self.tasks_text) self.new_archive: str | None = None @property def archived(self) -> set[str]: return tasks.archive_ids(self.new_archive or self.archive_text) def taken(self) -> set[str]: return self.doc.ids() | self.archived def resolve(self, given: str) -> str: """Exact id, or a unique prefix of an open id (read commands).""" known = self.doc.ids() if given in known: return given hits = sorted(i for i in known if i.startswith(given)) if len(hits) == 1: return hits[0] if hits: raise Failure(f"'{given}' matches {', '.join(hits)}") raise Failure(tasks.unknown_id(given, known)) def save(self, dry_run: bool) -> None: cfg = self.cfg new_tasks = tasks.render(self.doc) before = {p.key for p in checks.check(cfg, slow=False)[0]} after = checks.check(cfg, tasks_text=new_tasks, archive_text=self.new_archive, slow=False)[0] added = [p for p in after if p.key not in before] if added: raise Failure("refused: nothing written, the change adds problems:\n" + "\n".join(f" {p}" for p in added)) changes = [(cfg.tasks, self.tasks_text, new_tasks, self.tasks_stamp)] if self.new_archive is not None: changes.insert(0, (cfg.archive, self.archive_text, self.new_archive, self.archive_stamp)) if dry_run: for path, old, new, _ in changes: name = cfg.rel(path) sys.stdout.writelines(difflib.unified_diff( old.replace("\r\n", "\n").splitlines(keepends=True), new.replace("\r\n", "\n").splitlines(keepends=True), name, f"{name} (new)")) return written = [] for path, old, new, read_stamp in changes: if new == old: continue try: write_if_unchanged(path, new, read_stamp) except Failure as e: if written: raise Failure(f"{e}; {written[0]} was already written: remove its new top line(s)") from None raise written.append(cfg.rel(path)) def load_project(args, write: bool = False) -> Project: start = Path(args.project) if args.project else Path.cwd() if not start.is_dir(): raise Failure(f"--project {start}: no such folder") cfg = config.load_at(start) if cfg.format > config.FORMAT: raise Failure(f"project format {cfg.format} is newer than this wf ({config.FORMAT}): update /projects/public/workflow") if write and cfg.format < config.FORMAT: raise Failure(f"TASKS format {cfg.format}, wf needs {config.FORMAT}: " "run wf migrate --write (idle project, one commit)") return Project(cfg) # ------------------------------------------------------------------ output def show(item: tasks.Item) -> str: return "\n".join(item.lines()) def status_word(item: tasks.Item) -> str: if not item.status: return "-" return "blkd" if item.status.startswith("blocked") else "prog" def row(item: tasks.Item, width: int, cfg: config.Config) -> str: prio = "--" if item.prio is None else f"P{item.prio}" model = item.model if item.prio is not None else "-" lane = lanes.lane_of(item, cfg.lanes, cfg.slice_above) if item.prio is not None else "-" lane_w = max(len(l.name) for l in cfg.lanes) mark = f"[{item.sessions}] " if item.prio is not None and item.sessions != "parallel" else "" line = (f"{item.id:<{width}} {prio} {item.effort or '-':<4} {status_word(item):<4} {model:<6} " f"{lane:<{lane_w}} {mark}{item.title}") return line if len(line) <= WIDTH else line[:WIDTH - 1] + "…" def counts(doc: tasks.Doc) -> str: def n(key): return sum(len(s.items) for s in doc.sections if s.key == key) return " · ".join(f"{key} {n(key)}" for key in ("pending", "human", "awaiting", "deferred")) def resolve_refs(cfg: config.Config, item: tasks.Item) -> list[str]: out = [] for path, anchor in item.refs: out.append(refs.resolve(cfg.root, path, anchor)) if anchor and cfg.anchors_index and cfg.anchors_specs and cfg.anchors_specs.is_dir() \ and (cfg.root / path).resolve() == cfg.anchors_index.resolve(): for spec in sorted(cfg.anchors_specs.rglob("*.md")): text = read(spec) if any(anchor in refs.ANCHOR_RE.findall(l) for l in text.split("\n")): out.append(refs.resolve(cfg.root, cfg.rel(spec), anchor)) return out def inbox_path() -> Path: return Path(os.environ.get("WF_INBOX") or HERE / "inbox.md") def inbox_count() -> int: path = inbox_path() if not path.is_file(): return 0 return sum(1 for l in read(path).split("\n") if l.startswith("- ")) # ------------------------------------------------------------------ sessions (lanes) def sessions_dir(cfg: config.Config) -> Path: return cfg.root / ".wf" / "sessions" def alive(pid, socket) -> bool: try: os.kill(int(pid), 0) except (OSError, ValueError, TypeError): return False return bool(socket) and Path(socket).exists() def lane_names(cfg: config.Config) -> list[str]: return [l.name for l in cfg.lanes] def check_lane(cfg: config.Config, lane: str | None) -> None: if lane and lane not in lane_names(cfg): raise Usage(f"unknown lane '{lane}' ({', '.join(lane_names(cfg))})") def sessions(cfg: config.Config) -> dict[str, dict]: """Registered sessions keyed by lane name or "all" (no lane: every lane).""" out = {} folder = sessions_dir(cfg) for name in [*lane_names(cfg), "all"]: try: s = json.loads((folder / f"{name}.json").read_text()) except (OSError, ValueError): continue if isinstance(s, dict): s["alive"] = alive(s.get("pid"), s.get("socket")) out[name] = s return out def my_session() -> tuple[str, str] | None: """(socket, pid) of this agent session from Claude Code's env; None outside one.""" socket, pid = os.environ.get("CLAUDE_CODE_MESSAGING_SOCKET"), os.environ.get("CLAUDE_PID") return (socket, pid) if socket and pid else None def wf_folder(cfg: config.Config, name: str) -> Path: """.wf/<name> in the project root, created; .wf is git-ignored.""" folder = cfg.root / ".wf" / name folder.mkdir(parents=True, exist_ok=True) ignore = folder.parent / ".gitignore" if not ignore.exists(): ignore.write_text("*\n") return folder @contextlib.contextmanager def project_lock(root: Path): """Exclusive project lock (.wf/lock): every wf write and merge-back runs under it. Not reentrant.""" folder = root / ".wf" folder.mkdir(parents=True, exist_ok=True) ignore = folder / ".gitignore" if not ignore.exists(): ignore.write_text("*\n") with open(folder / "lock", "w") as f: fcntl.flock(f, fcntl.LOCK_EX) yield def register(cfg: config.Config, lane: str, model: str) -> str | None: """Record this agent session as the lane's session (lane "all" = every lane; env from Claude Code; skipped without it). Another live session already holding the lane keeps it; returns a warning line then.""" me = my_session() if not me: return None socket, pid = me other = sessions(cfg).get(lane) if other and other["alive"] and str(other.get("pid")) != pid: return (f"another live {lane} session holds this lane: uds:{other['socket']} " "(claims keep tasks apart; tell the owner if unintended)") try: record = {"lane": lane, "model": model, "socket": socket, "pid": int(pid) if pid.isdigit() else pid, "session": os.environ.get("CLAUDE_CODE_SESSION_ID", ""), "at": datetime.datetime.now().isoformat(timespec="seconds")} (wf_folder(cfg, "sessions") / f"{lane}.json").write_text(json.dumps(record) + "\n") except OSError as e: print(f"wf: session not registered: {e}", file=sys.stderr) return None def unregister(cfg: config.Config) -> None: """Drop session records whose pid is this session (orchestrator holds no lane).""" me = my_session() if not me: return for name in [*lane_names(cfg), "all"]: f = sessions_dir(cfg) / f"{name}.json" try: if str(json.loads(f.read_text()).get("pid")) == me[1]: f.unlink() except (OSError, ValueError, AttributeError): pass # ------------------------------------------------------------------ claims (wf status progress) def claim(cfg: config.Config, id: str, lane: str = "") -> None: """Mark task id as held by this session (skipped outside one). lane: the task's lane.""" me = my_session() if not me: return socket, pid = me model = next((s.get("model", name) for name, s in sessions(cfg).items() if str(s.get("pid")) == pid and s.get("socket") == socket), "") try: record = {"id": id, "lane": lane, "model": model, "socket": socket, "pid": int(pid) if pid.isdigit() else pid, "session": os.environ.get("CLAUDE_CODE_SESSION_ID", ""), "at": datetime.datetime.now().isoformat(timespec="seconds")} (wf_folder(cfg, "claims") / f"{id}.json").write_text(json.dumps(record) + "\n") except OSError as e: print(f"wf: claim not recorded: {e}", file=sys.stderr) def unclaim(cfg: config.Config, ids) -> None: for id in ids: try: (cfg.root / ".wf" / "claims" / f"{id}.json").unlink(missing_ok=True) except OSError: pass def held(p: "Project") -> dict[str, str]: """id → holder, for tasks in progress claimed by another live session.""" me = my_session() pending = {i.id: i for i in p.doc.section("pending").items} out = {} for path in sorted((p.cfg.root / ".wf" / "claims").glob("*.json")): try: c = json.loads(path.read_text()) except (OSError, ValueError): continue if not isinstance(c, dict) or (me and str(c.get("pid")) == me[1]) or not alive(c.get("pid"), c.get("socket")): continue item = pending.get(path.stem) if item is None or not (item.status or "").startswith("in progress"): continue out[item.id] = f"{c.get('model') + ' ' if c.get('model') else ''}session uds:{c.get('socket')}" return out def stale(p: "Project") -> list: """Items in progress whose claim is missing or held by a dead session.""" out = [] for item in p.doc.section("pending").items: if not (item.status or "").startswith("in progress"): continue try: c = json.loads((p.cfg.root / ".wf" / "claims" / f"{item.id}.json").read_text()) except (OSError, ValueError): c = None if not isinstance(c, dict) or not alive(c.get("pid"), c.get("socket")): out.append(item) return out def other_live(cfg: config.Config) -> int: """Live sessions in the project other than this one.""" me = my_session() return sum(1 for s in sessions(cfg).values() if s["alive"] and not (me and str(s.get("pid")) == me[1])) def multi_lines(p: Project, args, lane: str, extra: int) -> list[str]: """Worktree mode advice when more than one session is live in the project (git only).""" n = sum(1 for s in sessions(p.cfg).values() if s["alive"]) + extra main_top = config.git_top(p.cfg.root) if n < 2 or not main_top or not (main_top / ".git").is_dir(): return [] head = f"{n} live sessions here: " wt = config.linked_worktree(Path(args.project or Path.cwd()).resolve()) if wt: return [head + f"you are in worktree {p.cfg.rel(wt[0])} (branch {config.git_branch(wt[1])}); " "wf writes the main tree's TASKS.md"] if code := config.code_main(p.cfg): master = config.git_branch(code / ".git") or "master" cmd = f"wf start <task> --worktree {code / '.worktrees' / lane} --branch {lane}/<task>" return [head + f"work in your lane's code worktree, never on {master}:", f" {cmd} (run in {p.cfg.root}; wf there writes this TASKS.md)"] folder, master = p.cfg.root / ".worktrees" / lane, config.git_branch(main_top / ".git") or "master" cmd = (f"cd {p.cfg.rel(folder)} && git switch -c {lane}/<task> {master}" if folder.is_dir() else f"git worktree add {p.cfg.rel(folder)} -b {lane}/<task> {master}") if p.cfg.worktree_setup: cmd += " && wf setup" if folder.is_dir() else f" && cd {p.cfg.rel(folder)} && wf setup" return [head + f"work in your lane's worktree, never on {master}:", f" {cmd} (wf there writes this TASKS.md)"] def lane_lines(p: Project, me: str | None, runner: bool = False) -> list[str]: a = (p.doc, p.archived, p.cfg.lanes, p.cfg.slice_above, held(p)) return lanes.lanes_block(lanes.counts(*a, runner=runner), sessions(p.cfg), me, lanes.not_ready(*a) if runner else None) # ------------------------------------------------------------------ read commands def cmd_next(args) -> int: p = load_project(args) lane = args.lane check_lane(p.cfg, lane) model = args.as_ or "haiku" out: list[str] = [] warning = None if not args.as_: out.append("no --as: treated as haiku; pass --as " + "|".join(tasks.MODELS)) elif warning := register(p.cfg, lane or "all", model): out.append(warning) claims = held(p) if solo := tasks.solo_running(p.doc, claims): if out: print("\n".join(out)) raise Failure(f"solo {solo[0]} in progress by {solo[1]}: wait (its wf done notifies you)") item, skipped = lanes.pick(p.doc, p.archived, p.cfg.lanes, p.cfg.slice_above, lane, model, claims, other_live(p.cfg), args.owner) multi = multi_lines(p, args, lane or "all", 1 if warning else 0) if (marker := p.cfg.root / ".wf" / "push-failed").is_file(): out += [f"push failed: {marker.read_text().splitlines()[0]} → wf push (in {p.cfg.root})"] if not args.brief: awaiting = p.doc.section("awaiting").items if awaiting: out += ["===== Awaiting your decision (mention, don't block) =====", *(f"- {i.id}: {i.text}" for i in awaiting), ""] human = p.doc.section("human").items if human: out += ["===== Needs human (not picked) =====", *(f"- {i.id}: {i.title}" for i in human), ""] flight = ledgers.in_flight(p.cfg.root, p.cfg.ledgers) if p.cfg.ledgers and p.cfg.ledgers.is_dir() else [] if flight: out += ["===== In-flight plans =====", *flight, ""] if skipped: out += ["===== Skipped =====", *(f"- {i.id}: {why}" for i, why in skipped), ""] if multi: out += ["===== Multi-session =====", *multi, ""] lines = lane_lines(p, lane) if len(lines) > 1: out += ["===== Lanes =====", *lines, ""] waits = lanes.waiting_block(lanes.cross_waits(p.doc, p.archived, lane, p.cfg.lanes, p.cfg.slice_above), sessions(p.cfg), p.cfg.lanes, p.cfg.slice_above) if lane else [] if waits: out += ["===== Waiting on other lanes =====", *waits, ""] if item: if not args.brief: out.append("===== Next task =====") out.append(show(item)) if lanes.is_slice_job(item, p.cfg.slice_above): out += ["", "===== Slice job (no code) =====", f"effort {item.effort} > slice_above {p.cfg.slice_above}: split it, don't implement.", f' wf add "<title>. <goal>" -e <1h|1h> --parent {item.id} --model haiku|sonnet|opus' " (then wf body: Steps/Done/Ref)", f' wf note {item.id} "sliced into …" · wf status {item.id} clear · not wf done' " (the parent waits on its slices)"] if not args.brief: for block in resolve_refs(p.cfg, item): out += ["", block] n = 0 if args.brief else inbox_count() if n: out += ["", f"workflow inbox: {n} reports (triage: workflow session)"] if hint := ctx_hint_line(p.cfg): out += ["", hint] while out and not out[-1]: out.pop() if out: print("\n".join(out)) if not item: raise Failure(f"nothing pickable for {lane or 'all lanes'} ({model}) in Pending") return 0 def cmd_areas(args) -> int: p = load_project(args, write=bool(args.mark)) file = p.cfg.areas_file rel = p.cfg.areas_rel text = file.read_text(encoding="utf-8", errors="replace") if file.is_file() else "" found = areas_mod.parse(text) if not found: print(f"no areas in {rel}") return 0 names = ", ".join(a.name for a in found) wanted = args.name or args.mark if wanted and wanted not in [a.name for a in found]: hits = [a.name for a in found if a.name.startswith(wanted)] if len(hits) == 1: if args.name: args.name = hits[0] else: args.mark = hits[0] wanted = hits[0] elif hits: raise Usage(f"ambiguous area '{wanted}' in {rel}: " + ", ".join(hits)) else: raise Usage(f"no area '{wanted}' in {rel} ({names})") if args.mark: r = git_run(p.cfg.code, "rev-parse", "--short", "HEAD") if r.returncode != 0: raise Failure("git rev-parse HEAD failed: " + r.stderr.strip()) text = areas_mod.mark(text, args.mark, r.stdout.strip()) file.write_text(text, encoding="utf-8") found = areas_mod.parse(text) for a in found: if wanted and a.name != wanted: continue stale, miss, commits = areas_mod.status(p.cfg.code, a, p.cfg.area_stale_commits) line = " · ".join([f"{a.name}: " + ("stale" if stale else "ok")] + ([f"missing: {', '.join(miss)}"] if miss else []) + ([f"{commits} commits since {a.checked}"] if commits is not None else ["no Checked"] if not a.checked else [])) print(line) if not wanted: for folder, n in uncovered_areas(p, args): print(f"uncovered: {folder} ({n} file{'s' * (n != 1)}): no area's Paths covers it") return 0 def cmd_lanes(args) -> int: p = load_project(args) check_lane(p.cfg, args.lane) if args.unregister: unregister(p.cfg) elif args.lane or args.as_: register(p.cfg, args.lane or "all", args.as_ or "haiku") if args.wait is not None: import time poll = float(os.environ.get("WF_LANES_POLL", "30")) end = time.monotonic() + args.wait while True: p = load_project(args) if (p.cfg.root / "out" / "wf-batch.stop").exists(): print("stop requested: out/wf-batch.stop") return 2 if any(c[0] for c in lanes.counts(p.doc, p.archived, p.cfg.lanes, p.cfg.slice_above, held(p), runner=True).values()): break left = end - time.monotonic() if left <= 0: print("\n".join(lane_lines(p, args.lane, True))) return 1 time.sleep(min(poll, left)) print("\n".join(lane_lines(p, args.lane, True))) return 0 def runner_ids(p: Project) -> set[str]: """Pending ids a headless worker can take: not blocked, After done, runner-ready, no open slices, not sliced out.""" return {i.id for i in p.doc.section("pending").items if not i.error and not (i.status or "").startswith("blocked") and all(a in p.archived for a in i.after) and i.runner_ready and not tasks.open_slices(p.doc, i.id) and not lanes.sliced_out(i, p.cfg.slice_above)} def cmd_list(args) -> int: p = load_project(args) check_lane(p.cfg, args.lane) keys = list(tasks.SECTIONS) if args.section == "all" else [args.section] ready = None if args.ready: ready = {i.id for i in p.doc.section("pending").items if not i.error and not (i.status or "").startswith("blocked") and all(a in p.archived for a in i.after)} if args.runner: ready = runner_ids(p) stale_ids = {i.id for i in stale(p)} if args.stale else None want_ref = tuple(args.ref.split("#", 1)) if args.ref else None def keep(i: tasks.Item) -> bool: if args.prio is not None and (i.prio is None or i.prio > args.prio): return False if ready is not None and i.id not in ready: return False if args.blocked and status_word(i) != "blkd": return False if args.progress and status_word(i) != "prog": return False if stale_ids is not None and i.id not in stale_ids: return False if args.model and (i.prio is None or i.model != args.model): return False if args.lane and (i.prio is None or lanes.lane_of(i, p.cfg.lanes, p.cfg.slice_above) != args.lane): return False if want_ref and not any(path == want_ref[0] and (len(want_ref) == 1 or anchor == want_ref[1]) for path, anchor in i.refs): return False return True groups = [(s, [i for i in s.items if keep(i)]) for k in keys for s in p.doc.sections if s.key == k] if args.runner and args.lane: # pick order of the lane, then its fallback lane's order = lanes.ranked(p.doc, p.archived, p.cfg.lanes, p.cfg.slice_above, args.lane) keys, groups = ["pending"], [(None, [i for i in order if i.id in ready and (args.prio is None or i.prio <= args.prio)])] if args.n is not None: left = args.n for n, (s, items) in enumerate(groups): groups[n] = (s, items[:left]) left -= len(groups[n][1]) width = max((len(i.id) for _, items in groups for i in items), default=0) for s, items in groups: if len(keys) > 1: print(f"== {s.heading}") for i in items: print(row(i, width, p.cfg)) print(counts(p.doc)) return 0 def cmd_show(args) -> int: p = load_project(args) print("\n\n".join(show(p.doc.item(p.resolve(i))) for i in args.ids)) return 0 def cmd_ctx(args) -> int: p = load_project(args) target = args.target if "#" in target or "/" in target or (p.cfg.root / target).exists(): path, _, anchor = target.partition("#") users = [i.id for i in p.doc.all_items() if any(r == path and (not anchor or a == anchor) for r, a in i.refs)] print(refs.resolve(p.cfg.root, path, anchor or None)) print("\nTasks: " + (", ".join(users) if users else "none")) return 0 if target not in p.doc.ids() and target in p.archived: line = next(l for l in p.archive_text.replace("\r\n", "\n").split("\n") if (m := tasks.ARCHIVE_ID_RE.match(l)) and m.group(1) == target) print(f"done: {line}") return 0 id = p.resolve(target) section, _ = p.doc.find(id) item = p.doc.item(id) where = {i.id: s.heading for s in p.doc.sections for i in s.items} out = [show(item), "", f"Section: {section.heading}"] if item.after: out.append("After: " + ", ".join( f"{a} ({'open, ' + where[a] if a in where else 'done' if a in p.archived else 'unknown'})" for a in item.after)) needed = [i.id for i in p.doc.all_items() if id in i.after] blocks = [i.id for i in p.doc.all_items() if i.blocked_on == id] linked = [i.id for i in p.doc.all_items() if i.id != id and i.id not in needed and i.id not in blocks and id in refs.links("\n".join(i.lines()))] for label, ids in (("Needed by", needed), ("Blocks", blocks), ("Linked from", linked)): if ids: out.append(f"{label}: {', '.join(ids)}") if p.cfg.verify and not id.startswith("a-"): out.append("Verify (run before wf finish): " + " · ".join(p.cfg.verify)) for block in resolve_refs(p.cfg, item): out += ["", block] out += area_blocks(p, "\n".join(item.lines())) print("\n".join(out)) return 0 def area_blocks(p: Project, text: str) -> list[str]: """Notes of the areas the task text names (area name, anchor or a path no other area lists) → ctx lines.""" file = p.cfg.areas_file if not file.is_file(): return [] notes = file.read_text(encoding="utf-8", errors="replace") out = [] found = areas_mod.parse(notes) skip = areas_mod.shared_paths(found) for a in found: if areas_mod.matches(a, text, skip): out += ["", f"Area {a.name} ({p.cfg.areas_rel}):", areas_mod.block(notes, a.name)] return out def cmd_search(args) -> int: p = load_project(args) kinds = {k for k, on in (("task", args.tasks), ("archive", args.archive), ("doc", args.docs)) if on} or None docs = {} if kinds is None or "doc" in kinds: docs = {p.cfg.rel(f): read(f) for f in checks.doc_files(p.cfg)} hits = search.search(args.words, p.doc, p.archive_text, docs, kinds=kinds, limit=args.n, archive_name=p.cfg.rel(p.cfg.archive)) if not hits: raise Failure("no hits") for h in hits: if h.kind == "task": header = p.doc.item(h.where).text if h.where in p.doc.ids() else "" text = h.line if h.line == header else f"{h.label} — {h.line}" elif h.kind == "archive": text = re.sub(r"^\d{4}-\d\d-\d\d \*\*([^*]+)\*\*", r"\1", h.line) else: text = f"{h.label} — {h.line}" if h.line and h.line != h.label else h.label line = f"{h.where} {text}" print(line if len(line) <= 2 * WIDTH else line[:2 * WIDTH - 1] + "…") return 0 def cmd_log(args) -> int: p = load_project(args) lines = [l for l in p.archive_text.replace("\r\n", "\n").split("\n") if l.startswith("- ")] words = [w.lower() for w in args.words] lines = [l for l in lines if all(w in l.lower() for w in words)] print("\n".join(lines[:args.n])) return 0 def cmd_check(args) -> int: p = load_project(args) errors, warnings = checks.check(p.cfg) warnings += checks.merged_tool_worktrees(Path(os.environ.get("WF_TOOL_ROOT") or HERE)) for w in warnings: print(f"warn: {w}") for e in errors: print(f"ERROR: {e}") print(("OK: " if not errors else "") + f"{len(errors)} errors · {len(warnings)} warnings") return 1 if errors else 0 def find_projects(root: Path, depth: int = 3) -> list[Path]: out = [] def walk(folder: Path, left: int) -> None: if (folder / "wf.py").is_file() and (folder / "wflib").is_dir(): return # a wf checkout (its templates/ is no project) if (folder / config.NAME).is_file(): out.append(folder) return if left == 0: return try: children = sorted(c for c in folder.iterdir() if c.is_dir() and not c.is_symlink()) except OSError: return for c in children: if not c.name.startswith(".") and c.name not in SKIP_DIRS: walk(c, left - 1) walk(root, depth) return out def default_root(tool: Path) -> Path: """Folder `wf projects` scans: the tool's parent, else its grandparent (tool under e.g. public/), if it holds projects.""" return next((d for d in (tool.parent, tool.parent.parent) if find_projects(d)), tool.parent) def projects_root() -> Path: return Path(os.environ.get("WF_ROOT") or default_root(HERE)) def cmd_projects(args) -> int: root = projects_root() rows = [] for folder in find_projects(root): name = str(folder.relative_to(root)) try: cfg = config.load(folder) doc = tasks.parse(read(cfg.tasks)) archived = tasks.archive_ids(read(cfg.archive)) if cfg.archive.is_file() else set() errors = len(checks.check(cfg, slow=False)[0]) item = lanes.pick(doc, archived, cfg.lanes, cfg.slice_above, owner=True)[0] n = {k: sum(len(s.items) for s in doc.sections if s.key == k) for k in tasks.SECTIONS} rows.append((name, f"pending {n['pending']} human {n['human']} awaiting {n['awaiting']} " f"errors {errors} next: " + (f"{item.id} {item.title}" if item else "-"))) except (config.ConfigError, tasks.TaskError, OSError) as e: rows.append((name, f"broken: {e}")) width = max((len(n) for n, _ in rows), default=0) for name, text in rows: line = f"{name:<{width}} {text}" print(line if len(line) <= WIDTH else line[:WIDTH - 1] + "…") if not rows: raise Failure(f"no projects under {root}") return 0 def claude_projects() -> Path: return Path(os.environ.get("CLAUDE_CONFIG_DIR") or Path.home() / ".claude") / "projects" def since_time(value: str | None) -> str | None: """ISO (UTC) as is; 90m / 2h / 1d = that long ago.""" m = re.fullmatch(r"(\d+)([mhd])", value or "") if not m: return value delta = datetime.timedelta(**{{"m": "minutes", "h": "hours", "d": "days"}[m[2]]: int(m[1])}) return (datetime.datetime.now(datetime.timezone.utc) - delta).strftime("%Y-%m-%dT%H:%M:%S") def agent_label(path: Path) -> str: id = path.stem.removeprefix("agent-") try: desc = json.loads(path.with_name(path.stem + ".meta.json").read_text()).get("description") or "" except (OSError, ValueError, AttributeError): desc = "" return f"{id} {desc}".strip() def usage_transcripts(args) -> list[tuple[str, Path]]: """(label, jsonl): one agent, or a session's main transcript + its subagents.""" base = claude_projects() if args.agent: id = args.agent.removeprefix("agent-") found = sorted(base.glob(f"*/*/subagents/agent-{id}.jsonl")) if not found: raise Failure(f"no transcript for agent {args.agent}") return [(agent_label(found[0]), found[0])] sid = args.session or os.environ.get("CLAUDE_CODE_SESSION_ID") if sid: found = sorted(base.glob(f"*/{sid}*.jsonl")) if len(found) != 1: raise Failure(f"session {sid}: {len(found)} transcripts match") main = found[0] else: folder = base / re.sub(r"[^A-Za-z0-9]", "-", str(Path.cwd())) found = sorted(folder.glob("*.jsonl"), key=lambda p: p.stat().st_mtime) if not found: raise Failure(f"no transcripts in {folder} (--session ID)") main = found[-1] agents = sorted((main.parent / main.stem / "subagents").glob("agent-*.jsonl")) return [("main", main)] + [(agent_label(p), p) for p in agents] def cmd_usage(args) -> int: since = since_time(args.since) if args.report: return usage_report(since) if args.explore: return usage_explore(args, since) if args.log and not args.agent: raise Usage("--log needs --agent") if args.log: root = config.find_root(Path(args.project) if args.project else Path.cwd()) path = usage_transcripts(args)[0][1] now = datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") line = usage.log_line(now, root.name, args.log[0], args.effort or "-", args.log[1], path.stem.removeprefix("agent-"), usage.parse(read(path), since=since), lane=args.lane, dur=args.duration) (root / "out").mkdir(exist_ok=True) with open(root / "out" / "wf-cost.log", "a") as f: f.write(line + "\n") print(line) return 0 rows, total, dollars, unknown = [], usage.Usage(), 0.0, False for label, path in usage_transcripts(args): for model, u in usage.parse(read(path), since=since).items(): c = usage.cost(model, u) total.add(u) dollars += c or 0 unknown |= c is None rows.append((label, usage.short(model), u, "?" if c is None else f"{c:.2f}")) width = min(max([len(r[0]) for r in rows] + [5]), 40) def nums(u): out = ("~" if u.est else "") + usage.fmt(u.out) # ~ = estimated from content (subagent) return " ".join(f"{usage.fmt(n):>7}" for n in (u.inp, u.cw, u.cr)) + f" {out:>7}" print(f"{'agent':<{width}} {'model':<12} {'turns':>6} {'in':>7} {'cw':>7} {'cr':>7} {'out':>7} {'$':>7}") for label, model, u, c in rows: label = label if len(label) <= width else label[:width - 1] + "…" print(f"{label:<{width}} {model:<12} {u.turns:>6} {nums(u)} {c:>7}") print(f"{'total':<{width}} {'':<12} {total.turns:>6} {nums(total)} {dollars:>7.2f}{'+' if unknown else ''}") return 0 def usage_explore(args, since: str | None) -> int: root = str(Path.cwd()) agents = [(label, usage.explore(read(path), since=since, root=root)) for label, path in usage_transcripts(args) if label != "main"] r = usage.explore_report(agents) if not r["rows"]: print("no subagent with >= 4 calls") return 0 width = min(max([len(x[0]) for x in r["rows"]] + [5]), 40) print(f"{'agent':<{width}} {'calls':>5} {'1stEdit':>7} {'explore':>7} {'pre-edit':>8}") for label, calls, fe, exp, pre in r["rows"]: label = label if len(label) <= width else label[:width - 1] + "…" print(f"{label:<{width}} {calls:>5} {fe:>7} {exp:>7.0%} {pre:>8.0%}") tot = r["total"] or 1 print(f"ALL {r['agents']} agents, {usage.fmt(r['total'])} billed input: " + ", ".join(f"{k} {100 * v // tot}%" for k, v in sorted(r["cost"].items(), key=lambda kv: -kv[1]))) if r["pre_n"]: print(f"before 1st edit (median, {r['pre_n']} agents): {r['pre_share']:.0%} of input, " f"{r['pre_calls']:g} calls, +{usage.fmt(int(r['pre_ctx']))} context") rt = sum(r["results"].values()) or 1 print("result tokens: " + ", ".join(f"{k} {usage.fmt(v)} {100 * v // rt}%" for k, v in sorted(r["results"].items(), key=lambda kv: -kv[1]))) if r["files"]: print("files read by >= 2 agents (agents, tokens):") for f, a, n in r["files"]: print(f"{a:>4} {usage.fmt(n):>7} {f}") return 0 def fmt_dur(s) -> str: return "-" if s is None else f"{int(s) // 3600}h{int(s) % 3600 // 60:02d}m" if s >= 3600 else f"{int(s) // 60}m{int(s) % 60:02d}s" def usage_report(since: str | None) -> int: root = projects_root() entries = [] for folder in find_projects(root): log = folder / "out" / "wf-cost.log" if log.is_file(): entries += usage.parse_log(read(log), since=since) if not entries: raise Failure(f"no out/wf-cost.log lines in the projects under {root}") print(f"{'lane':<8} {'effort':<6} {'n':>4} {'done':>4} {'other':>5} {'med$':>7} {'total$':>8} {'$/done':>7} med_turns med_dur") for lane, effort, n, done, other, med, total, per, turns, dur in usage.report(entries, tasks.EFFORTS): per = "-" if per is None else f"{per:.2f}" print(f"{lane:<8} {effort:<6} {n:>4} {done:>4} {other:>5} {med:>7.2f} {total:>8.2f} {per:>7} {turns:>9} {fmt_dur(dur):>8}") return 0 # ------------------------------------------------------------------ write commands def split_list(value: str) -> list[str]: return tasks.split_refs(value) def stdin_lines() -> list[str]: return sys.stdin.read().replace("\r\n", "\n").split("\n") def set_slices_line(parent: tasks.Item, id: str) -> None: at = parent._line(tasks.SLICES_RE) if at is None: parent.body.insert(tasks._tail_start(parent), f"{tasks.INDENT}- Slices: [[{id}]]") else: parent.body[at] = parent.body[at].rstrip() + f", [[{id}]]" def cmd_add(args) -> int: p = load_project(args, write=True) doc = p.doc if args.title == "-": item = tasks.parse_block(sys.stdin.read()) key = args.section or ("awaiting" if item.id.startswith("a-") else "pending") else: key = args.section or "pending" prio, after = args.prio, split_list(args.after or "") text = " ".join(args.title.split()) leading = tasks.LEADING_ID_RE.match(text) if leading and not args.id: args.id, text = leading.groups() if args.parent: parent_section, _ = doc.find(args.parent) parent = doc.item(args.parent) key = args.section or parent_section.key prio = parent.prio if prio is None else prio id = args.id or tasks.slice_id(args.parent, p.taken()) m = re.match(re.escape(args.parent) + r"-(\d+)$", id) prev = parent.slices or ([f"{args.parent}-{int(m.group(1)) - 1}"] if m and int(m.group(1)) > 1 else []) if not after: # a deferred previous slice would block a live one forever: chain to the last live slice live = [s for s in prev if s in p.archived or key == "deferred" or (s in doc.ids() and doc.find(s)[0].key != "deferred")] after = live[-1:] task = key != "awaiting" if task and prio is None: raise Usage("add: a task needs -p 0-3") if task and args.effort is None: raise Usage("add: a task needs -e " + "|".join(tasks.EFFORTS)) if task and args.effort not in tasks.EFFORTS: raise Failure(f"effort '{args.effort}' (want {', '.join(tasks.EFFORTS)})") if not text: raise Usage("add: empty title") if not text.endswith((".", "?", "!")): text += "." if not args.parent: title = text.split(". ", 1)[0] id = args.id or tasks.make_id(title, p.taken(), "t-" if task else "a-") if not tasks.ID_RE.match(id): raise Failure(f"bad id '{id}' (want t-… or a-…, lowercase a-z 0-9 -)") if id in p.archived: raise Failure(f"id '{id}' already used in the archive (ids are never reused)") item = tasks.Item(id=id, prio=prio if task else None, effort=args.effort if task else None, text=text) if args.body: item.body = tasks.indent_body(stdin_lines()) if args.interactive: old_spelling("--interactive", "--sessions owner") args.sessions = args.sessions or "owner" if args.done and task: item.body.append(f"{tasks.INDENT}Done: {args.done}") if args.sessions and args.sessions != "parallel" and task: item.body.append(f"{tasks.INDENT}Sessions: {args.sessions}") if args.model and task: item.body.append(f"{tasks.INDENT}Model: {args.model}") if args.cloud and task: item.body.append(f"{tasks.INDENT}Cloud: {args.cloud}") if after: item.body.append(f"{tasks.INDENT}- After: " + ", ".join(f"[[{a}]]" for a in after)) if args.ref: item.body.append(f"{tasks.INDENT}Ref: " + ", ".join(split_list(args.ref))) if args.parent: set_slices_line(parent, id) tasks.insert(doc, item, key) p.save(args.dry_run) if not args.dry_run: print(item.header()) if item.prio is not None and item.done_text is None: if item.prio == 0: print("wf: warning: P0 without Done line is not runner-pickable; use --done \"<text>\"", file=sys.stderr) else: print("hint: no Done line; add one (--done) so runners can pick it", file=sys.stderr) return 0 RELEARN_LINE = "re-learned anything (>3 greps to find)? one anchor line → that area's code map" def area_tasks(p: Project) -> list[str]: """Add t-map-<area> refresh tasks for stale areas without an open one; returns output lines.""" file = p.cfg.areas_file if not file.is_file(): return [] out = [] for a in areas_mod.parse(file.read_text(encoding="utf-8", errors="replace")): if not areas_mod.status(p.cfg.code, a, p.cfg.area_stale_commits)[0]: continue prefix = f"t-map-{a.slug}" if any(i == prefix or i.startswith(prefix + "-") for i in p.doc.ids()): continue id, n = prefix, 1 while id in p.archived: n += 1 id = f"{prefix}-{n}" item = tasks.Item(id, prio=1, effort="<1h", text=f"Refresh {a.name} area map.") item.body = tasks.indent_body([ f"Steps: wf areas {a.name}; fix missing anchors + test recipe from git log --stat " f"<Checked>..HEAD -- <paths>; wf areas --mark {a.name}.", f"Done: wf areas {a.name} shows ok.", "Model: sonnet"]) tasks.insert(p.doc, item, "pending") out.append(f"added {id} (area map stale)") return out def uncovered_areas(p: Project, args) -> list[tuple[str, int]]: """Folders this task's diff (branch since its base + uncommitted) touches outside every area's Paths.""" file = p.cfg.areas_file if not file.is_file() or p.cfg.code_root: # code in another repo: this diff is not its code return [] start = Path(getattr(args, "project", None) or Path.cwd()).resolve() here = _here(start, p.cfg) found = [areas_mod.with_anchor_paths(here, a) for a in areas_mod.parse(file.read_text(encoding="utf-8", errors="replace"))] if not any(a.paths for a in found): return [] wt = config.linked_worktree(start) if wt: base = config.git_branch(wt[2] / ".git") or "master" else: base = next((b for b in ("master", "main") if git_run(here, "rev-parse", "-q", "--verify", f"refs/heads/{b}").returncode == 0), None) books = {p.cfg.rel(f) for f in (p.cfg.tasks, p.cfg.archive, p.cfg.root / config.NAME)} | {p.cfg.areas_rel} files = [f for f in areas_mod.changed_files(here, base) if f not in books] return areas_mod.uncovered(files, found, p.cfg.area_ignore) UNCOVERED_MAX = 3 # map tasks one wf done adds at most def uncovered_tasks(p: Project, args) -> list[str]: """Add t-map-<folder> tasks for folders the diff touches outside every area's Paths (once per id).""" out = [] for folder, n in uncovered_areas(p, args)[:UNCOVERED_MAX]: id = f"t-map-{areas_mod.slug(folder)}" if any(i == id or i.startswith(id + "-") for i in [*p.doc.ids(), *p.archived]): continue rel = p.cfg.areas_rel item = tasks.Item(id, prio=1, effort="<1h", text=f"Map {folder} area. No area's Paths covers it.") item.body = tasks.indent_body([ f"Steps: git log --stat -- {folder}; new ### section under ## Areas in {rel} (Code map anchors, " f"Test recipe, Paths: {folder}/) or add {folder}/ to an existing area's Paths; wf areas --mark <area>.", f"Done: wf areas lists the area ok and no longer names {folder} uncovered.", "Model: sonnet"]) tasks.insert(p.doc, item, "pending") out.append(f"added {id} (uncovered area: {folder}, {n} file{'s' * (n != 1)})") return out def cmd_done(args) -> int: p = load_project(args, write=True) doc = p.doc task_ids = [i for i in args.ids if not i.startswith("a-")] if len(task_ids) == 1 and not args.m: raise Usage('done: -m "<entry>" is required for a single task') if len(task_ids) > 1 and args.m: raise Usage("done: -m goes with one task; several ids use each task's goal") out = [] today = datetime.date.today().isoformat() for id in args.ids: doc.find(id) if not config.linked_worktree(Path(args.project or Path.cwd()).resolve()): for id in task_ids: if wt := claimed_worktree(p.cfg.root, doc.item(id)): raise Failure(f"{id} is in progress in worktree {wt[0]} (branch {wt[1]}): run wf finish/done there " f"(its bookkeeping goes in with wf merge), or wf status {id} clear first") before = lanes.pickable_ids(doc, p.archived) lane = lambda i: lanes.lane_of(i, p.cfg.lanes, p.cfg.slice_above) done_lanes = {lane(doc.item(i)) for i in task_ids} solo = [i for i in task_ids if doc.item(i).sessions == "solo"] for id in args.ids: if id.startswith("a-"): item = tasks.remove(doc, id) if args.m: line = tasks.archive_line(today, item, args.m) p.new_archive = tasks.archive_prepend(p.new_archive or p.archive_text, line) out.append(f"removed: {id}" + (f" → {p.cfg.rel(p.cfg.archive)}" if args.m else "")) freed = tasks.unblock(doc, id) if freed: out.append("unblocked: " + ", ".join(freed)) continue open_ = [s for s in tasks.open_slices(doc, id) if s not in args.ids] if open_: raise Failure(f"'{id}' has open slices: {', '.join(open_)}") item = tasks.remove(doc, id) line = tasks.archive_line(today, item, args.m or item.goal) p.new_archive = tasks.archive_prepend(p.new_archive or p.archive_text, line) out.append(f"done: {id} → {p.cfg.rel(p.cfg.archive)}") parent = tasks.parent_of(doc, id) if parent and not tasks.open_slices(doc, parent): out.append(f"last slice of {parent}: finish it with wf done {parent} -m \"…\"") if task_ids and not args.dry_run: out += area_tasks(p) out += uncovered_tasks(p, args) p.save(args.dry_run) if args.dry_run: return 0 unclaim(p.cfg, task_ids) after = lanes.pickable_ids(doc, p.archived | set(task_ids)) freed = [i for s in doc.sections if s.key == "pending" for i in s.items if i.id in after - before and lane(i) not in done_lanes] if task_ids: out += lanes.notify_block(freed, sessions(p.cfg), p.cfg.lanes, p.cfg.slice_above) me = my_session() out += lanes.solo_done_block(solo, sessions(p.cfg), me[1] if me else None) if task_ids: if p.cfg.verify: out += ["verify:", *(f" {v}" for v in p.cfg.verify)] out += ["checklist:", f" - {RELEARN_LINE}", *(f" - {d}" for d in p.cfg.done)] if not getattr(args, "finishing", False): out += worktree_done_lines(p, args) if hint := ctx_hint_line(p.cfg): out.append(hint) print("\n".join(out)) return 0 CTX_TAIL = 1 << 20 # bytes of the transcript read for the last request def ctx_hint_line(cfg) -> str | None: """This session's prompt size over cfg.ctx_hint → a /clear hint; no session id / transcript → None.""" sid = os.environ.get("CLAUDE_CODE_SESSION_ID") if not cfg.ctx_hint or not sid: return None found = sorted(claude_projects().glob(f"*/{glob.escape(sid)}.jsonl")) if len(found) != 1: return None try: with open(found[0], "rb") as f: f.seek(max(0, f.seek(0, 2) - CTX_TAIL)) text = f.read().decode("utf-8", "replace") except OSError: return None return usage.ctx_hint(usage.context_tokens(text), cfg.ctx_hint) def claimed_worktree(root: Path, item: tasks.Item) -> tuple[str, str] | None: """(worktree path, branch) of a linked worktree on the branch named by item's 'in progress: <branch>'.""" status = item.status or "" if not status.startswith("in progress: "): return None words = status.removeprefix("in progress: ").split() branch = words[0] if words else "" path = None for line in git_run(root, "worktree", "list", "--porcelain").stdout.splitlines(): if line.startswith("worktree "): path = line.removeprefix("worktree ") elif branch and line == f"branch refs/heads/{branch}" and path and Path(path).resolve() != root.resolve(): if config.linked_worktree(Path(path)): return path, branch return None def worktree_done_lines(p: Project, args) -> list[str]: """Merge-back steps when wf done runs in a linked worktree (multi-session mode).""" wt = config.linked_worktree(Path(args.project or Path.cwd()).resolve()) if not wt: return [] branch, main = config.git_branch(wt[1]), wt[2] master = config.git_branch(main / ".git") or "master" bmain = config.git_top(p.cfg.root) or main # books repo (private repo of a split project) files = " ".join(str(f.relative_to(bmain)) for f in (p.cfg.tasks, p.cfg.archive)) warn = [] if not branch and (n := unmerged_count(wt[0], master)): warn = [f"detached HEAD has {n} commit{'s' * (n != 1)} not in {master}: wf merge merges " f"{'them' if n != 1 else 'it'} (never leave them unmerged)"] return warn + [f"worktree mode (branch {branch or 'none: detached HEAD'}), after verify: commit your code here (explicit paths), then:", f" wf merge (rebase, ff-merge into {master}, commit {files}, push home; " f"conflict → git rebase {master}, resolve, verify, wf merge again)"] def git_run(cwd: Path, *args: str) -> subprocess.CompletedProcess: return subprocess.run(["git", *args], cwd=cwd, capture_output=True, text=True) def unmerged_count(top: Path, master: str) -> int: """Commits on this worktree's HEAD not in master (0 when git fails).""" r = git_run(top, "rev-list", "--count", f"{master}..HEAD") return int(r.stdout.strip()) if r.returncode == 0 and r.stdout.strip().isdigit() else 0 def cmd_merge(args) -> int: """Merge-back of a lane worktree's branch, under the project lock (main() takes it).""" start = Path(args.project or Path.cwd()).resolve() wt = config.linked_worktree(start) if not wt: raise Failure("merge runs inside a linked git worktree (lane worktree)") top, gitdir, main = wt branch = config.git_branch(gitdir) master = config.git_branch(main / ".git") or "master" cfg = config.load_at(start) bmain = config.git_top(cfg.root) or main # books repo: the private repo of a split project, else main private = [str(Path(x).resolve().relative_to(bmain.resolve())) for x in getattr(args, "private", None) or []] files = [str(f.relative_to(bmain)) for f in (cfg.tasks, cfg.archive)] + private if _dirty_outside(top, []): raise Failure("worktree has uncommitted changes: commit them first") had = unmerged_count(top, master) if not branch and not had: # after a merge: notes / follow-ups written since → bookkeeping commit only if not git_run(bmain, "status", "--porcelain", "--", *files).stdout.strip(): raise Failure("worktree is on a detached HEAD: nothing to merge") if private and (r := git_run(bmain, "add", "--", *private)).returncode: raise Failure(f"git add failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") r = git_run(bmain, "commit", "-q", "-m", args.m or "bookkeeping", "--", *files) if r.returncode: raise Failure(f"bookkeeping commit failed: {(r.stderr.strip() or r.stdout.strip() or 'git error').splitlines()[-1]}") print(f"committed {' '.join(files)}") return 0 out = [] # a branch, or a detached HEAD with commits not in master (never left behind) if git_run(top, "rebase", master).returncode: git_run(top, "rebase", "--abort") raise Failure(f"rebase onto {master} conflicts: git rebase {master}, resolve, verify, then wf merge again") out.append(f"rebased onto {master}") args.merged_sha = (git_run(top, "rev-parse", "--short", "HEAD").stdout.strip() if had or bmain == main else "-") r = git_run(main, "merge", "--ff-only", branch or git_run(top, "rev-parse", "HEAD").stdout.strip()) if r.returncode: raise Failure(f"ff-merge into {master} failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") out.append(f"fast-forwarded {master}") if private and (r := git_run(bmain, "add", "--", *private)).returncode: raise Failure(f"git add failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") if git_run(bmain, "status", "--porcelain", "--", *files).stdout.strip(): msg = (getattr(args, "private_msg", None) if private else None) or args.m \ or (f"{branch.rsplit('/', 1)[-1]} done" if branch else "bookkeeping") if bmain != main and args.merged_sha != "-": msg += f" (code {args.merged_sha})" r = git_run(bmain, "commit", "-q", "-m", msg, "--", *files) if r.returncode: raise Failure(f"bookkeeping commit failed: {(r.stderr.strip() or r.stdout.strip() or 'git error').splitlines()[-1]}") out.append(f"committed {' '.join(files)}") args.books_sha = git_run(bmain, "rev-parse", "--short", "HEAD").stdout.strip() git_run(top, "switch", "-q", "--detach", master) if branch: git_run(top, "branch", "-q", "-d", branch) args.push_rc = 0 if not args.no_push: args.push_rc, lines = run_push(cfg, [bmain, main]) out += lines out.append(f"merged {branch or 'detached HEAD'} into {master}") print("\n".join(out)) return 0 def run_push(cfg: config.Config, repos: list[Path]) -> tuple[int, list[str]]: """Push after a merge: workflow.toml push lines (cwd = root, env WF_MAIN; first nonzero stops and writes .wf/push-failed), else _push_home on each repo. (exit code, lines to print).""" marker = cfg.root / ".wf" / "push-failed" if not cfg.push: pushed = [r for r in dict.fromkeys(repos) if _push_home(r)] return 0, ["pushed home"] if pushed else [] env = {**os.environ, "WF_MAIN": str(cfg.root)} for line in cfg.push: r = subprocess.run(line, shell=True, cwd=cfg.root, env=env, capture_output=True, text=True) if r.returncode: tail = (r.stdout + r.stderr).strip().splitlines()[-5:] wf_folder(cfg, "orch") # creates .wf with its .gitignore marker.write_text(f"{line}\nexit {r.returncode}\n" + "\n".join(tail) + "\n") return r.returncode, [f"not pushed (exit {r.returncode}): {line}", *(f" {t}" for t in tail), "rerun: wf push"] marker.unlink(missing_ok=True) return 0, [f"pushed ({len(cfg.push)} push command{'s' * (len(cfg.push) != 1)})"] def cmd_push(args) -> int: """Run the push of a merge again (workflow.toml push, else push home), from the project or a code worktree.""" cfg = config.load_at(Path(args.project or Path.cwd()).resolve()) repos = [r for r in (config.git_top(cfg.root), config.code_main(cfg)) if r] rc, lines = run_push(cfg, repos) print("\n".join(lines) or "nothing to push (no push key, no remote home)") return rc def _push_home(main: Path) -> bool: """git push home --all/--tags when remote home exists; True if pushed.""" if "home" not in git_run(main, "remote").stdout.split(): return False for extra in ("--all", "--tags"): r = git_run(main, "push", "-q", "home", extra) if r.returncode: raise Failure(f"push home failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") return True def _dirty_outside(top: Path, keep: list[Path]) -> list[str]: """Uncommitted paths of the repo at `top` not under any of `keep` (absolute paths).""" stray = [] for line in git_run(top, "status", "--porcelain", "-uall").stdout.splitlines(): rel = line[3:].split(" -> ")[-1].strip('"').rstrip("/") path = (top / rel).resolve() if rel == config.WF_HOME: # wf start's pointer in a code worktree, never committed continue if not any(path == k or k in path.parents for k in keep): stray.append(rel) return stray GATE_RED_HELP = ( "\ngate red procedure (also for a red gate run after wf finish, on the merged sha):" "\n- culprit = your diff (git show --stat <sha>) explains the failure -> fix now (new commit, wf merge) or hand back" "\n- not your diff (earlier task broke it; see git log of the failing area) -> wf add -p 0 --model <m> -e <effort>" ' --done "gate ALL GREEN" -b "Fix gate red <test>. <goal>" with body line "Steps: <failing test, culprit sha/id>"' "\n- then report result: done+gate-red <culprit sha or id> <fix id> (your task is done, the lane goes on with the fix)" "\n- coalesced gate (wf res wait prints 'covers …; red: culprit is any commit in A^..B'): bisect that range for the" " culprit, else the P0 fix task's Steps name the range") def cmd_finish(args) -> int: """gate → done → commit the given paths → merge (lane worktree), in one call; checks first, so a refusal leaves the task open.""" if args.paths and not args.commit: raise Usage("finish: paths need --commit MSG") start = Path(args.project or Path.cwd()).resolve() root = config.find_root(start) cfg = config.load(root) wt = config.linked_worktree(start) top = Path(git_run(start, "rev-parse", "--show-toplevel").stdout.strip() or start) paths = [(start / f).resolve() for f in args.paths] books = [] if wt else [cfg.tasks.resolve(), cfg.archive.resolve()] btop = (config.git_top(cfg.root) or cfg.root).resolve() split = bool(wt) and btop != wt[2].resolve() # code worktree of a split project: books in another repo proot, priv, code = cfg.root.resolve(), [], [] for f, path in zip(args.paths, paths): if split and (path == proot or proot in path.parents): if not path.exists() and git_run(btop, "ls-files", "--error-unmatch", "--", str(path)).returncode: raise Failure(f"path {f} does not exist and is not tracked: nothing to commit") priv.append(path) continue code.append(path) if path != top and top not in path.parents: if split: raise Failure(f"path {f} is in neither repo ({top}, {cfg.root}): commit it in its own repo " "first, then wf finish without it") raise Failure(f"path {f} is outside this repo ({top}): commit it in its own repo first, " "then wf finish without it") if wt and not split and path in (top / cfg.tasks.relative_to(wt[2]), top / cfg.archive.relative_to(wt[2])): raise Failure(f"path {f} = this worktree's copy of the books: wf finish writes the main tree's " "and wf merge commits them; drop it") if not path.exists() and git_run(top, "ls-files", "--error-unmatch", "--", str(path)).returncode: raise Failure(f"path {f} does not exist and is not tracked: nothing to commit") paths = code if wt and paths and not git_run(top, "status", "--porcelain", "--", *map(str, paths)).stdout.strip(): raise Failure("nothing to commit in the --commit paths: drop them (wf finish without paths) or fix them") if stray := _dirty_outside(top, paths + books): raise Failure("uncommitted changes outside the --commit paths: " + " ".join(stray)) if cfg.quick_gate: _run_lines(cfg.quick_gate, _here(start, cfg), cfg, "quick_gate '{line}' red (exit {rc}): fix it, then wf finish again; later commands skipped" + GATE_RED_HELP) print(f"quick_gate: {len(cfg.quick_gate)} green", flush=True) with project_lock(root): p = load_project(args) if args.id not in p.doc.ids() and args.id in p.archived: print(f"{args.id} already done ({cfg.rel(cfg.archive)}): resuming with commit + merge", flush=True) else: done = argparse.Namespace(project=args.project, ids=[args.id], m=args.m, dry_run=False, finishing=True) cmd_done(done) sys.stdout.flush() resume = f" ({args.id} is done already: fix it, then rerun the same wf finish, it resumes)" commit = list(dict.fromkeys(map(str, [*paths, *(books if args.commit or not wt else [])]))) if commit: r = git_run(top, "add", "--", *map(str, paths)) if paths else None if r and r.returncode: raise Failure(f"git add failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}" + resume) msg = args.commit or f"{args.id} done" r = git_run(top, "commit", "-q", "-m", msg, "--", *commit) if r.returncode: raise Failure(f"commit failed: {(r.stderr.strip() or r.stdout.strip() or 'git error').splitlines()[-1]}" + resume) print("committed " + " ".join(os.path.relpath(c, top) for c in commit), flush=True) if not wt: sha = git_run(top, "rev-parse", "--short", "HEAD").stdout.strip() tool = getattr(args, "tool_commit", None) print(f"report: commit {sha}" + (f" tool {tool}" if tool else ""), flush=True) push_rc = 0 if not wt and not args.no_push: push_rc, lines = run_push(cfg, [top]) if lines: print("\n".join(lines), flush=True) if wt: ns = argparse.Namespace(project=args.project, m=None, no_push=args.no_push, private=priv, private_msg=args.commit) rc = cmd_merge(ns) if not rc and getattr(ns, "merged_sha", ""): tool = getattr(args, "tool_commit", None) books_sha = f" books {ns.books_sha}" if split else "" failed = f" push-failed {ns.push_rc}" if getattr(ns, "push_rc", 0) else "" print(f"report: commit {ns.merged_sha}{books_sha}" + (f" tool {tool}" if tool else "") + failed, flush=True) return rc return 0 def cmd_wip(args) -> int: """Wrap-up in one call (lane worktree): commit the given paths on the branch, note the state, status clear.""" start = Path(args.project or Path.cwd()).resolve() wt = config.linked_worktree(start) if not wt: raise Failure("wip runs inside a lane worktree") top = Path(git_run(start, "rev-parse", "--show-toplevel").stdout.strip() or start) p = load_project(args) id = p.resolve(args.id) btop = (config.git_top(p.cfg.root) or p.cfg.root).resolve() proot, priv, paths = p.cfg.root.resolve(), [], [] for f in args.paths: # split project: paths under the private root are committed there path = (start / f).resolve() (priv if btop != wt[2].resolve() and (path == proot or proot in path.parents) else paths).append(str(path)) if priv: r = git_run(btop, "add", "--", *priv) if r.returncode: raise Failure(f"git add failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") r = git_run(btop, "commit", "-q", "-m", f"{id} wip: {args.m}", "--", *priv) if r.returncode: raise Failure(f"commit failed: {(r.stderr.strip() or r.stdout.strip() or 'git error').splitlines()[-1]}") print("committed " + " ".join(os.path.relpath(c, btop) for c in priv), flush=True) if paths: r = git_run(top, "add", "--", *paths) if r.returncode: raise Failure(f"git add failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") r = git_run(top, "commit", "-q", "-m", args.commit or f"{args.id} WIP", "--", *paths) if r.returncode: raise Failure(f"commit failed: {(r.stderr.strip() or r.stdout.strip() or 'git error').splitlines()[-1]}") print("committed " + " ".join(os.path.relpath(c, top) for c in paths), flush=True) with project_lock(p.cfg.root): cmd_note(argparse.Namespace(project=args.project, id=id, line=args.m, dry_run=False)) cmd_status(argparse.Namespace(project=args.project, id=id, kind="clear", value=[], dry_run=False, clear_stale=False)) return 0 def _here(start: Path, cfg) -> Path: """Project folder to run config commands in: the worktree's copy of cfg.root, else cfg.root.""" wt = config.linked_worktree(start) if not wt: return cfg.root top, _, main = wt try: return top / cfg.root.relative_to(main) except ValueError: return top def _run_lines(lines: list[str], here: Path, cfg, fail: str) -> None: """Run shell lines in `here` with env WF_MAIN = main tree project folder; first nonzero → Failure.""" env = {**os.environ, "WF_MAIN": str(cfg.root)} for line in lines: print(f"$ {line}", flush=True) rc = subprocess.run(line, shell=True, cwd=here, env=env).returncode if rc: raise Failure(fail.format(line=line, rc=rc)) def cmd_setup(args) -> int: """Run workflow.toml worktree_setup in this lane worktree (cwd = its project folder, env WF_MAIN).""" start = Path(args.project or Path.cwd()).resolve() wt = config.linked_worktree(start) if not wt: raise Failure("setup runs inside a linked git worktree (lane worktree)") cfg = config.load_at(start) if not cfg.worktree_setup: print(f"no worktree_setup in {config.NAME}: nothing to do") return 0 _run_lines(cfg.worktree_setup, _here(start, cfg), cfg, "worktree_setup '{line}' failed (exit {rc}): later commands skipped") n = len(cfg.worktree_setup) print(f"worktree_setup: {n} command{'s' * (n != 1)} ok in {os.path.relpath(wt[0], cfg.root)}") return 0 def cmd_start(args) -> int: """Worker setup in one call (main tree): worktree + branch (clean check), worktree_setup, status progress, then wf ctx and a ready verify && wf finish line.""" start = Path(args.project or Path.cwd()).resolve() if config.linked_worktree(start): raise Failure("start runs in the main tree (it creates the worktree)") p = load_project(args) id = p.resolve(args.id) main = Path(git_run(start, "rev-parse", "--show-toplevel").stdout.strip() or p.cfg.root) repo = config.code_main(p.cfg) or main # split project: worktree + branch in the code repo master = config.git_branch(repo / ".git") or "master" wt = (Path.cwd() / args.worktree).resolve() shown = p.cfg.rel(wt) b = args.branch def git_ok(cwd, *a): r = git_run(cwd, *a) if r.returncode: raise Failure(f"git {a[0]} failed: {(r.stderr.strip() or r.stdout.strip() or 'git error').splitlines()[-1]}") return r.stdout has_branch = not git_run(repo, "rev-parse", "--verify", "-q", f"refs/heads/{b}").returncode extra = [] if not wt.exists(): if has_branch: git_ok(repo, "worktree", "add", "-q", str(wt), b) how = f"new, existing branch {b}: earlier WIP, read the notes" else: git_ok(repo, "worktree", "add", "-q", str(wt), "-b", b, master) how = f"new, branch {b} from {master}" else: own = git_run(wt, "rev-parse", "--path-format=absolute", "--git-common-dir").stdout.strip() want = git_run(repo, "rev-parse", "--path-format=absolute", "--git-common-dir").stdout.strip() if not own or Path(own).resolve() != Path(want).resolve(): raise Failure(f"{shown} is not a worktree of {repo}" + (f" (it belongs to {Path(own).resolve().parent})" if own else "") + f": pass a worktree path of {repo} (e.g. {repo}/.worktrees/<lane>)") dirty = git_run(wt, "status", "--porcelain").stdout.rstrip("\n") cur = config.git_branch(Path(git_run(wt, "rev-parse", "--absolute-git-dir").stdout.strip())) if dirty and not args.recovery: raise Failure(f"worktree {shown} has uncommitted changes: " + " ".join(l.strip() for l in dirty.splitlines()) + " (hand back, or --recovery when a dead worker left them)") if dirty and cur != b: raise Failure(f"worktree {shown} has uncommitted changes on another branch ({cur or 'detached'}), " f"not {b}: hand back") if cur == b: how = f"on {b}" elif has_branch: git_ok(wt, "switch", "-q", b) how = f"switched to existing branch {b}: earlier WIP, read the notes" else: git_ok(wt, "switch", "-q", "-c", b, master) how = f"switched to new branch {b} from {master}" if args.recovery: log = git_run(wt, "log", "--oneline", f"{master}..HEAD").stdout.strip() extra = [f"recovery: uncommitted:\n{dirty}" if dirty else "recovery: uncommitted: none", f"recovery: commits {master}..HEAD:" + (f"\n{log}" if log else " none")] print(f"worktree: {shown} ({how})", *extra, sep="\n", flush=True) if repo != main: (wt / config.WF_HOME).write_text(f"{p.cfg.root}\n") try: here = wt / p.cfg.root.relative_to(main) except ValueError: here = wt cmd_setup(argparse.Namespace(project=str(here))) sys.stdout.flush() with project_lock(p.cfg.root), contextlib.redirect_stdout(io.StringIO()): cmd_status(argparse.Namespace(project=args.project, id=id, kind="progress", value=[b], dry_run=False, clear_stale=False)) print(flush=True) cmd_ctx(argparse.Namespace(project=args.project, target=id)) steps = [f"cd {here}", *(f"({v})" for v in p.cfg.verify), f'python3 {HERE / "wf.py"} finish {id} -m "<entry>" --commit "<msg + footer>" <paths>'] print("\nFinish (after the work; fill in the quoted parts and the paths):\n " + " && ".join(steps)) return 0 def cmd_gate(args) -> int: """Run workflow.toml quick_gate (fast regression check) in this tree's project folder, env WF_MAIN.""" start = Path(args.project or Path.cwd()).resolve() cfg = config.load_at(start) if not cfg.quick_gate: print(f"no quick_gate in {config.NAME}: nothing to do") return 0 _run_lines(cfg.quick_gate, _here(start, cfg), cfg, "quick_gate '{line}' red (exit {rc}): fix it before wf done; later commands skipped" + GATE_RED_HELP) n = len(cfg.quick_gate) print(f"quick_gate: {n} command{'s' * (n != 1)} green") return 0 # ------------------------------------------------------------------ orchestrator GO_ON = ("done", "done+gate-red", "sliced") # outcomes after which the lane picks its next task PARK = ("handback", "awaiting", "needs-owner", "post-check-red") # block only that task (park_task), lane keeps picking def orch_main(args) -> tuple[Path, Project]: """(main tree, project) for wf orch; refuses a linked worktree.""" start = Path(args.project or Path.cwd()).resolve() if config.linked_worktree(start): raise Failure("orch runs in the main tree (the orchestrator's)") p = load_project(args) main = Path(git_run(start, "rev-parse", "--show-toplevel").stdout.strip() or p.cfg.root) return main, p def orch_record(cfg: config.Config, id: str) -> dict: try: r = json.loads((cfg.root / ".wf" / "orch" / f"{id}.json").read_text()) except (OSError, ValueError): return {} return r if isinstance(r, dict) else {} def live_session_cwds() -> list[Path]: """cwd of every live Claude Code session (~/.claude/sessions/*.json).""" out = [] for f in (claude_projects().parent / "sessions").glob("*.json"): try: s = json.loads(f.read_text()) os.kill(int(s["pid"]), 0) out.append(Path(s["cwd"]).resolve()) except (OSError, ValueError, KeyError, TypeError): continue return out def worktree_free(path: Path, branch: str, busy: set[Path], cwds: list[Path]) -> bool: """No live session in it, no other picked worker on it, and clean on a detached HEAD (or on branch = recovery).""" if path in busy: return False if not path.exists(): return True if any(c == path or path in c.parents for c in cwds): return False cur = config.git_branch(Path(git_run(path, "rev-parse", "--absolute-git-dir").stdout.strip() or path / ".git")) if cur == branch: return True return not cur and not git_run(path, "status", "--porcelain").stdout.strip() def lane_worktree(main: Path, p: Project, lane: str, id: str) -> Path: """First free main/.worktrees/<lane>[-n] for branch <lane>/<id> (worktree_free; other picked workers' trees busy).""" cfg, busy = p.cfg, set() for f in (cfg.root / ".wf" / "orch").glob("*.json"): r = orch_record(cfg, f.stem) other = p.doc.item(f.stem) if f.stem in p.doc.ids() else None if f.stem != id and r.get("worktree") and other and (other.status or "").startswith("in progress"): busy.add(Path(r["worktree"])) cwds = live_session_cwds() base = config.code_main(cfg) or main n = 1 while not worktree_free(wt := base / ".worktrees" / (lane if n == 1 else f"{lane}-{n}"), f"{lane}/{id}", busy, cwds): n += 1 return wt CLOUD = "cloud" # virtual lane: wf cloud send instead of a local worker (spec cloud-lane §4.6) def cloud_pick(main: Path, p: Project, id: str | None) -> int: """wf orch pick cloud: ledger check, first fitting task across the lanes (cloud.pick_key) -> wf cloud send.""" import wf_cloud from wflib import cloud cfg = p.cfg if not cfg.cloud: print(f"stop lane {CLOUD}: none fit (project not opted in: workflow.toml cloud = true)") return 0 with wf_cloud.locked(wf_cloud.state_dir()) as led: run, cap, bal, res = cloud.running(led), led["max_parallel"], cloud.balance(led), led["reserve_per_task"] if run >= cap: print(f"stop lane {CLOUD}: max parallel ({run} running >= max_parallel {cap})") return 0 if bal < res: print(f"stop lane {CLOUD}: ledger (balance ${bal:.2f} < reserve ${res:.2f})") return 0 if id: item = p.doc.item(p.resolve(id)) if why := cloud.fit(item, cfg): print(f"stop lane {CLOUD}: none fit ({item.id}: {why})") return 0 else: ready, seen = runner_ids(p), [] for l in cfg.lanes: seen += [i for i in lanes.ranked(p.doc, p.archived, cfg.lanes, cfg.slice_above, l.name) if i not in seen and i.id in ready and not (i.status or "").startswith("in progress")] fits = sorted((i for i in seen if cloud.fit(i, cfg) is None), key=cloud.pick_key) if not fits: print(f"stop lane {CLOUD}: none fit (wf cloud fit rules: wf set ID --cloud yes skips the regex list, opus only)") return 0 item = fits[0] buf = io.StringIO() with contextlib.redirect_stdout(buf): code = wf_cloud.main(["send", item.id, "--project", str(cfg.root)]) sent = buf.getvalue().strip() if code == 3: print(f"stop lane {CLOUD}: ledger (send refused, see above)") return 0 if code: if sent: print(sent) raise Failure(f"wf cloud send {item.id} failed (exit {code}): lane {CLOUD} stops, tell the owner") rec = json.loads(wf_cloud.record_path(cfg, item.id).read_text()) (wf_folder(cfg, "orch") / f"{item.id}.json").write_text(json.dumps( {"id": item.id, "lane": CLOUD, "task_lane": rec["lane"], "model": item.model, "effort": item.effort, "sid": rec["sid"], "branch": f"{rec['lane']}/{item.id}", "at": int(datetime.datetime.now().timestamp())}) + "\n") print(f"pick: {item.id} (lane {CLOUD}, model {item.model}, effort {item.effort or '-'}) · {sent.splitlines()[-1]}" f" · claimed (in progress: cloud:{rec['sid']})") print(f"agent: none (cloud session). On each wake (>= 10 min apart): wf cloud pull --all; per ended task " f"'<id>: <state> …' → wf orch post <id> {CLOUD} --result <state> [--commit <sha>]") return 0 def orch_pick(main: Path, p: Project, lane: str, id: str | None = None, recovery: str | None = None) -> int: cfg = p.cfg stop = cfg.root / "out" / "wf-batch.stop" if stop.exists() and not recovery: print(f"stop: {cfg.rel(stop)} exists: spawn nothing (let running workers finish)") return 0 if lane == CLOUD: return cloud_pick(main, p, id) if id: id = p.resolve(id) else: ready = runner_ids(p) order = [i for i in lanes.ranked(p.doc, p.archived, cfg.lanes, cfg.slice_above, lane) if i.id in ready and not (i.status or "").startswith("in progress")] if not order: print(f"none: lane {lane} has no runner-ready task (wf list --runner --lane {lane}; wf lanes)") return 0 id = order[0].id item = p.doc.item(id) branch = f"{lane}/{id}" wt = lane_worktree(main, p, lane, id) with project_lock(cfg.root), contextlib.redirect_stdout(io.StringIO()): cmd_status(argparse.Namespace(project=str(cfg.root), id=id, kind="progress", value=["worker"], dry_run=False, clear_stale=False)) model = item.model (wf_folder(cfg, "orch") / f"{id}.json").write_text(json.dumps( {"id": id, "lane": lane, "model": model, "effort": item.effort, "worktree": str(wt), "branch": branch, "at": int(datetime.datetime.now().timestamp())}) + "\n") kind = ", slice job: the worker only slices" if lanes.is_slice_job(item, cfg.slice_above) else "" print(f"pick: {id} (lane {lane}, model {model}, effort {item.effort or '-'}{kind}) · claimed (in progress: worker)") print(f"agent: subagent_type wf-worker, model {model}, no isolation; prompt:") print(f"Task: {id} Lane: {lane} Model: {model}\nMain tree: {main} Worktree: {wt} Branch: {branch}" + (f"\nRecovery: {recovery}" if recovery else "") + "\nFinal message: the 4 report lines only.") return 0 def cost_line(root: Path, agent: str, task: str, outcome: str, effort: str | None, lane: str, dur: int | None) -> str: """Append one usage line for a worker to <root>/out/wf-cost.log; returns it.""" path = usage_transcripts(argparse.Namespace(agent=agent, session=None))[0][1] now = datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") line = usage.log_line(now, root.name, task, effort or "-", outcome, path.stem.removeprefix("agent-"), usage.parse(read(path)), lane=lane, dur=dur) (root / "out").mkdir(exist_ok=True) with open(root / "out" / "wf-cost.log", "a") as f: f.write(line + "\n") return line def archived_meta(cfg: config.Config, id: str) -> tuple[str | None, str | None]: """(model, effort) of a task already archived: its block in TASKS.md just before the commit that removed it.""" rel = str(cfg.tasks.relative_to(cfg.root)) r = git_run(cfg.root, "log", "-n1", "--format=%H", f"-S**{id}**", "--", rel) if r.returncode or not r.stdout.strip(): return None, None old = git_run(cfg.root, "show", f"{r.stdout.strip()}^:{rel}") try: it = tasks.parse(old.stdout).item(id) if old.returncode == 0 and id in tasks.parse(old.stdout).ids() else None except Exception: it = None return (it.model, it.effort) if it else (None, None) ARCHIVED = "archived, nothing to mark" def park_task(root: Path, id: str, final: str, words: list[str]) -> str: """Block only this task after a handback/awaiting/needs-owner/post-check-red: needs-owner → After: its open Needs-human item (status cleared); awaiting → blocked on its open a-id; else Sessions: owner (+ status cleared on a handback) and a note. Returns what it did (alert text).""" p = load_project(argparse.Namespace(project=str(root)), write=True) if id not in p.doc.ids(): return ARCHIVED item = p.doc.item(id) why = " ".join(words[1:]) or "-" open_a = {i.id for s in p.doc.sections if s.key == "awaiting" for i in s.items} a = next((w for w in words[1:] if w.strip("[]") in open_a), None) if final == "awaiting" else None open_h = {i.id for i in p.doc.section("human").items} h = next((w.strip("[]") for w in words[1:] if w.strip("[]") in open_h), None) \ if final == "needs-owner" else None if h: # owner-only step: After: it, picked again once the owner's wf done archives it if h not in item.after: tasks.set_fields(p.doc, id, after=item.after + [h]) if item.status: tasks.set_status(p.doc, id, None) done = f"after {h}" elif (item.status or "").startswith("blocked:"): done = f"already {item.status}" elif a: tasks.set_status(p.doc, id, f"blocked: [[{a.strip('[]')}]]") done = f"blocked on {a.strip('[]')}" else: tasks.set_fields(p.doc, id, sessions="owner") done = "Sessions: owner" if final == "handback" and item.status: tasks.set_status(p.doc, id, None) done += ", status cleared" tasks.add_note(p.doc, id, f"orch {datetime.date.today().isoformat()}: {final} ({why}) → {done}; " "lane kept picking, owner decides") p.save(False) if done.endswith("status cleared") or h: unclaim(p.cfg, [id]) return done def cloud_park(root: Path, id: str, final: str, words: list[str]) -> str: """Lane cloud after a handback/lost/awaiting/needs-owner (wf cloud pull already noted/blocked it): awaiting/needs-owner → park_task; else Cloud: no (the cloud lane never re-sends it; local lanes may retry with pull's Recovery note) + status cleared + note. Returns what it did (alert text).""" if final in ("awaiting", "needs-owner"): return park_task(root, id, final, words) p = load_project(argparse.Namespace(project=str(root)), write=True) if id not in p.doc.ids(): return ARCHIVED item, done = p.doc.item(id), "Cloud: no" if item.cloud != "no": tasks.set_fields(p.doc, id, cloud="no") if (item.status or "").startswith("in progress"): tasks.set_status(p.doc, id, None) done += ", status cleared" tasks.add_note(p.doc, id, f"orch {datetime.date.today().isoformat()}: cloud {final} " f"({' '.join(words[1:]) or '-'}) → {done}; lane cloud kept picking, local lanes may retry") p.save(False) unclaim(p.cfg, [id]) return done + " (local lanes may retry)" def orch_post(main: Path, p: Project, args) -> int: cfg, id, lane = p.cfg, args.id, args.lane rec = orch_record(cfg, id) repo = config.code_main(cfg) or main wt = Path(rec.get("worktree") or repo / ".worktrees" / lane) item = p.doc.item(id) if id in p.doc.ids() else None if item is None and id not in p.archived: raise Failure(tasks.unknown_id(id, p.doc.ids())) line = (args.result or "").strip() if line[:7].lower() == "result:": line = line[7:].strip() # literal report line 'result: done' words = line.split() outcome = words[0] if words else ("done" if id in p.archived else "no-report") if outcome == "done" and item is not None and id not in p.archived and tasks.open_slices(p.doc, id): outcome = "sliced" # slice job: worker said done but the task is open with slices master = config.git_branch(repo / ".git") or "master" files = [str(f.relative_to(main)) for f in (cfg.tasks, cfg.archive)] problems, out = [], [] split = repo.resolve() != main.resolve() def report_commit() -> str: # after any merge below; split: a private sha → the code repo's master c = args.commit or (git_run(repo, "rev-parse", "--short", master).stdout.strip() if outcome.startswith("done") else "-") if split and c != "-" and git_run(repo, "cat-file", "-e", f"{c}^{{commit}}").returncode: fixed = git_run(repo, "rev-parse", "--short", master).stdout.strip() out.append(f"commit {c} not in code repo {repo} (a private sha?): using its {master} {fixed}") c = fixed return c commit, parked = None, None raised = bool(item and rec.get("model") and not outcome.startswith("done") and tasks.MODELS.index(item.model) > tasks.MODELS.index(rec["model"])) with project_lock(cfg.root): if outcome.startswith("done"): if id not in p.archived: problems.append(f"no archive line for {id}") if (wt / ".git").exists(): dirty = git_run(wt, "status", "--porcelain").stdout.strip() if dirty: problems.append(f"worktree {wt} dirty: " + " ".join(l.strip() for l in dirty.splitlines())) elif git_run(wt, "merge-base", "--is-ancestor", "HEAD", master).returncode: buf = io.StringIO() try: with contextlib.redirect_stdout(buf): cmd_merge(argparse.Namespace(project=str(wt), m=None, no_push=args.no_push)) out.append(f"merged worktree HEAD: {buf.getvalue().strip().splitlines()[-1]}") except Failure as e: problems.append(f"wf merge in {wt}: {e}") branch = rec.get("branch") or f"{lane}/{id}" if not git_run(repo, "rev-parse", "--verify", "-q", f"refs/heads/{branch}").returncode: problems.append(f"branch {branch} still there") errors = checks.check(config.load(cfg.root))[0] if errors: problems.append(f"wf check: {len(errors)} errors: {errors[0]}") commit = report_commit() if not problems and git_run(main, "status", "--porcelain", "--", *files).stdout.strip(): msg = f"{id} done" + (f" (code {commit})" if split and commit != "-" else " (orchestrator)") r = git_run(main, "commit", "-q", "-m", msg, "--", *files) out.append("committed leftover " + " ".join(files) if not r.returncode else f"leftover commit failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") if problems or not outcome.startswith("done"): if lane != CLOUD and (problems or outcome in PARK and not raised): parked = park_task(cfg.root, id, "post-check-red" if problems else outcome, words) elif lane == CLOUD and not problems and outcome in PARK + ("lost",) and not raised: # no re-send loop parked = cloud_park(cfg.root, id, outcome, words) if (not problems or parked and parked != ARCHIVED) \ and git_run(main, "status", "--porcelain", "--", *files).stdout.strip(): r = git_run(main, "commit", "-q", "-m", f"{id} {'post-check-red' if problems else outcome}" " (orchestrator)", "--", *files) out.append("committed leftover " + " ".join(files) if not r.returncode else f"leftover commit failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") final = "post-check-red" if problems else ("model-raised" if raised else outcome) if outcome.startswith("done") and not problems and (cfg.root / ".wf" / "push-failed").is_file(): final = "push-failed" am, ae = (None, None) if item or (rec.get("model") and rec.get("effort")) else archived_meta(cfg, id) model = rec.get("model") or (item.model if item else am) or "?" dur = args.duration if args.duration else ( # None/0 (orchestrator lacked duration_ms) → since the pick int(datetime.datetime.now().timestamp()) - rec["at"] if isinstance(rec.get("at"), int) else None) commit = commit or report_commit() now = datetime.datetime.now().isoformat(timespec="seconds") (cfg.root / "out").mkdir(exist_ok=True) with open(cfg.root / "out" / "wf-orch.log", "a") as f: f.write(f"{now} {lane} {model} {id} {final} {commit} {fmt_dur(dur) if dur is not None else '-'}" + (f" ({'; '.join(problems)})" if problems else "") + "\n") if args.agent: try: out.append("cost: " + cost_line(cfg.root, args.agent, id, final, rec.get("effort") or (item and item.effort) or ae, lane, dur)) except Failure as e: out.append(f"cost: not logged ({e})") if final != "post-check-red": # kept: the re-post after the fix still knows the pick time (cfg.root / ".wf" / "orch" / f"{id}.json").unlink(missing_ok=True) print(f"post: {id} {final}" + "".join(f"\n {l}" for l in problems + out)) if outcome == "done+gate-red" and len(words) >= 3: q = load_project(args) fix = q.doc.item(words[2]) if words[2] in q.doc.ids() else None if fix and not fix.runner_ready: print(f"fix {fix.id} not runner-ready: add its Done/Model (wf set {fix.id} --done … --model …), " f"then wf orch pick {lane} --id {fix.id}") return 0 if parked: print(f"alert: {id} {final} → tell the owner ({parked}); lane {lane} keeps picking") elif final not in GO_ON and final != "model-raised": print(f"stop lane {lane}: {final} → tell the owner" + (f" (wf push in {cfg.root})" if final == "push-failed" else "") + (f" (crash/no report: one fresh worker: wf orch pick {lane} --id {id} --recovery \"<why>\", then stop)" if final == "no-report" else "")) return 0 if args.no_pick: return 0 print() return orch_pick(main, load_project(args), lane) def cmd_orch(args) -> int: main, p = orch_main(args) if args.lane != CLOUD or CLOUD in lane_names(p.cfg): check_lane(p.cfg, args.lane) if args.action == "pick": return orch_pick(main, p, args.lane, args.id, args.recovery) return orch_post(main, p, args) def change(args, action) -> int: """Run one edit on one item, save, print its new header.""" p = load_project(args, write=True) note = action(p) p.save(args.dry_run) if not args.dry_run: print(note if note else p.doc.item(args.id).header()) return 0 def cmd_prio(args) -> int: return change(args, lambda p: tasks.set_prio(p.doc, args.id, args.prio)) def cmd_move(args) -> int: targets = [t for t in (args.section, args.before, args.after) if t] if len(targets) != 1: raise Usage("move: give a section, or --before ID, or --after ID") if args.section: return change(args, lambda p: tasks.move_to(p.doc, args.id, args.section)) return change(args, lambda p: tasks.move_rel(p.doc, args.id, args.before or args.after, before=bool(args.before), force=args.force)) def clear_stale(args) -> int: p = load_project(args, write=True) items = stale(p) for i in items: tasks.set_status(p.doc, i.id, None) p.save(args.dry_run) if not args.dry_run: unclaim(p.cfg, [i.id for i in items]) for i in items: print(f"cleared {i.id}") return 0 def cmd_status(args) -> int: if args.clear_stale: if args.id or args.kind or args.value: raise Usage("status: --clear-stale takes no id/kind/value") return clear_stale(args) if not args.id or not args.kind: raise Usage("status: ID and progress NOTE|blocked A-ID|clear (or --clear-stale)") if args.kind == "clear": if args.value: raise Usage("status: clear takes no value") status = None elif not args.value: raise Usage("status: progress needs a note, blocked needs an a-id") elif args.kind == "progress": status = "in progress: " + " ".join(args.value) else: status = f"blocked: [[{args.value[0]}]]" code = change(args, lambda p: tasks.set_status(p.doc, args.id, status)) if not args.dry_run: proj = load_project(args) cfg = proj.cfg if args.kind == "progress": claim(cfg, args.id, lanes.lane_of(proj.doc.item(args.id), cfg.lanes, cfg.slice_above)) else: unclaim(cfg, [args.id]) return code def cmd_set(args) -> int: if args.interactive is not None: old_spelling("--interactive", "--sessions owner") if args.sessions is None: args.sessions = "owner" if args.interactive == "yes" else "" fields = dict( title=args.title, effort=args.effort, after=None if args.after is None else split_list(args.after), refs=None if args.ref is None else split_list(args.ref), model=args.model, sessions=args.sessions, done=args.done, cloud=args.cloud) if args.model not in (None, "", *tasks.MODELS): raise Usage(f"set: --model {args.model} (want {', '.join(tasks.MODELS)}, or \"\" to remove)") if args.sessions not in (None, "", *tasks.SESSIONS): raise Usage(f"set: --sessions {args.sessions} (want {', '.join(tasks.SESSIONS)}, or \"\" to remove)") if args.cloud not in (None, "", *tasks.CLOUDS): raise Usage(f"set: --cloud {args.cloud} (want {', '.join(tasks.CLOUDS)}, or \"\" to remove)") if all(v is None for v in fields.values()): raise Usage("set: nothing to set (--title --effort --after --ref --model --sessions --cloud --done)") return change(args, lambda p: tasks.set_fields(p.doc, args.id, **fields)) def cmd_rename(args) -> int: def action(p): n = tasks.rename(p.doc, args.id, args.new, p.archived) return f"renamed: {args.id} → {args.new} ({n} link{'' if n == 1 else 's'})" return change(args, action) def cmd_note(args) -> int: return change(args, lambda p: tasks.add_note(p.doc, args.id, " ".join(args.line.split()))) def cmd_body(args) -> int: if args.text: raise Failure("body text goes on stdin (wf body ID <<'EOF' … EOF), not as an argument") lines = stdin_lines() return change(args, lambda p: tasks.set_body(p.doc, args.id, lines)) def cmd_tick(args) -> int: return change(args, lambda p: "ticked: " + tasks.tick(p.doc, args.id, args.which)) # ------------------------------------------------------------------ other commands def wf_version() -> str: try: out = subprocess.run(["git", "-C", str(HERE), "rev-parse", "--short=7", "HEAD"], capture_output=True, text=True, timeout=10) return out.stdout.strip() if out.returncode == 0 and out.stdout.strip() else "0000000" except (OSError, subprocess.SubprocessError): return "0000000" def cmd_report(args) -> int: text = " ".join(args.text.split()) if not text: raise Usage("report: say what happened") start = Path(args.project) if args.project else Path.cwd() try: project = config.find_root(start).name except config.ConfigError: project = start.resolve().name cmd = f" (cmd: {' '.join(args.cmd.split())})" if args.cmd else "" line = f"- {datetime.date.today().isoformat()} {project} {args.kind}: {text}{cmd} @{wf_version()}\n" fd = os.open(inbox_path(), os.O_WRONLY | os.O_APPEND | os.O_CREAT, 0o644) try: os.write(fd, line.encode("utf-8")) finally: os.close(fd) print("reported (workflow inbox); carry on") return 0 def cmd_init(args) -> int: root = (Path(args.project) if args.project else Path.cwd()).resolve() if (root / config.NAME).exists(): raise Failure(f"{root}/{config.NAME} exists already") templates = HERE / "templates" for target, source in ((config.NAME, "workflow.toml"), ("TASKS.md", "TASKS.md"), ("tasks/archive.md", "archive.md"), ("CLAUDE.md", "CLAUDE.md")): path = root / target if path.exists(): print(f"kept {target} (exists" + ("; old format → wf migrate)" if target == "TASKS.md" else ")")) continue path.parent.mkdir(parents=True, exist_ok=True) path.write_text(read(templates / source).replace("{name}", root.name), encoding="utf-8") print(f"wrote {target}") return 0 def cmd_migrate(args) -> int: from wflib import migrate p = load_project(args) cfg = p.cfg new, report = migrate.migrate(p.tasks_text, taken=p.archived) bump = cfg.format < config.FORMAT if new == p.tasks_text and not bump: print("nothing to migrate") return 0 name = cfg.rel(cfg.tasks) sys.stdout.writelines(difflib.unified_diff( p.tasks_text.replace("\r\n", "\n").splitlines(keepends=True), new.replace("\r\n", "\n").splitlines(keepends=True), name, f"{name} (new)")) if report.ids: print("\nids:") for title, id in report.ids.items(): print(f" {id} {title}") if report.notes: print("\nnotes:") for note in report.notes: print(f" {note}") errors = checks.check(cfg, tasks_text=new, slow=False)[0] for e in errors: print(f"ERROR: {e}") print(f"check after migrate: {len(errors)} errors") if not args.write: print("dry run: nothing written (wf migrate --write)") return 0 wrote = [] if new != p.tasks_text: write_if_unchanged(cfg.tasks, new, p.tasks_stamp) wrote.append(name) if bump: toml = cfg.root / config.NAME text = read(toml) changed = re.sub(r"(?m)^format(\s*)=(\s*)\d+", rf"format\g<1>=\g<2>{config.FORMAT}", text, count=1) write_if_unchanged(toml, changed, stamp(toml)) wrote.append(f"{config.NAME} (format = {config.FORMAT})") print("wrote " + ", ".join(wrote)) return 0 # ------------------------------------------------------------------ main def parser() -> argparse.ArgumentParser: ap = argparse.ArgumentParser(prog="wf", description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) ap.add_argument("--project", metavar="DIR", help="project folder (default: found from the working directory)") sub = ap.add_subparsers(dest="command", metavar="COMMAND") sections = list(tasks.SECTIONS) def cmd(name, func, help, write=False): sp = sub.add_parser(name, help=help, description=help) sp.set_defaults(func=func, locks=write) if write: sp.add_argument("--dry-run", action="store_true", help="print the diff, write nothing") return sp sp = cmd("next", cmd_next, "cold start: awaiting, needs human, plans in flight, next task with its refs") sp.add_argument("--brief", "-b", action="store_true", help="the task only") sp.add_argument("--lane", help="lane name (wf lanes); none = all lanes") sp.add_argument("--as", dest="as_", choices=tasks.MODELS, help="your model: takes tasks with Model ≤ it (default haiku)") sp.add_argument("--owner", action="store_true", help="the owner is present: Sessions: owner tasks pickable") sp = cmd("areas", cmd_areas, "area notes: anchors missing, commits since Checked (stale = missing or " "≥ area_stale_commits); no NAME → also uncovered: folders the diff since master touches outside " "every Paths (area_ignore skipped); --mark NAME stamps HEAD", write=True) sp.add_argument("name", nargs="?", help="one area only") sp.add_argument("--mark", metavar="NAME", help="set the area's Checked to HEAD (after a refresh)") sp = cmd("lanes", cmd_lanes, "lanes: pickable/waiting per lane, the session of each " "(--lane/--as registers yours; neither lane = all)") sp.add_argument("--lane", help="your lane (wf lanes lists them)") sp.add_argument("--as", dest="as_", choices=tasks.MODELS, help="your model (default haiku)") sp.add_argument("--unregister", action="store_true", help="drop this session's lane record (orchestrators hold none)") sp.add_argument("--wait", type=float, metavar="SECS", help="block until any lane has a pickable task (exit 0, prints the lanes line) or SECS pass " "(exit 1, prints the lanes line); out/wf-batch.stop exists (wf batch --stop) → exit 2 at once; " "polls every 30 s (env WF_LANES_POLL)") sp = cmd("list", cmd_list, "one line per item: id, priority, effort, status, model, lane, title") sp.add_argument("-s", "--section", choices=[*sections, "all"], default="pending") sp.add_argument("-p", "--prio", type=int, choices=range(4), help="this priority or higher") sp.add_argument("--ready", action="store_true", help="pickable now") sp.add_argument("--runner", action="store_true", help="--ready + Done line, no Sessions: owner, Done not about owner/confirming, no open slices") sp.add_argument("--blocked", action="store_true") sp.add_argument("--progress", action="store_true") sp.add_argument("--stale", action="store_true", help="in progress with no live claim (dead or no holder)") sp.add_argument("--model", choices=tasks.MODELS, help="tasks with this Model line (none = opus)") sp.add_argument("--lane", help="tasks of this lane (by effort); with --runner: in pick order") sp.add_argument("--ref", metavar="PATH[#ANCHOR]") sp.add_argument("-n", type=int, metavar="N") sp = cmd("show", cmd_show, "the raw item(s)") sp.add_argument("ids", nargs="+", metavar="ID") sp = cmd("ctx", cmd_ctx, "item with refs resolved, its relations and the areas it names; or a doc section and its tasks") sp.add_argument("target", metavar="ID|PATH#ANCHOR") sp = cmd("search", cmd_search, "ranked search over tasks, archive and docs") sp.add_argument("words", nargs="+") sp.add_argument("--tasks", action="store_true") sp.add_argument("--archive", action="store_true") sp.add_argument("--docs", action="store_true") sp.add_argument("-n", type=int, default=15, metavar="N") sp = cmd("log", cmd_log, "newest archive lines") sp.add_argument("words", nargs="*") sp.add_argument("-n", type=int, default=10, metavar="N") sp = cmd("usage", cmd_usage, "tokens and API-price $ per agent of a session, from Claude Code transcripts") sp.add_argument("--session", metavar="ID", help="session id or prefix (default: this session, else newest here)") sp.add_argument("--agent", metavar="ID", help="one subagent only (agent id from its completion notice)") sp.add_argument("--since", metavar="TIME", help="ISO UTC time, or 90m / 2h / 1d ago") sp.add_argument("--log", nargs=2, metavar=("TASK", "OUTCOME"), help="with --agent: append one line to <project>/out/wf-cost.log (orchestrator, per worker)") sp.add_argument("--effort", choices=tasks.EFFORTS, metavar="|".join(tasks.EFFORTS), help="the task's estimate (not reasoning effort), for --log") sp.add_argument("--duration", type=int, metavar="S", help="agent wall time in seconds (duration_ms / 1000), for --log") sp.add_argument("--lane", metavar="L", help="lane of the logged worker, for --log (default: its model)") sp.add_argument("--explore", action="store_true", help="exploring cost per subagent: pct billed input exploring / before the first edit, " "result tokens per tool, files read by >= 2 agents (same --session / --since)") sp.add_argument("--report", action="store_true", help="per lane and effort: n, done, $ (median, total, per done), turns, duration, from every project's log") cmd("projects", cmd_projects, "every wf project: counts, errors, next task") cmd("check", cmd_check, "validate ids, links, refs, order") sp = cmd("add", cmd_add, "new item; `add -` reads a whole item block from stdin", write=True) sp.add_argument("title", metavar='"Title. Goal."|-') sp.add_argument("-p", "--prio", type=int, choices=range(4)) sp.add_argument("-e", "--effort", metavar="|".join(tasks.EFFORTS)) sp.add_argument("-s", "--section", choices=sections) sp.add_argument("--id", help="the id (default: from the title, or <parent>-N with --parent); " "a leading 'ID: ' in the text works too") sp.add_argument("--after", metavar="ID,…") sp.add_argument("--ref", metavar="REF,…") sp.add_argument("--parent", metavar="ID", help="add as the next slice of ID, After: the previous one") sp.add_argument("--interactive", action="store_true", help=argparse.SUPPRESS) sp.add_argument("--model", choices=tasks.MODELS, help="Model line (default: none = opus)") sp.add_argument("--done", metavar="TEXT", help="Done line (makes the task runner-pickable)") sp.add_argument("--cloud", choices=tasks.CLOUDS, help="Cloud line: yes = cloud lane may take it (opus only), no = never") sp.add_argument("--sessions", choices=tasks.SESSIONS, help="Sessions line: solo = no other live session, owner = owner present (default parallel)") sp.add_argument("-b", "--body", action="store_true", help="read body lines from stdin") sp = cmd("done", cmd_done, "finish task(s): remove, archive line, print verify + checklist", write=True) sp.add_argument("ids", nargs="+", metavar="ID") sp.add_argument("-m", metavar="ENTRY", help="archive entry, ≤2 lines") sp = cmd("prio", cmd_prio, "set priority, reposition", write=True) sp.add_argument("id") sp.add_argument("prio", type=int, choices=range(4)) sp = cmd("move", cmd_move, "move to a section, or before/after another item", write=True) sp.add_argument("id") sp.add_argument("section", nargs="?", choices=sections) sp.add_argument("--before", metavar="ID") sp.add_argument("--after", metavar="ID") sp.add_argument("--force", action="store_true", help="allow breaking priority order") sp = cmd("status", cmd_status, "progress NOTE | blocked A-ID | clear", write=True) sp.add_argument("id", nargs="?") sp.add_argument("kind", nargs="?", choices=["progress", "blocked", "clear"]) sp.add_argument("value", nargs="*") sp.add_argument("--clear-stale", action="store_true", help="clear in progress of tasks with no live claim") sp = cmd("set", cmd_set, "change header fields, Done:, Model:, After:, Ref:", write=True) sp.add_argument("id") sp.add_argument("--title", help="new title, goal kept; 'Title. Goal.' or a question replaces the text") sp.add_argument("--effort") sp.add_argument("--after", metavar="ID,…", help='"" removes the line') sp.add_argument("--ref", metavar="REF,…", help='"" removes the line') sp.add_argument("--interactive", choices=["yes", "no"], help=argparse.SUPPRESS) sp.add_argument("--sessions", metavar="|".join(tasks.SESSIONS), help='Sessions line; "" removes it (= parallel)') sp.add_argument("--model", metavar="|".join(tasks.MODELS), help='Model line; "" removes it (= opus)') sp.add_argument("--cloud", metavar="|".join(tasks.CLOUDS), help='Cloud line (yes: skip the fit regex list, opus only; no: never cloud); "" removes it') sp.add_argument("--done", metavar="TEXT", help='Done line (makes it runner-pickable); "" removes it') sp = cmd("rename", cmd_rename, "change an open item's id and every [[link]] to it", write=True) sp.add_argument("id") sp.add_argument("new", metavar="NEW") sp = cmd("note", cmd_note, "append one line to the body", write=True) sp.add_argument("id") sp.add_argument("line") sp = cmd("body", cmd_body, "replace the body with stdin; Model:, Sessions:, After:, Ref: lines kept", write=True) sp.formatter_class = argparse.RawDescriptionHelpFormatter sp.epilog = "example (body text goes on stdin, not as an argument):\n wf body ID <<'EOF'\n - Steps: …\n - Done: …\n EOF" sp.add_argument("id") sp.add_argument("text", nargs="*", help=argparse.SUPPRESS) sp = cmd("tick", cmd_tick, "check a `- [ ]` box by number or text", write=True) sp.add_argument("id") sp.add_argument("which", metavar="N|TEXT") sp = cmd("merge", cmd_merge, "in a lane worktree: rebase, ff-merge into master, commit TASKS/archive, " "detach, push home (under the project lock)") sp.set_defaults(locks=True) sp.add_argument("-m", help="bookkeeping commit message (default '<id> done', id from the branch)") sp.add_argument("--no-push", action="store_true", help="skip git push home") sp = cmd("push", cmd_push, "rerun a merge's push: workflow.toml push lines, else git push home") sp.set_defaults(locks=True) sp = cmd("finish", cmd_finish, "worker's last step in one call: quick_gate → done → commit PATHS (explicit) → " "wf merge (lane worktree; main tree: TASKS/archive go into the commit). Run verify first. " "Refuses before done if the gate is red, files outside PATHS are uncommitted, or a PATH is outside " "this repo / missing / unchanged; an id already archived resumes (commit + merge). " "Prints the done output (notify / checklist) and the merge lines, then 'report: commit <sha> [tool <sha>]' " "(lane worktree) = the sha for your report") sp.add_argument("id") sp.add_argument("-m", required=True, metavar="ENTRY", help="archive entry, ≤2 lines") sp.add_argument("--commit", metavar="MSG", help="commit message for PATHS (incl. footer lines)") sp.add_argument("paths", nargs="*", metavar="PATH", help="files/folders to commit (relative to cwd)") sp.add_argument("--no-push", action="store_true", help="skip git push home") sp.add_argument("--tool-commit", metavar="SHA", help="tool-repo sha merged by hand; added to the report line") sp = cmd("wip", cmd_wip, "wrap-up in one call (lane worktree): commit PATHS on the branch, note MSG " "(state + next step), status clear") sp.add_argument("id") sp.add_argument("-m", required=True, metavar="NOTE", help="state + next step, one line") sp.add_argument("--commit", metavar="MSG", help="commit message (default '<id> WIP')") sp.add_argument("paths", nargs="*", metavar="PATH", help="files/folders to commit (relative to cwd)") cmd("setup", cmd_setup, "in a lane worktree: run workflow.toml worktree_setup (cwd = worktree, env WF_MAIN = " "main tree project folder); provides git-ignored inputs, commands must be idempotent") sp = cmd("start", cmd_start, "worker setup in one call, in the main tree: worktree at PATH on BRANCH " "(new: from master or the existing branch; existing: must be clean, switched to BRANCH), " "worktree_setup, status progress BRANCH, then wf ctx and a ready verify && wf finish line") sp.add_argument("id") sp.add_argument("--worktree", required=True, metavar="PATH", help="worktree folder (relative to cwd)") sp.add_argument("--branch", required=True, metavar="BRANCH") sp.add_argument("--recovery", action="store_true", help="a dead worker's WIP: uncommitted changes on BRANCH kept, printed with its commits") cmd("gate", cmd_gate, "run workflow.toml quick_gate (fast regression check a worker runs before wf done; " "cwd = this tree's project folder, env WF_MAIN = main tree project folder); red = exit 1, nothing = exit 0") sp = cmd("report", cmd_report, "report a workflow problem or idea to the workflow inbox") sp.add_argument("text") sp.add_argument("--kind", choices=["bug", "idea", "friction"], default="bug") sp.add_argument("--cmd", metavar="COMMAND") cmd("init", cmd_init, "make this folder a wf project") sp = cmd("migrate", cmd_migrate, "convert a numbered TASKS.md to the id format (dry run unless --write)") sp.add_argument("--write", action="store_true") sp = cmd("orch", cmd_orch, "orchestrator, main tree: pick LANE (pick + claim + free worktree + agent prompt) · " "post ID LANE (post-check, merge if needed, orch log + cost line, next pick)") osub = sp.add_subparsers(dest="action", required=True) op = osub.add_parser("pick", help="first runner-ready task of LANE: claim it, choose a free worktree " "(.worktrees/LANE, -2, …), print the wf-worker prompt; 'none:'/'stop:' line = spawn nothing. " "LANE cloud: ledger check, first cloud-fitting task of any lane (opus first) -> wf cloud send; " "'stop lane cloud: ledger|none fit|max parallel'") op.add_argument("lane") op.add_argument("--id", help="this task instead of the pick (a P0 fix, a recovery)") op.add_argument("--recovery", metavar="WHY", help="prompt gets 'Recovery: WHY' (crashed worker, one retry)") op = osub.add_parser("post", help="after the worker's report: post-check (done: archive line, branch gone, " "worktree clean and in master, else wf merge; wf check), leftover commit otherwise, " "out/wf-orch.log line, out/wf-cost.log line (--agent); awaiting/needs-owner/handback/post-check-red park only the task (blocked / After: h-task / Sessions: owner (lane cloud: Cloud: no) + note, 'alert:' line); then the next pick or 'stop lane'") op.add_argument("id") op.add_argument("lane") op.add_argument("--result", metavar="LINE", help="the report's result line (default: done if archived, else no-report)") op.add_argument("--commit", metavar="SHA", help="the report's commit (default: master's sha when done)") op.add_argument("--agent", metavar="ID", help="agent id from the completion notice: cost line (wf usage --log)") op.add_argument("--duration", type=int, metavar="S", help="agent duration_ms / 1000 (default or 0: since the pick)") op.add_argument("--no-pick", "--no-next", dest="no_pick", action="store_true", help="no next pick, claims nothing (batch end, stop file, cross-project rollout)") op.add_argument("--no-push", action="store_true", help="the wf merge it may run skips git push home") cmd("res", None, "shared memory/CPU ledger for agent jobs (wf res -h)") cmd("cloud", None, "cloud-lane ledger (wf cloud -h)") cmd("prep", None, "prep tasks without Done: wf prep [N|all] = wf batch --prep (wf batch -h)") cmd("batch", None, "start an unattended batch orchestrator (claude -p) under wf res (wf batch -h)") return ap def main(argv: list[str]) -> int: if argv[:1] == ["cloud"]: import wf_cloud return wf_cloud.main(argv[1:]) if argv[:1] == ["res"]: import wf_res return wf_res.main(argv[1:]) if argv[:1] == ["prep"]: import wf_res return wf_res.prep_main(argv[1:]) if argv[:1] == ["batch"]: import wf_res return wf_res.batch_main(argv[1:]) ap = parser() args = ap.parse_args(argv) if not args.command: ap.print_usage(sys.stderr) return 2 try: if getattr(args, "locks", False): root = config.find_root(Path(args.project) if args.project else Path.cwd()) with project_lock(root): return args.func(args) return args.func(args) except Failure as e: print(f"wf: {e}", file=sys.stderr) return e.code except (tasks.TaskError, config.ConfigError) as e: print(f"wf: {e}", file=sys.stderr) return 1 except BrokenPipeError: return 0 if __name__ == "__main__": sys.exit(main(sys.argv[1:])) |