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
« 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.
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.
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.
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"""
22from __future__ import annotations
24import copy
25from collections.abc import Callable, Mapping, Sequence
26from pathlib import Path
27from typing import Any, NamedTuple
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
43#: A cluster's outcome at landing.
44LANDED = "landed"
45HELD = "held"
46FAILED = "failed"
49class PullRequestLookup(NamedTuple):
50 """The pull request a cluster lands through, or why there is none to land."""
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
58class ClusterMerge(NamedTuple):
59 """What ``keel merge`` did with one cluster's pull request.
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 """
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, ...] = ()
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]
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")
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"
95_ALREADY_MERGED = "already merged"
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 ""
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``.
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")
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.
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 )
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.
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}")
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 }
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).
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:
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.
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
252def is_landing_record(record: Mapping[str, Any]) -> bool:
253 """Whether ``record`` is a landing's own record — one :func:`landing_record` built.
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
262def _output(res: object) -> str:
263 return str(getattr(res, "stdout", "") or getattr(res, "output", "")).strip()
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
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.
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 )
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 )
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)
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``.
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.
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`.
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.
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 )
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 )
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)
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 )
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)
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)
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 )