Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 6 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,8 +31,10 @@ board to a PR — or fork it as a starting point.
separate store, so the work graph can't drift out of sync.
- **The loop** pulls the top-priority `ready` feature → creates a disposable
`git worktree` off `origin/<base>` → dispatches a coder (`acp` delegate) scoped to
it → commits/pushes → opens a PR → `in_review`. A **merge webhook** is the single
edge that sets `done` (and reaps the worktree).
it → commits/pushes → opens a PR → `in_review`. A **merge webhook** sets `done`
(and reaps the worktree); where GitHub can't reach a webhook URL, a **merge poll**
(`merge_poll`, on by default) runs the same idempotent Done edge. Set
`max_concurrent > 1` to build several features in parallel, each in its own worktree.
- **DAG + gates** — `depends_on` are `blocks` edges; a dependent stays out of the
puller until its blocker is **merged** (foundation merge-gate). The **Ready gate**
requires a spec, EARS acceptance criteria, and explicit `files_to_modify`.
Expand Down Expand Up @@ -76,6 +78,8 @@ project_board:
repo: ~/dev/my-repo
base_branch: main
loop_enabled: false # flip true to start the background puller
max_concurrent: 1 # >1 builds features in parallel (each its own worktree)
merge_poll: true # poll merged PRs as a fallback to the webhook Done edge
# webhook_secret: "..." # set before exposing /webhook/pr publicly
```

Expand Down
4 changes: 1 addition & 3 deletions api.py
Original file line number Diff line number Diff line change
Expand Up @@ -138,9 +138,7 @@ async def _webhook_pr(request: Request):
try:
from . import worktree

repo = store_kw["repo"]
wt = os.path.join(repo, worktrees_root, f"feat-{f['id']}")
await worktree.remove_worktree(repo, wt, f"feat/{f['id']}")
await worktree.reap_feature_worktree(store_kw["repo"], worktrees_root, f["id"])
except Exception: # noqa: BLE001 — reaping is best-effort; done is already set
log.warning("[project_board] worktree reap for %s failed", f["id"], exc_info=True)
log.info("[project_board] merge webhook → done: %s (%s)", f["id"], pr_url)
Expand Down
135 changes: 99 additions & 36 deletions loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,17 +11,21 @@
└──▶ in_review ──delegate_to(reviewer)──▶ (CI + review on the PR)
│
merge webhook ▼ CI fail ▼ any failure ▼
done in_progress (bounce) blocked (flag + reason)
/merge poll in_progress (bounce) blocked (flag + reason)
done

CI status + merge arrive out-of-band via the board API (``api.py``), not here.
Concurrency is capped at 1 (token + merge-integration cost); the cap is where
parallelism lands later.
CI status arrives out-of-band via the board API (``api.py``). ``done`` is set by
the merge webhook (``api.record_merge``) — or, when no public webhook URL is
reachable, by the loop's **merge poll** (``merge_poll``), which asks ``gh`` whether
each ``in_review`` PR has merged and runs the same idempotent Done edge. Up to
``max_concurrent`` features build concurrently, each in its own worktree.
"""

from __future__ import annotations

import asyncio
import logging
import time

from . import worktree
from .store import BoardError, escalation_enabled, get_store
Expand Down Expand Up @@ -49,17 +53,27 @@ def __init__(self, cfg: dict):
# and we never write redundant tier/attempt labels.
self.coders = {str(k): str(v) for k, v in (self.cfg.get("coders") or {}).items()}
self.escalation_on = escalation_enabled(self.cfg)
# Concurrency: drive up to `max_concurrent` features at once, each in its own
# worktree. 1 (the default) = serial — the safe default for token + merge-
# integration cost; raise it on a repo that parallelizes cleanly.
self.max_concurrent = max(1, int(self.cfg.get("max_concurrent", 1)))
# Merge poll: a fallback to the /webhook/pr Done edge for deployments with no
# public webhook URL. On by default (cheap; only probes `in_review` PRs).
self.merge_poll = bool(self.cfg.get("merge_poll", True))
self.merge_poll_interval = float(self.cfg.get("merge_poll_interval_s", 60))
self._store_kw = dict(
db=self.cfg.get("db_path") or None,
repo=self.cfg.get("repo", "."),
base_branch=self.cfg.get("base_branch", "main"),
)
self._task: asyncio.Task | None = None
self._stop = asyncio.Event()
# The in-flight build's worktree, so shutdown can reap it (a cancel mid-drive
# would otherwise orphan the worktree; the coder subprocess is reaped by
# dispatch_coder's finally). (repo, worktree_path, branch) or None.
self._active: tuple[str, str, str] | None = None
# The running drive tasks, and the worktrees they hold (fid → (repo, wt,
# branch)) so shutdown can reap any a cancel mid-drive would orphan; the coder
# subprocess itself is reaped by dispatch_coder's finally.
self._drives: set[asyncio.Task] = set()
self._inflight: dict[str, tuple[str, str, str]] = {}
self._last_poll = 0.0 # monotonic ts of the last merge poll

def _store(self):
return get_store(**self._store_kw)
Expand All @@ -71,10 +85,12 @@ def start(self):
return None
self._task = asyncio.create_task(self._run(), name="project-board-loop")
log.info(
"[project_board] loop started (coder=%s reviewer=%s every %ss)",
"[project_board] loop started (coder=%s reviewer=%s every %ss, max_concurrent=%d, merge_poll=%s)",
self.coder_name,
self.reviewer_name,
self.interval,
self.max_concurrent,
self.merge_poll,
)
return self._task

Expand All @@ -86,11 +102,16 @@ async def stop(self):
await self._task
except (asyncio.CancelledError, Exception): # noqa: BLE001
pass
# Reap an interrupted build's worktree (a drive cancelled mid-flight leaves
# self._active set; a completed/blocked drive clears it). Best-effort.
act, self._active = self._active, None
if act:
repo, wt, branch = act
# Cancel any in-flight drives and await them out. A drive cancelled mid-flight
# can't run its own cleanup, so its worktree stays in self._inflight — reaped
# below. (A completed/blocked drive already popped itself.)
drives, self._drives = list(self._drives), set()
for t in drives:
t.cancel()
if drives:
await asyncio.gather(*drives, return_exceptions=True)
inflight, self._inflight = dict(self._inflight), {}
for fid, (repo, wt, branch) in inflight.items():
try:
await worktree.remove_worktree(repo, wt, branch or "")
log.info("[project_board] reaped in-flight worktree on shutdown: %s", wt)
Expand All @@ -100,26 +121,68 @@ async def stop(self):
# ── the puller ────────────────────────────────────────────────────────────
async def _run(self):
while not self._stop.is_set():
spawned = False
try:
worked = await self._tick()
await self._maybe_poll_merges()
spawned = self._spawn_ready()
except Exception: # noqa: BLE001 — a bad tick must never kill the loop
log.exception("[project_board] loop tick failed")
worked = False
# If we did work, loop again immediately (drain Ready); else sleep.
if not worked:
try:
await asyncio.wait_for(self._stop.wait(), timeout=self.interval)
except asyncio.TimeoutError:
pass

async def _tick(self) -> bool:
"""Claim and drive at most one Ready feature. Returns True if it claimed
one (so the runner can drain), False if Ready was empty."""
feature = self._store().claim_next_ready(assignee=self.coder_name)
if feature is None:
return False
await self._drive(feature)
return True
# Idle (nothing started, nothing running) → sleep the full interval. Busy
# → re-check soon so a freed concurrency slot refills and merges land
# promptly (the poll itself stays rate-limited by merge_poll_interval).
idle = not spawned and not self._drives
timeout = self.interval if idle else min(self.interval, 3.0)
try:
await asyncio.wait_for(self._stop.wait(), timeout=timeout)
except asyncio.TimeoutError:
pass

def _spawn_ready(self) -> bool:
"""Claim Ready features up to the concurrency cap and spawn a drive task for
each. Returns True if it started at least one (so the runner stays hot)."""
spawned = False
while len(self._drives) < self.max_concurrent:
feature = self._store().claim_next_ready(assignee=self.coder_name)
if feature is None:
break
task = asyncio.create_task(self._drive(feature), name=f"pb-drive-{feature['id']}")
self._drives.add(task)
task.add_done_callback(self._drives.discard)
spawned = True
return spawned

# ── the merge poll (Done-edge fallback to the webhook) ─────────────────────
async def _maybe_poll_merges(self):
"""Run the merge poll at most once per ``merge_poll_interval`` (and only when
enabled) — cheap, but no reason to hammer ``gh`` every busy tick."""
if not self.merge_poll:
return
now = time.monotonic()
if now - self._last_poll < self.merge_poll_interval:
return
self._last_poll = now
await self._poll_merges()

async def _poll_merges(self):
"""Ask ``gh`` whether each ``in_review`` PR has merged and run the idempotent
Done edge for any that have — the fallback for deployments GitHub can't post
a webhook to (otherwise a merged feature would sit in_review forever)."""
store = self._store()
repo = self._store_kw["repo"]
for f in store.list_features(state="in_review"):
pr_url = f.get("pr_url")
if not pr_url:
continue
try:
if not await worktree.pr_is_merged(pr_url, cwd=repo):
continue
done = store.record_merge(pr_url=pr_url)
except Exception: # noqa: BLE001 — a poll error must never kill the loop
log.warning("[project_board] merge poll for %s failed", f["id"], exc_info=True)
continue
if done:
await worktree.reap_feature_worktree(repo, self.root, f["id"])
log.info("[project_board] merge poll → done: %s (%s)", f["id"], pr_url)

async def _drive(self, feature: dict):
"""Drive one feature ready→in_review (or →blocked). `done` is set later by
Expand All @@ -143,7 +206,7 @@ async def _drive(self, feature: dict):
return
# Fresh worktree per attempt (a failed attempt may leave partial work).
wt, branch = await worktree.create_worktree(repo, base, fid, self.root)
self._active = (repo, wt, branch) # track for shutdown reaping
self._inflight[fid] = (repo, wt, branch) # track for shutdown reaping
try:
result = await worktree.dispatch_coder(coder, wt, prompt) # reaps subprocess
pr_url = await worktree.open_pr(wt, branch, base=base, title=title, body=(result or "")[:4000])
Expand All @@ -163,7 +226,7 @@ async def _drive(self, feature: dict):
store.flag_blocked(fid, str(exc))
if wt:
await worktree.remove_worktree(repo, wt, branch or "")
self._active = None
self._inflight.pop(fid, None)
return
# Built + PR opened. The fleet PR-review pipeline reviews it on open;
# only dispatch an explicit review when configured to (review_dispatch).
Expand All @@ -173,18 +236,18 @@ async def _drive(self, feature: dict):
await self._request_review(fid, pr_url)
# Keep the worktree (a CI-fail bounce re-dispatches); reaping happens
# on a terminal block above, and the coder subprocess is already reaped.
self._active = None # built OK — not an interrupted build to reap
self._inflight.pop(fid, None) # built OK — not an interrupted build to reap
return
except BoardError as exc:
log.warning("[project_board] %s blocked (board): %s", fid, exc)
store.flag_blocked(fid, str(exc))
self._active = None
self._inflight.pop(fid, None)
except Exception as exc: # noqa: BLE001 — unexpected; block, don't crash the loop
log.exception("[project_board] %s unexpected failure", fid)
store.flag_blocked(fid, f"unexpected: {type(exc).__name__}: {exc}")
if wt:
await worktree.remove_worktree(repo, wt, branch or "")
self._active = None
self._inflight.pop(fid, None)

async def _request_review(self, fid: str, pr_url: str):
"""Hand the PR to the reviewer (an a2a delegate, e.g. quinn). Best-effort:
Expand Down
11 changes: 9 additions & 2 deletions protoagent.plugin.yaml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
id: project_board
name: Project Board (coding orchestration)
version: 0.2.0
version: 0.3.0
description: >-
A board-driven coding-orchestration plugin: a lean 6-state board (backlog → ready
→ in_progress → in_review → done, + a blocked flag) backed by **beads** (`br`), an
Expand Down Expand Up @@ -40,7 +40,14 @@ config:
worktrees_root: .worktrees
loop_enabled: false # the background puller; off until a board exists
loop_interval_s: 30 # how often the puller wakes
max_concurrent: 1 # worktree/coder concurrency cap (token + merge cost)
max_concurrent: 1 # worktree/coder concurrency cap (token + merge cost).
# >1 drives that many features in parallel, each in its
# own worktree; 1 = serial (the safe default).
merge_poll: true # poll merged PRs as a fallback to the /webhook/pr Done
# edge — for deployments GitHub can't post a webhook to
# (no public URL). The Done edge stays idempotent either
# way; set false to rely on the webhook alone.
merge_poll_interval_s: 60 # how often the loop polls in_review PRs for a merge
db_path: "" # beads db path; blank → br auto-discovers .beads/*.db
webhook_secret: "" # GitHub webhook HMAC secret for /webhook/pr (the Done
# edge). Blank → signature NOT verified (dev only); set
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[project]
name = "project-board"
version = "0.2.0"
version = "0.3.0"
description = "Board-driven coding-orchestration plugin for protoAgent (beads board + ACP spawn loop + planning layer + console view)."
requires-python = ">=3.11"

Expand Down
Loading
Loading