"""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