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
7 changes: 4 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,9 +32,10 @@ board to a PR — or fork it as a starting point.
- **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** 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.
(and reaps the worktree); where GitHub can't reach a webhook URL, a **PR reconcile
poll** (`merge_poll`, on by default) drives the terminal edges itself — merged →
`done`, closed-unmerged → `blocked`. Set `max_concurrent > 1` to build several
features in parallel, each in its own worktree.
- **Resilience** — every `await` in a drive is bounded (a coder dispatch is hard-capped
by `coder_timeout_s`); **transient** failures (rate-limit / network / merge-conflict)
retry with backoff while **capability** failures (no diff / timeout) escalate a tier
Expand Down
49 changes: 28 additions & 21 deletions loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,10 @@

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.
reachable, by the loop's **PR reconcile** (``merge_poll``), which asks ``gh`` for
each ``in_review`` PR's state and drives the terminal edges: merged → done (the same
idempotent edge), closed-unmerged → blocked. Up to ``max_concurrent`` features build
concurrently, each in its own worktree.
"""

from __future__ import annotations
Expand Down Expand Up @@ -214,7 +215,7 @@ async def _run(self):
while not self._stop.is_set():
spawned = False
try:
await self._maybe_poll_merges()
await self._maybe_reconcile()
await self._maybe_sweep()
spawned = self._spawn_ready()
except Exception: # noqa: BLE001 — a bad tick must never kill the loop
Expand Down Expand Up @@ -273,38 +274,44 @@ def _cb(task: asyncio.Task):

return _cb

# ── 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
# ── the PR reconcile (terminal-edge fallback to the webhook) ───────────────
async def _maybe_reconcile(self):
"""Run the PR reconcile 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()
await self._reconcile_prs()

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)."""
async def _reconcile_prs(self):
"""Reconcile each ``in_review`` feature against its PR's real state — the
fallback to the webhook and the active half of the terminal edges (for
deployments GitHub can't post a webhook to, where a feature would otherwise
sit in_review forever): ``MERGED`` → done (+reap); ``CLOSED`` unmerged →
Blocked for triage (+reap; the work was rejected, don't silently re-dispatch);
``OPEN`` → leave it in review."""
store = self._store()
repo = self._store_kw["repo"]
for f in store.list_features(state="in_review"):
fid = f["id"]
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)
state = await worktree.pr_state(pr_url, cwd=repo)
if state == "MERGED":
if store.record_merge(pr_url=pr_url):
await worktree.reap_feature_worktree(repo, self.root, fid)
log.info("[project_board] reconcile → done: %s (%s)", fid, pr_url)
elif state == "CLOSED":
store.flag_blocked(fid, f"PR closed without merging — needs triage: {pr_url}")
await worktree.reap_feature_worktree(repo, self.root, fid)
log.info("[project_board] reconcile → blocked (PR closed): %s (%s)", fid, pr_url)
except Exception: # noqa: BLE001 — a reconcile error must never kill the loop
log.warning("[project_board] reconcile for %s failed", fid, exc_info=True)

async def _drive(self, feature: dict):
"""Drive one feature ready→in_review (or →blocked). `done` is set later by
Expand Down
9 changes: 5 additions & 4 deletions protoagent.plugin.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -58,11 +58,12 @@ config:
# in_review (more parallelism, but builds on un-merged code).
# Mark shared-structure features foundation=True; those
# always gate on merge regardless.
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
merge_poll: true # poll in_review PRs and reconcile the board to their real
# state (merged → done, closed-unmerged → blocked) — the
# fallback to /webhook/pr for deployments GitHub can't post
# a webhook to (no public URL). Done 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
merge_poll_interval_s: 60 # how often the loop reconciles in_review PRs
health_sweep_interval_s: 300 # periodic self-heal: reclaim slots from dead drives
# (in_progress with no live drive → ready) + reap orphaned
# feat-<id> worktrees (feature gone/done). 0 disables it.
Expand Down
57 changes: 34 additions & 23 deletions tests/test_loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -320,13 +320,14 @@ async def _quick(feature):
assert loop._drives == set()


# ── the merge poll (Done-edge fallback) ─────────────────────────────────────────
# ── the PR reconcile (terminal-edge fallback) ───────────────────────────────────


class _PollStore:
class _ReconcileStore:
def __init__(self, in_review):
self._in_review = in_review
self.merged = []
self.blocked = []

def list_features(self, state=None):
return self._in_review if state == "in_review" else []
Expand All @@ -335,57 +336,67 @@ def record_merge(self, *, pr_url):
self.merged.append(pr_url)
return {"id": "x", "board_state": "done"}

def flag_blocked(self, fid, reason):
self.blocked.append((fid, reason))


async def test_poll_merges_runs_done_edge_for_merged_only(monkeypatch):
store = _PollStore(
async def test_reconcile_drives_merged_to_done_and_closed_to_blocked(monkeypatch):
store = _ReconcileStore(
[
{"id": "bd-1", "pr_url": "https://example/pr/1"},
{"id": "bd-2", "pr_url": "https://example/pr/2"},
{"id": "bd-3", "pr_url": ""}, # no PR → skipped entirely
{"id": "bd-merged", "pr_url": "https://example/pr/1"},
{"id": "bd-closed", "pr_url": "https://example/pr/2"},
{"id": "bd-open", "pr_url": "https://example/pr/3"},
{"id": "bd-nopr", "pr_url": ""}, # no PR → skipped entirely
]
)
monkeypatch.setattr("project_board.loop.get_store", lambda **_kw: store)
states = {
"https://example/pr/1": "MERGED",
"https://example/pr/2": "CLOSED",
"https://example/pr/3": "OPEN",
}

async def _is_merged(url, *, cwd="."):
return url.endswith("/1") # only PR 1 has merged
async def _pr_state(url, *, cwd="."):
return states[url]

reaped = []

async def _reap(repo, root, fid):
reaped.append(fid)

monkeypatch.setattr(worktree, "pr_is_merged", _is_merged)
monkeypatch.setattr(worktree, "pr_state", _pr_state)
monkeypatch.setattr(worktree, "reap_feature_worktree", _reap)

await BoardLoop({})._poll_merges()
assert store.merged == ["https://example/pr/1"] # the unmerged + PR-less ones skipped
assert reaped == ["bd-1"] # the merged feature's worktree is reaped
await BoardLoop({})._reconcile_prs()
assert store.merged == ["https://example/pr/1"] # merged → done
assert [b[0] for b in store.blocked] == ["bd-closed"] # closed-unmerged → blocked
assert set(reaped) == {"bd-merged", "bd-closed"} # both terminal states reap; open kept


async def test_maybe_poll_is_rate_limited(monkeypatch):
async def test_maybe_reconcile_is_rate_limited(monkeypatch):
loop = BoardLoop({"merge_poll": True, "merge_poll_interval_s": 60})
calls = []

async def _poll():
async def _reconcile():
calls.append(1)

monkeypatch.setattr(loop, "_poll_merges", _poll)
monkeypatch.setattr(loop, "_reconcile_prs", _reconcile)
clock = {"t": 1000.0}
monkeypatch.setattr("project_board.loop.time.monotonic", lambda: clock["t"])

await loop._maybe_poll_merges() # first → polls
await loop._maybe_poll_merges() # immediately again → rate-limited, skipped
await loop._maybe_reconcile() # first → reconciles
await loop._maybe_reconcile() # immediately → rate-limited
clock["t"] += 61
await loop._maybe_poll_merges() # interval elapsed → polls again
await loop._maybe_reconcile() # interval elapsed → reconciles again
assert len(calls) == 2


async def test_merge_poll_off_never_polls(monkeypatch):
async def test_merge_poll_off_never_reconciles(monkeypatch):
loop = BoardLoop({"merge_poll": False})
called = []
monkeypatch.setattr(loop, "_poll_merges", lambda: called.append(1))
await loop._maybe_poll_merges()
assert called == [] # disabled → the poll is never reached
monkeypatch.setattr(loop, "_reconcile_prs", lambda: called.append(1))
await loop._maybe_reconcile()
assert called == [] # disabled → never reconciles


# ── crash recovery on boot ──────────────────────────────────────────────────────
Expand Down
12 changes: 6 additions & 6 deletions tests/test_worktree.py
Original file line number Diff line number Diff line change
Expand Up @@ -140,16 +140,16 @@ async def test_create_worktree_falls_back_to_local_base_without_a_remote(monkeyp
@pytest.mark.parametrize(
"gh_state,expected",
[
((0, "MERGED", ""), True),
((0, "OPEN", ""), False),
((0, "CLOSED", ""), False), # closed-unmerged is NOT done
((1, "", "no pr found"), False), # a gh failure → False, the poll retries
((0, "MERGED", ""), "MERGED"),
((0, "OPEN", ""), "OPEN"),
((0, "CLOSED", ""), "CLOSED"),
((1, "", "no pr found"), ""), # a gh failure → "", the reconcile retries
],
)
async def test_pr_is_merged(monkeypatch, gh_state, expected):
async def test_pr_state(monkeypatch, gh_state, expected):
gh = FakeGh({"view": gh_state})
_install(monkeypatch, FakeGit(), gh)
assert await worktree.pr_is_merged("https://example/pr/1", cwd="/repo") is expected
assert await worktree.pr_state("https://example/pr/1", cwd="/repo") == expected


# ── pr_url_for_branch: the crash-recovery probe ─────────────────────────────────
Expand Down
12 changes: 6 additions & 6 deletions worktree.py
Original file line number Diff line number Diff line change
Expand Up @@ -206,13 +206,13 @@ async def open_pr(worktree: str, branch: str, *, base: str = "main", title: str,
raise WorktreeError(f"gh pr create failed: {err.strip()[:300]}")


async def pr_is_merged(pr_url: str, *, cwd: str = ".") -> bool:
"""True iff the PR has merged — the merge poll's probe (a fallback to the
webhook for deployments with no public webhook URL). A non-zero ``gh`` /
transient failure returns False so the next poll simply retries; this never
raises into the loop."""
async def pr_state(pr_url: str, *, cwd: str = ".") -> str:
"""The PR's state — ``MERGED`` / ``CLOSED`` / ``OPEN`` — or ``""`` on a ``gh``
failure (the next poll just retries; this never raises into the loop). The PR
reconcile drives the board's Done/closed edges off this (the fallback to the
webhook for deployments with no public webhook URL)."""
rc, out, _err = await _gh("pr", "view", pr_url, "--json", "state", "--jq", ".state", cwd=cwd)
return rc == 0 and out.strip() == "MERGED"
return out.strip() if rc == 0 else ""


async def pr_url_for_branch(branch: str, *, cwd: str = ".") -> str:
Expand Down
Loading