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
3479 lines
141 KiB
Python
3479 lines
141 KiB
Python
"""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
|
||
|