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
« 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.
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.
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"""
15from __future__ import annotations
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
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)
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]
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
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.
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 )
131#: Called with each :data:`keel.swarm_worker.STAGES` entry as a live worker enters it.
132StageReporter = Callable[[str], None]
135def _no_stage(_stage: str) -> None:
136 """The stage reporter of a worker nobody watches."""
139def _now() -> str:
140 """The current UTC time, ISO 8601 — the timestamps a worker record carries."""
141 return datetime.datetime.now(datetime.UTC).isoformat()
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
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"
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"
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}"
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))
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))
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))
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))
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).
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))
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"
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."""
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
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
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.
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()
279def prune_worktrees(repo_root: Path, runner: SubprocessRunner | None = None) -> str:
280 """Run ``git worktree prune``; ``""``, or why it failed (#1278).
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}")
290def remove_empty_swarm_dirs(root: Path, swarm_id: str) -> None:
291 """Remove ``.keel/worktrees/<swarm_id>/``, then ``.keel/worktrees/``, when each is empty.
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
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.
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 }
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.
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)
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 )
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 }
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 ""
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.
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).
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``.
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.
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.
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
489 def stop(stage: str, reason: str, code: int = 1) -> dict[str, Any]:
490 return {**record, "ok": False, "code": code, "stage": stage, "output": reason}
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")
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 )
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 )
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 }
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 )
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)
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 )
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 )
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()]
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
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
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)
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 ""
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).
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 }
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)
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).
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 )
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
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 )
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).
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()}"
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]
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)
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)))
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, ""
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).
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]] = []
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))
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
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 )
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}
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.
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).
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).
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).
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``.
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] = []
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)
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)
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()
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)
1172 def _enter_stage(cluster_id: str, stage: str) -> None:
1173 report(cluster_id, stage=stage)
1175 passed_count = 0
1176 failed_count = 0
1177 wave_results: list[dict[str, Any]] = []
1178 warnings: list[str] = []
1179 current_plan = plan
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)
1195 cluster_tasks = list(wave.clusters)
1196 if not cluster_tasks:
1197 continue
1198 save_swarm_state(state, root=root_path)
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 }
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())
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)
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 )
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 }
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)
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)
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"
1378 state = replace(state, completed_at=_now())
1379 save_swarm_state(state, root=root_path)
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 )