From 8d57cb86482dfa2f5e72b6d48c492ca40e3320ae Mon Sep 17 00:00:00 2001 From: Ashen <310210685+ashen0x@users.noreply.github.com> Date: Fri, 18 Sep 2026 20:20:05 +0530 Subject: [PATCH] feat: connect stats collection and sharing to node runtime --- docs/analytics-v2-format.md | 10 +- routstr/core/admin.py | 51 +- routstr/core/ledger_analytics.py | 512 ++++++++++++ routstr/core/log_manager.py | 31 +- routstr/core/main.py | 14 +- routstr/core/settings.py | 72 +- routstr/core/usage_analytics_store.py | 9 +- routstr/nostr/analytics_runtime.py | 232 ++++++ routstr/nostr/listing.py | 38 + .../test_admin_settings_endpoint.py | 108 ++- tests/unit/test_analytics_runtime.py | 283 +++++++ tests/unit/test_ledger_analytics.py | 745 ++++++++++++++++++ tests/unit/test_nostr_stats_identity.py | 60 ++ tests/unit/test_settings.py | 10 +- tests/unit/test_terminal_outcomes.py | 54 ++ tests/unit/test_usage_analytics_store_utc.py | 34 + ui/app/page.tsx | 116 ++- ui/components/settings/admin-settings.tsx | 8 +- ui/components/top-models-usage-chart.tsx | 96 ++- ui/components/usage-metrics-chart.tsx | 81 +- ui/components/usage-summary-cards.tsx | 65 +- ui/lib/api/services/admin.ts | 50 +- ui/lib/usage-time.ts | 7 + 23 files changed, 2551 insertions(+), 135 deletions(-) create mode 100644 routstr/core/ledger_analytics.py create mode 100644 routstr/nostr/analytics_runtime.py create mode 100644 tests/unit/test_analytics_runtime.py create mode 100644 tests/unit/test_ledger_analytics.py create mode 100644 tests/unit/test_nostr_stats_identity.py create mode 100644 tests/unit/test_usage_analytics_store_utc.py create mode 100644 ui/lib/usage-time.ts diff --git a/docs/analytics-v2-format.md b/docs/analytics-v2-format.md index 5a388159..c94f9171 100644 --- a/docs/analytics-v2-format.md +++ b/docs/analytics-v2-format.md @@ -1,9 +1,17 @@ # Completed-day analytics reports -Public sharing is an explicit operator choice. Reports contain aggregate usage for completed UTC days. They do not contain prompts, customer identifiers, request identifiers, served upstream model identifiers, or private pricing diagnostics. +Public sharing follows the node's existing analytics setting. Reports contain aggregate usage for completed UTC days. They do not contain prompts, customer identifiers, request identifiers, served upstream model identifiers, or private pricing diagnostics. A public-sharing activation begins coverage on the next full UTC day. Toggle days and collection-loss days are excluded. Private terminal records remain stored when public sharing is disabled. An omitted day means unavailable coverage; an included all-zero day means the collector covered the day and recorded no completed requests. +## Operator controls and upgrades + +The existing "Share analytics to Nostr" switch controls public reports through `ENABLE_ANALYTICS_SHARING`. Its existing default remains `true`; an explicitly saved `false` remains off after upgrading or restarting. There is no separate format opt-in. + +Local collection starts automatically for the private dashboard and continues when public sharing is off. An upgraded node uses completed-day reports whenever its existing sharing setting is on. New coverage begins with the next full UTC day, so private history from before activation is not published. The runtime no longer starts the legacy publisher; consumers can still read legacy reports already on relays. + +A provider's signing identity must remain stable, and an identity change creates a new publication transition. Turning sharing off stops pending delivery but cannot retract reports already published. + ## Signed event - Nostr kind: `38422`. diff --git a/routstr/core/admin.py b/routstr/core/admin.py index 365e2f68..fbe66f82 100644 --- a/routstr/core/admin.py +++ b/routstr/core/admin.py @@ -1,7 +1,7 @@ import json import re import secrets -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from pathlib import Path from fastapi import APIRouter, Depends, HTTPException, Query, Request @@ -35,6 +35,7 @@ from .db import ( store_cashu_transaction_with_retry as store_cashu_transaction, ) from .exceptions import json_compliant +from .ledger_analytics import get_ledger_usage_dashboard from .log_manager import log_manager from .logging import get_logger from .provider_slugs import allocate_unique_provider_slug @@ -1597,6 +1598,42 @@ async def get_openrouter_presets() -> list[dict[str, object]]: return models_data +async def _usage_dashboard( + interval: int, + hours: int, + error_limit: int = 100, + model_limit: int = 20, + start_at: datetime | None = None, + end_at: datetime | None = None, +) -> dict: + if (start_at is None) != (end_at is None): + raise HTTPException(400, "Provide both start_at and end_at") + if start_at is not None and end_at is not None: + if start_at.tzinfo is None or end_at.tzinfo is None: + raise HTTPException(400, "Stats dates must include a timezone") + now = datetime.now(timezone.utc) + tomorrow = now.replace(hour=0, minute=0, second=0, microsecond=0) + timedelta( + days=1 + ) + if end_at > tomorrow or start_at >= min(end_at, now): + raise HTTPException(400, "Choose a non-empty stats period through today") + if end_at - start_at > timedelta(hours=MAX_USAGE_ANALYTICS_HOURS): + raise HTTPException(400, "Stats periods cannot exceed 365 days") + end_at = min(end_at, now) + dashboard = log_manager.get_usage_dashboard( + interval=interval, hours=hours, error_limit=error_limit, model_limit=model_limit + ) + return await get_ledger_usage_dashboard( + dashboard, + interval=interval, + hours=hours, + model_limit=model_limit, + session_factory=create_session, + start_at=start_at, + end_at=end_at, + ) + + @admin_router.get("/api/usage/metrics", dependencies=[Depends(require_admin_api)]) async def get_usage_metrics( request: Request, @@ -1611,7 +1648,7 @@ async def get_usage_metrics( ), ) -> dict: """Get usage metrics aggregated by time interval.""" - return log_manager.get_usage_metrics(interval=interval, hours=hours) + return (await _usage_dashboard(interval, hours))["metrics"] @admin_router.get("/api/usage/dashboard", dependencies=[Depends(require_admin_api)]) @@ -1632,16 +1669,20 @@ async def get_usage_dashboard( model_limit: int = Query( default=20, ge=1, le=100, description="Maximum number of models to return" ), + start_at: datetime | None = Query(default=None), + end_at: datetime | None = Query(default=None), ) -> dict: """ Get all dashboard analytics in one request. This runs one combined aggregation pass and avoids repeated scans. """ - return log_manager.get_usage_dashboard( + return await _usage_dashboard( interval=interval, hours=hours, error_limit=error_limit, model_limit=model_limit, + start_at=start_at, + end_at=end_at, ) @@ -1656,7 +1697,7 @@ async def get_usage_summary( ), ) -> dict: """Get summary statistics for the specified time period.""" - return log_manager.get_usage_summary(hours=hours) + return (await _usage_dashboard(15, hours))["summary"] @admin_router.get("/api/usage/error-details", dependencies=[Depends(require_admin_api)]) @@ -1694,7 +1735,7 @@ async def get_revenue_by_model( """ Get revenue breakdown by model. """ - return log_manager.get_revenue_by_model(hours=hours, limit=limit) + return (await _usage_dashboard(15, hours, model_limit=limit))["revenue_by_model"] @admin_router.get("/api/logs", dependencies=[Depends(require_admin_api)]) diff --git a/routstr/core/ledger_analytics.py b/routstr/core/ledger_analytics.py new file mode 100644 index 00000000..6308473c --- /dev/null +++ b/routstr/core/ledger_analytics.py @@ -0,0 +1,512 @@ +from __future__ import annotations + +from copy import deepcopy +from datetime import UTC, datetime, timedelta +from typing import Any + +from sqlalchemy import case +from sqlmodel import col, func, select + +from . import terminal_outcomes +from .db import ( + TerminalOutcome, + TerminalOutcomeEpoch, + TerminalOutcomeWriterRun, + create_session, +) +from .terminal_outcome_writer import SessionFactory + +_SETTLED_FIELDS = ( + "successful_chat_completions", + "revenue_msats", + "input_tokens", + "output_tokens", + "cache_read_input_tokens", + "cache_creation_input_tokens", + "total_tokens", +) +_DIAGNOSTIC_FIELDS = ( + "total_requests", + "failed_requests", + "errors", + "warnings", + "payment_processed", + "upstream_errors", + "refunds_msats", + "requests", +) +_DIMENSIONS = ("input", "output", "cache_read", "cache_creation") + + +def _measures() -> list[Any]: + input_tokens = ( + col(TerminalOutcome.input_tokens) + + col(TerminalOutcome.cache_read_input_tokens) + + col(TerminalOutcome.cache_creation_input_tokens) + ) + return [ + func.count(col(TerminalOutcome.outcome_id)).label( + "successful_chat_completions" + ), + func.sum(TerminalOutcome.revenue_msats).label("revenue_msats"), + func.sum(input_tokens).label("input_tokens"), + func.sum(TerminalOutcome.output_tokens).label("output_tokens"), + func.sum(TerminalOutcome.cache_read_input_tokens).label( + "cache_read_input_tokens" + ), + func.sum(TerminalOutcome.cache_creation_input_tokens).label( + "cache_creation_input_tokens" + ), + func.sum(input_tokens + col(TerminalOutcome.output_tokens)).label( + "total_tokens" + ), + ] + + +def _values(row: Any) -> dict[str, int]: + return {name: int(getattr(row, name) or 0) for name in _SETTLED_FIELDS} + + +def _timestamp(bucket: int, bucket_ms: int) -> str: + return datetime.fromtimestamp(bucket * bucket_ms / 1000, UTC).strftime( + "%Y-%m-%d %H:%M:%S" + ) + + +def _coverage_days( + start: datetime, end: datetime, epochs: list[TerminalOutcomeEpoch], clock: datetime +) -> list[str]: + missing = [] + day = start.date() + last = (end - timedelta(milliseconds=1)).date() + while day <= last: + day_end = datetime.combine(day + timedelta(days=1), datetime.min.time(), UTC) + if min(day_end, end) > clock or not any( + epoch.coverage_start_day <= day + and (epoch.coverage_end_day is None or day <= epoch.coverage_end_day) + for epoch in epochs + ): + missing.append(day.isoformat()) + day += timedelta(days=1) + return missing + + +async def get_ledger_usage_dashboard( + legacy_dashboard: dict[str, Any], + *, + interval: int, + hours: int, + model_limit: int = 20, + session_factory: SessionFactory = create_session, + now: datetime | None = None, + start_at: datetime | None = None, + end_at: datetime | None = None, +) -> dict[str, Any]: + """Overlay settled records while retaining log-only request diagnostics.""" + clock = now or datetime.now(UTC) + if clock.tzinfo is None: + clock = clock.replace(tzinfo=UTC) + clock = clock.astimezone(UTC) + end = end_at or clock + if end.tzinfo is None: + end = end.replace(tzinfo=UTC) + end = end.astimezone(UTC) + start = start_at or end - timedelta(hours=hours) + if start.tzinfo is None: + start = start.replace(tzinfo=UTC) + start = start.astimezone(UTC) + if end <= start: + raise ValueError("Stats period end must be after its start") + diagnostic_available = start_at is None and end_at is None + window_hours = (end - start).total_seconds() / 3600 + start_ms, end_ms = int(start.timestamp() * 1000), int(end.timestamp() * 1000) + bucket_ms = max(1, interval) * 60_000 + limit = max(1, min(model_limit, 100)) + top_limit = min(limit, 20) + bucket = (col(TerminalOutcome.terminal_at_ms) // bucket_ms).label("bucket") + model = func.nullif(TerminalOutcome.model_identifier, "").label("model") + bounds = ( + col(TerminalOutcome.terminal_day) >= start.date(), + col(TerminalOutcome.terminal_day) <= end.date(), + col(TerminalOutcome.terminal_at_ms) >= start_ms, + col(TerminalOutcome.terminal_at_ms) + < min(end_ms, int(clock.timestamp() * 1000)), + ) + measures = _measures() + total_tokens = measures[-1] + provenance = [] + sources = {} + for dimension in _DIMENSIONS: + source = col(getattr(TerminalOutcome, dimension + "_source")) + sources[dimension] = source + for status in ("reported", "estimated", "missing"): + provenance.append( + func.sum(case((source == status, 1), else_=0)).label( + f"{dimension}_{status}" + ) + ) + measured = (sources["input"] == "reported") & (sources["output"] == "reported") + for dimension in ("cache_read", "cache_creation"): + cache_tokens = col(getattr(TerminalOutcome, dimension + "_input_tokens")) + measured &= (sources[dimension] == "reported") | ( + (sources[dimension] == "missing") & (cache_tokens == 0) + ) + recorded_tokens = ( + col(TerminalOutcome.input_tokens) + + col(TerminalOutcome.output_tokens) + + col(TerminalOutcome.cache_read_input_tokens) + + col(TerminalOutcome.cache_creation_input_tokens) + ) + async with session_factory() as session: + bucket_rows = ( + await session.exec( + select(bucket, *measures) + .where(*bounds) + .group_by(bucket) + .order_by(bucket) + ) + ).all() + top_models: dict[str, list[str]] = {} + for metric, measure in ( + ("requests", measures[0]), + ("revenue", measures[1]), + ("tokens", total_tokens), + ): + rows = ( + await session.exec( + select(model, measure) + .where(*bounds, model.is_not(None)) + .group_by(model) + .having(measure > 0) + .order_by(measure.desc(), model) + .limit(top_limit) + ) + ).all() + top_models[metric] = [str(row[0]) for row in rows] + selected = sorted(set(name for names in top_models.values() for name in names)) + model_bucket_rows = ( + ( + await session.exec( + select(bucket, model, *measures) + .where(*bounds, model.in_(selected)) + .group_by(bucket, model) + .order_by(bucket, model) + ) + ).all() + if selected + else [] + ) + model_rows = ( + await session.exec( + select(model, *measures) + .where(*bounds) + .group_by(model) + .order_by(measures[1].desc(), model) + .limit(limit) + ) + ).all() + model_count = ( + await session.exec(select(func.count(func.distinct(model))).where(*bounds)) + ).one() + observation = ( + ( + await session.execute( + select( + func.max(TerminalOutcome.terminal_at_ms).label("latest"), + func.max(case((model.is_(None), 1), else_=0)).label( + "unattributed" + ), + func.sum(case((measured, 1), else_=0)).label( + "measured_token_requests" + ), + func.sum(case((measured, recorded_tokens), else_=0)).label( + "measured_tokens" + ), + *provenance, + ).where(*bounds) + ) + ) + .mappings() + .one() + ) + epochs = list( + ( + await session.exec( + select(TerminalOutcomeEpoch) + .where(col(TerminalOutcomeEpoch.coverage_start_day) <= end.date()) + .where( + col(TerminalOutcomeEpoch.coverage_end_day).is_(None) + | (col(TerminalOutcomeEpoch.coverage_end_day) >= start.date()) + ) + ) + ).all() + ) + + runs = list( + ( + await session.exec( + select(TerminalOutcomeWriterRun).where( + col(TerminalOutcomeWriterRun.status).in_( + ("active", "degraded", "lost") + ) + ) + ) + ).all() + ) + + # The log manager caches its dictionaries. Keep ledger overlays out of that cache. + result = deepcopy(legacy_dashboard) + if not diagnostic_available: + result["metrics"] = { + "metrics": [], + "totals": {name: 0 for name in _DIAGNOSTIC_FIELDS}, + } + summary = result.setdefault("summary", {}) + for name in ( + *_DIAGNOSTIC_FIELDS, + "total_entries", + "total_errors", + "total_warnings", + "refunds_sats", + "success_rate", + "refund_rate", + ): + summary[name] = 0 + summary["error_types"] = {} + result["error_details"] = {"errors": [], "total_count": 0} + result["revenue_by_model"] = {"models": []} + totals = {name: 0 for name in _SETTLED_FIELDS} + metrics_by_time: dict[str, dict[str, Any]] = {} + for point in result.get("metrics", {}).get("metrics", []): + at = datetime.fromisoformat(point["timestamp"]) + if at.tzinfo is None: + at = at.replace(tzinfo=UTC) + timestamp_ms = int(at.timestamp() * 1000) + if start_ms // bucket_ms * bucket_ms <= timestamp_ms < end_ms: + stamp = _timestamp(timestamp_ms // bucket_ms, bucket_ms) + metrics_by_time[stamp] = {**point, "timestamp": stamp, **totals} + mix: dict[str, dict[str, Any]] = {} + for row in bucket_rows: + stamp = _timestamp(int(row.bucket), bucket_ms) + values = _values(row) + for name in _SETTLED_FIELDS: + totals[name] += values[name] + point = metrics_by_time.setdefault( + stamp, {name: 0 for name in _DIAGNOSTIC_FIELDS} + ) + point.update(timestamp=stamp, **values) + mix[stamp] = { + "timestamp": stamp, + "total_successful": values["successful_chat_completions"], + "total_revenue_msats": values["revenue_msats"], + "total_tokens": values["total_tokens"], + "others": values["successful_chat_completions"], + "others_revenue_msats": values["revenue_msats"], + "others_tokens": values["total_tokens"], + "model_counts": {}, + "model_revenue_msats": {}, + "model_tokens": {}, + } + for row in model_bucket_rows: + point = mix[_timestamp(int(row.bucket), bucket_ms)] + for target, value, others in ( + ("model_counts", int(row.successful_chat_completions), "others"), + ("model_revenue_msats", int(row.revenue_msats), "others_revenue_msats"), + ("model_tokens", int(row.total_tokens), "others_tokens"), + ): + point[target][str(row.model)] = value + point[others] -= value + metric_points = sorted( + metrics_by_time.values(), key=lambda point: point["timestamp"] + ) + result["metrics"] = { + **result.get("metrics", {}), + "metrics": metric_points, + "interval_minutes": interval, + "hours_back": window_hours, + "total_buckets": len(metric_points), + "totals": {**result.get("metrics", {}).get("totals", {}), **totals}, + } + summary = result.setdefault("summary", {}) + completed = totals["successful_chat_completions"] + measured_requests = int(observation["measured_token_requests"] or 0) + measured_tokens = int(observation["measured_tokens"] or 0) + summary.update(totals) + summary.update( + unique_models_count=int(model_count), + unique_models=selected, + unique_models_truncated=len(selected) < int(model_count), + revenue_sats=totals["revenue_msats"] / 1000, + # Retained for older API clients; ledger revenue never subtracts hold releases. + net_revenue_msats=totals["revenue_msats"], + net_revenue_sats=totals["revenue_msats"] / 1000, + avg_input_tokens_per_completion=totals["input_tokens"] / completed + if completed + else 0, + avg_output_tokens_per_completion=totals["output_tokens"] / completed + if completed + else 0, + avg_total_tokens_per_completion=totals["total_tokens"] / completed + if completed + else 0, + measured_token_requests=measured_requests, + measured_tokens=measured_tokens, + avg_measured_tokens_per_completion=measured_tokens / measured_requests + if measured_requests + else None, + avg_revenue_per_request_msats=totals["revenue_msats"] / completed + if completed + else 0, + ) + legacy_models = { + row["model"]: row + for row in result.get("revenue_by_model", {}).get("models", []) + } + revenue_models = [] + for row in model_rows: + name = str(row.model) if row.model is not None else "unknown" + previous = legacy_models.get(name, {}) + sats = int(row.revenue_msats) / 1000 + count = int(row.successful_chat_completions) + revenue_models.append( + { + **previous, + "model": name, + "revenue_sats": sats, + "net_revenue_sats": sats, + "refunds_sats": previous.get("refunds_sats", 0), + "requests": previous.get("requests", 0), + "successful": count, + "failed": previous.get("failed", 0), + "avg_revenue_per_request": sats / count if count else 0, + } + ) + result["revenue_by_model"] = { + "models": revenue_models, + "total_revenue_sats": totals["revenue_msats"] / 1000, + "total_models": int(model_count) + int(bool(observation["unattributed"])), + } + result["model_usage_mix"] = { + "top_models": top_models["requests"], + "top_models_by_metric": top_models, + "metrics": list(mix.values()), + "interval_minutes": interval, + "hours_back": window_hours, + "total_buckets": len(mix), + } + incomplete_days = _coverage_days(start, end, epochs, clock) + checkpoints = [run.flushed_through_ms or run.started_at_ms for run in runs] + # A writer always trails the clock; only a checkpoint left in an earlier day + # is a gap. Today's buckets past the checkpoint are still updating. + unsafe_days = [ + day + for day in ( + datetime.fromtimestamp(checkpoint / 1000, UTC).date() + for checkpoint in checkpoints + ) + if day < clock.date() + ] + unsafe_days.extend(run.loss_day for run in runs if run.loss_day is not None) + writer = terminal_outcomes.terminal_outcome_writer + if writer.loss_pending and writer.loss_day is not None: + unsafe_days.append(writer.loss_day) + if not runs and any(epoch.current_slot == 1 for epoch in epochs): + unsafe_days.append(clock.date()) + if unsafe_days: + day = max(start.date(), min(unsafe_days)) + last = (end - timedelta(milliseconds=1)).date() + while day <= last: + incomplete_days.append(day.isoformat()) + day += timedelta(days=1) + incomplete_days = sorted(set(incomplete_days)) + result["analytics_source"] = "terminal_outcomes" + result["ledger_coverage"] = { + "from": start.isoformat(), + "to": end.isoformat(), + "complete": not incomplete_days, + "incomplete_days": incomplete_days, + "includes_current_day": start.date() + <= clock.date() + <= (end - timedelta(milliseconds=1)).date(), + "diagnostic_available": diagnostic_available, + "flushed_through": datetime.fromtimestamp( + min(checkpoints) / 1000, UTC + ).isoformat() + if checkpoints + else None, + "latest_outcome_at": datetime.fromtimestamp( + observation["latest"] / 1000, UTC + ).isoformat() + if observation["latest"] is not None + else None, + "token_sources": { + dimension: { + status: int(observation[f"{dimension}_{status}"] or 0) + for status in ("reported", "estimated", "missing") + } + for dimension in _DIMENSIONS + }, + } + first_bucket, last_bucket = start_ms // bucket_ms, (end_ms - 1) // bucket_ms + fill_complete = last_bucket - first_bucket + 1 <= 1000 + if fill_complete: + stamps = [ + _timestamp(value, bucket_ms) + for value in range(first_bucket, last_bucket + 1) + ] + else: + stamps = sorted(metrics_by_time.keys() | mix.keys()) + missing_days = set(incomplete_days) + for stamp in stamps: + bucket_start = datetime.fromisoformat(stamp).replace(tzinfo=UTC) + day = max(start, bucket_start).date() + last_day = ( + min(end, bucket_start + timedelta(milliseconds=bucket_ms)) + - timedelta(milliseconds=1) + ).date() + covered = True + while day <= last_day: + covered = covered and day.isoformat() not in missing_days + day += timedelta(days=1) + recorded = stamp in mix + bucket_end_ms = min(end_ms, int(bucket_start.timestamp() * 1000) + bucket_ms) + if not covered: + coverage = "partial" if recorded else "missing" + elif checkpoints and bucket_end_ms > min(checkpoints): + coverage = "updating" + else: + coverage = "complete" + empty = 0 if covered else None + point = metrics_by_time.setdefault( + stamp, + { + "timestamp": stamp, + **{name: 0 for name in _DIAGNOSTIC_FIELDS}, + }, + ) + if not recorded: + point.update({name: empty for name in _SETTLED_FIELDS}) + point["coverage"] = coverage + model_point = mix.setdefault( + stamp, + { + "timestamp": stamp, + "total_successful": empty, + "total_revenue_msats": empty, + "total_tokens": empty, + "others": empty, + "others_revenue_msats": empty, + "others_tokens": empty, + "model_counts": {}, + "model_revenue_msats": {}, + "model_tokens": {}, + }, + ) + model_point["coverage"] = coverage + for key, points in (("metrics", metrics_by_time), ("model_usage_mix", mix)): + result[key].update( + metrics=[points[stamp] for stamp in stamps], + total_buckets=len(stamps), + bucket_fill_complete=fill_complete, + ) + return result diff --git a/routstr/core/log_manager.py b/routstr/core/log_manager.py index 0444dcbf..b884648e 100644 --- a/routstr/core/log_manager.py +++ b/routstr/core/log_manager.py @@ -103,7 +103,10 @@ class LogManager: # If we only care about hours back, we can optimize file selection if hours_back is not None: - cutoff_date = datetime.now(timezone.utc) - timedelta(hours=hours_back) + # Log stamps and file dates are server-local. + cutoff_date = ( + datetime.now(timezone.utc) - timedelta(hours=hours_back) + ).astimezone() cutoff_timestamp_str = cutoff_date.strftime("%Y-%m-%d %H:%M:%S") filtered_files = [] for log_path in log_files: @@ -111,7 +114,7 @@ class LogManager: file_date_str = log_path.stem.split("_")[1] file_date = datetime.strptime( file_date_str, "%Y-%m-%d" - ).replace(tzinfo=timezone.utc) + ).replace(tzinfo=cutoff_date.tzinfo) # Include file if it's from the same day or after the cutoff day if file_date >= cutoff_date.replace( hour=0, minute=0, second=0, microsecond=0 @@ -270,22 +273,20 @@ class LogManager: def _bucket_key_for_timestamp( self, timestamp_str: str, interval_minutes: int ) -> str | None: - if len(timestamp_str) != 19: - return None - if timestamp_str[10] != " ": - return None - try: - hour = int(timestamp_str[11:13]) - minute = int(timestamp_str[14:16]) - except (TypeError, ValueError): + # Server-local stamps are bucketed in UTC, like the indexed store. + at = datetime.strptime(timestamp_str, "%Y-%m-%d %H:%M:%S").astimezone( + timezone.utc + ) + except ValueError: return None - total_minutes = hour * 60 + minute - rounded_minutes = (total_minutes // interval_minutes) * interval_minutes - rounded_hour = rounded_minutes // 60 - rounded_minute = rounded_minutes % 60 - return f"{timestamp_str[:10]} {rounded_hour:02d}:{rounded_minute:02d}:00" + rounded_minutes = ( + (at.hour * 60 + at.minute) // interval_minutes * interval_minutes + ) + return ( + f"{at:%Y-%m-%d} {rounded_minutes // 60:02d}:{rounded_minutes % 60:02d}:00" + ) def _extract_success_metrics( self, entry: dict[str, Any], message: str diff --git a/routstr/core/main.py b/routstr/core/main.py index 46aad1fc..2d8d3f29 100644 --- a/routstr/core/main.py +++ b/routstr/core/main.py @@ -26,7 +26,11 @@ from ..lightning import ( from ..nostr import ( announce_provider, providers_cache_refresher, - publish_usage_analytics, +) +from ..nostr.analytics_runtime import ( + prepare_analytics, + run_analytics, + shutdown_analytics, ) from ..nostr.discovery import providers_router from ..payment.models import models_router, update_sats_pricing @@ -118,6 +122,11 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: await reset_all_reserved_balances(session) + try: + await prepare_analytics() + except Exception: + logger.exception("Stats collection could not start; requests remain available") + # Apply app metadata from settings try: app.title = s.name @@ -161,7 +170,7 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: # it every iteration, so a key saved (or cleared) through the admin UI # takes effect without a restart. nip91_task = asyncio.create_task(announce_provider()) - analytics_task = asyncio.create_task(publish_usage_analytics()) + analytics_task = asyncio.create_task(run_analytics()) if global_settings.providers_refresh_interval_seconds > 0: providers_task = asyncio.create_task(providers_cache_refresher()) stale_reservation_task = asyncio.create_task(periodic_stale_reservation_sweep()) @@ -270,6 +279,7 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: "Error closing upstream HTTP connection pools", extra={"error": str(e), "error_type": type(e).__name__}, ) + await shutdown_analytics() class _ImmutableStaticFiles(StaticFiles): diff --git a/routstr/core/settings.py b/routstr/core/settings.py index d9324f13..1ce2afbb 100644 --- a/routstr/core/settings.py +++ b/routstr/core/settings.py @@ -531,22 +531,53 @@ class SettingsService: for k, v in _normalize_settings_data(partial).items() if k not in FIXED_FIELDS } - candidate_dict = {**current.dict(), **sanitized_partial} - candidate = Settings(**candidate_dict) from sqlmodel import text - # Ensure primary_mint reflects candidate mints if missing - if not candidate.primary_mint: - candidate.primary_mint = _compute_primary_mint(candidate.cashu_mints) - - await db_session.exec( # type: ignore - text( - "UPDATE settings SET data = :data, updated_at = :updated_at WHERE id = 1" - ).bindparams( - data=json.dumps(_strip_secret_fields(candidate.dict())), - updated_at=datetime.now(timezone.utc), + while True: + row = await db_session.exec( # type: ignore + text("SELECT data FROM settings WHERE id = 1") ) - ) + row = row.first() + seen = row[0] if row else None + # Build on the saved document: another worker may have saved + # since this one loaded, and its choices must survive this edit. + stored = ( + _strip_secret_fields(_normalize_settings_data(json.loads(seen))) + if isinstance(seen, str) + else {} + ) + candidate = Settings( + **{**current.dict(), **stored, **sanitized_partial} + ) + + # Ensure primary_mint reflects candidate mints if missing + if not candidate.primary_mint: + candidate.primary_mint = _compute_primary_mint( + candidate.cashu_mints + ) + + saved = await db_session.exec( # type: ignore + text( + "UPDATE settings SET data = :data, updated_at = :updated_at " + "WHERE id = 1 AND data = :seen" + ).bindparams( + data=json.dumps(_strip_secret_fields(candidate.dict())), + updated_at=datetime.now(timezone.utc), + seen=seen, + ) + ) + if seen is None or saved.rowcount == 1: + break + await db_session.rollback() + if ( + "enable_analytics_sharing" in sanitized_partial + and not candidate.enable_analytics_sharing + ): + # Saved with the opt-out, so a worker that read older flags + # cannot activate sharing after it. + from ..nostr.analytics_v2_delivery import fence_analytics_v2_opt_out + + await fence_analytics_v2_opt_out(db_session) await db_session.commit() # Update in-place. Env-only fields (e.g. DB pool sizing) are never # applied here: the engine pool is already built at boot from env, @@ -559,6 +590,21 @@ class SettingsService: cls._current = settings return settings + @classmethod + async def refresh(cls, db_session: AsyncSession, fields: tuple[str, ...]) -> None: + """Adopt saved fields; an update only reaches the worker that made it.""" + from sqlmodel import text + + async with cls._lock: + row = await db_session.exec(text("SELECT data FROM settings WHERE id = 1")) # type: ignore + row = row.first() + if row is None: + return + data = json.loads(row[0]) if isinstance(row[0], str) else dict(row[0]) + for name in fields: + if name in data: + setattr(settings, name, data[name]) + @classmethod async def reload_from_db(cls, db_session: AsyncSession) -> Settings: async with cls._lock: diff --git a/routstr/core/usage_analytics_store.py b/routstr/core/usage_analytics_store.py index 7ba90e24..6c6e5f97 100644 --- a/routstr/core/usage_analytics_store.py +++ b/routstr/core/usage_analytics_store.py @@ -783,7 +783,7 @@ class UsageAnalyticsStore: """ SELECT datetime( - (CAST(strftime('%s', minute_ts) AS INTEGER) / ?) * ?, + (CAST(strftime('%s', minute_ts, 'utc') AS INTEGER) / ?) * ?, 'unixepoch' ) AS bucket_ts, COALESCE(SUM(total_requests), 0) AS total_requests, @@ -1185,7 +1185,7 @@ class UsageAnalyticsStore: """ SELECT datetime( - (CAST(strftime('%s', minute_ts) AS INTEGER) / ?) * ?, + (CAST(strftime('%s', minute_ts, 'utc') AS INTEGER) / ?) * ?, 'unixepoch' ) AS bucket_ts, COALESCE(SUM(successful), 0) AS total_successful, @@ -1240,7 +1240,7 @@ class UsageAnalyticsStore: f""" SELECT datetime( - (CAST(strftime('%s', minute_ts) AS INTEGER) / ?) * ?, + (CAST(strftime('%s', minute_ts, 'utc') AS INTEGER) / ?) * ?, 'unixepoch' ) AS bucket_ts, model, @@ -1300,8 +1300,9 @@ class UsageAnalyticsStore: } def _cutoff_timestamp(self, hours_back: int) -> str: + # Stored minutes are server-local log stamps, so the cutoff must be too. cutoff = datetime.now(timezone.utc) - timedelta(hours=hours_back) - return cutoff.strftime("%Y-%m-%d %H:%M:%S") + return cutoff.astimezone().strftime("%Y-%m-%d %H:%M:%S") def _minute_key(self, timestamp: Any) -> str | None: if not isinstance(timestamp, str) or len(timestamp) != 19: diff --git a/routstr/nostr/analytics_runtime.py b/routstr/nostr/analytics_runtime.py new file mode 100644 index 00000000..0b7c507c --- /dev/null +++ b/routstr/nostr/analytics_runtime.py @@ -0,0 +1,232 @@ +from __future__ import annotations + +import asyncio +import time +from datetime import UTC, datetime + +from ..core.db import create_session +from ..core.logging import get_logger +from ..core.settings import SettingsService, settings +from ..core.terminal_outcomes import ( + start_terminal_outcome_writer, + stop_terminal_outcome_writer, + terminal_outcome_writer, +) +from .analytics_v2_delivery import ( + AnalyticsV2Delivery, + AnalyticsV2Producer, + DeliveryStateSnapshot, + SharingDisabledError, + activate_analytics_v2_sharing, + claim_analytics_v2_identity, + get_analytics_v2_delivery_state, + rotate_analytics_v2_identity, + run_analytics_v2_publisher, + transition_analytics_v2_sharing, +) +from .listing import DEFAULT_RELAY_URLS, nsec_to_keypair, resolve_provider_id_strict + +logger = get_logger(__name__) + + +async def _read_state() -> DeliveryStateSnapshot: + # Read the fence before the setting so a later opt-out refuses activation. + state = await get_analytics_v2_delivery_state(create_session) + async with create_session() as session: + await SettingsService.refresh(session, ("enable_analytics_sharing",)) + return state + + +class AnalyticsCoordinator: + def __init__(self) -> None: + self._task: asyncio.Task[None] | None = None + self._delivery: AnalyticsV2Delivery | None = None + self._identity: tuple[str, str, str] | None = None + self._relays: tuple[str, ...] = () + self._writer_started = False + self._retry_at = 0.0 + self._closed = False + + async def prepare_startup(self) -> None: + state = await get_analytics_v2_delivery_state(create_session) + self._writer_started = await start_terminal_outcome_writer() + if state.sharing_enabled and ( + not settings.enable_analytics_sharing or not self._writer_started + ): + await transition_analytics_v2_sharing(create_session, enabled=False) + + async def run(self) -> None: + try: + while True: + delay = 1 + try: + await self.sync_once() + except asyncio.CancelledError: + raise + except Exception: + logger.exception( + "Stats coordination failed; requests remain available" + ) + delay = 10 + await asyncio.sleep(delay) + finally: + await self.close() + + async def sync_once(self) -> None: + if self._closed: + return + state = await _read_state() + wants_public = settings.enable_analytics_sharing + if not wants_public: + await self._stop_public(disable=state.sharing_enabled) + if not terminal_outcome_writer.running: + if time.monotonic() < self._retry_at: + return + self._writer_started = await start_terminal_outcome_writer(serving=True) + if not self._writer_started: + await self._stop_public(disable=state.sharing_enabled) + self._retry_at = time.monotonic() + 10 + return + + if not wants_public: + return + + if time.monotonic() < self._retry_at: + return + keypair = nsec_to_keypair(settings.nsec) if settings.nsec else None + if keypair is None: + await self._stop_public(disable=state.sharing_enabled) + return + private_key, pubkey = keypair + relays = tuple(dict.fromkeys(settings.relays or DEFAULT_RELAY_URLS)) + if ( + state.identity_pubkey == pubkey + and state.provider_d + and (not settings.provider_id or settings.provider_id == state.provider_d) + ): + provider_d = state.provider_d + else: + try: + provider_d = await resolve_provider_id_strict(pubkey, list(relays)) + except Exception: + await self._stop_public(disable=state.sharing_enabled) + self._retry_at = time.monotonic() + 60 + logger.exception("Stats need a stable provider identity before sharing") + return + identity = (private_key, pubkey, provider_d) + if ( + self._identity == identity + and self._relays == relays + and state.sharing_enabled + and self._task is not None + and not self._task.done() + ): + return + + identity_changed = state.identity_pubkey is not None and ( + state.identity_pubkey != pubkey or state.provider_d != provider_d + ) + await self._stop_public(disable=identity_changed) + try: + if identity_changed: + await rotate_analytics_v2_identity( + create_session, pubkey=pubkey, provider_d=provider_d + ) + state = await _read_state() + if not settings.enable_analytics_sharing: + return + else: + claim = await claim_analytics_v2_identity( + create_session, pubkey=pubkey, provider_d=provider_d + ) + if claim == "mismatch": + self._retry_at = time.monotonic() + 10 + return + await activate_analytics_v2_sharing( + create_session, + coverage_day=datetime.now(UTC).date(), + expected_generation=state.generation, + ) + producer = AnalyticsV2Producer( + create_session, + private_key_hex=private_key, + public_key_hex=pubkey, + provider_d=provider_d, + ) + self._delivery = AnalyticsV2Delivery( + create_session, operator_relays=list(relays) + ) + self._identity = identity + self._relays = relays + self._task = asyncio.create_task( + run_analytics_v2_publisher(producer, self._delivery), + name="analytics-v2-publisher", + ) + except SharingDisabledError: + return + except Exception: + await self._stop_public(disable=True) + self._retry_at = time.monotonic() + 60 + raise + + async def _stop_task(self) -> None: + task, self._task = self._task, None + if task is not None: + task.cancel() + try: + await task + except asyncio.CancelledError: + pass + except Exception: + logger.exception("Stats publisher stopped with an error") + + async def _stop_public(self, *, disable: bool) -> None: + try: + if disable: + if self._delivery is not None: + await self._delivery.disable() + else: + await transition_analytics_v2_sharing(create_session, enabled=False) + elif self._delivery is not None: + await self._delivery.stop() + finally: + await self._stop_task() + self._delivery = None + self._identity = None + self._relays = () + + async def close(self) -> None: + if self._closed: + return + self._closed = True + try: + await self._stop_public(disable=False) + finally: + if self._writer_started: + await stop_terminal_outcome_writer() + self._writer_started = False + + +_coordinator: AnalyticsCoordinator | None = None + + +def _get_coordinator() -> AnalyticsCoordinator: + global _coordinator + if _coordinator is None or _coordinator._closed: + _coordinator = AnalyticsCoordinator() + return _coordinator + + +async def prepare_analytics() -> None: + await _get_coordinator().prepare_startup() + + +async def run_analytics() -> None: + await _get_coordinator().run() + + +async def shutdown_analytics() -> None: + global _coordinator + coordinator, _coordinator = _coordinator, None + if coordinator is not None: + await coordinator.close() diff --git a/routstr/nostr/listing.py b/routstr/nostr/listing.py index 02a668eb..6e39ce80 100644 --- a/routstr/nostr/listing.py +++ b/routstr/nostr/listing.py @@ -9,8 +9,11 @@ import json import os import random import time +import unicodedata from typing import Any, cast +from nostr_sdk import Event + from ..core import get_logger from ..core.settings import settings from .sdk import create_signed_event, fetch_events, parse_keypair, send_event @@ -251,6 +254,41 @@ async def _determine_provider_id(public_key_hex: str, relay_urls: list[str]) -> return fallback +async def resolve_provider_id_strict(public_key_hex: str, relay_urls: list[str]) -> str: + """Require a configured or unambiguous signed coordinate for durable stats.""" + explicit = settings.provider_id + if explicit: + if len(explicit) > 64 or any( + unicodedata.category(char) == "Cc" for char in explicit + ): + raise ValueError("PROVIDER_ID must contain 1 to 64 printable characters") + return explicit + results = await asyncio.gather( + *(query_listing_events(url, public_key_hex) for url in relay_urls), + return_exceptions=True, + ) + candidates: set[str] = set() + for result in results: + if isinstance(result, BaseException): + raise ValueError("Configure PROVIDER_ID while listing relays are unavailable") + events, ok = result + if not ok or len(events) >= 10: + raise ValueError("Configure PROVIDER_ID when listing history is incomplete") + for event in events: + try: + signed = Event.from_json(json.dumps(event)) + if not signed.verify() or event.get("pubkey") != public_key_hex: + continue + values = _get_tag_values(event, "d") + if event.get("kind") == 38421 and len(values) == 1 and values[0]: + candidates.add(values[0]) + except Exception: + continue + if len(candidates) != 1: + raise ValueError("Configure PROVIDER_ID to select one provider for public stats") + return next(iter(candidates)) + + async def publish_to_relay( relay_url: str, event: dict[str, Any], diff --git a/tests/integration/test_admin_settings_endpoint.py b/tests/integration/test_admin_settings_endpoint.py index 15d05159..c0bc8bf9 100644 --- a/tests/integration/test_admin_settings_endpoint.py +++ b/tests/integration/test_admin_settings_endpoint.py @@ -12,13 +12,14 @@ from __future__ import annotations import secrets import time from collections.abc import AsyncGenerator +from datetime import UTC, datetime, timedelta import pytest import pytest_asyncio from httpx import AsyncClient from routstr.core.admin import admin_sessions -from routstr.core.db import AsyncSession +from routstr.core.db import AsyncSession, TerminalOutcome, TerminalOutcomeEpoch from routstr.core.settings import SettingsService, settings @@ -34,6 +35,82 @@ async def admin_client( admin_sessions.pop(token, None) +@pytest.mark.integration +@pytest.mark.asyncio +async def test_dashboard_uses_exact_saved_dates_with_sharing_disabled( + admin_client: AsyncClient, + integration_session: AsyncSession, + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setattr(settings, "enable_analytics_sharing", False) + today = datetime.now(UTC).replace(hour=0, minute=0, second=0, microsecond=0) + start = today - timedelta(days=3) + end = start + timedelta(days=1) + integration_session.add( + TerminalOutcomeEpoch( + epoch=1, + coverage_start_day=start.date(), + coverage_end_day=end.date(), + ) + ) + for name, at, revenue in ( + ("selected-start", start, 1200), + ("selected-end-excluded", end, 9999), + ("recent-excluded", today, 7000), + ): + integration_session.add( + TerminalOutcome( + outcome_id=name, + terminal_at_ms=int(at.timestamp() * 1000), + terminal_day=at.date(), + model_identifier="fixture/model", + input_tokens=10, + output_tokens=5, + cache_read_input_tokens=0, + cache_creation_input_tokens=0, + revenue_msats=revenue, + ) + ) + await integration_session.commit() + response = await admin_client.get( + "/admin/api/usage/dashboard", + params={ + "hours": 24, + "start_at": start.isoformat(), + "end_at": end.isoformat(), + }, + ) + assert response.status_code == 200 + result = response.json() + assert result["analytics_source"] == "terminal_outcomes" + assert result["summary"]["successful_chat_completions"] == 1 + assert result["summary"]["revenue_msats"] == 1200 + assert result["summary"]["total_tokens"] == 15 + assert result["ledger_coverage"]["complete"] is True + assert result["ledger_coverage"]["diagnostic_available"] is False + assert result["error_details"]["errors"] == [] + assert result["model_usage_mix"]["metrics"][0]["total_revenue_msats"] == 1200 + + +@pytest.mark.integration +@pytest.mark.asyncio +@pytest.mark.parametrize( + "params", + [ + {"start_at": "2026-01-01T00:00:00Z"}, + {"start_at": "2026-01-01", "end_at": "2026-01-02"}, + {"start_at": "2026-01-02T00:00:00Z", "end_at": "2026-01-01T00:00:00Z"}, + {"start_at": "2024-01-01T00:00:00Z", "end_at": "2026-01-01T00:00:00Z"}, + {"start_at": "2099-01-01T00:00:00Z", "end_at": "2099-01-02T00:00:00Z"}, + ], +) +async def test_dashboard_rejects_ambiguous_date_ranges( + admin_client: AsyncClient, params: dict +) -> None: + response = await admin_client.get("/admin/api/usage/dashboard", params=params) + assert response.status_code == 400 + + @pytest.mark.integration @pytest.mark.asyncio async def test_get_settings_omits_admin_password_and_redacts_secrets( @@ -80,3 +157,32 @@ async def test_patch_settings_ignores_secret_fields( assert data["nsec"] == "[REDACTED]" # The live secret was not overwritten through the general settings endpoint. assert settings.nsec == "original-nsec" + + +@pytest.mark.integration +@pytest.mark.asyncio +async def test_existing_sharing_choice_is_saved( + admin_client: AsyncClient, + integration_session: AsyncSession, + monkeypatch: pytest.MonkeyPatch, +) -> None: + await SettingsService.initialize(integration_session) + monkeypatch.setattr(settings, "enable_analytics_sharing", False) + resp = await admin_client.patch( + "/admin/api/settings", + json={"enable_analytics_sharing": True}, + ) + assert resp.status_code == 200 + assert resp.json()["enable_analytics_sharing"] is True + saved = await admin_client.get("/admin/api/settings") + assert saved.json()["enable_analytics_sharing"] is True + assert "enable_analytics_collection" not in saved.json() + assert "enable_analytics_v2" not in saved.json() + stopped = await admin_client.patch( + "/admin/api/settings", + json={ + "enable_analytics_sharing": False, + }, + ) + assert stopped.status_code == 200 + assert stopped.json()["enable_analytics_sharing"] is False diff --git a/tests/unit/test_analytics_runtime.py b/tests/unit/test_analytics_runtime.py new file mode 100644 index 00000000..0da4229c --- /dev/null +++ b/tests/unit/test_analytics_runtime.py @@ -0,0 +1,283 @@ +from __future__ import annotations + +import asyncio +import json +from collections.abc import AsyncIterator +from contextlib import asynccontextmanager +from datetime import UTC, datetime, timedelta +from types import SimpleNamespace +from typing import Any + +import pytest +import pytest_asyncio +from sqlalchemy.ext.asyncio import create_async_engine +from sqlmodel import SQLModel, col, select, text +from sqlmodel.ext.asyncio.session import AsyncSession + +from routstr.core import terminal_outcomes +from routstr.core.db import TerminalOutcomeEpoch +from routstr.core.settings import SettingsService +from routstr.nostr import analytics_runtime as runtime + + +@pytest_asyncio.fixture +async def node(monkeypatch: pytest.MonkeyPatch) -> AsyncIterator[Any]: + engine = create_async_engine("sqlite+aiosqlite:///:memory:") + async with engine.begin() as connection: + await connection.run_sync(SQLModel.metadata.create_all) + # No saved row: each test drives the in-process flags directly. + await connection.exec_driver_sql( + "CREATE TABLE settings (id INTEGER PRIMARY KEY, data TEXT NOT NULL, " + "updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)" + ) + + @asynccontextmanager + async def session_factory() -> AsyncIterator[AsyncSession]: + async with AsyncSession(engine, expire_on_commit=False) as session: + yield session + + events: list[str] = [] + + async def publisher(*args: object) -> None: + events.append("start:daily") + try: + await asyncio.Event().wait() + finally: + events.append("stop:daily") + + writer = terminal_outcomes.TerminalOutcomeWriter(session_factory=session_factory) + monkeypatch.setattr(terminal_outcomes, "terminal_outcome_writer", writer) + monkeypatch.setattr(runtime, "terminal_outcome_writer", writer) + monkeypatch.setattr(runtime, "create_session", session_factory) + monkeypatch.setattr(runtime, "run_analytics_v2_publisher", publisher) + for key, value in { + "nsec": "11" * 32, + "provider_id": "stats-test-node", + "relays": ["wss://relay.example.com"], + "enable_analytics_sharing": False, + }.items(): + monkeypatch.setattr(runtime.settings, key, value) + coordinator = runtime.AnalyticsCoordinator() + try: + yield SimpleNamespace( + coordinator=coordinator, + writer=writer, + sessions=session_factory, + events=events, + ) + finally: + await coordinator.close() + await writer.stop() + await engine.dispose() + + +@pytest.mark.asyncio +async def test_collection_is_private_without_identity_or_sharing( + node: Any, monkeypatch: Any +) -> None: + monkeypatch.setattr(runtime.settings, "nsec", "") + await node.coordinator.prepare_startup() + await node.coordinator.sync_once() + assert node.writer.running + assert node.coordinator._task is None + state = await runtime.get_analytics_v2_delivery_state(node.sessions) + assert not state.sharing_enabled + assert state.identity_pubkey is None + + +@pytest.mark.asyncio +async def test_public_opt_out_keeps_private_collection_running( + node: Any, monkeypatch: Any +) -> None: + monkeypatch.setattr(runtime.settings, "enable_analytics_sharing", True) + await node.coordinator.prepare_startup() + await node.coordinator.sync_once() + await asyncio.sleep(0) + assert node.coordinator._task is not None + monkeypatch.setattr(runtime.settings, "enable_analytics_sharing", False) + await node.coordinator.sync_once() + assert node.writer.running + assert node.coordinator._task is None + assert not ( + await runtime.get_analytics_v2_delivery_state(node.sessions) + ).sharing_enabled + assert node.events == ["start:daily", "stop:daily"] + + +@pytest.mark.asyncio +async def test_opt_out_saved_by_another_worker_is_not_undone( + node: Any, monkeypatch: Any +) -> None: + monkeypatch.setattr(runtime.settings, "enable_analytics_sharing", True) + await node.coordinator.prepare_startup() + await node.coordinator.sync_once() + assert node.coordinator._task is not None + + # The other worker saved the opt-out and already disabled delivery. + async with node.sessions() as session: + await session.exec( # type: ignore[call-overload] + text("INSERT INTO settings (id, data) VALUES (1, :data)").bindparams( + data=json.dumps({"enable_analytics_sharing": False}) + ) + ) + await session.commit() + await runtime.transition_analytics_v2_sharing(node.sessions, enabled=False) + + await node.coordinator.sync_once() + assert node.coordinator._task is None + assert node.writer.running + assert not ( + await runtime.get_analytics_v2_delivery_state(node.sessions) + ).sharing_enabled + + +@pytest.mark.asyncio +async def test_missing_identity_does_not_stop_private_collection( + node: Any, monkeypatch: Any +) -> None: + monkeypatch.setattr(runtime.settings, "nsec", "") + monkeypatch.setattr(runtime.settings, "enable_analytics_sharing", True) + await node.coordinator.prepare_startup() + await node.coordinator.sync_once() + assert node.writer.running + assert node.coordinator._task is None + assert not ( + await runtime.get_analytics_v2_delivery_state(node.sessions) + ).sharing_enabled + + +@pytest.mark.asyncio +async def test_identity_rotation_cancels_previous_publisher( + node: Any, monkeypatch: Any +) -> None: + monkeypatch.setattr(runtime.settings, "enable_analytics_sharing", True) + await node.coordinator.prepare_startup() + await node.coordinator.sync_once() + await asyncio.sleep(0) + monkeypatch.setattr(runtime.settings, "nsec", "22" * 32) + await node.coordinator.sync_once() + await asyncio.sleep(0) + state = await runtime.get_analytics_v2_delivery_state(node.sessions) + keypair = runtime.nsec_to_keypair("22" * 32) + assert keypair is not None + assert state.identity_pubkey == keypair[1] + assert node.events == ["start:daily", "stop:daily", "start:daily"] + assert node.writer.running + + +@pytest.mark.asyncio +async def test_restart_writer_failure_disables_publication_before_serving( + node: Any, monkeypatch: Any +) -> None: + monkeypatch.setattr(runtime.settings, "enable_analytics_sharing", True) + await node.coordinator.prepare_startup() + await node.coordinator.sync_once() + await node.coordinator.close() + assert ( + await runtime.get_analytics_v2_delivery_state(node.sessions) + ).sharing_enabled + + async def failed_start() -> bool: + return False + + monkeypatch.setattr(runtime, "start_terminal_outcome_writer", failed_start) + restarted = runtime.AnalyticsCoordinator() + await restarted.prepare_startup() + assert not ( + await runtime.get_analytics_v2_delivery_state(node.sessions) + ).sharing_enabled + assert not node.writer.running + await restarted.close() + + +@pytest.mark.asyncio +async def test_saved_opt_out_survives_a_stale_worker_and_fences_activation( + node: Any, monkeypatch: Any +) -> None: + public = {"enable_analytics_sharing": True} + async with node.sessions() as session: + await SettingsService.initialize(session) + await SettingsService.update(public, session) + seen = (await runtime.get_analytics_v2_delivery_state(node.sessions)).generation + + # Another worker saves the opt-out; this one still holds the old flags. + async with node.sessions() as session: + row = (await session.exec(text("SELECT data FROM settings WHERE id = 1"))).one() # type: ignore[call-overload] + saved = {**json.loads(row[0]), "enable_analytics_sharing": False} + await session.exec( # type: ignore[call-overload] + text("UPDATE settings SET data = :data WHERE id = 1").bindparams( + data=json.dumps(saved) + ) + ) + await session.commit() + async with node.sessions() as session: + await SettingsService.update({"name": "renamed"}, session) + row = (await session.exec(text("SELECT data FROM settings WHERE id = 1"))).one() # type: ignore[call-overload] + assert json.loads(row[0])["enable_analytics_sharing"] is False + assert not runtime.settings.enable_analytics_sharing + + # Saving the opt-out itself moves the fence in the same commit. + async with node.sessions() as session: + await SettingsService.update({"enable_analytics_sharing": False}, session) + state = await runtime.get_analytics_v2_delivery_state(node.sessions) + assert not state.sharing_enabled + assert state.generation == seen + 1 + + +@pytest.mark.asyncio +@pytest.mark.parametrize("sharing_enabled", [True, False]) +async def test_upgrade_preserves_existing_sharing_choice_across_restart( + node: Any, sharing_enabled: bool +) -> None: + async with node.sessions() as session: + await session.exec( # type: ignore[call-overload] + text("INSERT INTO settings (id, data) VALUES (1, :data)").bindparams( + data=json.dumps( + { + "enable_analytics_sharing": sharing_enabled, + "provider_id": "stats-test-node", + "relays": ["wss://relay.example.com"], + } + ) + ) + ) + await session.commit() + restored = await SettingsService.initialize(session) + assert restored.enable_analytics_sharing is sharing_enabled + await node.coordinator.prepare_startup() + await node.coordinator.sync_once() + await asyncio.sleep(0) + assert node.events == (["start:daily"] if sharing_enabled else []) + assert node.writer.running + state = await runtime.get_analytics_v2_delivery_state(node.sessions) + assert state.sharing_enabled is sharing_enabled + if sharing_enabled: + async with node.sessions() as session: + epoch = ( + await session.exec( + select(TerminalOutcomeEpoch).where( + col(TerminalOutcomeEpoch.current_slot) == 1 + ) + ) + ).one() + assert epoch.coverage_start_day == datetime.now(UTC).date() + timedelta( + days=1 + ) + await node.coordinator.close() + + async with node.sessions() as session: + restored = await SettingsService.initialize(session) + assert restored.enable_analytics_sharing is sharing_enabled + restarted = runtime.AnalyticsCoordinator() + try: + await restarted.prepare_startup() + await restarted.sync_once() + await asyncio.sleep(0) + assert node.writer.running + state = await runtime.get_analytics_v2_delivery_state(node.sessions) + assert state.sharing_enabled is sharing_enabled + if not sharing_enabled: + assert restarted._task is None + assert node.events == [] + finally: + await restarted.close() diff --git a/tests/unit/test_ledger_analytics.py b/tests/unit/test_ledger_analytics.py new file mode 100644 index 00000000..1a08250d --- /dev/null +++ b/tests/unit/test_ledger_analytics.py @@ -0,0 +1,745 @@ +from __future__ import annotations + +from collections.abc import AsyncGenerator +from contextlib import asynccontextmanager +from copy import deepcopy +from datetime import UTC, date, datetime, timedelta +from pathlib import Path + +import pytest +from sqlalchemy.ext.asyncio import create_async_engine +from sqlmodel import SQLModel +from sqlmodel.ext.asyncio.session import AsyncSession + +from routstr.core.db import ( + TerminalOutcome, + TerminalOutcomeEpoch, + TerminalOutcomeWriterRun, +) +from routstr.core.ledger_analytics import get_ledger_usage_dashboard +from routstr.core.terminal_outcome_writer import SessionFactory + +NOW = datetime(2026, 9, 18, 12, tzinfo=UTC) + + +@pytest.fixture +async def sessions(tmp_path: Path) -> AsyncGenerator[SessionFactory, None]: + engine = create_async_engine(f"sqlite+aiosqlite:///{tmp_path / 'dashboard.db'}") + async with engine.begin() as connection: + await connection.run_sync(SQLModel.metadata.create_all) + + @asynccontextmanager + async def factory() -> AsyncGenerator[AsyncSession, None]: + async with AsyncSession(engine, expire_on_commit=False) as session: + yield session + + yield factory + await engine.dispose() + + +def _row( + name: str, + at: datetime, + *, + model: str | None = "reported/model", + revenue: int = 1000, + input_tokens: int = 10, + output_tokens: int = 5, + cache_read: int = 2, + cache_write: int = 1, + source: str = "reported", +) -> TerminalOutcome: + return TerminalOutcome( + outcome_id=name, + terminal_at_ms=int(at.timestamp() * 1000), + terminal_day=at.date(), + model_identifier=model, + revenue_msats=revenue, + input_tokens=input_tokens, + output_tokens=output_tokens, + cache_read_input_tokens=cache_read, + cache_creation_input_tokens=cache_write, + input_source=source, + output_source=source, + cache_read_source="reported" if cache_read else "missing", + cache_creation_source="reported" if cache_write else "missing", + ) + + +def _legacy() -> dict: + return { + "summary": { + "total_requests": 9, + "successful_chat_completions": 999, + "failed_requests": 2, + "total_errors": 3, + "refunds_msats": 9000, + "success_rate": 42, + "revenue_msats": 9999, + "net_revenue_msats": 999, + "input_tokens": 999, + "output_tokens": 999, + "total_tokens": 1998, + "avg_latency_ms": 45, + }, + "metrics": { + "metrics": [ + { + "timestamp": "2026-09-17 12:00:00", + "total_requests": 9, + "failed_requests": 2, + "errors": 3, + "refunds_msats": 9000, + "successful_chat_completions": 999, + "revenue_msats": 9999, + "input_tokens": 999, + "output_tokens": 999, + "total_tokens": 1998, + } + ], + "totals": { + "total_requests": 9, + "failed_requests": 2, + "refunds_msats": 9000, + "successful_chat_completions": 999, + "revenue_msats": 9999, + }, + }, + "error_details": { + "errors": [{"message": "upstream timeout"}], + "total_count": 3, + }, + "revenue_by_model": { + "models": [ + { + "model": "reported/model", + "requests": 7, + "failed": 2, + "refunds_sats": 9, + } + ] + }, + } + + +async def _seed(sessions: SessionFactory) -> None: + async with sessions() as session: + session.add_all( + [ + _row("before", NOW - timedelta(days=1, milliseconds=1), revenue=9999), + _row("start", NOW - timedelta(days=1)), + _row( + "free", + NOW - timedelta(hours=4), + model="free/model", + revenue=0, + input_tokens=12, + output_tokens=3, + cache_read=0, + cache_write=0, + source="estimated", + ), + _row( + "unknown", + NOW - timedelta(hours=2), + model=None, + revenue=250, + input_tokens=0, + output_tokens=0, + cache_read=0, + cache_write=0, + source="missing", + ), + _row("at-end", NOW, revenue=9999), + _row("future", NOW + timedelta(hours=1), revenue=9999), + TerminalOutcomeEpoch( + epoch=0, coverage_start_day=date(2026, 9, 17), current_slot=1 + ), + ] + ) + session.add( + TerminalOutcomeWriterRun( + run_id="active-run", + status="active", + started_at_ms=int((NOW - timedelta(days=1)).timestamp() * 1000), + heartbeat_at_ms=int(NOW.timestamp() * 1000), + flushed_through_ms=int(NOW.timestamp() * 1000), + ) + ) + await session.commit() + + +async def test_dashboard_replaces_log_totals_and_matches_every_chart( + sessions: SessionFactory, +) -> None: + await _seed(sessions) + legacy = _legacy() + before = deepcopy(legacy) + result = await get_ledger_usage_dashboard( + legacy, + interval=60, + hours=24, + model_limit=1, + session_factory=sessions, + now=NOW, + ) + summary = result["summary"] + assert summary["successful_chat_completions"] == 3 + assert summary["revenue_msats"] == summary["net_revenue_msats"] == 1250 + assert summary["input_tokens"] == 25 + assert summary["output_tokens"] == 8 + assert summary["total_tokens"] == 33 + assert summary["total_requests"] == 9 and summary["failed_requests"] == 2 + assert summary["refunds_msats"] == 9000 + assert summary["success_rate"] == 42 and summary["avg_latency_ms"] == 45 + assert result["error_details"] == before["error_details"] + assert legacy == before + assert result["analytics_source"] == "terminal_outcomes" + for field in ("successful_chat_completions", "revenue_msats", "total_tokens"): + assert ( + sum(point[field] for point in result["metrics"]["metrics"]) + == summary[field] + ) + assert result["metrics"]["totals"][field] == summary[field] + mix = result["model_usage_mix"] + assert sum(point["total_successful"] for point in mix["metrics"]) == 3 + assert sum(point["total_revenue_msats"] for point in mix["metrics"]) == 1250 + assert sum(point["total_tokens"] for point in mix["metrics"]) == 33 + assert sum(point["others"] for point in mix["metrics"]) == 1 + assert "free/model" in mix["top_models_by_metric"]["requests"] + assert summary["unique_models_count"] == 2 + assert result["revenue_by_model"]["models"][0]["successful"] == 1 + assert result["revenue_by_model"]["total_revenue_sats"] == 1.25 + coverage = result["ledger_coverage"] + assert coverage["complete"] and coverage["diagnostic_available"] + assert coverage["includes_current_day"] + assert coverage["token_sources"]["input"] == { + "reported": 1, + "estimated": 1, + "missing": 1, + } + assert coverage["token_sources"]["cache_read"] == { + "reported": 1, + "estimated": 0, + "missing": 2, + } + + +async def test_historical_dates_do_not_return_recent_ledger_or_log_data( + sessions: SessionFactory, +) -> None: + await _seed(sessions) + start, end = datetime(2026, 9, 17, tzinfo=UTC), datetime(2026, 9, 18, tzinfo=UTC) + result = await get_ledger_usage_dashboard( + _legacy(), + interval=60, + hours=24, + session_factory=sessions, + now=NOW, + start_at=start, + end_at=end, + ) + assert result["summary"]["successful_chat_completions"] == 2 + assert result["summary"]["revenue_msats"] == 10_999 + assert result["summary"]["total_requests"] == 0 + assert result["summary"]["refunds_msats"] == 0 + assert result["summary"]["success_rate"] == 0 + assert result["error_details"] == {"errors": [], "total_count": 0} + assert all( + point["timestamp"].startswith("2026-09-17") + for point in result["metrics"]["metrics"] + ) + coverage = result["ledger_coverage"] + assert coverage["from"] == start.isoformat() and coverage["to"] == end.isoformat() + assert not coverage["diagnostic_available"] and not coverage["includes_current_day"] + + +async def test_empty_ledger_has_no_fallback_to_legacy_successes( + sessions: SessionFactory, +) -> None: + result = await get_ledger_usage_dashboard( + _legacy(), + interval=60, + hours=24, + session_factory=sessions, + now=NOW, + ) + assert result["summary"]["successful_chat_completions"] == 0 + assert result["summary"]["total_tokens"] == 0 + assert result["summary"]["revenue_msats"] == 0 + assert result["summary"]["total_requests"] == 9 + assert len(result["model_usage_mix"]["metrics"]) == 24 + assert all( + point["coverage"] == "missing" and point["total_successful"] is None + for point in result["model_usage_mix"]["metrics"] + ) + assert result["ledger_coverage"]["incomplete_days"] == ["2026-09-17", "2026-09-18"] + assert not result["ledger_coverage"]["complete"] + assert result["ledger_coverage"]["latest_outcome_at"] is None + + +async def test_closed_epoch_preserves_history_and_exposes_collection_gap( + sessions: SessionFactory, +) -> None: + async with sessions() as session: + session.add_all( + [ + TerminalOutcomeEpoch( + epoch=0, + coverage_start_day=date(2026, 9, 15), + coverage_end_day=date(2026, 9, 16), + current_slot=None, + ), + TerminalOutcomeEpoch( + epoch=1, coverage_start_day=date(2026, 9, 19), current_slot=1 + ), + _row("historical", datetime(2026, 9, 16, 12, tzinfo=UTC)), + ] + ) + await session.commit() + result = await get_ledger_usage_dashboard( + _legacy(), + interval=60, + hours=24, + session_factory=sessions, + now=NOW, + start_at=datetime(2026, 9, 15, tzinfo=UTC), + end_at=datetime(2026, 9, 19, tzinfo=UTC), + ) + assert result["summary"]["successful_chat_completions"] == 1 + assert result["ledger_coverage"]["incomplete_days"] == ["2026-09-17", "2026-09-18"] + assert result["metrics"]["hours_back"] == 96 + + +async def test_model_limit_keeps_omitted_models_in_other_totals( + sessions: SessionFactory, +) -> None: + async with sessions() as session: + session.add_all( + [ + _row( + f"model-{i}", + NOW - timedelta(hours=1), + model=f"model/{i}", + revenue=i, + ) + for i in range(25) + ] + ) + await session.commit() + result = await get_ledger_usage_dashboard( + _legacy(), + interval=60, + hours=24, + model_limit=2, + session_factory=sessions, + now=NOW, + ) + mix = next( + point + for point in result["model_usage_mix"]["metrics"] + if point["total_successful"] + ) + assert result["summary"]["unique_models_count"] == 25 + assert len(result["revenue_by_model"]["models"]) == 2 + assert sum(mix["model_counts"].values()) + mix["others"] == 25 + assert sum(mix["model_revenue_msats"].values()) + mix[ + "others_revenue_msats" + ] == sum(range(25)) + assert sum(mix["model_tokens"].values()) + mix["others_tokens"] == 25 * 18 + + +async def test_future_part_of_custom_period_is_incomplete_and_excluded( + sessions: SessionFactory, +) -> None: + await _seed(sessions) + result = await get_ledger_usage_dashboard( + _legacy(), + interval=60, + hours=48, + session_factory=sessions, + now=NOW, + start_at=datetime(2026, 9, 18, tzinfo=UTC), + end_at=datetime(2026, 9, 20, tzinfo=UTC), + ) + assert result["summary"]["successful_chat_completions"] == 2 + assert result["summary"]["revenue_msats"] == 250 + assert result["ledger_coverage"]["incomplete_days"] == ["2026-09-18", "2026-09-19"] + assert not result["ledger_coverage"]["complete"] + + +@pytest.mark.parametrize( + ("lag", "incomplete_days", "idle_hour", "open_hour"), + [ + (timedelta(seconds=10), [], ("complete", 0), "updating"), + (timedelta(days=1), ["2026-09-17", "2026-09-18"], ("missing", None), "missing"), + ], +) +async def test_trailing_checkpoint_is_a_gap_only_when_left_in_an_earlier_day( + sessions: SessionFactory, + lag: timedelta, + incomplete_days: list[str], + idle_hour: tuple[str, int | None], + open_hour: str, +) -> None: + await _seed(sessions) + async with sessions() as session: + run = await session.get(TerminalOutcomeWriterRun, "active-run") + assert run is not None + run.flushed_through_ms = int((NOW - lag).timestamp() * 1000) + session.add(run) + await session.commit() + result = await get_ledger_usage_dashboard( + _legacy(), + interval=60, + hours=24, + session_factory=sessions, + now=NOW, + ) + assert result["summary"]["successful_chat_completions"] == 3 + assert result["ledger_coverage"]["incomplete_days"] == incomplete_days + points = {point["timestamp"]: point for point in result["metrics"]["metrics"]} + idle = points["2026-09-18 09:00:00"] + assert (idle["coverage"], idle["revenue_msats"]) == idle_hour + assert points["2026-09-18 11:00:00"]["coverage"] == open_hour + + +async def test_degraded_writer_exposes_gap_before_epoch_rotation( + sessions: SessionFactory, +) -> None: + await _seed(sessions) + async with sessions() as session: + run = await session.get(TerminalOutcomeWriterRun, "active-run") + assert run is not None + run.status = "degraded" + run.loss_day = date(2026, 9, 17) + session.add(run) + await session.commit() + result = await get_ledger_usage_dashboard( + _legacy(), + interval=60, + hours=24, + session_factory=sessions, + now=NOW, + ) + assert result["ledger_coverage"]["incomplete_days"] == ["2026-09-17", "2026-09-18"] + assert not result["ledger_coverage"]["complete"] + + +async def test_live_writer_loss_is_visible_before_it_can_persist_gap( + sessions: SessionFactory, + monkeypatch: pytest.MonkeyPatch, +) -> None: + from types import SimpleNamespace + + from routstr.core import terminal_outcomes + + await _seed(sessions) + monkeypatch.setattr( + terminal_outcomes, + "terminal_outcome_writer", + SimpleNamespace( + loss_pending=True, + loss_day=date(2026, 9, 17), + ), + ) + result = await get_ledger_usage_dashboard( + _legacy(), + interval=60, + hours=24, + session_factory=sessions, + now=NOW, + ) + assert result["ledger_coverage"]["incomplete_days"] == ["2026-09-17", "2026-09-18"] + assert not result["ledger_coverage"]["complete"] + + +async def test_chart_fills_covered_zero_buckets_without_fabricating_missing_days( + sessions: SessionFactory, +) -> None: + async with sessions() as session: + session.add_all( + [ + TerminalOutcomeEpoch( + epoch=0, + coverage_start_day=date(2026, 9, 14), + coverage_end_day=date(2026, 9, 16), + current_slot=None, + ), + _row("complete", datetime(2026, 9, 14, 12, tzinfo=UTC)), + _row("partial", datetime(2026, 9, 17, 12, tzinfo=UTC), revenue=250), + ] + ) + await session.commit() + result = await get_ledger_usage_dashboard( + _legacy(), + interval=1440, + hours=120, + session_factory=sessions, + now=NOW, + start_at=datetime(2026, 9, 13, tzinfo=UTC), + end_at=datetime(2026, 9, 18, tzinfo=UTC), + ) + points = result["metrics"]["metrics"] + mix = result["model_usage_mix"]["metrics"] + assert len(points) == len(mix) == 5 + assert [point["coverage"] for point in points] == [ + "missing", + "complete", + "complete", + "complete", + "partial", + ] + assert [point["revenue_msats"] for point in points] == [None, 1000, 0, 0, 250] + assert [point["timestamp"] for point in points] == [ + point["timestamp"] for point in mix + ] + assert [point["total_revenue_msats"] for point in mix] == [None, 1000, 0, 0, 250] + assert mix[0]["model_counts"] == mix[2]["model_counts"] == {} + assert result["summary"]["revenue_msats"] == 1250 + assert result["metrics"]["totals"]["successful_chat_completions"] == 2 + assert result["metrics"]["bucket_fill_complete"] + assert result["model_usage_mix"]["bucket_fill_complete"] + + +async def test_chart_keeps_requested_interval_and_bounds_bucket_filling( + sessions: SessionFactory, +) -> None: + await _seed(sessions) + result = await get_ledger_usage_dashboard( + _legacy(), interval=1, hours=24, session_factory=sessions, now=NOW + ) + assert result["metrics"]["interval_minutes"] == 1 + assert result["model_usage_mix"]["interval_minutes"] == 1 + assert not result["metrics"]["bucket_fill_complete"] + assert not result["model_usage_mix"]["bucket_fill_complete"] + assert len(result["metrics"]["metrics"]) == 3 + assert result["summary"]["revenue_msats"] == 1250 + + +async def test_measured_average_uses_paired_reported_requests_including_free_and_zero( + sessions: SessionFactory, +) -> None: + at = NOW - timedelta(hours=1) + input_only = _row( + "input-only", + at, + input_tokens=9, + output_tokens=0, + cache_read=0, + cache_write=0, + source="missing", + ) + input_only.input_source = "reported" + output_only = _row( + "output-only", + at, + input_tokens=0, + output_tokens=7, + cache_read=0, + cache_write=0, + source="missing", + ) + output_only.output_source = "reported" + async with sessions() as session: + session.add_all( + [ + _row("reported-cache", at), + _row( + "free", + at, + model=None, + revenue=0, + input_tokens=6, + output_tokens=4, + cache_read=0, + cache_write=0, + ), + _row( + "measured-zero", + at, + revenue=0, + input_tokens=0, + output_tokens=0, + cache_read=0, + cache_write=0, + ), + input_only, + output_only, + _row( + "estimated", + at, + input_tokens=20, + output_tokens=10, + cache_read=0, + cache_write=0, + source="estimated", + ), + ] + ) + await session.commit() + result = await get_ledger_usage_dashboard( + _legacy(), interval=60, hours=24, session_factory=sessions, now=NOW + ) + summary = result["summary"] + assert summary["measured_token_requests"] == 3 + assert summary["measured_tokens"] == 28 + assert summary["avg_measured_tokens_per_completion"] == pytest.approx(28 / 3) + assert summary["successful_chat_completions"] == 6 + assert summary["total_tokens"] == 74 + assert summary["avg_total_tokens_per_completion"] == pytest.approx(74 / 6) + assert summary["total_requests"] == 9 + assert summary["success_rate"] == 42 + assert "measured_tokens" not in result["metrics"]["totals"] + + +@pytest.mark.parametrize("dimension", ["cache_read", "cache_creation"]) +@pytest.mark.parametrize( + ("source", "count", "eligible"), + [ + ("reported", 3, True), + ("missing", 0, True), + ("missing", 3, False), + ("estimated", 0, False), + ("estimated", 3, False), + ], +) +async def test_measured_average_requires_reliable_nonzero_cache_counts( + sessions: SessionFactory, + dimension: str, + source: str, + count: int, + eligible: bool, +) -> None: + row = _row("cache", NOW - timedelta(hours=1), cache_read=0, cache_write=0) + setattr(row, dimension + "_source", source) + setattr(row, dimension + "_input_tokens", count) + async with sessions() as session: + session.add(row) + await session.commit() + result = await get_ledger_usage_dashboard( + _legacy(), interval=60, hours=24, session_factory=sessions, now=NOW + ) + summary = result["summary"] + assert summary["measured_token_requests"] == int(eligible) + assert summary["measured_tokens"] == (15 + count if eligible else 0) + assert summary["avg_measured_tokens_per_completion"] == ( + 15 + count if eligible else None + ) + assert summary["total_tokens"] == 15 + count + + +@pytest.mark.parametrize("scenario", ["empty", "unpaired", "measured-zero"]) +async def test_measured_average_distinguishes_unknown_from_measured_zero( + sessions: SessionFactory, + scenario: str, +) -> None: + rows = [] + if scenario == "unpaired": + input_only = _row( + "input", NOW - timedelta(hours=1), cache_read=0, cache_write=0 + ) + input_only.output_source = "missing" + output_only = _row( + "output", NOW - timedelta(hours=1), cache_read=0, cache_write=0 + ) + output_only.input_source = "missing" + rows = [input_only, output_only] + elif scenario == "measured-zero": + rows = [ + _row( + "zero", + NOW - timedelta(hours=1), + input_tokens=0, + output_tokens=0, + cache_read=0, + cache_write=0, + revenue=0, + ) + ] + async with sessions() as session: + session.add_all(rows) + await session.commit() + result = await get_ledger_usage_dashboard( + _legacy(), interval=60, hours=24, session_factory=sessions, now=NOW + ) + summary = result["summary"] + assert summary["measured_token_requests"] == ( + 1 if scenario == "measured-zero" else 0 + ) + assert summary["measured_tokens"] == 0 + assert summary["avg_measured_tokens_per_completion"] == ( + 0 if scenario == "measured-zero" else None + ) + + +@pytest.mark.parametrize("source", ["reported", "estimated", "missing"]) +async def test_measured_average_admits_only_reported_sources( + sessions: SessionFactory, + source: str, +) -> None: + row = _row( + "provenance", + NOW - timedelta(hours=1), + cache_read=0, + cache_write=0, + source=source, + ) + async with sessions() as session: + session.add(row) + await session.commit() + result = await get_ledger_usage_dashboard( + _legacy(), interval=60, hours=24, session_factory=sessions, now=NOW + ) + eligible = source == "reported" + assert result["summary"]["measured_token_requests"] == int(eligible) + assert result["summary"]["measured_tokens"] == (15 if eligible else 0) + assert result["ledger_coverage"]["token_sources"]["input"] == { + "reported": int(eligible), + "estimated": int(source == "estimated"), + "missing": int(not eligible and source != "estimated"), + } + + +async def test_measured_average_uses_exact_historical_bounds( + sessions: SessionFactory, +) -> None: + start, end = datetime(2026, 9, 10, tzinfo=UTC), datetime(2026, 9, 11, tzinfo=UTC) + async with sessions() as session: + session.add_all( + [ + _row("before", start - timedelta(milliseconds=1)), + _row( + "start", + start, + input_tokens=8, + output_tokens=4, + cache_read=0, + cache_write=0, + ), + _row( + "last", + end - timedelta(milliseconds=1), + input_tokens=4, + output_tokens=2, + cache_read=0, + cache_write=0, + ), + _row("end", end), + _row("recent", NOW - timedelta(hours=1)), + ] + ) + await session.commit() + result = await get_ledger_usage_dashboard( + _legacy(), + interval=60, + hours=24, + session_factory=sessions, + now=NOW, + start_at=start, + end_at=end, + ) + summary = result["summary"] + assert summary["measured_token_requests"] == 2 + assert summary["measured_tokens"] == 18 + assert summary["avg_measured_tokens_per_completion"] == 9 diff --git a/tests/unit/test_nostr_stats_identity.py b/tests/unit/test_nostr_stats_identity.py new file mode 100644 index 00000000..8adf4026 --- /dev/null +++ b/tests/unit/test_nostr_stats_identity.py @@ -0,0 +1,60 @@ +from __future__ import annotations + +from typing import Any + +import pytest + +from routstr.nostr import listing + +KEY = "11" * 32 +KEYPAIR = listing.nsec_to_keypair(KEY) +assert KEYPAIR is not None +PUBKEY = KEYPAIR[1] + + +@pytest.mark.asyncio +async def test_explicit_identity_works_without_relay_reads(monkeypatch: Any) -> None: + monkeypatch.setattr(listing.settings, "provider_id", "my-node") + assert await listing.resolve_provider_id_strict(PUBKEY, []) == "my-node" + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "scenario", ["outage", "empty", "ambiguous", "invalid", "truncated"] +) +async def test_unsafe_identity_history_requires_operator_choice( + monkeypatch: Any, scenario: str +) -> None: + first = listing.create_listing_event(KEY, "one", ["https://node.example"]) + second = listing.create_listing_event(KEY, "two", ["https://node.example"]) + forged = {**first, "sig": "00" * 64} + result = { + "outage": ([], False), + "empty": ([], True), + "ambiguous": ([first, second], True), + "invalid": ([forged], True), + "truncated": ([first] * 10, True), + }[scenario] + + async def query(*args: object) -> tuple[list, bool]: + return result + + monkeypatch.setattr(listing.settings, "provider_id", "") + monkeypatch.setattr(listing, "query_listing_events", query) + with pytest.raises(ValueError, match="PROVIDER_ID"): + await listing.resolve_provider_id_strict(PUBKEY, ["wss://relay.example"]) + + +@pytest.mark.asyncio +async def test_one_signed_listing_coordinate_is_reused(monkeypatch: Any) -> None: + event = listing.create_listing_event(KEY, "one", ["https://node.example"]) + + async def query(*args: object) -> tuple[list, bool]: + return [event], True + + monkeypatch.setattr(listing.settings, "provider_id", "") + monkeypatch.setattr(listing, "query_listing_events", query) + assert ( + await listing.resolve_provider_id_strict(PUBKEY, ["wss://relay.example"]) + == "one" + ) diff --git a/tests/unit/test_settings.py b/tests/unit/test_settings.py index 61f657ac..d7391033 100644 --- a/tests/unit/test_settings.py +++ b/tests/unit/test_settings.py @@ -5,9 +5,10 @@ from pathlib import Path import pytest from pydantic.v1 import ValidationError from sqlalchemy.ext.asyncio import create_async_engine -from sqlmodel import text +from sqlmodel import SQLModel, text from sqlmodel.ext.asyncio.session import AsyncSession +from routstr.core import db # noqa: F401 registers the app tables from routstr.core.settings import ENV_ONLY_FIELDS, Settings, SettingsService, settings NSEC_HEX = "1" * 64 @@ -41,6 +42,9 @@ async def test_settings_db_precedence_over_env() -> None: os.environ["ENABLE_ANALYTICS_SHARING"] = "true" engine = create_async_engine("sqlite+aiosqlite:///:memory:") + # Saving a sharing opt-out also fences delivery, which lives in the app schema. + async with engine.begin() as connection: + await connection.run_sync(SQLModel.metadata.create_all) async with AsyncSession(engine, expire_on_commit=False) as session: _ = await SettingsService.initialize(session) updated = await SettingsService.update( @@ -223,7 +227,7 @@ async def test_settings_initialize_discards_unknown_keys() -> None: # Simulate older persisted key name and an unknown key. await session.exec( # type: ignore text("UPDATE settings SET data = :data WHERE id = 1").bindparams( - data='{"name":"LegacyNode","nostr_analytics_enabled":false,"unknown_key":123}' + data='{"name":"LegacyNode","nostr_analytics_enabled":false,"unknown_key":123,"enable_analytics_v2":true,"enable_analytics_collection":true}' ) ) await session.commit() @@ -237,6 +241,8 @@ async def test_settings_initialize_discards_unknown_keys() -> None: assert '"enable_analytics_sharing": true' in stored_data assert "nostr_analytics_enabled" not in stored_data assert "unknown_key" not in stored_data + assert "enable_analytics_v2" not in stored_data + assert "enable_analytics_collection" not in stored_data # ── Secret fields are never written to the settings blob (issue #553) ──────── diff --git a/tests/unit/test_terminal_outcomes.py b/tests/unit/test_terminal_outcomes.py index 9ffed82c..b73fcc2f 100644 --- a/tests/unit/test_terminal_outcomes.py +++ b/tests/unit/test_terminal_outcomes.py @@ -760,6 +760,60 @@ async def test_unclean_restart_preserves_days_before_last_durable_flush( assert await writer.stop(timeout=1) +async def test_sharing_disabled_startup_resumes_private_coverage_after_gap( + ledger: tuple[AsyncEngine, SessionFactory], + monkeypatch: pytest.MonkeyPatch, +) -> None: + from nostr_sdk import Keys + + from routstr.nostr import analytics_runtime as runtime + + _, sessions = ledger + day = date(2026, 9, 1) + clock = MutableClock(_timestamp(day)) + writer = TerminalOutcomeWriter(session_factory=sessions, clock=clock) + assert await writer.start() + await runtime.claim_analytics_v2_identity( + sessions, + pubkey=Keys.parse("11" * 32).public_key().to_hex(), + provider_d="provider", + at_ms=clock.value, + ) + await runtime.activate_analytics_v2_sharing( + sessions, coverage_day=day, at_ms=clock.value + ) + clock.value = _timestamp(day + timedelta(days=2)) + assert await writer.flush(timeout=1) + assert await writer.stop(timeout=1) + clock.value = _timestamp(day + timedelta(days=5)) + stopped_writer = TerminalOutcomeWriter(session_factory=sessions, clock=clock) + monkeypatch.setattr(outcomes_module, "terminal_outcome_writer", stopped_writer) + monkeypatch.setattr(runtime, "terminal_outcome_writer", stopped_writer) + monkeypatch.setattr(runtime, "create_session", sessions) + monkeypatch.setattr(runtime.settings, "enable_analytics_sharing", False) + coordinator = runtime.AnalyticsCoordinator() + await coordinator.prepare_startup() + assert stopped_writer.running + await coordinator.close() + assert not (await runtime.get_analytics_v2_delivery_state(sessions)).sharing_enabled + async with sessions() as session: + epochs = ( + await session.exec( + select(TerminalOutcomeEpoch).order_by(col(TerminalOutcomeEpoch.epoch)) + ) + ).all() + assert all(epoch.current_slot is None for epoch in epochs[:-1]) + assert epochs[0].coverage_end_day == day + timedelta(days=1) + assert all( + epoch.coverage_end_day is not None + and epoch.coverage_start_day > epoch.coverage_end_day + for epoch in epochs[1:-1] + ) + assert epochs[-1].current_slot == 1 + assert epochs[-1].coverage_start_day == day + timedelta(days=6) + assert epochs[-1].coverage_end_day is None + + @pytest.mark.parametrize("fresh_process", [False, True]) async def test_failed_writer_start_cannot_backfill_missed_days_as_zero( ledger: tuple[AsyncEngine, SessionFactory], diff --git a/tests/unit/test_usage_analytics_store_utc.py b/tests/unit/test_usage_analytics_store_utc.py new file mode 100644 index 00000000..cbdffbbe --- /dev/null +++ b/tests/unit/test_usage_analytics_store_utc.py @@ -0,0 +1,34 @@ +import time +from pathlib import Path + +import pytest + +from routstr.core.usage_analytics_store import UsageAnalyticsStore + + +def test_local_log_minutes_are_bucketed_in_utc( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("TZ", "Asia/Kolkata") + time.tzset() + try: + store = UsageAnalyticsStore(logs_dir=tmp_path) + conn = store._get_connection_locked() + conn.execute( + "INSERT INTO analytics_minute (minute_ts, total_requests) VALUES (?, 1)", + ("2026-09-18 16:30:00",), + ) + metrics = store._query_metrics_locked( + conn, + cutoff_timestamp="2026-09-18 00:00:00", + interval_minutes=60, + hours_back=24, + ) + finally: + monkeypatch.delenv("TZ") + time.tzset() + + # 16:30 IST is 11:00 UTC; rounding the local hour first would give 10:00. + assert [point["timestamp"] for point in metrics["metrics"]] == [ + "2026-09-18 11:00:00" + ] diff --git a/ui/app/page.tsx b/ui/app/page.tsx index 134483dd..de67ef7e 100644 --- a/ui/app/page.tsx +++ b/ui/app/page.tsx @@ -109,8 +109,12 @@ function getRangeHours(range?: DateRange): number | null { } const normalized = normalizeDateRange(range); - const fromTime = normalized.from?.getTime(); - const toTime = normalized.to?.getTime(); + const from = normalized.from; + const to = normalized.to; + const fromTime = + from && Date.UTC(from.getFullYear(), from.getMonth(), from.getDate()); + const toTime = + to && Date.UTC(to.getFullYear(), to.getMonth(), to.getDate() + 1); if (fromTime === undefined || toTime === undefined) { return null; @@ -136,6 +140,9 @@ function formatCompactDateRangeLabel(range?: DateRange): string { } const sameMonth = format(from, 'yyyy-MM') === format(to, 'yyyy-MM'); + if (format(from, 'yyyy-MM-dd') === format(to, 'yyyy-MM-dd')) { + return format(from, 'MMM d, yyyy'); + } if (sameMonth) { return `${format(from, 'MMM d')} - ${format(to, 'd')}`; } @@ -543,7 +550,34 @@ export default function DashboardPage() { ? customRangeHours : activePreset.hours; const safeQueryHours = Math.min(queryHours, MAX_USAGE_RANGE_HOURS); + const queryRange = + isCustomRangeActive && customRange?.from && customRange.to + ? (() => { + const normalized = normalizeDateRange(customRange); + const from = normalized.from!; + const to = normalized.to!; + const end = Date.UTC( + to.getFullYear(), + to.getMonth(), + to.getDate() + 1 + ); + const start = Math.max( + Date.UTC(from.getFullYear(), from.getMonth(), from.getDate()), + end - MAX_USAGE_RANGE_HOURS * 3600000 + ); + return { + start: new Date(start).toISOString(), + end: new Date(end).toISOString(), + }; + })() + : undefined; const isUsageRangeCapped = safeQueryHours < queryHours; + const today = new Date(); + const latestCalendarDay = new Date( + today.getUTCFullYear(), + today.getUTCMonth(), + today.getUTCDate() + ); const autoInterval = getAutoIntervalMinutes(safeQueryHours); const usageRefetchIntervalMs = useMemo(() => { if (safeQueryHours > 90 * 24) { @@ -576,11 +610,24 @@ export default function DashboardPage() { const { data: usageDashboardData, isLoading: usageDashboardLoading, + error: usageDashboardError, refetch: refetchUsageDashboard, } = useQuery({ - queryKey: ['usage-dashboard', autoInterval, safeQueryHours], + queryKey: [ + 'usage-dashboard', + autoInterval, + safeQueryHours, + queryRange?.start, + queryRange?.end, + ], queryFn: () => - AdminService.getUsageDashboard(safeQueryHours, autoInterval, 100, 20), + AdminService.getUsageDashboard( + safeQueryHours, + autoInterval, + 100, + 20, + queryRange + ), enabled: isAuthenticated, refetchInterval: usageRefetchIntervalMs, staleTime: 30_000, @@ -590,6 +637,10 @@ export default function DashboardPage() { const summaryData = usageDashboardData?.summary; const errorData = usageDashboardData?.error_details; const modelUsageMixData = usageDashboardData?.model_usage_mix; + const ledgerMode = + usageDashboardData?.analytics_source === 'terminal_outcomes'; + const ledgerCoverage = usageDashboardData?.ledger_coverage; + const diagnosticsAvailable = ledgerCoverage?.diagnostic_available !== false; const hasModelUsageMixMetrics = Array.isArray(modelUsageMixData?.metrics) && modelUsageMixData.metrics.length > 0; @@ -620,11 +671,14 @@ export default function DashboardPage() { const revenuePoints = metricsData.metrics.map( (metric: UsageMetricData) => ({ ...metric, - revenue_display: convertRevenueMsats(metric.revenue_msats), + revenue_display: + metric.revenue_msats === null + ? null + : convertRevenueMsats(metric.revenue_msats), }) ) as ChartDatum[]; - return [ + const configs: ChartConfig[] = [ { id: 'revenue', title: 'Revenue Over Time', @@ -764,7 +818,17 @@ export default function DashboardPage() { ], }, ]; - }, [metricsData, metricsTotals, revenueDisplayUnit, usdPerSat]); + return configs.filter( + (config) => + diagnosticsAvailable || ['revenue', 'tokens'].includes(config.id) + ); + }, [ + metricsData, + metricsTotals, + revenueDisplayUnit, + usdPerSat, + diagnosticsAvailable, + ]); useEffect(() => { if (chartConfigs.length === 0) { @@ -859,8 +923,7 @@ export default function DashboardPage() { // DayPicker may emit from===to on the first click in range mode. // Keep waiting until the user explicitly picks a second (end) date. - const isSameDay = to ? from.getTime() === to.getTime() : false; - if (!hasPreviousStart || !to || isSameDay) { + if (!hasPreviousStart || !to) { setPendingCustomRange({ from, to: undefined }); return; } @@ -903,6 +966,16 @@ export default function DashboardPage() {

Usage Analytics

+ {(ledgerCoverage?.complete === false || + metricsData?.bucket_fill_complete === false || + modelUsageMixData?.bucket_fill_complete === false) && ( + + Partial data + + )}

All cards and charts in this section update from the selected @@ -914,6 +987,11 @@ export default function DashboardPage() { {MAX_USAGE_RANGE_HOURS / 24} days for server safety.

) : null} + {usageDashboardError && ( +

+ Stats could not be loaded. Try Refresh to retry this period. +

+ )}
@@ -929,6 +1007,7 @@ export default function DashboardPage() { variant='ghost' size='icon' id='dashboard-date-range' + disabled={!ledgerMode} className='text-muted-foreground hover:bg-muted/50 hover:text-foreground dark:hover:bg-input/50 h-full w-8 rounded-none border-0 bg-transparent p-0 sm:w-9' aria-label='Open custom date range' > @@ -941,6 +1020,7 @@ export default function DashboardPage() { > @@ -1044,16 +1125,21 @@ export default function DashboardPage() { {summaryLoading ? ( ) : summaryData ? ( - + ) : null} - {errorLoading ? ( - - ) : errorData ? ( - - ) : null} + {diagnosticsAvailable && + (errorLoading ? ( + + ) : errorData ? ( + + ) : null)}
diff --git a/ui/components/settings/admin-settings.tsx b/ui/components/settings/admin-settings.tsx index 68045b3f..085dacd5 100644 --- a/ui/components/settings/admin-settings.tsx +++ b/ui/components/settings/admin-settings.tsx @@ -129,8 +129,14 @@ export function AdminSettings() { settingsPayload = { ...settings, npub: result.npub }; } + // Send only what changed here, so a tab opened earlier cannot write back + // choices saved since, such as a stats opt-out. const updatedData = (await AdminService.updateSettings( - settingsPayload + Object.fromEntries( + Object.entries(settingsPayload).filter( + ([key, value]) => !areValuesEqual(value, initialSettings[key]) + ) + ) )) as SettingsData; setSettings(updatedData); setInitialSettings(updatedData); diff --git a/ui/components/top-models-usage-chart.tsx b/ui/components/top-models-usage-chart.tsx index addc722d..2453244f 100644 --- a/ui/components/top-models-usage-chart.tsx +++ b/ui/components/top-models-usage-chart.tsx @@ -14,9 +14,11 @@ import { useIsMobile } from '@/hooks/use-mobile'; import { type ModelUsageMix } from '@/lib/api/services/admin'; import type { DisplayUnit } from '@/lib/types/units'; import { cn } from '@/lib/utils'; +import { parseBucketDate } from '@/lib/usage-time'; interface TopModelsUsageChartProps { mix: ModelUsageMix; + completeCoverage?: boolean; displayUnit: DisplayUnit; usdPerSat: number | null; } @@ -39,22 +41,10 @@ interface LeaderboardRow { provider: string; rank: number; totalRaw: number; - trend: LeaderboardTrend; + trend: LeaderboardTrend | null; trendPercent: number | null; } -function parseBucketDate(value: string): Date | null { - const normalized = value.includes('T') - ? value - : `${value.replace(' ', 'T')}Z`; - const parsed = new Date(normalized); - if (!Number.isNaN(parsed.getTime())) { - return parsed; - } - const fallback = new Date(value); - return Number.isNaN(fallback.getTime()) ? null : fallback; -} - function hueFromString(input: string): number { let hash = 0; for (let i = 0; i < input.length; i += 1) { @@ -98,6 +88,7 @@ function formatTooltipTimestamp( const shouldShowTime = intervalMinutes <= 6 * 60 || hoursBack <= 48; if (shouldShowTime) { return date.toLocaleString([], { + timeZone: 'UTC', month: 'long', day: 'numeric', year: 'numeric', @@ -106,6 +97,7 @@ function formatTooltipTimestamp( }); } return date.toLocaleString([], { + timeZone: 'UTC', month: 'long', day: 'numeric', year: 'numeric', @@ -126,6 +118,7 @@ function formatAxisTimestamp( const shouldShowTime = intervalMinutes <= 6 * 60 || hoursBack <= 48; if (shouldShowTime && hasMultipleDays) { return date.toLocaleString([], { + timeZone: 'UTC', month: 'short', day: 'numeric', hour: '2-digit', @@ -135,6 +128,7 @@ function formatAxisTimestamp( if (shouldShowTime) { return date.toLocaleTimeString([], { + timeZone: 'UTC', hour: '2-digit', minute: '2-digit', }); @@ -142,6 +136,7 @@ function formatAxisTimestamp( if (intervalMinutes >= 24 * 60 && hoursBack >= 24 * 180) { return date.toLocaleDateString([], { + timeZone: 'UTC', month: 'short', year: '2-digit', }); @@ -149,12 +144,14 @@ function formatAxisTimestamp( if (hasMultipleDays) { return date.toLocaleDateString([], { + timeZone: 'UTC', month: 'short', day: 'numeric', }); } return date.toLocaleTimeString([], { + timeZone: 'UTC', hour: '2-digit', minute: '2-digit', }); @@ -240,6 +237,7 @@ function getModelPresentation(model: string): { export function TopModelsUsageChart({ mix, + completeCoverage, displayUnit, usdPerSat, }: TopModelsUsageChartProps) { @@ -304,8 +302,9 @@ export function TopModelsUsageChart({ const modelCounts = metric.model_counts ?? {}; const modelRevenue = metric.model_revenue_msats ?? {}; const modelTokens = metric.model_tokens ?? {}; - const point: Record = { + const point: Record = { timestamp: metric.timestamp, + coverage: metric.coverage ?? 'complete', total_successful: metric.total_successful, total_revenue_msats: metric.total_revenue_msats, total_tokens: metric.total_tokens, @@ -315,9 +314,18 @@ export function TopModelsUsageChart({ }; for (const item of series) { - point[item.requestsKey] = modelCounts[item.label] ?? 0; - point[item.revenueKey] = modelRevenue[item.label] ?? 0; - point[item.tokensKey] = modelTokens[item.label] ?? 0; + point[item.requestsKey] = + metric.total_successful === null + ? null + : (modelCounts[item.label] ?? 0); + point[item.revenueKey] = + metric.total_revenue_msats === null + ? null + : (modelRevenue[item.label] ?? 0); + point[item.tokensKey] = + metric.total_tokens === null + ? null + : (modelTokens[item.label] ?? 0); } return point; @@ -328,7 +336,7 @@ export function TopModelsUsageChart({ const hasMultipleDays = useMemo(() => { const daySet = new Set( chartData.map((item) => - parseBucketDate(String(item.timestamp))?.toDateString() + parseBucketDate(String(item.timestamp))?.toISOString().slice(0, 10) ) ); return daySet.size > 1; @@ -471,6 +479,15 @@ export function TopModelsUsageChart({ const rounded = abs >= 10 ? abs.toFixed(0) : abs.toFixed(1); return rounded.replace(/\.0$/, ''); }; + const canComparePeriods = + completeCoverage !== false && + mix.bucket_fill_complete !== false && + mixMetrics.every( + (metric) => + !metric.coverage || + metric.coverage === 'complete' || + metric.coverage === 'updating' + ); const leaderboardRows = useMemo(() => { if (leaderboardModels.length === 0 || mixMetrics.length === 0) { return []; @@ -504,12 +521,12 @@ export function TopModelsUsageChart({ 0 ); const trendPercent = - previousRaw > 0 + canComparePeriods && previousRaw > 0 ? ((currentRaw - previousRaw) / previousRaw) * 100 : null; - let trend: LeaderboardTrend = 'flat'; - if (previousRaw <= 0 && currentRaw > 0) { + let trend: LeaderboardTrend | null = canComparePeriods ? 'flat' : null; + if (canComparePeriods && previousRaw <= 0 && currentRaw > 0) { trend = 'new'; } else if (trendPercent !== null && trendPercent > 0.5) { trend = 'up'; @@ -547,7 +564,7 @@ export function TopModelsUsageChart({ })); return rows; - }, [leaderboardModels, mixMetrics, mode, series]); + }, [canComparePeriods, leaderboardModels, mixMetrics, mode, series]); if (chartData.length === 0) { return null; @@ -676,17 +693,22 @@ export function TopModelsUsageChart({ /> { if (!isChartPointerInside || !active || !payload?.length) { return null; } + const point = payload[0]?.payload; + const missing = + point?.coverage === 'missing' || + payload.every((entry) => entry.value === null); const rows = payload .map((entry) => { const value = typeof entry.value === 'number' ? entry.value - : Number(entry.value || 0); + : Number.NaN; return { color: String(entry.color || '#6b7280'), @@ -701,9 +723,6 @@ export function TopModelsUsageChart({ .sort((a, b) => b.value - a.value); const total = rows.reduce((sum, row) => sum + row.value, 0); - if (rows.length === 0) { - return null; - } return (
@@ -714,6 +733,17 @@ export function TopModelsUsageChart({ mix.hours_back )}

+ {missing ? ( +

+ No collection data for this interval. +

+ ) : point?.coverage === 'partial' ? ( +

+ Partial collection. Recorded activity only. +

+ ) : point?.coverage === 'updating' ? ( +

Still updating.

+ ) : null}
{rows.map((row) => (
Total - {formatValue(total)} + {missing ? 'Unavailable' : formatValue(total)}
@@ -782,7 +812,9 @@ export function TopModelsUsageChart({ Top models

- Change vs prior period + {canComparePeriods + ? 'Recent half vs earlier half' + : 'Recorded activity only'}

@@ -852,9 +884,11 @@ export function TopModelsUsageChart({ {formatLeaderboardTotal(row.totalRaw)} - - {trendLabel} - + {row.trend !== null ? ( + + {trendLabel} + + ) : null} ); })} diff --git a/ui/components/usage-metrics-chart.tsx b/ui/components/usage-metrics-chart.tsx index 3d8c8bf1..be3f1a41 100644 --- a/ui/components/usage-metrics-chart.tsx +++ b/ui/components/usage-metrics-chart.tsx @@ -1,5 +1,7 @@ 'use client'; +import { parseBucketDate } from '@/lib/usage-time'; + import { useEffect, useMemo, useRef, useState } from 'react'; import { Card, CardContent, CardHeader, CardTitle } from '@/components/ui/card'; import { Button } from '@/components/ui/button'; @@ -60,19 +62,22 @@ export function UsageMetricsChart({ const hasMultipleDays = useMemo(() => { const daySet = new Set( - data.map((item) => new Date(item.timestamp).toDateString()) + data.map((item) => + parseBucketDate(item.timestamp)?.toISOString().slice(0, 10) + ) ); return daySet.size > 1; }, [data]); const formatAxisTick = (timestamp: string): string => { - const date = new Date(timestamp); - if (Number.isNaN(date.getTime())) { + const date = parseBucketDate(timestamp); + if (!date) { return ''; } if (hasMultipleDays) { return date.toLocaleString([], { + timeZone: 'UTC', month: 'short', day: 'numeric', hour: '2-digit', @@ -80,12 +85,14 @@ export function UsageMetricsChart({ } return date.toLocaleTimeString([], { + timeZone: 'UTC', hour: '2-digit', minute: '2-digit', }); }; - const formatMetricValue = (value: number): string => { + const formatMetricValue = (value: number | null): string => { + if (value === null) return 'Unavailable'; const formatted = compactNumber.format(value); return metricType === 'currency' ? `${formatted} ${currencyUnitLabel}` @@ -93,9 +100,9 @@ export function UsageMetricsChart({ }; const metricTotals = useMemo(() => { - const fallbackTotals = dataKeys.reduce>( + const fallbackTotals = dataKeys.reduce>( (acc, dataKey) => { - acc[dataKey.key] = 0; + acc[dataKey.key] = null; return acc; }, {} @@ -104,10 +111,12 @@ export function UsageMetricsChart({ for (const point of data) { for (const dataKey of dataKeys) { const rawValue = point?.[dataKey.key]; + if (rawValue === null || rawValue === undefined) continue; const value = typeof rawValue === 'number' ? rawValue : Number(rawValue || 0); if (Number.isFinite(value)) { - fallbackTotals[dataKey.key] += value; + fallbackTotals[dataKey.key] = + (fallbackTotals[dataKey.key] ?? 0) + value; } } } @@ -119,7 +128,11 @@ export function UsageMetricsChart({ const mergedTotals = { ...fallbackTotals }; for (const dataKey of dataKeys) { const rawTotal = totals[dataKey.key]; - if (typeof rawTotal === 'number' && Number.isFinite(rawTotal)) { + if ( + mergedTotals[dataKey.key] !== null && + typeof rawTotal === 'number' && + Number.isFinite(rawTotal) + ) { mergedTotals[dataKey.key] = rawTotal; } } @@ -131,9 +144,7 @@ export function UsageMetricsChart({ () => dataKeys.map((dataKey) => ({ ...dataKey, - value: Number.isFinite(metricTotals[dataKey.key]) - ? metricTotals[dataKey.key] - : 0, + value: metricTotals[dataKey.key], })), [dataKeys, metricTotals] ); @@ -351,13 +362,45 @@ export function UsageMetricsChart({ /> - new Date(String(label)).toLocaleString() - } - /> - } + filterNull={false} + content={(props) => { + const timestamp = + parseBucketDate(String(props.label))?.toLocaleString([], { + timeZone: 'UTC', + timeZoneName: 'short', + }) ?? String(props.label ?? ''); + if ( + props.active && + props.payload?.length && + props.payload.every((entry) => entry.value === null) + ) { + return ( +
+

{timestamp}

+

+ No collection data for this interval. +

+
+ ); + } + return ( + { + const coverage = props.payload?.[0]?.payload?.coverage; + const note = + coverage === 'partial' + ? ' (partial collection)' + : coverage === 'updating' + ? ' (still updating)' + : ''; + return `${timestamp}${note}`; + }} + /> + ); + }} /> {visibleDataKeys.map((dataKey) => ( ))} diff --git a/ui/components/usage-summary-cards.tsx b/ui/components/usage-summary-cards.tsx index 7db46735..81148bc8 100644 --- a/ui/components/usage-summary-cards.tsx +++ b/ui/components/usage-summary-cards.tsx @@ -18,12 +18,19 @@ import { useCurrencyStore } from '@/lib/stores/currency'; import { useQuery } from '@tanstack/react-query'; import { fetchBtcUsdPrice, btcToSatsRate } from '@/lib/exchange-rate'; import { formatFromMsat } from '@/lib/currency'; +import { formatCost } from '@/lib/services/cost-validation'; interface UsageSummaryCardsProps { summary: UsageSummary; + ledgerMode?: boolean; + diagnosticsAvailable?: boolean; } -export function UsageSummaryCards({ summary }: UsageSummaryCardsProps) { +export function UsageSummaryCards({ + summary, + ledgerMode = false, + diagnosticsAvailable = true, +}: UsageSummaryCardsProps) { const { displayUnit } = useCurrencyStore(); const { data: btcUsdPrice } = useQuery({ queryKey: ['btc-usd-price'], @@ -33,16 +40,33 @@ export function UsageSummaryCards({ summary }: UsageSummaryCardsProps) { }); const usdPerSat = btcUsdPrice ? btcToSatsRate(btcUsdPrice) : null; - const formatAmount = (msat: number) => - formatFromMsat(msat, displayUnit, usdPerSat); + const formatAmount = (msat: number) => { + if (displayUnit === 'sat') { + if (msat > 0 && msat < 1) return '<0.001 sats'; + return `${(msat / 1000).toLocaleString(undefined, { maximumFractionDigits: 3 })} sats`; + } + if (displayUnit === 'usd' && usdPerSat !== null) { + return msat === 0 ? '$0.00' : formatCost((msat / 1000) * usdPerSat); + } + return formatFromMsat(msat, displayUnit, usdPerSat); + }; const totalTokens = Number(summary.total_tokens ?? 0); - const avgTotalTokensPerCompletion = Number( - summary.avg_total_tokens_per_completion ?? 0 - ); + const avgTotalTokensPerCompletion = ledgerMode + ? summary.avg_measured_tokens_per_completion + : Number(summary.avg_total_tokens_per_completion ?? 0); + const averageTokensValue = + summary.successful_chat_completions === 0 + ? 'No requests' + : avgTotalTokensPerCompletion == null + ? 'Not reported' + : avgTotalTokensPerCompletion.toLocaleString(undefined, { + maximumFractionDigits: 1, + }); const cards = [ { - title: 'Total Requests', + title: ledgerMode ? 'Logged Requests' : 'Total Requests', + diagnostic: true, value: summary.total_requests.toLocaleString(), icon: Activity, iconClassName: 'text-blue-600 dark:text-blue-300', @@ -61,9 +85,10 @@ export function UsageSummaryCards({ summary }: UsageSummaryCardsProps) { }, { title: 'Avg Tokens/Completion', - value: avgTotalTokensPerCompletion.toLocaleString(undefined, { - maximumFractionDigits: 1, - }), + description: ledgerMode + ? 'Average for completions with provider-reported input and output tokens.' + : undefined, + value: averageTokensValue, icon: Activity, iconClassName: 'text-indigo-600 dark:text-indigo-300', }, @@ -75,42 +100,48 @@ export function UsageSummaryCards({ summary }: UsageSummaryCardsProps) { }, { title: 'Operational Net', + hidden: ledgerMode, value: formatAmount(summary.net_revenue_msats), icon: DollarSign, iconClassName: 'text-lime-600 dark:text-lime-300', }, { title: 'Reverted Holds', + diagnostic: true, value: formatAmount(summary.refunds_msats), icon: TrendingDown, iconClassName: 'text-rose-600 dark:text-rose-300', }, { - title: 'Avg Revenue/Request', + title: ledgerMode ? 'Avg Revenue/Completion' : 'Avg Revenue/Request', value: formatAmount(summary.avg_revenue_per_request_msats), icon: CreditCard, iconClassName: 'text-violet-600 dark:text-violet-300', }, { - title: 'Success Rate', + title: ledgerMode ? 'Logged Success Rate' : 'Success Rate', + diagnostic: true, value: `${summary.success_rate.toFixed(1)}%`, icon: TrendingUp, iconClassName: 'text-teal-600 dark:text-teal-300', }, { title: 'Refund Rate', + diagnostic: true, value: `${summary.refund_rate.toFixed(1)}%`, icon: XCircle, iconClassName: 'text-fuchsia-600 dark:text-fuchsia-300', }, { title: 'Failed Requests', + diagnostic: true, value: summary.failed_requests.toLocaleString(), icon: XCircle, iconClassName: 'text-red-600 dark:text-red-300', }, { title: 'Errors', + diagnostic: true, value: summary.total_errors.toLocaleString(), icon: AlertTriangle, iconClassName: 'text-orange-600 dark:text-orange-300', @@ -123,18 +154,24 @@ export function UsageSummaryCards({ summary }: UsageSummaryCardsProps) { }, { title: 'Upstream Errors', + diagnostic: true, value: summary.upstream_errors.toLocaleString(), icon: AlertTriangle, iconClassName: 'text-pink-600 dark:text-pink-300', }, - ]; + ].filter( + (card) => !card.hidden && (!card.diagnostic || diagnosticsAvailable) + ); return (
{cards.map((card) => ( - + {card.title} diff --git a/ui/lib/api/services/admin.ts b/ui/lib/api/services/admin.ts index 15daa38d..f4aed30b 100644 --- a/ui/lib/api/services/admin.ts +++ b/ui/lib/api/services/admin.ts @@ -891,13 +891,18 @@ export class AdminService { hours: number = 24, interval: number = 15, errorLimit: number = 100, - modelLimit: number = 20 + modelLimit: number = 20, + range?: { start: string; end: string } ): Promise { const params = new URLSearchParams(); params.set('interval', String(interval)); params.set('hours', String(hours)); params.set('error_limit', String(errorLimit)); params.set('model_limit', String(modelLimit)); + if (range) { + params.set('start_at', range.start); + params.set('end_at', range.end); + } return await apiClient.get( `/admin/api/usage/dashboard?${params.toString()}` @@ -1123,22 +1128,24 @@ export interface TemporaryBalancesResponse { export interface UsageMetricData { timestamp: string; + coverage?: 'complete' | 'updating' | 'partial' | 'missing'; total_requests: number; - successful_chat_completions: number; + successful_chat_completions: number | null; failed_requests: number; errors: number; warnings: number; payment_processed: number; upstream_errors: number; - revenue_msats: number; + revenue_msats: number | null; refunds_msats: number; - input_tokens: number; - output_tokens: number; - total_tokens: number; + input_tokens: number | null; + output_tokens: number | null; + total_tokens: number | null; [key: string]: unknown; } export interface UsageMetrics { + bucket_fill_complete?: boolean; metrics: UsageMetricData[]; interval_minutes: number; hours_back: number; @@ -1177,6 +1184,9 @@ export interface UsageSummary { avg_input_tokens_per_completion: number; avg_output_tokens_per_completion: number; avg_total_tokens_per_completion: number; + measured_token_requests?: number; + measured_tokens?: number; + avg_measured_tokens_per_completion?: number | null; success_rate: number; revenue_msats: number; refunds_msats: number; @@ -1221,18 +1231,20 @@ export interface RevenueByModel { export interface ModelUsageMixMetric { timestamp: string; - total_successful: number; - total_revenue_msats: number; - total_tokens: number; - others: number; - others_revenue_msats: number; - others_tokens: number; + coverage?: 'complete' | 'updating' | 'partial' | 'missing'; + total_successful: number | null; + total_revenue_msats: number | null; + total_tokens: number | null; + others: number | null; + others_revenue_msats: number | null; + others_tokens: number | null; model_counts: Record; model_revenue_msats: Record; model_tokens: Record; } export interface ModelUsageMix { + bucket_fill_complete?: boolean; top_models: string[]; metrics: ModelUsageMixMetric[]; interval_minutes: number; @@ -1241,6 +1253,20 @@ export interface ModelUsageMix { } export interface UsageDashboardResponse { + analytics_source?: 'terminal_outcomes'; + ledger_coverage?: { + from: string; + to: string; + complete: boolean; + incomplete_days: string[]; + includes_current_day: boolean; + latest_outcome_at: string | null; + diagnostic_available: boolean; + token_sources: Record< + string, + { reported: number; estimated: number; missing: number } + >; + }; metrics: UsageMetrics; summary: UsageSummary; error_details: ErrorDetails; diff --git a/ui/lib/usage-time.ts b/ui/lib/usage-time.ts new file mode 100644 index 00000000..9047591a --- /dev/null +++ b/ui/lib/usage-time.ts @@ -0,0 +1,7 @@ +export function parseBucketDate(value: string): Date | null { + const normalized = value.includes('T') + ? value + : `${value.replace(' ', 'T')}Z`; + const parsed = new Date(normalized); + return Number.isNaN(parsed.getTime()) ? null : parsed; +}