Files
6krrt/src/metrics.py
adlee-was-taken 6b0c3e3e85 test: tidy the julianday window fix
Match the SQL indent of the two edited lines in metrics.py to their
neighbours (one space too deep, which made the diff noisier than the change),
drop test_src_has_no_unexpected_hits from the tripwire (an identical copy of
test_no_iso_t_datetime_comparison misfiled under the positive controls), and
replace an em dash in a comment with ASCII.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KkCGRantZsSwmcFpet6FTa
2026-10-05 21:32:12 -04:00

3479 lines
141 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Read-only aggregation helpers for the router dashboard / /health.
This module MUST NOT import ``dispatcher`` — it exists specifically to break
what would otherwise be a circular import (dispatcher wants /health metrics,
metrics wants the config and DB path that dispatcher already knows).
All functions take an ``sqlite3.Connection`` (with ``row_factory`` set) and
optionally a ``RouterConfig`` instance; none rely on module-level globals.
Functions
---------
quota_burn — kWh metered in the last 30 d and in the current billing
period, against the plan allowance; also reports per-provider
credit balance, burn rate and runway (from
allowance_remaining_usd telemetry or the polled
provider_balance_observations table) plus total_balance_usd
scoring_coverage — which scoring axes actually have data
capability_ceilings — vision / json_mode context-window sub-ceilings
capability_demand_warnings — demand-relative warnings for those sub-ceilings
rejection_warnings — grouped rejection-rate / new-pattern warnings
cache_rate_series — prefix-cache hit rate per (provider, model), trailing window
cache_rate_warnings — expiry check on objective.assumed_cache_rate
cost_estimate_calibration — estimate-vs-bill error per (provider, model),
REPORTED ONLY: nothing applies it to estimated_cost or to
ranking
latency_series — router-observed wall / TTFT p50+p95 per (provider, model),
a report and never an objective
unroutable_models — active rows that can never enter the candidate set
selection_coverage — which active models the router has actually picked
recent_decisions — last N rows from the route_decisions observability table
per_model — per-model aggregates over energy_observations (last 30 d)
verdict_mix — counts by verdict from verifications (last N days)
top_proficiency — top models by blended_score for a category
proficiency_matrix — every proficiency row with source / samples / provenance
recent_client_outcomes — latest POST /outcome reports, attribution included
pinch_summary — context-pruning savings over the last 30 d
"""
from __future__ import annotations
import re
import sqlite3
from datetime import date, datetime, timedelta, timezone
from typing import Any, Final, List, Optional, Sequence
import routing
SEED_CATEGORY = "seed_reference"
FALLBACK_CARBON_SOURCE = "static_fallback"
def _percentile(values: Sequence[float], p: float) -> Optional[float]:
"""Return the *p*-th percentile of *values* (0–100), or None for empty."""
if not values:
return None
if len(values) == 1:
return float(values[0])
s = sorted(values)
n = len(s)
k = p / 100.0 * (n - 1)
f = int(k)
frac = k - f
if f + 1 < n:
return s[f] + frac * (s[f + 1] - s[f])
return float(s[f])
def _next_reset_date(billing_reset_day: int, today: Optional[date] = None) -> str:
"""Return the next billing-reset date as an ISO date string.
When *today*'s day-of-month is before *billing_reset_day*, returns this
month's reset day; when today's day is on or past the reset day, returns
the next month's (handling December → January rollover).
"""
if today is None:
today = datetime.now(timezone.utc).date()
if today.day < billing_reset_day:
return date(today.year, today.month, billing_reset_day).isoformat()
if today.month == 12:
next_year = today.year + 1
next_month = 1
else:
next_year = today.year
next_month = today.month + 1
return date(next_year, next_month, billing_reset_day).isoformat()
def _billing_period_start(billing_reset_day: int, today: Optional[date] = None) -> str:
"""Return the most recent billing-period start as an ISO date string.
When *today*'s day-of-month is on or after *billing_reset_day*, returns
this month's reset day; otherwise returns the previous month's reset day
(handling January → December rollover).
"""
if today is None:
today = datetime.now(timezone.utc).date()
if today.day >= billing_reset_day:
return date(today.year, today.month, billing_reset_day).isoformat()
if today.month == 1:
prev_year = today.year - 1
prev_month = 12
else:
prev_year = today.year
prev_month = today.month - 1
return date(prev_year, prev_month, billing_reset_day).isoformat()
def _balance_series_stats(
samples: list[tuple[datetime, float]],
*,
latest: Optional[tuple[str, float]],
burn_window_hours: float,
warning_hours: float,
min_samples: int,
min_hours: float,
) -> dict[str, Any]:
"""Compute balance/burn/runway for one provider-scoped series.
``samples`` are the in-window (observed_at, balance) points, ordered
oldest-first; rows with unparseable timestamps are already skipped by the
caller. ``latest`` is the provider's most recent balance row regardless of
window (``(observed_at, balance_usd)``) and may be None.
The burn rate is intentionally conservative: only the latest
monotonically-decreasing balance segment is used, and two guards stop a
fresh credit top-up from producing a wild extrapolation.
"""
result: dict[str, Any] = {
"balance_usd": None,
"balance_at": None,
"burn_window_hours": burn_window_hours,
"burn_rate_usd_per_hour": None,
"projected_hours_remaining": None,
"runway_low_warning": False,
"runway_note": None,
}
if latest is not None:
result["balance_at"] = latest[0]
result["balance_usd"] = latest[1]
if not samples:
result["runway_note"] = (
"burn estimate unavailable: no decreasing balance samples in the current window"
)
return result
# Split into monotonically-decreasing segments at every balance INCREASE.
# A top-up (credit jump) starts a new segment; only the latest survives.
segments: list[list[tuple[datetime, float]]] = [[]]
for ts, value in samples:
current = segments[-1]
if not current:
current.append((ts, value))
elif value > current[-1][1]:
segments.append([(ts, value)])
else:
current.append((ts, value))
latest_segment = segments[-1]
if not latest_segment:
result["runway_note"] = (
"burn estimate unavailable: no decreasing balance samples in the current window"
)
return result
if len(latest_segment) < min_samples:
result["runway_note"] = (
f"burn estimate unavailable: segment after last balance increase has only "
f"{len(latest_segment)} sample(s), need {min_samples}"
)
return result
first_ts, first_balance = latest_segment[0]
last_ts, last_balance = latest_segment[-1]
elapsed_hours = (last_ts - first_ts).total_seconds() / 3600.0
if elapsed_hours < min_hours:
minutes = elapsed_hours * 60
if minutes < 60:
duration_str = f"{minutes:.0f} minutes"
else:
duration_str = f"{elapsed_hours:.1f} hours"
result["runway_note"] = (
f"burn estimate unavailable: segment after last balance increase spans only "
f"{duration_str}, need at least {min_hours} h"
)
return result
total_decrease = first_balance - last_balance
if total_decrease <= 0.0:
# Treat a flat or increasing-only segment like the no-burn case.
result["runway_note"] = (
"burn estimate unavailable: no decreasing balance samples in the current window"
)
return result
burn_rate = round(total_decrease / elapsed_hours, 6)
result["burn_rate_usd_per_hour"] = burn_rate
balance = result["balance_usd"]
if balance is not None and burn_rate > 0.0:
projected = balance / burn_rate
result["projected_hours_remaining"] = projected
result["runway_low_warning"] = projected < warning_hours
return result
def quota_balance_and_burn(conn: sqlite3.Connection, cfg: Any) -> dict:
"""Per-provider balance, burn and runway.
Returns ``{"by_provider": {provider: {balance_usd, balance_at,
balance_source, burn_window_hours, burn_rate_usd_per_hour,
projected_hours_remaining, runway_low_warning, runway_note}},
"total_balance_usd": float|None}`` with one entry per
``cfg.dispatch_providers`` key.
The two-source rule:
* ``has_energy_telemetry`` providers read
``energy_observations.allowance_remaining_usd`` scoped to their provider.
* Every other provider reads ``provider_balance_observations`` scoped to
its provider (the polled account-balance table), guarded with
``sqlite3.OperationalError`` so a live DB predating that table degrades
to an all-None entry instead of failing ``/metrics``.
``balance_source`` is CONFIG-derived and present even on all-None entries,
so a consumer can tell a prepaid-pool depletion story from an
overage-allowance accounting figure without re-deriving config:
``"telemetry"`` (per-completion allowance), ``"polled"`` (balance_url
configured), ``"unconfigured"`` (neither).
"""
# Read optional knobs from config; treat None as unset and use code defaults.
# Test configs use SimpleNamespace without these attributes, so getattr
# must have a fallback and then a second default when the attr is None.
burn_window_hours = getattr(cfg.objective, "quota_burn_window_hours", None)
if burn_window_hours is None:
burn_window_hours = 24
warning_hours = getattr(cfg.objective, "quota_runway_warning_hours", None)
if warning_hours is None:
warning_hours = 6
min_samples = getattr(cfg.objective, "quota_burn_min_segment_samples", None)
if min_samples is None:
min_samples = 3
min_hours = getattr(cfg.objective, "quota_burn_min_segment_hours", None)
if min_hours is None:
min_hours = 0.5
providers = getattr(cfg, "dispatch_providers", None) or {}
by_provider: dict[str, dict[str, Any]] = {}
for provider, pc in providers.items():
has_telemetry = getattr(pc, "has_energy_telemetry", False)
balance_url = getattr(pc, "balance_url", None)
if has_telemetry:
balance_source = "telemetry"
latest_sql = """
SELECT observed_at, allowance_remaining_usd
FROM energy_observations
WHERE provider = ?
AND allowance_remaining_usd IS NOT NULL
ORDER BY observed_at DESC, id DESC
LIMIT 1
"""
window_sql = """
SELECT observed_at, allowance_remaining_usd
FROM energy_observations
WHERE provider = ?
AND allowance_remaining_usd IS NOT NULL
AND julianday(observed_at) > julianday('now', '-' || ? || ' hours')
ORDER BY observed_at ASC, id ASC
"""
else:
balance_source = "polled" if balance_url else "unconfigured"
latest_sql = """
SELECT observed_at, balance_usd
FROM provider_balance_observations
WHERE provider = ?
ORDER BY observed_at DESC, id DESC
LIMIT 1
"""
window_sql = """
SELECT observed_at, balance_usd
FROM provider_balance_observations
WHERE provider = ?
AND julianday(observed_at) > julianday('now', '-' || ? || ' hours')
ORDER BY observed_at ASC, id ASC
"""
latest: Optional[tuple[str, float]] = None
raw_rows: list[sqlite3.Row] = []
try:
row = conn.execute(latest_sql, (provider,)).fetchone()
if row is not None and row[1] is not None:
latest = (row["observed_at"], float(row[1]))
raw_rows = conn.execute(
window_sql,
(provider, str(burn_window_hours)),
).fetchall()
except sqlite3.OperationalError:
# Live DBs that predate the provider_balance_observations table
# must not 500 /metrics; report an all-None entry instead.
latest = None
raw_rows = []
samples: list[tuple[datetime, float]] = []
for row in raw_rows:
try:
ts = datetime.fromisoformat(row["observed_at"])
except ValueError:
# Defensive: malformed timestamp would otherwise break /metrics.
continue
samples.append((ts, float(row[1])))
entry = _balance_series_stats(
samples,
latest=latest,
burn_window_hours=burn_window_hours,
warning_hours=warning_hours,
min_samples=min_samples,
min_hours=min_hours,
)
entry["balance_source"] = balance_source
by_provider[provider] = entry
known_balances = [
entry["balance_usd"]
for entry in by_provider.values()
if entry["balance_usd"] is not None
]
total_balance_usd = round(sum(known_balances), 6) if known_balances else None
return {"by_provider": by_provider, "total_balance_usd": total_balance_usd}
def _balance_delta_over_period(
conn: sqlite3.Connection, provider: str, period_start: str
) -> Optional[float]:
"""How far a prepaid balance fell over the period, in USD.
First reading minus last. Returns None when there are fewer than two
readings to compare, and a NEGATIVE value is returned as-is for the caller
to reject -- a balance that rose means a top-up landed mid-period, and the
drop across it is not this period's spend. Treating a top-up as negative
spend would understate the bill, which is the more dangerous direction.
"""
rows = conn.execute(
"""
SELECT balance_usd
FROM provider_balance_observations
WHERE provider = ?
AND julianday(observed_at) >= julianday(?)
ORDER BY observed_at ASC, id ASC
""",
(provider, period_start),
).fetchall()
if len(rows) < 2:
return None
return float(rows[0]["balance_usd"]) - float(rows[-1]["balance_usd"])
# Span and bucket, in seconds, for each selectable timeframe. Mirrors
# admin._ADMIN_HISTORY_RANGES so the quota graph and the dashboard's History
# card offer the same choices and bucket them identically.
QUOTA_HISTORY_RANGES: Final[dict[str, tuple[int, int]]] = {
"1h": (3600, 300),
"6h": (21600, 900),
"24h": (86400, 3600),
"7d": (604800, 21600),
"30d": (2592000, 86400),
}
# Below this share of priced calls, a provider's per-request spend is not a
# usable series and the account poll is used instead. OpenRouter sat at 0%
# before per-request capture started working, so a per-request spend line for
# it would have been a flat zero across the whole history -- the same lie the
# account cards told before it was fixed.
_SPEND_COVERAGE_FLOOR: Final[float] = 0.5
def quota_history(conn: sqlite3.Connection, cfg: Any, range_key: str = "24h") -> dict:
"""Per-provider usage over time, for the quota page's graph.
Three metrics, because they answer different questions and have different
coverage:
spend the billing question. Exact per request for a telemetry
provider; for a prepaid pool whose requests carry no cost, taken
from the movement of the polled balance instead -- coarser (the
poll interval, ~2h) but true, which a flat zero line is not.
calls perfect coverage for every provider, all the way back. The only
series with no caveat anywhere on it.
energy telemetry providers only. OpenRouter reports none, so its series
is legitimately empty rather than zero.
``sources`` names how each provider's spend was derived so the chart can
say so rather than implying one fidelity for both lines.
"""
if range_key not in QUOTA_HISTORY_RANGES:
raise ValueError(
f"unknown range {range_key!r}; expected one of "
f"{sorted(QUOTA_HISTORY_RANGES)}"
)
span, bucket = QUOTA_HISTORY_RANGES[range_key]
providers = list((getattr(cfg, "dispatch_providers", None) or {}).keys())
local = getattr(getattr(cfg, "local_dispatch", None), "provider", "ollama-local")
if getattr(getattr(cfg, "local_energy", None), "enabled", False):
providers.append(local)
def buckets(table: str, aggregate: str, provider: Optional[str]) -> list[list]:
where = "julianday(observed_at) >= julianday('now') - ? / 86400.0"
params: list[Any] = [bucket, bucket, span]
if provider is not None:
where += " AND provider = ?"
params.append(provider)
rows = conn.execute(
f"""
SELECT CAST((julianday(observed_at) - julianday('1970-01-01'))
* 86400.0 / ? AS INTEGER) * ? AS bucket_ts,
{aggregate} AS value
FROM {table}
WHERE {where}
GROUP BY 1 ORDER BY 1
""",
params,
).fetchall()
return [[int(r["bucket_ts"]), round(float(r["value"] or 0), 6)] for r in rows]
def balance_spend(provider: str) -> list[list]:
"""Spend per bucket from the fall of a polled balance.
A RISE means a top-up landed; that bucket is reported as zero spend
rather than negative, because a refill is not income.
"""
try:
rows = conn.execute(
"""
SELECT CAST((julianday(observed_at) - julianday('1970-01-01'))
* 86400.0 / ? AS INTEGER) * ? AS bucket_ts,
MIN(balance_usd) AS lo,
MAX(balance_usd) AS hi
FROM provider_balance_observations
WHERE provider = ?
AND julianday(observed_at) >= julianday('now') - ? / 86400.0
GROUP BY 1 ORDER BY 1
""",
(bucket, bucket, provider, span),
).fetchall()
except sqlite3.OperationalError:
return []
return [[int(r["bucket_ts"]), round(max(float(r["hi"]) - float(r["lo"]), 0.0), 6)]
for r in rows]
spend: dict[str, list] = {}
calls: dict[str, list] = {}
energy: dict[str, list] = {}
sources: dict[str, str] = {}
for name in providers:
is_local = name == local
table = "local_energy_observations" if is_local else "energy_observations"
scope = None if is_local else name
calls[name] = buckets(table, "COUNT(*)", scope)
energy[name] = buckets(table, "COALESCE(SUM(energy_kwh), 0)", scope)
cov = conn.execute(
f"""
SELECT COUNT(*) n, SUM(cost_usd IS NOT NULL) priced
FROM {table}
WHERE julianday(observed_at) >= julianday('now') - ? / 86400.0
{'' if is_local else 'AND provider = ?'}
""",
(span,) if is_local else (span, name),
).fetchone()
total, priced = cov["n"] or 0, cov["priced"] or 0
share = (priced / total) if total else 0.0
if total and share >= _SPEND_COVERAGE_FLOOR:
spend[name] = buckets(table, "COALESCE(SUM(cost_usd), 0)", scope)
sources[name] = "per_request"
else:
derived = balance_spend(name)
if derived:
spend[name] = derived
sources[name] = "balance_poll"
else:
spend[name] = []
sources[name] = "unavailable"
return {
"range": range_key,
"bucket_seconds": bucket,
"ranges": sorted(QUOTA_HISTORY_RANGES, key=lambda k: QUOTA_HISTORY_RANGES[k][0]),
"series": {"spend": spend, "calls": calls, "energy": energy},
"sources": sources,
}
def quota_accounts(conn: sqlite3.Connection, cfg: Any) -> dict:
"""Per-provider billing shape, spend, and alarm for the admin quota dashboard.
Returns the JSON shape from WI-3a: ``period``, ``accounts[]``,
``spend``, and ``alarm`` blocks. One function replaces
``quota_burn()`` + ``quota_balance_and_burn()``.
Shape derivation (4-case, first match wins):
1. provider name is ``\"ollama-local\"`` -> ``self_hosted``
2. ``has_energy_telemetry`` and ``plan_kwh_per_period`` set -> ``metered_plan``
3. ``balance_url`` set -> ``prepaid_credit``
4. none of the above -> ``unmetered``
"""
now = datetime.now(timezone.utc)
today = now.date()
# --- Period block ----------------------------------------------------------
reset_day = getattr(cfg.objective, "billing_reset_day", None)
if reset_day is not None:
period_start_str = _billing_period_start(reset_day, today)
next_reset_str = _next_reset_date(reset_day, today)
period_start_date = date.fromisoformat(period_start_str)
next_reset_date_obj = date.fromisoformat(next_reset_str)
period_days = (next_reset_date_obj - period_start_date).days
elapsed_days = (today - period_start_date).days
if period_days > 0 and elapsed_days >= 0:
raw_frac = elapsed_days / period_days
elapsed_fraction: Optional[float] = round(raw_frac, 4)
if elapsed_fraction is not None and elapsed_fraction < 0.02:
elapsed_fraction = None
else:
elapsed_fraction = None
period = {
"start": period_start_str,
"next_reset": next_reset_str,
"elapsed_fraction": elapsed_fraction,
"source": "billing_reset_day",
}
else:
period_start_str = (today - timedelta(days=30)).isoformat()
period = {
"start": period_start_str,
"next_reset": None,
"elapsed_fraction": 1.0,
"source": "30d_rolling",
}
providers = getattr(cfg, "dispatch_providers", None) or {}
accounts: list[dict[str, Any]] = []
plan_kwh = getattr(cfg.objective, "plan_kwh_per_period", None)
# Config knobs for burn estimation (mirrors quota_balance_and_burn defaults)
burn_window_hours = getattr(cfg.objective, "quota_burn_window_hours", None) or 24
warning_hours = getattr(cfg.objective, "quota_runway_warning_hours", None) or 6
min_samples = getattr(cfg.objective, "quota_burn_min_segment_samples", None) or 3
min_hours = getattr(cfg.objective, "quota_burn_min_segment_hours", None) or 0.5
# Pace alarm threshold from config
plan_pace_warn_ratio = getattr(cfg.objective, "plan_pace_warn_ratio", None) or 1.25
# --- Accounts — per configured provider ------------------------------------
for name, pc in providers.items():
has_telemetry = getattr(pc, "has_energy_telemetry", False)
balance_url = getattr(pc, "balance_url", None)
# 1. Determine shape via 4-case precedence
if name == "ollama-local":
shape = "self_hosted"
elif has_telemetry and plan_kwh is not None and plan_kwh > 0:
shape = "metered_plan"
elif balance_url:
shape = "prepaid_credit"
else:
shape = "unmetered"
account: dict[str, Any] = {"provider": name, "shape": shape}
# 2. Period spend + attribution coverage (all providers)
if shape == "self_hosted":
spend_row = conn.execute(
"""
SELECT COUNT(*) total_rows,
SUM(CASE WHEN cost_usd IS NOT NULL THEN 1 ELSE 0 END) cost_rows,
COALESCE(SUM(cost_usd), 0) total_cost
FROM local_energy_observations
WHERE julianday(observed_at) >= julianday(?)
""",
(period_start_str,),
).fetchone()
else:
spend_row = conn.execute(
"""
SELECT COUNT(*) total_rows,
SUM(CASE WHEN cost_usd IS NOT NULL THEN 1 ELSE 0 END) cost_rows,
COALESCE(SUM(cost_usd), 0) total_cost
FROM energy_observations
WHERE julianday(observed_at) >= julianday(?)
AND provider = ?
""",
(period_start_str, name),
).fetchone()
total_rows = spend_row["total_rows"] or 0
cost_rows = spend_row["cost_rows"] or 0
period_cost = round(float(spend_row["total_cost"] or 0.0), 6)
attribution_coverage = round(cost_rows / total_rows, 4) if total_rows > 0 else 1.0
spend_source = "per_request_billed"
# RULE 1 -- account truth for totals, request truth for attribution,
# never add them.
#
# Summing per-request cost is right only when the provider reports it
# on every request. OpenRouter did not report it at all until the
# usage-accounting flag shipped, so this sum was COALESCE'd to 0 and
# the page stated, in dollars, that a provider with real spend had
# cost nothing -- $0.00 against an account that had actually moved
# $10.25 over the period.
#
# For a prepaid pool the account poll IS the authoritative total: what
# the provider says it took out of the balance. Use it whenever the
# request rows cannot account for the period, and keep
# attribution_coverage to say how much of that total the per-model
# breakdown can explain.
if shape == "prepaid_credit" and attribution_coverage < 1.0:
delta = _balance_delta_over_period(conn, name, period_start_str)
if delta is not None and delta > 0:
period_cost = round(delta, 6)
spend_source = "account_poll_delta"
account["spend_usd"] = {
"period": period_cost,
"attribution_coverage": attribution_coverage,
"source": spend_source,
}
# 3. Plan block (metered_plan only)
plan_block: Optional[dict[str, Any]] = None
if shape == "metered_plan":
kwh_row = conn.execute(
"""
SELECT COALESCE(SUM(energy_kwh), 0) used_kwh
FROM energy_observations
WHERE julianday(observed_at) >= julianday(?)
AND provider = ?
""",
(period_start_str, name),
).fetchone()
used_kwh = round(float(kwh_row["used_kwh"]), 5)
used_fraction = round(used_kwh / plan_kwh, 4) if plan_kwh and plan_kwh > 0 else 0.0
plan_block = {
"kwh_per_period": plan_kwh,
"used_kwh": used_kwh,
"used_fraction": used_fraction,
}
account["plan"] = plan_block
# 4. Pool block (prepaid_credit only)
pool_block: Optional[dict[str, Any]] = None
if shape == "prepaid_credit":
try:
pool_row = conn.execute(
"""
SELECT balance_usd, total_credits_usd, total_usage_usd, observed_at
FROM provider_balance_observations
WHERE provider = ?
ORDER BY observed_at DESC, id DESC
LIMIT 1
""",
(name,),
).fetchone()
except sqlite3.OperationalError:
pool_row = None
if pool_row is not None:
bal = float(pool_row["balance_usd"])
tcu: Optional[float] = (
float(pool_row["total_credits_usd"])
if pool_row["total_credits_usd"] is not None
else None
)
tu: Optional[float] = (
float(pool_row["total_usage_usd"])
if pool_row["total_usage_usd"] is not None
else None
)
observed_at_str = pool_row["observed_at"]
try:
age = (now - datetime.fromisoformat(observed_at_str)).total_seconds()
except (ValueError, TypeError):
age = None
pool_block = {
"total_credits_usd": tcu,
"total_usage_usd": tu,
"balance_usd": bal,
"age_seconds": int(age) if age is not None else None,
}
account["pool"] = pool_block
# 5. Burn block (metered_plan or prepaid_credit with burn data)
burn_block: Optional[dict[str, Any]] = None
if shape in ("metered_plan", "prepaid_credit"):
if shape == "metered_plan":
balance_source = "telemetry"
latest_sql = """
SELECT observed_at, allowance_remaining_usd
FROM energy_observations
WHERE provider = ?
AND allowance_remaining_usd IS NOT NULL
ORDER BY observed_at DESC, id DESC
LIMIT 1
"""
window_sql = """
SELECT observed_at, allowance_remaining_usd
FROM energy_observations
WHERE provider = ?
AND allowance_remaining_usd IS NOT NULL
AND julianday(observed_at) > julianday('now', '-' || ? || ' hours')
ORDER BY observed_at ASC, id ASC
"""
else:
balance_source = "polled" if balance_url else "unconfigured"
latest_sql = """
SELECT observed_at, balance_usd
FROM provider_balance_observations
WHERE provider = ?
ORDER BY observed_at DESC, id DESC
LIMIT 1
"""
window_sql = """
SELECT observed_at, balance_usd
FROM provider_balance_observations
WHERE provider = ?
AND julianday(observed_at) > julianday('now', '-' || ? || ' hours')
ORDER BY observed_at ASC, id ASC
"""
latest_balance: Optional[tuple[str, float]] = None
raw_rows: list[sqlite3.Row] = []
try:
row = conn.execute(latest_sql, (name,)).fetchone()
if row is not None and row[1] is not None:
latest_balance = (row["observed_at"], float(row[1]))
raw_rows = conn.execute(
window_sql,
(name, str(burn_window_hours)),
).fetchall()
except sqlite3.OperationalError:
latest_balance = None
raw_rows = []
samples: list[tuple[datetime, float]] = []
for row in raw_rows:
try:
ts = datetime.fromisoformat(row["observed_at"])
except ValueError:
continue
samples.append((ts, float(row[1])))
stats = _balance_series_stats(
samples,
latest=latest_balance,
burn_window_hours=burn_window_hours,
warning_hours=warning_hours,
min_samples=min_samples,
min_hours=min_hours,
)
if stats["burn_rate_usd_per_hour"] is not None:
burn_block = {
"burn_rate_usd_per_hour": stats["burn_rate_usd_per_hour"],
"projected_hours_remaining": stats["projected_hours_remaining"],
}
account["burn"] = burn_block
# 6. Credit block.
#
# Two ways an account can carry a dollar balance, and both belong here:
# a polled prepaid pool (balance_url), and a telemetry provider's
# OVERAGE ALLOWANCE -- the credit a subscription bills against once the
# plan is spent. The latter was being fetched for the burn calculation
# above and then dropped on the floor, because this block was gated on
# balance_url alone, so a subscription card could never show what its
# overage is drawing on.
credit_block: Optional[dict[str, Any]] = None
if not balance_url and shape == "metered_plan" and latest_balance is not None:
credit_block = {
"balance_usd": latest_balance[1],
"balance_at": latest_balance[0],
"balance_source": "telemetry",
"meaning": "overage_allowance",
}
elif balance_url:
if shape not in ("metered_plan", "prepaid_credit"):
# Independent query for self_hosted/unmetered with a balance_url
try:
credit_row = conn.execute(
"""
SELECT balance_usd, observed_at
FROM provider_balance_observations
WHERE provider = ?
ORDER BY observed_at DESC, id DESC
LIMIT 1
""",
(name,),
).fetchone()
except sqlite3.OperationalError:
credit_row = None
credit_balance = (
float(credit_row["balance_usd"]) if credit_row is not None else None
)
credit_at = credit_row["observed_at"] if credit_row is not None else None
else:
# Reuse data already queried for the burn block
credit_balance = latest_balance[1] if latest_balance is not None else None
credit_at = latest_balance[0] if latest_balance is not None else None
if credit_balance is not None:
credit_block = {
"balance_usd": credit_balance,
"balance_at": credit_at,
"balance_source": "polled",
"meaning": "prepaid_pool",
}
account["credit"] = credit_block
# 6b. Where this account's money went, per model, over the period.
#
# Attached to the ACCOUNT rather than collected into a shared ledger:
# a blended provider-vs-provider table is exactly what made the first
# version unreadable, and a breakdown is only ever meaningful inside
# the account whose bill it explains.
#
# cost_usd is NULL wherever the provider did not report a per-request
# figure, and that is kept as NULL rather than coalesced to 0 -- the
# card says how much of the account total the list can account for, so
# a zero here would quietly contradict it.
# Rank by spend only when spend is actually known for most of the
# account. Ordering by cost under thin coverage puts nonsense on top:
# openrouter ranked xiaomi/mimo-v2.5 first on $0.00002 (two priced
# calls out of 666) above a model with 905 calls and no price at all.
# Under thin coverage, call count is the honest ordering, and
# models_ranked_by tells the card which sentence to print.
rank_by_cost = attribution_coverage >= 0.5
order_clause = (
"COALESCE(SUM(cost_usd), 0) DESC, COUNT(*) DESC"
if rank_by_cost
else "COUNT(*) DESC, COALESCE(SUM(cost_usd), 0) DESC"
)
account["models_ranked_by"] = "spend" if rank_by_cost else "calls"
if shape == "self_hosted":
model_sql = f"""
SELECT model_id,
COUNT(*) AS calls,
SUM(cost_usd) AS cost_usd,
SUM(CASE WHEN cost_usd IS NOT NULL THEN 1 ELSE 0 END) AS priced
FROM local_energy_observations
WHERE julianday(observed_at) >= julianday(?)
GROUP BY model_id
ORDER BY {order_clause}
LIMIT 12
"""
model_params: tuple[Any, ...] = (period_start_str,)
else:
model_sql = f"""
SELECT model_id,
COUNT(*) AS calls,
SUM(cost_usd) AS cost_usd,
SUM(CASE WHEN cost_usd IS NOT NULL THEN 1 ELSE 0 END) AS priced
FROM energy_observations
WHERE julianday(observed_at) >= julianday(?)
AND provider = ?
GROUP BY model_id
ORDER BY {order_clause}
LIMIT 12
"""
model_params = (period_start_str, name)
account["models"] = [
{
"model_id": r["model_id"],
"calls": r["calls"],
"cost_usd": round(float(r["cost_usd"]), 6) if r["cost_usd"] is not None else None,
"priced_calls": r["priced"] or 0,
}
for r in conn.execute(model_sql, model_params).fetchall()
]
# 7. Energy block (has_energy_telemetry or self_hosted)
energy_block: Optional[dict[str, Any]] = None
if has_telemetry:
energy_row = conn.execute(
"""
SELECT COALESCE(SUM(energy_kwh), 0) kwh, COUNT(*) n
FROM energy_observations
WHERE julianday(observed_at) > julianday('now', '-30 days')
AND provider = ?
""",
(name,),
).fetchone()
energy_block = {
"kwh_30d": round(float(energy_row["kwh"]), 5),
"calls_30d": energy_row["n"],
}
elif shape == "self_hosted":
energy_row = conn.execute(
"""
SELECT COALESCE(SUM(energy_kwh), 0) kwh, COUNT(*) n
FROM local_energy_observations
WHERE julianday(observed_at) > julianday('now', '-30 days')
"""
).fetchone()
energy_block = {
"kwh_30d": round(float(energy_row["kwh"]), 5),
"calls_30d": energy_row["n"],
}
account["energy"] = energy_block
accounts.append(account)
# --- Spend block -----------------------------------------------------------
by_provider_usd: dict[str, float] = {}
total_usd = 0.0
for acc in accounts:
period_val = acc["spend_usd"]["period"]
by_provider_usd[acc["provider"]] = period_val
total_usd += period_val
total_usd = round(total_usd, 6)
# Estimated cost from token × list-price (energy_observations + models join)
# The estimate worth comparing against the bill is the router's OWN --
# route_decisions.est_cost_usd, the figure the cost tiebreak ranks
# candidates on. Recomputing tokens x list price here instead measured
# something else entirely: a no-cache-discount price on traffic that runs
# ~92% cached, which came out at $430 against a $22.91 bill and rendered
# as "estimate_ratio 18.79". That number was not wrong about its own
# arithmetic; it was answering a question nobody asked.
#
# Per provider, so a provider whose estimate tracks well is not averaged
# in with one whose estimate is 4x out -- which is the whole finding.
estimated_by_provider: dict[str, float] = {}
for name, _pc in providers.items():
est_row = conn.execute(
"""
SELECT COALESCE(SUM(est_cost_usd), 0) AS estimated
FROM route_decisions
WHERE julianday(observed_at) >= julianday(?)
AND selected_provider = ?
AND est_cost_usd IS NOT NULL
""",
(period_start_str, name),
).fetchone()
if est_row and est_row["estimated"]:
estimated_by_provider[name] = round(float(est_row["estimated"]), 6)
estimated_usd = round(sum(estimated_by_provider.values()), 6)
# Per-provider ratios, and only where BOTH sides are real: a ratio against
# a zero or unknown actual is a division artifact, not a calibration.
estimate_ratio_by_provider: dict[str, float] = {}
for name, est in estimated_by_provider.items():
actual = by_provider_usd.get(name) or 0.0
if actual > 0 and est > 0:
estimate_ratio_by_provider[name] = round(est / actual, 2)
estimate_ratio: Optional[float] = (
round(estimated_usd / total_usd, 4) if total_usd > 0 and estimated_usd > 0 else None
)
spend = {
"by_provider_usd": by_provider_usd,
"total_usd": total_usd,
"estimated_usd": estimated_usd,
"estimated_by_provider_usd": estimated_by_provider,
"estimate_ratio": estimate_ratio,
"estimate_ratio_by_provider": estimate_ratio_by_provider,
}
# --- Alarm block -----------------------------------------------------------
alarm: Optional[dict[str, Any]] = None
# Check metered_plan pace (first metered_plan wins)
for acc in accounts:
if acc["shape"] == "metered_plan" and acc["plan"] is not None:
elapsed_frac = period["elapsed_fraction"]
pace_ratio: Optional[float] = None
if elapsed_frac is not None and elapsed_frac > 0:
used_fraction = acc["plan"]["used_fraction"]
pace_ratio = used_fraction / elapsed_frac
if pace_ratio is not None and pace_ratio > 2.0:
alarm = {
"kind": "plan_pace",
"provider": acc["provider"],
"severity": "critical",
"headline": f"{pace_ratio:.1f}x plan pace",
}
elif pace_ratio is not None and pace_ratio > plan_pace_warn_ratio:
alarm = {
"kind": "plan_pace",
"provider": acc["provider"],
"severity": "warning",
"headline": f"{pace_ratio:.1f}x plan pace",
}
elif acc["plan"]["used_fraction"] > 0.8:
# Absolute-usage floor, beneath the pace rule and not
# redundant with it. Two shapes need it:
#
# - No billing_reset_day: period is a rolling 30 days with
# elapsed_fraction pinned to 1.0, so pace_ratio equals
# used_fraction and can never exceed 1.0. Without this
# branch such a deployment burns 96% of its plan and
# NOTHING warns -- which is what the old >80% warning
# caught and what removing it would have silently lost.
# - Late in a real period: 90% used at 95% elapsed is a pace
# of 0.95, perfectly healthy by the pace rule, and still
# about to run out of plan.
alarm = {
"kind": "plan_usage",
"provider": acc["provider"],
"severity": "warning",
"headline": (
f"{acc['plan']['used_fraction'] * 100:.0f}% of plan used"
),
}
else:
alarm = {
"kind": "none",
"provider": acc["provider"],
"severity": "info",
"headline": f"${total_usd:.2f} period spend",
}
break # Only the first metered_plan drives the pace alarm
# Check prepaid_credit stale readings (only when no metered_plan alarm)
if alarm is None:
for acc in accounts:
if acc["shape"] == "prepaid_credit" and acc["pool"] is not None:
age = acc["pool"].get("age_seconds")
if age is not None and age > 3 * 3600: # poll_interval default 3600s
alarm = {
"kind": "stale_reading",
"provider": acc["provider"],
"severity": "warning",
"headline": f"balance reading {age:.0f}s old",
}
break
return {
"period": period,
"accounts": accounts,
"spend": spend,
"alarm": alarm,
}
def quota_burn(
conn: sqlite3.Connection,
cfg: Any,
) -> Optional[dict]:
"""Energy this router has metered, against the plan's allowance.
Accepts ``(conn, cfg)`` so the caller owns the connection and the config
— metrics.py never touches dispatcher's module-level ``cfg`` or its
``_db()`` helper, which is exactly why this module must never import
dispatcher.
Returns a dict with ``plan_kwh``, ``metered_kwh_30d``,
``metered_kwh_period``, ``metered_fraction_of_plan``,
``metered_calls_30d``, ``reset_date`` (the billing-period start),
``window_start_30d`` (the rolling 30-day window start), ``note``, plus the
per-provider ``by_provider`` mapping and ``total_balance_usd`` from
``quota_balance_and_burn``. When ``cfg.objective.billing_reset_day`` is
set, also returns ``next_reset_date`` — the upcoming billing-period reset
day.
"""
# This gate removes the report only; it is intentionally not used to refuse
# or alter request dispatch — routing decisions remain independent of quota.
if not cfg.objective.plan_kwh_per_period:
return None
plan = cfg.objective.plan_kwh_per_period
window_start_30d = (datetime.now(timezone.utc).date() - timedelta(days=30)).isoformat()
providers = getattr(cfg, "dispatch_providers", None)
if providers is None:
# SimpleNamespace fixture compatibility: cfg has no dispatch_providers,
# so keep the historical unscoped SUM/COUNT behavior.
telemetry_providers: Optional[list[str]] = None
else:
telemetry_providers = [
name
for name, pc in providers.items()
if getattr(pc, "has_energy_telemetry", False)
]
if telemetry_providers == []:
# No provider can meter energy. Both sums are honestly zero; do not
# run the queries with an empty IN () clause (SQLite syntax error).
metered_kwh_30d = 0.0
metered_calls_30d = 0
else:
if telemetry_providers:
telemetry_placeholders = ",".join("?" * len(telemetry_providers))
thirty_where = (
"julianday(observed_at) > julianday('now', '-30 days') "
f"AND provider IN ({telemetry_placeholders})"
)
thirty_params = tuple(telemetry_providers)
else:
thirty_where = "julianday(observed_at) > julianday('now', '-30 days')"
thirty_params = ()
row = conn.execute(
f"""
SELECT COALESCE(SUM(energy_kwh), 0) kwh, COUNT(*) n
FROM energy_observations
WHERE {thirty_where}
""",
thirty_params,
).fetchone()
metered_kwh_30d = round(float(row["kwh"]), 5)
metered_calls_30d = row["n"]
reset_day = getattr(cfg.objective, "billing_reset_day", None)
if reset_day is not None:
period_start = _billing_period_start(reset_day)
if telemetry_providers == []:
metered_kwh_period = 0.0
else:
if telemetry_providers:
period_placeholders = ",".join("?" * len(telemetry_providers))
period_where = (
"julianday(observed_at) >= julianday(?) "
f"AND provider IN ({period_placeholders})"
)
period_params = (period_start, *telemetry_providers)
else:
period_where = "julianday(observed_at) >= julianday(?)"
period_params = (period_start,)
period_row = conn.execute(
f"""
SELECT COALESCE(SUM(energy_kwh), 0) kwh
FROM energy_observations
WHERE {period_where}
""",
period_params,
).fetchone()
metered_kwh_period = round(float(period_row["kwh"]), 5)
metered_fraction_of_plan = round(metered_kwh_period / plan, 4)
reset_date = period_start
else:
metered_kwh_period = None
metered_fraction_of_plan = round(metered_kwh_30d / plan, 4)
reset_date = None
result = {
"plan_kwh": plan,
"metered_kwh_30d": metered_kwh_30d,
"metered_kwh_period": metered_kwh_period,
"metered_fraction_of_plan": metered_fraction_of_plan,
"metered_calls_30d": metered_calls_30d,
"reset_date": reset_date,
"window_start_30d": window_start_30d,
"note": "router-metered only; traffic bypassing the router is not counted",
}
# Merge the per-provider balance/burn/runway block. The old flat keys are
# deliberately gone: every consumer reads by_provider.
balance = quota_balance_and_burn(conn, cfg)
result["by_provider"] = balance["by_provider"]
result["total_balance_usd"] = balance["total_balance_usd"]
if reset_day is not None:
result["next_reset_date"] = _next_reset_date(reset_day)
return result
def scoring_coverage(
conn: sqlite3.Connection,
cfg: Any,
) -> dict:
"""Report which scoring axes actually have data behind them."""
placeholders = ",".join("?" * len(cfg.routing.allowed_access_levels))
routable = [
(r["model_id"], r["provider"])
for r in conn.execute(
f"""
SELECT model_id, provider FROM models
WHERE availability = 'active'
AND access_level IN ({placeholders})
""",
tuple(cfg.routing.allowed_access_levels),
)
]
with_energy = {
(r["model_id"], r["provider"])
for r in conn.execute(
"SELECT DISTINCT model_id, provider FROM energy_observations "
"WHERE task_category = ?",
(SEED_CATEGORY,),
)
}
with_proficiency = {
(r["model_id"], r["provider"])
for r in conn.execute(
"SELECT DISTINCT model_id, provider FROM proficiency "
"WHERE blended_score IS NOT NULL"
)
}
total = len(routable)
missing_energy = [m for m, p in routable if (m, p) not in with_energy]
missing_proficiency = [m for m, p in routable if (m, p) not in with_proficiency]
warnings: List[str] = []
# Sourced from quota_accounts' alarm, NOT from a second threshold of its
# own. The navbar bell and the quota chip read the same object, so they
# cannot disagree about whether anything is wrong -- which they did while
# this warning fired on ">80% of plan" (a bare percentage) and the chip
# fired on pace.
#
# Pace is the number that matters and a percentage alone actively misleads:
# 46% of the plan is comfortable on day 25 of a period and a ~3x overrun on
# day 5. The second sentence survives from the original warning because it
# is the correction this project has already had to make once -- the plan
# is a bill, not a wall, and nothing is refused when it is exceeded.
quota = quota_accounts(conn, cfg)
alarm = quota.get("alarm") or {}
if alarm.get("severity") in ("warning", "critical"):
provider = alarm.get("provider") or "a provider"
kind = alarm.get("kind")
if kind in ("plan_pace", "plan_usage"):
plan_block = next(
(
acc["plan"]
for acc in quota["accounts"]
if acc["provider"] == alarm.get("provider") and acc.get("plan")
),
None,
)
used_pct = (plan_block["used_fraction"] * 100) if plan_block else 0.0
kwh = plan_block["kwh_per_period"] if plan_block else "?"
if kind == "plan_pace":
elapsed = quota["period"].get("elapsed_fraction")
# Only when a real billing period exists. On the rolling-30d
# source elapsed_fraction is pinned to 1.0, and "at 100% of
# the period" reads as "the period is over" rather than "there
# is no period" -- the same species of misleading precision
# this warning was rewritten to stop emitting.
at_period = (
f" at {elapsed * 100:.0f}% of the period"
if elapsed and quota["period"].get("source") == "billing_reset_day"
else ""
)
lead = (
f"{provider} is running at {alarm['headline']} "
f"({used_pct:.0f}% of the {kwh} kWh plan{at_period})"
)
else:
source = quota["period"].get("source")
context = (
" (rolling 30d; no billing_reset_day configured, so pace is "
"unknown)"
if source == "30d_rolling"
else ""
)
lead = (
f"{provider} has used {used_pct:.0f}% of the {kwh} kWh "
f"plan{context}"
)
warnings.append(
f"{lead}. Overage is billed against the account's credit "
"balance; plan_kwh_per_period gates nothing."
)
else:
warnings.append(f"{provider}: {alarm['headline']}")
if missing_energy:
warnings.append(
f"{len(missing_energy)}/{total} routable models have no reference-workload "
f"observations — eco scores the neutral 0.5 for them and "
f"objective.max_energy_per_request cannot bound them. (Cost is "
f"unaffected: it is priced per request from catalog prices.) "
f"Run: python seed_energy.py --samples 7"
)
if missing_proficiency:
warnings.append(
f"{len(missing_proficiency)}/{total} routable models have no proficiency "
f"data — task_category cannot influence their ranking. "
f"Run: python eval_proficiency.py"
)
now = datetime.now(timezone.utc)
provider_rows = conn.execute(
"SELECT provider, MAX(last_updated) AS last_updated FROM models "
"WHERE availability = 'active' GROUP BY provider"
).fetchall()
multiple_providers = len(provider_rows) > 1
for row in provider_rows:
max_updated = row["last_updated"]
provider = row["provider"]
if not max_updated:
continue
try:
catalog_time = datetime.fromisoformat(max_updated)
age = now - catalog_time
if age > timedelta(days=1):
# One formatted number, not an int and a rounded fraction glued
# together with a dot: that produced "1.0.8 days ago" live,
# because f"{1}.{0.8}" keeps the fraction's own leading "0.".
days = age.total_seconds() / 86400
warning = (
f"catalog last polled {days:.1f} days ago — prices, context "
"windows and capability flags may be out of date. "
"Check: systemctl --user status llm-router-poller.timer"
)
if multiple_providers:
warning = f"[{provider}] {warning}"
warnings.append(warning)
except ValueError:
pass
ctx = context_ceilings(conn, cfg)
warnings.extend(
demand_ceiling_warnings(
conn, cfg, ctx, sorted({t for (t, _lt) in ctx})
)
)
for (tier, lat_tol), bucket in ctx.items():
if bucket["count"] == 0:
warnings.append(
f"tier {tier} has 0 eligible models for latency_tolerance={lat_tol}"
)
cap_ctx = capability_ceilings(conn, cfg)
warnings.extend(capability_demand_warnings(conn, cfg, cap_ctx))
warnings.extend(rejection_warnings(conn, cfg))
warnings.extend(content_fault_warnings(conn, cfg))
warnings.extend(classifier_degradation_warning(conn, cfg))
# Computed once and passed in: the warning is a reading of the series, not
# a second query that could disagree with the number on screen.
cache = cache_rate_series(conn, cfg)
warnings.extend(cache_rate_warnings(conn, cfg, cache))
# Premise-expiry checks: computed once and passed in, same pattern as the
# cache-rate warning so /metrics shows one consistent reading.
depth = proficiency_sample_depth_series(conn, cfg)
warnings.extend(proficiency_sample_depth_warnings(conn, cfg, depth))
spend = cumulative_spend_series(conn, cfg)
warnings.extend(cumulative_spend_warnings(conn, cfg, spend))
selection = selection_coverage(conn, cfg)
warnings.extend(selection_coverage_warnings(selection))
return {
"routable_models": total,
"with_energy_data": total - len(missing_energy),
"with_proficiency_data": total - len(missing_proficiency),
"quota": quota,
"warnings": warnings,
"flex_default": cfg.routing.default_flex_preference.value,
"selection": selection,
# The cache-rate series rides in `coverage` because /metrics assembles
# its payload in dispatcher.py, and this is a telemetry-coverage
# question anyway: which routes can even be measured, and at what rate.
"cache": cache,
# Premise-expiry series: same pattern, another telemetry-coverage
# reading — the operator needs the number that triggered the warning,
# not only the warning text.
"proficiency_sample_depth": depth,
"cumulative_spend": spend,
}
def demand_ceiling_warnings(
conn: sqlite3.Connection,
cfg: Any,
ctx: dict,
tiers: List[int],
) -> List[str]:
"""Demand-relative + escalation-hazard warnings for the given tiers.
Mirrors the logic previously inline in ``scoring_coverage`` so the admin
availability endpoint and ``scoring_coverage`` share one implementation.
``tiers`` is a list of integer tier numbers to inspect (typically the
affected tier plus the tier one step up). The caller supplies ``ctx`` from
``context_ceilings`` so the ceiling values reflect whichever model state
the caller needs evaluated.
"""
warnings: List[str] = []
# Determine which requested tiers have any observed chat demand.
tiers_with_demand: set[int] = set()
for tier in tiers:
row = conn.execute(
"SELECT MAX(required_context_tokens) AS mx FROM route_decisions"
" WHERE julianday(observed_at) > julianday('now', '-7 days')"
" AND kind = 'chat' AND task_tier = ?",
(tier,),
).fetchone()
mx_7d = row["mx"] if row else None
if mx_7d is None:
row = conn.execute(
"SELECT MAX(required_context_tokens) AS mx FROM route_decisions"
" WHERE julianday(observed_at) > julianday('now', '-30 days')"
" AND kind = 'chat' AND task_tier = ?",
(tier,),
).fetchone()
mx_7d = row["mx"] if row else None
if mx_7d is not None:
tiers_with_demand.add(tier)
for tier in tiers:
bucket = ctx.get((tier, "interactive"))
if bucket is None:
continue
ceiling = bucket["ceiling"]
if tier not in tiers_with_demand:
continue
row_cnt = conn.execute(
"SELECT COUNT(*) AS n FROM route_decisions"
" WHERE julianday(observed_at) > julianday('now', '-7 days')"
" AND kind = 'chat' AND task_tier = ?",
(tier,),
).fetchone()["n"]
window = 7 if row_cnt >= 10 else 30
row = conn.execute(
"SELECT MAX(required_context_tokens) AS mx FROM route_decisions"
f" WHERE julianday(observed_at) > julianday('now', '-{window} days')"
" AND kind = 'chat' AND task_tier = ?",
(tier,),
).fetchone()
observed_max = row["mx"] if row and row["mx"] is not None else None
if observed_max is not None and observed_max > ceiling:
warnings.append(
f"tier {tier} context ceiling ({ceiling}) is below observed max "
f"demand ({observed_max}) — requests above this may return 422"
)
esc_cfg = getattr(cfg, "escalation", None)
esc_enabled = esc_cfg and getattr(esc_cfg, "enabled", False) if esc_cfg else False
if esc_enabled:
max_tier = max(tiers) if tiers else 3
for t in tiers:
next_t = t + 1
if next_t > max_tier:
continue
t_key = (t, "interactive")
next_key = (next_t, "interactive")
if t_key not in ctx or next_key not in ctx:
continue
ceiling_next = ctx[next_key]["ceiling"]
if ceiling_next == 0:
continue
tokens = [
r["required_context_tokens"]
for r in conn.execute(
"SELECT required_context_tokens FROM route_decisions"
" WHERE julianday(observed_at) > julianday('now', '-7 days')"
" AND kind = 'chat' AND task_tier = ?",
(t,),
).fetchall()
if r["required_context_tokens"] is not None
]
p95 = _percentile(tokens, 95)
if p95 is not None and ceiling_next < p95:
warnings.append(
f"escalation hazard: tier {next_t} ceiling ({ceiling_next}) "
f"is below tier {t} p95 demand ({p95:g}) — "
"bumping a request up a tier could make it unservable"
)
return warnings
def _admin_excluded_models(conn: sqlite3.Connection) -> set[str]:
"""Model ids the operator has switched off through admin overrides.
Both ``deprecated`` and ``stale`` exclude; ``active`` clears. Must stay
in step with ``dispatcher._admin_excluded_models``, which carries the
full rationale -- the two are deliberately duplicated rather than shared,
because importing ``dispatcher`` here would reintroduce the circular
import this module exists to avoid.
Reads ``admin_model_overrides`` once per call. If the table does not
exist yet (fresh DB before the startup migration), returns an empty set.
"""
try:
return {
row["model_id"]
for row in conn.execute(
"SELECT model_id FROM admin_model_overrides "
"WHERE availability IN ('deprecated', 'stale', 'blocked')"
)
}
except sqlite3.OperationalError:
return set()
def context_ceilings(
conn: sqlite3.Connection,
cfg: Any,
) -> dict:
"""Context-window ceiling for each ``(tier, latency_tolerance)`` bucket.
Reads all model rows, groups them by ``(tier, latency_class)``, maps
the row-level ``latency_class`` to the ``latency_tolerance`` value that
``routing.select_candidates`` expects (``'standard' → INTERACTIVE``,
``'flex'` → BATCH`` — because the hard filter in
``rejection_reason`` rejects ``latency_class=="flex"`` only when
``latency_tolerance == INTERACTIVE``, see routing.py line 131), then
delegates to ``routing.select_candidates(required_context_tokens=0)``
so the count is neutralized and the ceiling is the raw max window
across each bucket's surviving models.
Admin deprecations from ``admin_model_overrides`` are applied via
``exclude_models`` so the ceiling matches the candidate set used by
live routing.
Returns ``{(tier, latency_tolerance): {"ceiling": int, "count": int}}``
for every distinct ``(tier, latency_class)`` pair present in the DB.
"""
rows = [dict(r) for r in conn.execute("SELECT * FROM models")]
exclude_models = _admin_excluded_models(conn)
return context_ceilings_with_rows(conn, cfg, rows, exclude_models=exclude_models)
def context_ceilings_with_rows(
conn: sqlite3.Connection,
cfg: Any,
rows: List[dict],
exclude_models: Optional[set[str]] = None,
require_vision: bool = False,
require_json_mode: bool = False,
) -> dict:
"""Like ``context_ceilings`` but over caller-supplied model rows.
``conn`` is still required because future versions may read energy or
proficiency and because the default ``exclude_models`` set is read from
``admin_model_overrides``. Callers that already know the exclusion set
(e.g. admin.py applying a just-committed override) may pass it explicitly.
``require_vision`` / ``require_json_mode`` gate the candidate set the
same way live routing gates it, so a capability-gated sub-ceiling goes
through this one path rather than a second copy of the filters.
"""
if exclude_models is None:
exclude_models = _admin_excluded_models(conn)
LATENCY_MAP = {
"standard": routing.INTERACTIVE,
"flex": routing.BATCH,
}
pairs: dict[tuple[int | None, str], str] = {}
for r in rows:
key = (r["tier"], r["latency_class"])
pairs[key] = LATENCY_MAP.get(r["latency_class"], routing.INTERACTIVE)
result: dict[tuple, dict[str, int]] = {}
for (tier, _), latency_tolerance in pairs.items():
candidates = routing.select_candidates(
rows,
required_context_tokens=0,
required_tier=tier,
latency_tolerance=latency_tolerance,
allowed_access_levels=cfg.routing.allowed_access_levels,
exclude_stale=True,
exclude_deprecated=True,
exclude_models=exclude_models,
min_tool_proficiency=None,
require_vision=require_vision,
require_json_mode=require_json_mode,
)
ceiling = max((r["effective_context_window"] for r in candidates), default=0)
result[(tier, latency_tolerance)] = {
"ceiling": ceiling,
"count": len(candidates),
}
return result
def capability_ceilings(
conn: sqlite3.Connection,
cfg: Any,
) -> dict:
"""Context-window ceiling per capability-gated subset, by tier.
The plain ``(tier, latency_tolerance)`` buckets from
``context_ceilings`` are blind to capability gates: during the
2026-09-04 incident the tier-1 interactive ceiling stayed at 782,324
while the vision-capable subset had collapsed to 192,500 (the only
vision rows large enough had been deprecated through admin overrides),
so every existing demand check was silent on requests that 422ed.
Vision and JSON mode are the two flags that can independently empty
the candidate set — both fail CLOSED on an unknown catalog flag, which
is what lets them shrink the set sharply — so each gets its own series
here. Deliberately NOT one bucket per capability combination; that is
a combinatorial explosion over dimensions that do not interact.
Delegates to ``context_ceilings_with_rows`` once per dimension, so every
sub-ceiling comes from the same ``routing.select_candidates`` path
live routing uses, admin deprecations included.
Returns
``{"vision": {(tier, latency_tolerance): {"ceiling", "count"}}, "json_mode":
{...}}``, one sub-dict per capability dimension, each keyed by the
``(tier, latency_tolerance)`` pair present in the DB.
"""
rows = [dict(r) for r in conn.execute("SELECT * FROM models")]
exclude_models = _admin_excluded_models(conn)
result: dict[str, dict[tuple, dict[str, int]]] = {}
for dimension in ("vision", "json_mode"):
sub_ceilings = context_ceilings_with_rows(
conn,
cfg,
rows,
exclude_models=exclude_models,
require_vision=dimension == "vision",
require_json_mode=dimension == "json_mode",
)
result[dimension] = sub_ceilings
return result
def capability_demand_warnings(
conn: sqlite3.Connection,
cfg: Any,
cap_ctx: dict,
) -> List[str]:
"""Demand-relative warnings for the capability-gated sub-ceilings.
Mirrors ``demand_ceiling_warnings`` for each dimension — 7d MAX probe
first with a 30d fallback, silence for a bucket no such request ever
hit, else window = 7 when the 7d row count is >= 10 (30 otherwise),
warning only when that window's observed max exceeds the sub-ceiling.
The one difference is what "demand" means: request rows must carry the
capability (``images`` / ``json_mode`` columns on route_decisions), so
a vision sub-ceiling is compared only against requests that actually
carry images, for the ``(tier, "interactive")`` bucket. ``cfg`` is
accepted for signature symmetry with ``demand_ceiling_warnings``; this
detector has no escalation counterpart to read from it.
The demand read is trailing-window, so a warning stays alive until a
big historical request ages out of it; the ``window: last {n}d``
annotation in the message is what keeps that self-explaining.
"""
warnings: List[str] = []
dimension_columns = {
"vision": "images",
"json_mode": "json_mode",
}
capability_nouns = {
"vision": "for requests carrying images",
"json_mode": "for requests requesting JSON mode",
}
for dimension in ("vision", "json_mode"):
buckets = cap_ctx.get(dimension) or {}
for tier in sorted({t for (t, _) in buckets}):
bucket = buckets.get((tier, "interactive"))
if bucket is None:
continue
column = dimension_columns.get(dimension)
if column is None:
continue
row = conn.execute(
"SELECT MAX(required_context_tokens) AS mx FROM route_decisions"
" WHERE julianday(observed_at) > julianday('now', '-7 days')"
f" AND kind = 'chat' AND task_tier = ? AND {column} = 1",
(tier,),
).fetchone()
mx_7d = row["mx"] if row else None
if mx_7d is None:
row = conn.execute(
"SELECT MAX(required_context_tokens) AS mx FROM route_decisions"
" WHERE julianday(observed_at) > julianday('now', '-30 days')"
f" AND kind = 'chat' AND task_tier = ? AND {column} = 1",
(tier,),
).fetchone()
mx_7d = row["mx"] if row else None
if mx_7d is None:
# No request carrying this capability was ever routed for this
# tier — an unexercised sub-ceiling warns about nothing.
continue
row_cnt = conn.execute(
"SELECT COUNT(*) AS n FROM route_decisions"
" WHERE julianday(observed_at) > julianday('now', '-7 days')"
f" AND kind = 'chat' AND task_tier = ? AND {column} = 1",
(tier,),
).fetchone()["n"]
window = 7 if row_cnt >= 10 else 30
row = conn.execute(
"SELECT MAX(required_context_tokens) AS mx FROM route_decisions"
f" WHERE julianday(observed_at) > julianday('now', '-{window} days')"
f" AND kind = 'chat' AND task_tier = ? AND {column} = 1",
(tier,),
).fetchone()
observed_max = row["mx"] if row and row["mx"] is not None else None
ceiling = bucket["ceiling"]
if observed_max is not None and observed_max > ceiling:
warnings.append(
f"{dimension}-capable tier {tier} context ceiling "
f"({ceiling}) is below observed max demand "
f"({observed_max}, window: last {window}d) "
f"{capability_nouns[dimension]} — requests above this may "
"return 422"
)
return warnings
# A novel (tier, normalized_reason) group must reach this count inside the
# alert window before warning. The 2026-09-04 incident produced exactly 2
# image rejections, so the floor must be <= 2; 1 would fire on any one-off
# genuinely-impossible request that correctly 422s.
NOVEL_GROUP_MIN_COUNT: Final = 2
# How many rejections a group needs inside the alert window before
# "new rejection pattern:" is reported. A single rejection is
# indistinguishable from a genuinely impossible request, which should 422 —
# the signal is a rate, not the existence of a rejection.
#
# Caveat: "novel" means ZERO occurrences in the prior baseline window, so
# a familiar-but-intermittent pattern can read as novel after a quiet gap
# and fire one spurious "new rejection pattern:" warning. It is
# self-clearing as the pattern's own rejections age into the baseline
# window; widen objective.rejection_warning_baseline_hours to suppress it.
def classifier_degradation_warning(conn: sqlite3.Connection, cfg: Any) -> list[str]:
"""Warn when most recent classifications came from a degraded source.
The cascade makes a local-classifier outage survivable, which is the
point — but survivable failures are the ones that go unnoticed for weeks.
Routing keeps working on borrowed categories while nothing says the
classifier has been down since Tuesday.
Silent below ``degraded_warn_min`` decisions: a share computed over three
requests is noise, and a warning that flaps on low traffic is one nobody
reads.
"""
min_total = getattr(cfg.classifier, "degraded_warn_min", 20)
threshold = getattr(cfg.classifier, "degraded_warn_threshold", 0.5)
row = conn.execute(
"""
SELECT
COUNT(*) AS total,
SUM(
CASE WHEN classification_source IN
('fallback', 'session_stale', 'session_history',
'classifier_cloud')
THEN 1 ELSE 0 END
) AS degraded
FROM route_decisions
WHERE classification_source IS NOT NULL
AND julianday(observed_at) >= julianday('now', '-24 hours')
"""
).fetchone()
total = (row["total"] if row else 0) or 0
degraded = (row["degraded"] if row else 0) or 0
if total < min_total:
return []
share = degraded / total
if share < threshold:
return []
return [
f"{share:.0%} of the last {total} classifications came from a degraded "
f"source ({degraded} of {total}). {_declined_for(conn)}Routing still "
f"works on borrowed categories, but their outcomes are excluded from "
f"proficiency."
]
def _declined_for(conn: sqlite3.Connection) -> str:
"""The recorded reasons the primary classifier's answer was not used.
The warning used to end "the local classifier has been failing", which is
wrong for the case that now dominates: under ``local_decision`` the
classifier answers and is declined for being unsure, and nothing is down.
Saying which reason is what tells the operator whether to look at the
endpoint (``transport_error``, ``timeout``) or at a threshold
(``below_confidence_min``). Empty when no reasons were recorded, which is
every row from before the column existed, and when the column is absent on
a database that has not been through dispatcher start-up.
"""
if not _has_column(conn, "route_decisions", "classifier_reject"):
return ""
rows = conn.execute(
"""
SELECT classifier_reject AS reason, COUNT(*) AS n
FROM route_decisions
WHERE classifier_reject IS NOT NULL
AND classification_source IN
('fallback', 'session_stale', 'session_history',
'classifier_cloud')
AND julianday(observed_at) >= julianday('now', '-24 hours')
GROUP BY classifier_reject
ORDER BY n DESC, classifier_reject
"""
).fetchall()
if not rows:
return ""
return "Declined for: " + ", ".join(f"{r['n']} {r['reason']}" for r in rows) + ". "
# Gate prefixes whose failure means a DERIVED field is broken, rather than the
# operator having switched something off.
#
# Everything else ``rejection_reason`` can return is either a choice
# (``access_level``, ``profile_excluded``), transient runtime state
# (``circuit_open``), or already covered by another warning (``stale`` and
# ``deprecated`` by the catalog-age check). Only ``context(...)`` and
# ``tier(...)`` say a value the poller computes has come out unusable, and only
# those can make a row permanently unselectable without anything noticing.
_STRUCTURAL_GATES: Final = ("context(", "tier(")
# The most permissive request expressible: one token of prompt, the lowest tier
# floor, batch latency (so a flex row is not counted out), no capability
# requirements. A model that cannot serve THIS cannot serve anything.
#
# One token, not zero. ``context_ceilings`` probes with zero, and that is
# precisely why it could not see the 2026-09-08 zero-context incident: a row
# with a zero-token window passes ``0 < 0`` and is counted a candidate, then
# contributes 0 to a max().
_MINIMAL_CONTEXT_TOKENS: Final = 1
def unroutable_models(conn: sqlite3.Connection, cfg: Any) -> List[dict]:
"""Active models that cannot serve even the most permissive request.
The gap this closes: **an ineligible candidate produces no signal at all.**
``capability_demand_warnings`` compares a ceiling, which is a ``max()``
over eligible models — a broken row contributes 0 to a max and is
invisible. ``rejection_warnings`` counts decisions with no selected model —
but nothing is rejected when the other rows serve the traffic perfectly
well. A row that silently fails a hard filter is an *absence*, and neither
detector can see an absence.
Found live on 2026-09-08: 12 of 30 active OpenRouter rows had an effective
context window of 0 and had never been selected once across 23,000+
decisions, two of them advertising 1M context. Nothing warned for weeks.
Delegates to ``routing.rejection_reason`` rather than re-testing the
fields, so this cannot drift away from the filters live routing applies —
the same reason ``rejection_reason`` returns a reason string instead of a
bool. Rows the operator has switched off through admin overrides are
dropped before the probe: those are deliberate, and reporting them would
make this warn about the thing it exists to distinguish from.
Returns ``[{"model_id", "provider", "reason"}]``, sorted.
"""
excluded = _admin_excluded_models(conn)
found: List[dict] = []
for r in conn.execute("SELECT * FROM models WHERE availability = 'active'"):
row = dict(r)
if row["model_id"] in excluded:
continue
# A restrict-only category gate is a narrowing, not ineligibility, so
# probe with a category the row itself admits.
eligible = routing.parse_eligible_categories(row.get("eligible_categories"))
reason = routing.rejection_reason(
row,
required_context_tokens=_MINIMAL_CONTEXT_TOKENS,
required_tier=1,
latency_tolerance=routing.BATCH,
allowed_access_levels=cfg.routing.allowed_access_levels,
exclude_stale=True,
exclude_deprecated=True,
task_category=sorted(eligible)[0] if eligible else None,
)
if reason is not None and reason.startswith(_STRUCTURAL_GATES):
found.append(
{
"model_id": row["model_id"],
"provider": row["provider"],
"reason": reason,
}
)
return sorted(found, key=lambda d: (d["provider"], d["model_id"]))
def selection_coverage(conn: sqlite3.Connection, cfg: Any) -> dict:
"""Which active models the router has actually picked, and which cannot be.
Two halves, deliberately reported differently.
``unroutable`` is a defect and warns. ``never_selected`` is a *count*, and
warns about nothing: most of a catalog legitimately never wins. Measured
on the live DB, 22 of 49 active models were selected across 24,044
decisions — enumerating the other 27 as a warning would fire permanently
and mean "you have more models than winners", which is normal. It is
reported as data so the portal can show it and the operator can prune an
allowlist, not as an alert nobody would read twice.
``by_provider`` is the same split per provider, and is also data only. The
tempting warning -- "a whole provider has zero selections, so it must be
misconfigured" -- was tried and dropped: ``ollama-local`` is dormant by
design under the default profile (docs/routing.md, local dispatch branch),
so it would fire permanently on a documented intentional state.
"""
window = getattr(cfg.objective, "selection_coverage_window_hours", None)
if window is None:
window = 168
unroutable = unroutable_models(conn, cfg)
dead = {(d["model_id"], d["provider"]) for d in unroutable}
# Access-gated rows are excluded, matching scoring_coverage's `routable`.
# A canary or preview row can never be selected under this config BY
# DESIGN, so counting it as routable-but-unselected both pads
# never_selected and makes the two routable_models figures in one /metrics
# payload disagree -- live, 51 here against 46 there, which is exactly the
# 1 canary + 4 preview rows.
placeholders = ",".join("?" * len(cfg.routing.allowed_access_levels))
active = [
(r["model_id"], r["provider"])
for r in conn.execute(
f"""
SELECT model_id, provider FROM models
WHERE availability = 'active'
AND access_level IN ({placeholders})
""",
tuple(cfg.routing.allowed_access_levels),
)
]
picked = {
(r["selected_model"], r["selected_provider"])
for r in conn.execute(
"""
SELECT DISTINCT selected_model, selected_provider
FROM route_decisions
WHERE selected_model IS NOT NULL
AND julianday(observed_at) > julianday('now', '-' || ? || ' hours')
""",
(str(window),),
)
}
decisions = conn.execute(
"""
SELECT COUNT(*) AS n FROM route_decisions
WHERE julianday(observed_at) > julianday('now', '-' || ? || ' hours')
""",
(str(window),),
).fetchone()["n"]
routable = [m for m in active if m not in dead]
never = sorted(m for m in routable if m not in picked)
by_provider: dict[str, dict] = {}
for _model_id, provider in routable:
by_provider.setdefault(provider, {"routable": 0, "selected": 0})
by_provider[provider]["routable"] += 1
for model_id, provider in routable:
if (model_id, provider) in picked:
by_provider[provider]["selected"] += 1
return {
"window_hours": window,
"decisions": decisions,
"active_models": len(active),
"routable_models": len(routable),
"selected_models": len([m for m in routable if m in picked]),
"never_selected": [{"model_id": m, "provider": p} for m, p in never],
"unroutable": unroutable,
"by_provider": by_provider,
}
# How many model ids an unroutable-group warning names before it summarizes.
# Enough to recognize the pattern, few enough to read on one line.
UNROUTABLE_MODELS_SHOWN: Final = 4
def selection_coverage_warnings(coverage: dict) -> List[str]:
"""Warnings from a ``selection_coverage`` result.
Takes the computed dict rather than the connection so ``/metrics`` does
not run the same three queries twice.
"""
warnings: List[str] = []
# One line per (provider, normalized gate), not per model. The live
# incident produced 11 rows failing the identical gate; 11 near-identical
# warnings is a wall nobody reads, and they share one root cause and one
# fix. Digit-normalized for the same reason rejection_warnings does it --
# context(0<1) and context(0<2) are the same defect.
groups: dict[tuple[str, str], List[str]] = {}
for row in coverage["unroutable"]:
key = (row["provider"], re.sub(r"\d+", "N", row["reason"]))
groups.setdefault(key, []).append(row["model_id"])
for (provider, gate), models in sorted(groups.items()):
shown = ", ".join(models[:UNROUTABLE_MODELS_SHOWN])
if len(models) > UNROUTABLE_MODELS_SHOWN:
shown += f", +{len(models) - UNROUTABLE_MODELS_SHOWN} more"
plural = "models" if len(models) > 1 else "model"
warnings.append(
f"[{provider}] {len(models)} active {plural} can never be "
f"selected: {gate} fails even for a one-token request "
f"({shown}). Nothing rejects them, because they never enter the "
f"candidate set — check the poller's derivation for those rows."
)
return warnings
def rejection_warnings(
conn: sqlite3.Connection,
cfg: Any,
) -> List[str]:
"""Reactive safety net over the rejections the router already logged.
The ceiling checks above only catch dimensions someone modelled; this
one catches everything, at the cost of firing after the first failures
rather than before. Rejections (``selected_model IS NULL AND
rejected_reason IS NOT NULL``) are pulled once across the alert plus
baseline window and split by exact age — alert bucket at
``rejection_warning_window_hours`` or younger, baseline bucket older
than the alert window but within ``rejection_warning_baseline_hours``
of it — then grouped by *normalized* reason plus ``task_tier``.
Two signals, one line per group:
* ``rejection rate:`` — the group reached
``objective.rejection_warning_min_count`` rejections in the alert
window.
* ``new rejection pattern:`` — the group is ABSENT from the baseline
window and reached ``NOVEL_GROUP_MIN_COUNT`` in the alert window.
This prefix wins when both would fire.
Zero alert-window rejections produce nothing at all: a rejection is
not inherently an error, and a genuinely impossible request should
422 — the signal is a rate.
"""
window = getattr(cfg.objective, "rejection_warning_window_hours", None)
if window is None:
window = 1
baseline = getattr(cfg.objective, "rejection_warning_baseline_hours", None)
if baseline is None:
baseline = 24
min_count = getattr(cfg.objective, "rejection_warning_min_count", None)
if min_count is None:
min_count = 6
now = datetime.now(timezone.utc)
rows = conn.execute(
"""
SELECT observed_at, task_tier, rejected_reason
FROM route_decisions
WHERE selected_model IS NULL
AND rejected_reason IS NOT NULL
AND julianday(observed_at) > julianday('now', '-' || ? || ' hours')
ORDER BY observed_at DESC
""",
(str(window + baseline),),
).fetchall()
alert_groups: dict[tuple[Optional[int], str], dict] = {}
baseline_keys: set[tuple[Optional[int], str]] = set()
for row in rows:
try:
observed = datetime.fromisoformat(row["observed_at"])
except ValueError:
# Defensive, same shape as quota_balance_and_burn: one malformed
# timestamp must not take the whole report down.
continue
# Group on digit-normalized reasons: the raw strings embed the
# request's own numbers ("context >= 242486 tokens"), so literal
# grouping yields N groups of one and hides the pattern entirely.
key = (row["task_tier"], re.sub(r"\d+", "N", row["rejected_reason"]))
age_hours = (now - observed).total_seconds() / 3600.0
if age_hours <= window:
group = alert_groups.get(key)
if group is None:
# Rows arrive newest-first, so the first touch per group is
# its most recent occurrence.
alert_groups[key] = {"count": 1, "most_recent": observed}
else:
group["count"] += 1
elif age_hours <= window + baseline:
baseline_keys.add(key)
if not alert_groups:
return []
warnings: List[str] = []
for (tier, normalized_reason), group in alert_groups.items():
count = group["count"]
is_novel = (tier, normalized_reason) not in baseline_keys
prefix: Optional[str] = None
if is_novel and count >= NOVEL_GROUP_MIN_COUNT:
prefix = "new rejection pattern:"
elif count >= min_count:
prefix = "rejection rate:"
if prefix is None:
continue
tier_label = "unknown" if tier is None else str(tier)
warnings.append(
f"{prefix} {count} rejection(s) in the last {window}h: "
f"{normalized_reason} (tier {tier_label}; "
f"most recent: {group['most_recent'].isoformat()})"
)
return warnings
def content_fault_warnings(
conn: sqlite3.Connection,
cfg: Any,
) -> List[str]:
"""Surface models with excessive malformed-output rates.
Follows the rejection_warnings pattern: alert window for fresh failures,
baseline window for comparison, novel-pattern detection for groups absent
from the baseline.
"""
window = getattr(cfg.objective, "rejection_warning_window_hours", None)
if window is None:
window = 1
baseline = getattr(cfg.objective, "rejection_warning_baseline_hours", None)
if baseline is None:
baseline = 24
min_count = getattr(cfg.objective, "rejection_warning_min_count", None)
if min_count is None:
min_count = 6
NOVEL_GROUP_MIN_COUNT = 2
# Pull malformed+attributable verifications across both windows.
rows = conn.execute(
"""
SELECT observed_at, model_id, provider
FROM verifications
WHERE verdict = 'malformed'
AND model_attributable = 1
AND julianday(observed_at) > julianday('now', '-' || ? || ' hours')
ORDER BY observed_at DESC
""",
(str(window + baseline),),
).fetchall()
alert_groups: dict[tuple[str, str], dict] = {}
baseline_keys: set[tuple[str, str]] = set()
for row in rows:
try:
observed = datetime.fromisoformat(row["observed_at"])
except ValueError:
continue
key = (row["model_id"], row["provider"])
age_hours = (datetime.now(timezone.utc) - observed).total_seconds() / 3600
if age_hours <= window:
group = alert_groups.get(key)
if group is None:
alert_groups[key] = {"count": 1, "most_recent": observed}
else:
group["count"] += 1
if observed > group["most_recent"]:
group["most_recent"] = observed
elif age_hours <= window + baseline:
baseline_keys.add(key)
if not alert_groups:
return []
warnings: List[str] = []
for (model, provider), group in alert_groups.items():
count = group["count"]
is_novel = (model, provider) not in baseline_keys
if is_novel and count >= NOVEL_GROUP_MIN_COUNT:
prefix = "new content fault:"
elif count >= min_count:
prefix = "content fault:"
else:
continue
warnings.append(
f"{prefix} {count} malformed output(s) "
f"in the last {window}h: {model} ({provider}; "
f"most recent: {group['most_recent'].isoformat()})"
)
return warnings
# Defaults for the cache-rate series, used when config omits the keys. The
# real values live in config/config.yaml with their reasoning; these exist so
# a SimpleNamespace config in a test, or a deployment on an older file, still
# produces a sane series rather than a crash.
_CACHE_RATE_WINDOW_HOURS: Final = 168
_CACHE_RATE_WARN_MARGIN: Final = 0.10
_CACHE_RATE_MIN_OBSERVATIONS: Final = 25
# Defaults for the premise-expiry checks, matching the pattern above. These
# live in config.yaml under objective:; these Final values exist so a
# SimpleNamespace or an older config still produces a sane series.
_PRF_DEPTH_WARN_MIN_SAMPLES: Final = 20
_PRF_DEPTH_WARN_MIN_ROWS: Final = 10
_CUMULATIVE_SPEND_WARN_USD: Final = 50.0
_CUMULATIVE_SPEND_WARN_MIN_ROWS: Final = 10
def _has_column(conn: sqlite3.Connection, table: str, column: str) -> bool:
"""True when *table* carries *column*, for additively-migrated schemas.
``cached_tokens_source`` arrives by ALTER at dispatcher start-up, and
``metrics`` is also imported by ``admin`` — which may open a database that
has not been through that start-up yet. A missing column must degrade the
cache series, not 500 the whole /metrics payload.
"""
return any(
r[1] == column for r in conn.execute(f"PRAGMA table_info({table})")
)
def cache_rate_series(conn: sqlite3.Connection, cfg: Any) -> dict:
"""Prefix-cache hit rate per (provider, model) over a trailing window.
``sum(cached_prompt_tokens) / sum(prompt_tokens)``, token-weighted rather
than a mean of per-request rates: the cost model prices tokens, and a
thousand-token turn and a two-hundred-thousand-token turn are not one
observation each as far as the bill is concerned.
Three filters, each of which has already been got wrong once:
* **Only rows the provider actually reported.** ``cached_prompt_tokens IS
NOT NULL`` is the operative test, because by the invariant ``af18009``
established that is exactly the ``cached_tokens_source = 'reported'``
set — and it is the test that still works on rows written *before* that
column existed, whose NULL source means "predates the column", not "the
provider said nothing". The source clause is added as a redundant
restatement of the invariant so that a future divergence between the two
fails closed rather than quietly widening the denominator.
* **No ``seed_reference`` rows.** ``seed_energy.py``'s sweep never carries
a cached count, it is not user traffic, and it was the entire reason
NeuralWatt's coverage read 12.8% when real dispatch traffic reads ~98%.
* **A trailing window, never a hardcoded start date.** Coverage begins
whenever the capture reached a given deployment; a window discovers that
instead of asserting it.
"""
window_hours = getattr(cfg.objective, "cache_rate_window_hours", None)
if window_hours is None:
window_hours = _CACHE_RATE_WINDOW_HOURS
assumed = getattr(cfg.objective, "assumed_cache_rate", None)
series: dict[str, Any] = {
"window_hours": window_hours,
"assumed_cache_rate": assumed,
"observations": 0,
"prompt_tokens": 0,
"cached_prompt_tokens": 0,
"cache_rate": None,
"by_model": [],
}
if not _has_column(conn, "energy_observations", "cached_prompt_tokens"):
return series
source_clause = ""
if _has_column(conn, "energy_observations", "cached_tokens_source"):
source_clause = (
" AND (cached_tokens_source IS NULL "
"OR cached_tokens_source = 'reported')"
)
rows = conn.execute(
f"""
SELECT provider, model_id,
COUNT(*) AS observations,
SUM(prompt_tokens) AS prompt_tokens,
SUM(cached_prompt_tokens) AS cached_prompt_tokens
FROM energy_observations
WHERE cached_prompt_tokens IS NOT NULL
AND prompt_tokens IS NOT NULL
AND prompt_tokens > 0
AND (task_category IS NULL OR task_category != ?)
{source_clause}
AND julianday(observed_at) > julianday('now', '-' || ? || ' hours')
GROUP BY provider, model_id
""",
(SEED_CATEGORY, str(window_hours)),
).fetchall()
total_obs = 0
total_prompt = 0
total_cached = 0
by_model: List[dict] = []
for row in rows:
prompt = row["prompt_tokens"] or 0
cached = row["cached_prompt_tokens"] or 0
obs = row["observations"] or 0
total_obs += obs
total_prompt += prompt
total_cached += cached
by_model.append(
{
"provider": row["provider"],
"model_id": row["model_id"],
"observations": obs,
"prompt_tokens": prompt,
"cached_prompt_tokens": cached,
"cache_rate": (cached / prompt) if prompt else None,
}
)
by_model.sort(key=lambda d: (d["provider"], d["model_id"]))
series["observations"] = total_obs
series["prompt_tokens"] = total_prompt
series["cached_prompt_tokens"] = total_cached
series["cache_rate"] = (total_cached / total_prompt) if total_prompt else None
series["by_model"] = by_model
return series
def cache_rate_warnings(
conn: sqlite3.Connection,
cfg: Any,
series: Optional[dict] = None,
) -> List[str]:
"""Expiry check on ``objective.assumed_cache_rate``.
The question is NOT "is the cache working" — a bare floor would answer
that one and it is the wrong one. ``assumed_cache_rate`` is a number the
cost model is *built on*: ``routing.estimated_cost`` prices a 100k-token
prompt almost entirely through it, and it was measured once, on
2026-08-23, and never re-measured. So the signal is DIVERGENCE from the
configured value, in either direction — a real rate of 0.99 mis-prices
just as surely as 0.60, it simply mis-prices the other way.
Two classes, the same shape ``rejection_warnings`` uses and for the same
reasons:
* ``cache rate:`` — the aggregate trailing rate is more than
``objective.cache_rate_warn_margin`` from the assumed one. This is the
premise-expiry check, and it is the one that decides whether the cost
model still describes this deployment.
* ``cache rate outlier:`` — one ``(provider, model_id)`` group diverges by
more than the margin. Grouped on the **structured columns**, never on a
formatted label: the aggregate can sit comfortably inside the margin
while a single route caches at half the rate, which is precisely the
model-switching leak the waves plan is chasing, and a formatted key
would hide a provider serving the same base model twice.
Neither fires on mere presence. Both require
``objective.cache_rate_warn_min_observations`` reported-cache rows —
aggregate or per group — the way ``classifier_degradation_warning``
requires ``degraded_warn_min`` before computing a share at all. A rate
over three requests is not a measurement, and a warning that flaps on low
traffic is one nobody reads.
"""
if series is None:
series = cache_rate_series(conn, cfg)
assumed = series.get("assumed_cache_rate")
if assumed is None:
# Nothing to be expired against.
return []
margin = getattr(cfg.objective, "cache_rate_warn_margin", None)
if margin is None:
margin = _CACHE_RATE_WARN_MARGIN
min_obs = getattr(cfg.objective, "cache_rate_warn_min_observations", None)
if min_obs is None:
min_obs = _CACHE_RATE_MIN_OBSERVATIONS
window = series["window_hours"]
warnings: List[str] = []
measured = series.get("cache_rate")
total_obs = series.get("observations") or 0
if measured is not None and total_obs >= min_obs:
delta = measured - assumed
if abs(delta) > margin:
direction = "above" if delta > 0 else "below"
warnings.append(
f"cache rate: measured {measured:.3f} over {total_obs} "
f"reported-cache observations (last {window}h) is "
f"{abs(delta):.3f} {direction} objective.assumed_cache_rate "
f"({assumed:.3f}) — routing.estimated_cost prices every "
f"request through that constant, so the ranking is running on "
f"a premise this deployment no longer supports. Re-measure it "
f"and update config.yaml."
)
for entry in series.get("by_model") or []:
rate = entry["cache_rate"]
obs = entry["observations"]
if rate is None or obs < min_obs:
continue
delta = rate - assumed
if abs(delta) <= margin:
continue
direction = "above" if delta > 0 else "below"
warnings.append(
f"cache rate outlier: {entry['model_id']} ({entry['provider']}) "
f"measured {rate:.3f} over {obs} reported-cache observations "
f"(last {window}h), {abs(delta):.3f} {direction} "
f"objective.assumed_cache_rate ({assumed:.3f})"
)
return warnings
def proficiency_sample_depth_series(
conn: sqlite3.Connection, cfg: Any
) -> dict:
"""Average sample depth (outcome_samples + self_eval_samples) per category.
Draws from the full proficiency table (cumulative, no trailing window). The
depth per row follows ``proficiency_matrix`` at the top of this file:
``samples = (self_eval_samples or 0) + (outcome_samples or 0)``.
Returns overall and per-category averages, plus thin row counts.
"""
min_samples = getattr(
cfg.objective, "proficiency_depth_warn_min_samples", None
)
if min_samples is None:
min_samples = _PRF_DEPTH_WARN_MIN_SAMPLES
series: dict = {
"overall_avg_depth": None,
"total_rows": 0,
"thin_rows": 0,
"by_category": [],
}
rows = conn.execute(
"""
SELECT category,
COUNT(*) AS cnt,
AVG(COALESCE(self_eval_samples, 0)
+ COALESCE(outcome_samples, 0)) AS avg_depth,
SUM(CASE WHEN COALESCE(self_eval_samples, 0)
+ COALESCE(outcome_samples, 0) < ? THEN 1 ELSE 0 END)
AS thin
FROM proficiency
GROUP BY category
""",
(min_samples,),
).fetchall()
total_rows = 0
total_sampled_rows = 0
depth_sum = 0.0
thin_total = 0
by_cat: list[dict] = []
for row in rows:
cnt = row["cnt"] or 0
avg = row["avg_depth"]
thin = row["thin"] or 0
total_rows += cnt
thin_total += thin
by_cat.append(
{
"category": row["category"],
"rows": cnt,
"avg_depth": avg,
"thin": thin,
}
)
if avg is not None:
total_sampled_rows += cnt
depth_sum += avg * cnt
by_cat.sort(key=lambda d: d["category"])
series["total_rows"] = total_rows
series["thin_rows"] = thin_total
series["by_category"] = by_cat
series["overall_avg_depth"] = (
(depth_sum / total_sampled_rows) if total_sampled_rows else None
)
return series
def proficiency_sample_depth_warnings(
conn: sqlite3.Connection,
cfg: Any,
series: Optional[dict] = None,
) -> List[str]:
"""Expiry check on ``objective.quality_tolerance``'s thin-data justification.
The config.yaml comment on ``quality_tolerance`` says "scores rest on 2-3
samples per category... Narrow it as samples accumulate." Once average
sample depth has grown past ``objective.proficiency_depth_warn_min_samples``
(default 20), that justification is clearly obsolete and /metrics warns that
the premise has expired.
Respects a minimum-observation floor
(``objective.proficiency_depth_warn_min_rows``) so a near-empty table does
not fire.
"""
if series is None:
series = proficiency_sample_depth_series(conn, cfg)
min_rows = getattr(cfg.objective, "proficiency_depth_warn_min_rows", None)
if min_rows is None:
min_rows = _PRF_DEPTH_WARN_MIN_ROWS
min_samples = getattr(
cfg.objective, "proficiency_depth_warn_min_samples", None
)
if min_samples is None:
min_samples = _PRF_DEPTH_WARN_MIN_SAMPLES
warnings: List[str] = []
total_rows = series.get("total_rows") or 0
avg_depth = series.get("overall_avg_depth")
if avg_depth is None or total_rows < min_rows:
return warnings
if avg_depth > min_samples:
warnings.append(
f"quality_tolerance premise expired: average proficiency sample "
f"depth {avg_depth:.1f} across {total_rows} rows has passed the "
f"warning threshold ({min_samples}), so the '2-3 samples per "
f"category' justification for objective.quality_tolerance={getattr(cfg.objective, 'quality_tolerance', '?')} "
f"is obsolete — the loose tolerance can be narrowed."
)
for entry in series.get("by_category") or []:
cat_avg = entry.get("avg_depth")
cat_rows = entry.get("rows") or 0
if cat_avg is None or cat_rows < min_rows:
continue
if cat_avg > min_samples:
warnings.append(
f"quality_tolerance premise expired: category "
f"{entry['category']} average depth {cat_avg:.1f} over "
f"{cat_rows} rows has passed the warning threshold "
f"({min_samples}) — objective.quality_tolerance premise "
f"is obsolete for this category."
)
return warnings
def cumulative_spend_series(conn: sqlite3.Connection, cfg: Any) -> dict:
"""All-time SUM(cost_usd) from energy_observations, excluding seed_reference.
Mirrors the cache-rate series' exclusion pattern: ``task_category !=
'seed_reference'`` keeps reference sweeps from inflating "real traffic".
Also surfaces per-provider breakdown and row counts. NULL cost_usd is
treated as 0.
"""
series: dict = {
"total_spend_usd": 0.0,
"priced_rows": 0,
"total_rows": 0,
"by_provider": [],
}
total_row = conn.execute(
"""
SELECT COUNT(*) AS total_rows,
COUNT(cost_usd) AS priced_rows,
COALESCE(SUM(cost_usd), 0) AS total_spend
FROM energy_observations
WHERE task_category IS NULL OR task_category != ?
""",
(SEED_CATEGORY,),
).fetchone()
total_rows = total_row["total_rows"] or 0
priced_rows = total_row["priced_rows"] or 0
total_spend = float(total_row["total_spend"] or 0.0)
provider_rows = conn.execute(
"""
SELECT provider,
COUNT(*) AS total_rows,
COUNT(cost_usd) AS priced_rows,
COALESCE(SUM(cost_usd), 0) AS provider_spend
FROM energy_observations
WHERE task_category IS NULL OR task_category != ?
GROUP BY provider
ORDER BY provider
""",
(SEED_CATEGORY,),
).fetchall()
series["total_spend_usd"] = round(total_spend, 4)
series["priced_rows"] = priced_rows
series["total_rows"] = total_rows
by_provider: list[dict] = []
for row in provider_rows:
by_provider.append(
{
"provider": row["provider"],
"total_rows": row["total_rows"] or 0,
"priced_rows": row["priced_rows"] or 0,
"provider_spend_usd": round(
float(row["provider_spend"] or 0.0), 4
),
}
)
series["by_provider"] = by_provider
return series
def cumulative_spend_warnings(
conn: sqlite3.Connection,
cfg: Any,
series: Optional[dict] = None,
) -> List[str]:
"""Expiry check on the cost-as-tiebreak comment's $0.07 claim.
The config.yaml comment on the blend section says "60% of every decision
adjudicated fractions of a cent (all real traffic to date totals $0.07)".
Once all-time spend exceeds ``objective.cumulative_spend_warn_usd``
(default $50.00), that justification no longer matches reality. Respects a
minimum-observation floor so a DB with no priced rows does not fire.
"""
if series is None:
series = cumulative_spend_series(conn, cfg)
warn_usd = getattr(cfg.objective, "cumulative_spend_warn_usd", None)
if warn_usd is None:
warn_usd = _CUMULATIVE_SPEND_WARN_USD
min_rows = getattr(cfg.objective, "cumulative_spend_warn_min_rows", None)
if min_rows is None:
min_rows = _CUMULATIVE_SPEND_WARN_MIN_ROWS
warnings: List[str] = []
priced_rows = series.get("priced_rows") or 0
total_spend = series.get("total_spend_usd") or 0.0
if priced_rows < min_rows:
return warnings
if total_spend > warn_usd:
warnings.append(
f"cost-as-tiebreak premise expired: cumulative spend "
f"${total_spend:.2f} across {priced_rows} priced rows has "
f"passed the warning threshold (${warn_usd:.2f}), so the "
f"blend comment's '$0.07 / fractions of a cent' justification "
f"in config.yaml no longer matches reality — it is off by "
f"{total_spend / 0.07:.0f}x."
)
for entry in series.get("by_provider") or []:
provider_spend = entry.get("provider_spend_usd") or 0.0
provider_priced = entry.get("priced_rows") or 0
if provider_spend > warn_usd and provider_priced >= min_rows:
warnings.append(
f"cost-as-tiebreak premise expired: provider "
f"{entry['provider']} spend ${provider_spend:.2f} has "
f"passed the warning threshold (${warn_usd:.2f}) — "
f"the blend comment's 'fractions of a cent' frame is "
f"stale for this provider."
)
return warnings
# Defaults for the two report-only series below, used when config omits the
# keys. Same purpose as the cache-rate defaults above: a SimpleNamespace config
# in a test, or a deployment on an older config file, still produces a series
# rather than a crash.
_COST_CALIBRATION_WINDOW_HOURS: Final = 168
_COST_CALIBRATION_MIN_OBSERVATIONS: Final = 25
_LATENCY_WINDOW_HOURS: Final = 168
_LATENCY_MIN_OBSERVATIONS: Final = 25
def cost_estimate_calibration(conn: sqlite3.Connection, cfg: Any) -> dict:
"""Per-``(provider, model_id)`` error between the cost ESTIMATE and the bill.
``routing.estimated_cost`` prices a request from catalog token prices
scaled to its shape. Measured against the provider's own billed figure at
~100% telemetry coverage, it is wrong by 1.6x-13x depending on the model.
**The scale error is not what this measures for.** A uniform 4x
overestimate reorders nothing — the ranking is a comparison, and every
candidate moves together. The 5x-plus SPREAD in that error does reorder:
on the 2026-09-12 reading the estimator priced
``deepseek/deepseek-v4-flash`` (3,028 micro-USD) below ``qwen3.6-35b``
(4,118) while the bill said the reverse (662 vs 511). So ``spread`` is the
headline figure here, not ``correction_factor``.
**Nothing applies any of this.** The factors are reported so their
STABILITY can be judged before anything is built on them; this project has
a standing lesson about a per-model figure that looked like a constant and
was a moment in time (CLAUDE.md, "Attribution drifts across hours"). A
factor computed once is one sample of a quantity that may move.
Four filters, each earning its place:
* **Join on ``request_id``, and require the model to match too.** The
request id alone is not enough: on the live database 5 of 4,143 joined
rows have a decision that selected ``qwen3.6-35b`` (neuralwatt) against
a completion billed by ``qwen/qwen3.6-35b-a3b`` (openrouter) — a
cross-provider failover. Those rows pair an estimate for one model with
a bill for another, which is exactly the contamination this series
cannot afford.
* **``cost_usd > 0``.** A zero or NULL bill is a row the provider did not
price, not a free request; a ratio against it is a division by
nothing.
* **No ``seed_reference``.** Sweep traffic is a fixed 400/400 shape that
no real request resembles, and it has skewed aggregates in this project
before. Checked on both sides of the join: ``route_decisions`` happens
to hold zero such rows today, which is a property of the current sweep
path and not a guarantee.
* **A trailing window, never a hardcoded start date**, for the same reason
``cache_rate_series`` uses one.
Groups below ``objective.cost_calibration_min_observations`` are still
listed — the count is itself information — but carry
``sufficient: False`` and are excluded from ``spread``, so a factor
computed over three requests cannot become the headline.
"""
window_hours = getattr(cfg.objective, "cost_calibration_window_hours", None)
if window_hours is None:
window_hours = _COST_CALIBRATION_WINDOW_HOURS
min_obs = getattr(cfg.objective, "cost_calibration_min_observations", None)
if min_obs is None:
min_obs = _COST_CALIBRATION_MIN_OBSERVATIONS
series: dict[str, Any] = {
"window_hours": window_hours,
"min_observations": min_obs,
"observations": 0,
"est_cost_usd": 0.0,
"billed_cost_usd": 0.0,
"overestimate_ratio": None,
"correction_factor": None,
"spread": None,
"spread_low": None,
"spread_high": None,
"by_model": [],
}
rows = conn.execute(
"""
SELECT e.provider AS provider,
e.model_id AS model_id,
COUNT(*) AS observations,
SUM(d.est_cost_usd) AS est_cost_usd,
SUM(e.cost_usd) AS billed_cost_usd
FROM route_decisions d
JOIN energy_observations e ON e.request_id = d.request_id
WHERE d.request_id IS NOT NULL
AND d.est_cost_usd IS NOT NULL
AND d.est_cost_usd > 0
AND e.cost_usd IS NOT NULL
AND e.cost_usd > 0
AND e.model_id = d.selected_model
AND e.provider = d.selected_provider
AND (e.task_category IS NULL OR e.task_category != ?)
AND (d.task_category IS NULL OR d.task_category != ?)
AND julianday(e.observed_at) > julianday('now', '-' || ? || ' hours')
GROUP BY e.provider, e.model_id
""",
(SEED_CATEGORY, SEED_CATEGORY, str(window_hours)),
).fetchall()
total_obs = 0
total_est = 0.0
total_billed = 0.0
by_model: List[dict] = []
for row in rows:
obs = row["observations"] or 0
est = row["est_cost_usd"] or 0.0
billed = row["billed_cost_usd"] or 0.0
total_obs += obs
total_est += est
total_billed += billed
by_model.append(
{
"provider": row["provider"],
"model_id": row["model_id"],
"observations": obs,
"est_cost_usd": est,
"billed_cost_usd": billed,
"est_per_request_usd": (est / obs) if obs else None,
"billed_per_request_usd": (billed / obs) if obs else None,
# est/billed, so it reads the same way round as the table in
# the brief: >1 means the estimator is charging the ranker
# more than the provider charged the account.
"overestimate_ratio": (est / billed) if billed else None,
# What would be multiplied INTO the estimate to land on the
# bill, if anything ever applied it. Nothing does.
"correction_factor": (billed / est) if est else None,
"sufficient": obs >= min_obs,
}
)
by_model.sort(key=lambda d: (d["provider"], d["model_id"]))
series["observations"] = total_obs
series["est_cost_usd"] = total_est
series["billed_cost_usd"] = total_billed
series["overestimate_ratio"] = (
(total_est / total_billed) if total_billed else None
)
series["correction_factor"] = (total_billed / total_est) if total_est else None
series["by_model"] = by_model
# The reordering figure: how far apart the best- and worst-estimated
# models are, over groups with enough observations to mean anything. One
# qualifying group has a spread of 1.0 by definition and names itself at
# both ends, which is correct — nothing is being mis-ordered relative to
# nothing.
trusted = [
entry
for entry in by_model
if entry["sufficient"] and entry["overestimate_ratio"] is not None
]
if trusted:
low = min(trusted, key=lambda d: d["overestimate_ratio"])
high = max(trusted, key=lambda d: d["overestimate_ratio"])
series["spread"] = high["overestimate_ratio"] / low["overestimate_ratio"]
series["spread_low"] = {
"provider": low["provider"],
"model_id": low["model_id"],
"overestimate_ratio": low["overestimate_ratio"],
}
series["spread_high"] = {
"provider": high["provider"],
"model_id": high["model_id"],
"overestimate_ratio": high["overestimate_ratio"],
}
return series
def latency_series(conn: sqlite3.Connection, cfg: Any) -> dict:
"""Router-observed wall clock and time-to-first-token, p50/p95 per model.
``energy_observations.router_wall_seconds`` and ``router_ttft_seconds``
are the router's OWN clock, kept deliberately apart from the provider's
``duration_seconds`` (see config/schema.sql for the exact span each
measures). That distinction is the whole reason this series can exist:
OpenRouter reports no ``duration_seconds`` at all, so before these columns
landed an OpenRouter model's latency was simply unobservable — and
``z-ai/glm-5.3-flash``, a model with *flash* in its name, was measured at
a p50 of 9.6s to first token against 1.6s for ``deepseek-v4-flash`` on
NeuralWatt.
**This is a report, not an objective.** Nothing in the ranking reads it.
Latency has never been a scored axis here, and making it one on the
strength of a first reading would repeat the mistake the eco sweep
documented — one sample of a quantity that tracks pool load and time of
day.
p50 and p95 rather than a mean, because the interesting failure is the
tail: the same model's wall clock spans 11s at the median and 54s at p95
on live data, and a mean hides which of those an operator is waiting on.
Two filters:
* **No ``seed_reference``.** Redundant today and kept anyway: the sweep
does not go through the dispatch path, so its rows carry no wall clock
at all (0 of 4,068 on the live database). Should that ever change, the
fixed 400/400 shape must not land in a latency percentile.
* **A trailing window**, matching the rest of the series here.
The two percentile pairs have INDEPENDENT counts. ``router_ttft_seconds``
is streaming-only by nature, so a deployment doing buffered work has
fewer TTFT rows than wall rows, and a group can be ``sufficient`` on one
and not the other. Reporting one count for both would let a handful of
streamed requests borrow the credibility of a large buffered sample.
"""
window_hours = getattr(cfg.objective, "latency_window_hours", None)
if window_hours is None:
window_hours = _LATENCY_WINDOW_HOURS
min_obs = getattr(cfg.objective, "latency_min_observations", None)
if min_obs is None:
min_obs = _LATENCY_MIN_OBSERVATIONS
series: dict[str, Any] = {
"window_hours": window_hours,
"min_observations": min_obs,
"wall_observations": 0,
"ttft_observations": 0,
"wall_p50": None,
"wall_p95": None,
"ttft_p50": None,
"ttft_p95": None,
"by_model": [],
}
# Both columns arrive by ALTER at dispatcher start-up, and `metrics` is
# also imported by `admin`, which may open a database that has not been
# through that start-up. A missing column degrades this series, never the
# whole /metrics payload -- the same contract `cache_rate_series` keeps.
if not _has_column(conn, "energy_observations", "router_wall_seconds"):
return series
rows = conn.execute(
"""
SELECT provider, model_id, router_wall_seconds, router_ttft_seconds
FROM energy_observations
WHERE (router_wall_seconds IS NOT NULL
OR router_ttft_seconds IS NOT NULL)
AND (task_category IS NULL OR task_category != ?)
AND julianday(observed_at) > julianday('now', '-' || ? || ' hours')
""",
(SEED_CATEGORY, str(window_hours)),
).fetchall()
grouped: dict[tuple[str, str], dict[str, list[float]]] = {}
all_wall: list[float] = []
all_ttft: list[float] = []
for row in rows:
key = (row["provider"], row["model_id"])
bucket = grouped.setdefault(key, {"wall": [], "ttft": []})
wall = row["router_wall_seconds"]
ttft = row["router_ttft_seconds"]
if wall is not None:
bucket["wall"].append(wall)
all_wall.append(wall)
if ttft is not None:
bucket["ttft"].append(ttft)
all_ttft.append(ttft)
by_model: List[dict] = []
for (provider, model_id), bucket in grouped.items():
wall = bucket["wall"]
ttft = bucket["ttft"]
by_model.append(
{
"provider": provider,
"model_id": model_id,
"wall_observations": len(wall),
"ttft_observations": len(ttft),
"wall_p50": _percentile(wall, 50),
"wall_p95": _percentile(wall, 95),
"ttft_p50": _percentile(ttft, 50),
"ttft_p95": _percentile(ttft, 95),
"wall_sufficient": len(wall) >= min_obs,
"ttft_sufficient": len(ttft) >= min_obs,
}
)
by_model.sort(key=lambda d: (d["provider"], d["model_id"]))
series["wall_observations"] = len(all_wall)
series["ttft_observations"] = len(all_ttft)
series["wall_p50"] = _percentile(all_wall, 50)
series["wall_p95"] = _percentile(all_wall, 95)
series["ttft_p50"] = _percentile(all_ttft, 50)
series["ttft_p95"] = _percentile(all_ttft, 95)
series["by_model"] = by_model
return series
# The prefix-stability probe's columns (pinch.prefix_probe), which arrive by
# ALTER at dispatcher start-up. Named here rather than inlined because
# recent_decisions has to both SELECT them conditionally and backfill the keys
# they would have produced, and the two lists drifting apart would surface as
# a silently missing field rather than an error.
_PREFIX_PROBE_COLUMNS: Final = (
"prefix_divergence_index",
"prefix_tokens_after_divergence",
"prefix_prev_message_count",
)
# Conversation-identity columns (agent, parent_key) that arrive by ALTER at
# dispatcher start-up. Probed and backfilled the same way as the prefix-probe
# columns above so an un-migrated DB still loads.
_CONVERSATION_IDENTITY_COLUMNS: Final = (
"agent",
"parent_key",
)
# The classifier's own attempt (confidence, coverage, why it was rejected).
# Added by ALTER at dispatcher start-up like the groups above, so probed and
# backfilled the same way: the live router.db has none of them until its next
# restart, and a missing column must read as NULL rather than 500 the payload.
_CLASSIFIER_ATTEMPT_COLUMNS: Final = (
"classifier_confidence",
"classifier_coverage",
"classifier_reject",
)
def recent_decisions(
conn: sqlite3.Connection,
limit: int = 50,
) -> List[dict]:
"""Last *N* rows from route_decisions, ordered by id DESC.
The row is the field set the TUI consumes; it is pinned by
tests/test_tui_schema_drift.py.
The prefix-probe columns are selected only when they exist. ``metrics`` is
imported by ``admin`` as well as by ``dispatcher``, so it can be handed a
database that has not been through dispatcher start-up and therefore never
ran the ALTER -- the live router.db is exactly that until its next restart.
A missing column must leave the field NULL, not 500 the whole payload.
"""
present = [
column
for column in _PREFIX_PROBE_COLUMNS
if _has_column(conn, "route_decisions", column)
]
identity_present = [
column
for column in _CONVERSATION_IDENTITY_COLUMNS
if _has_column(conn, "route_decisions", column)
]
attempt_present = [
column
for column in _CLASSIFIER_ATTEMPT_COLUMNS
if _has_column(conn, "route_decisions", column)
]
probe_select = "".join(f", {column}" for column in present)
identity_select = "".join(f", {column}" for column in identity_present)
attempt_select = "".join(f", {column}" for column in attempt_present)
rows = [
dict(row)
for row in conn.execute(
f"""
SELECT id, observed_at, kind, task_category, task_tier,
required_context_tokens, confidence, classifier_ms,
classification_source, latency_tolerance,
candidates_considered, selected_model, selected_provider,
runner_up_models, est_cost_usd, est_proficiency,
rejected_reason, session_key, tools, images, json_mode, streamed,
flex_preference, flex_swapped, flex_forced,
exploration, request_id,
pinch_original_tokens, pinch_final_tokens, profile{probe_select}{identity_select}{attempt_select}
FROM route_decisions
ORDER BY id DESC
LIMIT ?
""",
(limit,),
).fetchall()
]
# Present the same keys either way: a consumer reading a pre-migration
# database sees NULLs, which is what the column would have held anyway,
# rather than a KeyError or a quietly absent field.
for row in rows:
for column in _PREFIX_PROBE_COLUMNS:
row.setdefault(column, None)
for column in _CONVERSATION_IDENTITY_COLUMNS:
row.setdefault(column, None)
for column in _CLASSIFIER_ATTEMPT_COLUMNS:
row.setdefault(column, None)
return rows
def conversation_adoption(
conn: sqlite3.Connection,
cfg: Any,
) -> dict:
"""Conversation-adoption counters over route_decisions.
This is the (H1) adoption metric: of the routing decisions in the
window, how many carry a proper conversation session key (``c:...``)
rather than a simple fingerprint hash. Returns:
* ``n_c_conversations`` - COUNT(DISTINCT session_key) over the rows
whose session_key starts with ``c:`` (one per conversation).
* ``n_c_decisions`` - the row count over those same rows (one per turn;
a former version named this row count ``n_c_conversations``, which
over-counted multi-turn conversations).
* ``n_decisions`` - all route_decisions rows in the same window,
fingerprint rows included; this denominates ``share``.
* ``share`` - ``n_c_decisions / n_decisions``, or None when the window
holds no decisions at all.
The window is ``objective.adoption_window_seconds``; null or absent
means all time; 0 is rejected by the config validator. All three counts
read the same window, so ``share`` is a true in-window ratio.
The capability gate checks ``route_decisions.session_key`` only (PRAGMA
table_info): a database that has not run the identity migration yet
yields zeros, not an error.
SQLite LIKE is case-insensitive by default, fine here because the ``c:``
prefix is the only non-c delimiter and fingerprints never contain a colon.
Switch to ``LIKE 'c:%' ESCAPE '\'`` only if a future prefix needs escaping.
"""
window_seconds = getattr(cfg.objective, "adoption_window_seconds", None)
if not _has_column(conn, "route_decisions", "session_key"):
return {
"n_c_conversations": 0,
"n_c_decisions": 0,
"n_decisions": 0,
"share": None,
"window_seconds": window_seconds,
}
params: list = []
where_clause = "1 = 1"
if window_seconds and window_seconds > 0:
where_clause = (
"julianday(observed_at) > "
"julianday('now', '-' || ? || ' seconds')"
)
params.append(str(window_seconds))
sql = f"""
SELECT SUM(CASE WHEN session_key LIKE 'c:%' THEN 1 ELSE 0 END)
AS n_c_decisions,
COUNT(DISTINCT CASE WHEN session_key LIKE 'c:%'
THEN session_key END)
AS n_c_conversations,
COUNT(*) AS n_decisions
FROM route_decisions
WHERE {where_clause}
"""
if params:
row = conn.execute(sql, tuple(params)).fetchone()
else:
row = conn.execute(sql).fetchone()
# A no-GROUP-BY aggregate SELECT always yields exactly one row, but SUM
# over that empty set is NULL in SQLite -- coerce it, COUNT never is.
n_c_decisions = row["n_c_decisions"] or 0
n_decisions = row["n_decisions"]
return {
"n_c_conversations": row["n_c_conversations"],
"n_c_decisions": n_c_decisions,
"n_decisions": n_decisions,
"share": (n_c_decisions / n_decisions) if n_decisions else None,
"window_seconds": window_seconds,
}
def per_model(conn: sqlite3.Connection) -> List[dict]:
"""Per-model aggregates over the last 30 d of energy_observations."""
return [
dict(row)
for row in conn.execute(
"""
SELECT model_id,
provider,
COUNT(*) AS calls,
SUM(cost_usd) AS sum_cost_usd,
SUM(energy_kwh) AS sum_energy_kwh,
SUM(carbon_g_co2eq) AS sum_carbon_g_co2eq,
SUM(prompt_tokens) AS sum_prompt_tokens,
SUM(completion_tokens) AS sum_completion_tokens,
SUM(cached_prompt_tokens) AS sum_cached_prompt_tokens,
SUM(CASE WHEN cost_usd IS NOT NULL THEN 1 ELSE 0 END) AS cost_rows,
COUNT(*) AS total_rows,
AVG(completion_tokens) AS avg_completion_tokens,
AVG(attribution_ratio) AS avg_attribution_ratio
FROM energy_observations
WHERE julianday(observed_at) > julianday('now', '-30 days')
GROUP BY model_id, provider
"""
).fetchall()
]
def verdict_mix(
conn: sqlite3.Connection,
since_days: int = 7,
) -> dict:
"""Counts by verdict from verifications in the last *since_days*."""
rows = conn.execute(
"""
SELECT verdict, COUNT(*) n
FROM verifications
WHERE julianday(observed_at) > julianday('now', '-' || ? || ' days')
GROUP BY verdict
""",
(str(since_days),),
).fetchall()
return {row["verdict"]: row["n"] for row in rows}
def decision_outcome_summary(
conn: sqlite3.Connection,
since_days: int = 7,
) -> dict:
"""Route decisions and client-reported outcomes in the last *since_days*.
``decisions`` counts ``route_decisions`` rows (requests the router routed);
the ``client_*`` counts come only from ``verifications`` rows with
``kind = 'client_outcome'`` -- ground truth about the client's experience,
not the internal verification verdicts that :func:`verdict_mix` aggregates.
"""
decisions = conn.execute(
"""
SELECT COUNT(*) n
FROM route_decisions
WHERE julianday(observed_at) > julianday('now', '-' || ? || ' days')
""",
(str(since_days),),
).fetchone()["n"]
reports = conn.execute(
"""
SELECT
COUNT(*) client_reports,
SUM(CASE WHEN verdict = 'succeeded' THEN 1 ELSE 0 END) client_ok,
SUM(CASE WHEN verdict = 'failed' THEN 1 ELSE 0 END) client_failed
FROM verifications
WHERE kind = 'client_outcome'
AND julianday(observed_at) > julianday('now', '-' || ? || ' days')
""",
(str(since_days),),
).fetchone()
return {
"decisions": decisions,
"client_reports": reports["client_reports"],
"client_ok": reports["client_ok"] or 0,
"client_failed": reports["client_failed"] or 0,
}
def top_proficiency(
conn: sqlite3.Connection,
category: str,
) -> List[dict]:
"""Top models by blended_score for *category*, ordered DESC."""
return [
dict(row)
for row in conn.execute(
"""
SELECT model_id, provider, blended_score, source,
self_eval_samples
FROM proficiency
WHERE category = ?
AND blended_score IS NOT NULL
ORDER BY blended_score DESC
""",
(category,),
).fetchall()
]
def proficiency_matrix(conn: sqlite3.Connection, cfg) -> dict:
"""Every proficiency row, with the provenance the score's weight depends on.
``proficiency_score`` is the only category-dependent term in the ranking,
so this table decides routing — and the portal's entire surface for it was
a top-N list for one category on the dashboard. An operator could see
which model won and not why.
Four columns carry the weight and all four ship, because a bare
``blended_score`` is not comparable across rows:
* ``source`` — an ``outcome_prior`` is INHERITED from a trafficked
sibling, not measured here. Treating those as equal to
``outcome_blended`` is the mistake this project already made once with
``-flex`` rows, where an inherited score froze at its first copy while
the row it came from moved from 0.85 to 0.973.
* ``outcome_samples`` / ``self_eval_samples`` — a 1.00 at n=2 and a 0.97
at n=14 are different claims. ``docs_writing`` read 0.70-1.00 at n=2 and
0.66-0.97 at n=11-14, and the router paid for that ceiling.
* ``inherited_from`` — which family a variant borrowed from. NULL means
measured on this row.
* ``last_updated`` — recent traffic, or a months-old sweep.
``thin`` is computed here rather than in the page so the threshold comes
from ``proficiency.self_eval_min_samples`` and cannot drift into a
hardcoded number in JavaScript.
Read-only, and deliberately no write path exists: proficiency is derived
from evaluation and client outcomes, so a hand-edited score is a
fabricated measurement — the same failure as the empty leaderboards.yaml
and the provider's static_fallback carbon constant this project excludes.
"""
min_samples = cfg.proficiency.self_eval_min_samples
rows = [
dict(row)
for row in conn.execute(
"""
SELECT model_id, provider, category, blended_score,
leaderboard_score, self_eval_score, self_eval_samples,
outcome_score, outcome_samples, source, inherited_from,
last_updated
FROM proficiency
ORDER BY model_id, category
"""
).fetchall()
]
for row in rows:
samples = (row.get("self_eval_samples") or 0) + (
row.get("outcome_samples") or 0
)
row["total_samples"] = samples
row["thin"] = samples < min_samples
# Two independent claims, and the page colours on the union: a row can
# be inherited without being thin (it copies a well-sampled parent's
# count) and thin without being inherited.
row["inherited"] = row.get("inherited_from") is not None
by_source: dict[str, int] = {}
for row in rows:
key = row.get("source") or "unknown"
by_source[key] = by_source.get(key, 0) + 1
return {
"rows": rows,
"categories": list(cfg.proficiency.categories),
"models": sorted({(r["model_id"]) for r in rows}),
"self_eval_min_samples": min_samples,
"source_counts": by_source,
}
def recent_client_outcomes(
conn: sqlite3.Connection, limit: int = 50
) -> List[dict]:
"""The latest ``kind='client_outcome'`` verifications, attribution included.
``POST /outcome`` is the only ground truth the router gets — every other
signal is a proxy — and nothing in the portal showed what had been
reported. Worse, degraded-classification outcomes are recorded with
``model_attributable = 0`` and excluded from folding, so a silently
discarded report looked exactly like an applied one.
"Recorded but never surfaced" is a pattern this project has now hit four
times, so ``model_attributable`` and ``applied_at`` both ship: whether the
report counts, and whether it has been folded in yet.
"""
return [
dict(row)
for row in conn.execute(
"""
SELECT id, model_id, provider, request_id, task_category, verdict,
detail, completion_tokens, observed_at, applied_at,
model_attributable
FROM verifications
WHERE kind = 'client_outcome'
ORDER BY observed_at DESC, id DESC
LIMIT ?
""",
(int(limit),),
).fetchall()
]
def local_energy_summary(conn: sqlite3.Connection, cfg) -> Optional[dict]:
"""Aggregate local energy observations over the last 30 days.
Mirrors ``quota_burn``'s window so the two reports share the same reset
date. Returns None when local energy metering is disabled so callers can
omit the section entirely.
"""
if not cfg.local_energy.enabled:
return None
total = conn.execute(
"""
SELECT SUM(energy_kwh) kwh,
SUM(cost_usd) cost,
COUNT(*) n
FROM local_energy_observations
WHERE julianday(observed_at) > julianday('now', '-30 days')
"""
).fetchone()
by_type_rows = conn.execute(
"""
SELECT call_type,
COUNT(*) calls,
SUM(energy_kwh) kwh,
SUM(cost_usd) cost
FROM local_energy_observations
WHERE julianday(observed_at) > julianday('now', '-30 days')
GROUP BY call_type
"""
).fetchall()
by_type = {
row["call_type"]: {
"calls": row["calls"],
"kwh": round(float(row["kwh"] or 0), 5),
"cost_usd": round(float(row["cost"] or 0), 5),
}
for row in by_type_rows
}
result = {
"metered_kwh_30d": round(float(total["kwh"] or 0), 5),
"metered_cost_usd_30d": round(float(total["cost"] or 0), 5),
"calls_30d": total["n"],
"by_type": by_type,
"reset_date": None,
"window_start_30d": (datetime.now(timezone.utc).date() - timedelta(days=30)).isoformat(),
}
reset_day = getattr(cfg.objective, "billing_reset_day", None)
if reset_day is not None:
result["reset_date"] = _billing_period_start(reset_day)
result["next_reset_date"] = _next_reset_date(reset_day)
return result
def pinch_summary(conn: sqlite3.Connection, cfg: Any) -> Optional[dict]:
"""Aggregate context-pruning savings over the last 30 days.
Returns None when pinch is disabled so callers can omit the section.
Computes the number of pruned calls, share of pruned calls, total and
median tokens saved, and an estimated dollar saving priced at the blended
prompt rate the router uses for cost estimation.
"""
if not cfg.pinch.enabled:
return None
cache_rate = cfg.objective.assumed_cache_rate
totals = conn.execute(
"""
SELECT COUNT(*) AS calls_30d,
SUM(CASE
WHEN pinch_original_tokens IS NOT NULL
AND pinch_final_tokens IS NOT NULL
AND pinch_original_tokens > pinch_final_tokens
THEN 1 ELSE 0
END) AS pruned_calls_30d,
SUM(CASE
WHEN pinch_original_tokens IS NOT NULL
AND pinch_final_tokens IS NOT NULL
THEN pinch_original_tokens - pinch_final_tokens
ELSE 0
END) AS total_tokens_saved
FROM route_decisions
WHERE julianday(observed_at) > julianday('now', '-30 days')
"""
).fetchone()
calls_30d = totals["calls_30d"] or 0
pruned_calls_30d = totals["pruned_calls_30d"] or 0
share_pruned = pruned_calls_30d / calls_30d if calls_30d else 0.0
total_tokens_saved = totals["total_tokens_saved"] or 0
saved_rows = conn.execute(
"""
SELECT pinch_original_tokens - pinch_final_tokens AS tokens_saved
FROM route_decisions
WHERE julianday(observed_at) > julianday('now', '-30 days')
AND pinch_original_tokens IS NOT NULL
AND pinch_final_tokens IS NOT NULL
AND pinch_original_tokens > pinch_final_tokens
"""
).fetchall()
median_tokens_saved = _percentile(
[r["tokens_saved"] for r in saved_rows], 50
)
dollar_rows = conn.execute(
"""
SELECT rd.pinch_original_tokens - rd.pinch_final_tokens AS tokens_saved,
m.cost_per_1m_prompt,
m.cost_per_1m_prompt_cached
FROM route_decisions rd
JOIN models m
ON rd.selected_model = m.model_id
AND rd.selected_provider = m.provider
WHERE julianday(rd.observed_at) > julianday('now', '-30 days')
AND rd.pinch_original_tokens IS NOT NULL
AND rd.pinch_final_tokens IS NOT NULL
AND rd.pinch_original_tokens > rd.pinch_final_tokens
AND rd.selected_model IS NOT NULL
AND m.cost_per_1m_prompt IS NOT NULL
"""
).fetchall()
dollars_saved_usd_30d = 0.0
for r in dollar_rows:
saved = r["tokens_saved"]
prompt_price = r["cost_per_1m_prompt"]
cached_price = r["cost_per_1m_prompt_cached"]
if cached_price is None:
cached_price = prompt_price
blended = (1.0 - cache_rate) * prompt_price + cache_rate * cached_price
dollars_saved_usd_30d += saved * blended / 1_000_000
result = {
"calls_30d": calls_30d,
"pruned_calls_30d": pruned_calls_30d,
"share_pruned": round(share_pruned, 4),
"total_tokens_saved": total_tokens_saved,
"median_tokens_saved": median_tokens_saved,
"dollars_saved_usd_30d": round(dollars_saved_usd_30d, 6),
"reset_date": None,
"window_start_30d": (datetime.now(timezone.utc).date() - timedelta(days=30)).isoformat(),
}
reset_day = getattr(cfg.objective, "billing_reset_day", None)
if reset_day is not None:
result["reset_date"] = _billing_period_start(reset_day)
result["next_reset_date"] = _next_reset_date(reset_day)
return result