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

1"""Keel Swarm Review Runtime — each cluster pull request reviewed by its own seats (#1423). 

2 

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. 

10 

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. 

17 

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. 

21 

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""" 

25 

26from __future__ import annotations 

27 

28import concurrent.futures 

29import os 

30from collections.abc import Callable 

31from dataclasses import dataclass, replace 

32from pathlib import Path 

33from typing import Any, cast 

34 

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 

47 

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] 

56 

57 

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) 

62 

63 

64@dataclass(frozen=True) 

65class ReviewIO: 

66 """The seams ``swarm-review`` reaches the outside world through.""" 

67 

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" 

78 

79 

80def review_name(cluster_id: str, slot: str) -> str: 

81 return f"{cluster_id}.review-{slot}" 

82 

83 

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) 

87 

88 

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) 

92 

93 

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. 

96 

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})" 

108 

109 

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 ) 

166 

167 

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 

198 

199 

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 ) 

262 

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) 

320 

321 

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. 

325 

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 "" 

344 

345 

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. 

360 

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 )