diff options
| author | godosa <godosa@godosa.eu> | 2026-10-07 16:57:32 +0200 |
|---|---|---|
| committer | godosa <godosa@godosa.eu> | 2026-10-07 16:57:32 +0200 |
| commit | 61a98e767d7c45fb00b17f969c6fb42d2cbaca48 (patch) | |
| tree | 284e1c453f28a35ce12ce9f84f2053031b7c1020 | |
| parent | f23d7d865eb042512ebabad87e89b35b11cc1c64 (diff) | |
| download | workflow-61a98e767d7c45fb00b17f969c6fb42d2cbaca48.tar.gz workflow-61a98e767d7c45fb00b17f969c6fb42d2cbaca48.zip | |
orch: one stuck task parks itself, lane keeps picking
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JEAjUkQRrCYX5MZhWxdtj2
| -rw-r--r-- | CHANGES.md | 1 | ||||
| -rw-r--r-- | docs/manual.md | 4 | ||||
| -rw-r--r-- | docs/orchestrator.md | 4 | ||||
| -rw-r--r-- | shared/skills/wf-orchestrate/SKILL.md | 4 | ||||
| -rw-r--r-- | shared/skills/wf-pilot/SKILL.md | 2 | ||||
| -rw-r--r-- | templates/batch-prompt.md | 2 | ||||
| -rw-r--r-- | tests/test_orch.py | 33 | ||||
| -rwxr-xr-x | wf.py | 58 |
8 files changed, 87 insertions, 21 deletions
@@ -1,4 +1,5 @@ # Changes (newest first) +- 2026-10-07 One stuck task no longer freezes its lane: `wf orch post` with awaiting / handback / post-check-red (local lanes; lane cloud unchanged: pull re-queues) parks only that task (blocked on the a-id named in the result, else `Sessions: owner` + an `orch …` note; handback also clears its claim), commits, prints `alert: <id> <outcome> → tell the owner (…); lane <lane> keeps picking`, then the lane's next pick. `stop lane` only for none / stop file / push-failed / wip / no report. Batch prompt, wf-orchestrate, wf-pilot, docs: keep picking, never wait on the owner while tasks are pickable. Projects: nothing. - 2026-10-07 `wf orch post` done: leftover TASKS/archive changes in the main tree are committed (`<id> done (code <sha>)` in a split project, else `<id> done (orchestrator)`); split: a `--commit` sha not in the code repo (stale private HEAD) is replaced by the code repo's master sha (noted). wf-worker: the report's commit = first sha of `report: commit …`, never a looked-up HEAD. Projects: nothing. - 2026-10-07 wf start refuses an existing worktree of another repo; .wf-home ignored by finish/merge stray checks. Projects: nothing. - 2026-10-07 split projects 6/6: `wf check` split rules (cloud error; warns: no push, `.wf-home` not git-ignored, merged code worktrees); docs. Projects: split projects: push key, .wf-home in the code repo's .gitignore (migration task in the workflow project). diff --git a/docs/manual.md b/docs/manual.md index c9f7bd1..3d50ade 100644 --- a/docs/manual.md +++ b/docs/manual.md @@ -122,8 +122,8 @@ switches them to worktree mode (one branch per task, `wf finish` merges it). One After `/clear`, type `/wf-orchestrate` (`go` = start the proposed run without asking, `batch N` = prepare an unattended batch). The orchestrator holds no lane and never reads code. It spawns one fresh `wf-worker` per task (`wf orch pick` / `wf orch post`), one per lane in -parallel, on the task's model, and stops a lane on the first awaiting, handback or red -post-check. Use it when the backlog is runner-ready (section 3) and you want to steer. +parallel, on the task's model, and parks a task that ends awaiting, handback or red post-check +(blocked / `Sessions: owner`, alert) while its lane keeps going. Use it when the backlog is runner-ready (section 3) and you want to steer. ### Overnight: `/wf-pilot` `/wf-pilot [--batch 4] [--lanes a,b] [--for 8h]` in an idle session, which can be one you opened diff --git a/docs/orchestrator.md b/docs/orchestrator.md index f878855..f3e1432 100644 --- a/docs/orchestrator.md +++ b/docs/orchestrator.md @@ -30,7 +30,7 @@ only in sessions started afterwards: restart it. Two commands do steps 1-3 and 5-7 (one call each): `wf orch pick <lane> [--id ID] [--recovery WHY]` (stop file, pick, claim, free worktree, recorded in `.wf/orch/<id>.json`, prints the spawn line and the prompt below) and `wf orch post <id> <lane> --result "<report result line>" [--commit SHA] [--agent A] [--duration S] [--no-pick]` -(post-check, merge if needed, leftover commit, orch + cost log lines, then the next pick or `stop lane <lane>: <why>`). +(post-check, merge if needed, leftover commit, orch + cost log lines, then `alert: …` for a parked task, and the next pick or `stop lane <lane>: <why>` (none, stop file, push-failed, wip, no report)). The steps they automate: 1. Pick: `wf list --runner --lane <lane>` → first id (lanes: `wf lanes`). Claim: `wf status <id> progress "worker"`. No commit of @@ -63,7 +63,7 @@ The steps they automate: | model raised by the worker | continue; next pick spawns it on the new model | | sliced | next task in the lane | | done+gate-red <culprit> <fix-id> (async gate red, earlier task's break) | post-check as done; next pick in the lane = the P0 fix (no Done/Model → add them first); never stop the lane | - | awaiting / handback / post-check red / no report | stop the lane, tell the owner | + | awaiting / handback / post-check red | `wf orch post` parks only that task (blocked on its a-id, else `Sessions: owner` + note; handback clears the claim), prints `alert: …`; tell the owner, the lane keeps picking (never wait on the owner while tasks are pickable) | | wip (after a wrap-up) | stop the lane; reslice or rerun later | | crashed / killed (no report) | one fresh worker, prompt plus `Recovery: <why>` (it inspects git status, git log master..HEAD, `git diff -- TASKS.md` — often just the claim — and `wf show <id>`), then stop | 7. Log one line per task to `out/wf-orch.log` (git-ignored): time, lane, model, id, outcome, commit, duration. diff --git a/shared/skills/wf-orchestrate/SKILL.md b/shared/skills/wf-orchestrate/SKILL.md index 53a4785..c5b3d86 100644 --- a/shared/skills/wf-orchestrate/SKILL.md +++ b/shared/skills/wf-orchestrate/SKILL.md @@ -25,7 +25,9 @@ command: /projects/public/workflow/docs/orchestrator.md — read it only when a 3. Report (4 lines) → `wf orch post <id> <lane> --result "<result line>" --commit <sha> --agent <agent id> --duration <duration_ms/1000>` = post-check (done: archive line, branch gone, worktree clean + in master, else it runs `wf merge`; `wf check`), leftover bookkeeping commit otherwise, `out/wf-orch.log` + `out/wf-cost.log` lines, then the lane's next pick (spawn it) or `stop lane …`. -4. `stop lane` → tell the owner (awaiting/handback/post-check-red/wip). No report (crash) → it prints the one recovery pick +4. `alert: <id> <outcome> …` (awaiting/handback/post-check-red: post parked only that task, blocked or `Sessions: owner`) + → tell the owner, keep spawning the lane's next pick; never wait on the owner while tasks are pickable. + `stop lane` (none/stop file/push-failed/wip) → tell the owner. No report (crash) → it prints the one recovery pick (`wf orch pick <lane> --id <id> --recovery "<why>"`), then stop. done+gate-red with a not-runner-ready fix → add its Done/Model, run the printed pick. Cloud lane (project `cloud = true`): `wf orch pick cloud` → sends the first fitting task (`wf cloud send`), no agent to spawn; diff --git a/shared/skills/wf-pilot/SKILL.md b/shared/skills/wf-pilot/SKILL.md index 1f01f2b..86a207b 100644 --- a/shared/skills/wf-pilot/SKILL.md +++ b/shared/skills/wf-pilot/SKILL.md @@ -21,7 +21,7 @@ K = `--batch` (4); deadline = now + `--for` (8h); lanes = `--lanes` or all. `wf `wf res wait <rid> --timeout <same>`; end the turn. `fit: 0 of K …; nothing started` → stop (step 6, why: deadline). Exit 3 (busy) → background Bash `sleep 600`, then retry. 2. Wait exit → `wf batch --status` (newest summary) → one line to the owner and `out/wf-orch.log`: `<HH:MM> pilot batch <n> rc=<rc> done=<ids> stopped=<lanes: why> added=<ids>`. No task bodies, no logs. - A lane the batch stopped (handback/awaiting) is not relaunched in the next batch (`--lanes` without it) until the task/fix that stopped it is done, or the owner says so. + A handback/awaiting blocks only its task (parked by `wf orch post`, alert); its lane stays in the next batch while it has pickable tasks. Never wait on the owner while tasks are pickable. 3. Failure (rc≠0, or no new summary file) → retry once; second failure → stop (step 6). Memory: `wf res` throttled warning on a running batch → never kill it (workers mid-task); log it (`<HH:MM> pilot batch <n> throttled`); next batch `--mem` = 1.5 × last (cap per `wf res free`). Batch oom-killed → re-run the same N with 1.5 × `--mem`. diff --git a/templates/batch-prompt.md b/templates/batch-prompt.md index 1a0f338..9dfe564 100644 --- a/templates/batch-prompt.md +++ b/templates/batch-prompt.md @@ -1 +1 @@ -You are a batch orchestrator. Follow {here}/docs/orchestrator.md 'One worker' for up to {n} tasks (lanes: {lanes}), never implement. Before EVERY spawn (each round, a lane's P0 fix, a recovery worker) run `test -e {stop}`: it exists → spawn nothing, let running workers finish (never kill one), post-check them, append 'stopped: stop file' to {out}, delete the stop file, exit. Else, before each round run `wf batch --time-left {deadline}`: its line ends in ': stop' (less time left than a task takes, p90) → spawn nothing new, let running workers finish, post-check them, append 'stopped: deadline (<that line>)' to {out}, exit. Else each round: `wf orch pick <lane>` per lane, spawn one wf-worker per picked lane in ONE message (model and prompt: as printed), foreground (not background), wait for all, post-check each with `wf orch post <id> <lane> --result "<result line>" --commit <sha> --agent <agent id> --duration <s> --no-pick`, append one line per task to {out} (it exists with its '# wf-batch' header: only append, never rewrite) (time, lane, id, outcome, commit), repeat. If the project has `cloud = true`, the virtual lane cloud also runs each round (nothing to spawn): `wf orch pick cloud` until a `stop lane cloud` / `none fit` / `max parallel` line; never run `wf cloud pull` or `wf orch post <id> cloud` yourself: the batch job's sidecar pulls every 10 min (`wf cloud pull --all` → `wf orch post <id> cloud … --no-pick`, a 'cloud <id> <state> … (sidecar pull)' line in {out}), also while you wait on workers; a cloud stop never stops the local lanes. Stop a lane on its first stop outcome, and keep it stopped (never resume because some other task is pickable) until the task/fix that stopped it is done or the owner says so; done+gate-red <culprit> <fix-id> is not one (post-check as done, the lane picks the P0 fix next). Never ask the owner: follow-ups → wf add -p 2 --model <model>, decisions → wf add -s awaiting. At batch start note `git branch --list 'fast/*' 'slow/*'` (and worktrees): a branch already there is pre-existing; report it ONCE in {out} as 'pre-existing branch <name>' (+ the awaiting id when a task/awaiting names it), never per round or per task. End with the ids added in that file, then exit. +You are a batch orchestrator. Follow {here}/docs/orchestrator.md 'One worker' for up to {n} tasks (lanes: {lanes}), never implement. Before EVERY spawn (each round, a lane's P0 fix, a recovery worker) run `test -e {stop}`: it exists → spawn nothing, let running workers finish (never kill one), post-check them, append 'stopped: stop file' to {out}, delete the stop file, exit. Else, before each round run `wf batch --time-left {deadline}`: its line ends in ': stop' (less time left than a task takes, p90) → spawn nothing new, let running workers finish, post-check them, append 'stopped: deadline (<that line>)' to {out}, exit. Else each round: `wf orch pick <lane>` per lane, spawn one wf-worker per picked lane in ONE message (model and prompt: as printed), foreground (not background), wait for all, post-check each with `wf orch post <id> <lane> --result "<result line>" --commit <sha> --agent <agent id> --duration <s> --no-pick`, append one line per task to {out} (it exists with its '# wf-batch' header: only append, never rewrite) (time, lane, id, outcome, commit), repeat. If the project has `cloud = true`, the virtual lane cloud also runs each round (nothing to spawn): `wf orch pick cloud` until a `stop lane cloud` / `none fit` / `max parallel` line; never run `wf cloud pull` or `wf orch post <id> cloud` yourself: the batch job's sidecar pulls every 10 min (`wf cloud pull --all` → `wf orch post <id> cloud … --no-pick`, a 'cloud <id> <state> … (sidecar pull)' line in {out}), also while you wait on workers; a cloud stop never stops the local lanes. An `alert: <id> …` line from post (awaiting / handback / post-check-red) parks only that task: append 'alert <id> <outcome>' to {out} and keep picking in that lane (never wait on the owner while tasks are pickable). Stop a lane only on a `stop lane` / `none` line (none, stop file, push-failed, wip, no report after its one recovery) and keep it stopped until that is resolved or the owner says so; done+gate-red <culprit> <fix-id> is not one (post-check as done, the lane picks the P0 fix next). Never ask the owner: follow-ups → wf add -p 2 --model <model>, decisions → wf add -s awaiting. At batch start note `git branch --list 'fast/*' 'slow/*'` (and worktrees): a branch already there is pre-existing; report it ONCE in {out} as 'pre-existing branch <name>' (+ the awaiting id when a task/awaiting names it), never per round or per task. End with the ids added in that file, then exit. diff --git a/tests/test_orch.py b/tests/test_orch.py index 83216b2..157d7ef 100644 --- a/tests/test_orch.py +++ b/tests/test_orch.py @@ -157,9 +157,13 @@ class OrchTest(Cli): def test_post_done_without_archive_is_red(self): self.orch("pick", "fast") out = self.orch("post", "t-one", "fast", "--result", "done") - self.assertEqual(out, "post: t-one post-check-red\n no archive line for t-one\n" - "stop lane fast: post-check-red → tell the owner\n") + self.assertTrue(out.startswith( + "post: t-one post-check-red\n no archive line for t-one\n committed leftover TASKS.md tasks/archive.md\n" + "alert: t-one post-check-red → tell the owner (Sessions: owner); lane fast keeps picking\n\n" + "pick: t-two "), out) + self.assertNotIn("stop lane", out) self.assertIn(" post-check-red ", self.log()) + self.assertIn(" Sessions: owner\n", self.wf("show", "t-one")[1]) def backdate_pick(self, id="t-one", secs=125): f = self.root / ".wf" / "orch" / f"{id}.json" @@ -184,13 +188,32 @@ class OrchTest(Cli): self.assertRegex(self.log().splitlines()[-1], r" t-one done \S+ 2m0[5-9]s$") self.assertFalse((self.root / ".wf" / "orch" / "t-one.json").exists()) - def test_post_handback_commits_leftovers_and_stops(self): + def test_post_handback_parks_task_and_picks_next(self): self.orch("pick", "fast") out = self.orch("post", "t-one", "fast", "--result", "handback unclear spec") - self.assertEqual(out, "post: t-one handback\n committed leftover TASKS.md tasks/archive.md\n" - "stop lane fast: handback → tell the owner\n") + self.assertTrue(out.startswith( + "post: t-one handback\n committed leftover TASKS.md tasks/archive.md\n" + "alert: t-one handback → tell the owner (Sessions: owner, status cleared); lane fast keeps picking\n\n" + "pick: t-two (lane fast"), out) + self.assertNotIn("stop lane", out) + shown = self.wf("show", "t-one")[1] + self.assertIn(" Sessions: owner\n", shown) + self.assertIn("handback (unclear spec)", shown) + self.assertNotIn("in progress", shown.splitlines()[0]) self.assertEqual(subprocess.run(["git", "-C", str(self.root), "log", "-1", "--format=%s"], capture_output=True, text=True).stdout, "t-one handback (orchestrator)\n") + self.assertIn("Sessions: owner", subprocess.run(["git", "-C", str(self.root), "show", "HEAD:TASKS.md"], + capture_output=True, text=True).stdout) + + def test_post_awaiting_blocks_task_and_picks_next(self): + self.orch("pick", "fast") + out = self.wf("add", "-s", "awaiting", "Which way?")[1] + aid = next(w.strip("*[]:") for w in out.split() if w.strip("*[]:").startswith("a-")) + out = self.orch("post", "t-one", "fast", "--result", f"awaiting {aid}", "--no-pick") + self.assertEqual(out, f"post: t-one awaiting\n committed leftover TASKS.md tasks/archive.md\n" + f"alert: t-one awaiting → tell the owner (blocked on {aid}); lane fast keeps picking\n") + self.assertIn(f"blocked: [[{aid}]]", self.wf("show", "t-one")[1]) + self.assertIn("pick: t-two", self.orch("pick", "fast")) def test_post_done_on_slice_job_counts_as_sliced(self): self.orch("pick", "slow", "--id", "t-big") @@ -1597,6 +1597,7 @@ def cmd_gate(args) -> int: # ------------------------------------------------------------------ orchestrator GO_ON = ("done", "done+gate-red", "sliced") # outcomes after which the lane picks its next task +PARK = ("handback", "awaiting", "post-check-red") # block only that task (park_task), the lane keeps picking def orch_main(args) -> tuple[Path, Project]: @@ -1778,6 +1779,38 @@ def archived_meta(cfg: config.Config, id: str) -> tuple[str | None, str | 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/post-check-red: 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 + if (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"): + unclaim(p.cfg, [id]) + return done + + def orch_post(main: Path, p: Project, args) -> int: cfg, id, lane = p.cfg, args.id, args.lane rec = orch_record(cfg, id) @@ -1803,7 +1836,9 @@ def orch_post(main: Path, p: Project, args) -> int: out.append(f"commit {c} not in code repo {repo} (a private sha?): using its {master} {fixed}") c = fixed return c - commit = None + 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: @@ -1832,12 +1867,15 @@ def orch_post(main: Path, p: Project, args) -> int: 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]}") - elif git_run(main, "status", "--porcelain", "--", *files).stdout.strip(): - r = git_run(main, "commit", "-q", "-m", f"{id} {outcome} (orchestrator)", "--", *files) - out.append("committed leftover " + " ".join(files) if not r.returncode - else f"leftover commit failed: {(r.stderr.strip() or 'git error').splitlines()[-1]}") - raised = bool(item and rec.get("model") and not outcome.startswith("done") - and tasks.MODELS.index(item.model) > tasks.MODELS.index(rec["model"])) + if problems or not outcome.startswith("done"): + if lane != CLOUD and (problems or outcome in PARK and not raised): # cloud: pull already re-queued + parked = park_task(cfg.root, id, "post-check-red" if problems else 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" @@ -1867,7 +1905,9 @@ def orch_post(main: Path, p: Project, args) -> int: print(f"fix {fix.id} not runner-ready: add its Done/Model (wf set {fix.id} --done … --model …), " f"then wf orch pick {lane} --id {fix.id}") return 0 - if final not in GO_ON and final != "model-raised": + if final in PARK and lane != CLOUD: + 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)" @@ -2306,7 +2346,7 @@ def parser() -> argparse.ArgumentParser: op.add_argument("--recovery", metavar="WHY", help="prompt gets 'Recovery: WHY' (crashed worker, one retry)") op = osub.add_parser("post", help="after the worker's report: post-check (done: archive line, branch gone, " "worktree clean and in master, else wf merge; wf check), leftover commit otherwise, " - "out/wf-orch.log line, out/wf-cost.log line (--agent), then the next pick or 'stop lane'") + "out/wf-orch.log line, out/wf-cost.log line (--agent); awaiting/handback/post-check-red park only the task (blocked / Sessions: owner + 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)") |
