Coverage for src/keel/swarm_worker.py: 100%
352 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"""What a live swarm worker is, what it may do, and what it leaves behind (#1400).
3A live ``swarm-run`` worker used to be a child ``keel ship`` — an assessment that never
4writes code, commits or opens a pull request, in any mode — so a live swarm had nothing to
5land. A live worker is now the cluster's **implementer seat**, dispatched by keel through
6the same machinery ``keel delegate run`` uses for one issue (:mod:`keel.delegate`), in the
7cluster's own worktree. keel then commits what the implementer wrote, runs the project's
8implementation gates, pushes the cluster branch and opens one pull request for the cluster.
10This module holds every *decision* in that flow, and nothing else:
12- **Consent.** The parent obtains the operator's consent exactly the way every other live
13 keel command does (:func:`keel.consent.build_consent_contract`), and
14 :func:`delegate_consent` turns an approved contract into a :class:`ConsentDelegation`: who
15 consented, which scopes, which run and clusters, when. Each worker is handed
16 :func:`worker_scopes` — the parent's scopes, for a cluster the delegation names, and
17 nothing for any other — and asks :func:`consent_refusal` before its first mutation, so a
18 worker never performs a mutation outside the scopes the parent handed it.
19- **The seat.** :func:`plan_implementer` resolves the cluster's implementer seat
20 (``knobs.team.implement``, a bench, or ``--delegate``, already resolved by
21 :func:`keel.team.resolve_assignment`) into a :class:`keel.delegate.RunPlan`, and refuses a
22 seat keel cannot run as a worker: a host subagent, and a transport that cannot edit a
23 worktree.
24- **What is written.** The brief, the commit message, the pull request's title and body,
25 the ship-provenance comment that arms ``keel merge``'s evidence gate on it, the attribution
26 labels it carries, and the ``ship_run`` record of the gates the worker ran for its head
27 (#1420) — the two things besides review evidence that ``keel merge`` asks of it.
28- **What is left behind** (#1278). :func:`worktree_disposal` decides what becomes of a
29 worker's worktree and branch when it ends, and :func:`classify_leftovers` which of a run's
30 worktrees, directories and branches ``keel swarm-status --clean`` may remove.
32Pure and deterministic: no subprocess, no filesystem, no clock. The runtime
33(:mod:`keel.swarm_runtime`) performs the mutations, behind :func:`consent_refusal`.
34"""
36from __future__ import annotations
38import json
39import re
40from collections import Counter
41from collections.abc import Iterable, Mapping
42from dataclasses import dataclass, field, replace
43from typing import Any
45from . import artifacts, consent, delegate, ledger, ship
46from . import findings as fnd
47from . import providers as providers_mod
48from .gates import GateOutcome, collect_findings
50#: Every mutation a live worker performs, in the order it performs them. The parent's
51#: consent contract is built over exactly these, so the scopes the operator approves are
52#: the scopes a worker needs — ``filesystem``, ``git`` and ``github`` — and no more.
53#: ``labels`` is the attribution pair the worker applies to its pull request (#1420).
54WORKER_SIDE_EFFECTS = (
55 "git_worktree",
56 "file_edit",
57 "git_commit",
58 "git_push",
59 "pull_request",
60 "labels",
61)
63#: Where a live worker can stop, in order. A worker reports the stage it failed at; the
64#: branch is pushed only past ``gates`` and the pull request opened only past ``push``.
65#: ``tamper`` is asked twice — after the implementer, and again after the gates, which run
66#: the implementer's code — and stops the worker before its first git step that follows.
67STAGES = (
68 "consent",
69 "worktree",
70 "implement",
71 "tamper",
72 "commit",
73 "gates",
74 "push",
75 "pull_request",
76 "done",
77)
79#: The consent a parent may have read from its own environment. It is the parent's to
80#: delegate, explicitly, and never reaches a worker's children through inheritance: the
81#: implementer is an agent CLI that can run ``keel`` itself, and an inherited
82#: ``KEEL_APPROVE_SCOPE`` would let it approve its own mutations.
83CONSENT_ENV_VARS = ("KEEL_APPROVE_SCOPE", "KEEL_OPERATOR", "KEEL_CONSENT_MODE")
85#: The forge credentials ``gh`` (and most GitHub API clients) read from the environment.
86#: The implementer is an agent CLI with tools, steered by issue text nobody vetted, and
87#: keel pushes the branch and opens the pull request itself once the implementer has
88#: exited — so the implementer is never handed these. Nor are the gates, which run the
89#: code the implementer wrote, or keel's own local git steps, which need no forge: only
90#: keel's push and ``gh pr create`` hold them, running in keel's process with the
91#: operator's environment. A model provider's key (``ANTHROPIC_API_KEY``, ``OPENAI_API_KEY``,
92#: ``GEMINI_API_KEY`` …) is not a forge credential and passes through: the seat needs it
93#: to reach its model.
94FORGE_TOKEN_ENV_VARS = (
95 "GH_TOKEN",
96 "GITHUB_TOKEN",
97 "GH_ENTERPRISE_TOKEN",
98 "GITHUB_ENTERPRISE_TOKEN",
99)
101#: The git configuration the implementer runs under, handed over through git's
102#: environment channel (``GIT_CONFIG_COUNT`` / ``GIT_CONFIG_KEY_<n>`` /
103#: ``GIT_CONFIG_VALUE_<n>``, git 2.31+), which outranks every config file. Each entry is
104#: there for one reason:
105#:
106#: - ``credential.helper`` set empty clears the helper list every config file built
107#: (``osxkeychain``, ``manager``, ``store`` …), including a helper scoped to a URL —
108#: the form ``gh auth setup-git`` writes — so git has no stored password to offer.
109#: - ``protocol.allow=never`` and ``never`` for each network transport by name make git
110#: refuse to reach any remote at all — fetch included, HTTPS and SSH alike. This is the
111#: entry that holds when a credential still exists somewhere (an SSH key, a helper the
112#: list above does not reach): the transport never opens. The transports are named as
113#: well as defaulted so a user's own ``protocol.https.allow=always`` cannot outrank it.
114#: - ``protocol.file.allow=user`` restores git's own default for local paths, which
115#: ``protocol.allow=never`` would otherwise also cover: a project whose tests clone or
116#: push to a repository on disk keeps working, and a path on disk is not the forge.
117IMPLEMENTER_GIT_CONFIG = (
118 ("credential.helper", ""),
119 ("protocol.allow", "never"),
120 ("protocol.http.allow", "never"),
121 ("protocol.https.allow", "never"),
122 ("protocol.ssh.allow", "never"),
123 ("protocol.git.allow", "never"),
124 ("protocol.file.allow", "user"),
125)
127#: Git configuration the parent's environment may already carry. The implementer does
128#: not inherit it: ``GIT_CONFIG_PARAMETERS`` is read after ``GIT_CONFIG_COUNT`` and would
129#: outrank :data:`IMPLEMENTER_GIT_CONFIG`, and an inherited count would renumber it.
130_INHERITED_GIT_CONFIG_VARS = ("GIT_CONFIG_COUNT", "GIT_CONFIG_PARAMETERS")
131_INHERITED_GIT_CONFIG_PREFIXES = ("GIT_CONFIG_KEY_", "GIT_CONFIG_VALUE_")
133#: The gate phases a worker runs after its commit: the ones the s4 loop judges an
134#: implementation by. The jury is a review, deferred with the rest of review to the steps
135#: that land the cluster (#1287).
136GATE_PHASES = "guard,test"
138#: The only delegate transport that runs an agent with tools in a working directory.
139#: ``api``/``ollama`` return text, and a generic ``profile`` is not known to edit files
140#: (the rule ``/keel:ship`` s4 already applies), so a worker built on one would commit
141#: nothing.
142TOOL_TRANSPORT = "cli"
145@dataclass(frozen=True)
146class ConsentDelegation:
147 """The operator's live consent, as the parent hands it to one swarm run's workers."""
149 swarm_id: str
150 clusters: tuple[str, ...]
151 scopes: tuple[str, ...]
152 operator: str
153 source: str
154 mode: str
155 #: When the operator's consent was recorded — the parent's consent record timestamp.
156 delegated_at: str
157 #: The parent's own consent record, verbatim.
158 consent_record: Mapping[str, Any] = field(default_factory=dict)
160 def for_run(self, swarm_id: str, clusters: Iterable[str]) -> ConsentDelegation:
161 """The same consent, delegated to one run's clusters — the scopes do not change."""
162 return replace(self, swarm_id=swarm_id, clusters=tuple(clusters))
164 def to_dict(self) -> dict[str, Any]:
165 return {
166 "swarm_id": self.swarm_id,
167 "clusters": list(self.clusters),
168 "scopes": list(self.scopes),
169 "operator": self.operator,
170 "source": self.source,
171 "mode": self.mode,
172 "delegated_at": self.delegated_at,
173 "consent_record": dict(self.consent_record),
174 }
177def delegate_consent(
178 contract: Mapping[str, Any], *, swarm_id: str, cluster_ids: Iterable[str]
179) -> tuple[ConsentDelegation | None, str]:
180 """``(delegation, "")`` for an approved live contract, else ``(None, reason)``.
182 Refused, with the reason, when the contract is not approved (a scope is missing), when
183 it approved nothing a worker could be handed (``consent_mode: agent`` leaves approval
184 to a host agent, and ``swarm-run`` dispatches its workers itself), and when it does not
185 say who consented: a delegation records its operator, so an anonymous one is not made.
186 """
187 ok, message = consent.assert_operator_consent(dict(contract))
188 if not ok:
189 return None, message
190 record = contract.get("consent_record")
191 scopes = tuple(contract.get("delegated_agent_scope", {}).get("approved_mutation_scopes", ()))
192 if not isinstance(record, Mapping) or not scopes:
193 return None, (
194 f"operator consent is {contract.get('status')!r}, which approves no scope a "
195 "worker can be handed: swarm-run --live dispatches its workers itself, so the "
196 "operator approves the scopes explicitly (--approve-scope "
197 f"{','.join(required_scopes())} --operator NAME, or KEEL_APPROVE_SCOPE with "
198 "KEEL_OPERATOR under consent_mode: standing)"
199 )
200 operator = record.get("operator")
201 if not isinstance(operator, str) or not operator.strip():
202 return None, (
203 "operator consent names no operator, and a delegation records who consented: "
204 "pass --operator NAME beside --approve-scope"
205 )
206 delegation = ConsentDelegation(
207 swarm_id=swarm_id,
208 clusters=tuple(cluster_ids),
209 scopes=consent.normalize_scopes(scopes),
210 operator=operator,
211 source=str(record.get("source", "none")),
212 mode=str(contract.get("mode", "explicit")),
213 delegated_at=str(record.get("timestamp", "")),
214 consent_record=dict(record),
215 )
216 return delegation, ""
219def required_scopes() -> tuple[str, ...]:
220 """The consent scopes a live worker needs: those :data:`WORKER_SIDE_EFFECTS` map to."""
221 return consent.side_effect_scopes(WORKER_SIDE_EFFECTS)
224def worker_scopes(delegation: ConsentDelegation, cluster_id: str) -> tuple[str, ...]:
225 """The scopes one worker is handed: exactly the parent's, for a delegated cluster.
227 A cluster the delegation does not name gets nothing, so a worker started for a cluster
228 the operator's consent was not delegated to cannot mutate anything.
229 """
230 return delegation.scopes if cluster_id in delegation.clusters else ()
233def may(scopes: Iterable[str], side_effect: str) -> tuple[bool, str]:
234 """Whether a worker holding ``scopes`` may perform ``side_effect``, and why not."""
235 held = tuple(scopes)
236 missing = [s for s in consent.side_effect_scopes((side_effect,)) if s not in held]
237 if not missing:
238 return True, ""
239 return False, (
240 f"{side_effect} needs the {', '.join(missing)} consent scope, which this worker was "
241 f"not handed (it holds: {', '.join(held) or 'none'}); a worker never widens the "
242 "scopes the operator delegated"
243 )
246def consent_refusal(scopes: Iterable[str]) -> str:
247 """Why a worker holding ``scopes`` may not start, or ``""`` when it may.
249 Asked once, before the worker's first mutation, over every mutation it will make
250 (:data:`WORKER_SIDE_EFFECTS`): a worker that could commit and push but not open its
251 pull request would leave a pushed branch nobody asked for, so it does none of them.
252 """
253 held = tuple(scopes)
254 for side_effect in WORKER_SIDE_EFFECTS:
255 allowed, reason = may(held, side_effect)
256 if not allowed:
257 return reason
258 return ""
261def child_env(environ: Mapping[str, str]) -> dict[str, str]:
262 """``environ`` without the parent's consent variables, for a worker's children."""
263 return {key: value for key, value in environ.items() if key not in CONSENT_ENV_VARS}
266def worker_env(environ: Mapping[str, str]) -> dict[str, str]:
267 """keel's own local git steps and the gates: :func:`child_env` without the forge.
269 The gates run the code the implementer wrote, and ``keel run-gates --phases
270 guard,test --defer-jury`` reads nothing from GitHub, so they are not handed
271 :data:`FORGE_TOKEN_ENV_VARS`; nor are keel's ``git status``/``add``/``commit``, which
272 are local. git keeps the operator's own configuration here — only the implementer is
273 locked out of the remote (:func:`implementer_env`).
274 """
275 return {k: v for k, v in child_env(environ).items() if k not in FORGE_TOKEN_ENV_VARS}
278def _inherited_git_config(key: str) -> bool:
279 return key in _INHERITED_GIT_CONFIG_VARS or key.startswith(_INHERITED_GIT_CONFIG_PREFIXES)
282def implementer_env(environ: Mapping[str, str], *, gh_config_dir: str) -> dict[str, str]:
283 """The implementer seat's environment: the worker's, with no way to the forge.
285 :func:`worker_env` (no consent, no :data:`FORGE_TOKEN_ENV_VARS`), with git and ``gh``
286 locked out of the remote:
288 - ``GH_CONFIG_DIR`` is ``gh_config_dir``, a directory holding no login, so ``gh``
289 finds no stored account — neither in its config nor, having no host listed there,
290 in the system keyring — and says it is not logged in.
291 - ``GIT_TERMINAL_PROMPT=0`` and an empty ``GIT_ASKPASS`` leave git no one to ask for a
292 password (an empty ``GIT_ASKPASS`` also shadows ``core.askPass``, ``SSH_ASKPASS``
293 and an editor's askpass bridge inherited from the operator's terminal).
294 - :data:`IMPLEMENTER_GIT_CONFIG` clears the credential helpers and refuses every
295 network transport.
297 This removes the ambient credentials, so a brief that talks the seat into ``git
298 push`` or ``gh pr create`` fails. It is not a sandbox: the seat runs as the operator's
299 OS user and can read what that user can. See ``docs/keel/swarm.md``, trust notes.
300 """
301 env = {
302 key: value for key, value in worker_env(environ).items() if not _inherited_git_config(key)
303 }
304 env["GH_CONFIG_DIR"] = gh_config_dir
305 env["GIT_TERMINAL_PROMPT"] = "0"
306 env["GIT_ASKPASS"] = ""
307 env["GIT_CONFIG_COUNT"] = str(len(IMPLEMENTER_GIT_CONFIG))
308 for index, (key, value) in enumerate(IMPLEMENTER_GIT_CONFIG):
309 env[f"GIT_CONFIG_KEY_{index}"] = key
310 env[f"GIT_CONFIG_VALUE_{index}"] = value
311 return env
314#: What keel's own git steps after the implementer run under, besides a ``core.hooksPath``
315#: of an empty directory keel makes once the implementer has exited: the worktree shares
316#: the operator's repository, whose config and hooks the implementer could write.
317#: ``core.fsmonitor`` names a program ``git status``/``add``/``commit`` would run, and
318#: ``commit.gpgsign`` would run ``gpg.program``; both are off, so keel's commit is not
319#: signed. ``--no-verify`` on the commit and the push is the hooks' second lock.
320KEEL_GIT_CONFIG = (("core.fsmonitor", "false"), ("commit.gpgsign", "false"))
323def keel_git(args: Iterable[str], *, hooks_dir: str) -> list[str]:
324 """A git argv for one of keel's own steps after the implementer: no hooks, no
325 fsmonitor, no signing program, whatever the repository's config now says."""
326 argv = ["git", "-c", f"core.hooksPath={hooks_dir}"]
327 for key, value in KEEL_GIT_CONFIG:
328 argv += ["-c", f"{key}={value}"]
329 return [*argv, *args]
332#: Config the tamper check does not compare. ``branch.<name>.*`` is what git itself writes
333#: when a sibling worker cuts its worktree under ``branch.autoSetupMerge``, and none of
334#: keel's steps reads it: the push names its URL and its refspec.
335TAMPER_IGNORED_CONFIG_PREFIXES = ("branch.",)
338@dataclass(frozen=True)
339class GitSnapshot:
340 """What of the repository's git setup keel's own steps would run under.
342 ``config`` is ``git config --list --show-origin --show-scope`` — every scope, with the
343 files ``include.path``/``includeIf`` pull in — one line per entry, in order; ``hooks``
344 maps each file under the hooks directory to its digest.
345 """
347 git_dir: str
348 common_dir: str
349 hooks_dir: str
350 config: tuple[str, ...]
351 hooks: Mapping[str, str] = field(default_factory=dict)
354#: Config sections whose subsection is a URL or a condition — ``url.<base>.insteadOf``,
355#: ``credential.<url>.helper``, ``http.<url>.extraHeader``, ``includeIf.<cond>.path`` —
356#: which a finding names without it: a URL can carry a credential.
357_URL_SUBSECTIONS = ("url", "credential", "http", "includeif")
360def _config_entry(line: str) -> tuple[str, str]:
361 """``(scope, key)`` of one ``--show-scope --show-origin`` line; never its value."""
362 parts = line.split("\t", 2)
363 key = parts[-1].split("=", 1)[0]
364 section, _, rest = key.partition(".")
365 if section in _URL_SUBSECTIONS and "." in rest:
366 key = f"{section}.<{section}>.{rest.rsplit('.', 1)[1]}"
367 return (parts[0] if len(parts) == 3 else ""), key
370def _compared(config: tuple[str, ...]) -> list[str]:
371 return [
372 line
373 for line in config
374 if not _config_entry(line)[1].startswith(TAMPER_IGNORED_CONFIG_PREFIXES)
375 ]
378def tamper_findings(before: GitSnapshot, after: GitSnapshot) -> tuple[str, ...]:
379 """What changed between two snapshots, named without the values (a config value can
380 be a credential); ``()`` when nothing did."""
381 findings = [
382 f"the {label} moved"
383 for label, old, new in (
384 ("git directory", before.git_dir, after.git_dir),
385 ("common git directory", before.common_dir, after.common_dir),
386 ("hooks directory", before.hooks_dir, after.hooks_dir),
387 )
388 if old != new
389 ]
390 old, new = _compared(before.config), _compared(after.config)
391 if old != new:
392 changed = (Counter(old) - Counter(new)) + (Counter(new) - Counter(old))
393 keys = sorted({" ".join(_config_entry(line)).strip() for line in changed.elements()})
394 findings.append(
395 f"git config changed: {', '.join(keys)}" if keys else "git config was reordered"
396 )
397 for name in sorted(set(before.hooks) | set(after.hooks)):
398 if name not in before.hooks:
399 findings.append(f"hook {name} added")
400 elif name not in after.hooks:
401 findings.append(f"hook {name} removed")
402 elif before.hooks[name] != after.hooks[name]:
403 findings.append(f"hook {name} changed")
404 return tuple(findings)
407def tamper_reason(findings: Iterable[str], *, during: str) -> str:
408 """Why a worker stops at ``tamper``."""
409 return (
410 f"the repository's git setup changed while {during} ({'; '.join(findings)}); keel's "
411 "own git steps run with the operator's credentials and would run under it, so the "
412 "worker stops and pushes nothing"
413 )
416def plan_implementer(
417 assignment: Mapping[str, Any] | None,
418 *,
419 config: Any,
420 registry: providers_mod.Registry | None,
421 prompt_path: str,
422 cwd: str,
423 timeout: int = delegate.DEFAULT_TIMEOUT_S,
424) -> tuple[delegate.RunPlan | None, str]:
425 """The cluster's implementer seat as a delegate plan, or ``(None, reason)``.
427 The seat is the one :func:`keel.team.resolve_assignment` resolved for the cluster —
428 ``--delegate`` over a ``--team`` or difficulty bench over ``knobs.team.implement`` over
429 the host default — so ``swarm-plan`` shows the seat ``swarm-run --live`` dispatches. It
430 is resolved and planned exactly as ``keel delegate run --provider <seat> --role
431 implement`` would, with the cluster's worktree as ``cwd``.
432 """
433 if not assignment:
434 return None, "the cluster has no resolved assignment, so it has no implementer seat"
435 seat = assignment["implementer"]
436 if seat["kind"] != "provider":
437 return None, (
438 f"the implementer seat {seat['provider']!r} (from {seat.get('source')}) is a host "
439 "subagent, which only an agent host can spawn; pass --delegate <provider> or name "
440 "a provider in knobs.team.implement for a live swarm"
441 )
442 token = seat["provider"] + (f":{seat['model']}" if seat.get("model") else "")
443 try:
444 resolution = delegate.resolve_provider(config, registry, token)
445 plan = delegate.plan_run(
446 resolution.provider,
447 "implement",
448 prompt_path,
449 cwd,
450 timeout,
451 seat.get("effort"),
452 resolution.model,
453 profile=resolution.profile,
454 )
455 except delegate.DelegateError as exc:
456 return None, f"the implementer seat {token!r} cannot be planned ({exc.code}): {exc.message}"
457 if plan.transport != TOOL_TRANSPORT:
458 return None, (
459 f"the implementer seat {token!r} runs over the {plan.transport!r} transport, "
460 "which cannot edit a worktree; a live worker needs an agent CLI "
461 "(claude, codex, agy) — pass --delegate"
462 )
463 return plan, ""
466def _issue_line(number: int, scopes: Mapping[int, Any]) -> tuple[str, str]:
467 scope = scopes.get(number)
468 title = getattr(scope, "title", "") or ""
469 body = getattr(scope, "body", "") or ""
470 return title.strip(), body.strip()
473def render_brief(
474 cluster: Any,
475 issue_scopes: Mapping[int, Any],
476 *,
477 swarm_id: str,
478 branch: str,
479 base_branch: str,
480) -> str:
481 """The implementer's brief for one cluster: its issues, its scope, and its limits."""
482 scope = ", ".join(cluster.combined_scope) or "*"
483 lines = [
484 f"You are the implementer for cluster {cluster.cluster_id} of keel swarm {swarm_id}.",
485 (
486 f"Your working directory is this cluster's own git worktree, on branch {branch}, "
487 f"cut from {base_branch}."
488 ),
489 "",
490 "Implement the issue(s) below completely, in this working tree:",
491 f"- Change only files inside this working tree, within the cluster's scope: {scope}.",
492 "- Add or update tests for what you change, and leave the project's tests passing.",
493 (
494 "- Do not commit, push, open a pull request, or run any command that changes the "
495 "repository's history or its remote. keel commits your changes, runs the "
496 "project's gates, pushes the branch and opens the pull request itself."
497 ),
498 (
499 "- You have no GitHub credentials, and git here cannot reach any remote: `git "
500 "push`, `git fetch` and `gh` fail by design, so do not work around them."
501 ),
502 (
503 "- Do not change git's configuration or hooks: keel neither commits nor pushes "
504 "the work of an implementer that does."
505 ),
506 (
507 "- If an issue cannot be implemented as written, leave it unchanged and say why "
508 "in your final message."
509 ),
510 ]
511 for number in cluster.issues:
512 title, body = _issue_line(number, issue_scopes)
513 lines += ["", f"## Issue #{number}: {title or '(title unavailable)'}", ""]
514 lines.append(body or "(The issue body could not be read; work from the title.)")
515 return "\n".join(lines) + "\n"
518def _issue_refs(cluster: Any) -> str:
519 return ", ".join(f"#{n}" for n in cluster.issues)
522def commit_message(cluster: Any, *, swarm_id: str, plan: delegate.RunPlan) -> str:
523 """The cluster commit: what it implements, and who wrote it."""
524 system = plan.attribution.get("system") or plan.provider
525 lines = [
526 f"feat(swarm): implement {_issue_refs(cluster)} ({swarm_id}/{cluster.cluster_id})",
527 "",
528 f"Written by the implementer seat {system} through keel swarm-run --live.",
529 "",
530 ]
531 lines += [f"Refs #{n}" for n in cluster.issues]
532 return "\n".join(lines) + "\n"
535def pull_request_title(cluster: Any, issue_scopes: Mapping[int, Any], *, swarm_id: str) -> str:
536 """One pull request per cluster; a one-issue cluster is titled after its issue."""
537 if len(cluster.issues) == 1:
538 title, _body = _issue_line(cluster.issues[0], issue_scopes)
539 if title:
540 return f"{title} (#{cluster.issues[0]}, swarm {swarm_id}/{cluster.cluster_id})"
541 return f"swarm {swarm_id}/{cluster.cluster_id}: implement {_issue_refs(cluster)}"
544def pull_request_body(
545 cluster: Any,
546 *,
547 swarm_id: str,
548 branch: str,
549 base_branch: str,
550 commit: str,
551 plan: delegate.RunPlan,
552 seat_source: str | None,
553 delegation: ConsentDelegation,
554) -> str:
555 """The cluster's pull request body: the work, the seat, the gates and the consent.
557 ``Refs``, never ``Closes``: the pull request carries no review evidence yet, and an
558 issue is closed by the steps that review and land the cluster, not by opening it.
559 """
560 system = plan.attribution.get("system") or plan.provider
561 lines = [f"Implements cluster `{cluster.cluster_id}` of keel swarm `{swarm_id}`.", ""]
562 lines += [f"Refs #{n}" for n in cluster.issues]
563 lines += [
564 "",
565 f"- **Branch:** `{branch}`, cut from `{base_branch}`, at `{commit}`",
566 f"- **Implementer:** `{system}` (seat from `{seat_source or 'unknown'}`)",
567 (
568 f"- **Gates:** `keel run-gates --phases {GATE_PHASES} --defer-jury` passed in the "
569 "cluster worktree at that commit"
570 ),
571 (
572 f"- **Consent:** delegated by `{delegation.operator}` ({delegation.source}, "
573 f"{delegation.delegated_at}) for `{', '.join(delegation.scopes)}`"
574 ),
575 "",
576 (
577 "This pull request carries no review evidence yet: review and landing are not "
578 "part of `swarm-run`."
579 ),
580 ]
581 return "\n".join(lines) + "\n"
584#: ``https://<host>/<owner>/<repo>/pull/<n>``: the URL ``gh pr create`` prints last.
585_PULL_REQUEST_REPO_URL = re.compile(r"^https?://[^/\s]+/([^/\s]+/[^/\s]+)/pull/[1-9][0-9]*/?$")
588def pull_request_repo(url: str) -> str | None:
589 """The ``owner/repo`` a ``gh pr create`` URL names, or ``None`` when it names none."""
590 match = _PULL_REQUEST_REPO_URL.match(url.strip())
591 return match.group(1) if match else None
594def provenance_run_id(swarm_id: str, cluster_id: str) -> str:
595 """The run id a cluster's ship-provenance comment carries: one per cluster of a run."""
596 return f"{swarm_id}/{cluster_id}"
599def ship_provenance_body(
600 cluster: Any, *, swarm_id: str, commit: str, plan: delegate.RunPlan
601) -> str:
602 """The ship-provenance comment for a cluster's pull request — the one a live ``keel
603 ship`` run posts on its own (:func:`keel.artifacts.render_ship_provenance`).
605 A keel-made pull request arms ``keel merge``'s evidence gate itself, the way a ship run's
606 does: the branch ``swarm/<id>/<cluster>`` matches no ship-branch pattern, so without this
607 comment the gate reads the pull request as not a keel run, and ``swarm-land`` is held on
608 "evidence gate is not enforced" rather than on the review evidence it lacks. The
609 attribution is the seat's, verbatim from :mod:`keel.agents` (``plan.attribution``).
610 """
611 return artifacts.render_ship_provenance(
612 run_id=provenance_run_id(swarm_id, cluster.cluster_id),
613 issue=cluster.issues[0] if cluster.issues else None,
614 head_sha=commit,
615 implementer_attribution=dict(plan.attribution),
616 )
619def provenance_warning(pull_request: str, why: str) -> str:
620 """What a worker reports when its pull request is open but not stamped."""
621 return (
622 f"the ship-provenance comment was not posted on {pull_request} ({why}); keel merge "
623 'will hold it as "evidence gate is not enforced" until the comment is posted '
624 "(`keel post-comment --artifact ship-provenance`) or its review verdicts are"
625 )
628def attribution_labels(plan: delegate.RunPlan) -> tuple[str, ...]:
629 """The labels a cluster pull request carries for the seat that implemented it (#1420).
631 ``agent_label`` and ``model_label`` verbatim from ``plan.attribution`` — the record
632 :mod:`keel.agents` computed for the seat, the one ``keel attribution`` prints — never
633 composed here. A seat with no model has no ``model_label``, and only what exists is
634 applied.
635 """
636 return tuple(
637 label
638 for label in (plan.attribution.get("agent_label"), plan.attribution.get("model_label"))
639 if isinstance(label, str) and label.strip()
640 )
643def label_names(listing: str) -> tuple[str, ...]:
644 """The label names in a ``gh label list --json name`` answer; ``()`` when unreadable."""
645 try:
646 rows = json.loads(listing)
647 except ValueError:
648 return ()
649 if not isinstance(rows, list):
650 return ()
651 return tuple(
652 row["name"] for row in rows if isinstance(row, dict) and isinstance(row.get("name"), str)
653 )
656def labels_warning(pull_request: str, labels: Iterable[str], why: str) -> str:
657 """What a worker reports when its pull request is open but not labelled (#1420)."""
658 wanted = ", ".join(labels) or "agent:<vendor>"
659 return (
660 f"the attribution labels ({wanted}) were not applied to {pull_request} ({why}); keel "
661 "merge will hold it on attribution-label until its agent:<vendor> label is applied — "
662 "apply the labels `keel attribution` prints for the implementer seat"
663 )
666#: What a live worker's ``ship_run`` record says about the merge (#1420). The worker never
667#: assesses one: it opened the pull request, and its review evidence, CI and ``keel merge``
668#: decide whether it lands. ``defer`` is the one merge action no reader takes as a merge
669#: (:func:`keel.closeorder.record_attests_merge` reads ``merge`` as one, and ``keel status``
670#: reads anything but ``defer``/``block``/``skip`` as shipped).
671WORKER_MERGE_REASON = (
672 "recorded by a swarm worker when it opened the pull request: the gates it ran, no "
673 "merge assessment; review evidence and keel merge decide whether it lands"
674)
676#: The ``schema_version`` of the report ``keel run-gates --json`` prints.
677RUN_GATES_SCHEMA = "keel.run-gates.v1"
680def parse_gate_report(output: str) -> tuple[GateOutcome, ...] | None:
681 """The gate outcomes in ``keel run-gates --json`` output, or ``None`` when it has none.
683 The worker's runner folds the child's stderr into its stdout, so the report may follow
684 an extension or capability notice: it is read from the first line that opens a JSON
685 object and decodes to a ``keel.run-gates.v1`` report. Each outcome is restored with the
686 fields :func:`keel.ledger.build_ship_run_record` records and
687 :func:`keel.ledger.record_gates_passed` judges. Anything unreadable — a missing gate id,
688 a finding with no known severity — makes the whole report unreadable, never a partial
689 one: a gate dropped from the record would be a gate the record silently passed.
690 """
691 decoder = json.JSONDecoder()
692 lines = output.splitlines(keepends=True)
693 offset = 0
694 for line in lines:
695 start, offset = offset, offset + len(line)
696 if line.rstrip() != "{":
697 continue
698 try:
699 report, _end = decoder.raw_decode(output, start)
700 except ValueError:
701 continue
702 if isinstance(report, dict) and report.get("schema_version") == RUN_GATES_SCHEMA:
703 return _gate_outcomes(report.get("gate_outcomes"))
704 return None
707def _gate_outcomes(entries: Any) -> tuple[GateOutcome, ...] | None:
708 if not isinstance(entries, list):
709 return None
710 outcomes: list[GateOutcome] = []
711 for entry in entries:
712 outcome = _gate_outcome(entry)
713 if outcome is None:
714 return None
715 outcomes.append(outcome)
716 return tuple(outcomes)
719def _gate_outcome(entry: Any) -> GateOutcome | None:
720 """One ``run-gates`` outcome as a :class:`keel.gates.GateOutcome`, or ``None``."""
721 if not isinstance(entry, dict) or not isinstance(entry.get("gate"), str):
722 return None
723 found: list[fnd.Finding] = []
724 for raw in entry.get("findings") or ():
725 severity = raw.get("severity") if isinstance(raw, dict) else None
726 if severity not in fnd.SEVERITIES:
727 return None
728 source = raw.get("source")
729 found.append(
730 fnd.Finding(
731 severity,
732 str(raw.get("message") or ""),
733 source if isinstance(source, str) and source else entry["gate"],
734 )
735 )
736 error = entry.get("error")
737 on_fail = entry.get("on_fail")
738 return GateOutcome(
739 entry["gate"],
740 entry.get("ok") is True,
741 tuple(found),
742 error=error if isinstance(error, str) else None,
743 skipped=entry.get("skipped") is True,
744 timed_out=entry.get("timed_out") is True,
745 not_run=entry.get("not_run") is True,
746 # Strict when absent, as `record_gates_passed` reads it.
747 on_fail=on_fail if isinstance(on_fail, str) else "block",
748 unconfigured=entry.get("unconfigured") is True,
749 )
752def gate_report_summary(outcomes: Iterable[GateOutcome]) -> str:
753 """The gates a red run judged, one line each, then the findings — as ``run-gates``
754 prints them without ``--json``, for the reason a worker stopped at ``gates``."""
755 lines: list[str] = []
756 found: list[fnd.Finding] = []
757 for outcome in outcomes:
758 lines.append(f" {_gate_word(outcome):>7} {outcome.gate}")
759 found.extend(outcome.findings)
760 lines += [f" [{f.severity}] {f.source}: {f.message}" for f in found]
761 return "\n".join(lines)
764def _gate_word(outcome: GateOutcome) -> str:
765 if outcome.not_run:
766 return "NOT-RUN"
767 if outcome.skipped:
768 return "SKIPPED"
769 if outcome.ok:
770 return "ok"
771 return "TIMEOUT" if outcome.timed_out else "FAIL"
774def gates_record(
775 cluster: Any,
776 *,
777 swarm_id: str,
778 plan: delegate.RunPlan,
779 branch: str,
780 base_branch: str,
781 head: str,
782 pull_request: int,
783 outcomes: Iterable[GateOutcome],
784 changed_files: list[str] | None,
785 config: Any = None,
786) -> dict[str, Any]:
787 """The ``ship_run`` record of the gates a worker ran on the head it pushed (#1420).
789 ``keel merge`` lands a pull request only with a ``ship_run`` record whose gates passed
790 for its current head (:func:`keel.ledger.gates_pass_for_head`). The worker knows all of
791 it once ``gh pr create`` returns: the gates it just judged, the pull request, the head.
792 Built by :func:`keel.ledger.build_ship_run_record`, as ``keel ship --live
793 --append-ledger`` builds its own, from the outcomes exactly as the worker's run reported
794 them — so a blocking gate that run did not execute (the deferred jury, a ``pre-merge``
795 gate outside ``--phases``) is recorded ``not_run``, and the record is not a pass.
797 What it says and does not say:
799 - ``capture.not_run``: the run never reached capture, so no capture reader counts the
800 record as a merged pull request (``--capture-status not-run``).
801 - ``assessment.merge.action: defer`` (:data:`WORKER_MERGE_REASON`), and no tier, review
802 count, window or CI: the worker assessed none of them.
803 - ``actors.implementer``: the seat's ``system``, the string its attribution labels are
804 derived from, so the evidence gate's cross-check agrees with the labels applied.
805 - No consent: a live worker's consent is the run's delegation, which the
806 ``consent_delegation`` record carries pinned to the pushed head. A consent status here
807 would win over it in ``keel consent-verify`` without that pin.
808 - The cluster's first issue, as the provenance comment names it: the record has one.
809 """
810 listed = list(outcomes)
811 return ledger.build_ship_run_record(
812 command="swarm-run",
813 base_branch=base_branch,
814 changed_files=changed_files,
815 outcomes=listed,
816 verdict=fnd.summarize(collect_findings(listed)),
817 assessment=ship.ShipAssessment(
818 tier=None, # type: ignore[arg-type]
819 reviewers=None, # type: ignore[arg-type]
820 window_open=None, # type: ignore[arg-type]
821 ci_ok=None,
822 merge=ship.MergeDecision("defer", WORKER_MERGE_REASON),
823 ),
824 target=f"PR #{pull_request}",
825 run_id=provenance_run_id(swarm_id, cluster.cluster_id),
826 issue_number=cluster.issues[0] if cluster.issues else None,
827 pr_number=pull_request,
828 branch=branch,
829 head_sha=head,
830 capture_status=None,
831 capture_not_run=True,
832 config=config,
833 implementer=plan.attribution.get("system") or plan.provider,
834 )
837def gates_not_a_pass(record: Mapping[str, Any]) -> tuple[str, ...]:
838 """The gates that keep a worker's recorded run from counting as a pass, by name.
840 Empty when :func:`keel.ledger.record_gates_passed` accepts the record. A gate is named
841 when it did not run clean or skipped, or when it is a blocking gate the run did not
842 execute; a record with no gate at all names ``(no gate)``.
843 """
844 if ledger.record_gates_passed(dict(record)):
845 return ()
846 gates = [g for g in record.get("gates") or () if isinstance(g, dict)]
847 named = tuple(
848 str(g.get("gate"))
849 for g in gates
850 if g.get("error")
851 or not (g.get("ok") is True or g.get("skipped") is True)
852 or (g.get("not_run") is True and g.get("on_fail") not in ("warn", "suggest"))
853 )
854 return named or ("(no gate)",)
857def gates_not_a_pass_warning(pull_request: str, gates: Iterable[str]) -> str:
858 """What a worker reports when the gates it recorded do not count as a pass."""
859 return (
860 f"the gates recorded for {pull_request} do not count as a gates-pass "
861 f"({', '.join(gates)} did not pass or did not run in the worker); keel merge will hold "
862 "it on no gates-pass until a run of every blocking gate passes on its head (`keel "
863 "ship --live --append-ledger --capture-status not-run --pull-request <n> --head-sha "
864 "<sha>`)"
865 )
868def gates_record_warning(pull_request: str, why: str) -> str:
869 """What a worker reports when it could not record its gates-pass (#1420)."""
870 return (
871 f"the gates-pass for {pull_request} is not in the run ledger ({why}); keel merge will "
872 "hold it on no gates-pass until one is recorded for its head (`keel ship --live "
873 "--append-ledger --capture-status not-run --pull-request <n> --head-sha <sha>`)"
874 )
877@dataclass(frozen=True)
878class WorktreeDisposal:
879 """What becomes of a live worker's worktree and branch when the worker ends (#1278)."""
881 #: ``remove``, ``keep``, or ``none`` when the worker created no worktree to act on.
882 worktree: str
883 delete_branch: bool
884 #: Why, in words the run result and the state file can carry.
885 reason: str
888def worktree_disposal(
889 *, ok: bool, worktree_created: bool, implementer_ran: bool
890) -> WorktreeDisposal:
891 """Decide what a live worker leaves behind (#1278).
893 - A worker that never created its worktree touches nothing: what sits at its path, and
894 the branch of that name, are not its own — a previous run's kept worktree, say.
895 - A worker that succeeded has an open pull request headed by its branch: the worktree
896 goes, the branch stays.
897 - A worker that failed after its implementer ran keeps both, so the operator can
898 inspect what the seat wrote, the commit the gates rejected, or the branch a pull
899 request could not be opened for. ``keel swarm-status --clean`` removes them later.
900 - A worker that failed before its implementer ran (its push URL or git setup could not
901 be read) pushed nothing and has nothing to inspect: the worktree goes, and so does
902 the branch keel cut for it.
903 """
904 if not worktree_created:
905 return WorktreeDisposal("none", False, "no worktree was created")
906 if ok:
907 return WorktreeDisposal("remove", False, "its branch heads the open pull request")
908 if implementer_ran:
909 return WorktreeDisposal("keep", False, "kept for inspection: the worker failed")
910 return WorktreeDisposal(
911 "remove", True, "the worker failed before its implementer ran; nothing was pushed"
912 )
915@dataclass(frozen=True)
916class WorktreeEntry:
917 """One block of ``git worktree list --porcelain``."""
919 path: str
920 branch: str | None
921 #: git's own verdict that the worktree's directory is gone.
922 prunable: bool
925def parse_worktree_list(porcelain: str) -> tuple[WorktreeEntry, ...]:
926 """``git worktree list --porcelain`` as entries; text outside a block is ignored."""
927 entries: list[WorktreeEntry] = []
928 for block in porcelain.replace("\r\n", "\n").split("\n\n"):
929 lines = [line for line in block.split("\n") if line]
930 if not lines or not lines[0].startswith("worktree "):
931 continue
932 branch = next((ln[len("branch ") :] for ln in lines if ln.startswith("branch ")), None)
933 entries.append(
934 WorktreeEntry(
935 path=lines[0][len("worktree ") :],
936 branch=branch,
937 prunable=any(ln == "prunable" or ln.startswith("prunable ") for ln in lines),
938 )
939 )
940 return tuple(entries)
943#: The namespace keel cuts a live worker's branch in: ``swarm/<swarm_id>/<cluster_id>``.
944SWARM_BRANCH_PREFIX = "swarm/"
947def swarm_branch_ids(branch: str) -> tuple[str, str] | None:
948 """``(swarm_id, cluster_id)`` of a branch keel's swarm cut, or ``None`` for any other.
950 ``branch`` may be the short name or ``refs/heads/…``. Only the exact shape
951 ``swarm/<id>/<cluster>`` counts, so a branch someone keeps under ``swarm/`` at another
952 depth is never taken for keel's.
953 """
954 name = branch.removeprefix("refs/heads/")
955 if not name.startswith(SWARM_BRANCH_PREFIX):
956 return None
957 parts = name[len(SWARM_BRANCH_PREFIX) :].split("/")
958 if len(parts) != 2 or not all(parts):
959 return None
960 return parts[0], parts[1]
963@dataclass(frozen=True)
964class SwarmRunRecord:
965 """What a run's state file says about its leftovers."""
967 #: No ``completed_at`` yet, or a state file keel cannot read: the run may be running.
968 unfinished: bool
969 #: ``cluster_id`` -> the pull request its worker opened.
970 pull_requests: Mapping[str, int] = field(default_factory=dict)
971 #: Clusters whose worker pushed its branch (with or without a pull request).
972 pushed: frozenset[str] = frozenset()
975@dataclass(frozen=True)
976class SwarmLeftover:
977 """One thing a swarm run left under keel's own paths or branch namespace (#1278)."""
979 #: ``worktree`` (registered), ``directory`` (on disk, not registered), ``registration``
980 #: (registered, its directory gone) or ``branch``.
981 kind: str
982 swarm_id: str
983 #: Empty for a run's own ``.keel/worktrees/<swarm_id>/`` directory.
984 cluster_id: str
985 #: The path, or the branch name.
986 target: str
987 #: ``remove`` or ``keep``.
988 action: str
989 reason: str
991 def to_dict(self) -> dict[str, str]:
992 return {
993 "kind": self.kind,
994 "swarm_id": self.swarm_id,
995 "cluster_id": self.cluster_id,
996 "target": self.target,
997 "action": self.action,
998 "reason": self.reason,
999 }
1002LEFTOVER_KINDS = ("registration", "worktree", "directory", "branch")
1005def _branch_keep_reason(record: SwarmRunRecord | None, cluster_id: str) -> str:
1006 if record is None:
1007 return "no run state says whether it was pushed; delete it by hand if it was not"
1008 if cluster_id in record.pull_requests:
1009 return f"it heads pull request #{record.pull_requests[cluster_id]}"
1010 if cluster_id in record.pushed:
1011 return "it was pushed"
1012 return ""
1015def classify_leftovers(
1016 *,
1017 worktrees: Iterable[tuple[str, str, str, bool]],
1018 directories: Iterable[tuple[str, str, str]],
1019 branches: Iterable[str],
1020 runs: Mapping[str, SwarmRunRecord],
1021 named: str | None = None,
1022) -> tuple[SwarmLeftover, ...]:
1023 """Which of the swarm's leftovers ``keel swarm-status --clean`` may remove (#1278).
1025 ``worktrees`` are the registered worktrees under ``.keel/worktrees/`` as
1026 ``(swarm_id, cluster_id, path, prunable)``; ``directories`` what sits there on disk and
1027 is not registered, as ``(swarm_id, cluster_id, path)`` (``cluster_id`` empty for an
1028 empty run directory); ``branches`` local branch names, of which only keel's swarm
1029 namespace is considered; ``runs`` what each run's state file says; ``named`` the run
1030 the operator named with ``--swarm-id``.
1032 keel cannot tell a killed run from a running one: both have a state file with no
1033 ``completed_at``. Such a run's leftovers are kept unless the operator names the run —
1034 naming it is their statement that no ``swarm-run`` of it is still running. A branch is
1035 also kept when its worker recorded a push or a pull request (the pull request is headed
1036 by it), and when no state file says either way.
1037 """
1039 def in_flight(swarm_id: str) -> str:
1040 record = runs.get(swarm_id)
1041 if record is None or not record.unfinished or swarm_id == named:
1042 return ""
1043 return (
1044 f"swarm run {swarm_id!r} has not completed; if no swarm-run of it is still "
1045 f"running, name it with --swarm-id {swarm_id}"
1046 )
1048 found: list[SwarmLeftover] = []
1049 for swarm_id, cluster_id, path, prunable in worktrees:
1050 if prunable:
1051 why = "its directory is gone; git worktree prune drops the registration"
1052 found.append(SwarmLeftover("registration", swarm_id, cluster_id, path, "remove", why))
1053 continue
1054 why = in_flight(swarm_id)
1055 action = "keep" if why else "remove"
1056 why = why or "no running swarm run owns it"
1057 found.append(SwarmLeftover("worktree", swarm_id, cluster_id, path, action, why))
1058 for swarm_id, cluster_id, path in directories:
1059 why = in_flight(swarm_id)
1060 action = "keep" if why else "remove"
1061 why = why or ("an empty run directory" if not cluster_id else "not a registered worktree")
1062 found.append(SwarmLeftover("directory", swarm_id, cluster_id, path, action, why))
1063 for branch in branches:
1064 ids = swarm_branch_ids(branch)
1065 if ids is None:
1066 continue
1067 swarm_id, cluster_id = ids
1068 why = in_flight(swarm_id) or _branch_keep_reason(runs.get(swarm_id), cluster_id)
1069 action = "keep" if why else "remove"
1070 why = why or "nothing was pushed and no pull request was opened"
1071 name = branch.removeprefix("refs/heads/")
1072 found.append(SwarmLeftover("branch", swarm_id, cluster_id, name, action, why))
1073 return tuple(
1074 sorted(
1075 found,
1076 key=lambda x: (x.swarm_id, x.cluster_id, LEFTOVER_KINDS.index(x.kind), x.target),
1077 )
1078 )