Construção de um pipeline auditável de previsões multiagente
Resumo
Este notebook monta um sistema de previsão com quatro etapas: um painel de agentes de pesquisa, agregação de probabilidades, um debate adversarial e um supervisor que identifica divergências, busca evidências esclarecedoras e emite uma probabilidade final. A influência do supervisor depende do nível de confiança declarado por ele: pode ser ignorado, combinado com o resultado anterior ou autorizado a substituí-lo. Os valores de configuração para as premissas de correlação, os pesos das etapas, os limites de busca e o tratamento da confiança são declarados de forma centralizada.
O pipeline registra probabilidades intermediárias, prompts, evidências e resultados de cada etapa para que um analista possa inspecionar como a previsão final foi formada. Ele reproduz previsões registradas para duas perguntas ainda sem resposta e mostra como elas evoluem ao longo das etapas. Como nenhum dos resultados era conhecido quando as previsões foram registradas, não é possível avaliar os exemplos, e o notebook não apresenta evidências de que o debate ou a supervisão melhorem a agregação dos agentes.
Os rótulos de confiança e os pesos de combinação são convenções de projeto, não estimativas calibradas. Os preços de mercado foram incluídos nos prompts dos agentes, portanto a distância entre as previsões e esses preços não constitui uma comparação independente. O notebook apresenta a auditabilidade como uma propriedade do sistema e deixa a avaliação preditiva para perguntas já respondidas e trabalhos posteriores de calibração.
Ideias principais
- A arquitetura combina pesquisa dos agentes, agregação, debate e revisão do supervisor em um processo de previsão rastreável.
- O supervisor investiga divergências, e seu efeito no conjunto depende de um critério baseado na confiança.
- As probabilidades e evidências de cada etapa devem ser mantidas para que uma previsão possa ser revisada posteriormente.
- As previsões de exemplo ainda não têm resultado conhecido e, por isso, não fornecem evidência de precisão nem provam que as etapas posteriores melhoram os resultados.
- A confiança declarada e os pesos de combinação precisam de avaliação empírica antes de serem interpretados como grandezas calibradas.
Tags
Texto completo
# 08_forecasting_pipeline.py
```py
# ---
# jupyter:
# jupytext:
# cell_metadata_filter: tags,-all
# text_representation:
# extension: .py
# format_name: percent
# format_version: '1.3'
# jupytext_version: 1.19.3
# kernelspec:
# display_name: Python 3 (ipykernel)
# language: python
# name: python3
# ---
# %% [markdown]
# # Full Forecasting Pipeline
#
# **Docker image**: `ml4t`
#
# The last four notebooks each built one piece. This one runs them as a sequence: a panel of
# research agents, an aggregate over them, a debate between two sides of the aggregate, and a
# supervisor that reads the panel, searches where the agents disagreed, and decides whether to
# override. That is the **AIA Forecaster** architecture, and the point of assembling it is not
# that four stages beat one; it is that a forecast produced this way can be taken apart when it
# is wrong.
#
# Both questions the pipeline runs on were unresolved when the capture was taken. It produces
# probabilities and cannot be scored, and neither can any claim that a later stage improved on
# an earlier one. Scoring against known outcomes is
# [`09_evaluation_and_governance`](09_evaluation_and_governance.ipynb).
#
# **Learning Objectives**:
# - Build a supervisor that finds where a panel disagreed, searches on those points, and
# returns its own probability
# - Gate that supervisor's influence on the confidence it states, so it can adjust the ensemble
# without silently replacing it
# - Compose four stages into one class whose settings are all declared in one place
# - Read a stage-by-stage probability path and see which stage moved the answer
# - Record a run in a form a scoring pipeline could later consume
#
# **Book Reference**: Chapter 24, Sections 24.7 (Multi-agent forecasting systems) and 24.9
# (Preparing for production)
#
# **Prerequisites**: [`04_research_agent`](04_research_agent.ipynb),
# [`05_aggregation_math`](05_aggregation_math.ipynb),
# [`06_multi_agent_research`](06_multi_agent_research.ipynb),
# [`07_adversarial_debate`](07_adversarial_debate.ipynb).
# %%
"""Full Forecasting Pipeline: agent-debate-supervisor end-to-end."""
import json
import textwrap
import time
from datetime import date, datetime
import matplotlib.pyplot as plt
import polars as pl
from agent_fixtures import get_chapter_clear_question, get_chapter_contested_question
from agent_observability import (
TRACES_DIR,
RunTrace,
show_agents,
show_debate_transcript,
show_supervisor,
trace_llm,
)
from agent_pipeline import neyman_extremize
from agent_providers import ChatMessage, TokenUsage, create_llm_client
from agent_research import ResearchAgent, format_agent_summary, parse_json
from agent_schemas import (
AgentForecastArtifact,
ForecastQuestion,
ForecastResult,
SearchResult,
SupervisorArtifact,
)
from agent_specialists import DebateAgent
from agent_tools import (
SearchClient,
ToolExecutor,
create_search_client,
)
from IPython.display import Markdown, display
from utils.style import COLORS, add_message_title, show_with_alt
# %% [markdown]
# ## Settings
#
# `RUN_LIVE` left at `False` replays the two pinned captures, one per question, and makes no
# API calls.
#
# `N_AGENTS`, `MAX_STEPS` and `MAX_SEARCH_RESULTS` configure the research phase exactly as in
# [`06_multi_agent_research`](06_multi_agent_research.ipynb). `DEBATE_ROUNDS` caps the
# argument, and `NEYMAN_CORRELATION` is the pairwise correlation assumed when the panel is
# aggregated.
#
# Three weights decide how much each later stage can move the answer, and they are the numbers
# to argue with - on a live run. `DEBATE_WEIGHT` is the debate midpoint's share of the
# post-debate probability. `SUPERVISOR_MEDIUM_WEIGHT` is the supervisor's share when it states
# medium confidence; at high confidence it replaces the value outright and at low confidence it
# is ignored. A replay reports the probabilities its capture recorded, so changing either weight
# with `RUN_LIVE` left at `False` changes nothing: the run that produced the numbers is over.
#
# The three `CONFIDENCE_WHEN_*` values are what the pipeline reports as its own confidence in
# each of those three cases. They are an ordering, not an estimate: a forecast the supervisor
# overrode is marked as having had more reconciliation than one where it was ignored, and
# nothing here measures whether either is more likely to be right.
# %% tags=["parameters"]
RUN_LIVE = False
PINNED_TRACES = [
"08_forecasting_pipeline_20260609T141954Z_ef5bca7b95bd.json", # recession (clear)
"08_forecasting_pipeline_20260609T142158Z_24e083e7fe54.json", # rate hike (contested)
]
LLM_PROVIDER = ""
N_AGENTS = 3
DEBATE_ROUNDS = 3
MAX_STEPS = 5
MAX_SEARCH_RESULTS = 5
NEYMAN_CORRELATION = 0.3
DEBATE_WEIGHT = 0.3
SUPERVISOR_MEDIUM_WEIGHT = 0.4
CONFIDENCE_WHEN_OVERRIDDEN = 0.8
CONFIDENCE_WHEN_BLENDED = 0.6
CONFIDENCE_WHEN_IGNORED = 0.5
# %% [markdown]
# ## Supervisor Prompts
#
# The supervisor operates in two phases:
# 1. **Identify disagreements** among agents and propose clarifying searches
# 2. **Finalize** with updated probability, incorporating new evidence
#
# These prompts are shown inline from the AIA Forecaster's production templates.
# %%
SUPERVISOR_DISAGREEMENTS_PROMPT = """\
You are the SUPERVISOR agent.
You receive M agent forecasts and rationales for the same question.
Your job is NOT to average them directly.
Step 1: Identify key disagreements, ambiguities, missing base rates, or claims that should be fact-checked.
Step 2: Propose up to {max_queries} clarifying search queries that would resolve these disagreements.
Output JSON only with:
{{"disagreements": ["..."], "queries": ["..."]}}
AGENT INPUTS:
{agent_summaries}"""
# %% [markdown]
# The finalize prompt receives the original question, agent panel, and bounded
# follow-up evidence. It requests a probability, confidence label, and rationale.
# %%
SUPERVISOR_FINALIZE_PROMPT = """\
You are the SUPERVISOR agent.
Given:
1) The original question
2) The set of agent forecasts and rationales
3) Additional evidence from your follow-up searches
You must output:
1) Updated forecast p_yes in [0,1]
2) Confidence in whether your update direction is correct: "high" | "medium" | "low"
3) A short rationale
Output JSON only:
{{"p_yes": 0.0, "confidence": "high", "rationale": "..."}}
QUESTION:
{question}
AGENT INPUTS:
{agent_summaries}
SUPERVISOR SEARCH EVIDENCE:
{supervisor_evidence}"""
# %% [markdown]
# ## Supervisor phase helpers
#
# The supervisor's three phases (identify, search, finalize) factor
# naturally into free functions. The class then becomes a thin shell that
# owns the LLM + search clients and threads them through the helpers.
# %%
def _supervisor_identify_disagreements(
llm,
agent_summaries: str,
max_queries: int,
) -> tuple[list[str], list[str], TokenUsage]:
"""Phase 1: LLM call asking for disagreements and clarifying queries."""
prompt = SUPERVISOR_DISAGREEMENTS_PROMPT.format(
max_queries=max_queries, agent_summaries=agent_summaries
)
raw, tokens = llm.complete_with_usage(
[ChatMessage(role="user", content=prompt)], json_mode=True
)
parsed = parse_json(raw)
disagreements = [str(x) for x in parsed.get("disagreements", [])][:20]
queries = [str(x) for x in parsed.get("queries", [])][:max_queries]
return disagreements, queries, tokens
# %% [markdown]
# The search phase executes only the supervisor's bounded query list and
# applies the question's point-in-time cutoff to every request.
# %%
def _supervisor_run_searches(
search,
queries: list[str],
max_search_results: int,
cutoff_date: date | None,
) -> dict[str, list[SearchResult]]:
"""Phase 2: execute the clarifying-search queries via ToolExecutor."""
if search is None:
return {}
executor = ToolExecutor(search=search)
return {
q: executor.execute_search(q, max_results=max_search_results, cutoff_date=cutoff_date)
for q in queries
}
# %% [markdown]
# Search results become a compact evidence block for the final model call.
# Dates and source URLs remain visible for audit.
# %%
def _format_supervisor_evidence(sr: dict[str, list[SearchResult]]) -> str:
"""Render the search-result dict into the supervisor-finalize prompt's evidence block."""
lines: list[str] = []
for q, results in sr.items():
lines.append(f"QUERY: {q}")
for i, r in enumerate(results, start=1):
lines.append(f"{i}. {r.title}")
if r.url:
lines.append(f" URL: {r.url}")
if r.snippet:
lines.append(f" {r.snippet}")
if r.published:
lines.append(f" Published: {r.published}")
lines.append("")
return "\n".join(lines) if lines else "No additional search evidence."
# %% [markdown]
# The final phase parses and clamps the supervisor probability. Invalid
# confidence labels fall back to `medium` instead of triggering an override.
# %%
def _supervisor_finalize(
llm,
question: str,
agent_summaries: str,
search_results: dict[str, list[SearchResult]],
) -> tuple[float | None, str | None, str | None, TokenUsage]:
"""Phase 3: LLM call asking for final p_yes / confidence / rationale."""
evidence_text = _format_supervisor_evidence(search_results)
prompt = SUPERVISOR_FINALIZE_PROMPT.format(
question=question,
agent_summaries=agent_summaries,
supervisor_evidence=evidence_text,
)
raw, tokens = llm.complete_with_usage(
[ChatMessage(role="user", content=prompt)], json_mode=True
)
parsed = parse_json(raw)
p_yes_raw = parsed.get("p_yes")
confidence = parsed.get("confidence")
rationale = parsed.get("rationale")
if confidence is not None:
conf_str = str(confidence).lower()
if conf_str not in ("high", "medium", "low"):
conf_str = "medium"
confidence = conf_str
p_yes = float(p_yes_raw) if p_yes_raw is not None else None
if p_yes is not None:
p_yes = max(0.0, min(1.0, p_yes))
return p_yes, confidence, rationale, tokens
# %% [markdown]
# The driver combines the three phases and returns both the artifact and its
# token count.
# %%
def _run_supervisor(
supervisor,
question: str,
agent_summaries: str,
cutoff_date: date | None,
) -> tuple[SupervisorArtifact, TokenUsage]:
"""Run identify, search, and finalize in sequence."""
disagreements, queries, identify_tokens = _supervisor_identify_disagreements(
supervisor.llm, agent_summaries, supervisor.max_queries
)
search_results = _supervisor_run_searches(
supervisor.search, queries, supervisor.max_search_results, cutoff_date
)
p_yes, confidence, rationale, finalize_tokens = _supervisor_finalize(
supervisor.llm, question, agent_summaries, search_results
)
tokens = identify_tokens + finalize_tokens
artifact = SupervisorArtifact(
disagreements=disagreements,
queries=queries,
search_results=search_results,
p_yes=p_yes,
confidence=confidence,
rationale=str(rationale) if rationale is not None else None,
token_usage=tokens,
)
return artifact, tokens
# %% [markdown]
# ## The SupervisorAgent Class
#
# Three phases: (1) identify disagreements, (2) run clarifying searches,
# (3) finalize with evidence. The supervisor only overrides the ensemble
# when its confidence is "high", preserving agent diversity by default.
# %%
class SupervisorAgent:
"""Supervisor that reconciles agent ensemble via clarifying searches."""
def __init__(
self,
llm,
search: SearchClient | None = None,
max_queries: int = 3,
max_search_results: int = 5,
) -> None:
self.llm = llm
self.search = search
self.max_queries = max_queries
self.max_search_results = max_search_results
self.token_usage = TokenUsage()
def run(
self,
question: str,
agent_summaries: str,
cutoff_date: date | None = None,
) -> SupervisorArtifact:
"""Run supervisor reconciliation. Returns SupervisorArtifact."""
artifact, self.token_usage = _run_supervisor(self, question, agent_summaries, cutoff_date)
return artifact
# %% [markdown]
# ## Pipeline helpers
#
# The debate stage is the implementation from
# [`07_adversarial_debate`](07_adversarial_debate.ipynb), imported from
# `agent_specialists`. What is local to this notebook is the research phase and the rule that
# decides how much of the supervisor's opinion reaches the final number.
# %%
def _run_research_agents(
llm,
search,
question: ForecastQuestion,
n_agents: int,
max_steps: int,
) -> list[AgentForecastArtifact]:
"""Phase 1: run N identical ResearchAgents and collect their artifacts."""
artifacts: list[AgentForecastArtifact] = []
for i in range(n_agents):
agent = ResearchAgent(llm=llm, search=search, agent_id=f"agent_{i}", max_steps=max_steps)
artifacts.append(agent.run(question, market_price=question.current_market_price))
return artifacts
# %% [markdown]
# The supervisor has seen the agents' summaries and one round of clarifying searches; the
# agents each did their own research. So the supervisor gets a say proportional to the
# confidence it states, and never an unconditional one: it replaces the post-debate probability
# only at high confidence, is mixed in at `SUPERVISOR_MEDIUM_WEIGHT` at medium, and is ignored
# at low.
#
# The confidence values the pipeline attaches to its own output are stated conventions, not
# measurements. They rank three outcomes - the supervisor overrode, it contributed, it was
# ignored - so a consumer can order forecasts by how much reconciliation they received. Nothing
# estimates them, and [`09_evaluation_and_governance`](09_evaluation_and_governance.ipynb) is
# where a confidence that means something has to come from.
# %%
def _blend_final_probability(
post_debate: float,
supervisor_artifact: SupervisorArtifact,
*,
medium_weight: float = SUPERVISOR_MEDIUM_WEIGHT,
) -> tuple[float, float]:
"""Phase 4 to final: confidence-gated supervisor override. Returns (p_yes, confidence)."""
final_p = post_debate
final_confidence = CONFIDENCE_WHEN_IGNORED
if supervisor_artifact.p_yes is not None and supervisor_artifact.confidence == "high":
final_p = supervisor_artifact.p_yes
final_confidence = CONFIDENCE_WHEN_OVERRIDDEN
elif supervisor_artifact.confidence == "medium":
if supervisor_artifact.p_yes is not None:
final_p = (1 - medium_weight) * post_debate + (
medium_weight * supervisor_artifact.p_yes
)
final_confidence = CONFIDENCE_WHEN_BLENDED
return max(0.01, min(0.99, final_p)), final_confidence
# %% [markdown]
# The execution helper composes the four phases and records their artifacts.
# Keeping orchestration outside the class leaves the reader-facing class as a
# small configuration object.
# %%
def _forecast_one(forecaster, question: ForecastQuestion) -> ForecastResult:
started = time.time()
cutoff = date.fromisoformat(question.cutoff_date) if question.cutoff_date else None
agents = _run_research_agents(
forecaster.llm, forecaster.search, question, forecaster.n_agents, forecaster.max_steps
)
answered = [agent for agent in agents if agent.forecast_produced]
if not answered:
raise RuntimeError(f"no agent produced a forecast for: {question.question}")
summaries = "\n\n---\n\n".join(format_agent_summary(agent) for agent in answered)
aggregation = neyman_extremize(
[agent.p_yes for agent in answered], base=0.5, correlation=forecaster.correlation
)
aggregate_p = (
aggregation.extremized_probability
if aggregation.extremized_probability is not None
else aggregation.raw_probability
)
debate = DebateAgent(
llm=forecaster.llm,
max_rounds=forecaster.debate_rounds,
consensus_threshold=forecaster.consensus_threshold,
).run(question.question, summaries, aggregate_p)
midpoint = (
(debate.bull_final_probability + debate.bear_final_probability) / 2
if debate.bull_final_probability is not None
else aggregate_p
)
supervisor = SupervisorAgent(llm=forecaster.llm, search=forecaster.search).run(
question.question, summaries, cutoff_date=cutoff
)
post_debate = (1 - DEBATE_WEIGHT) * aggregate_p + DEBATE_WEIGHT * midpoint
final_p, confidence = _blend_final_probability(
post_debate, supervisor, medium_weight=SUPERVISOR_MEDIUM_WEIGHT
)
tokens = sum((agent.token_usage for agent in agents), start=TokenUsage())
tokens = tokens + debate.token_usage + supervisor.token_usage
return ForecastResult(
question=question,
agents=agents,
aggregation=aggregation,
debate=debate,
supervisor=supervisor,
final_probability=round(final_p, 4),
final_confidence=round(confidence, 3),
total_token_usage=tokens,
duration_seconds=round(time.time() - started, 2),
)
# %% [markdown]
# ## The AIAForecaster Class
#
# The complete four-phase pipeline:
# 1. **Research agents**: N parallel agents produce forecasts
# 2. **Aggregation**: Neyman extremization combines agent probabilities
# 3. **Debate**: Bull/bear stress-test the aggregate
# 4. **Supervisor**: Reconcile with clarifying searches, confidence-gated override
# %%
class AIAForecaster:
"""Complete AIA Forecaster pipeline: agents → aggregate → debate → supervisor."""
def __init__(
self,
llm,
search: SearchClient | None = None,
n_agents: int = 3,
max_steps: int = 5,
debate_rounds: int = 3,
consensus_threshold: float = 0.05,
correlation: float = 0.3,
) -> None:
self.llm = llm
self.search = search
self.n_agents = n_agents
self.max_steps = max_steps
self.debate_rounds = debate_rounds
self.consensus_threshold = consensus_threshold
self.correlation = correlation
def forecast(self, question: ForecastQuestion) -> ForecastResult:
"""Run the full pipeline on a single question."""
return _forecast_one(self, question)
# %% [markdown]
# ## The Two Questions
#
# The pipeline runs end to end on both of the chapter's pinned questions:
# `CHAPTER_CLEAR_QUESTION`, the recession question where the research agents landed close
# together in [`06_multi_agent_research`](06_multi_agent_research.ipynb), and
# `CHAPTER_CONTESTED_QUESTION`, the rate-hike question where they spread out in
# [`07_adversarial_debate`](07_adversarial_debate.ipynb). Running both shows what each stage
# does when the panel already agrees and when it does not.
#
# Both were unresolved when the captures were taken, which is what makes them honest forecasts
# and also what makes them unscoreable; the replay cell reports each capture's date from the
# record. Replayed by default; `RUN_LIVE = True` with `ANTHROPIC_API_KEY` and `TAVILY_API_KEY`
# forecasts current questions instead.
# %%
questions = [get_chapter_clear_question(), get_chapter_contested_question()]
print(f"Forecasting {len(questions)} questions:")
for q in questions:
market = f"{q.current_market_price:.0%}" if q.current_market_price is not None else "?"
print(f" • {q.question}")
if q.resolution_date:
print(f" Resolves: {q.resolution_date} | Market: {market}")
else:
print(f" Market: {market}")
# %% [markdown]
# ## Running the Pipeline
# %% [markdown]
# Live execution is isolated in one helper. The publication path below never
# calls it while `RUN_LIVE` remains false.
# %%
def _run_live_questions(questions_to_run: list[ForecastQuestion]) -> tuple[list, list]:
"""Run and persist fresh provider-backed forecasts."""
llm = create_llm_client(LLM_PROVIDER)
search = create_search_client(LLM_PROVIDER)
live_results, live_traces = [], []
for q in questions_to_run:
tracer = trace_llm(llm, label="pipeline")
forecaster = AIAForecaster(
llm=tracer,
search=search,
n_agents=N_AGENTS,
max_steps=MAX_STEPS,
debate_rounds=DEBATE_ROUNDS,
correlation=NEYMAN_CORRELATION,
)
result = forecaster.forecast(q)
run = RunTrace.from_result(
result,
notebook="08_forecasting_pipeline",
provider=llm.model_name,
params={
"n_agents": N_AGENTS,
"max_steps": MAX_STEPS,
"debate_rounds": DEBATE_ROUNDS,
"correlation": NEYMAN_CORRELATION,
"debate_weight": DEBATE_WEIGHT,
"supervisor_medium_weight": SUPERVISOR_MEDIUM_WEIGHT,
},
llm_calls=tracer.calls,
notes="Full AIA pipeline: research, aggregate, debate, supervisor.",
)
path = run.save()
print(
f" ✓ {q.question[:50]}... → {result.final_probability:.2f} "
f"({result.duration_seconds:.1f}s) | {len(run.llm_calls)} calls → {path.name}"
)
live_results.append(result)
live_traces.append(run)
return live_results, live_traces
# %% [markdown]
# Replay loads only the two committed trace names. No provider client or search
# client is created on this path.
# %%
def _load_pinned_questions() -> tuple[list, list]:
"""Rehydrate the two committed pipeline traces."""
replay_results, replay_traces = [], []
for pinned_name in PINNED_TRACES:
run = RunTrace.load(TRACES_DIR / pinned_name)
result = run.forecast_result()
recorded = datetime.fromisoformat(run.created_at).date().isoformat()
print(
f" ✓ {result.question.question[:50]}... → {result.final_probability:.2f} "
f"(replay of a capture recorded {recorded}) | {len(run.llm_calls)} calls"
)
replay_results.append(result)
replay_traces.append(run)
return replay_results, replay_traces
# %%
if RUN_LIVE:
results, run_traces = _run_live_questions(questions)
else:
results, run_traces = _load_pinned_questions()
# %% [markdown]
# ## Results Summary
# %%
grand_total = TokenUsage()
for r in results:
grand_total = grand_total + r.total_token_usage
summary_df = pl.DataFrame(
[
{
"question": r.question.question[:80],
"final": round(r.final_probability, 3),
"market": (
round(r.question.current_market_price, 3)
if r.question.current_market_price is not None
else None
),
"confidence": round(r.final_confidence, 2),
"duration_s": (
round(r.duration_seconds, 1) if r.duration_seconds is not None else None
),
}
for r in results
]
)
print(f"Total tokens across {len(results)} questions: {grand_total.total_tokens:,}\n")
summary_df
# %% [markdown]
# ## The Full Trace, One Question
#
# The untruncated record for the first question, stage by stage, through the same
# `agent_observability` helpers used across the chapter: each research agent's queries,
# documents and rationale; the debate transcript with both sides' complete arguments; and the
# supervisor's reconciliation, including the disagreements it flagged, the clarifying searches
# it ran, and the probability it returned. All of it is in the saved JSON, so this readout can
# be rebuilt from disk with `RunTrace.load` long after the run.
# %%
r = results[0]
print(f"Question: {r.question.question}\n")
print(show_agents(r.agents))
# %% [markdown]
# ### Aggregation
# %%
if r.aggregation:
print(f" method: {r.aggregation.method}")
print(f" inputs: {r.aggregation.input_probabilities}")
print(f" raw mean: {r.aggregation.raw_probability:.2f}")
print(f" extremized: {r.aggregation.extremized_probability:.2f}")
print(f" d={r.aggregation.extremization_factor:.2f}, n_eff={r.aggregation.effective_n:.1f}")
# %% [markdown]
# ### Debate
# %%
if r.debate:
print(show_debate_transcript(r.debate))
# %% [markdown]
# ### Supervisor, and what the pipeline returned
# %%
if r.supervisor:
print(show_supervisor(r.supervisor))
print(f" probability: {r.final_probability:.2f}")
print(f" confidence: {r.final_confidence:.2f}")
if r.duration_seconds is not None:
print(f" duration: {r.duration_seconds:.1f}s")
else:
print(" duration: n/a (replayed from pinned trace)")
# %% [markdown]
# ## Where the Probability Went
#
# Three of the numbers a run produces are the same quantity at different points in the
# pipeline: the aggregate over the research agents, the post-debate blend, and the final
# probability. Those are what the line below joins. Everything else a stage produced is an
# input to one of them - the market price the agents were shown, the agents' own answers, the
# debate midpoint that enters the post-debate blend at the run's own debate weight, and the
# supervisor's own probability - and is drawn as an open marker at the stage that read it.
#
# The distinction decides what a gap on this chart means. Between two carried values it is
# movement, and a stage that never moves anything on any question is being paid for and not
# used. Between two research agents it is disagreement: they answer in parallel and neither
# saw the other.
# %% [markdown]
# The post-debate value is the one number on the line that no stage stores: it is rebuilt from
# the aggregate and the debate midpoint. On a replay it has to be rebuilt with the weights the
# capture was recorded with rather than the ones set in this notebook now. Mixing the two
# recomputes the middle point of the line and leaves the points either side of it at their
# recorded values. Raising `DEBATE_WEIGHT` far enough would pull the post-debate point above
# both of its neighbours and draw a large supervisor correction that never happened. So the
# weights take effect on a live run, and are read back from the trace on a replay.
#
# Neither weight was recorded when the two committed captures were taken, so they fall back to
# the values in use then. `carried_probabilities` re-derives each capture's final probability
# from its own parts and raises if the fallback does not reproduce it, which is what keeps the
# fallback from becoming an assumption nobody checks.
# %%
CAPTURED_DEBATE_WEIGHT = 0.3
CAPTURED_SUPERVISOR_MEDIUM_WEIGHT = 0.4
def carried_probabilities(result: ForecastResult, run: RunTrace) -> tuple[float, float, float]:
"""Return (aggregate, debate midpoint, post-debate) under the run's own weights."""
debate_weight = float(run.params.get("debate_weight", CAPTURED_DEBATE_WEIGHT))
medium_weight = float(
run.params.get("supervisor_medium_weight", CAPTURED_SUPERVISOR_MEDIUM_WEIGHT)
)
aggregate_p = (
result.aggregation.extremized_probability
if result.aggregation.extremized_probability is not None
else result.aggregation.raw_probability
)
midpoint = (
(result.debate.bull_final_probability + result.debate.bear_final_probability) / 2
if result.debate and result.debate.bull_final_probability is not None
else aggregate_p
)
post_debate = (1 - debate_weight) * aggregate_p + debate_weight * midpoint
if result.supervisor is not None:
rebuilt, _ = _blend_final_probability(
post_debate, result.supervisor, medium_weight=medium_weight
)
if abs(round(rebuilt, 4) - result.final_probability) > 1e-4:
raise RuntimeError(
f"Rebuilt final probability {rebuilt:.4f} does not match the recorded "
f"{result.final_probability:.4f} for {result.question.question[:50]}: "
f"the weights this line is drawn with are not the ones the run used."
)
return aggregate_p, midpoint, post_debate
# %%
STAGE_X = {
"Market": 0,
"Agents": 1,
"Aggregate": 2,
"Post-debate": 3,
"Supervisor": 4,
"Final": 5,
}
carried_rows, input_rows = [], []
for r, run in zip(results, run_traces, strict=True):
q_short = textwrap.shorten(r.question.question, width=44, placeholder="...")
aggregate_p, midpoint, post_debate = carried_probabilities(r, run)
carried_rows += [
{"question": q_short, "stage": "Aggregate", "p_yes": aggregate_p},
{"question": q_short, "stage": "Post-debate", "p_yes": post_debate},
{"question": q_short, "stage": "Final", "p_yes": r.final_probability},
]
if r.question.current_market_price is not None:
input_rows.append(
{"question": q_short, "stage": "Market", "p_yes": r.question.current_market_price}
)
for a in r.agents:
if a.p_yes is not None:
input_rows.append({"question": q_short, "stage": "Agents", "p_yes": a.p_yes})
if r.debate and r.debate.bull_final_probability is not None:
input_rows.append({"question": q_short, "stage": "Post-debate", "p_yes": midpoint})
if r.supervisor and r.supervisor.p_yes is not None:
input_rows.append({"question": q_short, "stage": "Supervisor", "p_yes": r.supervisor.p_yes})
carried_df = pl.DataFrame(carried_rows)
inputs_df = pl.DataFrame(input_rows)
fig, ax = plt.subplots()
questions = carried_df["question"].unique(maintain_order=True).to_list()
palette = dict(zip(questions, [COLORS["blue"], COLORS["amber"]]))
for q in questions:
carried = carried_df.filter(pl.col("question") == q)
ax.plot(
[STAGE_X[stage] for stage in carried["stage"]],
carried["p_yes"].to_list(),
marker="o",
linewidth=2,
color=palette[q],
label=q,
zorder=3,
)
stage_inputs = inputs_df.filter(pl.col("question") == q)
ax.scatter(
[STAGE_X[stage] for stage in stage_inputs["stage"]],
stage_inputs["p_yes"].to_list(),
facecolors="none",
edgecolors=palette[q],
s=55,
zorder=2,
)
ax.set_xticks(list(STAGE_X.values()), list(STAGE_X.keys()))
ax.set_xlabel("Pipeline stage")
ax.set_ylabel("Probability of yes")
ax.set_ylim(0, 1)
add_message_title(
ax,
"The probability the pipeline carries, and what each stage read",
subtitle="Filled markers joined by a line are the running answer; open markers are inputs",
)
ax.legend(loc="upper right")
show_with_alt(
fig,
"Chart of probability against pipeline stage, one colour per question. A line joins three "
"filled markers on each - the aggregate, the post-debate blend and the final probability - "
"passing over the supervisor stage without a marker there. Open markers of the same colour "
"sit at the market price, at each research agent's answer, at the debate midpoint and at "
"the supervisor's own probability. On one question the three agent markers coincide; on "
"the other, two coincide and the third sits far above them. Both lines stay in the lower "
"half of the range, and one ends higher than it started while the other ends lower.",
)
# %% [markdown]
# The steps below are taken over the carried values alone, so each one is the same quantity
# before and after a stage weighted something into it.
# %%
largest_move = (
carried_df.with_columns(
pl.col("p_yes").diff().over("question").alias("move"),
pl.col("stage").shift().over("question").alias("from_stage"),
)
.drop_nulls("move")
.with_columns(pl.col("move").abs().alias("size"))
.sort("size", descending=True)
.group_by("question", maintain_order=True)
.first()
.select(
"question",
pl.format("{} to {}", "from_stage", "stage").alias("largest step"),
pl.col("move").round(3),
)
)
largest_move
# %% [markdown]
# The largest carried step falls at a different stage on each question: the post-debate blend
# on one, the supervisor blend on the other. No stage in this pipeline decides the answer and
# none of them is idle.
#
# The open markers carry a reading the line cannot. On the recession question all three
# research agents returned the same probability and the aggregate came out below every one of
# them, because extremizing away from a base of even odds treats agreement between agents as
# evidence - which is the assumption
# [`05_aggregation_math`](05_aggregation_math.ipynb) derives and
# [`06_multi_agent_research`](06_multi_agent_research.ipynb) tests against agents that share a
# prompt. Whether any of this movement is an improvement is a scoring question, and neither of
# these questions had resolved.
# %% [markdown]
# ## Token use
# %%
print(f" Questions forecasted: {len(results)}")
print(f" Total tokens: {grand_total.total_tokens:,}")
print(f" Input tokens: {grand_total.input_tokens:,}")
print(f" Output tokens: {grand_total.output_tokens:,}")
# %% [markdown]
# ## The Record a Scoring Pipeline Would Read
#
# Everything above lives in memory. What a scoring or monitoring system needs is a flat record
# per question: the forecast, the cutoff, the outcome if it is known, and enough of the
# intermediate stages to attribute a bad forecast to one of them.
#
# Neither of these questions had resolved when the capture was taken, so `resolved_outcome` is
# empty and nothing here can be scored.
# [`09_evaluation_and_governance`](09_evaluation_and_governance.ipynb) builds the scoring rules
# against a panel of questions whose outcomes are known, rather than against these two.
# %%
serialized = [
{
"question": r.question.question,
"cutoff_date": r.question.cutoff_date,
"final_probability": r.final_probability,
"resolved_outcome": r.question.resolved_outcome,
"agent_probs": [a.p_yes for a in r.agents],
"aggregate": r.aggregation.extremized_probability if r.aggregation else None,
"debate_consensus": r.debate.consensus_reached if r.debate else None,
"supervisor_p_yes": r.supervisor.p_yes if r.supervisor else None,
"supervisor_confidence": r.supervisor.confidence if r.supervisor else None,
"tokens": r.total_token_usage.total_tokens,
"duration_s": r.duration_seconds,
}
for r in results
]
print(json.dumps(serialized[0], indent=2))
print(f"\n({len(serialized)} records; the second has the same shape)")
# %% [markdown]
# The pipeline moves a probability through four stages and records where it went. What it does
# not establish is that any stage improved the estimate. The aggregation credits the panel with
# an independence nobody measured. The debate narrows the two sides by a few points on the
# contested question, which is a smaller disagreement and not a more accurate one. The
# supervisor's confidence label is its own assertion about itself. Reading the stage-to-stage
# movement as progressive refinement is the mistake this record exists to prevent, and it is
# why the figure above draws the whole path rather than the endpoint.
# %%
recession, rate_hike = results
display(
Markdown(
"**This replay.** "
f"The recession forecast ends at {recession.final_probability:.1%} and the rate-hike "
f"forecast at {rate_hike.final_probability:.1%}, against market prices of "
f"{recession.question.current_market_price:.1%} and "
f"{rate_hike.question.current_market_price:.1%}. Neither distance is a comparison: "
"both market prices were in every research agent's prompt, so the pipeline was told "
"where the market stood before it looked at anything. And both questions were "
"unresolved when the capture was taken, so neither forecast can be scored."
)
)
# %% [markdown]
# ## Key Takeaways
#
# 1. **The value of a pipeline is that each stage is inspectable, not that each stage improves
# the answer.** Four stages give four places to look when a forecast is wrong. Whether any
# of them made it better is a scoring question, and scoring needs resolved questions.
# 2. **An override needs a gate, and the gate needs a rule.** The supervisor can replace the
# ensemble only at high stated confidence, blends at medium, and is ignored at low. Without
# that, one model's second opinion silently outranks three agents' evidence.
# 3. **A stated confidence is not a measured one.** Both the supervisor's own label and the
# scalar this pipeline attaches to the final probability are conventions. They order
# outcomes; they do not estimate anything.
# 4. **Declare every constant that moves the final number in one place.** The blend
# weights, the correlation and the override thresholds decide the output, and a reader who
# cannot find them cannot evaluate the pipeline.
# 5. **The stage-by-stage record is the deliverable.** One JSON file per question holds every
# prompt, every document, every intermediate probability, and it is what makes a forecast
# reviewable months later.
#
# **Known limitations of what is built here.** Both questions were unresolved when captured, so
# nothing in this notebook can be scored and no claim about accuracy is available from it. The
# market price is passed to every research agent, so the pipeline's distance from the market
# is not an independent comparison. The blend weights and confidence scalars are conventions,
# and no experiment here shows the four-stage output is better than the three-agent mean.
#
# **Next**: [`09_evaluation_and_governance`](09_evaluation_and_governance.ipynb) builds the
# scoring rules, calibration curves and security controls a pipeline like this needs before
# anyone acts on it.
#
# **Book**: Section 24.7 covers the pipeline architecture and section 24.9 the production
# considerations that follow from it.
```Exibido na íntegra, com atribuição conforme a licença da fonte. Licença: MIT
Este resumo foi escrito pelo agente de pesquisa da Stratmill com base no original; não é uma cópia da fonte.