Coverage for src/keel/swarm_landing.py: 100%

178 statements  

« prev     ^ index     » next       coverage.py v7.16.2, created at 2026-10-02 20:26 +0000

1"""Keel Swarm Landing — a wave's clusters land as their pull requests, through keel merge. 

2 

3Each cluster's pull request is merged by the code path ``keel merge`` runs (#1287): the 

4merge window, the merge lock, the CI rollup, the evidence gate, the gates-pass pin, the 

5head-pinned squash and the post-merge drift check apply to it exactly as to any other 

6pull request. Nothing here merges, rebases or checks out a branch locally — the base 

7advances on the host and each pull request closes as merged — so the operator's 

8checkout is left where it was. 

9 

10Once a cluster's pull request merges, its issues are closed the way ``/keel:ship`` closes 

11them at s11–s12 (#1422): the closure comment, rendered by 

12:func:`keel.closure.render_closure_comment` from the run ledger's ``ship_run`` record, on the 

13pull request and on each issue, then each issue closed. A failure there is a warning: the 

14merge already happened and is never undone, and the next cluster still lands. 

15 

16This module is the per-wave loop and the reading of the answers it gets. The three things 

17it does to the outside world — find a cluster's pull request, hand one to keel merge, and 

18close a landed cluster's issues — are injected by the CLI, so the loop is tested against 

19recorded answers. 

20""" 

21 

22from __future__ import annotations 

23 

24import copy 

25from collections.abc import Callable, Mapping, Sequence 

26from pathlib import Path 

27from typing import Any, NamedTuple 

28 

29from .swarm import ( 

30 ClusterClosure, 

31 SwarmCluster, 

32 SwarmLandingResult, 

33 SwarmPlan, 

34 SwarmWave, 

35 evaluate_wave_landing_mode, 

36 load_swarm_state, 

37 render_dependent_wave_refusal, 

38 save_swarm_state, 

39 update_worker_state, 

40) 

41from .swarm_runtime import SubprocessRunner, cluster_branch, default_runner 

42 

43#: A cluster's outcome at landing. 

44LANDED = "landed" 

45HELD = "held" 

46FAILED = "failed" 

47 

48 

49class PullRequestLookup(NamedTuple): 

50 """The pull request a cluster lands through, or why there is none to land.""" 

51 

52 number: int | None 

53 reason: str = "" 

54 #: The cluster's pull request has merged already: nothing is left to land or review. 

55 merged: bool = False 

56 

57 

58class ClusterMerge(NamedTuple): 

59 """What ``keel merge`` did with one cluster's pull request. 

60 

61 ``outcome`` is :data:`LANDED`, :data:`HELD` or :data:`FAILED`; ``warning`` is set 

62 when the merge landed but needs a look (``keel merge`` found drift). 

63 """ 

64 

65 outcome: str 

66 reason: str 

67 warning: str = "" 

68 #: The head ``keel merge`` verified and merged — the one every pin was taken against. 

69 #: Empty when the payload names none (a dry run that stopped early, a refusal). 

70 head_sha: str = "" 

71 #: The heads ``head_sha`` descends from by capture commits alone, which the evidence 

72 #: gate and the gates-pass counted for it (``keel merge``'s ``covered_heads``). 

73 covered_heads: tuple[str, ...] = () 

74 

75 

76#: ``(branch, recorded pull request)`` -> the pull request to land. 

77FindPullRequest = Callable[[str, int | None], PullRequestLookup] 

78#: ``(pull request, cluster id, dry run)`` -> what keel merge did with it. 

79MergePullRequest = Callable[[int, str, bool], ClusterMerge] 

80#: ``(cluster, pull request, merge)`` -> how the landed cluster's issues were closed (#1422). 

81CloseCluster = Callable[[SwarmCluster, int, ClusterMerge], ClusterClosure] 

82 

83#: What a live ``swarm-land`` does, for its operator consent (#1422): ``keel merge``'s own 

84#: side effects, then the closure — a comment on the pull request and on each issue, and the 

85#: issue closed. The comment and the close need ``github``, which the merge needs already. 

86LANDING_SIDE_EFFECTS: tuple[str, ...] = ("git_worktree", "merge", "comments", "issue_close") 

87 

88#: The ``command`` of the ``ship_run`` record a landing appends after the merge. 

89LANDING_COMMAND = "swarm-land" 

90#: What that record says about the merge: it happened, through keel merge. Read as a merge 

91#: by :func:`keel.closeorder.record_attests_merge`, so closing the issue is not premature. 

92LANDING_MERGE_REASON = "merged by keel swarm-land through keel merge" 

93 

94 

95_ALREADY_MERGED = "already merged" 

96 

97 

98def _open_pull_request(reply: Mapping[str, Any], *, branch: str, base_branch: str) -> str: 

99 """Why the pull request in ``reply`` cannot be landed for ``branch``; ``""`` when it can.""" 

100 number = reply.get("number") 

101 state = reply.get("state") 

102 if state == "MERGED": 

103 return f"PR #{number} is {_ALREADY_MERGED} — this cluster landed in an earlier run" 

104 if state != "OPEN": 

105 return ( 

106 f"PR #{number} is {str(state).lower() or 'not open'}; reopen it, or open a new " 

107 f"one for {branch}" 

108 ) 

109 head = reply.get("headRefName") 

110 if head is not None and head != branch: 

111 return f"PR #{number} is for branch {head}, not the cluster's branch {branch}" 

112 base = reply.get("baseRefName") 

113 if base is not None and base != base_branch: 

114 return f"PR #{number} targets {base}, not the configured base branch {base_branch}" 

115 return "" 

116 

117 

118def pull_request_from_view( 

119 recorded: int, reply: object, *, branch: str, base_branch: str 

120) -> PullRequestLookup: 

121 """The pull request ``swarm-run`` recorded for a cluster, checked against ``gh pr view``. 

122 

123 The record says which pull request the worker opened; the host says whether it is 

124 still the open one for this branch and base. Anything else holds the cluster. 

125 """ 

126 if not isinstance(reply, Mapping) or reply.get("number") != recorded: 

127 return PullRequestLookup(None, f"PR #{recorded} (recorded by swarm-run) could not be read") 

128 why = _open_pull_request(reply, branch=branch, base_branch=base_branch) 

129 if not why: 

130 return PullRequestLookup(recorded) 

131 return PullRequestLookup(None, why, merged=reply.get("state") == "MERGED") 

132 

133 

134def pull_request_from_list(reply: object, *, branch: str) -> PullRequestLookup: 

135 """The one open pull request ``gh pr list --head <branch>`` names, or why there is none. 

136 

137 The fallback for a cluster whose run recorded no pull request (a state file from 

138 before #1287, or a pull request opened by hand). 

139 """ 

140 prs = [p for p in reply if isinstance(p, Mapping)] if isinstance(reply, list) else [] 

141 open_prs = [p for p in prs if p.get("state") == "OPEN" and isinstance(p.get("number"), int)] 

142 if len(open_prs) > 1: 

143 numbers = ", ".join(f"#{p['number']}" for p in open_prs) 

144 return PullRequestLookup( 

145 None, 

146 f"ambiguous: {len(open_prs)} open PRs for {branch} ({numbers}) — close the " 

147 "strays before landing", 

148 ) 

149 if open_prs: 

150 return PullRequestLookup(open_prs[0]["number"]) 

151 merged = [p for p in prs if p.get("state") == "MERGED" and isinstance(p.get("number"), int)] 

152 if merged: 

153 return PullRequestLookup( 

154 None, 

155 f"PR #{merged[0]['number']} is {_ALREADY_MERGED} — this cluster landed earlier", 

156 merged=True, 

157 ) 

158 return PullRequestLookup( 

159 None, 

160 f"no open pull request for {branch} — swarm-run --live opens one per cluster", 

161 ) 

162 

163 

164def merge_outcome(code: int, payload: Mapping[str, Any] | None) -> ClusterMerge: 

165 """Read ``keel merge``'s exit code and payload as one cluster's landing outcome. 

166 

167 ``None`` means ``keel merge`` refused before it reached the pull request — consent, 

168 a missing ``git``/``gh``, the config — and printed why on stderr. A payload with a 

169 ``merge_output`` is one where the merge call was made: when that failed, the 

170 cluster *failed*; every other refusal *holds* it with ``keel merge``'s own reason. 

171 """ 

172 if payload is None: 

173 return ClusterMerge( 

174 HELD, "keel merge refused before reaching the pull request; its reason is on stderr" 

175 ) 

176 reason = str(payload.get("reason") or "no reason given") 

177 if code == 0: 

178 return ClusterMerge(LANDED, reason, **_merged_head(payload)) 

179 if code == 3: 

180 return ClusterMerge( 

181 LANDED, 

182 reason, 

183 f"PR #{payload.get('pull_request')} merged, but keel merge detected drift — " 

184 "run keel verify-merge on it", 

185 **_merged_head(payload), 

186 ) 

187 if "merge_output" in payload: 

188 return ClusterMerge(FAILED, f"keel merge: {reason}: {payload['merge_output']}"[:300]) 

189 return ClusterMerge(HELD, f"keel merge: {reason}") 

190 

191 

192def _merged_head(payload: Mapping[str, Any]) -> dict[str, Any]: 

193 """The head ``keel merge``'s evidence gate verified — the head it merged — and the heads 

194 that gate counted for it, as :class:`ClusterMerge` fields.""" 

195 block = payload.get("evidence") 

196 block = block if isinstance(block, Mapping) else {} 

197 head = block.get("head_sha") 

198 covered = block.get("covered_heads") 

199 return { 

200 "head_sha": head if isinstance(head, str) else "", 

201 "covered_heads": tuple(str(sha) for sha in covered) if isinstance(covered, list) else (), 

202 } 

203 

204 

205def landing_record( 

206 prior: Mapping[str, Any], 

207 *, 

208 head_sha: str, 

209 reviewers: Sequence[str], 

210 run_context: Mapping[str, Any], 

211) -> dict[str, Any]: 

212 """The ``ship_run`` record of a cluster ``swarm-land`` merged (#1422). 

213 

214 ``/keel:ship`` renders its closure comment from the ``ship_run`` record it appends at 

215 s11, and ``evidence-verify`` holds the posted comment to that record's render. A landing 

216 does the same with ``prior`` — the record whose gates-pass ``keel merge`` accepted for 

217 ``head_sha``, which a live worker wrote when it opened the pull request (#1420) — and 

218 changes only what the landing knows: 

219 

220 - ``command`` is :data:`LANDING_COMMAND`, and ``assessment.merge`` is ``merge`` with 

221 :data:`LANDING_MERGE_REASON`: the pull request merged. Every other assessment field 

222 stays as recorded — the landing assessed no tier, window or CI of its own. 

223 - ``git.head_sha`` is the head keel merge merged. 

224 - ``actors.reviewers`` are the reviewers whose verdicts the evidence gate counted for 

225 that head, when any are named; otherwise what the record already said. 

226 - ``run_context`` is this landing's: its host agent, transport and operator consent. 

227 

228 The gates, the changed files, the implementer, the issue and the capture block are 

229 carried as recorded, so the record still passes for the head 

230 (:func:`keel.ledger.gates_pass_for_head`) and its attribution still matches the 

231 labels. Pure: ``prior`` is copied, never changed. 

232 """ 

233 record = copy.deepcopy(dict(prior)) 

234 record["command"] = LANDING_COMMAND 

235 git = record.get("git") 

236 record["git"] = {**(git if isinstance(git, dict) else {}), "head_sha": head_sha} 

237 assessment = record.get("assessment") 

238 record["assessment"] = { 

239 **(assessment if isinstance(assessment, dict) else {}), 

240 "merge": {"action": "merge", "reason": LANDING_MERGE_REASON}, 

241 } 

242 if reviewers: 

243 actors = record.get("actors") 

244 record["actors"] = { 

245 **(actors if isinstance(actors, dict) else {}), 

246 "reviewers": list(reviewers), 

247 } 

248 record["run_context"] = copy.deepcopy(dict(run_context)) 

249 return record 

250 

251 

252def is_landing_record(record: Mapping[str, Any]) -> bool: 

253 """Whether ``record`` is a landing's own record — one :func:`landing_record` built. 

254 

255 A landing that finds the record for the merged head is already its own renders from it 

256 rather than appending a second one, so a closure retried for the same merge posts the 

257 same comment again (edited in place) instead of a new one. 

258 """ 

259 return record.get("command") == LANDING_COMMAND 

260 

261 

262def _output(res: object) -> str: 

263 return str(getattr(res, "stdout", "") or getattr(res, "output", "")).strip() 

264 

265 

266def _read_head(repo_root: Path, runner: SubprocessRunner | None) -> tuple[str, str] | None: 

267 """Where the checkout is: ``("branch", name)``, ``("detached", sha)`` or ``None``.""" 

268 run = runner or default_runner 

269 res = run(["git", "symbolic-ref", "--quiet", "--short", "HEAD"], repo_root) 

270 if res.ok and _output(res): 

271 return ("branch", _output(res)) 

272 res = run(["git", "rev-parse", "--verify", "--quiet", "HEAD"], repo_root) 

273 if res.ok and _output(res): 

274 return ("detached", _output(res)) 

275 return None 

276 

277 

278def _return_to( 

279 repo_root: Path, head: tuple[str, str] | None, runner: SubprocessRunner | None 

280) -> str | None: 

281 """Put the checkout back where the operator had it; a warning when that fails. 

282 

283 Landing merges on the host and checks nothing out, so this is a postcondition, not 

284 an undo: the promise #1279 made — the operator ends where they started — is checked 

285 after every wave rather than assumed of every step keel merge takes. 

286 """ 

287 if head is None or _read_head(repo_root, runner) == head: 

288 return None 

289 ref = head[1] 

290 run = runner or default_runner 

291 # A branch name checks that branch out; a sha detaches at it, as it started. 

292 # `--` so a file that happens to share the name is never read as a path. 

293 res = run(["git", "checkout", ref, "--"], repo_root) 

294 if res.ok: 

295 return None 

296 return ( 

297 f"could not return the checkout to {ref} ({_output(res) or 'no output'}); " 

298 f"run `git checkout {ref}` by hand" 

299 ) 

300 

301 

302def _close_landed( 

303 close_cluster: CloseCluster, cluster: SwarmCluster, number: int, merged: ClusterMerge 

304) -> ClusterClosure: 

305 """``close_cluster`` for one landed cluster; whatever it raises is a warning (#1422).""" 

306 try: 

307 return close_cluster(cluster, number, merged) 

308 except Exception as exc: # noqa: BLE001 - the merge happened; closing it never undoes it 

309 return ClusterClosure( 

310 cluster.cluster_id, 

311 number, 

312 cluster.issues, 

313 warnings=(f"closing the issues raised {type(exc).__name__}: {exc}",), 

314 ) 

315 

316 

317def _target_wave(plan: SwarmPlan, wave_index: int) -> SwarmWave | None: 

318 return next((w for w in plan.waves if w.wave_index == wave_index), None) 

319 

320 

321def land_wave_clusters( 

322 plan: SwarmPlan, 

323 wave_index: int, 

324 *, 

325 dry_run: bool, 

326 find_pull_request: FindPullRequest, 

327 merge_pull_request: MergePullRequest, 

328 root: str | Path = ".", 

329 runner: SubprocessRunner | None = None, 

330 close_cluster: CloseCluster | None = None, 

331) -> SwarmLandingResult: 

332 """Land one wave: each cluster's pull request, in order, through ``keel merge``. 

333 

334 A ``sequential_dependent`` wave is refused before anything is looked up (#1276). 

335 Otherwise every cluster is independent of the others: its pull request is found 

336 (the number the run recorded, else the one open pull request for its branch) and 

337 handed to ``merge_pull_request`` — keel merge, which claims the merge lock for that 

338 one merge. A cluster with no pull request, or one keel merge refuses, is held with 

339 the reason; the next cluster is tried all the same. A dry run asks keel merge for 

340 its own dry run, so it reports what the live landing would do and merges nothing. 

341 

342 ``runner`` is an injection seam for tests, like the ``_run`` seams elsewhere in keel: 

343 ``swarm-land`` never passes it. It runs the git commands that read where the checkout 

344 is and put it back (:func:`_read_head`, :func:`_return_to`); left out, they run through 

345 :func:`keel.swarm_runtime.default_runner`. 

346 

347 ``close_cluster`` closes a landed cluster's issues (#1422): called once per cluster 

348 that merged, never for a held or failed one, and never in a dry run — there each cluster 

349 that would land reports, with an empty :class:`ClusterClosure`, what would be closed. 

350 What it could not do is a warning, as is anything it raises: the merge stands, and the 

351 next cluster lands all the same. Left out, nothing is closed and nothing is reported. 

352 

353 The wave's ``mode`` is judged on its clusters' planned scopes. A ``pr_diff_map`` of 

354 real diffs used to be accepted here and no caller passed one; it could only relabel 

355 ``mode``, because every cluster lands through its own ``keel merge`` whatever the mode 

356 says, so it is gone (#1280). 

357 """ 

358 root_path = Path(root).resolve() 

359 target_wave = _target_wave(plan, wave_index) 

360 if target_wave is None: 

361 return SwarmLandingResult( 

362 swarm_id=plan.swarm_id, 

363 wave_index=wave_index, 

364 mode="none", 

365 landed_clusters=(), 

366 failed_clusters=(), 

367 status="failed", 

368 ) 

369 

370 decision = evaluate_wave_landing_mode(target_wave, {}) 

371 if decision.mode == "refused": 

372 # A dependent wave's branches predate the landing it depends on. It is refused 

373 # before any lookup or merge, dry run or live, so a preview reports exactly 

374 # what the live run would do (#1276). 

375 return SwarmLandingResult( 

376 swarm_id=plan.swarm_id, 

377 wave_index=wave_index, 

378 mode=decision.mode, 

379 landed_clusters=(), 

380 failed_clusters=(), 

381 status="failed", 

382 refused=render_dependent_wave_refusal(target_wave), 

383 ) 

384 

385 landed: list[str] = [] 

386 failed: list[str] = [] 

387 held: list[tuple[str, str]] = [] 

388 warnings: list[str] = [] 

389 prs: list[tuple[str, int]] = [] 

390 closures: list[ClusterClosure] = [] 

391 state = load_swarm_state(plan.swarm_id, root=root_path) 

392 # Only the pull request each worker opened is read from the run state. The review record 

393 # `swarm-review` leaves beside it (#1440) is a report for `swarm-status`, and is not read 

394 # here: whether a cluster is reviewed is `keel merge`'s evidence gate's question, asked of 

395 # the verdicts on the pull request at its current head — the only authority. 

396 recorded = {w.cluster_id: w.pull_request for w in state.workers} if state else {} 

397 home = _read_head(root_path, runner) 

398 

399 def record(cluster_id: str, status: str, details: str, number: int | None = None) -> None: 

400 nonlocal state 

401 # Kept in memory only: a dry run never saves it (below). 

402 if state: 

403 state = update_worker_state( 

404 state, 

405 cluster_id, 

406 step="s10", 

407 status=status, 

408 details=details, 

409 pull_request=number, 

410 ) 

411 

412 try: 

413 for c in target_wave.clusters: 

414 branch = cluster_branch(plan.swarm_id, c.cluster_id) 

415 found = find_pull_request(branch, recorded.get(c.cluster_id)) 

416 if found.number is None: 

417 held.append((c.cluster_id, found.reason)) 

418 record(c.cluster_id, "held", found.reason) 

419 continue 

420 number = found.number 

421 prs.append((c.cluster_id, number)) 

422 try: 

423 merged = merge_pull_request(number, c.cluster_id, dry_run) 

424 except Exception as exc: # noqa: BLE001 - one cluster fails, the wave goes on 

425 merged = ClusterMerge(FAILED, f"keel merge raised {type(exc).__name__}: {exc}") 

426 if merged.outcome == LANDED: 

427 landed.append(c.cluster_id) 

428 record(c.cluster_id, "merged", f"PR #{number}: {merged.reason}", number) 

429 if merged.warning: 

430 warnings.append(f"{c.cluster_id}: {merged.warning}") 

431 if close_cluster is None: 

432 continue 

433 if dry_run: 

434 closures.append(ClusterClosure(c.cluster_id, number, c.issues, dry_run=True)) 

435 continue 

436 closure = _close_landed(close_cluster, c, number, merged) 

437 closures.append(closure) 

438 warnings.extend(f"{c.cluster_id}: {warning}" for warning in closure.warnings) 

439 elif merged.outcome == FAILED: 

440 failed.append(c.cluster_id) 

441 record(c.cluster_id, "failed", f"PR #{number}: {merged.reason}", number) 

442 else: 

443 held.append((c.cluster_id, f"PR #{number}: {merged.reason}")) 

444 record(c.cluster_id, "held", f"PR #{number}: {merged.reason}", number) 

445 finally: 

446 warning = _return_to(root_path, home, runner) 

447 if warning: 

448 warnings.append(warning) 

449 

450 # A dry run changes nothing, the run's state included. 

451 if state and not dry_run: 

452 save_swarm_state(state, root=root_path) 

453 

454 # A held cluster is not landed, so it can never leave the wave "success"; it is 

455 # also not a failure of the merge itself, so it degrades the status exactly like a 

456 # failed cluster without being reported as one. 

457 not_landed = len(failed) + len(held) 

458 overall_status = ( 

459 "success" 

460 if not_landed == 0 and len(landed) > 0 

461 else ("partial_failure" if len(landed) > 0 else "failed") 

462 ) 

463 return SwarmLandingResult( 

464 swarm_id=plan.swarm_id, 

465 wave_index=wave_index, 

466 mode=decision.mode, 

467 landed_clusters=tuple(landed), 

468 failed_clusters=tuple(failed), 

469 status=overall_status, 

470 held_clusters=tuple(held), 

471 warnings=tuple(warnings), 

472 pull_requests=tuple(prs), 

473 closures=tuple(closures), 

474 )