Coverage for src/ai_jury/orchestrator.py: 99%

575 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-30 06:29 +0000

1"""Jury orchestration: review -> debate -> synthesis. 

2 

3The orchestrator owns the round structure and prompt assembly; adapters only run 

4their CLI. Rounds run agents concurrently (thread pool) because each call is an 

5independent, IO-bound subprocess. 

6""" 

7 

8from __future__ import annotations 

9 

10import random 

11import shutil 

12import string 

13import time 

14from concurrent.futures import ThreadPoolExecutor 

15from dataclasses import dataclass, field, replace 

16 

17from . import convergence, injection, largediff, panel, prompts, routing 

18from .adapters import RETRYABLE_ERROR_CODES, Adapter, AgentResult, make_adapter 

19from .config import JuryConfig, agy_opt_in_hint 

20from .consensus import FindingGroup, demote_local_only_groups, group_findings 

21from .diffprofile import profile_diff 

22from .findings import ( 

23 Finding, 

24 Verdict, 

25 emitted_findings_block, 

26 parse_findings, 

27 parse_verdicts, 

28) 

29from .policy import ReviewPolicy, render_policy_section 

30from .privilege import audit_privilege 

31from .redaction import redact 

32 

33 

34class RunBudget: 

35 """Wall-clock budget for one jury run (issue #30). 

36 

37 Tracks the elapsed time since construction and derives the timeout to pass a 

38 single agent call from the optional total-run and per-phase budgets. ``None`` 

39 for either budget means uncapped; when both are unset ``call_timeout`` 

40 returns ``None`` so adapters fall back to their own per-agent timeout and 

41 behaviour is identical to having no budget at all. 

42 """ 

43 

44 def __init__(self, total_timeout: int | None, phase_timeout: int | None): 

45 self.total = total_timeout 

46 self.phase = phase_timeout 

47 self._start = time.monotonic() 

48 

49 def elapsed(self) -> float: 

50 return time.monotonic() - self._start 

51 

52 def remaining(self) -> float | None: 

53 if self.total is None: 

54 return None 

55 return max(0.0, self.total - self.elapsed()) 

56 

57 def expired(self) -> bool: 

58 return self.total is not None and self.elapsed() >= self.total 

59 

60 def call_timeout(self) -> int | None: 

61 """Per-call timeout: the min of the phase budget and remaining total. 

62 

63 The agent's own per-agent timeout is applied by the adapter (it takes the 

64 min with this value), so it is not needed here. Returns ``None`` when 

65 neither budget caps the call, leaving the adapter to use its configured 

66 per-agent timeout. 

67 """ 

68 caps: list[float] = [] 

69 if self.phase is not None: 

70 caps.append(float(self.phase)) 

71 remaining = self.remaining() 

72 if remaining is not None: 

73 caps.append(remaining) 

74 if not caps: 

75 return None 

76 return max(1, int(min(caps))) 

77 

78 

79def _run_with_retry( 

80 adapter: Adapter, 

81 prompt: str, 

82 phase: str, 

83 budget: RunBudget, 

84 retries: int, 

85 log, 

86) -> AgentResult: 

87 """Run one agent for one phase, retrying transient failures (issue #30). 

88 

89 Retries only failures whose typed error code is in 

90 ``RETRYABLE_ERROR_CODES`` (timeout/rate-limit/spawn), up to ``retries`` extra 

91 attempts. A deterministic failure (auth, missing CLI, empty output, generic 

92 nonzero exit) is returned immediately. The returned result's ``attempts`` 

93 records how many tries were made. Retrying stops early when the run budget is 

94 exhausted so a retry never overruns the total timeout. 

95 

96 The result is also stamped with the model id the adapter sent (issue #709). 

97 This is one of exactly two places a seat is invoked — the other is ``jury 

98 run-agent`` — so stamping here covers every phase of every run, and the 

99 ballot quotes :attr:`AgentResult.model` instead of deriving a second answer 

100 from the spec. 

101 """ 

102 max_attempts = max(1, retries + 1) 

103 result = adapter.run(prompt, phase=phase, timeout=budget.call_timeout()) 

104 attempts = 1 

105 while ( 

106 not result.ok 

107 and result.error_code in RETRYABLE_ERROR_CODES 

108 and attempts < max_attempts 

109 and not budget.expired() 

110 ): 

111 log(f"{adapter.name}: {phase} attempt {attempts} failed ({result.error_code}); retrying") 

112 result = adapter.run(prompt, phase=phase, timeout=budget.call_timeout()) 

113 attempts += 1 

114 result.attempts = attempts 

115 result.model = adapter.resolved_model() 

116 return result 

117 

118 

119def _order_by_agents(results: list[AgentResult], order: list[str]) -> list[AgentResult]: 

120 """Reorder phase results into the configured/enabled agent order. 

121 

122 Round phases run agents concurrently (ThreadPoolExecutor.map), so the order 

123 in which results arrive is not guaranteed across runs. The report and all 

124 downstream consumers must NOT depend on thread-completion order, so we sort 

125 every phase's results by each agent's index in ``order`` (the stable 

126 enabled-agent list). Agents not present in ``order`` (should not happen) 

127 sort to the end, preserving their relative arrival order as a stable 

128 tiebreak so the sort is total and deterministic. 

129 """ 

130 index = {name: i for i, name in enumerate(order)} 

131 fallback = len(order) 

132 return sorted(results, key=lambda r: index.get(r.agent, fallback)) 

133 

134 

135@dataclass 

136class JuryOutcome: 

137 reviews: list[AgentResult] 

138 debate: list[AgentResult] 

139 synthesis: AgentResult | None 

140 chair: str 

141 findings: list[Finding] = field(default_factory=list) 

142 warnings: list[str] = field(default_factory=list) 

143 groups: list[FindingGroup] = field(default_factory=list) 

144 verify: AgentResult | None = None 

145 verdicts: list[Verdict] = field(default_factory=list) 

146 context_mode: str = "diff-only" 

147 redact_secrets: bool = True 

148 redaction_count: int = 0 

149 injection_hits: list = field(default_factory=list) 

150 # Execution/partial-result signals (issue #30): agents skipped because their 

151 # CLI was unavailable (name, reason), and whether the run budget was 

152 # exhausted before all phases completed. 

153 skipped: list = field(default_factory=list) 

154 budget_exhausted: bool = False 

155 # Adaptive-rounds signals (issue #40): rounds actually executed and a short 

156 # human-readable reason for why the debate ran / stopped. 

157 rounds_executed: int = 1 

158 stop_reason: str = "" 

159 # Set when this outcome was served from the local result cache (issue #33), 

160 # so the report/metadata can mark it as cached rather than freshly computed. 

161 from_cache: bool = False 

162 # The paths and symbols in the change the panel was shown (issue #710), so a 

163 # ballot's ``Checked:`` line can be resolved against something instead of 

164 # being accepted on shape alone. ``None`` means "this outcome was not built 

165 # from a diff" — a hand-assembled outcome, a caller that never had one — and 

166 # the scope rule then falls back to the structural test it applied before. 

167 changed: largediff.ChangeIndex | None = None 

168 # What tiered routing decided (#714): the plan's dict, or the standard record. 

169 routing: dict = field(default_factory=lambda: routing.standard_plan([]).as_dict()) 

170 

171 

172def _run_phase( 

173 adapters: list[Adapter], 

174 prompt_for: dict[str, str], 

175 phase: str, 

176 parallel: bool, 

177 *, 

178 budget: RunBudget, 

179 retries: int, 

180 log, 

181) -> list[AgentResult]: 

182 def task(a: Adapter) -> AgentResult: 

183 return _run_with_retry(a, prompt_for[a.name], phase, budget, retries, log) 

184 

185 if parallel and len(adapters) > 1: 

186 with ThreadPoolExecutor(max_workers=len(adapters)) as pool: 

187 return list(pool.map(task, adapters)) 

188 return [task(a) for a in adapters] 

189 

190 

191def _others(reviews: list[AgentResult], me: str) -> str: 

192 """Identity-labeled peer reviews (legacy path; ``anonymize_debate = false``). 

193 

194 Renders each *other* reviewer's round-1 output with its real agent/vendor 

195 identity in the stable enabled-agent order. This is the pre-#37 behaviour and 

196 leaks both identity and position; the anonymizing path below is the default. 

197 """ 

198 chunks = [ 

199 f"### {r.agent} ({r.vendor})\n{r.output}" 

200 for r in reviews 

201 if r.agent != me and r.ok and r.output 

202 ] 

203 return "\n\n".join(chunks) if chunks else "_(no other reviews available)_" 

204 

205 

206def _anon_label(i: int) -> str: 

207 """Stable anonymous reviewer label: 0->'A', 1->'B', ... 26->'AA'.""" 

208 letters = string.ascii_uppercase 

209 label = "" 

210 i += 1 

211 while i > 0: 

212 i, rem = divmod(i - 1, 26) 

213 label = letters[rem] + label 

214 return label 

215 

216 

217def _anonymize_peers( 

218 reviews: list[AgentResult], me: str, rng: random.Random 

219) -> tuple[str, dict[str, str]]: 

220 """Chatham House peer view for a debater (#37). 

221 

222 Returns ``(prompt_text, label_to_agent)`` where the prompt text renders each 

223 *other* successful reviewer's round-1 output under an anonymous 

224 ``### Reviewer A`` / ``### Reviewer B`` heading — NO vendor or agent name. 

225 The debater's OWN review is excluded (it is passed separately as 

226 ``own_review``). Presentation order is shuffled DETERMINISTICALLY using the 

227 shared run RNG so neither identity nor position is a stable signal; the same 

228 seed yields the same order, different seeds may differ. 

229 

230 ``label_to_agent`` keeps the anonymous-label -> real-agent mapping internal so 

231 callers can still recover authorship (the report attributes by real name). 

232 """ 

233 peers = [r for r in reviews if r.agent != me and r.ok and r.output] 

234 if not peers: 

235 return "_(no other reviews available)_", {} 

236 # Deterministic per-debater shuffle from the shared run RNG. We shuffle a 

237 # copy so the caller's review list (used elsewhere) is untouched. 

238 order = list(peers) 

239 rng.shuffle(order) 

240 chunks: list[str] = [] 

241 label_to_agent: dict[str, str] = {} 

242 for i, r in enumerate(order): 

243 label = f"Reviewer {_anon_label(i)}" 

244 label_to_agent[label] = r.agent 

245 chunks.append(f"### {label}\n{r.output}") 

246 return "\n\n".join(chunks), label_to_agent 

247 

248 

249def _debate_round( 

250 debaters: list[Adapter], 

251 reviews: list[AgentResult], 

252 diff: str, 

253 config: JuryConfig, 

254 run_rng: random.Random, 

255 agent_order: list[str], 

256 prior: list[AgentResult], 

257 budget: RunBudget, 

258 retries: int, 

259 log, 

260 round_no: int, 

261 template: str = prompts.DEBATE, 

262) -> list[AgentResult]: 

263 """Run one debate round and return its results in stable agent order. 

264 

265 ``prior`` holds the previous round's debate outputs (empty for the first 

266 debate round); when present they are appended to each debater's prompt as a 

267 "prior debate" addendum so later rounds in an adaptive run (issue #40) build 

268 on, rather than repeat, earlier cross-examination. The peer-review anonymizing 

269 path (#37) is preserved unchanged. 

270 """ 

271 log(f"round {round_no}: {len(debaters)} agents cross-examining") 

272 own = {r.agent: r.output for r in reviews if r.ok} 

273 # Prior-round debate output quotes attacker-controlled diff text, so it is 

274 # untrusted: neutralize sentinels (issue #316/L-1) before it is fenced and 

275 # appended below, matching every other peer-output slot. 

276 prior_txt = prompts.neutralize_sentinels( 

277 "\n\n".join(f"### {r.agent}\n{r.output}" for r in prior if r.ok and r.output) 

278 ) 

279 debate_prompt: dict[str, str] = {} 

280 for a in debaters: 

281 if config.anonymize_debate: 

282 # Per-debater deterministic shuffle: derive a child RNG from the 

283 # shared run RNG so each debater gets an independent but reproducible 

284 # peer ordering (same seed -> same order). 

285 peer_rng = random.Random(run_rng.random()) 

286 other_reviews, _label_map = _anonymize_peers(reviews, a.name, peer_rng) 

287 else: 

288 other_reviews = _others(reviews, a.name) 

289 text = template.format( 

290 name=a.name, 

291 diff=prompts.neutralize_sentinels(diff), 

292 own_review=prompts.neutralize_sentinels( 

293 own.get(a.name, "_(your review was unavailable)_") 

294 ), 

295 other_reviews=prompts.neutralize_sentinels(other_reviews), 

296 notice=prompts._UNTRUSTED_NOTICE, 

297 ) 

298 if prior_txt: 

299 text += ( 

300 "\n\n=== PRIOR DEBATE (earlier round) ===\n" 

301 "Build on this; do not just repeat it. Only keep a DISPUTE or " 

302 "MISSED item if it is still unresolved.\n\n" 

303 "<<<UNTRUSTED_REVIEW\n" + prior_txt + "\nUNTRUSTED_REVIEW>>>\n" 

304 ) 

305 debate_prompt[a.name] = text 

306 results = _run_phase( 

307 debaters, 

308 debate_prompt, 

309 "debate", 

310 config.parallel, 

311 budget=budget, 

312 retries=retries, 

313 log=log, 

314 ) 

315 # Same stable-ordering guarantee as round 1: independent of thread-pool 

316 # completion order. 

317 return _order_by_agents(results, agent_order) 

318 

319 

320def _skip_reason(spec) -> str: 

321 """Why an unavailable seat was skipped, worded for its transport (#831). 

322 

323 The one hardcoded ``CLI not found ({command})`` rendered ``CLI not found ()`` for a 

324 commandless local or hosted-API seat, whose ``command`` is empty. Name what is actually 

325 missing instead — a CLI on PATH, a reachable endpoint, or an API key — with the value 

326 passed through :func:`redact` so an endpoint's credentials never reach the log. 

327 """ 

328 command = getattr(spec, "command", "") or "" 

329 if command: 

330 return f"CLI not found ({redact(command)[0]})" 

331 endpoint = getattr(spec, "endpoint", "") or "" 

332 if endpoint: 

333 return f"endpoint not reachable ({redact(endpoint)[0]})" 

334 vendor = redact(getattr(spec, "vendor", "") or "hosted-API")[0] 

335 return f"{vendor} seat unavailable (no API key, or the API is unreachable)" 

336 

337 

338def run_jury( 

339 config: JuryConfig, 

340 diff: str, 

341 *, 

342 context: str = "", 

343 hints: str = "", 

344 mock: bool = False, 

345 strict: bool = False, 

346 seed: int | None = None, 

347 policy: ReviewPolicy | None = None, 

348 log=lambda _msg: None, 

349 budget: RunBudget | None = None, 

350 on_event=None, 

351 mode: str = "code", 

352 risk: str | None = None, 

353) -> JuryOutcome: 

354 # Jury mode (issue #221): "code" (default) reviews a diff with the code-review 

355 # rubric; "issue" reviews a GitHub issue's prose for completeness/clarity. 

356 # Only the prompt TEMPLATES differ — the round structure, consensus, voting, 

357 # verification, ordering, and determinism are identical. ``tmpl`` selects the 

358 # four phase templates; each is threaded into the phase that uses it so the 

359 # call sites are otherwise unchanged. 

360 tmpl = prompts.for_mode(mode) 

361 # Live play-by-play hook (issue #210): an optional callback fired after each 

362 # phase result is produced — ``on_event(kind, result, round_no=None)`` with 

363 # kind in {"review", "debate", "verify", "synthesis"}. It lets a caller stream 

364 # the deliberation as it happens (CLI ``--live``) without the orchestrator 

365 # doing any I/O itself. Fired in stable per-phase order (not thread-completion 

366 # order) so the event sequence is deterministic. Defaults to a no-op. 

367 emit = on_event or (lambda *_a, **_k: None) 

368 # Repository review policy (optional, #8): maintainer-authored, TRUSTED 

369 # content rendered into each REVIEW prompt in a clearly separated section. 

370 # When ``policy`` is None a sentinel placeholder is used, so the prompt is 

371 # unchanged except for that section. The policy is distinct from the 

372 # agent-runtime ``config`` and never enters the untrusted diff/context fences. 

373 policy_section = render_policy_section(policy) 

374 # Run reproducibility: a single shared RNG seeds every randomized 

375 # orchestration decision (future: anonymized-rebuttal order, rotating 

376 # chair, tie-breaks). The seed comes from the explicit ``seed`` argument if 

377 # given, else from ``config.seed``. We construct a dedicated 

378 # ``random.Random`` instance rather than touching the global ``random`` 

379 # module so seeding a jury run never perturbs unrelated global state. 

380 # When the seed is None the RNG is unseeded (still deterministic 

381 # orchestration; randomness, if any, is just not reproducible run-to-run). 

382 # LLM output itself is never made deterministic by this — only the 

383 # orchestration around it. ``run_rng`` is the shared run RNG: pass it to 

384 # any feature that needs reproducible randomness instead of using ``random``. 

385 run_seed = seed if seed is not None else config.seed 

386 run_rng = random.Random(run_seed) # shared run RNG (see docstring) 

387 

388 # Run budget (issue #30): a single wall-clock budget threaded through every 

389 # phase. Defaults (both None) leave behaviour identical to no budget, with 

390 # each agent bounded only by its own per-agent timeout. ``retries`` is the 

391 # number of extra attempts for transient (retryable) failures. A caller may 

392 # pass a SHARED budget so ``total_timeout`` spans a whole chunked review 

393 # rather than resetting per chunk (issue #31 / review finding). 

394 if budget is None: 

395 budget = RunBudget(config.total_timeout, config.phase_timeout) 

396 retries = config.retries 

397 

398 # Context policy: diff-only sends only the diff; expanded includes context. 

399 ctx_cfg = getattr(config, "context", None) 

400 context_mode = getattr(ctx_cfg, "mode", "diff-only") if ctx_cfg else "diff-only" 

401 redact_on = getattr(ctx_cfg, "redact_secrets", True) if ctx_cfg else True 

402 if context_mode == "diff-only": 

403 context = "" 

404 redaction_count = 0 

405 if redact_on: 

406 diff, _n1 = redact(diff) 

407 context, _n2 = redact(context) 

408 redaction_count = _n1 + _n2 

409 if redaction_count: 

410 log(f"redacted {redaction_count} secret(s) before sending to agents") 

411 

412 # Static-analysis pre-pass block (#523), joined into the Round 1 prompt AFTER 

413 # the context-mode filter above (#715). It is not user context: it is produced 

414 # locally by this run's linters, so the "diff-only" mode — the default — must 

415 # not discard it. Carrying it inside ``context`` did exactly that, and 

416 # `hints = true` reached no reviewer under any default configuration. 

417 review_context = "\n\n".join(part for part in (context, hints) if part.strip()) 

418 

419 # Prompt-injection heuristic (OWASP LLM01): scan untrusted diff/context for 

420 # patterns that try to override instructions, then SURFACE them as a synthetic 

421 # finding/warning. We never act on them; the CI gate is derived from 

422 # structured consensus (see ci.evaluate_ci), so an injected "APPROVE" 

423 # cannot flip the verdict. 

424 # Index the change AFTER redaction, so the index describes the same bytes the 

425 # panel is shown and never carries a secret the diff no longer has (#710). 

426 changed = largediff.change_index(diff) 

427 

428 injection_hits = injection.scan_inputs(diff, context) 

429 injection_findings: list[Finding] = [] 

430 if injection_hits: 

431 log(f"prompt-injection heuristic: {len(injection_hits)} suspicious pattern(s) flagged") 

432 syn = injection.hits_to_finding(injection_hits) 

433 if syn is not None: 433 ↛ 438line 433 didn't jump to line 438 because the condition on line 433 was always true

434 injection_findings.append(syn) 

435 

436 # Least-privilege audit: warn when a configured agent could perform 

437 # write/tool actions while reviewing attacker-controlled content. 

438 privilege_warnings = audit_privilege(config.enabled_agents) 

439 for w in privilege_warnings: 

440 log(f"least-privilege warning: {w}") 

441 if strict and privilege_warnings: 

442 raise RuntimeError( 

443 "least-privilege check failed (--strict): " + "; ".join(privilege_warnings) 

444 ) 

445 

446 specs = config.enabled_agents 

447 adapters = [make_adapter(s, mock=mock) for s in specs] 

448 

449 # Filter to available agents (unless strict, where a missing CLI is fatal). 

450 # Skipped agents are recorded (name, reason) so the report can state exactly 

451 # which agents never ran — part of the partial-result policy (issue #30). 

452 usable: list[Adapter] = [] 

453 skipped: list[tuple[str, str]] = [] 

454 for a in adapters: 

455 if a.available(): 

456 usable.append(a) 

457 elif strict: 

458 raise RuntimeError(f"agent '{a.name}' CLI not available: {a.spec.command}") 

459 else: 

460 reason = _skip_reason(a.spec) 

461 log(f"skipping '{a.name}': {reason}") 

462 skipped.append((a.name, reason)) 

463 if not usable: 

464 message = ( 

465 "no usable agents — install an agent CLI (claude / codex), run a local " 

466 "model, or set a hosted-API key (e.g. ANTHROPIC_API_KEY with `jury init --agents " 

467 "claude-api`) — or use --mock" 

468 ) 

469 # An agy-only machine: the one CLI it has is opt-in, so say why it was 

470 # not used rather than only telling it to install another. 

471 hint = agy_opt_in_hint((s.adapter_key for s in specs), shutil.which) 

472 if hint: 

473 message += f". Note: {hint}." 

474 raise RuntimeError(message) 

475 

476 usable_names = [a.name for a in usable] 

477 # Tiered routing (#714): a pure plan over the enabled bench, the usable 

478 # names and the diff's risk band decides who sits in round 1, who anchors 

479 # and who is benched; the standard mode records a plan too so the report 

480 # always says what happened. Benched seats stay available for escalation. 

481 if config.routing == routing.MODE_TIERED: 

482 plan = routing.plan_panel( 

483 specs, 

484 usable_names, 

485 # The band of the WHOLE change when the caller has one. A chunked 

486 # review calls this once per chunk (#31), and a chunk of a large or 

487 # security-touching diff can look routine on its own — routing the 

488 # panel off it would bench the frontier seats on exactly the change 

489 # they were kept for. `review_diff` profiles the filtered diff once 

490 # and threads the answer through every chunk. 

491 risk if risk is not None else profile_diff(diff).risk, 

492 chair=config.chair, 

493 min_vendors=int(getattr(config.ci, "min_vendors", 0) or 0), 

494 min_reviews=int(getattr(config.ci, "min_reviews", 0) or 0), 

495 ) 

496 else: 

497 plan = routing.standard_plan(usable_names) 

498 log(routing.describe(plan)) 

499 seated = set(plan.panel) 

500 round1 = [a for a in usable if a.name in seated] 

501 benched = [a for a in usable if a.name in set(plan.benched)] 

502 

503 # State the number a downstream consumer will actually receive, BEFORE the 

504 # panel is paid for (#699). "3 agents reviewing" is not that number: it is one 

505 # review per agent that answers *and names what it read*, and the chair's 

506 # synthesis record rides along beside them without being one. What is knowable 

507 # here is only the ceiling — the line says "at most" — because a seat that 

508 # returns nothing, or returns prose naming nothing checkable, abstains and is 

509 # not a review (#700, round 2). When a consumer's minimum is configured and 

510 # the bench cannot reach even the ceiling, this is a shortfall the run should 

511 # name here rather than at the consumer, an hour and three CLI invocations 

512 # later. 

513 log(panel.describe(len(round1), available=len(usable))) 

514 too_small = panel.shortfall(len(round1), config.ci.min_reviews, stage="before the panel runs") 

515 if too_small: 

516 raise RuntimeError(too_small) 

517 

518 # Round 1: independent reviews. 

519 log(f"round 1: {len(round1)} agents reviewing") 

520 review_prompt = { 

521 a.name: tmpl["review"].format( 

522 name=a.name, 

523 context=prompts.neutralize_sentinels(review_context or "_(none)_"), 

524 diff=prompts.neutralize_sentinels(diff), 

525 policy=policy_section, 

526 notice=prompts._UNTRUSTED_NOTICE, 

527 ) 

528 for a in round1 

529 } 

530 reviews = _run_phase( 

531 round1, 

532 review_prompt, 

533 "review", 

534 config.parallel, 

535 budget=budget, 

536 retries=retries, 

537 log=log, 

538 ) 

539 # Stable ordering: the thread pool can return results in any completion 

540 # order. Reorder to the enabled-agent order so the report (and every 

541 # downstream consumer) is independent of which thread finished first. 

542 agent_order = [a.name for a in round1] 

543 reviews = _order_by_agents(reviews, agent_order) 

544 

545 # Parse structured findings from each successful review and aggregate them. 

546 # Seed with the synthetic injection finding/warnings so they surface in the 

547 # report and outcome.warnings without ever influencing agent behaviour. 

548 all_findings: list[Finding] = list(injection_findings) 

549 all_warnings: list[str] = injection.hits_to_warnings(injection_hits) 

550 all_warnings.extend(privilege_warnings) 

551 for r in reviews: 

552 if not r.ok: 

553 continue 

554 found, warns = parse_findings(r.output, r.agent) 

555 r.findings = found 

556 r.warnings = warns 

557 r.structured = emitted_findings_block(r.output) 

558 all_findings.extend(found) 

559 all_warnings.extend(warns) 

560 

561 # Stream round-1 reviews as they're now finalized (stable order). 

562 for r in reviews: 

563 emit("review", r) 

564 

565 # Deterministic consensus grouping across reviewers. 

566 groups = group_findings(all_findings, len(reviews)) 

567 

568 # Names of agents whose round-1 review succeeded — the chair resolver uses 

569 # this to (optionally) prefer a non-reviewer chair (#38). 

570 reviewer_names = [r.agent for r in reviews if r.ok] 

571 

572 # Resolve the chair ONCE for the whole run so verify and synthesis use the 

573 # SAME chair. ``chair = "rotate"`` and prefer-non-reviewer both consume the 

574 # shared run RNG / reviewer info, so resolving once (rather than recomputing 

575 # per phase) is what keeps a rotating chair stable within a run (#38). 

576 # Escalation (#714): a critical or major finding after round 1 brings the 

577 # benched frontier seats into the debate and draws the chair from the 

578 # frontier seats. Decided once, here, and recorded on the plan — including 

579 # what it can actually do in THIS run: a single-round run has no debate to 

580 # join, and ``--auto`` sets exactly that on the routine band that benched 

581 # the seats, so the record names the chair-only case rather than promising 

582 # a debate that will not happen (review round 1). 

583 round1_debaters = [a for a in round1 if any(r.agent == a.name and r.ok for r in reviews)] 

584 if benched: 

585 plan.escalated, plan.escalation_reason = routing.should_escalate(groups) 

586 if plan.escalated: 

587 log(f"tiered routing: escalating — {plan.escalation_reason}") 

588 chair_pool = [a.name for a in round1] 

589 if plan.escalated: 

590 chair_pool = routing.frontier_names(specs, usable_names) or chair_pool 

591 agent_order = agent_order + [a.name for a in benched] 

592 chair_name = resolve_chair(config, chair_pool, reviewer_names, run_rng) 

593 

594 # Round 2+: debate. Only agents whose round-1 review succeeded participate. 

595 # Two modes (issue #40): 

596 # - fixed (early_stop = false): honour ``rounds`` exactly — run one debate 

597 # round iff rounds >= 2. Reproducible fixed-N behaviour for benchmarking. 

598 # - adaptive (early_stop = true): skip the debate when round-1 reviewers 

599 # already agree, otherwise run debate up to ``max_rounds`` rounds and stop 

600 # as soon as a round resolves all disputes. 

601 debate: list[AgentResult] = [] 

602 rounds_executed = 1 

603 stop_reason = "" 

604 budget_exhausted = False 

605 debaters = round1_debaters 

606 if plan.escalated: 

607 # Benched frontier seats cross-examine: they receive every round-1 

608 # review as "the others" and answer in the debate format. 

609 debaters = debaters + benched 

610 can_debate = len(debaters) >= 2 

611 

612 if config.early_stop: 

613 max_rounds = config.effective_max_rounds 

614 if not can_debate: 

615 stop_reason = "stopped after round 1: need >=2 successful reviews to debate" 

616 log(stop_reason) 

617 elif max_rounds < 2: 

618 stop_reason = "stopped after round 1: max_rounds < 2" 

619 log(stop_reason) 

620 else: 

621 converged, why = convergence.review_convergence(groups, len(reviews)) 

622 if converged: 

623 stop_reason = f"early stop after round 1: {why}" 

624 log(stop_reason) 

625 else: 

626 log(f"early stop active: {why}; running debate up to {max_rounds} round(s)") 

627 prior: list[AgentResult] = [] 

628 round_no = 1 

629 while round_no < max_rounds: 

630 if budget.expired(): 

631 budget_exhausted = True 

632 stop_reason = f"stopped at round {rounds_executed}: run budget exhausted" 

633 log(stop_reason) 

634 break 

635 round_no += 1 

636 debate = _debate_round( 

637 debaters, 

638 reviews, 

639 diff, 

640 config, 

641 run_rng, 

642 agent_order, 

643 prior, 

644 budget, 

645 retries, 

646 log, 

647 round_no, 

648 template=tmpl["debate"], 

649 ) 

650 rounds_executed = round_no 

651 for r in debate: 

652 emit("debate", r, round_no) 

653 dconv, dwhy = convergence.debate_convergence(debate) 

654 if dconv: 

655 stop_reason = f"converged after round {round_no}: {dwhy}" 

656 log(stop_reason) 

657 break 

658 prior = debate 

659 stop_reason = f"ran {round_no} rounds: {dwhy}" 

660 else: 

661 stop_reason = stop_reason or ( 

662 f"reached max_rounds ({max_rounds}) with disagreement remaining" 

663 ) 

664 else: 

665 # Fixed-N: exactly the historical behaviour. 

666 if config.rounds >= 2 and can_debate: 

667 if budget.expired(): 

668 budget_exhausted = True 

669 stop_reason = "round 2 skipped: run budget exhausted" 

670 log(stop_reason) 

671 else: 

672 debate = _debate_round( 

673 debaters, 

674 reviews, 

675 diff, 

676 config, 

677 run_rng, 

678 agent_order, 

679 [], 

680 budget, 

681 retries, 

682 log, 

683 2, 

684 template=tmpl["debate"], 

685 ) 

686 rounds_executed = 2 

687 for r in debate: 

688 emit("debate", r, 2) 

689 elif config.rounds >= 2: 

690 stop_reason = "round 2 skipped: need >=2 successful reviews to debate" 

691 log(stop_reason) 

692 else: 

693 stop_reason = "single round (rounds = 1)" 

694 

695 # Verification: the chair judges candidate findings to reduce false 

696 # positives. Skipped when the run budget is exhausted (issue #30) so a 

697 # partial run still returns what completed instead of overrunning. 

698 verify_result: AgentResult | None = None 

699 verdicts: list[Verdict] = [] 

700 if config.verify: 

701 if budget.expired(): 

702 budget_exhausted = True 

703 msg = "verification skipped: run budget exhausted" 

704 log(msg) 

705 all_warnings.append(msg) 

706 else: 

707 verify_result, verdicts, verify_warnings = _verify( 

708 chair_name, 

709 usable, 

710 all_findings, 

711 diff, 

712 context, 

713 budget, 

714 retries, 

715 log, 

716 template=tmpl["verify"], 

717 ) 

718 all_warnings.extend(verify_warnings) 

719 _apply_verdicts(groups, verdicts) 

720 if verify_result is not None: 

721 emit("verify", verify_result) 

722 

723 # Local-only demotion (issue #442) runs AFTER verification, never before: 

724 # _reject_targets' member-tier guard (orchestrator._reject_targets) assumes 

725 # group.severity == max(member severities) to decide whether a rejecting 

726 # verdict may suppress the whole group. Demoting group.severity earlier 

727 # would desync it from that invariant and could let a verdict aimed at a 

728 # minor local-only duplicate collateral-reject a genuinely critical, 

729 # never-verified co-located finding merged into the same group. 

730 if config.demote_local_only: 

731 vendor_by_reviewer = {a.name: a.vendor for a in config.agents} 

732 demote_local_only_groups(groups, vendor_by_reviewer) 

733 

734 # Synthesis: the chair consolidates. When the resolved chair is ALSO a 

735 # round-1 reviewer, feed it an anonymized view of the reviews (#38 guardrail) 

736 # so it cannot preferentially weight its own findings; the report still 

737 # attributes by real name because it renders the real outcome data, not this 

738 # synthesis prompt. 

739 synthesis: AgentResult | None = None 

740 if budget.expired(): 

741 budget_exhausted = True 

742 msg = "synthesis skipped: run budget exhausted" 

743 log(msg) 

744 if msg not in all_warnings: 744 ↛ 766line 744 didn't jump to line 766 because the condition on line 744 was always true

745 all_warnings.append(msg) 

746 else: 

747 chair_is_reviewer = chair_name in reviewer_names 

748 anonymize_synthesis = config.anonymize_debate and chair_is_reviewer 

749 synthesis = _synthesize( 

750 chair_name, 

751 usable, 

752 reviews, 

753 debate, 

754 diff, 

755 budget, 

756 retries, 

757 log, 

758 verdicts=verdicts, 

759 anonymize_reviews=anonymize_synthesis, 

760 rng=run_rng, 

761 template=tmpl["synthesis"], 

762 ) 

763 if synthesis is not None: 

764 emit("synthesis", synthesis) 

765 

766 if plan.escalated: 

767 # What escalation DID, read off the run (#714, review rounds 1 and 2). 

768 # Predicting it was wrong twice: a single-round run has no debate to 

769 # join — `--auto` sets exactly one round on the `low` band that benched 

770 # the seats — and an adaptive run can converge after round 1 and skip 

771 # the debate it was going to have. The benched seats that actually 

772 # produced a debate result are the answer, whatever the reason. 

773 bench_names = {a.name for a in benched} 

774 # `ok` matters: a benched seat whose debate call failed produced a row 

775 # but no cross-examination, and counting it would put the frontier seat 

776 # in the record for work it did not do (review round 3). 

777 joined = [r.agent for r in debate if r.agent in bench_names and r.ok] 

778 effect = routing.escalation_effect(joined, debate_ran=bool(debate)) 

779 plan.escalation_reason = f"{plan.escalation_reason}; {effect}" 

780 log(f"tiered routing: {plan.escalation_reason}") 

781 

782 return JuryOutcome( 

783 reviews=reviews, 

784 debate=debate, 

785 synthesis=synthesis, 

786 chair=chair_name, 

787 findings=all_findings, 

788 warnings=all_warnings, 

789 groups=groups, 

790 verify=verify_result, 

791 verdicts=verdicts, 

792 context_mode=context_mode, 

793 redact_secrets=redact_on, 

794 redaction_count=redaction_count, 

795 injection_hits=injection_hits, 

796 skipped=skipped, 

797 budget_exhausted=budget_exhausted, 

798 rounds_executed=rounds_executed, 

799 stop_reason=stop_reason, 

800 changed=changed, 

801 routing=plan.as_dict(), 

802 ) 

803 

804 

805def resolve_chair( 

806 config: JuryConfig, 

807 usable: list[str], 

808 reviewers: list[str], 

809 rng: random.Random, 

810) -> str: 

811 """Resolve the chair for a run as a PURE function of its inputs (#38). 

812 

813 Precedence: 

814 1. ``chair = "rotate"`` — pick deterministically from the usable agents 

815 using the shared run ``rng``. Same seed -> same chair; different seeds 

816 may differ. Falls back to the first usable agent when none are usable. 

817 2. An explicit ``config.chair`` that names a usable agent — honoured as-is 

818 (an operator-chosen chair always wins). 

819 3. ``prefer_non_reviewer_chair`` — when set and a usable agent that was NOT 

820 a successful round-1 reviewer exists, prefer the first such agent 

821 (neutral chair). This only applies when the configured chair is not 

822 itself a usable agent. 

823 4. Fallback to the first usable agent (legacy behaviour). 

824 

825 Keeping this pure (no Adapter objects, no I/O) makes it directly 

826 unit-testable and guarantees ``_verify`` and ``_synthesize`` agree because 

827 the caller resolves it ONCE and threads the result through both. 

828 """ 

829 if not usable: 

830 return config.chair 

831 names = set(usable) 

832 

833 if config.chair == "rotate": 

834 # Deterministic rotation: sort for a stable candidate order independent 

835 # of dict/thread ordering, then index with the shared run RNG. Sorting 

836 # the candidate list (not iterating the set) makes the pick a pure 

837 # function of (seed, usable-name set): same seed + same agents -> same 

838 # chair, regardless of RNG-consumption order elsewhere. 

839 candidates = sorted(names) 

840 return candidates[rng.randrange(len(candidates))] 

841 

842 if config.chair in names: 

843 return config.chair 

844 

845 if config.prefer_non_reviewer_chair: 

846 reviewer_set = set(reviewers) 

847 non_reviewers = [n for n in usable if n not in reviewer_set] 

848 if non_reviewers: 

849 return non_reviewers[0] 

850 

851 return usable[0] 

852 

853 

854def _format_findings_for_verify(findings: list[Finding]) -> str: 

855 """Render candidate findings for the chair's verification prompt. 

856 

857 Reviewer identity is omitted (#250) so the chair can't favour its own 

858 findings while judging them — parity with the #37/#38 anonymization. 

859 Verdicts match back by file/line/claim, so dropping it is safe. 

860 """ 

861 if not findings: 

862 return "_(no candidate findings)_" 

863 lines = [] 

864 for f in findings: 

865 loc = f.file or "?" 

866 if f.line is not None: 

867 loc = f"{loc}:{f.line}" 

868 lines.append(f"- [{f.severity}] {loc} — {f.claim}") 

869 return "\n".join(lines) 

870 

871 

872def _format_verdicts(verdicts: list[Verdict]) -> str: 

873 if not verdicts: 

874 return "_(no verification verdicts)_" 

875 lines = [] 

876 for v in verdicts: 

877 loc = v.file or "?" 

878 if v.line is not None: 

879 loc = f"{loc}:{v.line}" 

880 lines.append(f"- [{v.status}] {loc} — {v.claim}: {v.reasoning}") 

881 return "\n".join(lines) 

882 

883 

884def _verify( 

885 chair_name, 

886 usable, 

887 findings, 

888 diff, 

889 context, 

890 budget, 

891 retries, 

892 log, 

893 template=prompts.VERIFY, 

894) -> tuple[AgentResult | None, list[Verdict], list[str]]: 

895 chair = next((a for a in usable if a.name == chair_name), None) 

896 if chair is None: 

897 return None, [], [] 

898 log(f"verification: chair '{chair_name}' judging {len(findings)} candidate findings") 

899 prompt = template.format( 

900 diff=prompts.neutralize_sentinels(diff), 

901 findings=prompts.neutralize_sentinels(_format_findings_for_verify(findings)), 

902 context=prompts.neutralize_sentinels(context or "_(none)_"), 

903 notice=prompts._UNTRUSTED_NOTICE, 

904 ) 

905 result = _run_with_retry(chair, prompt, "verify", budget, retries, log) 

906 if not result.ok: 

907 return result, [], [f"verification failed: {result.error}"] 

908 verdicts, warnings = parse_verdicts(result.output, chair_name) 

909 return result, verdicts, warnings 

910 

911 

912def _verdict_matches_group(verdict: Verdict, group: FindingGroup) -> bool: 

913 from .consensus import _normalize_claim, _normalize_path 

914 

915 rep = group.representative 

916 # Case-EXACT path match (fold_case=False): on a case-sensitive filesystem 

917 # ``Config.py`` != ``config.py``, so a verdict must not reject a finding it 

918 # only case-collapses onto (audit 2026-06-13 r6/M). 

919 if _normalize_path(verdict.file, fold_case=False) != _normalize_path(rep.file, fold_case=False): 

920 return False 

921 if verdict.line is not None and rep.line is not None and abs(verdict.line - rep.line) > 3: 

922 return False 

923 v_claim = _normalize_claim(verdict.claim) 

924 r_claim = _normalize_claim(rep.claim) 

925 if not v_claim: 

926 # An empty verdict claim is allowed to match the finding *at this 

927 # location* (the verifier may omit the claim and refer to it by 

928 # position). But a verdict with NEITHER a claim NOR a line has no 

929 # location precision at all: it would otherwise match — and, when 

930 # ``unsupported``, REJECT — every finding group in the file, including 

931 # unrelated criticals, flipping the CI gate from FAIL to PASS. Such a 

932 # claim-less, line-less verdict is a file-wide wildcard and must not 

933 # match (security audit 2026-06-13 r6/M). Require a concrete line that 

934 # actually pins the finding before honoring an empty-claim match. 

935 return verdict.line is not None and rep.line is not None 

936 if v_claim == r_claim: 

937 return True 

938 v_tokens, r_tokens = set(v_claim.split()), set(r_claim.split()) 

939 if not v_tokens or not r_tokens: 

940 return False 

941 inter = len(v_tokens & r_tokens) 

942 union = len(v_tokens) + len(r_tokens) - inter 

943 return (inter / union if union else 0.0) >= 0.5 

944 

945 

946def _claim_sim(a_claim: str, b_claim: str) -> float: 

947 """Token-set similarity between two claims: 1.0 exact, else Jaccard, 0.0 if 

948 either side is empty.""" 

949 from .consensus import _normalize_claim 

950 

951 a = _normalize_claim(a_claim) 

952 b = _normalize_claim(b_claim) 

953 if not a or not b: 

954 return 0.0 

955 if a == b: 

956 return 1.0 

957 at, bt = set(a.split()), set(b.split()) 

958 inter = len(at & bt) 

959 union = len(at) + len(bt) - inter 

960 return (inter / union) if union else 0.0 

961 

962 

963# Verdict statuses that move a finding into a non-blocking bucket (suppress it). 

964_REJECTING_STATUSES = frozenset({"unsupported", "needs_human_decision"}) 

965# Apply most-blocking statuses first so a contradictory verdict pair on one 

966# finding is fail-closed: a `verified` (blocking) judgement is recorded before 

967# any `unsupported`/`needs_human_decision` and cannot then be flipped to 

968# non-blocking by verdict array ordering (audit 2026-06-13 r8/M). 

969_STATUS_PRIORITY = {"verified": 0, "needs_human_decision": 1, "unsupported": 2} 

970# Minimum claim similarity for a verdict to be considered "about" a finding at 

971# all. The PRIMARY defence against a verdict dismissing a co-located *distinct* 

972# finding is that a rejection attaches to AT MOST the single best-matching group 

973# (`_best_reject_target`), so a verdict whose claim copies a benign neighbour 

974# routes to that neighbour, not to the co-located critical (audit 2026-06-13 

975# r8/M). The threshold stays moderate so the verifier's legitimate paraphrased 

976# rejections (it drops the reviewer-name prefix etc.) still apply. 

977_REJECT_CLAIM_THRESHOLD = 0.5 

978 

979 

980def _reject_targets(verdict: Verdict, groups: list[FindingGroup]) -> list[FindingGroup]: 

981 """Un-statused groups a rejecting verdict may suppress (fail-closed, r7/r8). 

982 

983 Defences (each closes a distinct collateral-rejection vector found across 

984 audit rounds 6-9): 

985 

986 0. **Line required.** A rejecting verdict must pin a concrete line. A 

987 line-less verdict is too imprecise to safely suppress a finding and would 

988 act as a file-wide-by-claim wildcard (audit r9/M, the claim-ful 

989 counterpart of the round-6 line-less-wildcard fix). 

990 1. **Member-tier guard.** A group may merge findings of different severities 

991 (consensus keeps the max). A verdict is "about" the member whose claim it 

992 best matches; if that member is *less severe* than the group's max, the 

993 verdict is dismissing a lesser co-located finding and must NOT suppress 

994 the (e.g. critical) group. 

995 2. **Best-tier only.** Across candidate groups, suppress only those at the 

996 highest match similarity — a verdict copying a benign neighbour rejects 

997 that neighbour (and its duplicate phrasings, which tie) but not a 

998 separate, less-similar critical group. 

999 3. **Least-severe within a tie.** If the best-similarity tier still spans 

1000 severities (an exact `_claim_sim` tie between a critical and a benign 

1001 decoy), suppress only the *least*-severe groups — a tie must never drag a 

1002 critical down alongside a decoy (audit r9/M). 

1003 """ 

1004 from .findings import SEVERITY_ORDER 

1005 

1006 if verdict.line is None: 

1007 return [] 

1008 scored: list[tuple[float, FindingGroup]] = [] 

1009 for group in groups: 

1010 if group.status: 

1011 continue 

1012 if not _verdict_matches_group(verdict, group): 

1013 continue 

1014 members = getattr(group, "members", None) or [group.representative] 

1015 best_sim, best_member = max( 

1016 ((_claim_sim(verdict.claim, m.claim), m) for m in members), 

1017 key=lambda t: t[0], 

1018 ) 

1019 if best_sim < _REJECT_CLAIM_THRESHOLD: 

1020 continue 

1021 # Member-tier guard: refuse if the verdict best-names a member less 

1022 # severe than the group's max severity (lower rank = more severe). 

1023 if SEVERITY_ORDER.get(best_member.severity, 99) > SEVERITY_ORDER.get(group.severity, 99): 

1024 continue 

1025 scored.append((best_sim, group)) 

1026 if not scored: 

1027 return [] 

1028 best = max(sim for sim, _ in scored) 

1029 tier = [(sim, group) for sim, group in scored if sim >= best] 

1030 # Within the top-similarity tier, keep only the least-severe groups. 

1031 least_rank = max(SEVERITY_ORDER.get(g.severity, 99) for _, g in tier) 

1032 return [g for _, g in tier if SEVERITY_ORDER.get(g.severity, 99) == least_rank] 

1033 

1034 

1035def _apply_verdicts(groups: list[FindingGroup], verdicts: list[Verdict]) -> None: 

1036 """Attach verification statuses to consensus groups. 

1037 

1038 unsupported -> bucket 'rejected'; needs_human_decision -> bucket 'disputed'; 

1039 verified -> status recorded, bucket unchanged. 

1040 """ 

1041 # Stable-sort by blocking priority so contradictions resolve fail-closed. 

1042 for verdict in sorted(verdicts, key=lambda v: _STATUS_PRIORITY.get(v.status, 3)): 

1043 if verdict.status in _REJECTING_STATUSES: 

1044 # Suppress only the best-similarity tier this verdict names — never 

1045 # collaterally a co-located, less-similar distinct finding. 

1046 bucket = "rejected" if verdict.status == "unsupported" else "disputed" 

1047 for target in _reject_targets(verdict, groups): 

1048 target.status = verdict.status 

1049 target.status_reasoning = verdict.reasoning 

1050 target.bucket = bucket 

1051 continue 

1052 # A verifying (non-suppressing) verdict may attach to every matching 

1053 # group: when reviewers phrase the same issue differently it can land in 

1054 # more than one group, and all should carry the judgement. 

1055 for group in groups: 

1056 if group.status: 

1057 continue 

1058 if _verdict_matches_group(verdict, group): 

1059 group.status = verdict.status 

1060 group.status_reasoning = verdict.reasoning 

1061 

1062 

1063def _synthesize( 

1064 chair_name, 

1065 usable, 

1066 reviews, 

1067 debate, 

1068 diff, 

1069 budget, 

1070 retries, 

1071 log, 

1072 verdicts=None, 

1073 anonymize_reviews=False, 

1074 rng=None, 

1075 template=prompts.SYNTHESIS, 

1076) -> AgentResult | None: 

1077 chair = next((a for a in usable if a.name == chair_name), None) 

1078 if chair is None: 

1079 return None 

1080 log(f"synthesis: chair '{chair_name}' consolidating verdict") 

1081 if anonymize_reviews: 

1082 # Chair self-preference guardrail (#38): present round-1 reviews to the 

1083 # chair under anonymous labels (no agent/vendor identity, no stable 

1084 # order) so it cannot tell which review is "its own". Uses the shared run 

1085 # RNG for deterministic-but-unstable ordering. ``me=None`` keeps ALL 

1086 # reviews (we are not excluding a debater here, only stripping identity). 

1087 peer_rng = random.Random(rng.random()) if rng is not None else random.Random() 

1088 reviews_txt, _label_map = _anonymize_peers(reviews, None, peer_rng) 

1089 else: 

1090 reviews_txt = ( 

1091 "\n\n".join( 

1092 f"### {r.agent} ({r.vendor})\n{r.output}" for r in reviews if r.ok and r.output 

1093 ) 

1094 or "_(no reviews)_" 

1095 ) 

1096 debate_txt = ( 

1097 "\n\n".join(f"### {r.agent}\n{r.output}" for r in debate if r.ok and r.output) 

1098 or "_(no debate round)_" 

1099 ) 

1100 prompt = template.format( 

1101 diff=prompts.neutralize_sentinels(diff), 

1102 reviews=prompts.neutralize_sentinels(reviews_txt), 

1103 debate=prompts.neutralize_sentinels(debate_txt), 

1104 notice=prompts._UNTRUSTED_NOTICE, 

1105 ) 

1106 if verdicts: 

1107 # The verdicts quote candidate findings, which transitively quote 

1108 # untrusted diff text (issue v1.5.0/M-1: this addendum was the one slot 

1109 # the #316/L-1 fix missed). Fence + neutralize it like every other 

1110 # peer-output slot so an embedded closing token can't break out. 

1111 prompt += ( 

1112 "\n\n=== VERIFICATION VERDICTS (may quote UNTRUSTED text) ===\n" 

1113 "<<<UNTRUSTED_FINDINGS\n" 

1114 + prompts.neutralize_sentinels(_format_verdicts(verdicts)) 

1115 + "\nUNTRUSTED_FINDINGS>>>\n" 

1116 ) 

1117 return _run_with_retry(chair, prompt, "synthesis", budget, retries, log) 

1118 

1119 

1120def _first_sent_model(parts: list[AgentResult]) -> str: 

1121 """The first model id *parts* recorded sending (issue #722). 

1122 

1123 Every chunk of one diff is one seat invoked repeatedly through one adapter, 

1124 so the parts normally all carry the same stamped id (#709) and "the first 

1125 non-empty one" is simply "the one". Where they disagree — which only an 

1126 adapter that consults a live model listing between invocations can produce 

1127 — there is no id that is true of the whole merge, and inventing one 

1128 ("mixed") would be a string no invocation sent. The first is a model this 

1129 run really did send, and it is the one whose output opens the merged body, 

1130 so the ballot's id and the text a reader checks it against come from the 

1131 same invocation. 

1132 

1133 **Callers pass the parts whose output is in the merged body**, because that 

1134 is what the id has to be true of. A failed part is not excluded by carrying 

1135 ``""``: :func:`_run_with_retry` stamps ``result.model`` from 

1136 ``adapter.resolved_model()`` after ``adapter.run`` returns, whether the run 

1137 succeeded or not, so a chunk that failed *does* normally record an id — and 

1138 it can be a different one from its siblings', since the fallback an adapter 

1139 performs against a live listing is exactly the kind of failure that changes 

1140 it. Scanning all the parts would then let a chunk that contributed no text 

1141 name the model for text it did not produce. 

1142 

1143 A part carries ``""`` only when nothing stamped it: a hand-assembled 

1144 outcome, a pre-#709 record, or a failure raised before the adapter returned. 

1145 A merge in which no scanned part recorded an id stays empty, which 

1146 :func:`ai_jury.ballots.describe_model` labels ``recomputed`` rather than 

1147 ``requested``. That fallback is unchanged. 

1148 """ 

1149 for p in parts: 

1150 model = (getattr(p, "model", "") or "").strip() 

1151 if model: 

1152 return model 

1153 return "" 

1154 

1155 

1156def _merge_results_by_agent(phase_lists: list[list[AgentResult]]) -> list[AgentResult]: 

1157 """Merge per-chunk results for the same agent into one result (issue #31). 

1158 

1159 Outputs are concatenated under per-chunk headers, durations summed, ``ok`` is 

1160 true if the agent succeeded on any chunk, and ``attempts`` keeps the max so a 

1161 retried chunk is still visible. Agent order follows first appearance. 

1162 

1163 The model id the invocation sent (:attr:`AgentResult.model`, #709) rides 

1164 along too — see :func:`_first_sent_model`. Dropping it here made a chunked 

1165 review's provenance strictly weaker than a single-chunk one's (#722). It is 

1166 read from the same parts the body is built from: the id has to be one that 

1167 produced the text under it, and a chunk can fail *after* recording a 

1168 different id than its siblings (an adapter that fell back against a live 

1169 listing). Only when no part contributed body text at all — a seat that 

1170 failed on every chunk — is the whole set scanned, so a failed seat still 

1171 reports what it sent. 

1172 """ 

1173 order: list[str] = [] 

1174 by_agent: dict[str, list[AgentResult]] = {} 

1175 for lst in phase_lists: 

1176 for r in lst: 

1177 if r.agent not in by_agent: 

1178 by_agent[r.agent] = [] 

1179 order.append(r.agent) 

1180 by_agent[r.agent].append(r) 

1181 

1182 merged: list[AgentResult] = [] 

1183 for name in order: 

1184 parts = by_agent[name] 

1185 

1186 # bolt: Consolidate multiple metrics (ok, body, total_duration, max_attempts) 

1187 # into a single-pass O(N) explicit loop to bypass multiple generator instantiations 

1188 ok = False 

1189 body_parts = [] 

1190 body_sources: list[AgentResult] = [] 

1191 first_err = None 

1192 total_duration = 0.0 

1193 max_attempts = 0 

1194 

1195 for i, p in enumerate(parts, 1): 

1196 if p.ok: 

1197 ok = True 

1198 if p.output: 

1199 body_parts.append(f"#### chunk {i}\n{p.output}") 

1200 body_sources.append(p) 

1201 elif first_err is None: 

1202 first_err = p 

1203 

1204 total_duration += p.duration_s 

1205 if p.attempts > max_attempts: 

1206 max_attempts = p.attempts 

1207 

1208 body = "\n\n".join(body_parts) 

1209 

1210 merged.append( 

1211 AgentResult( 

1212 name, 

1213 parts[0].vendor, 

1214 ok, 

1215 body, 

1216 round(total_duration, 3), 

1217 error=None if ok else (first_err.error if first_err else None), 

1218 error_code=None if ok else (first_err.error_code if first_err else None), 

1219 attempts=max_attempts, 

1220 model=_first_sent_model(body_sources or parts), 

1221 ) 

1222 ) 

1223 return merged 

1224 

1225 

1226def _combine_chair_results(results: list[AgentResult], chair: str) -> AgentResult | None: 

1227 """Combine per-chunk chair results (verify/synthesis) into one labelled result. 

1228 

1229 The combined record keeps the model id the chair's invocation sent (#722), 

1230 taken from the same population as ``vendor`` — the parts whose output is in 

1231 the body — by :func:`_first_sent_model`. 

1232 """ 

1233 ok_parts = [r for r in results if r.ok and r.output] 

1234 if not ok_parts: 

1235 return results[0] if results else None 

1236 vendor = ok_parts[0].vendor 

1237 

1238 # bolt: Consolidate body text concatenation and duration sum into a single-pass O(N) loop 

1239 body_parts = [] 

1240 total_duration = 0.0 

1241 for i, r in enumerate(ok_parts, 1): 

1242 body_parts.append(f"### chunk {i}\n{r.output}") 

1243 total_duration += r.duration_s 

1244 

1245 body = "\n\n".join(body_parts) 

1246 return AgentResult( 

1247 chair, vendor, True, body, round(total_duration, 3), model=_first_sent_model(ok_parts) 

1248 ) 

1249 

1250 

1251def _merge_chunk_outcomes(outcomes: list[JuryOutcome], config: JuryConfig) -> JuryOutcome: 

1252 """Fold per-chunk outcomes (issue #31) into one renderable JuryOutcome. 

1253 

1254 Findings are unioned and re-grouped across all chunks so the consensus view 

1255 is global; verdicts are re-applied to the merged groups. Review/debate/chair 

1256 outputs are merged per agent with chunk labels so the report stays coherent. 

1257 """ 

1258 if len(outcomes) == 1: 

1259 return outcomes[0] 

1260 base = outcomes[0] 

1261 

1262 reviews = _merge_results_by_agent([o.reviews for o in outcomes]) 

1263 debate = ( 

1264 _merge_results_by_agent([o.debate for o in outcomes]) 

1265 if any(o.debate for o in outcomes) 

1266 else [] 

1267 ) 

1268 findings = [f for o in outcomes for f in o.findings] 

1269 groups = group_findings(findings, len(reviews)) 

1270 verdicts = [v for o in outcomes for v in o.verdicts] 

1271 # Scope each chunk's verdicts to that chunk's own findings. A verdict is 

1272 # produced while verifying ONE chunk (whose prompt held only that chunk's 

1273 # findings, but whose attacker-controlled diff text could steer it); after 

1274 # the global merge an unscoped verdict could reject a *different* chunk's 

1275 # structured critical and flip the CI gate (audit 2026-06-13 r7/M). Chunks 

1276 # are file-disjoint, so apply each chunk's verdicts only to groups whose 

1277 # location is one of that chunk's files. 

1278 from .consensus import _normalize_path 

1279 

1280 for o in outcomes: 

1281 chunk_files = {_normalize_path(f.file, fold_case=False) for f in o.findings if f.file} 

1282 chunk_groups = [ 

1283 g 

1284 for g in groups 

1285 if _normalize_path(g.representative.file, fold_case=False) in chunk_files 

1286 ] 

1287 _apply_verdicts(chunk_groups, o.verdicts) 

1288 

1289 # Runs AFTER verdicts are applied — see the matching comment in run_jury for 

1290 # why (the _reject_targets member-tier guard assumes group.severity is the 

1291 # true max member severity; demoting earlier would desync that invariant). 

1292 if config.demote_local_only: 

1293 vendor_by_reviewer = {a.name: a.vendor for a in config.agents} 

1294 demote_local_only_groups(groups, vendor_by_reviewer) 

1295 

1296 warnings = [w for o in outcomes for w in o.warnings] 

1297 

1298 # ONE chair name for the whole merged record (#714, r5). A chunked tiered 

1299 # run escalates per chunk, so the chunks can have different chairs; the run 

1300 # publishes the chair of the first chunk that escalated, and the combined 

1301 # synthesis and verify must carry that same name — labelling them with the 

1302 # quiet first chunk's seat while the outcome names another is a report that 

1303 # says one seat synthesised a body another seat's chunk opens. 

1304 merged_chair = next( 

1305 (o.chair for o in outcomes if (o.routing or {}).get("escalated")), base.chair 

1306 ) 

1307 synthesis = _combine_chair_results([o.synthesis for o in outcomes if o.synthesis], merged_chair) 

1308 verify = _combine_chair_results([o.verify for o in outcomes if o.verify], merged_chair) 

1309 

1310 # bolt: Consolidate collection aggregations (sum, extend, max, any) into a single-pass O(N) explicit loop 

1311 redaction_count = 0 

1312 injection_hits = [] 

1313 budget_exhausted = False 

1314 rounds_executed = 0 

1315 

1316 for o in outcomes: 

1317 redaction_count += o.redaction_count 

1318 injection_hits.extend(o.injection_hits) 

1319 if o.budget_exhausted: 

1320 budget_exhausted = True 

1321 if o.rounds_executed > rounds_executed: 

1322 rounds_executed = o.rounds_executed 

1323 

1324 return JuryOutcome( 

1325 reviews=reviews, 

1326 debate=debate, 

1327 synthesis=synthesis, 

1328 # A chunked run escalates per chunk, so a quiet first chunk must not 

1329 # publish its economical chair for a run that escalated later: the 

1330 # chair of the first chunk that escalated is the run's (#714, r4), and 

1331 # the combined synthesis and verify above carry the same name (r5). 

1332 chair=merged_chair, 

1333 findings=findings, 

1334 warnings=warnings, 

1335 groups=groups, 

1336 verify=verify, 

1337 verdicts=verdicts, 

1338 context_mode=base.context_mode, 

1339 redact_secrets=base.redact_secrets, 

1340 redaction_count=redaction_count, 

1341 injection_hits=injection_hits, 

1342 skipped=base.skipped, 

1343 budget_exhausted=budget_exhausted, 

1344 rounds_executed=rounds_executed, 

1345 stop_reason=f"chunked review across {len(outcomes)} part(s)", 

1346 # The whole change, not one chunk of it (#710): a reviewer that named a 

1347 # file from another chunk named a file in this change. 

1348 changed=largediff.merge_change_indexes([o.changed for o in outcomes]), 

1349 # Every chunk routed off the same band and the same bench, so the first 

1350 # chunk's plan describes the run; escalation is true when ANY chunk 

1351 # escalated, because the benched seats really did review by then. 

1352 routing=_merged_routing(outcomes, base, debate), 

1353 ) 

1354 

1355 

1356def _merged_routing(outcomes: list[JuryOutcome], base: JuryOutcome, debate: list) -> dict: 

1357 """The routing record of a chunked run, read off the merged run (#714, r3). 

1358 

1359 Every chunk routed off the same band and the same bench, so the first 

1360 chunk's plan describes the panel. Escalation is per chunk — round 1 of one 

1361 chunk can carry a major finding while another is quiet — so the record says 

1362 on how many it escalated, and what escalation then did is recomputed from 

1363 the **merged** debate rather than copied from whichever chunk escalated 

1364 first. Copying it was wrong the way predicting was wrong: the sentence 

1365 described one chunk and was published as the run's. 

1366 """ 

1367 escalated = [o for o in outcomes if (o.routing or {}).get("escalated")] 

1368 merged = {**base.routing, "escalated": bool(escalated)} 

1369 if escalated: 

1370 bench = set(base.routing.get("benched") or []) 

1371 joined = [r.agent for r in debate if r.agent in bench and r.ok] 

1372 merged["escalation_reason"] = ( 

1373 f"escalated on {len(escalated)} of {len(outcomes)} chunk(s); " 

1374 f"{routing.escalation_effect(joined, debate_ran=bool(debate))}" 

1375 ) 

1376 return merged 

1377 

1378 

1379def plan_for(config: JuryConfig, diff: str) -> largediff.DiffPlan: 

1380 """The diff plan *config* selects for *diff* (pure). 

1381 

1382 The one place that maps ``[jury.diff]`` onto :func:`largediff.plan_diff`, so 

1383 every caller that needs to know which files the panel will actually be shown 

1384 — :func:`review_diff` below, and the CLI's static-hints pre-pass (#737) — 

1385 reads the same ``kept`` list off the same filters. ``plan_diff`` is pure, so 

1386 planning the same diff twice returns the same answer. 

1387 """ 

1388 dc = config.diff 

1389 return largediff.plan_diff( 

1390 diff, 

1391 max_bytes=dc.max_bytes, 

1392 chunk=dc.chunk, 

1393 chunk_max_bytes=dc.chunk_max_bytes, 

1394 exclude_generated=dc.exclude_generated, 

1395 exclude=dc.exclude, 

1396 include=dc.include, 

1397 ) 

1398 

1399 

1400def review_diff( 

1401 config: JuryConfig, 

1402 diff: str, 

1403 *, 

1404 context: str = "", 

1405 hints: str = "", 

1406 mock: bool = False, 

1407 strict: bool = False, 

1408 seed: int | None = None, 

1409 policy: ReviewPolicy | None = None, 

1410 log=lambda _msg: None, 

1411 on_event=None, 

1412) -> tuple[JuryOutcome, largediff.DiffPlan]: 

1413 """Plan a diff (filter + size + mode) then run the jury (issue #31). 

1414 

1415 The single entry point the CLI uses: it measures and filters the diff, 

1416 reports the size and the selected handling mode, and dispatches: 

1417 

1418 - ``full`` — review the filtered diff in one ``run_jury`` pass; 

1419 - ``chunked`` — review each chunk and merge the outcomes; 

1420 - ``too_large`` — raise ``RuntimeError`` with an actionable message. 

1421 

1422 Returns ``(outcome, plan)`` so the caller can surface the plan. Existing 

1423 callers of :func:`run_jury` are unaffected. 

1424 """ 

1425 plan = plan_for(config, diff) 

1426 log( 

1427 f"diff size: {plan.total_bytes} B total, {plan.kept_bytes} B after filters " 

1428 f"({len(plan.kept)} file(s) kept, {len(plan.excluded)} excluded); " 

1429 f"mode: {plan.mode}" 

1430 ) 

1431 if plan.excluded: 

1432 log("excluded: " + ", ".join(f"{p} [{why}]" for p, why in plan.excluded)) 

1433 log(plan.reason) 

1434 

1435 if plan.mode == largediff.MODE_TOO_LARGE: 

1436 raise RuntimeError(f"diff too large to review: {plan.reason}") 

1437 if not plan.chunks: 

1438 raise RuntimeError( 

1439 "nothing to review after filters — all files were excluded " 

1440 "(check [jury.diff] include/exclude patterns)" 

1441 ) 

1442 

1443 # One shared budget across all chunks so ``total_timeout`` bounds the WHOLE 

1444 # review, not each chunk independently (review finding). ``phase_timeout`` and 

1445 # per-agent timeouts still apply per call via the same budget. 

1446 shared_budget = RunBudget(config.total_timeout, config.phase_timeout) 

1447 

1448 # Redact the shared context ONCE here, before fan-out (#249). The same context 

1449 # is reviewed against every chunk; letting each per-chunk ``run_jury`` redact 

1450 # it would count its secrets once per chunk and ``_merge_chunk_outcomes`` would 

1451 # sum them, inflating ``redaction_count`` (e.g. a 1-secret context over 8 

1452 # chunks reported 8). Pre-redacting makes each chunk's re-redaction a no-op — 

1453 # the ``[REDACTED:…]`` placeholders no longer match — so we add the one-time 

1454 # context count back at the end. Diff/chunk redactions are still counted 

1455 # per chunk and summed, which is correct (each chunk's diff is distinct). 

1456 # `config.context` is the right path: `_from_dict` flattens the `[jury]` 

1457 # table onto JuryConfig, so the `[jury.context]` sub-table is `config.context` 

1458 # (a ContextConfig), NOT `config.jury.context` — there is no `config.jury`. 

1459 # This mirrors how run_jury() reads it. 

1460 ctx_cfg = getattr(config, "context", None) 

1461 ctx_mode = getattr(ctx_cfg, "mode", "diff-only") if ctx_cfg else "diff-only" 

1462 redact_on = getattr(ctx_cfg, "redact_secrets", True) if ctx_cfg else True 

1463 context_redactions = 0 

1464 if redact_on and ctx_mode != "diff-only" and context: 

1465 context, context_redactions = redact(context) 

1466 

1467 # One band for the whole change, computed on the FILTERED diff the panel 

1468 # will actually see — the excluded files are not part of the change under 

1469 # review, and profiling the raw argument would route off files the panel is 

1470 # never shown (#714, r4). Every chunk then routes off the same answer. 

1471 whole_risk = ( 

1472 profile_diff("".join(plan.chunks)).risk if config.routing == routing.MODE_TIERED else None 

1473 ) 

1474 

1475 def _run(chunk: str) -> JuryOutcome: 

1476 return run_jury( 

1477 config, 

1478 chunk, 

1479 context=context, 

1480 hints=hints, 

1481 mock=mock, 

1482 strict=strict, 

1483 seed=seed, 

1484 policy=policy, 

1485 log=log, 

1486 budget=shared_budget, 

1487 on_event=on_event, 

1488 risk=whole_risk, 

1489 ) 

1490 

1491 def _finalize(outcome: JuryOutcome) -> JuryOutcome: 

1492 # Add the one-time context redaction count (per-chunk runs saw an already- 

1493 # redacted context and counted 0 for it). 

1494 if not context_redactions: 

1495 return outcome 

1496 return replace(outcome, redaction_count=outcome.redaction_count + context_redactions) 

1497 

1498 if plan.mode == largediff.MODE_FULL: 

1499 return _finalize(_run(plan.chunks[0])), plan 

1500 

1501 outcomes = [] 

1502 for i, chunk in enumerate(plan.chunks, 1): 

1503 log(f"reviewing chunk {i}/{len(plan.chunks)}") 

1504 outcomes.append(_run(chunk)) 

1505 return _finalize(_merge_chunk_outcomes(outcomes, config)), plan