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

499 statements  

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

1"""Keel Swarm Runtime — Isolated multi-worktree execution & cluster orchestration. 

2 

3Thin I/O execution layer for running parallel swarm workers in isolated Git worktrees, 

4handling worker state persistence, and managing fail-soft rebalancing across waves. 

5 

6A dry run's worker is a ``keel ship --dry-run`` assessment per cluster. A live run's worker 

7(#1400) is the cluster's implementer seat, dispatched through :mod:`keel.delegaterun` in 

8the cluster's own worktree, followed by a commit, the project's gates, a push and one pull 

9request, stamped with the ship-provenance comment a live ``keel ship`` run posts on its own, 

10labelled with the seat's attribution, and with the gates it ran recorded in the run ledger 

11for the head it pushed (#1420) — every mutation behind the consent scopes the parent 

12delegated (:mod:`keel.swarm_worker`, where each of those decisions is made). 

13""" 

14 

15from __future__ import annotations 

16 

17import concurrent.futures 

18import datetime 

19import functools 

20import hashlib 

21import os 

22import shutil 

23import subprocess # nosec B404 

24import sys 

25import tempfile 

26import threading 

27import traceback 

28from collections.abc import Callable, Mapping 

29from dataclasses import dataclass, field, replace 

30from pathlib import Path 

31from typing import Any 

32 

33from . import delegaterun, doctor, github, ledger, redaction, runner, swarm_worker 

34from .delegate import RunPlan 

35from .runner import CommandResult 

36from .swarm import ( 

37 SwarmCluster, 

38 SwarmPlan, 

39 SwarmRunResult, 

40 SwarmRunState, 

41 SwarmWorkerStatus, 

42 load_swarm_state, 

43 pull_request_number, 

44 rebalance_swarm_plan, 

45 save_swarm_state, 

46 ship_handoff_args, 

47 tail_child_output, 

48 update_worker_state, 

49 worker_seed, 

50) 

51 

52SubprocessRunner = Callable[[list[str], Path], CommandResult] 

53#: Runs one implementer plan with the given child environment; returns the delegate result. 

54Implementer = Callable[[RunPlan, dict[str, str]], dict[str, Any]] 

55#: ``(argv, cwd)`` -> the push's result. The worker builds the argv — its hardened ``git 

56#: push --no-verify <url> <commit>:<ref>`` (:func:`keel.swarm_worker.keel_git`), to the URL 

57#: it read before the implementer ran — and the seat runs it in keel's own process, with the 

58#: operator's environment, which is where the credentials are. 

59Pusher = Callable[[list[str], Path], CommandResult] 

60#: ``(title, body, base, head, cwd)`` -> ``gh pr create``'s result. 

61PullRequestOpener = Callable[[str, str, str, str, Path], CommandResult] 

62#: ``(owner_repo, number, body, cwd)`` -> the result of posting ``body`` on that pull request. 

63CommentPoster = Callable[[str, int, str, Path], CommandResult] 

64#: ``(owner_repo, number, labels, cwd)`` -> the result of applying ``labels`` to that pull 

65#: request, creating any the repository lacks first (#1420). 

66PullRequestLabeler = Callable[[str, int, list[str], Path], CommandResult] 

67 

68#: The runner's own wall-clock limit, for the short git commands swarm and canary run 

69#: through it (worktree add/remove, checkout, merge). A worker's child ``keel ship`` is 

70#: not one of them: it runs the gate suite, and gets the swarm's worker timeout (#1279). 

71DEFAULT_RUNNER_TIMEOUT_S = 300 

72 

73 

74def default_runner( 

75 cmd: list[str], 

76 cwd: Path, 

77 timeout_s: int = DEFAULT_RUNNER_TIMEOUT_S, 

78 env: dict[str, str] | None = None, 

79) -> CommandResult: 

80 """Run a subprocess command in cwd and return a CommandResult. 

81 

82 A command still running after ``timeout_s`` seconds is killed and reported 

83 ``timed_out`` with ``code=124``, the gate runner's code for the same case. 

84 ``env`` replaces the child's environment when given (``None`` inherits it). 

85 """ 

86 try: 

87 proc = subprocess.run( # nosec B603 

88 cmd, 

89 cwd=cwd, 

90 env=env, 

91 stdin=subprocess.DEVNULL, 

92 stdout=subprocess.PIPE, 

93 stderr=subprocess.STDOUT, 

94 text=True, 

95 timeout=timeout_s, 

96 check=False, 

97 ) 

98 return CommandResult( 

99 ok=(proc.returncode == 0), 

100 code=proc.returncode, 

101 output=proc.stdout or "", 

102 timed_out=False, 

103 stdout=proc.stdout or "", 

104 stderr="", 

105 ) 

106 except subprocess.TimeoutExpired as exc: 

107 out = ( 

108 exc.output 

109 if isinstance(exc.output, str) 

110 else (exc.stdout if isinstance(exc.stdout, str) else "") 

111 ) 

112 return CommandResult( 

113 ok=False, 

114 code=124, 

115 output=out or "", 

116 timed_out=True, 

117 stdout=out or "", 

118 stderr="", 

119 ) 

120 except Exception as exc: # noqa: BLE001 

121 return CommandResult( 

122 ok=False, 

123 code=1, 

124 output=str(exc), 

125 timed_out=False, 

126 stdout="", 

127 stderr=str(exc), 

128 ) 

129 

130 

131#: Called with each :data:`keel.swarm_worker.STAGES` entry as a live worker enters it. 

132StageReporter = Callable[[str], None] 

133 

134 

135def _no_stage(_stage: str) -> None: 

136 """The stage reporter of a worker nobody watches.""" 

137 

138 

139def _now() -> str: 

140 """The current UTC time, ISO 8601 — the timestamps a worker record carries.""" 

141 return datetime.datetime.now(datetime.UTC).isoformat() 

142 

143 

144def build_worktree_path(swarm_id: str, cluster_id: str, root: str | Path = ".") -> Path: 

145 """Return the isolated worktree filesystem path for a specific swarm cluster worker.""" 

146 return Path(root) / ".keel" / "worktrees" / swarm_id / cluster_id 

147 

148 

149def build_brief_path(swarm_id: str, cluster_id: str, root: str | Path = ".") -> Path: 

150 """Where a live worker's implementer brief is written: beside the run's state file, 

151 outside the cluster's worktree, so the brief is never committed with the work.""" 

152 return Path(root) / ".keel" / "state" / "swarm" / swarm_id / f"{cluster_id}.brief.md" 

153 

154 

155def build_gh_config_path(swarm_id: str, cluster_id: str, root: str | Path = ".") -> Path: 

156 """The ``GH_CONFIG_DIR`` a live worker's implementer runs with: a directory beside the 

157 brief, outside the worktree, that holds no ``gh`` login (#1400).""" 

158 return Path(root) / ".keel" / "state" / "swarm" / swarm_id / f"{cluster_id}.no-gh-login" 

159 

160 

161def cluster_branch(swarm_id: str, cluster_id: str) -> str: 

162 """The branch a cluster's worktree is cut on, and the head of its pull request.""" 

163 return f"swarm/{swarm_id}/{cluster_id}" 

164 

165 

166def _default_implement(plan: RunPlan, env: dict[str, str]) -> dict[str, Any]: 

167 """``keel delegate run``'s executor, with the worker's scrubbed environment.""" 

168 return delegaterun.execute(plan, _run=functools.partial(runner.run_argv, env=env)) 

169 

170 

171def _default_push(argv: list[str], cwd: Path) -> CommandResult: 

172 """The worker's push, with the operator's environment (``env`` not passed).""" 

173 return runner.run_argv(argv, cwd=str(cwd)) 

174 

175 

176def _default_open_pr(title: str, body: str, base: str, head: str, cwd: Path) -> CommandResult: 

177 return github.open_pr(title, body, base, head, cwd=str(cwd)) 

178 

179 

180def _default_post_comment(owner_repo: str, number: int, body: str, cwd: Path) -> CommandResult: 

181 """The call ``keel post-comment`` makes for a live ship run's provenance, with the 

182 operator's environment — so the comment's author is the account keel's ``gh`` is.""" 

183 return github.post_issue_comment(owner_repo, number, body, cwd=str(cwd)) 

184 

185 

186def _default_label_pr(owner_repo: str, number: int, labels: list[str], cwd: Path) -> CommandResult: 

187 """Apply ``labels`` to the pull request with the operator's environment, creating each 

188 one the repository lacks first, as ``keel doctor --fix`` creates them (#1420). 

189 

190 The answer is the label post's: a failed listing or create is not, since a sibling 

191 worker may have created the same label a moment earlier — the post then still succeeds, 

192 and when the label really is missing, the post is what fails and says so. 

193 """ 

194 listed = github.list_labels(repo=owner_repo, cwd=str(cwd)) 

195 existing = swarm_worker.label_names(listed.stdout or listed.output) if listed.ok else () 

196 for name in doctor.missing_labels(labels, existing): 

197 github.create_label(name, repo=owner_repo, cwd=str(cwd)) 

198 return github.add_issue_labels(owner_repo, number, labels, cwd=str(cwd)) 

199 

200 

201#: The key a live worker hands its built ``ship_run`` record to the run under (#1420). The 

202#: run appends it on the one thread that collects results, then drops the key. 

203GATES_RECORD_KEY = "gates_record" 

204 

205 

206@dataclass(frozen=True) 

207class LiveRun: 

208 """What a live run's parent hands its workers: the operator's consent, delegated, and 

209 each cluster's planned implementer seat. The five callables are the I/O seams.""" 

210 

211 consent: swarm_worker.ConsentDelegation 

212 #: ``cluster_id`` -> the cluster's implementer :class:`keel.delegate.RunPlan`. 

213 dispatches: Mapping[str, RunPlan] 

214 #: ``issue`` -> its scope, for the brief and the pull request's title. 

215 issue_scopes: Mapping[int, Any] = field(default_factory=dict) 

216 #: ``cluster_id`` -> where its implementer seat came from (``assignment`` source). 

217 seat_sources: Mapping[str, str] = field(default_factory=dict) 

218 remote: str = "origin" 

219 #: The run ledger the parent recorded the delegation in (``swarm-run --live`` writes 

220 #: its ``delegated`` event before any worker starts); each cluster pull request a 

221 #: worker opens is recorded there too (#1400). ``None`` records nothing. 

222 ledger_path: Path | None = None 

223 #: The project config, for the capture redaction a ledger record passes before it is 

224 #: written (:func:`keel.ledger.sanitize_record`); ``None`` applies the default policy. 

225 config: Any = None 

226 implement: Implementer = _default_implement 

227 push: Pusher = _default_push 

228 open_pr: PullRequestOpener = _default_open_pr 

229 post_comment: CommentPoster = _default_post_comment 

230 label_pr: PullRequestLabeler = _default_label_pr 

231 

232 

233def create_swarm_worktree( 

234 repo_root: Path, 

235 worktree_path: Path, 

236 branch_name: str, 

237 base_branch: str = "main", 

238 runner: SubprocessRunner | None = None, 

239) -> bool: 

240 """Create an isolated git worktree for a swarm cluster worker.""" 

241 worktree_path.parent.mkdir(parents=True, exist_ok=True) 

242 run = runner or default_runner 

243 cmd = [ 

244 "git", 

245 "worktree", 

246 "add", 

247 "-B", 

248 branch_name, 

249 str(worktree_path), 

250 base_branch, 

251 ] 

252 res = run(cmd, repo_root) 

253 return res.ok 

254 

255 

256def remove_swarm_worktree( 

257 repo_root: Path, 

258 worktree_path: Path, 

259 runner: SubprocessRunner | None = None, 

260) -> bool: 

261 """Remove a previously created isolated git worktree; ``True`` only when it is gone. 

262 

263 ``git worktree remove --force`` first. Whatever it leaves on disk is removed with 

264 ``rmtree``, and when git refused, ``git worktree prune`` then drops the registration of 

265 a directory that is now gone — a registration left behind makes the next ``git 

266 worktree add`` of the same branch fail. The answer is the directory's own state 

267 afterwards: this used to return ``True`` unconditionally, even when nothing was 

268 removed (#1278). 

269 """ 

270 run = runner or default_runner 

271 res = run(["git", "worktree", "remove", "--force", str(worktree_path)], repo_root) 

272 if worktree_path.exists(): 

273 shutil.rmtree(worktree_path, ignore_errors=True) 

274 if not res.ok: 

275 run(["git", "worktree", "prune"], repo_root) 

276 return not worktree_path.exists() 

277 

278 

279def prune_worktrees(repo_root: Path, runner: SubprocessRunner | None = None) -> str: 

280 """Run ``git worktree prune``; ``""``, or why it failed (#1278). 

281 

282 It drops the registration of every worktree whose directory is gone — git's own 

283 ``gc`` does the same after ``gc.worktreePruneExpire`` — and touches no directory, no 

284 branch, and no worktree that is locked. 

285 """ 

286 res = (runner or default_runner)(["git", "worktree", "prune"], repo_root) 

287 return "" if res.ok else (res.output.strip() or f"exit {res.code}") 

288 

289 

290def remove_empty_swarm_dirs(root: Path, swarm_id: str) -> None: 

291 """Remove ``.keel/worktrees/<swarm_id>/``, then ``.keel/worktrees/``, when each is empty. 

292 

293 ``rmdir`` only — a directory with anything in it (a sibling worker's worktree, one kept 

294 for inspection) stays, and so does a symlink (#1278). 

295 """ 

296 base = Path(root) / ".keel" / "worktrees" 

297 for directory in (base / swarm_id, base): 

298 try: 

299 # Refuses a directory with anything in it, a missing one, and a symlink. 

300 directory.rmdir() 

301 except OSError: 

302 continue 

303 

304 

305def settle_live_worktree( 

306 root: Path, 

307 worktree_path: Path, 

308 branch: str, 

309 disposal: swarm_worker.WorktreeDisposal, 

310 runner: SubprocessRunner | None = None, 

311) -> dict[str, Any]: 

312 """Carry out a :class:`keel.swarm_worker.WorktreeDisposal`, and say what happened. 

313 

314 Returns the keys a worker's result gains: ``worktree`` (where it is still on disk, or 

315 ``None``), ``worktree_state`` (``removed``, ``kept``, ``remove-failed`` or ``none``), 

316 ``branch_deleted`` and ``warnings`` — so a removal that failed is recorded, not 

317 swallowed (#1278). 

318 """ 

319 warnings: list[str] = [] 

320 branch_deleted = False 

321 state = "kept" if disposal.worktree == "keep" else "none" 

322 if disposal.worktree == "remove": 

323 removed = remove_swarm_worktree(root, worktree_path, runner=runner) 

324 if not removed: 

325 warnings.append( 

326 f"the worktree {worktree_path} could not be removed; " 

327 "`keel swarm-status --clean` retries it" 

328 ) 

329 elif disposal.delete_branch: 

330 deleted = (runner or default_runner)(["git", "branch", "-D", branch], root) 

331 branch_deleted = deleted.ok 

332 if not deleted.ok: 

333 warnings.append( 

334 f"the branch {branch} could not be deleted: {deleted.output.strip()}" 

335 ) 

336 state = "removed" if removed else "remove-failed" 

337 return { 

338 # A kept worktree, or one whose removal failed (`remove_swarm_worktree` answers 

339 # from the directory's own state, so it is still there). 

340 "worktree": str(worktree_path) if state in ("kept", "remove-failed") else None, 

341 "worktree_state": state, 

342 "branch_deleted": branch_deleted, 

343 "warnings": warnings, 

344 } 

345 

346 

347def execute_cluster_worker( 

348 project_yaml: str, 

349 issue: int, 

350 root: Path, 

351 worktree_dir: Path, 

352 *, 

353 dry_run: bool = True, 

354 role: str = "core", 

355 extra_args: list[str] | None = None, 

356 runner: SubprocessRunner | None = None, 

357 timeout_s: int = DEFAULT_RUNNER_TIMEOUT_S, 

358) -> dict[str, Any]: 

359 """Execute a single cluster worker pipeline (keel ship) within its worktree. 

360 

361 ``timeout_s`` bounds the child when keel runs it (no ``runner`` passed): the child 

362 runs the project's gate suite even in a dry run, so the runner's own 300 s cut off 

363 a suite the project allows longer (#1279). A passed ``runner`` owns its own bound. 

364 """ 

365 run = runner or functools.partial(default_runner, timeout_s=timeout_s) 

366 cmd = [ 

367 sys.executable, 

368 "-m", 

369 "keel", 

370 "ship", 

371 project_yaml, 

372 "--root", 

373 str(worktree_dir if worktree_dir.exists() else root), 

374 "--issue", 

375 str(issue), 

376 "--json", 

377 ] 

378 # `keel ship` gates every live-only path on `--live`; leaving out `--dry-run` does 

379 # not turn it on, so a live swarm's children ran the dry assessment (#1269). 

380 cmd.append("--dry-run" if dry_run else "--live") 

381 if extra_args: 

382 cmd.extend(extra_args) 

383 

384 target_cwd = worktree_dir if worktree_dir.exists() else root 

385 res = run(cmd, target_cwd) 

386 output = res.stdout or res.output 

387 if res.timed_out: 

388 # Said last, where the tail the swarm keeps of the output still holds it: a 

389 # killed child printed no verdict, and its partial output read like a failure 

390 # of the change rather than of the budget (#1279). 

391 prefix = f"{output.rstrip()}\n" if output.strip() else "" 

392 output = ( 

393 f"{prefix}keel ship timed out after {timeout_s}s (exit 124) and returned no " 

394 "verdict; raise the limit with swarm-run --worker-timeout if the gate suite " 

395 "legitimately needs longer." 

396 ) 

397 

398 return { 

399 "issue": issue, 

400 "role": role, 

401 "ok": res.ok, 

402 "code": res.code, 

403 "timed_out": res.timed_out, 

404 "output": output, 

405 } 

406 

407 

408def _last_line(result: CommandResult) -> str: 

409 """The last line a successful command printed; ``""`` for a failed one.""" 

410 lines = (result.stdout or result.output or "").strip().splitlines() if result.ok else [] 

411 return lines[-1].strip() if lines else "" 

412 

413 

414def execute_live_cluster_worker( 

415 cluster: SwarmCluster, 

416 *, 

417 swarm_id: str, 

418 root: Path, 

419 worktree_dir: Path, 

420 project_yaml: str, 

421 base_branch: str, 

422 live: LiveRun, 

423 runner: SubprocessRunner | None = None, 

424 timeout_s: int = DEFAULT_RUNNER_TIMEOUT_S, 

425 progress: dict[str, bool] | None = None, 

426 on_stage: StageReporter = _no_stage, 

427) -> dict[str, Any]: 

428 """Implement one cluster live: its implementer seat, a commit, the gates, a push, a PR. 

429 

430 ``progress``, when given, is filled in as the worker goes — ``worktree_created`` once 

431 its worktree exists, ``implementer_ran`` once the seat is started — so a caller can 

432 settle the worktree correctly even when the worker raises part-way (#1278). 

433 

434 ``on_stage`` is called with each stage of :data:`keel.swarm_worker.STAGES` as the 

435 worker enters it, so ``keel swarm-status`` can show how far a worker has got while it 

436 runs (#1280); the stage it ends at is the result's ``stage``. 

437 

438 Before its first mutation the worker checks that the scopes the parent delegated to 

439 *this* cluster cover every mutation it will make 

440 (:func:`keel.swarm_worker.consent_refusal`), and does nothing when they do not. It then 

441 stops at the first stage that fails — so a failed implementer or a red gate leaves the 

442 branch unpushed and no pull request opened, and the result says which stage stopped it 

443 and why. ``timeout_s`` bounds the gate run and the git commands; the implementer runs 

444 under its own plan's timeout. 

445 

446 The worker's children — the implementer, git, the gates — run without the parent's 

447 consent variables or forge tokens (:func:`keel.swarm_worker.worker_env`): consent 

448 reaches a worker as the delegation it is handed, never through inheritance. The 

449 implementer also runs without the forge (:func:`keel.swarm_worker.implementer_env`): no 

450 ``gh`` login, no git credential helper and no network transport, so only keel's own 

451 push, pull request, provenance comment and attribution labels — made here, in keel's 

452 process, after the implementer has exited — reach the remote. 

453 

454 The worktree shares the operator's repository, so the implementer could also write its 

455 config and hooks, which keel's own credentialed steps would then run. The worker reads 

456 the push URL and the git setup before the implementer runs, stops at ``tamper`` when 

457 that setup changed — after the implementer, and again after the gates — and runs its 

458 own git steps with no hooks, fsmonitor or signing program 

459 (:func:`keel.swarm_worker.keel_git`), pushing to the URL it read rather than a remote's 

460 name. 

461 """ 

462 env = swarm_worker.worker_env(os.environ) 

463 run = runner or functools.partial(default_runner, timeout_s=timeout_s, env=env) 

464 scopes = swarm_worker.worker_scopes(live.consent, cluster.cluster_id) 

465 branch = cluster_branch(swarm_id, cluster.cluster_id) 

466 plan = live.dispatches[cluster.cluster_id] 

467 system = plan.attribution.get("system") or plan.provider 

468 record: dict[str, Any] = { 

469 "issue": cluster.issues[0] if cluster.issues else 0, 

470 "issues": list(cluster.issues), 

471 "role": cluster.role, 

472 "branch": branch, 

473 "scopes": list(scopes), 

474 "implementer": system, 

475 "commit": None, 

476 "pr_url": None, 

477 "pushed": False, 

478 # Whether the pull request carries the ship-provenance comment that arms keel 

479 # merge's evidence gate; an open pull request without it is held at landing. 

480 "provenance_posted": False, 

481 # Whether it carries the seat's attribution labels, and whether the gates this 

482 # worker ran are in the run ledger as a ship run for its head: without either, 

483 # keel merge holds it (#1420). 

484 "labels_applied": False, 

485 "gates_recorded": False, 

486 } 

487 progress = {} if progress is None else progress 

488 

489 def stop(stage: str, reason: str, code: int = 1) -> dict[str, Any]: 

490 return {**record, "ok": False, "code": code, "stage": stage, "output": reason} 

491 

492 on_stage("consent") 

493 if why := swarm_worker.consent_refusal(scopes): 

494 return stop("consent", why) 

495 # A registration whose directory is gone — a crashed run's, after its directory was 

496 # deleted — makes `git worktree add -B` of the same branch fail, so it is dropped 

497 # first (#1278). A sibling's worktree still being added is locked, and prune skips it. 

498 on_stage("worktree") 

499 pruned = prune_worktrees(root, runner=run) 

500 created = create_swarm_worktree(root, worktree_dir, branch, base_branch=base_branch, runner=run) 

501 if not created or not worktree_dir.exists(): 

502 # Most often a previous run of this --swarm-id left its worktree or branch behind. 

503 prune_note = f" (git worktree prune failed first: {pruned})" if pruned else "" 

504 return stop( 

505 "worktree", 

506 f"failed to create isolated worktree at {worktree_dir}{prune_note}; if a " 

507 f"previous run of this swarm left it behind, `keel swarm-status <project.yaml> " 

508 f"--swarm-id {swarm_id} --clean` removes it", 

509 ) 

510 progress["worktree_created"] = True 

511 # Read before the implementer runs, from configuration it has not touched yet: where 

512 # keel will push, the commit the work must descend from, and the git setup keel's own 

513 # steps will run under. 

514 start = _last_line(run(["git", "rev-parse", "HEAD"], worktree_dir)) 

515 push_url = _last_line(run(["git", "remote", "get-url", "--push", live.remote], root)) 

516 if not push_url: 

517 return stop("push", f"the push URL of remote {live.remote!r} could not be read") 

518 before = _git_snapshot(run, worktree_dir) 

519 if before is None or not start: 

520 return stop("worktree", f"the git setup of {worktree_dir} could not be read") 

521 

522 on_stage("implement") 

523 brief = Path(plan.prompt_path) 

524 brief.parent.mkdir(parents=True, exist_ok=True) 

525 brief.write_text( 

526 swarm_worker.render_brief( 

527 cluster, live.issue_scopes, swarm_id=swarm_id, branch=branch, base_branch=base_branch 

528 ), 

529 encoding="utf-8", 

530 ) 

531 gh_config = build_gh_config_path(swarm_id, cluster.cluster_id, root) 

532 # From here the worktree holds what the seat did, and a failure keeps it (#1278). 

533 progress["implementer_ran"] = True 

534 result = live.implement( 

535 plan, swarm_worker.implementer_env(os.environ, gh_config_dir=str(gh_config)) 

536 ) 

537 if not result.get("ok"): 

538 return stop( 

539 "implement", 

540 f"the implementer {system} failed ({result.get('error_code')}): {result.get('error')}", 

541 ) 

542 

543 on_stage("tamper") 

544 if findings := _tampered(run, worktree_dir, before): 

545 return stop("tamper", swarm_worker.tamper_reason(findings, during="the implementer ran")) 

546 # Every git step from here is keel's own, and runs with no hooks (an empty directory 

547 # made only now, so nothing could be planted in it), no fsmonitor and no signing program. 

548 with tempfile.TemporaryDirectory(prefix="keel-no-hooks-") as no_hooks: 

549 return _commit_gate_push_open( 

550 cluster, 

551 record, 

552 stop, 

553 git=functools.partial(swarm_worker.keel_git, hooks_dir=no_hooks), 

554 run=run, 

555 before=before, 

556 start=start, 

557 push_url=push_url, 

558 swarm_id=swarm_id, 

559 worktree_dir=worktree_dir, 

560 project_yaml=project_yaml, 

561 base_branch=base_branch, 

562 branch=branch, 

563 plan=plan, 

564 live=live, 

565 on_stage=on_stage, 

566 ) 

567 

568 

569def _hook_digests(hooks_dir: Path) -> dict[str, str]: 

570 """Each file under ``hooks_dir`` -> the digest of what it holds (through a symlink).""" 

571 paths = sorted(hooks_dir.rglob("*")) if hooks_dir.is_dir() else [] 

572 return { 

573 path.relative_to(hooks_dir).as_posix(): hashlib.sha256(path.read_bytes()).hexdigest() 

574 for path in paths 

575 if path.is_file() 

576 } 

577 

578 

579def _git_snapshot(run: SubprocessRunner, worktree_dir: Path) -> swarm_worker.GitSnapshot | None: 

580 """The worktree's git setup as keel's own steps would meet it, or ``None`` when git 

581 could not say. Read with plain ``git`` — ``rev-parse`` and ``config`` run no hook.""" 

582 where = run( 

583 ["git", "rev-parse", "--absolute-git-dir", "--git-common-dir", "--git-path", "hooks"], 

584 worktree_dir, 

585 ) 

586 config = run(["git", "config", "--list", "--show-origin", "--show-scope"], worktree_dir) 

587 paths = (where.stdout or where.output).splitlines() if where.ok else [] 

588 if len(paths) != 3 or not config.ok: 

589 return None 

590 # `--git-common-dir` and `--git-path` may answer relative to the worktree. 

591 git_dir, common_dir, hooks_dir = (str(worktree_dir / p) for p in paths) 

592 return swarm_worker.GitSnapshot( 

593 git_dir=git_dir, 

594 common_dir=common_dir, 

595 hooks_dir=hooks_dir, 

596 config=tuple((config.stdout or config.output).splitlines()), 

597 hooks=_hook_digests(Path(hooks_dir)), 

598 ) 

599 

600 

601def _tampered( 

602 run: SubprocessRunner, worktree_dir: Path, before: swarm_worker.GitSnapshot 

603) -> tuple[str, ...]: 

604 """What changed in the worktree's git setup since ``before``; unreadable is a change.""" 

605 after = _git_snapshot(run, worktree_dir) 

606 if after is None: 

607 return ("the git setup could no longer be read",) 

608 return swarm_worker.tamper_findings(before, after) 

609 

610 

611def _commit_gate_push_open( 

612 cluster: SwarmCluster, 

613 record: dict[str, Any], 

614 stop: Callable[..., dict[str, Any]], 

615 *, 

616 git: Callable[[list[str]], list[str]], 

617 run: SubprocessRunner, 

618 before: swarm_worker.GitSnapshot, 

619 start: str, 

620 push_url: str, 

621 swarm_id: str, 

622 worktree_dir: Path, 

623 project_yaml: str, 

624 base_branch: str, 

625 branch: str, 

626 plan: RunPlan, 

627 live: LiveRun, 

628 on_stage: StageReporter, 

629) -> dict[str, Any]: 

630 """The worker's steps after the implementer: commit, gates, push, pull request, and the 

631 pull request's ship-provenance comment.""" 

632 system = plan.attribution.get("system") or plan.provider 

633 on_stage("commit") 

634 status = run(git(["status", "--porcelain"]), worktree_dir) 

635 if not status.ok: 

636 return stop("commit", f"git status failed in {worktree_dir}: {status.output.strip()}") 

637 if status.output.strip(): 

638 message = swarm_worker.commit_message(cluster, swarm_id=swarm_id, plan=plan) 

639 for args in (["add", "-A"], ["commit", "--no-verify", "-m", message]): 

640 step = run(git(args), worktree_dir) 

641 if not step.ok: 

642 return stop("commit", f"git {args[0]} failed: {step.output.strip()}") 

643 head = _last_line(run(git(["rev-parse", "HEAD"]), worktree_dir)) 

644 if not head or head == start: 

645 return stop( 

646 "implement", 

647 f"the implementer {system} changed nothing in the worktree, so there is no work " 

648 "to commit or push", 

649 ) 

650 record["commit"] = head 

651 # A seat may commit its own work, but not replace the branch's history: what keel 

652 # pushes must descend from the commit the worktree was cut at. 

653 if not run(git(["merge-base", "--is-ancestor", start, head]), worktree_dir).ok: 

654 return stop( 

655 "tamper", 

656 f"{head} does not descend from {start}, the commit the worktree was cut at; the " 

657 "branch's history was replaced, so nothing is pushed", 

658 ) 

659 

660 on_stage("gates") 

661 gates = run( 

662 [ 

663 sys.executable, 

664 "-m", 

665 "keel", 

666 "run-gates", 

667 project_yaml, 

668 "--root", 

669 str(worktree_dir), 

670 "--phases", 

671 swarm_worker.GATE_PHASES, 

672 "--defer-jury", 

673 # The per-gate outcome, which the gates-pass record is written from (#1420): 

674 # the exit code alone cannot say that a blocking gate did not run. 

675 "--json", 

676 ], 

677 worktree_dir, 

678 ) 

679 printed = gates.stdout or gates.output 

680 outcomes = swarm_worker.parse_gate_report(printed) 

681 if not gates.ok: 

682 detail = printed.strip() if outcomes is None else swarm_worker.gate_report_summary(outcomes) 

683 return stop( 

684 "gates", 

685 f"the project's gates failed at {head}; the branch is not pushed:\n{detail}", 

686 code=gates.code, 

687 ) 

688 

689 # The gates ran the implementer's code, which could have written the git setup too. 

690 if findings := _tampered(run, worktree_dir, before): 

691 return stop("tamper", swarm_worker.tamper_reason(findings, during="the gates ran")) 

692 # What the pushed head changes against the commit the worktree was cut at, for the 

693 # gates-pass record; `None` (unreadable), never `[]`, when git cannot say. 

694 diff = run(git(["diff", "--no-ext-diff", "--name-only", start, head, "--"]), worktree_dir) 

695 listed = (diff.stdout or diff.output).splitlines() if diff.ok else None 

696 changed = None if listed is None else [path for path in listed if path.strip()] 

697 

698 # To the URL read before the implementer ran, never to the remote's name. 

699 on_stage("push") 

700 pushed = live.push( 

701 git(["push", "--no-verify", push_url, f"{head}:refs/heads/{branch}"]), worktree_dir 

702 ) 

703 if not pushed.ok: 

704 return stop("push", f"git push {live.remote} {branch} failed: {pushed.output.strip()}") 

705 record["pushed"] = True 

706 

707 on_stage("pull_request") 

708 opened = live.open_pr( 

709 swarm_worker.pull_request_title(cluster, live.issue_scopes, swarm_id=swarm_id), 

710 swarm_worker.pull_request_body( 

711 cluster, 

712 swarm_id=swarm_id, 

713 branch=branch, 

714 base_branch=base_branch, 

715 commit=head, 

716 plan=plan, 

717 seat_source=live.seat_sources.get(cluster.cluster_id), 

718 delegation=live.consent, 

719 ), 

720 base_branch, 

721 branch, 

722 worktree_dir, 

723 ) 

724 if not opened.ok: 

725 return stop( 

726 "pull_request", 

727 f"{branch} is pushed, but gh pr create failed: {opened.output.strip()}", 

728 ) 

729 url = _last_line(opened) 

730 # `/keel:ship` applies the seat's attribution labels to every pull request it opens, 

731 # and keel merge holds one without its `agent:<vendor>` label (#1420). 

732 unlabelled = _apply_labels(live, url, swarm_worker.attribution_labels(plan), worktree_dir) 

733 # A keel-made pull request arms keel merge's evidence gate at creation, as a live 

734 # `keel ship` run's does: its branch matches no ship-branch pattern, and without the 

735 # stamp `swarm-land` is held on "evidence gate is not enforced" (found on the first 

736 # end-to-end live run). Posted by keel, with the operator's credentials, after the 

737 # implementer exited — like the push and the pull request. 

738 missing = _post_provenance( 

739 live, 

740 url, 

741 swarm_worker.ship_provenance_body(cluster, swarm_id=swarm_id, commit=head, plan=plan), 

742 worktree_dir, 

743 ) 

744 gates_record, gates_note = _gates_record( 

745 cluster, 

746 live=live, 

747 plan=plan, 

748 swarm_id=swarm_id, 

749 branch=branch, 

750 base_branch=base_branch, 

751 head=head, 

752 url=url, 

753 outcomes=outcomes, 

754 changed=changed, 

755 ) 

756 done = { 

757 **record, 

758 "ok": True, 

759 "code": 0, 

760 "stage": "done", 

761 "pr_url": url, 

762 "provenance_posted": not missing, 

763 "labels_applied": not unlabelled, 

764 "output": f"pull request opened: {url}", 

765 } 

766 if gates_record is not None: 

767 # Appended by the run, on the thread that collects results (`record_worker_gates`). 

768 done[GATES_RECORD_KEY] = gates_record 

769 # Not a failed worker: the pull request is open and stays open. It is held at landing 

770 # until what is missing is posted, so the run says what and how. 

771 notes = [note for note in (unlabelled, missing, gates_note) if note] 

772 if notes: 

773 done["output"] += "".join(f"\n{note}" for note in notes) 

774 done["warnings"] = notes 

775 return done 

776 

777 

778def _apply_labels(live: LiveRun, url: str, labels: tuple[str, ...], cwd: Path) -> str: 

779 """Apply the seat's attribution ``labels`` to the pull request at ``url`` (#1420); 

780 ``""`` once applied, else the warning.""" 

781 owner_repo = swarm_worker.pull_request_repo(url) 

782 number = pull_request_number(url) 

783 if not labels: 

784 why = "keel's attribution for the implementer seat names no agent label" 

785 elif owner_repo is None or number is None: 

786 why = "gh pr create printed no pull request URL keel can read" 

787 else: 

788 applied = live.label_pr(owner_repo, number, list(labels), cwd) 

789 if applied.ok: 

790 return "" 

791 why = f"the label post failed: {applied.output.strip() or 'no output'}" 

792 return swarm_worker.labels_warning(url or "the pull request", labels, why) 

793 

794 

795def _gates_record( 

796 cluster: SwarmCluster, 

797 *, 

798 live: LiveRun, 

799 plan: RunPlan, 

800 swarm_id: str, 

801 branch: str, 

802 base_branch: str, 

803 head: str, 

804 url: str, 

805 outcomes: tuple[Any, ...] | None, 

806 changed: list[str] | None, 

807) -> tuple[dict[str, Any] | None, str]: 

808 """The worker's gates-pass record, redacted as every ledger record is, and the 

809 warning the run gives about it: ``(record, "")``, ``(record, why it is no pass)`` or 

810 ``(None, why there is none)`` (#1420).""" 

811 number = pull_request_number(url) 

812 shown = url or "the pull request" 

813 if number is None: 

814 why = "gh pr create printed no pull request number keel can read" 

815 return None, swarm_worker.gates_record_warning(shown, why) 

816 if outcomes is None: 

817 why = "keel run-gates printed no report keel can read, so what each gate did is unknown" 

818 return None, swarm_worker.gates_record_warning(shown, why) 

819 built = swarm_worker.gates_record( 

820 cluster, 

821 swarm_id=swarm_id, 

822 plan=plan, 

823 branch=branch, 

824 base_branch=base_branch, 

825 head=head, 

826 pull_request=number, 

827 outcomes=outcomes, 

828 changed_files=changed, 

829 config=live.config, 

830 ) 

831 try: 

832 gates_record = ledger.sanitize_record(built, live.config) 

833 except redaction.RedactionError as exc: 

834 why = f"the capture redaction policy is invalid: {exc}" 

835 return None, swarm_worker.gates_record_warning(shown, why) 

836 failed = swarm_worker.gates_not_a_pass(gates_record) 

837 return gates_record, swarm_worker.gates_not_a_pass_warning(shown, failed) if failed else "" 

838 

839 

840def record_worker_gates(live: LiveRun, outcome: Mapping[str, Any]) -> dict[str, Any]: 

841 """Append the gates-pass record a live worker built to the run ledger (#1420). 

842 

843 Returns the worker's result without the record, saying ``gates_recorded``, and with a 

844 warning when the append failed — never a failed worker: the pull request is open either 

845 way. Appended here, on the run's one collecting thread, as the consent records are, so 

846 no two workers write to the ledger at once. The ledger is the one ``keel merge --root 

847 <root>`` and ``swarm-land`` read, so :func:`keel.ledger.gates_pass_for_head` finds the 

848 record for the pull request's head. A run that records nothing (no ledger) appends 

849 nothing. 

850 """ 

851 result = {key: value for key, value in outcome.items() if key != GATES_RECORD_KEY} 

852 gates_record = outcome.get(GATES_RECORD_KEY) 

853 if live.ledger_path is None or gates_record is None: 

854 return result 

855 why = "" 

856 try: 

857 ledger.append_record(live.ledger_path, gates_record) 

858 except (OSError, ledger.LedgerError) as exc: 

859 why = f"{type(exc).__name__}: {exc}" 

860 if not why: 

861 return {**result, "gates_recorded": True} 

862 shown = str(outcome.get("pr_url") or "its pull request") 

863 warning = swarm_worker.gates_record_warning(shown, why) 

864 return { 

865 **result, 

866 "output": f"{result.get('output', '')}\n{warning}", 

867 "warnings": [*result.get("warnings", ()), warning], 

868 } 

869 

870 

871def _post_provenance(live: LiveRun, url: str, body: str, cwd: Path) -> str: 

872 """Post ``body`` on the pull request at ``url``; ``""`` once posted, else the warning.""" 

873 owner_repo = swarm_worker.pull_request_repo(url) 

874 number = pull_request_number(url) 

875 if owner_repo is None or number is None: 

876 why = "gh pr create printed no pull request URL keel can read" 

877 else: 

878 posted = live.post_comment(owner_repo, number, body, cwd) 

879 if posted.ok: 

880 return "" 

881 why = f"the comment post failed: {posted.output.strip() or 'no output'}" 

882 return swarm_worker.provenance_warning(url or "the pull request", why) 

883 

884 

885def record_pull_request_consent(live: LiveRun, cluster_id: str, outcome: Mapping[str, Any]) -> str: 

886 """Append the ``pull_request`` delegation event for a worker's open pull request (#1400). 

887 

888 Returns ``""`` once it is in the run ledger (or when the run records nothing), else a 

889 warning. Never a failed worker: the pull request is open either way. Without this 

890 event ``keel consent-verify`` reports no delegated consent for that pull request — it 

891 matches a delegation only by the pull request's number and pushed head, never by its 

892 branch name (a fork can name a branch anything). 

893 """ 

894 if live.ledger_path is None: 

895 return "" 

896 url = str(outcome.get("pr_url") or "") 

897 number = pull_request_number(url) 

898 why = "gh pr create printed no pull request number keel can read" 

899 if number is not None: 

900 why = "" 

901 try: 

902 ledger.append_record( 

903 live.ledger_path, 

904 ledger.build_consent_delegation_record( 

905 live.consent.to_dict(), 

906 recorded_at=_now(), 

907 pull_request={ 

908 "cluster": cluster_id, 

909 "number": number, 

910 "branch": outcome.get("branch"), 

911 "head_sha": outcome.get("commit"), 

912 "url": url, 

913 }, 

914 ), 

915 ) 

916 except (OSError, ledger.LedgerError) as exc: 

917 why = f"{type(exc).__name__}: {exc}" 

918 if not why: 

919 return "" 

920 return ( 

921 f"the consent delegation for {url or 'its pull request'} is not in the run ledger " 

922 f"({why}); keel consent-verify will report no delegated consent for it — record the " 

923 "consent by hand or re-run the cluster" 

924 ) 

925 

926 

927def _swarm_ids_under(base: Path, path: str) -> tuple[str, str] | None: 

928 """``(swarm_id, cluster_id)`` of a worktree path exactly two levels under ``base``.""" 

929 try: 

930 parts = Path(path).resolve().relative_to(base).parts 

931 except ValueError: 

932 return None 

933 return (parts[0], parts[1]) if len(parts) == 2 else None 

934 

935 

936def _run_record(swarm_id: str, root: Path) -> swarm_worker.SwarmRunRecord | None: 

937 """What ``swarm_id``'s state file says about its leftovers; ``None`` with no file.""" 

938 state = load_swarm_state(swarm_id, root=root) 

939 if state is None: 

940 exists = (root / ".keel" / "state" / "swarm" / f"{swarm_id}.json").exists() 

941 # A state file keel cannot read may belong to a running run. 

942 return swarm_worker.SwarmRunRecord(unfinished=True) if exists else None 

943 return swarm_worker.SwarmRunRecord( 

944 unfinished=state.completed_at is None, 

945 pull_requests={ 

946 w.cluster_id: w.pull_request for w in state.workers if w.pull_request is not None 

947 }, 

948 pushed=frozenset(w.cluster_id for w in state.workers if w.pushed), 

949 ) 

950 

951 

952def find_swarm_leftovers( 

953 root: str | Path, 

954 *, 

955 swarm_id: str | None = None, 

956 runner: SubprocessRunner | None = None, 

957) -> tuple[tuple[swarm_worker.SwarmLeftover, ...], str]: 

958 """Every worktree, directory and branch a swarm run left behind, and what to do (#1278). 

959 

960 Looks only under keel's own paths — worktrees registered, or directories present, 

961 exactly at ``.keel/worktrees/<swarm_id>/<cluster_id>`` — and its own branch namespace, 

962 ``swarm/<swarm_id>/<cluster_id>``; ``swarm_id`` narrows it to one run. 

963 :func:`keel.swarm_worker.classify_leftovers` decides what may be removed. Returns 

964 ``(leftovers, "")``, or ``((), why)`` when git could not list them. 

965 """ 

966 root_path = Path(root).resolve() 

967 run = runner or default_runner 

968 base = root_path / ".keel" / "worktrees" 

969 listed = run(["git", "worktree", "list", "--porcelain"], root_path) 

970 refs = run(["git", "for-each-ref", "--format=%(refname)", "refs/heads/swarm/"], root_path) 

971 if not listed.ok or not refs.ok: 

972 failed = listed if not listed.ok else refs 

973 return (), f"git could not list the worktrees or branches: {failed.output.strip()}" 

974 

975 worktrees: list[tuple[str, str, str, bool]] = [] 

976 entries = swarm_worker.parse_worktree_list(listed.stdout or listed.output) 

977 for entry in entries: 

978 ids = _swarm_ids_under(base, entry.path) 

979 if ids is not None: 

980 # In the platform's own spelling: git prints `C:/…` on Windows. 

981 worktrees.append((*ids, str(Path(entry.path).resolve()), entry.prunable)) 

982 # Every worktree registered under `.keel/worktrees/`, of keel's shape or not: a 

983 # directory that is one, holds one or sits inside one is never taken for a leftover 

984 # directory. (The main worktree holds all of it, so only those below the base count.) 

985 resolved = (Path(entry.path).resolve() for entry in entries) 

986 registered = [p for p in resolved if base in p.parents] 

987 

988 def unregistered(path: Path) -> bool: 

989 return not any(p == path or p in path.parents or path in p.parents for p in registered) 

990 

991 directories: list[tuple[str, str, str]] = [] 

992 runs_dirs = sorted(base.iterdir()) if base.is_dir() and not base.is_symlink() else [] 

993 for run_dir in runs_dirs: 

994 if run_dir.is_symlink() or not run_dir.is_dir(): 

995 continue 

996 children = sorted(run_dir.iterdir()) 

997 if not children: 

998 directories.append((run_dir.name, "", str(run_dir))) 

999 for child in children: 

1000 if child.is_dir() and not child.is_symlink() and unregistered(child): 

1001 directories.append((run_dir.name, child.name, str(child))) 

1002 

1003 branches: dict[str, str] = {} # branch -> its run 

1004 for line in (refs.stdout or refs.output).splitlines(): 

1005 ids = swarm_worker.swarm_branch_ids(line.strip()) 

1006 if ids is not None: 

1007 branches[line.strip()] = ids[0] 

1008 if swarm_id is not None: 

1009 worktrees = [w for w in worktrees if w[0] == swarm_id] 

1010 directories = [d for d in directories if d[0] == swarm_id] 

1011 branches = {b: sid for b, sid in branches.items() if sid == swarm_id} 

1012 seen = {w[0] for w in worktrees} | {d[0] for d in directories} | set(branches.values()) 

1013 runs = {sid: record for sid in sorted(seen) if (record := _run_record(sid, root_path))} 

1014 leftovers = swarm_worker.classify_leftovers( 

1015 worktrees=worktrees, 

1016 directories=directories, 

1017 branches=branches, 

1018 runs=runs, 

1019 named=swarm_id, 

1020 ) 

1021 return leftovers, "" 

1022 

1023 

1024def clean_swarm_leftovers( 

1025 root: str | Path, 

1026 leftovers: tuple[swarm_worker.SwarmLeftover, ...], 

1027 *, 

1028 runner: SubprocessRunner | None = None, 

1029) -> tuple[list[swarm_worker.SwarmLeftover], list[tuple[swarm_worker.SwarmLeftover, str]]]: 

1030 """Remove every leftover marked ``remove``; ``(removed, [(failed, why)])`` (#1278). 

1031 

1032 Registrations go with ``git worktree prune``, worktrees with 

1033 :func:`remove_swarm_worktree`, unregistered directories with ``rmtree`` — before the 

1034 branches, which git will not delete while a worktree has them checked out — and each 

1035 run directory left empty with ``rmdir``. Nothing marked ``keep`` is touched. 

1036 """ 

1037 root_path = Path(root).resolve() 

1038 run = runner or default_runner 

1039 todo = [x for x in leftovers if x.action == "remove"] 

1040 removed: list[swarm_worker.SwarmLeftover] = [] 

1041 failed: list[tuple[swarm_worker.SwarmLeftover, str]] = [] 

1042 

1043 def settle(item: swarm_worker.SwarmLeftover, gone: bool, why: str) -> None: 

1044 if gone: 

1045 removed.append(item) 

1046 else: 

1047 failed.append((item, why)) 

1048 

1049 registrations = [x for x in todo if x.kind == "registration"] 

1050 if registrations: 

1051 why = prune_worktrees(root_path, runner=run) 

1052 for item in registrations: 

1053 settle(item, not why, f"git worktree prune failed: {why}") 

1054 for item in todo: 

1055 if item.kind == "worktree": 

1056 gone = remove_swarm_worktree(root_path, Path(item.target), runner=run) 

1057 settle(item, gone, "the directory is still there") 

1058 elif item.kind == "directory" and item.cluster_id: 

1059 shutil.rmtree(item.target, ignore_errors=True) 

1060 settle(item, not Path(item.target).exists(), "the directory is still there") 

1061 for item in todo: 

1062 if item.kind == "branch": 

1063 deleted = run(["git", "branch", "-D", item.target], root_path) 

1064 settle(item, deleted.ok, deleted.output.strip() or f"exit {deleted.code}") 

1065 for swarm_id in sorted({x.swarm_id for x in todo}): 

1066 remove_empty_swarm_dirs(root_path, swarm_id) 

1067 for item in todo: 

1068 if item.kind == "directory" and not item.cluster_id: 

1069 settle(item, not Path(item.target).exists(), "the directory is not empty") 

1070 return removed, failed 

1071 

1072 

1073def _disposal(ok: bool, progress: Mapping[str, bool]) -> swarm_worker.WorktreeDisposal: 

1074 return swarm_worker.worktree_disposal( 

1075 ok=ok, 

1076 worktree_created=progress.get("worktree_created", False), 

1077 implementer_ran=progress.get("implementer_ran", False), 

1078 ) 

1079 

1080 

1081def _with_settlement( 

1082 outcome: dict[str, Any], settled: dict[str, Any], branch: str 

1083) -> dict[str, Any]: 

1084 """A live worker's result with what became of its worktree; a kept worktree is named 

1085 at the end of ``output``, where the tail the swarm keeps still holds it (#1278).""" 

1086 output = outcome.get("output", "") 

1087 if settled["worktree_state"] == "kept": 

1088 output = ( 

1089 f"{output}\nthe worktree is kept for inspection at {settled['worktree']} (branch " 

1090 f"{branch}); `keel swarm-status <project.yaml> --clean` removes it" 

1091 ) 

1092 # The worker's own warnings (an unstamped pull request) come first, then the settle's: 

1093 # a plain merge would let the settle's list replace the worker's. 

1094 warnings = [*outcome.get("warnings", ()), *settled["warnings"]] 

1095 return {**outcome, **settled, "output": output, "warnings": warnings} 

1096 

1097 

1098def run_swarm_orchestration( 

1099 plan: SwarmPlan, 

1100 project_yaml: str, 

1101 *, 

1102 root: str | Path = ".", 

1103 dry_run: bool = True, 

1104 max_workers: int = 4, 

1105 runner: SubprocessRunner | None = None, 

1106 base_branch: str, 

1107 timeout_s: int = DEFAULT_RUNNER_TIMEOUT_S, 

1108 live: LiveRun | None = None, 

1109) -> SwarmRunResult: 

1110 """Execute the waves and clusters of a SwarmPlan with fail-soft isolation. 

1111 

1112 ``base_branch`` is required, as it is for ``land_wave_clusters``: the worktrees 

1113 are branched from it and the wave is later landed onto it, and a default of 

1114 ``main`` branched a ``develop`` project's clusters off the wrong history (#1262). 

1115 

1116 ``timeout_s`` is each worker's gate budget — the child ``keel ship`` of a dry run, the 

1117 ``keel run-gates`` of a live one; ``swarm-run`` passes 

1118 :func:`keel.swarm.worker_timeout_s` or its ``--worker-timeout`` (#1279). 

1119 

1120 A live run (``dry_run=False``) needs ``live``: the operator's consent as the parent 

1121 delegated it, and each cluster's planned implementer seat (#1400). Without it nothing 

1122 starts — there is no live worker that runs without delegated consent. A live worker 

1123 always gets its own worktree, and a dry run's never does (#1288); the 

1124 ``create_worktrees`` switch that used to sit beside ``dry_run`` could only ever agree 

1125 with it, and is gone (#1280). 

1126 

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

1128 ``swarm-run`` never passes it. It replaces the subprocess runner of every git command, 

1129 gate run and dry-run child ``keel ship`` a worker makes; left out, each runs through 

1130 :func:`default_runner` under ``timeout_s``. 

1131 

1132 Each worker's record carries its wave, when it started and ended, and — a live worker 

1133 — the stage it is in, written as it enters it, so ``keel swarm-status`` shows how far a 

1134 run has got while it is in flight (#1280). A worker is ``running`` from the moment it 

1135 starts, not from the moment its wave does: a dry run's workers take turns. 

1136 """ 

1137 if not dry_run and live is None: 

1138 raise ValueError("a live swarm run needs the operator's delegated consent (#1400)") 

1139 root_path = Path(root).resolve() 

1140 workers_list: list[SwarmWorkerStatus] = [] 

1141 

1142 # Initialize workers from the plan's own staffing, so the board reports the team the 

1143 # planner resolved rather than the record's placeholder defaults (#1017). A live 

1144 # worker's record also says which consent scopes it was handed (#1400). 

1145 for w in plan.waves: 

1146 for c in w.clusters: 

1147 seed = worker_seed(c, wave=w.wave_index, updated_at=_now()) 

1148 if live is not None: 

1149 seed = replace(seed, scopes=swarm_worker.worker_scopes(live.consent, c.cluster_id)) 

1150 workers_list.append(seed) 

1151 

1152 state = SwarmRunState( 

1153 swarm_id=plan.swarm_id, 

1154 total_workers=len(workers_list), 

1155 active_wave=1, 

1156 workers=tuple(workers_list), 

1157 started_at=_now(), 

1158 consent=None if live is None else live.consent.to_dict(), 

1159 ) 

1160 save_swarm_state(state, root=root_path) 

1161 

1162 # Live workers report from their own threads while the loop below records the ones that 

1163 # finished, so every update of the run's state — and its save — is made under one lock. 

1164 state_lock = threading.Lock() 

1165 

1166 def report(cluster_id: str, **fields: Any) -> None: 

1167 nonlocal state 

1168 with state_lock: 

1169 state = update_worker_state(state, cluster_id, **fields) 

1170 save_swarm_state(state, root=root_path) 

1171 

1172 def _enter_stage(cluster_id: str, stage: str) -> None: 

1173 report(cluster_id, stage=stage) 

1174 

1175 passed_count = 0 

1176 failed_count = 0 

1177 wave_results: list[dict[str, Any]] = [] 

1178 warnings: list[str] = [] 

1179 current_plan = plan 

1180 

1181 # Waves are followed by their index, not by position: a failure makes 

1182 # rebalance_swarm_plan drop a wave, and a position counter over the shrunken 

1183 # plan then stepped past the next, unrelated wave without running it (#1268). 

1184 last_wave: int | None = None 

1185 running = "implementing the cluster" if live is not None else "executing ship pipeline" 

1186 while True: 

1187 remaining = [w for w in current_plan.waves if last_wave is None or w.wave_index > last_wave] 

1188 if not remaining: 

1189 break 

1190 wave = remaining[0] 

1191 last_wave = wave.wave_index 

1192 # No worker thread runs between waves, so this update needs no lock. 

1193 state = replace(state, active_wave=wave.wave_index) 

1194 

1195 cluster_tasks = list(wave.clusters) 

1196 if not cluster_tasks: 

1197 continue 

1198 save_swarm_state(state, root=root_path) 

1199 

1200 wave_record: dict[str, Any] = { 

1201 "wave_index": wave.wave_index, 

1202 "mode": wave.mode, 

1203 "eligible_direct_landing": wave.eligible_direct_landing, 

1204 "cluster_results": {}, 

1205 } 

1206 

1207 def _worker_fn(cluster: Any) -> tuple[str, dict[str, Any]]: 

1208 c_id = cluster.cluster_id 

1209 wt_path = build_worktree_path(plan.swarm_id, c_id, root=root_path) 

1210 # Running from when this worker starts, not when its wave did: a dry run's 

1211 # workers run one at a time, and the rest are still queued (#1280). 

1212 report(c_id, step="s4", status="running", details=running, started_at=_now()) 

1213 

1214 # A live worker implements its cluster in a worktree of its own, cut here and 

1215 # settled here, whichever way the worker ends — a returned failure or a raise 

1216 # (#1400, #1278): see `swarm_worker.worktree_disposal` for what stays. 

1217 if not dry_run: 

1218 branch = cluster_branch(plan.swarm_id, c_id) 

1219 progress: dict[str, bool] = {} 

1220 try: 

1221 outcome = execute_live_cluster_worker( 

1222 cluster, 

1223 swarm_id=plan.swarm_id, 

1224 root=root_path, 

1225 worktree_dir=wt_path, 

1226 project_yaml=project_yaml, 

1227 base_branch=base_branch, 

1228 live=live, 

1229 runner=runner, 

1230 timeout_s=timeout_s, 

1231 progress=progress, 

1232 on_stage=functools.partial(_enter_stage, c_id), 

1233 ) 

1234 except BaseException: 

1235 # What the seat touched is kept for inspection; nothing is lost to 

1236 # a crash in keel's own code. 

1237 settle_live_worktree( 

1238 root_path, wt_path, branch, _disposal(False, progress), runner 

1239 ) 

1240 raise 

1241 else: 

1242 settled = settle_live_worktree( 

1243 root_path, 

1244 wt_path, 

1245 branch, 

1246 _disposal(bool(outcome.get("ok")), progress), 

1247 runner, 

1248 ) 

1249 return c_id, _with_settlement(outcome, settled, branch) 

1250 

1251 # A dry run's worker is an assessment in the operator's checkout. 

1252 return c_id, execute_cluster_worker( 

1253 project_yaml=project_yaml, 

1254 issue=cluster.issues[0] if cluster.issues else 0, 

1255 root=root_path, 

1256 worktree_dir=wt_path, 

1257 dry_run=True, 

1258 role=cluster.role, 

1259 # The lead hands its cluster's team to the child ship. Without this the 

1260 # child re-resolved from config alone and dropped both the difficulty 

1261 # bench and the operator's per-run overrides. 

1262 extra_args=list(ship_handoff_args(cluster.assignment)), 

1263 runner=runner, 

1264 timeout_s=timeout_s, 

1265 ) 

1266 

1267 # Run wave clusters in parallel thread pool 

1268 # Without a worktree each child runs in the operator's own checkout (a dry run 

1269 # creates none), and its gate suite is not written to share a tree with a 

1270 # sibling's: `.coverage`, `.pytest_cache`, build output. Such children run one 

1271 # at a time; only isolated workers run in parallel (#1288). 

1272 pool_workers = min(max_workers, len(cluster_tasks)) if not dry_run else 1 

1273 with concurrent.futures.ThreadPoolExecutor(max_workers=pool_workers) as executor: 

1274 future_to_cluster = { 

1275 executor.submit(_worker_fn, cluster): cluster for cluster in cluster_tasks 

1276 } 

1277 for future in concurrent.futures.as_completed(future_to_cluster): 

1278 try: 

1279 c_id, worker_res = future.result() 

1280 except Exception as exc: # noqa: BLE001 - one worker must not end the run 

1281 # A worker that raised (a malformed assignment, an OSError making 

1282 # its path) is a failed cluster, not a failed run. Left unguarded, 

1283 # the raise discarded the other workers' results and froze them 

1284 # as `running` in the state file (#1271). 

1285 cluster = future_to_cluster[future] 

1286 c_id = cluster.cluster_id 

1287 # The traceback goes to stderr: the failure may be keel's own bug, 

1288 # and the one-line reason alone would hide where it came from. 

1289 sys.stderr.write( 

1290 f"swarm worker {c_id} raised:\n" + "".join(traceback.format_exception(exc)) 

1291 ) 

1292 worker_res = { 

1293 "issue": cluster.issues[0] if cluster.issues else 0, 

1294 "role": cluster.role, 

1295 "ok": False, 

1296 "code": 1, 

1297 "output": f"worker raised {type(exc).__name__}: {exc}", 

1298 } 

1299 # A raising live worker keeps what its seat touched (#1278). 

1300 kept = build_worktree_path(plan.swarm_id, c_id, root=root_path) 

1301 if not dry_run and kept.exists(): 

1302 worker_res["worktree"] = str(kept) 

1303 worker_res["output"] += ( 

1304 f"\nthe worktree is kept for inspection at {kept}; " 

1305 "`keel swarm-status <project.yaml> --clean` removes it" 

1306 ) 

1307 # The gates-pass a live worker built goes into the run ledger here, on this 

1308 # one thread, before its result is kept anywhere (#1420). 

1309 if live is not None and worker_res.get("ok", False): 

1310 worker_res = record_worker_gates(live, worker_res) 

1311 # Bounded once, here, where every origin of `output` — the child's 

1312 # stdout, a worktree failure, a raised worker — meets both places it is 

1313 # kept: this wave record and, on failure, the state file's `details` 

1314 # (#1280). 

1315 worker_res = {**worker_res, "output": tail_child_output(worker_res["output"])} 

1316 wave_record["cluster_results"][c_id] = worker_res 

1317 issue_val = worker_res.get("issue", 0) 

1318 warnings.extend(f"{c_id}: {w}" for w in worker_res.get("warnings", ())) 

1319 # Where the worker's worktree still is, and whether its branch was pushed: 

1320 # what `keel swarm-status --clean` reads (#1278). 

1321 left = { 

1322 "worktree": worker_res.get("worktree") or "", 

1323 "pushed": bool(worker_res.get("pushed")), 

1324 } 

1325 

1326 # The stage a live worker ended at: `done`, or the one that stopped it. A 

1327 # dry run's worker, or one that raised, reports none, and keeps the last 

1328 # stage it entered (#1280). 

1329 ended = {"stage": worker_res.get("stage") or None, "finished_at": _now()} 

1330 if worker_res.get("ok", False): 

1331 passed_count += 1 

1332 # Recorded here, on the one thread that collects results, so no two 

1333 # workers append to the ledger at once (#1400). 

1334 if live is not None and ( 

1335 missed := record_pull_request_consent(live, c_id, worker_res) 

1336 ): 

1337 warnings.append(f"{c_id}: {missed}") 

1338 # A live worker stops at an open pull request, which CI (s6) and review 

1339 # take from there; a dry assessment ran the whole backbone. 

1340 report( 

1341 c_id, 

1342 step="s10" if live is None else "s6", 

1343 status="passed", 

1344 details="pipeline completed" if live is None else worker_res["output"], 

1345 # The pull request `swarm-land` merges (#1287); a dry run opens none. 

1346 pull_request=pull_request_number(str(worker_res.get("pr_url") or "")), 

1347 **left, 

1348 **ended, 

1349 ) 

1350 else: 

1351 failed_count += 1 

1352 report( 

1353 c_id, 

1354 step="s4", 

1355 status="failed", 

1356 details=worker_res.get("output", ""), 

1357 **left, 

1358 **ended, 

1359 ) 

1360 # Dynamically rebalance subsequent waves if needed 

1361 current_plan = rebalance_swarm_plan(current_plan, issue_val) 

1362 

1363 # Once the wave's workers are done, the run's directory goes when nothing is left 

1364 # in it — it used to accumulate, one per run (#1278). 

1365 if not dry_run: 

1366 remove_empty_swarm_dirs(root_path, plan.swarm_id) 

1367 wave_results.append(wave_record) 

1368 

1369 # Finalize state 

1370 overall_status = ( 

1371 "success" 

1372 if failed_count == 0 and passed_count > 0 

1373 else ("partial_failure" if passed_count > 0 else "failed") 

1374 ) 

1375 if passed_count == 0 and failed_count == 0: 

1376 overall_status = "success" 

1377 

1378 state = replace(state, completed_at=_now()) 

1379 save_swarm_state(state, root=root_path) 

1380 

1381 return SwarmRunResult( 

1382 swarm_id=plan.swarm_id, 

1383 status=overall_status, 

1384 total_workers=len(workers_list), 

1385 passed_count=passed_count, 

1386 failed_count=failed_count, 

1387 dry_run=dry_run, 

1388 wave_results=tuple(wave_results), 

1389 consent=state.consent, 

1390 warnings=tuple(warnings), 

1391 )