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

907 statements  

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

1"""Keel Swarm — Deterministic static dependency analysis & conflict clustering. 

2 

3Pure, stdlib-first dependency graph analysis for multi-agent parallel execution. 

4Partitions candidate issues into orthogonal (disjoint) Waves and independent Clusters, 

5enabling Direct Batch Landing for non-overlapping diff trees. A wave that depends on 

6an earlier one is refused at landing until that wave lands and the rest is re-planned. 

7 

8All functions here are pure and deterministic — no subprocess, no network. 

9""" 

10 

11from __future__ import annotations 

12 

13import datetime 

14import fnmatch 

15import json 

16import posixpath 

17import re 

18from collections.abc import Mapping, Sequence 

19from dataclasses import dataclass, field, replace 

20from pathlib import Path 

21from typing import Any 

22 

23from . import agents, classify, ship, swarm_worker, workspace 

24from . import team as team_policy 

25from .config import ProjectConfig 

26 

27#: Difficulty points by resolved risk tier. A tier-3 cluster is not merely riskier, it 

28#: is more *work*: it touches the globs the project declared load-bearing, so the change 

29#: has to be argued for as well as written. 

30_TIER_POINTS = {1: 0, 2: 2, 3: 4} 

31 

32#: ``(at least this many predicted files, points)``, widest first. Predicted breadth is 

33#: the only size signal available before a line is written, and it is the one that decides 

34#: whether a cluster is one edit or a refactor. 

35_FILE_COUNT_POINTS = ((8, 3), (4, 2), (2, 1)) 

36 

37#: Label -> points. Only labels a human deliberately set count: ``priority:`` says how 

38#: much it matters (and so how much review attention it will draw), ``size:`` is the 

39#: estimate a human already made and is worth more than any heuristic here. 

40_LABEL_POINTS = { 

41 "priority:critical": 2, 

42 "priority:high": 1, 

43 "size:m": 1, 

44 "size:l": 2, 

45 "size:xl": 3, 

46} 

47 

48#: Cap on the points a cluster earns for depending on already-scheduled work. Depth says 

49#: the cluster lands on a tree someone else just moved; past a few dependencies that is 

50#: the same problem, not a worse one. 

51MAX_DEPENDENCY_POINTS = 3 

52 

53#: ``(highest score still in this band, band)``, lightest first; anything above the last 

54#: ceiling is the heaviest band. The bands themselves are :data:`keel.team.DIFFICULTY_BANDS`, 

55#: because a band's whole purpose is to name a ``knobs.team.by_difficulty`` bench. 

56_BAND_CEILINGS = ((2, "easy"), (5, "standard")) 

57 

58#: Regex to capture backticked file paths or common file paths in issue text 

59_PATH_BACKTICK_RE = re.compile(r"`([a-zA-Z0-9_\-./]+\.[a-zA-Z0-9_\-]+)`") 

60_PATH_GENERAL_RE = re.compile( 

61 r"(?:^|[\s(\[])([a-zA-Z0-9_\-./]+/(?:[a-zA-Z0-9_\-./]+\.[a-zA-Z0-9_\-]+|[a-zA-Z0-9_\-]+/|\*))" 

62) 

63 

64#: The scope of an issue nobody described: everything. It intersects every other path 

65#: (:func:`_normalized_paths_intersect`), so an undescribed issue is serialised against 

66#: every other issue instead of being assumed disjoint from all of them (#1274). The 

67#: fallback used to be a unique ``scope/issue-N/*`` per issue, which by construction 

68#: conflicted with nothing and put every undescribed issue in wave 1. 

69SCOPE_ANY = "*" 

70 

71#: Where an issue's scope came from, strongest first. ``override`` is ``--issue-scope``; 

72#: ``issue-body`` a ``Scope:`` declaration in the issue; ``area-label`` an ``area:<name>`` 

73#: label mapped through ``policy_pack.scan.areas``; ``declared-file`` the operator's 

74#: ``--declared-file`` (plus the paths the text names); ``default`` means nothing declared 

75#: a scope, so it holds :data:`SCOPE_ANY` beside whatever paths the text happened to name. 

76SCOPE_SOURCES = ("override", "issue-body", "area-label", "declared-file", "default") 

77 

78#: ``Scope: a, b`` — the one-line declaration. 

79_SCOPE_LINE_RE = re.compile(r"^\s*scope\s*:(.*)$", re.IGNORECASE) 

80#: ``## Scope`` — the heading form, any level; the bullet list under it is the scope. 

81_SCOPE_HEADING_RE = re.compile(r"^\s{0,3}#{1,6}\s*scope\s*#*\s*$", re.IGNORECASE) 

82_SCOPE_BULLET_RE = re.compile(r"^\s*[-*+]\s+(.*)$") 

83_FENCE_RE = re.compile(r"^\s{0,3}(```|~~~)") 

84#: A token counts as a glob only if it looks like one — a bare word is the prose around 

85#: the list (`src/a.py — the parser`), not a path. ``./Makefile`` names a bare file. 

86_GLOB_MARKS = frozenset("/.*?[") 

87 

88#: The two modes a wave is planned in (:func:`wave_landing_mode`). 

89WAVE_MODES = ("orthogonal_parallel", "sequential_dependent") 

90 

91 

92class SwarmPlanError(ValueError): 

93 """A persisted swarm plan keel will not land from (#1275). 

94 

95 Unreadable, not JSON, the wrong shape, or a schema version this keel does not know. 

96 The loader refuses rather than guessing: a plan read half-right lands a wave nobody 

97 ran, which is the defect persisting the plan exists to remove. 

98 """ 

99 

100 

101#: A plan key that is an issue number exactly as :meth:`SwarmPlan.to_dict` writes one — 

102#: ``str(int)`` — so a key that parses but would be written back differently (``"07"``, 

103#: ``"+7"``, ``" 7"``) is refused instead of silently renamed. 

104_ISSUE_KEY_RE = re.compile(r"\A(0|-?[1-9][0-9]*)\Z") 

105 

106 

107def _plan_object(value: Any, where: str) -> Mapping[str, Any]: 

108 if not isinstance(value, dict): 

109 raise SwarmPlanError(f"{where}: expected an object, found {type(value).__name__}") 

110 return value 

111 

112 

113def _plan_field(data: Mapping[str, Any], key: str, where: str) -> Any: 

114 if key not in data: 

115 raise SwarmPlanError(f"{where}: missing '{key}'") 

116 return data[key] 

117 

118 

119def _plan_int(data: Mapping[str, Any], key: str, where: str) -> int: 

120 value = _plan_field(data, key, where) 

121 # `bool` is an `int` to isinstance; `true` in a count is a malformed file, not 1. 

122 if isinstance(value, bool) or not isinstance(value, int): 

123 raise SwarmPlanError(f"{where}.{key}: expected an integer, found {value!r}") 

124 return value 

125 

126 

127def _plan_str(data: Mapping[str, Any], key: str, where: str) -> str: 

128 value = _plan_field(data, key, where) 

129 if not isinstance(value, str): 

130 raise SwarmPlanError(f"{where}.{key}: expected a string, found {value!r}") 

131 return value 

132 

133 

134def _plan_bool(data: Mapping[str, Any], key: str, where: str) -> bool: 

135 value = _plan_field(data, key, where) 

136 if not isinstance(value, bool): 

137 raise SwarmPlanError(f"{where}.{key}: expected true or false, found {value!r}") 

138 return value 

139 

140 

141def _plan_list(data: Mapping[str, Any], key: str, where: str) -> list[Any]: 

142 value = _plan_field(data, key, where) 

143 if not isinstance(value, list): 

144 raise SwarmPlanError(f"{where}.{key}: expected a list, found {type(value).__name__}") 

145 return value 

146 

147 

148def _plan_strs(data: Mapping[str, Any], key: str, where: str) -> tuple[str, ...]: 

149 items = _plan_list(data, key, where) 

150 if not all(isinstance(item, str) for item in items): 

151 raise SwarmPlanError(f"{where}.{key}: expected a list of strings") 

152 return tuple(items) 

153 

154 

155def _plan_ints(data: Mapping[str, Any], key: str, where: str) -> tuple[int, ...]: 

156 items = _plan_list(data, key, where) 

157 if not all(isinstance(item, int) and not isinstance(item, bool) for item in items): 

158 raise SwarmPlanError(f"{where}.{key}: expected a list of integers") 

159 return tuple(items) 

160 

161 

162def _plan_issue_key(key: Any, where: str) -> int: 

163 if not isinstance(key, str) or not _ISSUE_KEY_RE.match(key): 

164 raise SwarmPlanError(f"{where}: key {key!r} is not an issue number") 

165 return int(key) 

166 

167 

168@dataclass(frozen=True) 

169class IssueScope: 

170 """Normalized predicted blast radius and role for a single backlog issue.""" 

171 

172 issue: int 

173 title: str = "" 

174 body: str = "" 

175 labels: tuple[str, ...] = () 

176 declared_files: tuple[str, ...] = () 

177 predicted_files: tuple[str, ...] = () 

178 role: str = "core" 

179 #: Where ``predicted_files`` came from — one of :data:`SCOPE_SOURCES` — so a plan 

180 #: that serialises everything can be read back to the issue that declared nothing. 

181 scope_source: str = "default" 

182 

183 def to_dict(self) -> dict[str, Any]: 

184 return { 

185 "issue": self.issue, 

186 "title": self.title, 

187 "role": self.role, 

188 "labels": list(self.labels), 

189 "declared_files": list(self.declared_files), 

190 "predicted_files": list(self.predicted_files), 

191 "scope_source": self.scope_source, 

192 } 

193 

194 @classmethod 

195 def from_dict(cls, data: Any, where: str = "issue scope") -> IssueScope: 

196 """The inverse of :meth:`to_dict` (#1275); ``body`` is not serialised, so it is empty. 

197 

198 Raises :class:`SwarmPlanError` for a record of the wrong shape. 

199 """ 

200 obj = _plan_object(data, where) 

201 source = _plan_str(obj, "scope_source", where) 

202 if source not in SCOPE_SOURCES: 

203 raise SwarmPlanError( 

204 f"{where}.scope_source: {source!r} is not one of {', '.join(SCOPE_SOURCES)}" 

205 ) 

206 return cls( 

207 issue=_plan_int(obj, "issue", where), 

208 title=_plan_str(obj, "title", where), 

209 labels=_plan_strs(obj, "labels", where), 

210 declared_files=_plan_strs(obj, "declared_files", where), 

211 predicted_files=_plan_strs(obj, "predicted_files", where), 

212 role=_plan_str(obj, "role", where), 

213 scope_source=source, 

214 ) 

215 

216 

217@dataclass(frozen=True) 

218class Difficulty: 

219 """How much work a cluster is, and the evidence that says so (#1017). 

220 

221 Deliberately not the risk tier. A tier answers *how dangerous is this change* and is 

222 read off the files it touches; a band answers *how much work is it*, which is what 

223 decides whether the strong implementer is worth spending on it. A one-line fix to a 

224 tier-3 glob is dangerous and trivial; a twelve-file docs migration is safe and long. 

225 Keeping them apart is what lets ``knobs.team`` staff on one and gate on the other. 

226 """ 

227 

228 score: int 

229 band: str 

230 tier: int 

231 file_count: int 

232 dependency_depth: int 

233 #: ``(name, points)`` for every input that contributed, so a surprising band can be 

234 #: read back rather than guessed at. 

235 signals: tuple[tuple[str, int], ...] = () 

236 

237 def to_dict(self) -> dict[str, Any]: 

238 return { 

239 "score": self.score, 

240 "band": self.band, 

241 "tier": self.tier, 

242 "file_count": self.file_count, 

243 "dependency_depth": self.dependency_depth, 

244 "signals": [{"name": name, "points": points} for name, points in self.signals], 

245 } 

246 

247 @classmethod 

248 def from_dict(cls, data: Any, where: str = "difficulty") -> Difficulty: 

249 """The inverse of :meth:`to_dict` (#1275); raises :class:`SwarmPlanError`.""" 

250 obj = _plan_object(data, where) 

251 signals = [] 

252 for n, raw in enumerate(_plan_list(obj, "signals", where)): 

253 signal = _plan_object(raw, f"{where}.signals[{n}]") 

254 signals.append( 

255 ( 

256 _plan_str(signal, "name", f"{where}.signals[{n}]"), 

257 _plan_int(signal, "points", f"{where}.signals[{n}]"), 

258 ) 

259 ) 

260 return cls( 

261 score=_plan_int(obj, "score", where), 

262 band=_plan_str(obj, "band", where), 

263 tier=_plan_int(obj, "tier", where), 

264 file_count=_plan_int(obj, "file_count", where), 

265 dependency_depth=_plan_int(obj, "dependency_depth", where), 

266 signals=tuple(signals), 

267 ) 

268 

269 

270@dataclass(frozen=True) 

271class SwarmCluster: 

272 """A unit of work within a Swarm Wave consisting of one or more related issues.""" 

273 

274 cluster_id: str 

275 issues: tuple[int, ...] 

276 role: str 

277 combined_scope: tuple[str, ...] 

278 depends_on_issues: tuple[int, ...] = () 

279 #: How much work this cluster is (:func:`score_difficulty`). 

280 difficulty: Difficulty | None = None 

281 #: Who runs it (:func:`keel.team.resolve_assignment`) — lead, implementer, reviewers. 

282 assignment: dict[str, Any] | None = None 

283 

284 def to_dict(self) -> dict[str, Any]: 

285 return { 

286 "cluster_id": self.cluster_id, 

287 "issues": list(self.issues), 

288 "role": self.role, 

289 "combined_scope": list(self.combined_scope), 

290 "depends_on_issues": list(self.depends_on_issues), 

291 "difficulty": None if self.difficulty is None else self.difficulty.to_dict(), 

292 "assignment": self.assignment, 

293 } 

294 

295 @classmethod 

296 def from_dict(cls, data: Any, where: str = "cluster") -> SwarmCluster: 

297 """The inverse of :meth:`to_dict` (#1275); raises :class:`SwarmPlanError`.""" 

298 obj = _plan_object(data, where) 

299 raw_difficulty = _plan_field(obj, "difficulty", where) 

300 assignment = _plan_field(obj, "assignment", where) 

301 if assignment is not None and not isinstance(assignment, dict): 

302 raise SwarmPlanError(f"{where}.assignment: expected an object or null") 

303 return cls( 

304 cluster_id=_plan_str(obj, "cluster_id", where), 

305 issues=_plan_ints(obj, "issues", where), 

306 role=_plan_str(obj, "role", where), 

307 combined_scope=_plan_strs(obj, "combined_scope", where), 

308 depends_on_issues=_plan_ints(obj, "depends_on_issues", where), 

309 difficulty=( 

310 None 

311 if raw_difficulty is None 

312 else Difficulty.from_dict(raw_difficulty, f"{where}.difficulty") 

313 ), 

314 assignment=assignment, 

315 ) 

316 

317 

318@dataclass(frozen=True) 

319class SwarmWave: 

320 """A parallel execution wave containing mutually independent or sequenced clusters.""" 

321 

322 wave_index: int 

323 mode: str # "orthogonal_parallel" or "sequential_dependent" 

324 eligible_direct_landing: bool 

325 clusters: tuple[SwarmCluster, ...] 

326 

327 def to_dict(self) -> dict[str, Any]: 

328 return { 

329 "wave_index": self.wave_index, 

330 "mode": self.mode, 

331 "eligible_direct_landing": self.eligible_direct_landing, 

332 "clusters": [c.to_dict() for c in self.clusters], 

333 } 

334 

335 @classmethod 

336 def from_dict(cls, data: Any, where: str = "wave") -> SwarmWave: 

337 """The inverse of :meth:`to_dict` (#1275); raises :class:`SwarmPlanError`.""" 

338 obj = _plan_object(data, where) 

339 mode = _plan_str(obj, "mode", where) 

340 if mode not in WAVE_MODES: 

341 raise SwarmPlanError(f"{where}.mode: {mode!r} is not one of {', '.join(WAVE_MODES)}") 

342 return cls( 

343 wave_index=_plan_int(obj, "wave_index", where), 

344 mode=mode, 

345 eligible_direct_landing=_plan_bool(obj, "eligible_direct_landing", where), 

346 clusters=tuple( 

347 SwarmCluster.from_dict(raw, f"{where}.clusters[{n}]") 

348 for n, raw in enumerate(_plan_list(obj, "clusters", where)) 

349 ), 

350 ) 

351 

352 

353@dataclass(frozen=True) 

354class SwarmPlan: 

355 """Complete deterministic partition and execution plan for a swarm run.""" 

356 

357 swarm_id: str 

358 total_issues: int 

359 waves: tuple[SwarmWave, ...] 

360 conflict_map: dict[int, tuple[int, ...]] = field(default_factory=dict) 

361 issue_scopes: dict[int, IssueScope] = field(default_factory=dict) 

362 

363 def to_dict(self) -> dict[str, Any]: 

364 return { 

365 "swarm_id": self.swarm_id, 

366 "total_issues": self.total_issues, 

367 "waves": [w.to_dict() for w in self.waves], 

368 "conflict_map": {str(k): list(v) for k, v in self.conflict_map.items()}, 

369 "issue_scopes": {str(k): v.to_dict() for k, v in self.issue_scopes.items()}, 

370 } 

371 

372 @classmethod 

373 def from_dict(cls, data: Any, where: str = "plan") -> SwarmPlan: 

374 """The inverse of :meth:`to_dict`: ``from_dict(p.to_dict()).to_dict() == p.to_dict()``. 

375 

376 Strict (#1275). Every field :meth:`to_dict` writes must be present with its type, 

377 an issue-number key must be written exactly as ``to_dict`` writes one, and an 

378 issue scope must sit under its own issue's key; anything else raises 

379 :class:`SwarmPlanError` naming where. A plan is what ``swarm-land`` lands, so a 

380 field read by guesswork would land a wave the run never planned. 

381 """ 

382 obj = _plan_object(data, where) 

383 raw_conflicts = _plan_object( 

384 _plan_field(obj, "conflict_map", where), f"{where}.conflict_map" 

385 ) 

386 conflict_map = { 

387 _plan_issue_key(key, f"{where}.conflict_map"): _plan_ints( 

388 raw_conflicts, key, f"{where}.conflict_map" 

389 ) 

390 for key in raw_conflicts 

391 } 

392 raw_scopes = _plan_object(_plan_field(obj, "issue_scopes", where), f"{where}.issue_scopes") 

393 issue_scopes: dict[int, IssueScope] = {} 

394 for key, raw in raw_scopes.items(): 

395 issue = _plan_issue_key(key, f"{where}.issue_scopes") 

396 scope = IssueScope.from_dict(raw, f"{where}.issue_scopes[{key!r}]") 

397 if scope.issue != issue: 

398 raise SwarmPlanError( 

399 f"{where}.issue_scopes[{key!r}]: holds the scope of issue #{scope.issue}" 

400 ) 

401 issue_scopes[issue] = scope 

402 return cls( 

403 swarm_id=_plan_str(obj, "swarm_id", where), 

404 total_issues=_plan_int(obj, "total_issues", where), 

405 waves=tuple( 

406 SwarmWave.from_dict(raw, f"{where}.waves[{n}]") 

407 for n, raw in enumerate(_plan_list(obj, "waves", where)) 

408 ), 

409 conflict_map=conflict_map, 

410 issue_scopes=issue_scopes, 

411 ) 

412 

413 

414@dataclass(frozen=True) 

415class SwarmWorkerStatus: 

416 """Live state for an individual worker operating on a swarm cluster.""" 

417 

418 cluster_id: str 

419 issue: int 

420 role: str 

421 agent: str = "claude" 

422 model: str = "default" 

423 step: str = "s0" 

424 status: str = "queued" # queued, running, passed, failed, merged, held 

425 updated_at: str = "" 

426 details: str = "" 

427 #: The team lead this worker reports through — the CTO/lead/worker hierarchy is only 

428 #: legible on the board if a worker record says which lead owns it (#1017). 

429 lead: str = "" 

430 #: The difficulty band the worker was staffed from, so a status board shows *why* 

431 #: this cluster drew this provider. 

432 difficulty: str = "" 

433 #: The consent scopes a live worker was handed by its parent (#1400) — exactly the 

434 #: parent's delegated scopes, or none. Empty for a dry run, which mutates nothing. 

435 scopes: tuple[str, ...] = () 

436 #: The pull request a live worker opened for its cluster (#1287): the one 

437 #: ``swarm-land`` merges. ``None`` until a worker opens one, and for a dry run. 

438 pull_request: int | None = None 

439 #: Where a live worker's worktree is still on disk when the worker ended — kept for 

440 #: inspection after a failure, or left by a removal that failed (#1278). Empty when 

441 #: the worktree was removed or never created. 

442 worktree: str = "" 

443 #: Whether a live worker pushed its branch; ``keel swarm-status --clean`` keeps a 

444 #: pushed branch (#1278). 

445 pushed: bool = False 

446 #: The wave the worker's cluster runs in, so the board can group by it (#1280). ``0`` 

447 #: in a record written before the field existed. 

448 wave: int = 0 

449 #: How far a live worker has got: one of :data:`keel.swarm_worker.STAGES`, recorded as 

450 #: the worker enters it, and once it ends the stage it stopped at — ``done`` when it 

451 #: opened its pull request (#1280). Empty until a live worker starts, and for a dry 

452 #: run's worker, whose child ``keel ship`` reports no stages. 

453 stage: str = "" 

454 #: When the worker started and when it ended, ISO 8601 (#1280). Empty until it does. 

455 started_at: str = "" 

456 finished_at: str = "" 

457 #: What the last live ``keel swarm-review`` did with this cluster's pull request (#1440): 

458 #: :func:`keel.swarm_review.review_record` — its status, the head it reviewed, each 

459 #: seat's outcome, when, and the run id. ``None`` until a live review runs, and in a 

460 #: record from before #1440. A report, never an input: ``swarm-land`` lands only what 

461 #: ``keel merge``'s evidence gate accepts, and does not read this. 

462 review: dict[str, Any] | None = None 

463 

464 def elapsed_s(self, now: str = "") -> int | None: 

465 """Whole seconds the worker has run: to ``finished_at``, else — still running — 

466 to ``now``. ``None`` when it has not started, or a time does not parse (#1280). 

467 

468 ``now`` is the caller's: this module reads no clock for a derived value. 

469 """ 

470 start = _parse_iso(self.started_at) 

471 end = _parse_iso(self.finished_at or now) 

472 if start is None or end is None: 

473 return None 

474 return max(0, int((end - start).total_seconds())) 

475 

476 def to_dict(self) -> dict[str, Any]: 

477 return { 

478 "cluster_id": self.cluster_id, 

479 "issue": self.issue, 

480 "role": self.role, 

481 "agent": self.agent, 

482 "model": self.model, 

483 "step": self.step, 

484 "status": self.status, 

485 "updated_at": self.updated_at, 

486 "details": self.details, 

487 "lead": self.lead, 

488 "difficulty": self.difficulty, 

489 "scopes": list(self.scopes), 

490 "pull_request": self.pull_request, 

491 "worktree": self.worktree, 

492 "pushed": self.pushed, 

493 "wave": self.wave, 

494 "stage": self.stage, 

495 "started_at": self.started_at, 

496 "finished_at": self.finished_at, 

497 "review": self.review, 

498 } 

499 

500 

501def _parse_iso(value: str) -> datetime.datetime | None: 

502 """``value`` as an aware datetime (a naive one is read as UTC), or ``None``.""" 

503 try: 

504 parsed = datetime.datetime.fromisoformat(value) 

505 except ValueError: 

506 return None 

507 return parsed if parsed.tzinfo is not None else parsed.replace(tzinfo=datetime.UTC) 

508 

509 

510@dataclass(frozen=True) 

511class SwarmRunState: 

512 """State tracking for a live or completed swarm execution.""" 

513 

514 swarm_id: str 

515 total_workers: int 

516 active_wave: int = 1 

517 workers: tuple[SwarmWorkerStatus, ...] = () 

518 started_at: str = "" 

519 completed_at: str | None = None 

520 #: The operator consent a live run's parent delegated to its workers (#1400): who 

521 #: consented, which scopes, which clusters, when — 

522 #: :meth:`keel.swarm_worker.ConsentDelegation.to_dict`. ``None`` for a dry run. 

523 consent: dict[str, Any] | None = None 

524 

525 def to_dict(self) -> dict[str, Any]: 

526 return { 

527 "swarm_id": self.swarm_id, 

528 "total_workers": self.total_workers, 

529 "active_wave": self.active_wave, 

530 "workers": [w.to_dict() for w in self.workers], 

531 "started_at": self.started_at, 

532 "completed_at": self.completed_at, 

533 "consent": self.consent, 

534 } 

535 

536 

537#: How much of one child `keel ship`'s output a swarm keeps per cluster (#1280). The 

538#: whole stdout used to go into `SwarmRunResult.wave_results` — which `swarm-run --json` 

539#: re-emits — and a failing cluster's into the state file as its `details`, so a 

540#: 20-cluster run wrote a multi-megabyte blob. The tail is what is kept: a failing 

541#: child states why at the end. 

542CHILD_OUTPUT_TAIL_CHARS = 4096 

543 

544 

545def tail_child_output(text: str) -> str: 

546 """Return ``text`` bounded to its last ``CHILD_OUTPUT_TAIL_CHARS`` characters. 

547 

548 Output within the cap comes back unchanged. Longer output keeps its tail behind a 

549 one-line marker that says how many earlier characters were dropped, so a reader 

550 never mistakes the kept part for the whole — or tries to parse it as the child's 

551 JSON document, which it no longer is. 

552 """ 

553 if len(text) <= CHILD_OUTPUT_TAIL_CHARS: 

554 return text 

555 dropped = len(text) - CHILD_OUTPUT_TAIL_CHARS 

556 return ( 

557 f"[keel: {dropped} earlier chars of child output dropped; " 

558 f"last {CHILD_OUTPUT_TAIL_CHARS} kept]\n{text[-CHILD_OUTPUT_TAIL_CHARS:]}" 

559 ) 

560 

561 

562def worker_timeout_s(config: ProjectConfig) -> int: 

563 """Return the wall-clock seconds one swarm worker's child ``keel ship`` may run. 

564 

565 Derived from the two budgets the project already sets for what the child runs 

566 (#1279): ``knobs.gate_timeout_s`` for its command gates and ``knobs.jury_timeout_s`` 

567 for the ``jury`` builtin. The child runs the gate suite in a dry run too, so a 

568 fixed 300 s killed a suite the project itself allows ten minutes. The sum covers 

569 one gate at its full budget plus a panel at its full budget; a suite whose gates 

570 add up to more than that is raised per run with ``swarm-run --worker-timeout``. 

571 """ 

572 return config.knobs.gate_timeout_s + config.knobs.jury_timeout_s 

573 

574 

575@dataclass(frozen=True) 

576class SwarmRunResult: 

577 """Outcome summary for a complete or partial swarm execution.""" 

578 

579 swarm_id: str 

580 status: str # "success", "partial_failure", "failed" 

581 total_workers: int 

582 passed_count: int 

583 failed_count: int 

584 dry_run: bool 

585 wave_results: tuple[dict[str, Any], ...] = () 

586 #: The consent delegation a live run's workers ran under (#1400); ``None`` when dry. 

587 consent: dict[str, Any] | None = None 

588 #: What the operator must look at by hand (#1278): a worktree or branch keel could not 

589 #: remove, a ``git worktree prune`` that failed, a worktree kept for inspection. 

590 warnings: tuple[str, ...] = () 

591 

592 def to_dict(self) -> dict[str, Any]: 

593 return { 

594 "swarm_id": self.swarm_id, 

595 "status": self.status, 

596 "total_workers": self.total_workers, 

597 "passed_count": self.passed_count, 

598 "failed_count": self.failed_count, 

599 "dry_run": self.dry_run, 

600 "wave_results": list(self.wave_results), 

601 "consent": self.consent, 

602 "warnings": list(self.warnings), 

603 } 

604 

605 

606@dataclass(frozen=True) 

607class LandingDecision: 

608 """Evaluation of whether a wave can directly land or requires sequential funneling.""" 

609 

610 mode: str # "direct_batch", "sequential_funnel" or "refused" 

611 eligible: bool 

612 cluster_ids: tuple[str, ...] 

613 reason: str = "orthogonal_diffs" 

614 

615 def to_dict(self) -> dict[str, Any]: 

616 return { 

617 "mode": self.mode, 

618 "eligible": self.eligible, 

619 "cluster_ids": list(self.cluster_ids), 

620 "reason": self.reason, 

621 } 

622 

623 

624@dataclass(frozen=True) 

625class ClusterClosure: 

626 """How a landed cluster's issues were closed after its merge (#1422). 

627 

628 ``/keel:ship`` closes the loop at s11–s12: the closure comment on the pull request and on 

629 the issue, then the issue closed. ``swarm-land`` does the same for every cluster it 

630 merges, and this is its report. ``closure_posted`` is one ``(kind, number, action)`` 

631 triple per comment — ``kind`` is ``pr`` or ``issue``, ``action`` is ``posted``, or 

632 ``edited`` when a closure comment of this run was already there. ``closed_issues`` were 

633 closed by this landing; ``already_closed`` were found closed and left alone. 

634 ``warnings`` say what did not happen and why — the merge is never undone for any of it. 

635 

636 A dry run (``dry_run``) posts and closes nothing: the record says what the live landing 

637 would close, with every list empty. 

638 """ 

639 

640 cluster_id: str 

641 pull_request: int 

642 issues: tuple[int, ...] 

643 dry_run: bool = False 

644 closure_posted: tuple[tuple[str, int, str], ...] = () 

645 closed_issues: tuple[int, ...] = () 

646 already_closed: tuple[int, ...] = () 

647 warnings: tuple[str, ...] = () 

648 

649 def to_dict(self) -> dict[str, Any]: 

650 return { 

651 "pull_request": self.pull_request, 

652 "issues": list(self.issues), 

653 "dry_run": self.dry_run, 

654 "closure_posted": [ 

655 {"kind": kind, "number": number, "action": action} 

656 for kind, number, action in self.closure_posted 

657 ], 

658 "closed_issues": list(self.closed_issues), 

659 "already_closed": list(self.already_closed), 

660 "warnings": list(self.warnings), 

661 } 

662 

663 

664def closure_target(kind: str, number: int) -> str: 

665 return f"PR #{number}" if kind == "pr" else f"issue #{number}" 

666 

667 

668def render_cluster_closure(closure: ClusterClosure) -> str: 

669 """One line: what a landed cluster's closure did, or — in a dry run — would do.""" 

670 issues = ", ".join(f"#{n}" for n in closure.issues) or "none" 

671 if closure.dry_run: 

672 return ( 

673 f"{closure.cluster_id}: would post the closure comment on PR #{closure.pull_request} " 

674 f"and each issue, then close {issues}" 

675 ) 

676 posted = ", ".join( 

677 f"{closure_target(kind, number)} {action}" 

678 for kind, number, action in closure.closure_posted 

679 ) 

680 parts = [f"closure {posted}" if posted else "no closure comment posted"] 

681 closed = ", ".join(f"#{n}" for n in closure.closed_issues) 

682 parts.append(f"closed {closed}" if closed else "closed none") 

683 if closure.already_closed: 

684 parts.append(f"already closed {', '.join(f'#{n}' for n in closure.already_closed)}") 

685 return f"{closure.cluster_id}: {'; '.join(parts)}" 

686 

687 

688@dataclass(frozen=True) 

689class SwarmLandingResult: 

690 """Outcome report for landing a swarm wave. 

691 

692 Landing merges each cluster's pull request through ``keel merge`` (#1287), so every 

693 cluster ends in one of three lists: landed (merged; in a dry run, would merge), 

694 held (``keel merge`` or the pull-request lookup refused it, with the reason) or 

695 failed (the merge call itself failed). 

696 """ 

697 

698 swarm_id: str 

699 wave_index: int 

700 mode: str 

701 landed_clusters: tuple[str, ...] 

702 failed_clusters: tuple[str, ...] 

703 status: str # "success", "partial_failure", "failed" 

704 #: Clusters that did not land and why — (cluster_id, reason) pairs: no open pull 

705 #: request, or one ``keel merge`` refused (window, merge state, CI, evidence, 

706 #: gates-pass, lock). Held is not failed: nothing was attempted that could fail. 

707 held_clusters: tuple[tuple[str, str], ...] = () 

708 #: Why a landing did not start at all: a ``sequential_dependent`` wave (dry run or 

709 #: live, ``mode`` is then ``"refused"``, #1276). Empty when it started. 

710 refused: str = "" 

711 #: What the operator must still look at by hand: a merge ``keel merge`` reported as 

712 #: drifted, or a checkout that could not be returned to where it started. 

713 warnings: tuple[str, ...] = () 

714 #: The pull request each cluster was landed through — (cluster_id, number) pairs. 

715 pull_requests: tuple[tuple[str, int], ...] = () 

716 #: How each landed cluster's issues were closed (#1422), in landing order. Only a 

717 #: landed cluster has one: a held or failed cluster posts and closes nothing. 

718 closures: tuple[ClusterClosure, ...] = () 

719 

720 def to_dict(self) -> dict[str, Any]: 

721 return { 

722 "swarm_id": self.swarm_id, 

723 "wave_index": self.wave_index, 

724 "mode": self.mode, 

725 "landed_clusters": list(self.landed_clusters), 

726 "failed_clusters": list(self.failed_clusters), 

727 "held_clusters": [list(pair) for pair in self.held_clusters], 

728 "pull_requests": {cluster: number for cluster, number in self.pull_requests}, 

729 "closures": {closure.cluster_id: closure.to_dict() for closure in self.closures}, 

730 "status": self.status, 

731 "refused": self.refused, 

732 "warnings": list(self.warnings), 

733 } 

734 

735 

736#: Punctuation prose puts around a path. A dot is not in it: a leading one is part of 

737#: `.github/…` or `.keel/…` (stripping it made those `github/…`, #1279), and a trailing 

738#: one is handled by :func:`_rstrip_path_punctuation`. 

739_PATH_LEAD = "`'\" \t\r\n,;:" 

740 

741 

742def _rstrip_path_punctuation(p: str) -> str: 

743 """Trailing prose punctuation, dots included — except a final `.` or `..` segment, 

744 which is a path step (`src/a/..` is `src`), not the end of a sentence.""" 

745 while True: 

746 before = p 

747 p = p.rstrip(_PATH_LEAD) 

748 # `\\` counts as a separator here: this runs before it becomes `/`. 

749 steps = p.replace("\\", "/") 

750 if p.endswith(".") and not ("/" in steps and steps.rsplit("/", 1)[-1] in (".", "..")): 

751 p = p[:-1] 

752 if p == before: 

753 return p 

754 

755 

756def _strip_path_punctuation(p: str) -> str: 

757 """Prose punctuation off both ends. A bracket goes only when it is unbalanced — the 

758 leftover of prose like `(see src/a.py)` — and balanced ones are the path's own: 

759 `docs/(draft)/` and `(docs/draft)/` keep theirs. Removing a pair that *looks* like 

760 it wraps the path cannot be both idempotent and consistent (`(docs/draft)/` wraps 

761 only once `normpath` drops the `/`), and the path-extraction patterns never 

762 capture brackets, so only a `--declared-file` value, which is literal, has any.""" 

763 while True: 

764 before = p 

765 p = _rstrip_path_punctuation(p.lstrip(_PATH_LEAD)) 

766 if p.startswith("(") and p.count("(") > p.count(")"): 

767 p = p[1:] 

768 elif p.endswith(")") and p.count(")") > p.count("("): 

769 p = p[:-1] 

770 if p == before: 

771 return p 

772 

773 

774def _normalize_path(p: str) -> str: 

775 """One canonical spelling of a path, the same however often it is applied: the 

776 plan normalises twice, and a second pass used to eat the `)` the first exposed 

777 (`docs/(draft)/` → `docs/(draft)` → `docs/(draft`, #1279). The whole pass is 

778 repeated to a fixed point, since `normpath` can itself expose new end 

779 punctuation (`(docs/draft)/` → `(docs/draft)`). It ends: after the first round 

780 no `\\` is left, and every later change only shortens the string.""" 

781 while True: 

782 cleaned = _strip_path_punctuation(p) 

783 cleaned = cleaned.replace("\\", "/").removeprefix("./").removeprefix("/") 

784 cleaned = posixpath.normpath(cleaned) if cleaned else "" 

785 if cleaned == p: 

786 return cleaned 

787 p = cleaned 

788 

789 

790def extract_predicted_paths(text: str) -> list[str]: 

791 """Extract candidate file and directory paths from issue title or markdown body.""" 

792 found: set[str] = set() 

793 for match in _PATH_BACKTICK_RE.finditer(text): 

794 found.add(_normalize_path(match.group(1))) 

795 for match in _PATH_GENERAL_RE.finditer(text): 

796 found.add(_normalize_path(match.group(1))) 

797 found.discard("") 

798 return sorted(found) 

799 

800 

801def _scope_globs(text: str) -> list[str]: 

802 """The globs in one declaration item: comma- or space-separated, prose skipped.""" 

803 globs: list[str] = [] 

804 for token in re.split(r"[\s,]+", text): 

805 if _GLOB_MARKS.isdisjoint(token): 

806 continue 

807 glob = _normalize_path(token) 

808 if glob and glob != ".": 

809 globs.append(glob) 

810 return globs 

811 

812 

813def parse_scope_declaration(body: str) -> tuple[str, ...]: 

814 """The path globs an issue body declares as its scope, in order, de-duplicated. 

815 

816 Two spellings, both line-based, both ignored inside a fenced code block: 

817 

818 * a line ``Scope: src/keel/swarm*.py, docs/keel/swarm.md``; 

819 * a heading ``## Scope`` (any level) followed by a bullet list, one or more globs 

820 per bullet, ending at the first line that is neither a bullet nor blank. 

821 

822 Globs are separated by commas or whitespace and may be backticked. A token with no 

823 ``/``, ``.`` or wildcard is prose and is skipped. Several declarations add up. An 

824 empty tuple means the body declares nothing. 

825 """ 

826 found: list[str] = [] 

827 in_fence = False 

828 in_heading = False 

829 for line in body.splitlines(): 

830 if _FENCE_RE.match(line): 

831 in_fence = not in_fence 

832 in_heading = False 

833 continue 

834 if in_fence: 

835 continue 

836 if in_heading: 

837 bullet = _SCOPE_BULLET_RE.match(line) 

838 if bullet: 

839 found.extend(_scope_globs(bullet.group(1))) 

840 continue 

841 if not line.strip(): 

842 continue 

843 in_heading = False 

844 if _SCOPE_HEADING_RE.match(line): 

845 in_heading = True 

846 continue 

847 declared = _SCOPE_LINE_RE.match(line) 

848 if declared: 

849 found.extend(_scope_globs(declared.group(1))) 

850 return tuple(dict.fromkeys(found)) 

851 

852 

853def scope_areas(config: ProjectConfig | None) -> dict[str, tuple[str, ...]]: 

854 """``policy_pack.scan.areas`` — area name to path globs — or ``{}`` without one. 

855 

856 The mapping an ``area:<name>`` label is read through. It is the project's own 

857 grouping of its tree (the scan commands use it too); keel adds no second one. 

858 """ 

859 pack = config.policy_pack if config is not None else {} 

860 scan = pack.get("scan") 

861 areas = scan.get("areas") if isinstance(scan, dict) else None 

862 if not isinstance(areas, dict): 

863 return {} 

864 return { 

865 str(name): tuple(str(glob) for glob in globs) 

866 for name, globs in areas.items() 

867 if isinstance(globs, list) 

868 } 

869 

870 

871def area_label_scope(labels: Sequence[str], areas: Mapping[str, Sequence[str]]) -> tuple[str, ...]: 

872 """The globs the ``area:<name>`` labels map to; an area not in ``areas`` adds none.""" 

873 found: list[str] = [] 

874 for label in labels: 

875 if label.startswith("area:"): 

876 for glob in areas.get(label.removeprefix("area:"), ()): 

877 normalized = _normalize_path(glob) 

878 if normalized: 

879 found.append(normalized) 

880 return tuple(dict.fromkeys(found)) 

881 

882 

883def parse_issue_scope_override(text: str) -> tuple[int, tuple[str, ...]]: 

884 """``N=glob[,glob…]`` → ``(N, globs)``, for ``--issue-scope``. 

885 

886 Raises ``ValueError`` with the reason when ``N`` is not a positive issue number or 

887 no glob survives normalisation — an override that names nothing is a typo, not a 

888 request to plan the issue as :data:`SCOPE_ANY`. 

889 """ 

890 number, sep, rest = text.partition("=") 

891 number = number.strip().lstrip("#") 

892 if not sep: 

893 raise ValueError(f"expected N=glob[,glob…], got {text!r}") 

894 if not (number.isascii() and number.isdigit()) or int(number) < 1: 

895 raise ValueError(f"the issue number must be a positive integer, got {number!r}") 

896 globs = tuple(dict.fromkeys(g for g in (_normalize_path(p) for p in rest.split(",")) if g)) 

897 if not globs: 

898 raise ValueError(f"issue #{int(number)} needs at least one glob after '='") 

899 return int(number), globs 

900 

901 

902def issue_facts_from_json(payload: str) -> tuple[str, str, tuple[str, ...]] | None: 

903 """``(title, body, labels)`` from ``gh issue view --json title,body,labels``. 

904 

905 ``None`` when the payload is not a JSON object — the caller treats that exactly like 

906 an issue it could not read. A field of the wrong type reads as empty. 

907 """ 

908 try: 

909 data = json.loads(payload) 

910 except json.JSONDecodeError: 

911 return None 

912 if not isinstance(data, dict): 

913 return None 

914 title = data.get("title") if isinstance(data.get("title"), str) else "" 

915 body = data.get("body") if isinstance(data.get("body"), str) else "" 

916 raw_labels = data.get("labels") if isinstance(data.get("labels"), list) else [] 

917 labels = tuple( 

918 item["name"] 

919 for item in raw_labels 

920 if isinstance(item, dict) and isinstance(item.get("name"), str) 

921 ) 

922 return title, body, labels 

923 

924 

925def extract_issue_scope( 

926 issue: int, 

927 *, 

928 title: str = "", 

929 body: str = "", 

930 labels: list[str] | tuple[str, ...] | None = None, 

931 declared_files: list[str] | tuple[str, ...] | None = None, 

932 config: ProjectConfig | None = None, 

933 scope: Sequence[str] | None = None, 

934) -> IssueScope: 

935 """Extract and normalize the predicted files and role for an issue. 

936 

937 The declared scope is the first of these that names anything 

938 (:data:`SCOPE_SOURCES`): ``scope`` (the ``--issue-scope`` override), the body's 

939 ``Scope:`` declaration (:func:`parse_scope_declaration`), and the ``area:<name>`` 

940 labels mapped through ``policy_pack.scan.areas``. ``declared_files`` are added to it. 

941 Without one, the paths the title and body name and the role hints are kept, and an 

942 issue with no ``declared_files`` either also gets :data:`SCOPE_ANY`, so it conflicts 

943 with every other issue. 

944 """ 

945 norm_labels = tuple(sorted(set(labels or ()))) 

946 norm_declared = tuple( 

947 sorted(set(_normalize_path(f) for f in (declared_files or ()) if _normalize_path(f))) 

948 ) 

949 

950 candidates = ( 

951 ("override", tuple(g for g in (_normalize_path(p) for p in (scope or ())) if g)), 

952 ("issue-body", parse_scope_declaration(body)), 

953 ("area-label", area_label_scope(norm_labels, scope_areas(config))), 

954 ) 

955 source, declared_scope = next(((s, g) for s, g in candidates if g), ("", ())) 

956 

957 predicted = set(norm_declared) | set(declared_scope) 

958 if not declared_scope: 

959 predicted.update(extract_predicted_paths(f"{title}\n{body}")) 

960 

961 # Resolve role from labels or config 

962 resolved_role = "core" 

963 for label in norm_labels: 

964 if label.startswith("role:"): 

965 resolved_role = label.removeprefix("role:") 

966 break 

967 if label.startswith("area:"): 

968 resolved_role = label.removeprefix("area:") 

969 break 

970 

971 known_roles = agents.known_roles(config) if config else frozenset() 

972 if known_roles and resolved_role not in known_roles: 

973 resolved_role = "core" if "core" in known_roles else resolved_role 

974 

975 # Default directory hints based on role or labels if no specific files found 

976 if not predicted: 

977 if "visual" in resolved_role or any("visual" in lbl for lbl in norm_labels): 

978 predicted.add("keel-visual/*") 

979 elif "docs" in resolved_role or any("docs" in lbl for lbl in norm_labels): 

980 predicted.add("docs/*") 

981 elif "website" in resolved_role or any("website" in lbl for lbl in norm_labels): 

982 predicted.add("website/*") 

983 elif "cli" in resolved_role or any("cli" in lbl for lbl in norm_labels): 

984 predicted.add("src/keel/cli.py") 

985 

986 # A path the text mentions, or a role hint, is a guess about the issue, not a statement 

987 # of what it touches: an issue that names `a.py` and also edits `cli.py` is the 

988 # collision #1274 describes. So the guesses stay in the scope — the tier and the 

989 # difficulty read them — but unless something *declared* a scope, the issue also gets 

990 # SCOPE_ANY and is serialised rather than assumed disjoint on the strength of a mention. 

991 if not declared_scope: 

992 if norm_declared: 

993 source = "declared-file" 

994 else: 

995 predicted.add(SCOPE_ANY) 

996 source = "default" 

997 

998 return IssueScope( 

999 issue=issue, 

1000 title=title.strip(), 

1001 body=body.strip(), 

1002 labels=norm_labels, 

1003 declared_files=norm_declared, 

1004 predicted_files=tuple(sorted(predicted)), 

1005 role=resolved_role, 

1006 scope_source=source, 

1007 ) 

1008 

1009 

1010def _normalized_paths_intersect(a: str, b: str) -> bool: 

1011 """:func:`paths_intersect` for two paths already passed through ``_normalize_path``. 

1012 

1013 The matching rules live here and nowhere else. ``paths_intersect`` normalizes and 

1014 delegates; the swarm plan normalizes each scope once and calls this directly. 

1015 """ 

1016 if not a or not b: 

1017 return False 

1018 if a == b: 

1019 return True 

1020 if a == "*" or b == "*": 

1021 return True 

1022 # Directory prefix overlap 

1023 a_dir = a if a.endswith("/") else a + "/" 

1024 b_dir = b if b.endswith("/") else b + "/" 

1025 if b.startswith(a_dir) or a.startswith(b_dir): 

1026 return True 

1027 # Glob matching 

1028 if "*" in a and fnmatch.fnmatch(b, a): 

1029 return True 

1030 if "*" in b and fnmatch.fnmatch(a, b): 

1031 return True 

1032 return False 

1033 

1034 

1035def paths_intersect(path_a: str, path_b: str) -> bool: 

1036 """True if path_a and path_b refer to the same file or overlapping glob/directory.""" 

1037 return _normalized_paths_intersect(_normalize_path(path_a), _normalize_path(path_b)) 

1038 

1039 

1040def scopes_intersect(scope_a: IssueScope, scope_b: IssueScope) -> tuple[str, ...]: 

1041 """Return common overlapping paths between two issue scopes.""" 

1042 overlaps: set[str] = set() 

1043 for fa in scope_a.predicted_files: 

1044 for fb in scope_b.predicted_files: 

1045 if paths_intersect(fa, fb): 

1046 overlaps.add(fa if len(fa) <= len(fb) else fb) 

1047 return tuple(sorted(overlaps)) 

1048 

1049 

1050def _normalized_files(scope: IssueScope) -> tuple[str, ...]: 

1051 """The scope's predicted paths as the matcher sees them, normalized once.""" 

1052 return tuple(_normalize_path(f) for f in scope.predicted_files) 

1053 

1054 

1055def _normalized_scopes_conflict(files_a: tuple[str, ...], files_b: tuple[str, ...]) -> bool: 

1056 """:func:`scopes_have_conflict` over paths already normalized by ``_normalized_files``.""" 

1057 # Fast path: an identical path is a conflict under the ``a == b`` rule, and a set 

1058 # intersection finds one without the pairwise loop. The empty string is dropped 

1059 # because the matcher rejects it on either side — ``("",)`` vs ``("",)`` is not a 

1060 # conflict, and must not become one here. 

1061 common = set(files_a) & set(files_b) 

1062 common.discard("") 

1063 if common: 

1064 return True 

1065 for fa in files_a: 

1066 for fb in files_b: 

1067 if _normalized_paths_intersect(fa, fb): 

1068 return True 

1069 return False 

1070 

1071 

1072def scopes_have_conflict(scope_a: IssueScope, scope_b: IssueScope) -> bool: 

1073 """True if any predicted path of one scope overlaps any path of the other. 

1074 

1075 Equivalent to ``bool(scopes_intersect(scope_a, scope_b))`` — the same matcher and 

1076 the same normalization — but returns at the first overlap instead of collecting 

1077 and sorting them all. ``build_swarm_plan`` asks this O(N²) times and only needs 

1078 the boolean. 

1079 """ 

1080 return _normalized_scopes_conflict(_normalized_files(scope_a), _normalized_files(scope_b)) 

1081 

1082 

1083def _file_count_points(count: int) -> int: 

1084 """Points for a scope of ``count`` predicted files (the widest threshold it reaches).""" 

1085 for threshold, points in _FILE_COUNT_POINTS: 

1086 if count >= threshold: 

1087 return points 

1088 return 0 

1089 

1090 

1091def difficulty_band(score: int) -> str: 

1092 """The band a difficulty ``score`` falls in.""" 

1093 for ceiling, band in _BAND_CEILINGS: 

1094 if score <= ceiling: 

1095 return band 

1096 return team_policy.DIFFICULTY_BANDS[-1] 

1097 

1098 

1099def cluster_scopes( 

1100 cluster: SwarmCluster, scopes: Mapping[int, IssueScope] 

1101) -> tuple[IssueScope, ...]: 

1102 """The issue scopes a cluster covers, in the cluster's own order. 

1103 

1104 An issue with no scope in the map is skipped rather than faked: a scoreable cluster 

1105 is one whose issues the planner actually analysed, and inventing an empty scope would 

1106 quietly lower the band of a cluster whose largest issue went missing. 

1107 """ 

1108 return tuple(scopes[issue] for issue in cluster.issues if issue in scopes) 

1109 

1110 

1111def score_difficulty( 

1112 scopes: Sequence[IssueScope], 

1113 *, 

1114 tier3_globs: tuple[str, ...] = (), 

1115 docs_globs: tuple[str, ...] = (), 

1116 allowlist_globs: tuple[str, ...] = (), 

1117 dependency_depth: int = 0, 

1118) -> Difficulty: 

1119 """How much work one cluster of issues is — pure, and the same answer every run. 

1120 

1121 Takes the cluster's scopes rather than one issue's, because a cluster is the unit a 

1122 lead is handed and a cluster of three small issues is not three small pieces of work. 

1123 Every input is already deterministic: the risk tier from ``knobs.tier3_globs``, the 

1124 predicted blast radius the planner computed, labels a human set, and how much 

1125 already-scheduled work this cluster lands on top of. 

1126 """ 

1127 files = sorted({path for scope in scopes for path in scope.predicted_files}) 

1128 labels = sorted({label.lower() for scope in scopes for label in scope.labels}) 

1129 tier = classify.tier_for_files( 

1130 files, 

1131 tier3_globs=tier3_globs, 

1132 docs_globs=docs_globs, 

1133 allowlist_globs=allowlist_globs, 

1134 ) 

1135 # The observed depth is recorded and the *points* are what the cap bites: a cluster 

1136 # sitting on nine earlier issues and one sitting on three are not the same situation, 

1137 # and reporting both as `depends-on:3` threw away the difference at the only place a 

1138 # reader could have seen it. Capping the points still says "past a few dependencies 

1139 # it is the same problem, not a worse one". 

1140 depth = max(dependency_depth, 0) 

1141 depth_points = min(depth, MAX_DEPENDENCY_POINTS) 

1142 candidates: list[tuple[str, int]] = [ 

1143 (f"tier-{tier}", _TIER_POINTS.get(tier, 2)), 

1144 (f"files:{len(files)}", _file_count_points(len(files))), 

1145 *((label, _LABEL_POINTS[label]) for label in labels if label in _LABEL_POINTS), 

1146 (f"depends-on:{depth}", depth_points), 

1147 ] 

1148 # Only what actually moved the score is recorded: a signal worth zero points is not 

1149 # evidence, and a reader who has to filter them out is reading noise. 

1150 signals = tuple((name, points) for name, points in candidates if points) 

1151 score = sum(points for _name, points in signals) 

1152 return Difficulty( 

1153 score=score, 

1154 band=difficulty_band(score), 

1155 tier=tier, 

1156 file_count=len(files), 

1157 dependency_depth=depth, 

1158 signals=signals, 

1159 ) 

1160 

1161 

1162@dataclass(frozen=True) 

1163class AssignmentOverrides: 

1164 """Per-run staffing a batch runner passes down to every cluster it launches. 

1165 

1166 The command-line half of the team contract: ``--delegate``, ``--review-delegate``, 

1167 ``--effort``, ``--team`` and ``--reviewers`` as one value, so ``swarm-plan``, 

1168 ``swarm-run`` and a work block hand the same object to the same resolver instead of 

1169 each threading five arguments and dropping a different one. 

1170 """ 

1171 

1172 delegate: str | None = None 

1173 review_delegates: tuple[str, ...] = () 

1174 effort: str | None = None 

1175 team_profile: str | None = None 

1176 reviewers: int | None = None 

1177 host_agent: str = team_policy.HOST_DEFAULT 

1178 

1179 def to_dict(self) -> dict[str, Any]: 

1180 return { 

1181 "delegate": self.delegate, 

1182 "review_delegates": list(self.review_delegates), 

1183 "effort": self.effort, 

1184 "team": self.team_profile, 

1185 "reviewers": self.reviewers, 

1186 "host_agent": self.host_agent, 

1187 } 

1188 

1189 

1190def resolve_cluster_assignment( 

1191 cluster: SwarmCluster, 

1192 difficulty: Difficulty, 

1193 *, 

1194 config: ProjectConfig | None = None, 

1195 overrides: AssignmentOverrides | None = None, 

1196 jury_availability: Mapping[str, Any] | None = None, 

1197) -> dict[str, Any]: 

1198 """Who runs this cluster: :func:`keel.team.resolve_assignment`, per cluster. 

1199 

1200 The swarm does not get its own idea of who implements. It calls the resolver every 

1201 other command calls, with *this cluster's* role, risk tier and difficulty band — so a 

1202 cluster's lead can hand the assignment straight to ``keel ship`` and the child run 

1203 resolves the same team from the same config. 

1204 

1205 ``jury_availability`` is the s7 panel probe (#1066), measured once per plan by 

1206 :func:`keel.providerprobe.jury_availability_for_any_tier` and passed in rather than taken 

1207 here: this module is pure, and the measurement is the one machine-dependent input the 

1208 resolver has. Passing it is not optional in spirit — without it a swarm on a panel 

1209 project publishes ``review_panel: jury`` for a tier-3 cluster while the child ``keel 

1210 ship`` it launches on the same machine seats three host reviewers, which is exactly the 

1211 in-process disagreement #1066 closed one layer down. 

1212 """ 

1213 settings = overrides or AssignmentOverrides() 

1214 policy = config.knobs.team if config is not None else team_policy.TeamPolicy() 

1215 assignment = team_policy.resolve_assignment( 

1216 policy, 

1217 tier=difficulty.tier, 

1218 role=cluster.role, 

1219 default_count=ship.reviewer_count(difficulty.tier), 

1220 reviewer_override=settings.reviewers, 

1221 delegate=settings.delegate, 

1222 review_delegates=settings.review_delegates, 

1223 host_agent=settings.host_agent, 

1224 legacy=None if config is None else agents.legacy_team_seats(config), 

1225 difficulty=difficulty.band, 

1226 team_profile=settings.team_profile, 

1227 effort=settings.effort, 

1228 jury_availability=jury_availability, 

1229 ) 

1230 # A role keel will not hand to a child is an operator-visible fact, not a silent 

1231 # omission: `assignment.warnings` is where a lead is told to look before launching. 

1232 assignment["warnings"] = [*assignment["warnings"], *safe_role(assignment["role"])[1]] 

1233 return assignment 

1234 

1235 

1236#: Characters a role may contribute to a child's argv. A role is read off an issue label, 

1237#: which anyone with triage rights can write; `--role` hands it to a child `argparse`, 

1238#: where a value opening with `-` is parsed as a flag rather than as the option's value. 

1239#: `role:--live` on one issue would therefore break every child ship in the batch. 

1240_SAFE_ROLE = re.compile(r"\A[A-Za-z0-9][A-Za-z0-9._-]*\Z") 

1241 

1242 

1243def safe_role(role: str | None) -> tuple[str | None, tuple[str, ...]]: 

1244 """``(role, warnings)`` — the role if it can be an argv token, else ``None`` and why. 

1245 

1246 Dropped rather than escaped or quoted: there is no quoting that makes `--live` stop 

1247 being a flag to `argparse`, and a role keel cannot pass is a role the child resolves 

1248 from config on its own — a smaller loss than a batch that dies on argument parsing. 

1249 The warning is deterministic so two runs of the same plan report it identically. 

1250 """ 

1251 if role is None or not role: 

1252 return None, () 

1253 if _SAFE_ROLE.match(role): 

1254 return role, () 

1255 return None, ( 

1256 f"role {role!r} is not passed to the child ship: a role becomes a `--role` argv " 

1257 "token, and this one is not [A-Za-z0-9][A-Za-z0-9._-]*. The child resolves its " 

1258 "role from the issue's own labels instead", 

1259 ) 

1260 

1261 

1262def ship_handoff_args(assignment: dict[str, Any] | None) -> tuple[str, ...]: 

1263 """A resolved assignment as ``keel ship`` flags for the child run. 

1264 

1265 One place, because every batch runner has the same job — a swarm lead launching its 

1266 cluster's ships, a work block handing over the next issue — and a child that resolves 

1267 its own team from config alone would silently drop the per-run overrides the operator 

1268 passed to the parent. Only flags ``keel ship`` actually defines are emitted; a seat 

1269 that is a host subagent rather than a provider is left to the adapter, which is the 

1270 layer that knows how to spawn one. 

1271 

1272 The vocabulary is :data:`keel.workblock.DELEGATION_FLAGS` — the same five a work block 

1273 hands down — because the child is the same `keel ship` either way. The two differ only 

1274 in *what* they carry: a work block passes the operator's flags through unresolved, 

1275 while a lead has already resolved its cluster's bench and passes the seats that came 

1276 out. ``--effort`` and ``--team`` ride along so the child's own resolution reproduces 

1277 the parent's rather than re-deriving a different one from config alone. 

1278 """ 

1279 if assignment is None: 

1280 return () 

1281 args: list[str] = [] 

1282 implementer = assignment["implementer"] 

1283 if implementer["kind"] == "provider": 

1284 model = implementer["model"] 

1285 args += ["--delegate", implementer["provider"] + (f":{model}" if model else "")] 

1286 for reviewer in assignment["reviewers"]: 

1287 if reviewer["kind"] == "provider": 

1288 model = reviewer["model"] 

1289 args += ["--review-delegate", reviewer["provider"] + (f":{model}" if model else "")] 

1290 role, _warnings = safe_role(assignment["role"]) 

1291 if role: 

1292 args += ["--role", role] 

1293 if assignment["effort"]: 

1294 args += ["--effort", assignment["effort"]] 

1295 if assignment["team_profile"]: 

1296 args += ["--team", assignment["team_profile"]] 

1297 return tuple(args) 

1298 

1299 

1300def worker_seed(cluster: SwarmCluster, *, wave: int = 0, updated_at: str = "") -> SwarmWorkerStatus: 

1301 """The queued worker record for a cluster, staffed from the cluster's assignment. 

1302 

1303 Pure, and separate from the runtime that persists it, because *which provider is 

1304 running this cluster and which lead owns it* is a fact of the plan. The runtime used 

1305 to seed every worker with the ``claude``/``default`` field defaults, so a swarm that 

1306 had resolved a real team still reported the placeholder one on the board. ``wave`` is 

1307 the plan's too: the index of the wave the cluster sits in (#1280). 

1308 """ 

1309 assignment = cluster.assignment 

1310 implementer = None if assignment is None else assignment["implementer"] 

1311 model = None if implementer is None else implementer["model"] 

1312 return SwarmWorkerStatus( 

1313 cluster_id=cluster.cluster_id, 

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

1315 role=cluster.role, 

1316 agent="claude" if implementer is None else implementer["name"], 

1317 model=model or "default", 

1318 lead="" if assignment is None else assignment["lead"]["name"], 

1319 difficulty="" if cluster.difficulty is None else cluster.difficulty.band, 

1320 step="s0", 

1321 status="queued", 

1322 updated_at=updated_at, 

1323 wave=wave, 

1324 ) 

1325 

1326 

1327def build_swarm_plan( 

1328 issue_scopes: list[IssueScope] | tuple[IssueScope, ...], 

1329 *, 

1330 swarm_id: str | None = None, 

1331 config: ProjectConfig | None = None, 

1332 overrides: AssignmentOverrides | None = None, 

1333 jury_availability: Mapping[str, Any] | None = None, 

1334) -> SwarmPlan: 

1335 """Deterministically partition candidate issues into orthogonal Waves and Clusters. 

1336 

1337 Every cluster comes out scored (:func:`score_difficulty`) and staffed 

1338 (:func:`resolve_cluster_assignment`). Scoring and staffing happen *after* the 

1339 partition and never feed back into it, which is the property that makes the two 

1340 independently reviewable: re-staffing a backlog — a new ``team.by_difficulty`` row, a 

1341 ``--team`` profile, a ``--delegate`` — moves who runs a cluster and cannot move which 

1342 wave it lands in. 

1343 

1344 ``jury_availability`` rides through to every cluster's staffing (#1066); it is measured 

1345 once, by the caller, because it is a fact about the machine and not about a cluster. 

1346 """ 

1347 if not issue_scopes: 

1348 now_id = swarm_id or ( 

1349 "swarm-" + datetime.datetime.now(datetime.UTC).strftime("%Y%m%d-%H%M%S") 

1350 ) 

1351 return SwarmPlan( 

1352 swarm_id=now_id, 

1353 total_issues=0, 

1354 waves=(), 

1355 conflict_map={}, 

1356 issue_scopes={}, 

1357 ) 

1358 

1359 tier3_globs = () if config is None else config.knobs.tier3_globs 

1360 docs_globs = () if config is None else config.knobs.docs_gate_paths 

1361 allowlist_globs = () if config is None else config.knobs.docs_only_allowlist 

1362 

1363 # Sort issues for determinism 

1364 sorted_scopes = sorted(issue_scopes, key=lambda s: s.issue) 

1365 scope_by_id = {s.issue: s for s in sorted_scopes} 

1366 

1367 # Compute pairwise conflict graph 

1368 conflict_map: dict[int, list[int]] = {s.issue: [] for s in sorted_scopes} 

1369 # Normalize every scope once, outside the O(N²) pair loop: the loop used to pay 

1370 # for ``_normalize_path`` on both sides of every one of the |A|·|B| path pairs, 

1371 # for every pair of issues — about half the per-pair cost. Kept as a list aligned 

1372 # with ``sorted_scopes``, not a dict by issue number: two scopes carrying the same 

1373 # issue number are two scopes, and keying by number would silently compare one 

1374 # of them with the other's files. 

1375 normalized = [_normalized_files(s) for s in sorted_scopes] 

1376 for i, sa in enumerate(sorted_scopes): 

1377 for j in range(i + 1, len(sorted_scopes)): 

1378 sb = sorted_scopes[j] 

1379 if _normalized_scopes_conflict(normalized[i], normalized[j]): 

1380 conflict_map[sa.issue].append(sb.issue) 

1381 conflict_map[sb.issue].append(sa.issue) 

1382 

1383 frozen_conflicts: dict[int, tuple[int, ...]] = { 

1384 k: tuple(sorted(v)) for k, v in conflict_map.items() 

1385 } 

1386 

1387 # Greedy wave partitioning (independent set partitioning) 

1388 remaining_issues = [s.issue for s in sorted_scopes] 

1389 waves: list[SwarmWave] = [] 

1390 wave_idx = 1 

1391 

1392 assigned_prior_issues: set[int] = set() 

1393 

1394 while remaining_issues: 

1395 current_wave_issues: list[int] = [] 

1396 current_wave_scope_set: set[str] = set() 

1397 

1398 for issue_num in list(remaining_issues): 

1399 scope = scope_by_id[issue_num] 

1400 # Check if this issue conflicts with anything already placed in current wave 

1401 has_wave_conflict = False 

1402 for placed_num in current_wave_issues: 

1403 if issue_num in frozen_conflicts.get(placed_num, ()): 

1404 has_wave_conflict = True 

1405 break 

1406 

1407 if not has_wave_conflict: 

1408 current_wave_issues.append(issue_num) 

1409 current_wave_scope_set.update(scope.predicted_files) 

1410 remaining_issues.remove(issue_num) 

1411 

1412 # Build clusters for this wave 

1413 clusters: list[SwarmCluster] = [] 

1414 for issue_num in current_wave_issues: 

1415 sc = scope_by_id[issue_num] 

1416 # Find dependencies on previously assigned waves 

1417 deps = tuple( 

1418 sorted( 

1419 dep 

1420 for dep in frozen_conflicts.get(issue_num, ()) 

1421 if dep in assigned_prior_issues 

1422 ) 

1423 ) 

1424 cluster = SwarmCluster( 

1425 cluster_id=f"cluster-{wave_idx}-{issue_num}", 

1426 issues=(issue_num,), 

1427 role=sc.role, 

1428 combined_scope=sc.predicted_files, 

1429 depends_on_issues=deps, 

1430 ) 

1431 # Scored from the cluster's own issue list rather than from the issue this 

1432 # loop happens to be on. The partition emits one issue per cluster today, so 

1433 # the two are the same tuple — but the scorer's contract is "how much work is 

1434 # *this cluster*", and a caller that hand-rolled the single-issue case would 

1435 # be the thing to fix on the day clusters group. 

1436 difficulty = score_difficulty( 

1437 cluster_scopes(cluster, scope_by_id), 

1438 tier3_globs=tier3_globs, 

1439 docs_globs=docs_globs, 

1440 allowlist_globs=allowlist_globs, 

1441 dependency_depth=len(deps), 

1442 ) 

1443 cluster = replace(cluster, difficulty=difficulty) 

1444 clusters.append( 

1445 replace( 

1446 cluster, 

1447 assignment=resolve_cluster_assignment( 

1448 cluster, 

1449 difficulty, 

1450 config=config, 

1451 overrides=overrides, 

1452 jury_availability=jury_availability, 

1453 ), 

1454 ) 

1455 ) 

1456 

1457 assigned_prior_issues.update(current_wave_issues) 

1458 

1459 mode, eligible = wave_landing_mode(wave_idx, clusters) 

1460 waves.append( 

1461 SwarmWave( 

1462 wave_index=wave_idx, 

1463 mode=mode, 

1464 eligible_direct_landing=eligible, 

1465 clusters=tuple(clusters), 

1466 ) 

1467 ) 

1468 wave_idx += 1 

1469 

1470 now_id = swarm_id or ("swarm-" + datetime.datetime.now(datetime.UTC).strftime("%Y%m%d-%H%M%S")) 

1471 return SwarmPlan( 

1472 swarm_id=now_id, 

1473 total_issues=len(sorted_scopes), 

1474 waves=tuple(waves), 

1475 conflict_map=frozen_conflicts, 

1476 issue_scopes=scope_by_id, 

1477 ) 

1478 

1479 

1480def wave_landing_mode(wave_index: int, clusters: Sequence[SwarmCluster]) -> tuple[str, bool]: 

1481 """A wave's plan mode and whether it may land as a direct batch. 

1482 

1483 The first wave is ``orthogonal_parallel``: nothing lands before it. A later wave is 

1484 too only when none of its clusters depends on an issue from an earlier wave; 

1485 otherwise it is ``sequential_dependent`` and not eligible for the direct batch, 

1486 because the base it branched from moves when the wave it depends on lands. The 

1487 mode used to be computed as "the wave is not empty", which every wave is, so every 

1488 wave claimed direct landing (#1276). 

1489 """ 

1490 orthogonal = wave_index == 1 or not any(c.depends_on_issues for c in clusters) 

1491 return ("orthogonal_parallel" if orthogonal else "sequential_dependent", orthogonal) 

1492 

1493 

1494def seat_label(seat: dict[str, Any] | None) -> str: 

1495 """A resolved seat as ``provider:model@effort`` — the short form both renderers use.""" 

1496 if seat is None: 

1497 return "unassigned" 

1498 label = seat["name"] 

1499 if seat["model"]: 

1500 label += f":{seat['model']}" 

1501 if seat["effort"]: 

1502 label += f"@{seat['effort']}" 

1503 return label 

1504 

1505 

1506def cluster_staffing_lines(cluster: SwarmCluster) -> list[str]: 

1507 """The difficulty and team rows for a cluster, or nothing when it has neither. 

1508 

1509 Shared by both plan renderers so the tree and the table cannot describe the same 

1510 cluster differently — the tree existing to be read *instead of* the table is exactly 

1511 why they must agree. 

1512 """ 

1513 rows: list[str] = [] 

1514 if cluster.difficulty is not None: 

1515 d = cluster.difficulty 

1516 rows.append( 

1517 f"Difficulty: {d.band} (score {d.score}, tier {d.tier}, " 

1518 f"{d.file_count} file(s), depth {d.dependency_depth})" 

1519 ) 

1520 if cluster.assignment is not None: 

1521 assignment = cluster.assignment 

1522 reviewers = ", ".join(seat_label(seat) for seat in assignment["reviewers"]) 

1523 panel = reviewers or assignment["review_panel"] 

1524 rows.append( 

1525 f"Team: lead {seat_label(assignment['lead'])} → " 

1526 f"implementer {seat_label(assignment['implementer'])}, review {panel}" 

1527 ) 

1528 return rows 

1529 

1530 

1531def render_swarm_plan_text(plan: SwarmPlan) -> str: 

1532 """Render human-readable tabular text summary of the SwarmPlan.""" 

1533 lines: list[str] = [ 

1534 f"keel swarm plan — {plan.swarm_id}", 

1535 f" total issues : {plan.total_issues}", 

1536 f" total waves : {len(plan.waves)}", 

1537 "", 

1538 ] 

1539 for w in plan.waves: 

1540 # What swarm-land does with the wave (#1276): a dependent wave is refused. 

1541 landing = ( 

1542 "eligible for direct batch landing" 

1543 if w.eligible_direct_landing 

1544 else "dependent on an earlier wave — swarm-land refuses it; " 

1545 "land the earlier wave, then re-plan" 

1546 ) 

1547 lines.append(f"Wave {w.wave_index} [{w.mode}] — {landing}:") 

1548 for c in w.clusters: 

1549 issues_str = ", ".join(f"#{i}" for i in c.issues) 

1550 scope_prefix = c.combined_scope[:3] 

1551 scope_str = ", ".join(scope_prefix) + ("..." if len(c.combined_scope) > 3 else "") 

1552 dep_str = ( 

1553 f" (depends on: {', '.join(f'#{d}' for d in c.depends_on_issues)})" 

1554 if c.depends_on_issues 

1555 else "" 

1556 ) 

1557 lines.append( 

1558 f" • Cluster {c.cluster_id}: {issues_str} [{c.role}] → {scope_str}{dep_str}" 

1559 ) 

1560 for detail in cluster_staffing_lines(c): 

1561 lines.append(f" {detail}") 

1562 lines.append("") 

1563 return "\n".join(lines).strip() 

1564 

1565 

1566def render_swarm_plan_tree(plan: SwarmPlan) -> str: 

1567 """Render a visual ASCII/Unicode DAG dependency tree of the SwarmPlan.""" 

1568 if plan.total_issues == 0: 

1569 return f"keel swarm plan — {plan.swarm_id} (0 issues)" 

1570 

1571 direct_count = sum(1 for w in plan.waves if w.eligible_direct_landing) 

1572 hdr_a = f"│ 🐝 Keel Swarm Plan — {plan.swarm_id:<38} │" 

1573 hdr_b = ( 

1574 f"│ Issues: {plan.total_issues:<3} │ Waves: {len(plan.waves):<3} " 

1575 f"│ Direct Landing Waves: {direct_count:<2} │" 

1576 ) 

1577 lines: list[str] = [ 

1578 "╭" + "─" * 62 + "╮", 

1579 hdr_a, 

1580 hdr_b, 

1581 "╰" + "─" * 62 + "╯", 

1582 "", 

1583 ] 

1584 

1585 for w in plan.waves: 

1586 mode_icon = "⚡" if w.eligible_direct_landing else "⏳" 

1587 landing_label = ( 

1588 "Direct Batch Landing" 

1589 if w.eligible_direct_landing 

1590 else "Dependent — refused until re-planned" 

1591 ) 

1592 lines.append(f"{mode_icon} Wave {w.wave_index} [{w.mode}] — {landing_label}") 

1593 

1594 num_clusters = len(w.clusters) 

1595 for i, c in enumerate(w.clusters): 

1596 is_last_cluster = i == num_clusters - 1 

1597 c_prefix = "└── " if is_last_cluster else "├── " 

1598 c_indent = " " if is_last_cluster else "│ " 

1599 

1600 issues_str = ", ".join(f"#{num}" for num in c.issues) 

1601 lines.append(f"{c_prefix}📦 Cluster {c.cluster_id} ({issues_str}) [{c.role}]") 

1602 

1603 scope_items = list(c.combined_scope) 

1604 children = [ 

1605 f"Scope: {', '.join(scope_items[:3])}{'...' if len(scope_items) > 3 else ''}", 

1606 *cluster_staffing_lines(c), 

1607 ] 

1608 if c.depends_on_issues: 

1609 dep_issues = ", ".join(f"#{d}" for d in c.depends_on_issues) 

1610 children.append(f"⛓️ Depends on: {dep_issues}") 

1611 # The connector is decided from the assembled list rather than from each 

1612 # optional row in turn: every added row used to mean another place that had to 

1613 # know whether something came after it, and one that guessed wrong drew a tree 

1614 # with two last branches. 

1615 for index, child in enumerate(children): 

1616 branch = "└── " if index == len(children) - 1 else "├── " 

1617 lines.append(f"{c_indent}{branch}{child}") 

1618 lines.append("") 

1619 

1620 return "\n".join(lines).rstrip() 

1621 

1622 

1623def workers_by_wave( 

1624 workers: Sequence[SwarmWorkerStatus], 

1625) -> tuple[tuple[int, tuple[SwarmWorkerStatus, ...]], ...]: 

1626 """``workers`` grouped by wave, lowest wave first, in their recorded order within one 

1627 (#1280). Wave ``0`` — a record from before workers carried their wave — sorts first.""" 

1628 waves: dict[int, list[SwarmWorkerStatus]] = {} 

1629 for w in workers: 

1630 waves.setdefault(w.wave, []).append(w) 

1631 return tuple((wave, tuple(waves[wave])) for wave in sorted(waves)) 

1632 

1633 

1634def format_elapsed(seconds: int | None) -> str: 

1635 """``seconds`` the way the board prints it — ``45s``, ``12m05s``, ``1h02m`` — or ``-``.""" 

1636 if seconds is None: 

1637 return "-" 

1638 if seconds < 60: 

1639 return f"{seconds}s" 

1640 if seconds < 3600: 

1641 return f"{seconds // 60}m{seconds % 60:02d}s" 

1642 return f"{seconds // 3600}h{seconds % 3600 // 60:02d}m" 

1643 

1644 

1645def _status_counts(workers: Sequence[SwarmWorkerStatus]) -> dict[str, int]: 

1646 counts: dict[str, int] = {} 

1647 for w in workers: 

1648 counts[w.status] = counts.get(w.status, 0) + 1 

1649 return counts 

1650 

1651 

1652#: The width of the board's ``Review`` column (#1440). 

1653REVIEW_COLUMN = 60 

1654 

1655 

1656def review_summary(record: Mapping[str, Any] | None) -> str: 

1657 """One cell of the board: what the last live ``swarm-review`` did with the cluster (#1440). 

1658 

1659 ``posted 2/2 APPROVE @ 596e3d8a``, ``changes-requested 1/2 APPROVE, 1 REQUEST_CHANGES @ 

1660 …`` (status ``posted-changes-requested``), ``held 2/3 APPROVE, 1 failed @ …``, 

1661 ``already-merged`` — or ``not reviewed`` with no record. The head is the one the seats 

1662 reviewed; ``swarm-status`` reads no pull request, so whether it is still the pull 

1663 request's head is the reader's to check. 

1664 """ 

1665 if not record: 

1666 return "not reviewed" 

1667 seats = [s for s in record.get("seats") or () if isinstance(s, Mapping)] 

1668 status = str(record.get("status")) 

1669 # The one status too long for the cell; `--json` carries it whole. 

1670 parts = ["changes-requested" if status == "posted-changes-requested" else status] 

1671 if seats: 

1672 outcomes = [s.get("outcome") for s in seats] 

1673 tally = f"{outcomes.count('APPROVE')}/{len(seats)} APPROVE" 

1674 for outcome in ("REQUEST_CHANGES", "failed", "refused", "not-run"): 

1675 if outcome in outcomes: 

1676 tally += f", {outcomes.count(outcome)} {outcome}" 

1677 parts.append(tally) 

1678 head = record.get("head_sha") 

1679 if isinstance(head, str) and head: 

1680 parts.append(f"@ {head[:8]}") 

1681 return " ".join(parts) 

1682 

1683 

1684def swarm_status_payload(state: SwarmRunState, *, now: str = "") -> dict[str, Any]: 

1685 """What ``keel swarm-status --json`` prints for a run it read (#1280). 

1686 

1687 The run's state, each worker with its derived ``elapsed_s`` (to ``now`` while it is 

1688 still running), and ``waves``: the workers grouped by wave — each wave's clusters and 

1689 how many are in each status. ``workers`` stays the flat list it always was. 

1690 """ 

1691 return { 

1692 **state.to_dict(), 

1693 "as_of": now, 

1694 "workers": [ 

1695 { 

1696 **w.to_dict(), 

1697 "elapsed_s": w.elapsed_s(now), 

1698 "review_summary": review_summary(w.review), 

1699 } 

1700 for w in state.workers 

1701 ], 

1702 "waves": [ 

1703 { 

1704 "wave": wave, 

1705 "clusters": [w.cluster_id for w in members], 

1706 "statuses": _status_counts(members), 

1707 } 

1708 for wave, members in workers_by_wave(state.workers) 

1709 ], 

1710 } 

1711 

1712 

1713def render_swarm_status_dashboard(state: SwarmRunState | None, *, now: str = "") -> str: 

1714 """Render a terminal ASCII matrix status board of the swarm run, grouped by wave. 

1715 

1716 ``now`` is when the board is drawn, for a running worker's elapsed time; without it a 

1717 running worker's elapsed time is ``-`` (#1280). 

1718 """ 

1719 if state is None: 

1720 return "keel swarm status — no active or recent swarm run found." 

1721 

1722 status_badges = { 

1723 "queued": "[QUEUED ⏳]", 

1724 "running": "[RUNNING ⚙️]", 

1725 "passed": "[PASSED ✓]", 

1726 "failed": "[FAILED ✗]", 

1727 "merged": "[MERGED 🚢]", 

1728 # Landing kept the cluster back — its review evidence did not verify — so it 

1729 # is neither failed nor merged (#1280). 

1730 "held": "[HELD ⏸️]", 

1731 } 

1732 

1733 start_str = state.started_at[:19] if state.started_at else "pending" 

1734 hdr_info = ( 

1735 f"│ Active Wave: {state.active_wave:<3} │ Total Workers: {state.total_workers:<3} " 

1736 f"│ Started: {start_str:<24} │" 

1737 ) 

1738 # The lead and the band it staffed from sit beside the worker, because the board is 

1739 # where an operator asks "who is running this, and why that provider?" — the answer 

1740 # was in the plan JSON and nowhere a running swarm could be watched (#1017). The 

1741 # stage and elapsed time answer "how far has it got?" while it runs (#1280). 

1742 cols_hdr = ( 

1743 f"│ {'Cluster':<16} │ {'Issue':<6} │ {'Role':<8} │ {'Step':<5} │ {'Stage':<12} " 

1744 f"│ {'Elapsed':<7} │ {'Lead':<12} │ {'Band':<8} │ {'Agent / Model':<16} " 

1745 f"│ {'Status':<10} │ {'Review':<{REVIEW_COLUMN}} │" 

1746 ) 

1747 width = len(cols_hdr) - 2 

1748 lines = [ 

1749 "╭" + "─" * width + "╮", 

1750 f"│ 🐝 Keel Swarm Live Status — {state.swarm_id:<{width - 30}} │", 

1751 f"{hdr_info[:-1]}{' ' * (width - len(hdr_info) + 2)}│", 

1752 "├" + "─" * width + "┤", 

1753 cols_hdr, 

1754 ] 

1755 

1756 for wave, members in workers_by_wave(state.workers): 

1757 counts = ", ".join(f"{n} {status}" for status, n in _status_counts(members).items()) 

1758 label = f"Wave {wave}" if wave else "Wave ? (not recorded)" 

1759 heading = f"{label} — {len(members)} worker{'s' if len(members) != 1 else ''}: {counts}" 

1760 lines.append("├" + "─" * width + "┤") 

1761 lines.append(f"│ {heading:<{width - 2}} │") 

1762 for w in members: 

1763 badge = status_badges.get(w.status, f"[{w.status.upper()}]") 

1764 agent_str = f"{w.agent}:{w.model}"[:16] 

1765 elapsed = format_elapsed(w.elapsed_s(now)) 

1766 lines.append( 

1767 f"│ {w.cluster_id:<16} │ #{w.issue:<5} │ {w.role:<8} │ {w.step:<5} " 

1768 f"│ {w.stage or '-':<12} │ {elapsed:<7} │ {w.lead[:12]:<12} " 

1769 f"│ {w.difficulty[:8]:<8} │ {agent_str:<16} │ {badge:<10} " 

1770 f"│ {review_summary(w.review)[:REVIEW_COLUMN]:<{REVIEW_COLUMN}} │" 

1771 ) 

1772 

1773 lines.append("╰" + "─" * width + "╯") 

1774 return "\n".join(lines) 

1775 

1776 

1777def resolve_swarm_state_dir(root: str | Path = ".") -> Path: 

1778 """Ensure and return the path to `.keel/state/swarm/` directory.""" 

1779 p = Path(root) / ".keel" / "state" / "swarm" 

1780 p.mkdir(parents=True, exist_ok=True) 

1781 return p 

1782 

1783 

1784def save_swarm_state(state: SwarmRunState, root: str | Path = ".") -> Path: 

1785 """Persist a SwarmRunState JSON snapshot to `.keel/state/swarm/<swarm_id>.json`.""" 

1786 state_dir = resolve_swarm_state_dir(root) 

1787 file_path = state_dir / f"{state.swarm_id}.json" 

1788 # Atomic + durable, like the checkpoint and activity records. This was a bare 

1789 # `write_text`: the same torn-file-on-interruption bug #869 fixed in its two 

1790 # named files — `checkpoint.py` and `activity.py` — and never reached here, 

1791 # because each writer carried its own copy of the dance instead of sharing 

1792 # one. The shared `write_text_atomic` arrived with #932. (This comment said 

1793 # #872 until #1290; that issue is an unrelated `gh api` fix, and a reader 

1794 # sent there finds nothing about torn files.) 

1795 workspace.write_text_atomic(file_path, json.dumps(state.to_dict(), indent=2)) 

1796 return file_path 

1797 

1798 

1799def load_swarm_state(swarm_id: str, root: str | Path = ".") -> SwarmRunState | None: 

1800 """Load a SwarmRunState from disk if present.""" 

1801 state_dir = Path(root) / ".keel" / "state" / "swarm" 

1802 file_path = state_dir / f"{swarm_id}.json" 

1803 if not file_path.exists(): 

1804 return None 

1805 try: 

1806 data = json.loads(file_path.read_text(encoding="utf-8")) 

1807 # JSON that parses but has the wrong shape — a hand edit, a torn write — is 

1808 # as unreadable as JSON that does not parse; it must not kill swarm-status, 

1809 # the recovery tool (#1273). 

1810 if not isinstance(data, dict): 

1811 return None 

1812 raw_workers = data.get("workers", []) 

1813 if not isinstance(raw_workers, list) or not all(isinstance(w, dict) for w in raw_workers): 

1814 return None 

1815 workers = tuple( 

1816 SwarmWorkerStatus( 

1817 cluster_id=str(w.get("cluster_id", "")), 

1818 issue=int(w.get("issue", 0)), 

1819 role=str(w.get("role", "core")), 

1820 agent=str(w.get("agent", "claude")), 

1821 model=str(w.get("model", "default")), 

1822 step=str(w.get("step", "s0")), 

1823 status=str(w.get("status", "queued")), 

1824 updated_at=str(w.get("updated_at", "")), 

1825 details=str(w.get("details", "")), 

1826 lead=str(w.get("lead", "")), 

1827 difficulty=str(w.get("difficulty", "")), 

1828 scopes=_stored_scopes(w.get("scopes")), 

1829 pull_request=_stored_pull_request(w.get("pull_request")), 

1830 worktree=str(w.get("worktree") or ""), 

1831 # Only a `true` that was written is a push: `--clean` keeps what it says. 

1832 pushed=w.get("pushed") is True, 

1833 # A record from before #1280 has none of these, and still loads. 

1834 wave=_stored_wave(w.get("wave")), 

1835 stage=_stored_stage(w.get("stage")), 

1836 started_at=_stored_text(w.get("started_at")), 

1837 finished_at=_stored_text(w.get("finished_at")), 

1838 # A record from before #1440 has none: the cluster reads as not reviewed. 

1839 review=_stored_review(w.get("review")), 

1840 ) 

1841 for w in raw_workers 

1842 ) 

1843 stored_consent = data.get("consent") 

1844 return SwarmRunState( 

1845 swarm_id=str(data.get("swarm_id", swarm_id)), 

1846 total_workers=int(data.get("total_workers", len(workers))), 

1847 active_wave=int(data.get("active_wave", 1)), 

1848 workers=workers, 

1849 started_at=str(data.get("started_at", "")), 

1850 completed_at=data.get("completed_at"), 

1851 consent=stored_consent if isinstance(stored_consent, dict) else None, 

1852 ) 

1853 # OverflowError: `1e999` is valid JSON, parses to infinity, and `int()` refuses it. 

1854 # OSError: a file that exists but cannot be opened — no read permission, or a 

1855 # directory named `<id>.json` — is unreadable in the same sense (#1280). 

1856 except (json.JSONDecodeError, ValueError, KeyError, TypeError, OverflowError, OSError): 

1857 return None 

1858 

1859 

1860def _stored_scopes(value: Any) -> tuple[str, ...]: 

1861 """A worker record's ``scopes`` as written, or none — a record from before #1400 has 

1862 no such field, and a malformed one must not read as a scope the worker held.""" 

1863 if not isinstance(value, list): 

1864 return () 

1865 return tuple(scope for scope in value if isinstance(scope, str)) 

1866 

1867 

1868def _stored_pull_request(value: Any) -> int | None: 

1869 """A worker record's ``pull_request`` as written, or none — a record from before 

1870 #1287 has no such field, and anything but a positive number is not a pull request.""" 

1871 if isinstance(value, bool) or not isinstance(value, int) or value < 1: 

1872 return None 

1873 return value 

1874 

1875 

1876def _stored_wave(value: Any) -> int: 

1877 """A worker record's ``wave`` as written, or ``0`` — "not recorded" (#1280).""" 

1878 if isinstance(value, bool) or not isinstance(value, int) or value < 0: 

1879 return 0 

1880 return value 

1881 

1882 

1883def _stored_stage(value: Any) -> str: 

1884 """A worker record's ``stage`` when it is one a live worker has, else ``""`` (#1280).""" 

1885 return value if value in swarm_worker.STAGES else "" 

1886 

1887 

1888def _stored_review(value: Any) -> dict[str, Any] | None: 

1889 """A worker's ``review`` record as written, or ``None`` (#1440). 

1890 

1891 Kept only when it names a status; its seats are kept only as objects. A record of 

1892 another shape — a hand edit, a future writer — reads as "not reviewed" rather than 

1893 failing the load of the whole run. 

1894 """ 

1895 if not isinstance(value, dict): 

1896 return None 

1897 status = value.get("status") 

1898 if not isinstance(status, str) or not status: 

1899 return None 

1900 seats = value.get("seats") 

1901 seats = [s for s in seats if isinstance(s, dict)] if isinstance(seats, list) else [] 

1902 return {**value, "seats": seats} 

1903 

1904 

1905def _stored_text(value: Any) -> str: 

1906 """A worker record's text field as written, or ``""`` when it is not text (#1280).""" 

1907 return value if isinstance(value, str) else "" 

1908 

1909 

1910#: The number at the end of a pull request URL, as ``gh pr create`` prints it. 

1911_PULL_REQUEST_URL = re.compile(r"/pull/([1-9][0-9]*)/?$") 

1912 

1913 

1914def pull_request_number(url: str) -> int | None: 

1915 """The pull request number in a ``gh pr create`` URL, or ``None`` when it names none.""" 

1916 match = _PULL_REQUEST_URL.search(url.strip()) 

1917 return int(match.group(1)) if match else None 

1918 

1919 

1920#: What a persisted plan file says it is (#1275), and the one layout this keel reads. 

1921#: The version is bumped whenever :meth:`SwarmPlan.to_dict` changes shape; a file of any 

1922#: other version is refused rather than read by guesswork. 

1923SWARM_PLAN_SCHEMA = "keel.swarm-plan" 

1924SWARM_PLAN_VERSION = 1 

1925#: A plan lives beside its run's state file: ``<swarm_id>.json`` is the state, 

1926#: ``<swarm_id>.plan.json`` the plan that run executed. 

1927SWARM_PLAN_SUFFIX = ".plan.json" 

1928 

1929 

1930def swarm_plan_payload(plan: SwarmPlan) -> dict[str, Any]: 

1931 """The persisted form of ``plan``: its :meth:`SwarmPlan.to_dict` under a versioned envelope.""" 

1932 return {"schema": SWARM_PLAN_SCHEMA, "version": SWARM_PLAN_VERSION, "plan": plan.to_dict()} 

1933 

1934 

1935def swarm_plan_from_payload(data: Any) -> SwarmPlan: 

1936 """The plan in a :func:`swarm_plan_payload`, or :class:`SwarmPlanError` saying why not.""" 

1937 envelope = _plan_object(data, "the file") 

1938 if envelope.get("schema") != SWARM_PLAN_SCHEMA: 

1939 raise SwarmPlanError(f"not a keel swarm plan (its 'schema' is not {SWARM_PLAN_SCHEMA!r})") 

1940 version = envelope.get("version") 

1941 if isinstance(version, bool) or version != SWARM_PLAN_VERSION: 

1942 raise SwarmPlanError( 

1943 f"schema version {version!r} is not one this keel reads ({SWARM_PLAN_VERSION}); " 

1944 "land it with the keel that wrote it, or re-plan with --issues after removing it" 

1945 ) 

1946 return SwarmPlan.from_dict(_plan_field(envelope, "plan", "the file"), "plan") 

1947 

1948 

1949def swarm_plan_path(swarm_id: str, root: str | Path = ".") -> Path: 

1950 """Where ``swarm-run`` persists the plan of run ``swarm_id``.""" 

1951 return Path(root) / ".keel" / "state" / "swarm" / f"{swarm_id}{SWARM_PLAN_SUFFIX}" 

1952 

1953 

1954def save_swarm_plan(plan: SwarmPlan, root: str | Path = ".") -> Path: 

1955 """Persist the plan a run executes beside its state file (#1275). 

1956 

1957 ``swarm-land`` lands this file's waves rather than re-planning from the issues, which 

1958 can partition differently once an issue's text, scope or labels change. Written with 

1959 the same atomic, durable writer as the state. 

1960 """ 

1961 resolve_swarm_state_dir(root) 

1962 file_path = swarm_plan_path(plan.swarm_id, root) 

1963 workspace.write_text_atomic(file_path, json.dumps(swarm_plan_payload(plan), indent=2)) 

1964 return file_path 

1965 

1966 

1967def load_swarm_plan(swarm_id: str, root: str | Path = ".") -> SwarmPlan | None: 

1968 """The plan run ``swarm_id`` persisted, ``None`` when it persisted none. 

1969 

1970 A file that exists and cannot be used — unreadable, not JSON, malformed, an unknown 

1971 schema version, or the plan of another run — raises :class:`SwarmPlanError`. Unlike 

1972 the state (:func:`load_swarm_state`), a plan is not fail-soft: ``None`` would send the 

1973 caller to re-plan, which is exactly the silent switch this file exists to prevent. 

1974 """ 

1975 file_path = swarm_plan_path(swarm_id, root) 

1976 if not file_path.exists(): 

1977 return None 

1978 try: 

1979 data = json.loads(file_path.read_text(encoding="utf-8")) 

1980 except OSError as exc: 

1981 raise SwarmPlanError(f"{file_path}: cannot be read ({exc})") from exc 

1982 except ValueError as exc: 

1983 raise SwarmPlanError(f"{file_path}: not JSON ({exc})") from exc 

1984 try: 

1985 plan = swarm_plan_from_payload(data) 

1986 except SwarmPlanError as exc: 

1987 raise SwarmPlanError(f"{file_path}: {exc}") from exc 

1988 if plan.swarm_id != swarm_id: 

1989 raise SwarmPlanError( 

1990 f"{file_path}: holds the plan of swarm {plan.swarm_id!r}, not {swarm_id!r}" 

1991 ) 

1992 return plan 

1993 

1994 

1995def latest_swarm_id(root: str | Path = ".") -> str | None: 

1996 """The most recently written run under ``.keel/state/swarm/``, or ``None``. 

1997 

1998 Only state files count: a ``<id>.plan.json`` beside them is a run's plan, and its 

1999 stem (``<id>.plan``) names no run (#1275). 

2000 """ 

2001 state_dir = Path(root) / ".keel" / "state" / "swarm" 

2002 if not state_dir.exists(): 

2003 return None 

2004 files = sorted( 

2005 (p for p in state_dir.glob("*.json") if not p.name.endswith(SWARM_PLAN_SUFFIX)), 

2006 key=lambda p: p.stat().st_mtime, 

2007 reverse=True, 

2008 ) 

2009 return files[0].stem if files else None 

2010 

2011 

2012def _plan_issues(plan: SwarmPlan) -> tuple[int, ...]: 

2013 return tuple(sorted({i for w in plan.waves for c in w.clusters for i in c.issues})) 

2014 

2015 

2016def _scope_text(scope: IssueScope) -> str: 

2017 return f"{', '.join(scope.predicted_files) or '(none)'} ({scope.scope_source})" 

2018 

2019 

2020def swarm_plan_drift(persisted: SwarmPlan, current: SwarmPlan) -> tuple[str, ...]: 

2021 """How the plan the issues give *now* differs from the plan a run executed (#1275). 

2022 

2023 One line per difference, in a stable order: the issue set, then each wave whose 

2024 clusters differ, then each issue whose scope moved. Staffing is not compared — it 

2025 follows the machine's providers, not the issues. Empty when the two agree. 

2026 """ 

2027 lines: list[str] = [] 

2028 ran, now = _plan_issues(persisted), _plan_issues(current) 

2029 if ran != now: 

2030 lines.append( 

2031 f"issues: the run planned {', '.join(f'#{i}' for i in ran) or 'none'}; " 

2032 f"the issues named now are {', '.join(f'#{i}' for i in now) or 'none'}" 

2033 ) 

2034 ran_waves = {w.wave_index: tuple(c.cluster_id for c in w.clusters) for w in persisted.waves} 

2035 now_waves = {w.wave_index: tuple(c.cluster_id for c in w.clusters) for w in current.waves} 

2036 for index in sorted(set(ran_waves) | set(now_waves)): 

2037 before, after = ran_waves.get(index, ()), now_waves.get(index, ()) 

2038 if before != after: 

2039 lines.append( 

2040 f"wave {index}: the run planned {', '.join(before) or 'nothing'}; " 

2041 f"the issues now plan {', '.join(after) or 'nothing'}" 

2042 ) 

2043 for issue in sorted(set(persisted.issue_scopes) & set(current.issue_scopes)): 

2044 before_scope, after_scope = persisted.issue_scopes[issue], current.issue_scopes[issue] 

2045 if (before_scope.predicted_files, before_scope.scope_source) != ( 

2046 after_scope.predicted_files, 

2047 after_scope.scope_source, 

2048 ): 

2049 lines.append( 

2050 f"issue #{issue}: the run's scope was {_scope_text(before_scope)}; " 

2051 f"it is now {_scope_text(after_scope)}" 

2052 ) 

2053 return tuple(lines) 

2054 

2055 

2056def update_worker_state( 

2057 state: SwarmRunState, 

2058 cluster_id: str, 

2059 *, 

2060 step: str | None = None, 

2061 status: str | None = None, 

2062 details: str | None = None, 

2063 pull_request: int | None = None, 

2064 worktree: str | None = None, 

2065 pushed: bool | None = None, 

2066 stage: str | None = None, 

2067 started_at: str | None = None, 

2068 finished_at: str | None = None, 

2069 review: dict[str, Any] | None = None, 

2070) -> SwarmRunState: 

2071 """Return a new SwarmRunState with the specified worker's fields updated. 

2072 

2073 ``pull_request`` is recorded when given and otherwise kept, so a later status update 

2074 never forgets the pull request ``swarm-land`` has to merge (#1287). ``worktree`` and 

2075 ``pushed`` likewise (#1278), ``stage``, ``started_at`` and ``finished_at`` (#1280), and 

2076 ``review`` (#1440) — which, given, replaces the cluster's previous review record whole: 

2077 the latest review wins. 

2078 """ 

2079 # `replace` rather than a field-by-field rebuild: the rebuild had to name every 

2080 # field, so each field added to the record (the lead and difficulty band a worker 

2081 # reports, #1017) was silently reset to its default by the first status update. 

2082 updated_workers = [ 

2083 replace( 

2084 w, 

2085 step=step if step is not None else w.step, 

2086 status=status if status is not None else w.status, 

2087 updated_at=datetime.datetime.now(datetime.UTC).isoformat(), 

2088 details=details if details is not None else w.details, 

2089 pull_request=pull_request if pull_request is not None else w.pull_request, 

2090 worktree=worktree if worktree is not None else w.worktree, 

2091 pushed=pushed if pushed is not None else w.pushed, 

2092 stage=stage if stage is not None else w.stage, 

2093 started_at=started_at if started_at is not None else w.started_at, 

2094 finished_at=finished_at if finished_at is not None else w.finished_at, 

2095 review=review if review is not None else w.review, 

2096 ) 

2097 if w.cluster_id == cluster_id 

2098 else w 

2099 for w in state.workers 

2100 ] 

2101 

2102 # `replace` for the same reason: a run's `consent` record (#1400) must survive an update. 

2103 return replace(state, workers=tuple(updated_workers)) 

2104 

2105 

2106def rebalance_swarm_plan(plan: SwarmPlan, failed_issue: int) -> SwarmPlan: 

2107 """Dynamically recalculate a SwarmPlan when an issue fails during execution. 

2108 

2109 Any subsequent wave clusters that depended on ``failed_issue`` will have 

2110 the failed dependency omitted, while independent disjoint clusters 

2111 proceed without interruption. The edge is inferred from file overlap, and 

2112 work that will not land no longer overlaps anything, so the edge goes; it 

2113 used to stay, and the plan asserted a dependency on an issue it no longer 

2114 held — in ``depends_on_issues``, ``conflict_map`` and ``issue_scopes`` (#1277). 

2115 """ 

2116 new_waves = [] 

2117 for w in plan.waves: 

2118 new_clusters = [] 

2119 for c in w.clusters: 

2120 if failed_issue in c.issues: 

2121 continue 

2122 if failed_issue in c.depends_on_issues: 

2123 c = replace( 

2124 c, depends_on_issues=tuple(i for i in c.depends_on_issues if i != failed_issue) 

2125 ) 

2126 new_clusters.append(c) 

2127 if new_clusters: 

2128 # Re-derived, not copied: dropping the failed issue can remove a wave's 

2129 # last dependency, and then nothing it waits on is going to land. 

2130 mode, eligible = wave_landing_mode(w.wave_index, new_clusters) 

2131 new_waves.append( 

2132 SwarmWave( 

2133 wave_index=w.wave_index, 

2134 mode=mode, 

2135 eligible_direct_landing=eligible, 

2136 clusters=tuple(new_clusters), 

2137 ) 

2138 ) 

2139 

2140 return SwarmPlan( 

2141 swarm_id=plan.swarm_id, 

2142 total_issues=sum(len(c.issues) for w in new_waves for c in w.clusters), 

2143 waves=tuple(new_waves), 

2144 conflict_map={ 

2145 issue: tuple(other for other in others if other != failed_issue) 

2146 for issue, others in plan.conflict_map.items() 

2147 if issue != failed_issue 

2148 }, 

2149 issue_scopes={i: scope for i, scope in plan.issue_scopes.items() if i != failed_issue}, 

2150 ) 

2151 

2152 

2153def render_swarm_run_result(result: SwarmRunResult) -> str: 

2154 """Render a human-readable text summary of a SwarmRunResult.""" 

2155 status_icon = ( 

2156 "✓" if result.status == "success" else ("⚠️" if result.status == "partial_failure" else "✗") 

2157 ) 

2158 lines = [ 

2159 f"keel swarm run — {result.swarm_id}", 

2160 f" status : {result.status} {status_icon}", 

2161 f" total workers : {result.total_workers}", 

2162 f" passed : {result.passed_count}", 

2163 f" failed : {result.failed_count}", 

2164 f" dry-run : {'true' if result.dry_run else 'false'}", 

2165 f" total waves : {len(result.wave_results)}", 

2166 ] 

2167 if result.consent is not None: 

2168 lines.append( 

2169 f" consent : delegated by {result.consent.get('operator')} " 

2170 f"({', '.join(result.consent.get('scopes', ()))})" 

2171 ) 

2172 for wave in result.wave_results: 

2173 for cluster_id, res in wave.get("cluster_results", {}).items(): 

2174 if res.get("pr_url"): 

2175 lines.append(f" {cluster_id:<13} : pull request {res['pr_url']}") 

2176 elif res.get("stage") and not res.get("ok"): 

2177 lines.append(f" {cluster_id:<13} : stopped at {res['stage']} — {res['output']}") 

2178 lines.extend(f" warning : {warning}" for warning in result.warnings) 

2179 return "\n".join(lines) 

2180 

2181 

2182def evaluate_wave_landing_mode( 

2183 wave: SwarmWave, 

2184 pr_diff_map: dict[str, list[str] | tuple[str, ...]], 

2185) -> LandingDecision: 

2186 """Evaluate whether a wave can directly land, needs sequential funneling, or is refused. 

2187 

2188 A wave the plan marked ``sequential_dependent`` is **refused** whatever its size: 

2189 it exists because its clusters overlap work an earlier wave lands, so its 

2190 branches were cut from a base that has since moved (#1276). Rebasing them 

2191 through the funnel would lean on an overlap check nothing feeds real diffs yet 

2192 (#1266), so the operator lands the earlier wave and re-plans instead. Only a 

2193 wave the plan found orthogonal is then judged on its own clusters' diffs; the 

2194 funnel stays reachable for a caller that supplies ``pr_diff_map``. 

2195 """ 

2196 cluster_ids = tuple(c.cluster_id for c in wave.clusters) 

2197 if not wave.eligible_direct_landing: 

2198 return LandingDecision( 

2199 mode="refused", 

2200 eligible=False, 

2201 cluster_ids=cluster_ids, 

2202 reason="depends_on_earlier_wave", 

2203 ) 

2204 if len(cluster_ids) <= 1: 

2205 return LandingDecision( 

2206 mode="direct_batch", 

2207 eligible=True, 

2208 cluster_ids=cluster_ids, 

2209 reason="single_cluster", 

2210 ) 

2211 

2212 # Check pairwise disjointness of actual diff files 

2213 has_conflict = False 

2214 for i, c1 in enumerate(wave.clusters): 

2215 diff1 = pr_diff_map.get(c1.cluster_id, c1.combined_scope) 

2216 for c2 in wave.clusters[i + 1 :]: 

2217 diff2 = pr_diff_map.get(c2.cluster_id, c2.combined_scope) 

2218 for f1 in diff1: 

2219 for f2 in diff2: 

2220 if paths_intersect(f1, f2): 

2221 has_conflict = True 

2222 break 

2223 if has_conflict: 

2224 break 

2225 if has_conflict: 

2226 break 

2227 if has_conflict: 

2228 break 

2229 

2230 if not has_conflict: 

2231 return LandingDecision( 

2232 mode="direct_batch", 

2233 eligible=True, 

2234 cluster_ids=cluster_ids, 

2235 reason="orthogonal_diff_trees", 

2236 ) 

2237 

2238 return LandingDecision( 

2239 mode="sequential_funnel", 

2240 eligible=False, 

2241 cluster_ids=cluster_ids, 

2242 reason="overlapping_diff_trees", 

2243 ) 

2244 

2245 

2246def render_dependent_wave_refusal(wave: SwarmWave) -> str: 

2247 """The reason ``swarm-land`` gives for refusing a ``sequential_dependent`` wave. 

2248 

2249 Its branches were cut, and its pull requests reviewed, before the earlier wave it 

2250 depends on landed, so the way through is to land the earlier wave and re-plan the 

2251 rest (#1276). 

2252 """ 

2253 deps = sorted({d for c in wave.clusters for d in c.depends_on_issues}) 

2254 listed = f" ({', '.join(f'#{d}' for d in deps)})" if deps else "" 

2255 return ( 

2256 f"wave {wave.wave_index} depends on issues landed by an earlier wave{listed}; " 

2257 "its branches were cut before that landing. Land the earlier wave, then re-plan " 

2258 "the remaining issues (keel swarm-plan / swarm-run without the landed ones) and " 

2259 "land again; no pull request was merged" 

2260 ) 

2261 

2262 

2263def render_swarm_landing_result(result: SwarmLandingResult) -> str: 

2264 """Render human-readable summary of a SwarmLandingResult.""" 

2265 status_icon = ( 

2266 "✓" if result.status == "success" else ("⚠️" if result.status == "partial_failure" else "✗") 

2267 ) 

2268 lines = [ 

2269 f"keel swarm land — {result.swarm_id} (wave {result.wave_index})", 

2270 f" status : {result.status} {status_icon}", 

2271 f" mode : {result.mode}", 

2272 f" landed : {', '.join(result.landed_clusters) if result.landed_clusters else 'none'}", 

2273 f" failed : {', '.join(result.failed_clusters) if result.failed_clusters else 'none'}", 

2274 ] 

2275 if result.pull_requests: 

2276 prs = ", ".join(f"{cluster_id} #{number}" for cluster_id, number in result.pull_requests) 

2277 lines.append(f" PRs : {prs}") 

2278 if result.held_clusters: 

2279 lines.append(" held : not landed") 

2280 for cluster_id, reason in result.held_clusters: 

2281 lines.append(f" {cluster_id}: {reason}") 

2282 if result.closures: 

2283 lines.append(" closure : the landed clusters' issues") 

2284 for closure in result.closures: 

2285 lines.append(f" {render_cluster_closure(closure)}") 

2286 if result.refused: 

2287 lines.append(f" refused : {result.refused}") 

2288 for warning in result.warnings: 

2289 lines.append(f" warning : {warning}") 

2290 return "\n".join(lines)