Coverage for src/keel/swarm_review_runtime.py: 100%
157 statements
« prev ^ index » next coverage.py v7.16.2, created at 2026-10-02 20:26 +0000
« prev ^ index » next coverage.py v7.16.2, created at 2026-10-02 20:26 +0000
1"""Keel Swarm Review Runtime — each cluster pull request reviewed by its own seats (#1423).
3The thin I/O half of ``keel swarm-review``. For each cluster of a wave of the plan
4``swarm-run`` persisted, it finds the cluster's pull request, reads its head, diff and
5review contract, plans the cluster's reviewer seats (:mod:`keel.swarm_review`, where every
6decision is made), and — live — runs each seat read-only in its own temporary checkout of
7the head, reads each answer, re-reads the head, and hands every verdict that parsed —
8approvals and change requests — to ``keel review`` to post, pinned to the head the seats
9reviewed.
11Every seat runs through :mod:`keel.delegaterun` in the read-only ``review`` role, under the
12implementer's lockdown environment (:func:`keel.swarm_worker.implementer_env`: no forge
13token, no ``gh`` login, no git credential helper, no network transport). Its checkout is a
14detached worktree at the head under ``.keel/worktrees/<swarm_id>/<cluster_id>.review-<slot>``,
15removed when the seat ends, whichever way; one a crash leaves behind is a leftover
16``keel swarm-status --clean`` removes.
18Each live cluster's outcome is recorded in its worker record of the run state
19(:func:`record_review`, #1440), which ``keel swarm-status`` shows; ``swarm-land`` does not
20read it.
22The outside world — the pull-request lookup, the reads, the seat itself and the post — is
23injected (:class:`ReviewIO`), so the loop is tested against recorded answers.
24"""
26from __future__ import annotations
28import concurrent.futures
29import os
30from collections.abc import Callable
31from dataclasses import dataclass, replace
32from pathlib import Path
33from typing import Any, cast
35from . import swarm_review, swarm_runtime, swarm_worker
36from .delegate import RunPlan
37from .swarm import (
38 SwarmCluster,
39 SwarmPlan,
40 load_swarm_state,
41 save_swarm_state,
42 update_worker_state,
43)
44from .swarm_landing import FindPullRequest
45from .swarm_review import ClusterReview, PullRequestFacts, ReviewSeat, SeatVerdict
46from .swarm_runtime import Implementer, SubprocessRunner, default_runner
48#: ``(pull request, with_diff)`` -> what keel reads before any seat runs; raises on failure.
49ReadPullRequest = Callable[[int, bool], PullRequestFacts]
50#: ``pull request`` -> its head right now; raises when it cannot be read.
51ReadHead = Callable[[int], str]
52#: ``issue`` -> ``(title, body)``, empty strings when it cannot be read.
53ReadIssue = Callable[[int], tuple[str, str]]
54#: ``(pull request, head, bundle, run id)`` -> ``""`` once posted, else why not.
55PostVerdicts = Callable[[int, str, list[dict[str, Any]], str], str]
58def _default_seat(plan: RunPlan, env: dict[str, str]) -> dict[str, Any]:
59 """``keel delegate run``'s executor with the seat's locked-down environment — the one a
60 live implementer runs through, resolved at call time."""
61 return swarm_runtime._default_implement(plan, env)
64@dataclass(frozen=True)
65class ReviewIO:
66 """The seams ``swarm-review`` reaches the outside world through."""
68 find_pull_request: FindPullRequest
69 read_pull_request: ReadPullRequest
70 read_head: ReadHead
71 read_issue: ReadIssue
72 post_verdicts: PostVerdicts
73 #: Runs one seat's plan with its environment: ``keel delegate run``'s executor.
74 run_seat: Implementer = _default_seat
75 #: The git commands keel itself runs (fetch, worktree add/remove, the tamper reads).
76 runner: SubprocessRunner | None = None
77 remote: str = "origin"
80def review_name(cluster_id: str, slot: str) -> str:
81 return f"{cluster_id}.review-{slot}"
84def review_checkout_path(swarm_id: str, cluster_id: str, slot: str, root: str | Path) -> Path:
85 """A seat's checkout of the head, beside the run's worker worktrees."""
86 return swarm_runtime.build_worktree_path(swarm_id, review_name(cluster_id, slot), root)
89def review_brief_path(swarm_id: str, cluster_id: str, slot: str, root: str | Path) -> Path:
90 """A seat's brief, beside the run's state and outside every checkout."""
91 return swarm_runtime.build_brief_path(swarm_id, review_name(cluster_id, slot), root)
94def _ensure_commit(root: Path, sha: str, pr: int, remote: str, run: SubprocessRunner) -> str:
95 """``""`` once ``sha`` is in the local repository, else why it is not.
97 A head the cluster's worker pushed from this repository is already here; one pushed
98 since (a fix on the pull request) is fetched from the pull request's own ref.
99 """
100 probe = ["git", "cat-file", "-e", f"{sha}^{{commit}}"]
101 if run(probe, root).ok:
102 return ""
103 fetched = run(["git", "fetch", "--no-tags", remote, f"refs/pull/{pr}/head"], root)
104 if run(probe, root).ok:
105 return ""
106 why = fetched.output.strip()[:200] or f"exit {fetched.code}"
107 return f"the head {sha} is not in this repository and could not be fetched ({why})"
110def _run_seat(
111 seat: ReviewSeat,
112 brief: str,
113 *,
114 root: Path,
115 swarm_id: str,
116 cluster_id: str,
117 facts: PullRequestFacts,
118 io: ReviewIO,
119 run: SubprocessRunner,
120) -> tuple[SeatVerdict, list[str]]:
121 """One seat, in its own checkout of the head, removed afterwards whatever happened."""
122 # Only an eligible seat runs, and an eligible seat is planned with its checkout.
123 plan = cast(RunPlan, seat.plan)
124 checkout = Path(str(plan.cwd))
125 warnings: list[str] = []
126 created = False
127 try:
128 added = run(["git", "worktree", "add", "--detach", str(checkout), facts.head_sha], root)
129 created = added.ok and checkout.exists()
130 if not created:
131 why = added.output.strip()[:200] or f"exit {added.code}"
132 return swarm_review.failed_verdict(
133 seat,
134 f"its checkout of the head could not be made at {checkout} ({why}); a previous "
135 "run may have left it — `keel swarm-status <project.yaml> --clean` removes it",
136 ), warnings
137 before = swarm_runtime._git_snapshot(run, checkout)
138 if before is None:
139 return swarm_review.failed_verdict(
140 seat, f"the git setup of {checkout} could not be read"
141 ), warnings
142 brief_file = Path(plan.prompt_path)
143 brief_file.parent.mkdir(parents=True, exist_ok=True)
144 brief_file.write_text(brief, encoding="utf-8")
145 gh_config = swarm_runtime.build_gh_config_path(
146 swarm_id, review_name(cluster_id, seat.slot), root
147 )
148 result = io.run_seat(
149 plan, swarm_worker.implementer_env(os.environ, gh_config_dir=str(gh_config))
150 )
151 if findings := swarm_runtime._tampered(run, checkout, before):
152 return swarm_review.failed_verdict(
153 seat,
154 swarm_worker.tamper_reason(findings, during="the reviewer ran"),
155 tampered=True,
156 ), warnings
157 return swarm_review.read_seat_verdict(
158 seat, result, head_sha=facts.head_sha, pr_title=facts.title
159 ), warnings
160 finally:
161 if created and not swarm_runtime.remove_swarm_worktree(root, checkout, runner=run):
162 warnings.append(
163 f"the review checkout {checkout} could not be removed; "
164 "`keel swarm-status <project.yaml> --clean` retries it"
165 )
168def _run_seats(
169 seats: tuple[ReviewSeat, ...],
170 briefs: dict[str, str],
171 *,
172 max_workers: int,
173 **kwargs: Any,
174) -> tuple[tuple[SeatVerdict, ...], list[str]]:
175 """Every eligible seat, at most ``max_workers`` at once; results in seat order."""
176 eligible = [s for s in seats if s.eligible]
177 results: dict[str, tuple[SeatVerdict, list[str]]] = {}
178 with concurrent.futures.ThreadPoolExecutor(
179 max_workers=max(1, min(max_workers, len(eligible)))
180 ) as pool:
181 futures = {
182 pool.submit(_run_seat, seat, briefs[seat.slot], **kwargs): seat for seat in eligible
183 }
184 for future in concurrent.futures.as_completed(futures):
185 seat = futures[future]
186 try:
187 results[seat.slot] = future.result()
188 except Exception as exc: # noqa: BLE001 - one seat must not end the review
189 results[seat.slot] = (
190 swarm_review.failed_verdict(
191 seat, f"the seat raised {type(exc).__name__}: {exc}"
192 ),
193 [],
194 )
195 verdicts = tuple(results[s.slot][0] for s in eligible)
196 warnings = [w for s in eligible for w in results[s.slot][1]]
197 return verdicts, warnings
200def _review_cluster(
201 cluster: SwarmCluster,
202 *,
203 swarm_id: str,
204 recorded: int | None,
205 dry_run: bool,
206 io: ReviewIO,
207 root: Path,
208 config: Any,
209 registry: Any,
210 host_agent: str,
211 seat_timeout: int,
212 max_workers: int,
213) -> ClusterReview:
214 report = ClusterReview(cluster.cluster_id, swarm_review.ERROR)
215 found = io.find_pull_request(
216 swarm_runtime.cluster_branch(swarm_id, cluster.cluster_id), recorded
217 )
218 if found.number is None:
219 status = swarm_review.MERGED if found.merged else swarm_review.SKIPPED
220 return swarm_review.with_status(report, status, found.reason)
221 number = found.number
222 report = replace(report, pull_request=number)
223 try:
224 facts = io.read_pull_request(number, not dry_run)
225 except Exception as exc: # noqa: BLE001 - one cluster fails, the wave goes on
226 return swarm_review.with_status(
227 report, swarm_review.ERROR, f"PR #{number} could not be read: {exc}"
228 )
229 required = swarm_review.required_count(facts.contract)
230 assignment = cluster.assignment or {}
231 seats = swarm_review.review_seats(
232 assignment,
233 config=config,
234 registry=registry,
235 implementer_vendor=swarm_review.seat_vendor(
236 assignment.get("implementer"), config=config, registry=registry, host_agent=host_agent
237 ),
238 brief_path=lambda slot: str(review_brief_path(swarm_id, cluster.cluster_id, slot, root)),
239 checkout_path=lambda slot: str(
240 review_checkout_path(swarm_id, cluster.cluster_id, slot, root)
241 ),
242 timeout=seat_timeout,
243 )
244 report = replace(
245 report, head_sha=facts.head_sha, tier=facts.tier, required=required, seats=seats
246 )
247 refusal = swarm_review.cluster_refusal(
248 seats,
249 panel=str(assignment.get("review_panel") or "reviewers"),
250 required=required,
251 require_distinct=swarm_review.distinct_vendors_required(facts.contract),
252 )
253 if refusal:
254 return swarm_review.with_status(report, swarm_review.REFUSED, refusal)
255 if dry_run:
256 return swarm_review.with_status(
257 report,
258 swarm_review.PLANNED,
259 f"a live run reviews PR #{number} at {facts.head_sha} with "
260 f"{sum(1 for s in seats if s.eligible)} seat(s)",
261 )
263 run = io.runner or default_runner
264 if why := _ensure_commit(root, facts.head_sha, number, io.remote, run):
265 return swarm_review.with_status(report, swarm_review.ERROR, why)
266 issues = [(n, *io.read_issue(n)) for n in cluster.issues]
267 briefs = {
268 seat.slot: swarm_review.render_review_brief(
269 swarm_id=swarm_id,
270 cluster_id=cluster.cluster_id,
271 pull_request=number,
272 facts=facts,
273 seat=seat,
274 focus=swarm_review.focus_for(facts.contract, seat.slot, index),
275 issues=issues,
276 )
277 for index, seat in enumerate(seats)
278 if seat.eligible
279 }
280 verdicts, warnings = _run_seats(
281 seats,
282 briefs,
283 max_workers=max_workers,
284 root=root,
285 swarm_id=swarm_id,
286 cluster_id=cluster.cluster_id,
287 facts=facts,
288 io=io,
289 run=run,
290 )
291 report = replace(report, verdicts=verdicts, warnings=tuple(warnings))
292 if why := swarm_review.posting_decision(verdicts, required=required):
293 return swarm_review.with_status(report, swarm_review.HELD, why)
294 try:
295 current = io.read_head(number)
296 except Exception as exc: # noqa: BLE001 - an unreadable head posts nothing
297 return swarm_review.with_status(
298 report, swarm_review.HELD, f"the head could not be re-read ({exc}); nothing is posted"
299 )
300 if why := swarm_review.head_moved(facts.head_sha, current):
301 return swarm_review.with_status(report, swarm_review.HELD, why)
302 items = swarm_review.posted_items(verdicts)
303 why = io.post_verdicts(
304 number, facts.head_sha, items, swarm_review.review_run_id(swarm_id, cluster.cluster_id)
305 )
306 if why:
307 return swarm_review.with_status(
308 report, swarm_review.HELD, f"keel review did not post: {why}"
309 )
310 status = swarm_review.posted_status(verdicts)
311 reason = f"{len(items)} verdict(s) posted on PR #{number}, pinned to {facts.head_sha}"
312 if status == swarm_review.POSTED_CHANGES_REQUESTED:
313 rejecting = ", ".join(v.slot for v in verdicts if v.outcome == swarm_review.REQUEST_CHANGES)
314 reason += (
315 f"; seat(s) {rejecting} request changes, so keel merge holds it on "
316 "review-verdict-not-approved until the findings are addressed and the head is "
317 "reviewed again"
318 )
319 return swarm_review.with_status(report, status, reason)
322def record_review(root: Path, swarm_id: str, review: ClusterReview, *, now: str) -> str:
323 """Write ``review`` into its cluster's worker record in the run state (#1440); ``""``
324 once written, else why it was not.
326 The file ``swarm-run`` writes, through the same atomic writer. Re-read just before the
327 write and written once per cluster, as each cluster ends, so a run that stops part-way
328 keeps the clusters it reviewed, and a record replaces the cluster's previous one whole.
329 Like ``swarm-land``, it takes no lock: do not run it beside a ``swarm-run`` of the same
330 swarm, whose writes would race it.
331 """
332 state = load_swarm_state(swarm_id, root=root)
333 if state is None:
334 return f"the run state of {swarm_id} could not be read; the review is not recorded"
335 if not any(w.cluster_id == review.cluster_id for w in state.workers):
336 return f"the run state of {swarm_id} has no worker {review.cluster_id} to record it on"
337 record = swarm_review.review_record(
338 review,
339 run_id=swarm_review.review_run_id(swarm_id, review.cluster_id),
340 reviewed_at=now,
341 )
342 save_swarm_state(update_worker_state(state, review.cluster_id, review=record), root=root)
343 return ""
346def review_wave_clusters(
347 plan: SwarmPlan,
348 wave_index: int,
349 *,
350 dry_run: bool,
351 io: ReviewIO,
352 root: str | Path = ".",
353 config: Any = None,
354 registry: Any = None,
355 host_agent: str = "claude",
356 seat_timeout: int = 1800,
357 max_workers: int = 3,
358) -> swarm_review.SwarmReviewResult:
359 """Review every cluster of one wave of ``plan``, in order.
361 A dry run reads each pull request and plans its seats, and runs, checks out and posts
362 nothing. A cluster that raises is a failed cluster; the next one is reviewed all the
363 same. The run's state file is only read, for the pull request each worker recorded.
364 """
365 root_path = Path(root).resolve()
366 wave = next((w for w in plan.waves if w.wave_index == wave_index), None)
367 if wave is None:
368 return swarm_review.SwarmReviewResult(
369 plan.swarm_id,
370 wave_index,
371 dry_run,
372 warnings=(f"the plan has no wave {wave_index}",),
373 )
374 state = load_swarm_state(plan.swarm_id, root=root_path)
375 recorded = {w.cluster_id: w.pull_request for w in state.workers} if state else {}
376 reviews: list[ClusterReview] = []
377 warnings: list[str] = []
378 for cluster in wave.clusters:
379 try:
380 reviewed = _review_cluster(
381 cluster,
382 swarm_id=plan.swarm_id,
383 recorded=recorded.get(cluster.cluster_id),
384 dry_run=dry_run,
385 io=io,
386 root=root_path,
387 config=config,
388 registry=registry,
389 host_agent=host_agent,
390 seat_timeout=seat_timeout,
391 max_workers=max_workers,
392 )
393 except Exception as exc: # noqa: BLE001 - one cluster must not end the wave
394 reviewed = ClusterReview(
395 cluster.cluster_id,
396 swarm_review.ERROR,
397 f"reviewing it raised {type(exc).__name__}: {exc}",
398 )
399 reviews.append(reviewed)
400 # A dry run records nothing; a live one records each cluster as it ends.
401 if not dry_run and (
402 why := record_review(root_path, plan.swarm_id, reviewed, now=swarm_runtime._now())
403 ):
404 warnings.append(f"{cluster.cluster_id}: {why}")
405 if not dry_run:
406 swarm_runtime.remove_empty_swarm_dirs(root_path, plan.swarm_id)
407 return swarm_review.SwarmReviewResult(
408 plan.swarm_id, wave_index, dry_run, tuple(reviews), tuple(warnings)
409 )