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

345 statements  

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

1"""Structured run ledger helpers for keel workflows.""" 

2 

3from __future__ import annotations 

4 

5import json 

6from collections.abc import Callable, Collection, Mapping 

7from pathlib import Path 

8from typing import Any 

9 

10from . import capture, consent, redaction, workspace 

11from . import config as cfg 

12 

13LEDGER_SCHEMA_VERSION = "keel.run-ledger.v1" 

14CAPTURE_HEALTH_SCHEMA_VERSION = "keel.capture-health.v1" 

15DEFAULT_LEDGER_PATH = ".keel/state/run-ledger.jsonl" 

16RECORD_TYPE_SHIP_RUN = "ship_run" 

17#: The operator's live consent as ``swarm-run --live`` delegated it to its workers (#1400): 

18#: one record when the run delegates it, and one more per cluster pull request its 

19#: workers open. See :func:`build_consent_delegation_record`. 

20RECORD_TYPE_CONSENT_DELEGATION = "consent_delegation" 

21# The record kinds this keel reads. A reader skips any other kind (#1400): the ledger is 

22# append-only and shared by every keel on a checkout, so a kind a newer keel adds must not 

23# make an older keel refuse the history it can still read. The writer stays strict — 

24# :func:`encode_record` only ever writes a kind named here. 

25#: :data:`KNOWN_RECORD_TYPES` in its published order (the ledger contract's 

26#: ``record_types``). A keel older than 1.26.0 refuses a ledger holding any kind but the 

27#: first: those readers did not yet skip a kind they did not know. 

28RECORD_TYPES: tuple[str, ...] = (RECORD_TYPE_SHIP_RUN, RECORD_TYPE_CONSENT_DELEGATION) 

29KNOWN_RECORD_TYPES: frozenset[str] = frozenset(RECORD_TYPES) 

30#: What :func:`parse_records` and :func:`read_records` return unless a reader asks for 

31#: more. Every reader keel had before #1400 reads ship runs and nothing else; it gets ship 

32#: runs and nothing else, so a delegation record — which names a pull request and a head 

33#: — can never be read as a ship run by one that forgot to check the kind. 

34SHIP_RUN_RECORDS: tuple[str, ...] = (RECORD_TYPE_SHIP_RUN,) 

35#: ``delegated``: written once, before any worker starts. ``pull_request``: written once 

36#: per cluster whose worker opened its pull request, naming it and the head it pushed. 

37CONSENT_DELEGATION_EVENTS: tuple[str, ...] = ("delegated", "pull_request") 

38UNKNOWN_RECORD_TYPE_HANDLING = "skip-with-warning" 

39# A skipped kind is named in a warning; a hostile or corrupt line must not flood it. 

40_KIND_DISPLAY_LIMIT = 80 

41 

42 

43class LedgerError(ValueError): 

44 """Raised when a ledger file cannot be decoded as the stable schema.""" 

45 

46 

47def ledger_contract_as_dict(config: cfg.ProjectConfig) -> dict[str, Any]: 

48 """Return the project-neutral ledger storage and schema contract.""" 

49 path, source = configured_ledger_path(config) 

50 return { 

51 "schema_version": LEDGER_SCHEMA_VERSION, 

52 "format": "jsonl", 

53 "path": path, 

54 "path_source": source, 

55 "missing_handling": "treat-as-empty", 

56 "append_owner": ["ship", "swarm-run"], 

57 "readers": [ 

58 "morning", 

59 "wrap", 

60 "overnight", 

61 "capture-verification", 

62 "consent-verify", 

63 "ledger", 

64 ], 

65 "consumer_neutral": True, 

66 "capture_redaction": redaction.contract_as_dict(config), 

67 "capture_contract": capture.contract_as_dict(config), 

68 "capture_health": capture_health_contract_as_dict(), 

69 "record_types": list(RECORD_TYPES), 

70 "default_read_record_types": list(SHIP_RUN_RECORDS), 

71 "unknown_record_types": UNKNOWN_RECORD_TYPE_HANDLING, 

72 } 

73 

74 

75def capture_health_contract_as_dict() -> dict[str, Any]: 

76 """Return the ledger-derived capture-health summary contract.""" 

77 return { 

78 "schema_version": CAPTURE_HEALTH_SCHEMA_VERSION, 

79 "source": "run-ledger ship_run records", 

80 "readers": ["morning", "wrap", "status", "ledger"], 

81 "consumer_neutral": True, 

82 "missing_ledger_handling": "clean-empty-history", 

83 "dry_run": { 

84 "no_mutations": True, 

85 "safe_reconcile_actions_only": True, 

86 }, 

87 "states": ["clean", "needs-reconcile"], 

88 "item_statuses": ["applied", "deferred", "skipped", "missing-marker"], 

89 } 

90 

91 

92def configured_ledger_path(config: cfg.ProjectConfig) -> tuple[str, str]: 

93 """Return the configured ledger path and the config source that supplied it.""" 

94 pack = config.policy_pack or {} 

95 reports = pack.get("reports") if isinstance(pack.get("reports"), dict) else {} 

96 value = reports.get("run_ledger") 

97 if isinstance(value, str) and value.strip(): 

98 return value, "policy_pack.reports.run_ledger" 

99 return DEFAULT_LEDGER_PATH, "default" 

100 

101 

102def resolve_path(root: str | Path, config: cfg.ProjectConfig) -> Path: 

103 """Resolve the configured ledger path under ``root`` and reject escapes.""" 

104 raw, _ = configured_ledger_path(config) 

105 path = Path(raw) 

106 if workspace.is_root_anchored(raw): 

107 raise LedgerError("run ledger path must be relative to the project root") 

108 root_path = Path(root).resolve() 

109 resolved = (root_path / path).resolve() 

110 try: 

111 resolved.relative_to(root_path) 

112 except ValueError as exc: 

113 raise LedgerError("run ledger path escapes the project root") from exc 

114 return resolved 

115 

116 

117def build_ship_run_record( 

118 *, 

119 command: str, 

120 base_branch: str, 

121 changed_files: list[str] | None, 

122 declared_files: list[str] | None = None, 

123 outcomes: list[Any], 

124 verdict: Any, 

125 assessment: Any, 

126 issue_intake: dict[str, Any] | None = None, 

127 target: str | None = None, 

128 run_id: str | None = None, 

129 issue_number: int | None = None, 

130 pr_number: int | None = None, 

131 branch: str | None = None, 

132 head_sha: str | None = None, 

133 capture_status: str | None = None, 

134 capture_reason: str | None = None, 

135 capture_artifact: str | None = None, 

136 #: The files the *capture* is about, when they are not the files this run's git 

137 #: diff reported. On a post-merge s11 the local diff is empty and the sink reads 

138 #: the PR's files from the host; the ledger has to hash the same list or the 

139 #: dedupe compares two fingerprints of one lesson. ``None`` keeps 

140 #: ``changed_files``, which is every other caller. 

141 capture_changed_files: list[str] | tuple[str, ...] | None = None, 

142 capture_retrieved: list[str] | tuple[str, ...] = (), 

143 capture_not_run: bool = False, 

144 issue_title: str | None = None, 

145 issue_labels: list[str] | tuple[str, ...] = (), 

146 existing_records: list[dict[str, Any]] | None = None, 

147 config: cfg.ProjectConfig | None = None, 

148 implementer: str | None = None, 

149 reviewer_agents: list[str] | None = None, 

150 tester: str | None = None, 

151 host_agent: str | None = None, 

152 transport: str | None = None, 

153 profile: str | None = None, 

154 jury_mode: str | None = None, 

155 jury_panel: dict[str, Any] | None = None, 

156 implement_mode: str | None = None, 

157 implement_phases: list[dict[str, Any]] | None = None, 

158 implement_loop: dict[str, Any] | None = None, 

159 consent_status: str | None = None, 

160 consent_scopes: list[str] | tuple[str, ...] | None = None, 

161 run_controls: dict[str, Any] | None = None, 

162 fix_attribution: dict[str, Any] | None = None, 

163) -> dict[str, Any]: 

164 """Build one deterministic consumer-neutral ship ledger record.""" 

165 return { 

166 "schema_version": LEDGER_SCHEMA_VERSION, 

167 "record_type": RECORD_TYPE_SHIP_RUN, 

168 "command": command, 

169 "run_id": run_id, 

170 "target": target, 

171 "issue": {"number": issue_number}, 

172 "pull_request": {"number": pr_number}, 

173 "git": { 

174 "base_branch": base_branch, 

175 "branch": branch, 

176 "head_sha": head_sha, 

177 }, 

178 # `None` means git could not be read, and stays distinct from `[]` all the way 

179 # into the record: a consumer of the ledger (or of the closure comment rendered 

180 # from it) must not read "we could not see the diff" as "the diff was empty" — 

181 # that is the conflation this whole family of fixes exists to remove, and a 

182 # record claiming TIER-3 with zero files is self-contradictory besides. 

183 "changes": { 

184 "file_count": None if changed_files is None else len(changed_files), 

185 "files": None if changed_files is None else list(changed_files), 

186 "unreadable": changed_files is None, 

187 }, 

188 "declared": _declared_block(declared_files), 

189 "gates": [ 

190 { 

191 "gate": outcome.gate, 

192 "ok": outcome.ok, 

193 "skipped": outcome.skipped, 

194 "timed_out": outcome.timed_out, 

195 # getattr: outcomes reach here from adapters and fixtures that predate 

196 # these fields. Defaulting not_run to False keeps an older producer's 

197 # record readable; on_fail defaults to the strict value so a record 

198 # that *does* carry not_run without a severity fails closed. 

199 "not_run": getattr(outcome, "not_run", False), 

200 "on_fail": getattr(outcome, "on_fail", "block"), 

201 "error": outcome.error, 

202 "finding_count": len(outcome.findings), 

203 # Which panel the jury gate's verdict is (#1437): one keel convened for 

204 # this run, or the one already posted for this head and reused. 

205 **_jury_source(outcome), 

206 } 

207 for outcome in outcomes 

208 ], 

209 "verdict": { 

210 "blocked": verdict.blocked, 

211 "counts": dict(verdict.counts), 

212 }, 

213 "assessment": { 

214 "tier": assessment.tier, 

215 "reviewers": assessment.reviewers, 

216 "window_open": assessment.window_open, 

217 "ci_ok": assessment.ci_ok, 

218 "merge": { 

219 "action": assessment.merge.action, 

220 "reason": assessment.merge.reason, 

221 }, 

222 "halted": assessment.halted, 

223 "bypassed_window": assessment.bypassed_window, 

224 }, 

225 "actors": { 

226 "implementer": implementer, 

227 "reviewers": list(reviewer_agents or ()), 

228 "tester": tester, 

229 # Who took each s9 fix round, from the run-events file (#1016). An escalated 

230 # round was not fixed by the implementer, and a closure comment rendered from 

231 # `implementer` alone says it was. 

232 "fixers": _fixers(fix_attribution), 

233 "attribution_sentence": _attribution_sentence(fix_attribution), 

234 }, 

235 "run_context": _run_context( 

236 host_agent=host_agent, 

237 transport=transport, 

238 profile=profile, 

239 jury_mode=jury_mode, 

240 jury_panel=jury_panel, 

241 implement_mode=implement_mode, 

242 implement_phases=implement_phases, 

243 implement_loop=implement_loop, 

244 consent_status=consent_status, 

245 consent_scopes=consent_scopes, 

246 ), 

247 "run_controls": run_controls, 

248 "issue_intake": issue_intake, 

249 "capture": capture.record_marker( 

250 pr_number=pr_number, 

251 status=capture_status, 

252 reason=capture_reason, 

253 artifact=capture_artifact, 

254 retrieved=capture_retrieved, 

255 title=issue_title, 

256 labels=issue_labels, 

257 # Not `changed_files`: `changes.files` above records what this run's git 

258 # diff reported, which must stay None-preserving, while the capture 

259 # fingerprint has to be the one the document on disk was written with. 

260 changed_files=( 

261 capture_changed_files if capture_changed_files is not None else changed_files 

262 ), 

263 existing_records=existing_records or [], 

264 config=config, 

265 not_run=capture_not_run, 

266 ), 

267 } 

268 

269 

270def _fixers(fix_attribution: dict[str, Any] | None) -> list[dict[str, Any]]: 

271 """The per-round fixer records from a ``keel.fix-attribution.v1`` document.""" 

272 rounds = fix_attribution.get("rounds") if isinstance(fix_attribution, dict) else None 

273 if not isinstance(rounds, list): 

274 return [] 

275 return [item for item in rounds if isinstance(item, dict)] 

276 

277 

278def _attribution_sentence(fix_attribution: dict[str, Any] | None) -> str | None: 

279 """The rendered *"implemented by agy, fixed by opus in round 2"* phrase, when recorded.""" 

280 if not isinstance(fix_attribution, dict): 

281 return None 

282 sentence = fix_attribution.get("sentence") 

283 return sentence if isinstance(sentence, str) and sentence.strip() else None 

284 

285 

286def _declared_block(declared_files: list[str] | None) -> dict[str, Any] | None: 

287 """Build the implementer's declared-scope block, or ``None`` when unset. 

288 

289 ``declared_files`` is the implementer's contract of which files the change is 

290 *supposed* to touch (distinct from the observed ``changes`` diff). When the 

291 implementer does not declare a scope, the block is omitted so readers can 

292 degrade to advisory back-compat behavior. 

293 """ 

294 if declared_files is None: 

295 return None 

296 files = [str(path) for path in declared_files] 

297 return {"file_count": len(files), "files": files} 

298 

299 

300def declared_files_for_record(record: dict[str, Any]) -> list[str] | None: 

301 """Return the implementer's declared file list from a ship_run ``record``. 

302 

303 Returns ``None`` when no declared scope was recorded (a missing or malformed 

304 ``declared`` block), letting ``scope-verify`` degrade to an advisory pass. 

305 """ 

306 declared = record.get("declared") 

307 if not isinstance(declared, dict): 

308 return None 

309 files = declared.get("files") 

310 if not isinstance(files, list): 

311 return None 

312 return [str(path) for path in files] 

313 

314 

315def build_run_context( 

316 *, 

317 host_agent: str | None, 

318 transport: str | None, 

319 consent_status: str | None, 

320 consent_scopes: list[str] | tuple[str, ...] | None, 

321) -> dict[str, Any]: 

322 """The ``run_context`` block :func:`build_ship_run_record` writes, for a writer that 

323 builds its record another way — ``swarm-land``'s landing record (#1422). 

324 

325 The profile, the jury and the s4 fields are left unset: a landing knows none of them. 

326 """ 

327 return _run_context( 

328 host_agent=host_agent, 

329 transport=transport, 

330 profile=None, 

331 jury_mode=None, 

332 jury_panel=None, 

333 consent_status=consent_status, 

334 consent_scopes=consent_scopes, 

335 ) 

336 

337 

338def _run_context( 

339 *, 

340 host_agent: str | None, 

341 transport: str | None, 

342 profile: str | None, 

343 jury_mode: str | None, 

344 jury_panel: dict[str, Any] | None, 

345 consent_status: str | None, 

346 consent_scopes: list[str] | tuple[str, ...] | None, 

347 implement_mode: str | None = None, 

348 implement_phases: list[dict[str, Any]] | None = None, 

349 implement_loop: dict[str, Any] | None = None, 

350) -> dict[str, Any]: 

351 """Build the deterministic consumer-neutral preflight run-context block. 

352 

353 Every field is optional; a missing scalar degrades to ``None`` so the 

354 closure renderer can present ``unknown``/``none`` without a schema change. 

355 Consent is a small summary: a status and the approved mutation scopes, 

356 reusing the operator/approve-scope inputs already resolved by the caller. 

357 

358 ``implement_mode``/``implement_phases`` are the s4 profile (#1020). A ``default`` run 

359 records ``None`` and an empty phase list, so a record written before the knob existed 

360 reads identically to one written by a run that did not use it. 

361 """ 

362 scopes = [str(scope) for scope in (consent_scopes or ()) if str(scope).strip()] 

363 return { 

364 "host_agent": host_agent if _nonblank(host_agent) else None, 

365 "transport": transport if _nonblank(transport) else None, 

366 "profile": profile if _nonblank(profile) else None, 

367 "jury_mode": jury_mode if _nonblank(jury_mode) else None, 

368 # The panel-availability probe this run measured, or `None` when the tier 

369 # never named a panel (#1066). Recorded so a reader of the ledger can tell 

370 # a jury-reviewed change from one a host bench reviewed because the panel 

371 # could not be staffed, without re-deriving it from a machine that has 

372 # since changed. 

373 "jury_panel": dict(jury_panel) if isinstance(jury_panel, dict) else None, 

374 # The s4 profile and, under `tdd`, one record per phase — the ledger is where a 

375 # closure comment and a later audit learn that this change was written test-first 

376 # and which commit each half of s4 produced. 

377 "implement_mode": implement_mode if _nonblank(implement_mode) else None, 

378 "implement_phases": [dict(phase) for phase in implement_phases or ()], 

379 # The s4 iteration loop (#1165): its policy and one record per iteration, or 

380 # `None` for a run that neither configured nor recorded one. 

381 "implement_loop": dict(implement_loop) if isinstance(implement_loop, dict) else None, 

382 "consent": { 

383 "status": consent_status if _nonblank(consent_status) else None, 

384 "scopes": scopes, 

385 }, 

386 } 

387 

388 

389def _nonblank(value: Any) -> bool: 

390 return isinstance(value, str) and bool(value.strip()) 

391 

392 

393def encode_record(record: dict[str, Any]) -> str: 

394 """Encode one ledger record as stable JSONL.""" 

395 _validate_record(record) 

396 return json.dumps(record, sort_keys=True, separators=(",", ":")) + "\n" 

397 

398 

399def parse_records( 

400 text: str, 

401 *, 

402 warn: Callable[[str], None] | None = None, 

403 kinds: Collection[str] = SHIP_RUN_RECORDS, 

404) -> list[dict[str, Any]]: 

405 """Parse ledger JSONL text into validated records of the kinds this keel knows. 

406 

407 **Every record of a known kind is validated; only those of ``kinds`` are returned** 

408 — ship runs unless the reader asks for more (:data:`SHIP_RUN_RECORDS`). A reader that 

409 wants a ``consent_delegation`` record (``consent-verify``) or every kind (``keel 

410 ledger``) names it; every other reader is handed ship runs alone, as it always was. 

411 

412 **Forward compatible** (#1400): a well-formed record whose ``record_type`` names a 

413 kind outside :data:`KNOWN_RECORD_TYPES` is skipped, never returned and never a 

414 reason to refuse the ledger — so no reader can mistake it for a ship run, and a 

415 kind a newer keel appends does not lock an older keel out of its own history. 

416 ``warn`` is told once per skipped kind, with the line it was first seen on. 

417 

418 Everything else is refused exactly as before: invalid JSON, a non-object, a wrong 

419 ``schema_version``, a record with no ``record_type`` at all (or one that is not a 

420 non-blank string), and — through :func:`_validate_record` — any record of a known 

421 kind. Only a record that *names* a kind can be set aside as one. 

422 """ 

423 records: list[dict[str, Any]] = [] 

424 skipped: set[str] = set() 

425 for line_number, raw in enumerate(text.splitlines(), start=1): 

426 if not raw.strip(): 

427 continue 

428 try: 

429 record = json.loads(raw) 

430 except json.JSONDecodeError as exc: 

431 raise LedgerError(f"line {line_number}: invalid JSON") from exc 

432 kind = _unknown_kind(record) 

433 if kind is not None: 

434 if warn is not None and kind not in skipped: 

435 warn(unknown_kind_warning(kind, line_number)) 

436 skipped.add(kind) 

437 continue 

438 _validate_record(record, line_number=line_number) 

439 if record["record_type"] in kinds: 

440 records.append(record) 

441 return records 

442 

443 

444def unknown_kind_warning(kind: str, line_number: int) -> str: 

445 """The one warning a reader gives for a record kind it skips.""" 

446 shown = kind if len(kind) <= _KIND_DISPLAY_LIMIT else kind[:_KIND_DISPLAY_LIMIT] + "..." 

447 return ( 

448 f"run ledger line {line_number}: skipped record_type {shown!r}, which this keel " 

449 "does not know (a newer keel may have written it); later records of that type " 

450 "are skipped too" 

451 ) 

452 

453 

454def read_records( 

455 path: str | Path, 

456 *, 

457 warn: Callable[[str], None] | None = None, 

458 kinds: Collection[str] = SHIP_RUN_RECORDS, 

459) -> list[dict[str, Any]]: 

460 """Read a ledger file; a missing ledger is a valid empty history. 

461 

462 Record kinds this keel does not know are skipped, and only ``kinds`` are returned, as 

463 :func:`parse_records` does. 

464 """ 

465 ledger_path = Path(path) 

466 if not ledger_path.exists(): 

467 return [] 

468 return parse_records(ledger_path.read_text(encoding="utf-8"), warn=warn, kinds=kinds) 

469 

470 

471def build_consent_delegation_record( 

472 delegation: Mapping[str, Any], 

473 *, 

474 recorded_at: str, 

475 pull_request: Mapping[str, Any] | None = None, 

476) -> dict[str, Any]: 

477 """A ``consent_delegation`` record for ``delegation`` (#1400). 

478 

479 ``delegation`` is :meth:`keel.swarm_worker.ConsentDelegation.to_dict`: who consented, 

480 to which scopes, how, when, and for which run and clusters. Without ``pull_request`` 

481 the record is the ``delegated`` event, written before any worker starts; with it, the 

482 ``pull_request`` event for one cluster — ``{"cluster", "number", "branch", 

483 "head_sha", "url"}`` — written once that cluster's worker has opened its pull request. 

484 

485 Two events, never an update: the ledger is append-only, and the pull request number is 

486 known only after ``gh pr create``. Each record carries the whole delegation, so a 

487 reader needs one record to answer for a pull request. The parent's verbatim consent 

488 record is not copied: every field a reader needs is named here. 

489 """ 

490 record: dict[str, Any] = { 

491 "schema_version": LEDGER_SCHEMA_VERSION, 

492 "record_type": RECORD_TYPE_CONSENT_DELEGATION, 

493 "event": "delegated" if pull_request is None else "pull_request", 

494 "swarm_id": delegation.get("swarm_id"), 

495 "clusters": list(delegation.get("clusters") or ()), 

496 "scopes": list(delegation.get("scopes") or ()), 

497 "operator": delegation.get("operator"), 

498 "source": delegation.get("source"), 

499 "mode": delegation.get("mode"), 

500 "delegated_at": delegation.get("delegated_at"), 

501 "recorded_at": recorded_at, 

502 "pull_request": None if pull_request is None else dict(pull_request), 

503 } 

504 _validate_record(record) 

505 return record 

506 

507 

508def consent_delegation_for_pr( 

509 records: list[dict[str, Any]], pr_number: int 

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

511 """The latest ``pull_request`` delegation record naming ``pr_number``, or ``None``. 

512 

513 The pull request number is the only key: keel wrote it from ``gh pr create``'s answer, 

514 so it names the pull request the worker opened. A branch name is not a key — anyone can 

515 open a pull request from a fork with a ``swarm/<id>/<cluster>`` branch — and a reader 

516 that gates on the record must still check the pull request's head against the record's 

517 ``head_sha`` (:func:`keel.consentverify.consent_for_pr` does). 

518 """ 

519 match: dict[str, Any] | None = None 

520 for record in records: 

521 if record.get("record_type") != RECORD_TYPE_CONSENT_DELEGATION: 

522 continue 

523 if (record.get("pull_request") or {}).get("number") == pr_number: 

524 match = record 

525 return match 

526 

527 

528def latest_ship_run_for_pr( 

529 records: list[dict[str, Any]], 

530 pr_number: int, 

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

532 """Return the last ship_run record whose pull_request matches ``pr_number``. 

533 

534 Records are appended in chronological order, so the last match is the most 

535 recent ship run for that PR. Returns ``None`` when no record matches. 

536 

537 **Scoped by pull request, not by head, and a caller that gates on the record must 

538 say so itself.** A pull request outlives its heads, so the record this returns may 

539 have been written for a commit that is no longer the head — right for a reader 

540 asking "what happened on this PR" (capture health, scope verification), wrong for 

541 anything a merge decision hangs on. The two callers that gate check the head 

542 themselves: :func:`gates_pass_for_head` selects on ``git.head_sha`` here, and 

543 :func:`keel.juryavail.shipped` refuses a record whose head is not the one being 

544 verified — without which an earlier head's fallback weakened the current head's 

545 contract (#1068 round 2). 

546 """ 

547 match: dict[str, Any] | None = None 

548 for record in records: 

549 if record.get("record_type") != RECORD_TYPE_SHIP_RUN: 

550 continue 

551 pull_request = record.get("pull_request") 

552 number = pull_request.get("number") if isinstance(pull_request, dict) else None 

553 if number == pr_number: 

554 match = record 

555 return match 

556 

557 

558def _jury_source(outcome: Any) -> dict[str, Any]: 

559 """The jury gate's provenance fields for a ledger gate entry (#1437). 

560 

561 ``source: reused`` with the posted comment it was read from, or ``source: ran`` for a 

562 jury this run executed; nothing for any other gate, or for a jury nobody ran. 

563 """ 

564 reused = getattr(outcome, "reused_from", None) 

565 if reused is not None: 

566 return {"source": "reused", "reused_from": reused.as_dict()} 

567 if outcome.gate == "jury" and not getattr(outcome, "not_run", False): 

568 return {"source": "ran"} 

569 return {} 

570 

571 

572def record_gates_passed(record: dict[str, Any]) -> bool: 

573 """Return whether a ship_run record's gates count as a clean pass. 

574 

575 A pass requires that the run was not blocked by findings and that every 

576 recorded gate either ran clean (``ok``) or was deliberately skipped, with no 

577 gate reporting an error. A missing or malformed ``gates``/``verdict`` block 

578 degrades to "not a pass" so a corrupt record can never authorize a merge. 

579 

580 A gate marked ``not_run`` was never executed by the runner that wrote the 

581 record — an ``agentic`` gate reaching the command-only runner. For a 

582 ``block``-severity gate that is **not** a pass: certifying "gates passed" for a 

583 blocking review nobody ran is exactly the fail-open this check exists to stop. 

584 Soft (``warn``/``suggest``) not-run gates are tolerated, matching ``skipped``. 

585 Records written before this field existed have no ``not_run`` key and are 

586 unaffected. 

587 """ 

588 verdict = record.get("verdict") 

589 if not isinstance(verdict, dict) or verdict.get("blocked") is not False: 

590 return False 

591 gates = record.get("gates") 

592 if not isinstance(gates, list) or not gates: 

593 return False 

594 for gate in gates: 

595 if not isinstance(gate, dict): 

596 return False 

597 if gate.get("error"): 

598 return False 

599 if not (gate.get("ok") is True or gate.get("skipped") is True): 

600 return False 

601 # Missing `on_fail` defaults to the strict value *here*, at read time. Defaulting 

602 # only at write time protects records keel wrote and no others: any producer that 

603 # learns `not_run` without its sibling key would otherwise fail open in exactly 

604 # the certification path this check exists to close. 

605 # Fail closed on anything not explicitly soft: a missing key, a JSON-round-tripped 

606 # `None`, or a severity name keel does not know all mean "we cannot tell this was 

607 # optional", and this is the certification path. 

608 if gate.get("not_run") is True and gate.get("on_fail") not in ("warn", "suggest"): 

609 return False 

610 return True 

611 

612 

613def gates_pass_for_head( 

614 records: list[dict[str, Any]], 

615 pr_number: int, 

616 head_sha: str, 

617 covered_heads: Collection[str] = (), 

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

619 """Find a passing gates run recorded against ``head_sha`` for ``pr_number``. 

620 

621 Returns ``(matched, record)``. ``matched`` is ``True`` only when the **latest** 

622 ship_run record for the PR carrying the exact current ``head_sha`` passed its 

623 gates (see :func:`record_gates_passed`); the record is returned on a match and 

624 ``None`` otherwise. A blank ``head_sha`` never matches — an unknown head must 

625 not be authorized by a stale green run. This is a pure function: it reads only 

626 its arguments and performs no I/O. 

627 

628 Latest-wins, not any-pass. Re-gating the same head is ordinary — a flaky suite 

629 settles, a fix-loop re-runs, a dependency lands — so the same ``head_sha`` can 

630 carry a green record followed by a red one. Scanning for *any* green would let 

631 the superseded pass authorize the merge and the later red never be consulted: 

632 a fail-open in exactly the gate that exists to hold the merge closed. Only the 

633 most recent verdict for the head counts. 

634 

635 ``covered_heads`` are heads the current one descends from by capture commits alone 

636 (#1203): the learning lands on the pull request's own branch after the gates ran, so 

637 the gates-pass is recorded against a head the branch has moved one lesson past. A 

638 record for one of them counts, and **latest-wins still holds across the whole set** — 

639 a red record for any of them after a green one is the verdict. The set is produced by 

640 `capture.capture_only_descent` and nothing else, so it admits a markdown file inside 

641 the sink, never code; CI still runs on the new head and `keel merge` still requires 

642 that rollup to be green. 

643 """ 

644 if not isinstance(head_sha, str) or not head_sha.strip(): 

645 return False, None 

646 accepted = {head_sha, *covered_heads} 

647 latest: dict[str, Any] | None = None 

648 for record in records: 

649 if record.get("record_type") != RECORD_TYPE_SHIP_RUN: 

650 continue 

651 pull_request = record.get("pull_request") 

652 number = pull_request.get("number") if isinstance(pull_request, dict) else None 

653 if number != pr_number: 

654 continue 

655 git = record.get("git") 

656 record_sha = git.get("head_sha") if isinstance(git, dict) else None 

657 if record_sha not in accepted: 

658 continue 

659 latest = record 

660 if latest is None or not record_gates_passed(latest): 

661 return False, None 

662 return True, latest 

663 

664 

665def _latest_per_pr(records: list[dict[str, Any]]) -> list[dict[str, Any]]: 

666 """One record per pull request — the last — keeping records that name none. 

667 

668 A pull request has one capture, and since #1157 it can leave more than one 

669 marked record: the writer's clash is keyed by ``(pull request, head)``, so a 

670 superseded head's marker survives beside the merged head's. Counting rows 

671 rather than pull requests then reported one merged pull request twice — 

672 ``applied: 1, skipped: 1`` for a single capture — in morning, wrap and status. 

673 Insertion order is preserved so these readers still list captures in the order 

674 the ledger recorded them. 

675 """ 

676 latest: dict[int, dict[str, Any]] = {} 

677 unkeyed: list[dict[str, Any]] = [] 

678 for record in records: 

679 pull_request = record.get("pull_request") 

680 number = pull_request.get("number") if isinstance(pull_request, dict) else None 

681 if isinstance(number, int): 

682 latest[number] = record 

683 else: 

684 unkeyed.append(record) 

685 # Identity, not equality: two equal rows are still two rows. A set of id()s makes 

686 # the filter one pass; `any(record is kept for kept in keep)` was records × kept, 

687 # and the ledger is append-only, so that product only grows (#1355). Every id is 

688 # of an object `records` still holds, so none can be reused mid-comprehension. 

689 kept_ids = {id(kept) for kept in [*latest.values(), *unkeyed]} 

690 return [record for record in records if id(record) in kept_ids] 

691 

692 

693def capture_health_summary(records: list[dict[str, Any]]) -> dict[str, Any]: 

694 """Summarize capture visibility for morning, wrap, status, and ledger readers.""" 

695 merged_records = _latest_per_pr([r for r in records if _is_merged_ship_run(r)]) 

696 items = [_capture_health_item(record) for record in merged_records] 

697 counts = { 

698 "applied": 0, 

699 "marker_only": 0, 

700 "create_learning": 0, 

701 "duplicate_learning": 0, 

702 "deferred": 0, 

703 "skipped": 0, 

704 "missing_marker": 0, 

705 "needs_reconcile": 0, 

706 } 

707 skipped_by_reason: dict[str, int] = {} 

708 for item in items: 

709 status = item["status"] 

710 learning = item["learning_decision"] 

711 if status == "applied": 

712 counts["applied"] += 1 

713 elif status == "deferred": 

714 counts["deferred"] += 1 

715 elif status == "skipped": 

716 counts["skipped"] += 1 

717 skipped_by_reason[item["reason"] or "unspecified"] = ( 

718 skipped_by_reason.get(item["reason"] or "unspecified", 0) + 1 

719 ) 

720 else: 

721 counts["missing_marker"] += 1 

722 if learning == "marker-only": 

723 counts["marker_only"] += 1 

724 elif learning == "create-learning": 

725 counts["create_learning"] += 1 

726 elif learning == "duplicate": 

727 counts["duplicate_learning"] += 1 

728 if item["needs_reconcile"]: 

729 counts["needs_reconcile"] += 1 

730 return { 

731 "schema_version": CAPTURE_HEALTH_SCHEMA_VERSION, 

732 "status": "needs-reconcile" if counts["needs_reconcile"] else "clean", 

733 "record_count": len(merged_records), 

734 "counts": counts, 

735 "skipped_by_reason": dict(sorted(skipped_by_reason.items())), 

736 "items": items, 

737 "reconcile_actions": [action for item in items for action in item["reconcile_actions"]], 

738 "dry_run": { 

739 "no_mutations": True, 

740 "description": "Morning and wrap surface these actions; they do not mutate " 

741 "ledger, GitHub, or capture destinations.", 

742 }, 

743 } 

744 

745 

746def _is_merged_ship_run(record: Any) -> bool: 

747 if not isinstance(record, dict): 

748 return False 

749 if record.get("record_type") != RECORD_TYPE_SHIP_RUN: 

750 return False 

751 capture_block = record.get("capture") 

752 # The operator declared this run never reached capture, which settles the 

753 # question before the assessment is consulted: that block records what ship 

754 # *recommended*, and a recommendation to merge is read here as proof the PR 

755 # merged. For a run re-recording gates after a rebase both are true at once — 

756 # gates green, nothing merged — and the record would otherwise be counted as a 

757 # merged PR whose marker went missing, a capture gap that is not there (#945). 

758 # Only an explicit declaration does this; a plain ``status: None`` still falls 

759 # through, so a marker that genuinely went missing is still reported. 

760 if isinstance(capture_block, dict) and capture_block.get("not_run") is True: 

761 return False 

762 assessment = record.get("assessment") 

763 if isinstance(assessment, dict): 

764 merge = assessment.get("merge") 

765 if isinstance(merge, dict) and merge.get("action") == "merge": 

766 return True 

767 if isinstance(capture_block, dict) and capture_block.get("status") in ( 

768 "applied", 

769 "deferred", 

770 "skipped", 

771 ): 

772 return True 

773 return False 

774 

775 

776def _capture_marker(record: dict[str, Any]) -> str | None: 

777 capture = record.get("capture") 

778 marker = capture.get("marker") if isinstance(capture, dict) else None 

779 return marker if isinstance(marker, str) and marker.strip() else None 

780 

781 

782def record_head_sha(record: Mapping[str, Any]) -> str | None: 

783 """The head a ship_run record was written for, or ``None`` when it names none.""" 

784 git = record.get("git") 

785 head = git.get("head_sha") if isinstance(git, Mapping) else None 

786 return head.strip() if isinstance(head, str) and head.strip() else None 

787 

788 

789def capture_marker_for_head( 

790 records: list[dict[str, Any]], 

791 *, 

792 pr_number: int | None, 

793 head_sha: str | None, 

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

795 """The recorded capture marker for this ``(pull request, head)``, if any. 

796 

797 The same rule :func:`existing_capture_marker` enforces, asked **before** a 

798 record exists. The clash is keyed on the pair and nothing else, so a caller 

799 about to do durable work for this run can find out whether its append will 

800 land — a learning file written ahead of an append that then no-ops is a 

801 document on disk that no ledger record will ever point at. 

802 """ 

803 if not isinstance(pr_number, int): 

804 return None 

805 head = head_sha.strip() if isinstance(head_sha, str) and head_sha.strip() else None 

806 for existing in records: 

807 if existing.get("record_type") != RECORD_TYPE_SHIP_RUN: 

808 continue 

809 other = existing.get("pull_request") 

810 if (other.get("number") if isinstance(other, dict) else None) != pr_number: 

811 continue 

812 if record_head_sha(existing) != head: 

813 continue 

814 if _capture_marker(existing) is not None: 

815 return existing 

816 return None 

817 

818 

819def existing_capture_marker( 

820 records: list[dict[str, Any]], record: dict[str, Any] 

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

822 """The already-recorded capture marker ``record`` would duplicate, if any. 

823 

824 One capture marker per **(pull request, head)**, enforced at write time. It was 

825 once only *detected*, and only afterwards: :func:`keel.capture.verify_session` 

826 refuses the whole session on a second one ("multiple capture markers found for 

827 merged PR"), ``capture-reconcile`` returns ``blocked`` with no actions to offer, 

828 and nothing in this module can remove a line — so the recovery was editing the 

829 ledger by hand, which is forging audit history to make a gate pass. 

830 

831 Re-running the same append is the most natural thing to do after a crash mid-s11, 

832 which made the obvious recovery the very action that bricks the run. Checking here 

833 costs one pass over records the caller already holds. Returns the conflicting 

834 record so the caller can name it; ``None`` when the append is new. 

835 

836 **Scoped to the head, because that is how the merge gate reads it** (#1157). Keyed 

837 by pull request alone, this refused every later head once *any* record carried a 

838 marker — including the very first run, red. :func:`gates_pass_for_head` and 

839 :func:`keel.juryavail.is_ship_run_for_head` then asked for a passing record on the 

840 *current* head, which no run was any longer allowed to write. A pull request whose 

841 first ship run failed became permanently unmergeable, and the documented exit was 

842 editing an append-only audit ledger to make a gate pass, which is the one action 

843 this design exists to prevent. The writer now keys the record the way every reader 

844 that gates does: retrying the same head still clashes, a new head is a new run. 

845 

846 A record naming no head is compared to other records naming no head — the same 

847 rule, applied to the value they have, rather than an exemption from it. 

848 """ 

849 if _capture_marker(record) is None: 

850 return None 

851 pull_request = record.get("pull_request") 

852 pr = pull_request.get("number") if isinstance(pull_request, dict) else None 

853 if not isinstance(pr, int): 

854 return None 

855 return capture_marker_for_head(records, pr_number=pr, head_sha=record_head_sha(record)) 

856 

857 

858def append_record(path: str | Path, record: dict[str, Any]) -> None: 

859 """Append one validated JSONL record, creating parent directories as needed.""" 

860 ledger_path = Path(path) 

861 ledger_path.parent.mkdir(parents=True, exist_ok=True) 

862 workspace.ensure_runtime_gitignore_for(ledger_path) 

863 with ledger_path.open("a", encoding="utf-8") as handle: 

864 handle.write(encode_record(record)) 

865 

866 

867def sanitize_record( 

868 record: dict[str, Any], 

869 config: cfg.ProjectConfig | None = None, 

870) -> dict[str, Any]: 

871 """Apply capture redaction before a ledger record becomes durable.""" 

872 result = redaction.sanitize(record, redaction.policy_from_config(config)) 

873 sanitized = dict(result.value) 

874 sanitized["redaction"] = result.audit 

875 return sanitized 

876 

877 

878def _unknown_kind(record: Any) -> str | None: 

879 """The ``record_type`` of a record of a kind this keel does not know, else ``None``. 

880 

881 Only a record that is otherwise in this ledger's schema and *names* its kind 

882 qualifies: a missing, blank or non-string ``record_type`` is not a kind to skip but 

883 a malformed record, and :func:`_validate_record` refuses it as it always has. 

884 """ 

885 if not isinstance(record, dict) or record.get("schema_version") != LEDGER_SCHEMA_VERSION: 

886 return None 

887 kind = record.get("record_type") 

888 if not isinstance(kind, str) or not kind.strip() or kind in KNOWN_RECORD_TYPES: 

889 return None 

890 return kind 

891 

892 

893def _validate_record(record: Any, *, line_number: int | None = None) -> None: 

894 prefix = f"line {line_number}: " if line_number is not None else "" 

895 if not isinstance(record, dict): 

896 raise LedgerError(f"{prefix}record must be an object") 

897 if record.get("schema_version") != LEDGER_SCHEMA_VERSION: 

898 raise LedgerError(f"{prefix}unsupported schema_version") 

899 kind = record.get("record_type") 

900 if kind == RECORD_TYPE_CONSENT_DELEGATION: 

901 if problem := _consent_delegation_problem(record): 

902 raise LedgerError(f"{prefix}invalid consent_delegation record: {problem}") 

903 elif kind != RECORD_TYPE_SHIP_RUN: 

904 raise LedgerError(f"{prefix}unsupported record_type") 

905 

906 

907def _consent_delegation_problem(record: dict[str, Any]) -> str: 

908 """What is wrong with a ``consent_delegation`` record, or ``""`` when nothing is. 

909 

910 Strict about every field a reader relies on; a field this keel does not name is 

911 tolerated, so a later keel can add one without making this one refuse the ledger. 

912 """ 

913 event = record.get("event") 

914 if event not in CONSENT_DELEGATION_EVENTS: 

915 return f"event must be one of {', '.join(CONSENT_DELEGATION_EVENTS)}" 

916 for key in ("swarm_id", "operator", "source", "mode", "delegated_at", "recorded_at"): 

917 if not _nonblank(record.get(key)): 

918 return f"{key} must be a non-blank string" 

919 clusters = record.get("clusters") 

920 if ( 

921 not isinstance(clusters, list) 

922 or not clusters 

923 or not all(_nonblank(c) for c in clusters) 

924 or len(set(clusters)) != len(clusters) 

925 ): 

926 return "clusters must be a non-empty list of distinct non-blank strings" 

927 scopes = record.get("scopes") 

928 if not isinstance(scopes, list) or not scopes or not all(_nonblank(s) for s in scopes): 

929 return "scopes must be a non-empty list of non-blank strings" 

930 try: 

931 normalized = list(consent.normalize_scopes(scopes)) 

932 except ValueError as exc: 

933 return str(exc) 

934 if normalized != scopes: 

935 return f"scopes must be normalised: {normalized}" 

936 return _delegated_pull_request_problem(event, record.get("pull_request"), clusters) 

937 

938 

939def _delegated_pull_request_problem(event: Any, pull_request: Any, clusters: list[Any]) -> str: 

940 """What is wrong with a delegation record's ``pull_request`` for its ``event``.""" 

941 if event == "delegated": 

942 return "" if pull_request is None else "a delegated event names no pull request" 

943 if not isinstance(pull_request, dict): 

944 return "a pull_request event names its pull request" 

945 if pull_request.get("cluster") not in clusters: 

946 return "pull_request.cluster must be one of the delegation's clusters" 

947 number = pull_request.get("number") 

948 if isinstance(number, bool) or not isinstance(number, int) or number < 1: 

949 return "pull_request.number must be a positive integer" 

950 for key in ("branch", "head_sha"): 

951 if not _nonblank(pull_request.get(key)): 

952 return f"pull_request.{key} must be a non-blank string" 

953 url = pull_request.get("url") 

954 if url is not None and not isinstance(url, str): 

955 return "pull_request.url must be a string" 

956 return "" 

957 

958 

959def _capture_health_item(record: dict[str, Any]) -> dict[str, Any]: 

960 block = record.get("capture") if isinstance(record.get("capture"), dict) else {} 

961 status = block.get("status") 

962 marker = block.get("marker") 

963 reason = block.get("marker_reason") or block.get("reason") 

964 learning = block.get("learning") if isinstance(block.get("learning"), dict) else {} 

965 item_status = _capture_health_status(status, marker) 

966 needs_reconcile = item_status in {"missing-marker", "deferred"} 

967 pr_number = (record.get("pull_request") or {}).get("number") 

968 item = { 

969 "run_id": record.get("run_id"), 

970 "issue": (record.get("issue") or {}).get("number"), 

971 "pull_request": pr_number, 

972 "status": item_status, 

973 "capture_status": status, 

974 "marker": marker if isinstance(marker, str) and marker.strip() else None, 

975 "reason": reason if isinstance(reason, str) and reason.strip() else None, 

976 "learning_decision": learning.get("decision"), 

977 "learning_reason": learning.get("reason"), 

978 "needs_reconcile": needs_reconcile, 

979 "reconcile_actions": [], 

980 } 

981 if needs_reconcile: 

982 item["reconcile_actions"].append(_capture_reconcile_action(pr_number, item_status)) 

983 return item 

984 

985 

986def _capture_health_status(status: Any, marker: Any) -> str: 

987 if not isinstance(marker, str) or not marker.strip(): 

988 return "missing-marker" 

989 if status == "deferred": 

990 return "deferred" 

991 if status == "skipped": 

992 return "skipped" 

993 return "applied" 

994 

995 

996def _capture_reconcile_action(pr_number: Any, status: str) -> dict[str, Any]: 

997 return { 

998 "type": "capture-reconcile", 

999 "reason": status, 

1000 "pr": pr_number if isinstance(pr_number, int) else None, 

1001 "command": ( 

1002 f"keel capture-reconcile .keel/project.yaml --root . --merged-pr {pr_number}" 

1003 if isinstance(pr_number, int) 

1004 else "keel capture-reconcile .keel/project.yaml --root ." 

1005 ), 

1006 "dry_run": True, 

1007 }