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
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-30 06:29 +0000
1"""Jury orchestration: review -> debate -> synthesis.
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"""
8from __future__ import annotations
10import random
11import shutil
12import string
13import time
14from concurrent.futures import ThreadPoolExecutor
15from dataclasses import dataclass, field, replace
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
34class RunBudget:
35 """Wall-clock budget for one jury run (issue #30).
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 """
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()
49 def elapsed(self) -> float:
50 return time.monotonic() - self._start
52 def remaining(self) -> float | None:
53 if self.total is None:
54 return None
55 return max(0.0, self.total - self.elapsed())
57 def expired(self) -> bool:
58 return self.total is not None and self.elapsed() >= self.total
60 def call_timeout(self) -> int | None:
61 """Per-call timeout: the min of the phase budget and remaining total.
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)))
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).
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.
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
119def _order_by_agents(results: list[AgentResult], order: list[str]) -> list[AgentResult]:
120 """Reorder phase results into the configured/enabled agent order.
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))
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())
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)
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]
191def _others(reviews: list[AgentResult], me: str) -> str:
192 """Identity-labeled peer reviews (legacy path; ``anonymize_debate = false``).
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)_"
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
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).
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.
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
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.
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)
320def _skip_reason(spec) -> str:
321 """Why an unavailable seat was skipped, worded for its transport (#831).
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)"
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)
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
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")
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())
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)
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)
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 )
446 specs = config.enabled_agents
447 adapters = [make_adapter(s, mock=mock) for s in specs]
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)
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)]
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)
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)
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)
561 # Stream round-1 reviews as they're now finalized (stable order).
562 for r in reviews:
563 emit("review", r)
565 # Deterministic consensus grouping across reviewers.
566 groups = group_findings(all_findings, len(reviews))
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]
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)
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
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)"
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)
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)
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)
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}")
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 )
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).
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).
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)
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))]
842 if config.chair in names:
843 return config.chair
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]
851 return usable[0]
854def _format_findings_for_verify(findings: list[Finding]) -> str:
855 """Render candidate findings for the chair's verification prompt.
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)
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)
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
912def _verdict_matches_group(verdict: Verdict, group: FindingGroup) -> bool:
913 from .consensus import _normalize_claim, _normalize_path
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
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
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
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
980def _reject_targets(verdict: Verdict, groups: list[FindingGroup]) -> list[FindingGroup]:
981 """Un-statused groups a rejecting verdict may suppress (fail-closed, r7/r8).
983 Defences (each closes a distinct collateral-rejection vector found across
984 audit rounds 6-9):
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
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]
1035def _apply_verdicts(groups: list[FindingGroup], verdicts: list[Verdict]) -> None:
1036 """Attach verification statuses to consensus groups.
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
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)
1120def _first_sent_model(parts: list[AgentResult]) -> str:
1121 """The first model id *parts* recorded sending (issue #722).
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.
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.
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 ""
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).
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.
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)
1182 merged: list[AgentResult] = []
1183 for name in order:
1184 parts = by_agent[name]
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
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
1204 total_duration += p.duration_s
1205 if p.attempts > max_attempts:
1206 max_attempts = p.attempts
1208 body = "\n\n".join(body_parts)
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
1226def _combine_chair_results(results: list[AgentResult], chair: str) -> AgentResult | None:
1227 """Combine per-chunk chair results (verify/synthesis) into one labelled result.
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
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
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 )
1251def _merge_chunk_outcomes(outcomes: list[JuryOutcome], config: JuryConfig) -> JuryOutcome:
1252 """Fold per-chunk outcomes (issue #31) into one renderable JuryOutcome.
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]
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
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)
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)
1296 warnings = [w for o in outcomes for w in o.warnings]
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)
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
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
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 )
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).
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
1379def plan_for(config: JuryConfig, diff: str) -> largediff.DiffPlan:
1380 """The diff plan *config* selects for *diff* (pure).
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 )
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).
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:
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.
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)
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 )
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)
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)
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 )
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 )
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)
1498 if plan.mode == largediff.MODE_FULL:
1499 return _finalize(_run(plan.chunks[0])), plan
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