diff --git a/.env.example b/.env.example index d093e006..d4c768d5 100644 --- a/.env.example +++ b/.env.example @@ -14,6 +14,7 @@ UPSTREAM_API_KEY=your-upstream-api-key # HTTP_URL=https://api.mynode.com # ONION_URL=http://mynode.onion (auto fetched from compose) # RELAYS="wss://relay.damus.io,wss://relay.nostr.band,wss://eden.nostr.land,wss://relay.routstr.com" +# ENABLE_ANALYTICS_SHARING=true # CASHU_MINTS="https://mint.minibits.cash/Bitcoin,https://mint.cubabitcoin.org,https://ecashmint.otrta.me" # RECEIVE_LN_ADDRESS= diff --git a/docs/provider/configuration.md b/docs/provider/configuration.md index 633ccd10..20c0e70c 100644 --- a/docs/provider/configuration.md +++ b/docs/provider/configuration.md @@ -97,6 +97,7 @@ Announce your node on the network: | **Npub** | Your Nostr public key | | **Nsec** | Your Nostr private key (for signing) | | **Relays** | Relays to publish announcements | +| **Share Analytics** | Publish aggregate usage stats to Nostr | See [Discovery](discovery.md) for details. @@ -122,6 +123,7 @@ Use environment variables for: | `DESCRIPTION` | Node description | `A Routstr Node` | | `NPUB` | Nostr public key (bech32) | — | | `NSEC` | Nostr private key | — | +| `ENABLE_ANALYTICS_SHARING` | Enable usage analytics sharing to Nostr | `true` | | `CASHU_MINTS` | Comma-separated mint URLs | `https://mint.minibits.cash/Bitcoin` | | `RECEIVE_LN_ADDRESS` | Lightning address for withdrawals | — | | `TOR_PROXY_URL` | SOCKS5 proxy for Tor | `socks5://127.0.0.1:9050` | diff --git a/docs/provider/dashboard.md b/docs/provider/dashboard.md index aaa860bf..66852845 100644 --- a/docs/provider/dashboard.md +++ b/docs/provider/dashboard.md @@ -152,6 +152,7 @@ Manage which mints you accept payments from: |-------|-------------| | **Nsec** | Private key for signing announcements | | **Relays** | Where to publish your node advertisement | +| **Share Analytics** | Toggle publishing aggregate usage stats to Nostr | ### Security diff --git a/logs/.gitkeep b/logs/.gitkeep index e69de29b..8b137891 100644 --- a/logs/.gitkeep +++ b/logs/.gitkeep @@ -0,0 +1 @@ + diff --git a/routstr/auth.py b/routstr/auth.py index 002f71da..f844f1e7 100644 --- a/routstr/auth.py +++ b/routstr/auth.py @@ -776,6 +776,8 @@ async def adjust_payment_for_tokens( "key_hash": key.hashed_key[:8] + "...", "billing_key_hash": billing_key.hashed_key[:8] + "...", "charged_amount": cost.total_msats, + "input_tokens": cost.input_tokens, + "output_tokens": cost.output_tokens, "new_balance": billing_key.balance, "model": model, }, @@ -799,6 +801,8 @@ async def adjust_payment_for_tokens( "cost_difference": cost_difference, "input_msats": cost.input_msats, "output_msats": cost.output_msats, + "input_tokens": cost.input_tokens, + "output_tokens": cost.output_tokens, }, ) diff --git a/routstr/core/admin.py b/routstr/core/admin.py index ed161fbe..46589164 100644 --- a/routstr/core/admin.py +++ b/routstr/core/admin.py @@ -34,6 +34,8 @@ admin_router = APIRouter(prefix="/admin", include_in_schema=False) admin_sessions: dict[str, int] = {} ADMIN_SESSION_DURATION = 3600 +# Usage analytics remain queryable up to 12 months. +MAX_USAGE_ANALYTICS_HOURS = 365 * 24 def require_admin_api(request: Request) -> None: @@ -1147,16 +1149,57 @@ async def get_usage_metrics( interval: int = Query( default=15, ge=1, le=1440, description="Time interval in minutes" ), - hours: int = Query(default=24, ge=1, description="Hours of history to analyze"), + hours: int = Query( + default=24, + ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, + description="Hours of history to analyze", + ), ) -> dict: """Get usage metrics aggregated by time interval.""" return log_manager.get_usage_metrics(interval=interval, hours=hours) +@admin_router.get("/api/usage/dashboard", dependencies=[Depends(require_admin_api)]) +async def get_usage_dashboard( + request: Request, + interval: int = Query( + default=15, ge=1, le=1440, description="Time interval in minutes" + ), + hours: int = Query( + default=24, + ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, + description="Hours of history to analyze", + ), + error_limit: int = Query( + default=100, ge=1, le=1000, description="Maximum number of errors to return" + ), + model_limit: int = Query( + default=20, ge=1, le=100, description="Maximum number of models to return" + ), +) -> dict: + """ + Get all dashboard analytics in one request. + This runs one combined aggregation pass and avoids repeated scans. + """ + return log_manager.get_usage_dashboard( + interval=interval, + hours=hours, + error_limit=error_limit, + model_limit=model_limit, + ) + + @admin_router.get("/api/usage/summary", dependencies=[Depends(require_admin_api)]) async def get_usage_summary( request: Request, - hours: int = Query(default=24, ge=1, description="Hours of history to analyze"), + hours: int = Query( + default=24, + ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, + description="Hours of history to analyze", + ), ) -> dict: """Get summary statistics for the specified time period.""" return log_manager.get_usage_summary(hours=hours) @@ -1165,7 +1208,12 @@ async def get_usage_summary( @admin_router.get("/api/usage/error-details", dependencies=[Depends(require_admin_api)]) async def get_error_details( request: Request, - hours: int = Query(default=24, ge=1, description="Hours of history to analyze"), + hours: int = Query( + default=24, + ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, + description="Hours of history to analyze", + ), limit: int = Query( default=100, ge=1, le=1000, description="Maximum number of errors to return" ), @@ -1179,7 +1227,12 @@ async def get_error_details( ) async def get_revenue_by_model( request: Request, - hours: int = Query(default=24, ge=1, description="Hours of history to analyze"), + hours: int = Query( + default=24, + ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, + description="Hours of history to analyze", + ), limit: int = Query( default=20, ge=1, le=100, description="Maximum number of models to return" ), diff --git a/routstr/core/log_manager.py b/routstr/core/log_manager.py index 39ac8e7c..0444dcbf 100644 --- a/routstr/core/log_manager.py +++ b/routstr/core/log_manager.py @@ -1,17 +1,73 @@ import json +import time from collections import defaultdict from datetime import datetime, timedelta, timezone +from heapq import heappush, heapreplace from pathlib import Path -from typing import Any, Iterator +from threading import Lock +from typing import Any, Callable, Iterator, TypeVar from .logging import get_logger +from .usage_analytics_store import UsageAnalyticsStore logger = get_logger(__name__) +T = TypeVar("T") class LogManager: def __init__(self, logs_dir: Path = Path("logs")): self.logs_dir = logs_dir + self._usage_store = UsageAnalyticsStore(logs_dir=logs_dir) + self._analytics_cache_ttl_seconds = 30.0 + self._analytics_cache: dict[tuple[Any, ...], tuple[float, Any]] = {} + self._analytics_cache_lock = Lock() + self._cache_miss = object() + + def _get_cached(self, key: tuple[Any, ...]) -> Any: + now = time.time() + with self._analytics_cache_lock: + cached = self._analytics_cache.get(key) + if cached is None: + return self._cache_miss + + expires_at, value = cached + if expires_at <= now: + self._analytics_cache.pop(key, None) + return self._cache_miss + + return value + + def _set_cached( + self, key: tuple[Any, ...], value: Any, ttl_seconds: float | None = None + ) -> None: + ttl = ( + self._analytics_cache_ttl_seconds + if ttl_seconds is None + else max(1.0, ttl_seconds) + ) + expires_at = time.time() + ttl + with self._analytics_cache_lock: + self._analytics_cache[key] = (expires_at, value) + + def _cache_call( + self, + key: tuple[Any, ...], + compute: Callable[[], T], + ttl_seconds: float | None = None, + ) -> T: + cached = self._get_cached(key) + if cached is not self._cache_miss: + return cached + + value = compute() + self._set_cached(key, value, ttl_seconds=ttl_seconds) + return value + + def _get_cached_entries(self, hours: int) -> list[dict[str, Any]]: + return self._cache_call( + ("usage_entries", hours), + lambda: list(self._yield_log_entries(hours_back=hours)), + ) def _yield_log_entries( self, @@ -19,7 +75,6 @@ class LogManager: specific_date: str | None = None, reverse_files: bool = False, max_files: int | None = None, - window_center: datetime | None = None, ) -> Iterator[dict[str, Any]]: """ Yields log entries from files. @@ -29,7 +84,6 @@ class LogManager: specific_date: specific date string (YYYY-MM-DD) to look at. reverse_files: if True, process files in reverse order (newest first). max_files: maximum number of log files to process (most recent if reverse_files is True). - window_center: datetime object to center a 5-month window around. """ if not self.logs_dir.exists(): return @@ -44,36 +98,6 @@ class LogManager: log_files.append(log_file) else: log_files = sorted(self.logs_dir.glob("app_*.log")) - - if window_center: - # Calculate the 5 months: [center-2, center-1, center, center+1, center+2] - allowed_month_years = [] - cur_m = window_center.month - cur_y = window_center.year - - for offset in range(-2, 3): - m = cur_m + offset - y = cur_y - while m <= 0: - m += 12 - y -= 1 - while m > 12: - m -= 12 - y += 1 - allowed_month_years.append(f"{y}-{m:02d}") - - filtered_files = [] - for log_path in log_files: - try: - # Stem is "app_YYYY-MM-DD" - file_date_str = log_path.stem.split("_")[1] - file_month_year = file_date_str[:7] # YYYY-MM - if file_month_year in allowed_month_years: - filtered_files.append(log_path) - except Exception: - continue - log_files = filtered_files - if reverse_files: log_files.reverse() @@ -303,128 +327,208 @@ class LogManager: return 0 def get_usage_summary(self, hours: int = 24) -> dict: - entries = list( - self._yield_log_entries( - hours_back=hours, window_center=datetime.now(timezone.utc) - ) + def compute() -> dict: + try: + return self._usage_store.get_summary(hours_back=hours) + except Exception as e: + logger.error( + f"Usage analytics index failed, falling back to log scan: {e}" + ) + return self._calculate_summary_stats(self._get_cached_entries(hours)) + + return self._cache_call( + ("usage_summary", hours), + compute, ) - return self._calculate_summary_stats(entries) def get_usage_metrics(self, interval: int = 15, hours: int = 24) -> dict: - entries = list( - self._yield_log_entries( - hours_back=hours, window_center=datetime.now(timezone.utc) - ) + def compute() -> dict: + try: + return self._usage_store.get_metrics( + interval_minutes=interval, + hours_back=hours, + ) + except Exception as e: + logger.error( + f"Usage analytics index failed, falling back to log scan: {e}" + ) + return self._aggregate_metrics_by_time( + self._get_cached_entries(hours), interval, hours + ) + + return self._cache_call( + ("usage_metrics", interval, hours), + compute, + ) + + def get_usage_dashboard( + self, + interval: int = 15, + hours: int = 24, + error_limit: int = 100, + model_limit: int = 20, + ) -> dict: + # Large ranges are expensive to scan; keep cached longer. + if hours <= 24: + cache_ttl = 60.0 + elif hours <= 7 * 24: + cache_ttl = 300.0 + elif hours <= 30 * 24: + cache_ttl = 1800.0 + elif hours <= 90 * 24: + cache_ttl = 7200.0 + else: + cache_ttl = 21600.0 + + def compute() -> dict: + try: + return self._usage_store.get_dashboard( + interval_minutes=interval, + hours_back=hours, + error_limit=error_limit, + model_limit=model_limit, + ) + except Exception as e: + logger.error( + f"Usage analytics index failed, falling back to log scan: {e}" + ) + return self._aggregate_dashboard( + interval_minutes=interval, + hours_back=hours, + error_limit=error_limit, + model_limit=model_limit, + ) + + return self._cache_call( + ("usage_dashboard", interval, hours, error_limit, model_limit), + compute, + ttl_seconds=cache_ttl, ) - return self._aggregate_metrics_by_time(entries, interval, hours) def get_error_details(self, hours: int = 24, limit: int = 100) -> dict: - errors: list[dict[str, Any]] = [] + def compute() -> dict: + try: + return self._usage_store.get_error_details(hours_back=hours, limit=limit) + except Exception as e: + logger.error( + f"Usage analytics index failed, falling back to log scan: {e}" + ) - for entry in self._yield_log_entries(hours_back=hours): - if str(entry.get("levelname", "")).upper() != "ERROR": - continue + errors: list[dict] = [] + for entry in self._get_cached_entries(hours): + if str(entry.get("levelname", "")).upper() == "ERROR": + timestamp_str = entry.get("asctime", "") + errors.append( + { + "timestamp": timestamp_str, + "message": entry.get("message", ""), + "error_type": entry.get("error_type", "unknown"), + "pathname": entry.get("pathname", ""), + "lineno": entry.get("lineno", 0), + "request_id": entry.get("request_id", ""), + } + ) - errors.append( - { - "timestamp": entry.get("asctime", ""), - "message": entry.get("message", ""), - "error_type": entry.get("error_type", "unknown"), - "pathname": entry.get("pathname", ""), - "lineno": entry.get("lineno", 0), - "request_id": entry.get("request_id", ""), - } - ) + errors.sort(key=lambda x: x["timestamp"], reverse=True) + return {"errors": errors[:limit], "total_count": len(errors)} - errors.sort(key=lambda x: str(x["timestamp"]), reverse=True) - return {"errors": errors[:limit], "total_count": len(errors)} + return self._cache_call(("error_details", hours, limit), compute) def get_revenue_by_model(self, hours: int = 24, limit: int = 20) -> dict: - entries = list( - self._yield_log_entries( - hours_back=hours, window_center=datetime.now(timezone.utc) - ) - ) - - model_stats: dict[str, dict[str, int | float]] = defaultdict( - lambda: { - "revenue_msats": 0, - "refunds_msats": 0, - "requests": 0, - "successful": 0, - "failed": 0, - } - ) - - for entry in entries: + def compute() -> dict: try: - model = entry.get("model", "unknown") - if not isinstance(model, str): - model = "unknown" - - message = str(entry.get("message", "")).lower() - - completed, revenue_msats, _, _ = self._extract_success_metrics( - entry, message + return self._usage_store.get_revenue_by_model( + hours_back=hours, limit=limit ) - if completed: - model_stats[model]["requests"] += 1 - model_stats[model]["successful"] += 1 - if revenue_msats > 0: - model_stats[model]["revenue_msats"] += revenue_msats - - failed = ( - "revert payment" in message or "upstream request failed" in message + except Exception as e: + logger.error( + f"Usage analytics index failed, falling back to log scan: {e}" ) - if failed: - model_stats[model]["requests"] += 1 - model_stats[model]["failed"] += 1 - if "revert payment" in message: - max_cost = entry.get("max_cost_for_model", 0) - if isinstance(max_cost, (int, float)) and max_cost > 0: - model_stats[model]["refunds_msats"] += max_cost - except Exception: - continue + entries = self._get_cached_entries(hours) - models: list[dict[str, Any]] = [] - total_revenue = 0.0 - - for model, stats in model_stats.items(): - revenue_msats = float(stats["revenue_msats"]) - refunds_msats = float(stats["refunds_msats"]) - - revenue_sats = revenue_msats / 1000 - refunds_sats = refunds_msats / 1000 - net_revenue_sats = revenue_sats - refunds_sats - - total_revenue += net_revenue_sats - - requests = int(stats["requests"]) - successful = int(stats["successful"]) - - models.append( - { - "model": model, - "revenue_sats": revenue_sats, - "refunds_sats": refunds_sats, - "net_revenue_sats": net_revenue_sats, - "requests": requests, - "successful": successful, - "failed": int(stats["failed"]), - "avg_revenue_per_request": ( - revenue_sats / successful if successful > 0 else 0 - ), + model_stats: dict[str, dict[str, int | float]] = defaultdict( + lambda: { + "revenue_msats": 0, + "refunds_msats": 0, + "requests": 0, + "successful": 0, + "failed": 0, } ) - models.sort(key=lambda x: float(x["net_revenue_sats"]), reverse=True) + for entry in entries: + try: + model = entry.get("model", "unknown") + if not isinstance(model, str): + model = "unknown" - return { - "models": models[:limit], - "total_revenue_sats": total_revenue, - "total_models": len(models), - } + message = str(entry.get("message", "")).lower() + + completed, revenue_msats, _, _ = self._extract_success_metrics( + entry, message + ) + if completed: + model_stats[model]["requests"] += 1 + model_stats[model]["successful"] += 1 + if revenue_msats > 0: + model_stats[model]["revenue_msats"] += revenue_msats + + failed = ( + "revert payment" in message + or "upstream request failed" in message + ) + if failed: + model_stats[model]["requests"] += 1 + model_stats[model]["failed"] += 1 + if "revert payment" in message: + max_cost = entry.get("max_cost_for_model", 0) + if isinstance(max_cost, (int, float)) and max_cost > 0: + model_stats[model]["refunds_msats"] += max_cost + + except Exception: + continue + + models: list[dict[str, Any]] = [] + total_revenue = 0.0 + + for model, stats in model_stats.items(): + revenue_msats = float(stats["revenue_msats"]) + refunds_msats = float(stats["refunds_msats"]) + + revenue_sats = revenue_msats / 1000 + refunds_sats = refunds_msats / 1000 + net_revenue_sats = revenue_sats - refunds_sats + + total_revenue += net_revenue_sats + + requests = int(stats["requests"]) + successful = int(stats["successful"]) + + models.append( + { + "model": model, + "revenue_sats": revenue_sats, + "refunds_sats": refunds_sats, + "net_revenue_sats": net_revenue_sats, + "requests": requests, + "successful": successful, + "failed": int(stats["failed"]), + "avg_revenue_per_request": ( + revenue_sats / successful if successful > 0 else 0 + ), + } + ) + + models.sort(key=lambda x: float(x["net_revenue_sats"]), reverse=True) + + return { + "models": models[:limit], + "total_revenue_sats": total_revenue, + "total_models": len(models), + } + + return self._cache_call(("revenue_by_model", hours, limit), compute) def _build_summary_response(self, stats: dict[str, Any]) -> dict[str, Any]: revenue_sats = stats["revenue_msats"] / 1000 @@ -555,6 +659,344 @@ class LogManager: return self._build_summary_response(stats) + def _aggregate_dashboard( + self, + interval_minutes: int, + hours_back: int, + error_limit: int, + model_limit: int, + ) -> dict[str, Any]: + time_buckets: dict[str, dict[str, Any]] = defaultdict( + lambda: { + "total_requests": 0, + "successful_chat_completions": 0, + "failed_requests": 0, + "errors": 0, + "warnings": 0, + "payment_processed": 0, + "upstream_errors": 0, + "revenue_msats": 0.0, + "refunds_msats": 0.0, + "input_tokens": 0, + "output_tokens": 0, + "total_tokens": 0, + } + ) + summary_stats: dict[str, Any] = { + "total_entries": 0, + "total_requests": 0, + "successful_chat_completions": 0, + "failed_requests": 0, + "total_errors": 0, + "total_warnings": 0, + "payment_processed": 0, + "upstream_errors": 0, + "unique_models": set(), + "error_types": defaultdict(int), + "revenue_msats": 0.0, + "refunds_msats": 0.0, + "input_tokens": 0, + "output_tokens": 0, + "total_tokens": 0, + } + model_stats: dict[str, dict[str, int | float]] = defaultdict( + lambda: { + "revenue_msats": 0, + "refunds_msats": 0, + "requests": 0, + "successful": 0, + "failed": 0, + } + ) + model_mix_buckets: dict[str, dict[str, int]] = defaultdict( + lambda: defaultdict(int) + ) + model_mix_revenue_buckets: dict[str, dict[str, float]] = defaultdict( + lambda: defaultdict(float) + ) + model_mix_token_buckets: dict[str, dict[str, int]] = defaultdict( + lambda: defaultdict(int) + ) + model_mix_totals: dict[str, int] = defaultdict(int) + model_mix_revenue_totals: dict[str, float] = defaultdict(float) + model_mix_token_totals: dict[str, int] = defaultdict(int) + latest_errors_heap: list[tuple[str, dict[str, Any]]] = [] + total_error_count = 0 + + for entry in self._yield_log_entries(hours_back=hours_back): + try: + summary_stats["total_entries"] += 1 + + timestamp_str = entry.get("asctime", "") + message = str(entry.get("message", "")).lower() + level = str(entry.get("levelname", "")).upper() + model = entry.get("model", "unknown") + if not isinstance(model, str): + model = "unknown" + + bucket_key = ( + self._bucket_key_for_timestamp(timestamp_str, interval_minutes) + if isinstance(timestamp_str, str) + else None + ) + bucket = time_buckets[bucket_key] if bucket_key else None + + if level == "ERROR": + summary_stats["total_errors"] += 1 + if bucket: + bucket["errors"] += 1 + if "error_type" in entry: + summary_stats["error_types"][str(entry["error_type"])] += 1 + + total_error_count += 1 + error_item = { + "timestamp": timestamp_str, + "message": entry.get("message", ""), + "error_type": entry.get("error_type", "unknown"), + "pathname": entry.get("pathname", ""), + "lineno": entry.get("lineno", 0), + "request_id": entry.get("request_id", ""), + } + if len(latest_errors_heap) < error_limit: + heappush(latest_errors_heap, (timestamp_str, error_item)) + elif timestamp_str > latest_errors_heap[0][0]: + heapreplace(latest_errors_heap, (timestamp_str, error_item)) + elif level == "WARNING": + summary_stats["total_warnings"] += 1 + if bucket: + bucket["warnings"] += 1 + + completed, revenue_msats, input_tokens, output_tokens = ( + self._extract_success_metrics(entry, message) + ) + if completed: + summary_stats["total_requests"] += 1 + summary_stats["successful_chat_completions"] += 1 + summary_stats["input_tokens"] += input_tokens + summary_stats["output_tokens"] += output_tokens + summary_stats["total_tokens"] += input_tokens + output_tokens + model_stats[model]["requests"] += 1 + model_stats[model]["successful"] += 1 + model_mix_totals[model] += 1 + if bucket: + bucket["total_requests"] += 1 + bucket["successful_chat_completions"] += 1 + bucket["input_tokens"] += input_tokens + bucket["output_tokens"] += output_tokens + bucket["total_tokens"] += input_tokens + output_tokens + if bucket_key: + model_mix_buckets[bucket_key][model] += 1 + if revenue_msats > 0: + model_mix_revenue_buckets[bucket_key][model] += revenue_msats + model_mix_revenue_totals[model] += revenue_msats + if input_tokens > 0 or output_tokens > 0: + token_total = input_tokens + output_tokens + model_mix_token_buckets[bucket_key][model] += token_total + model_mix_token_totals[model] += token_total + + if revenue_msats > 0: + summary_stats["revenue_msats"] += revenue_msats + model_stats[model]["revenue_msats"] += revenue_msats + if bucket: + bucket["revenue_msats"] += revenue_msats + + failed = ( + "upstream request failed" in message + or "revert payment" in message + ) + if failed: + summary_stats["total_requests"] += 1 + summary_stats["failed_requests"] += 1 + model_stats[model]["requests"] += 1 + model_stats[model]["failed"] += 1 + if bucket: + bucket["total_requests"] += 1 + bucket["failed_requests"] += 1 + + if "payment processed successfully" in message: + summary_stats["payment_processed"] += 1 + if bucket: + bucket["payment_processed"] += 1 + + if "upstream" in message and level == "ERROR": + summary_stats["upstream_errors"] += 1 + if bucket: + bucket["upstream_errors"] += 1 + + if model != "unknown": + summary_stats["unique_models"].add(model) + + if "revert payment" in message: + max_cost = entry.get("max_cost_for_model", 0) + if isinstance(max_cost, (int, float)) and max_cost > 0: + max_cost_float = float(max_cost) + summary_stats["refunds_msats"] += max_cost_float + model_stats[model]["refunds_msats"] += max_cost_float + if bucket: + bucket["refunds_msats"] += max_cost_float + except Exception: + continue + + metrics_result = [] + for bucket_key in sorted(time_buckets.keys()): + bucket = dict(time_buckets[bucket_key]) + bucket["requests"] = bucket["total_requests"] + metrics_result.append({"timestamp": bucket_key, **bucket}) + + models: list[dict[str, Any]] = [] + total_revenue = 0.0 + for model_name, stats in model_stats.items(): + revenue_msats = float(stats["revenue_msats"]) + refunds_msats = float(stats["refunds_msats"]) + revenue_sats = revenue_msats / 1000 + refunds_sats = refunds_msats / 1000 + net_revenue_sats = revenue_sats - refunds_sats + total_revenue += net_revenue_sats + + successful = int(stats["successful"]) + models.append( + { + "model": model_name, + "revenue_sats": revenue_sats, + "refunds_sats": refunds_sats, + "net_revenue_sats": net_revenue_sats, + "requests": int(stats["requests"]), + "successful": successful, + "failed": int(stats["failed"]), + "avg_revenue_per_request": ( + revenue_sats / successful if successful > 0 else 0 + ), + } + ) + + models.sort(key=lambda x: float(x["net_revenue_sats"]), reverse=True) + latest_errors = [ + item + for _, item in sorted( + latest_errors_heap, key=lambda x: x[0], reverse=True + ) + ] + top_model_limit = max(1, min(model_limit, 20)) + top_models_requests = [ + model_name + for model_name, _ in sorted( + ( + (name, count) + for name, count in model_mix_totals.items() + if name != "unknown" + ), + key=lambda item: item[1], + reverse=True, + )[:top_model_limit] + ] + top_models_revenue = [ + model_name + for model_name, _ in sorted( + ( + (name, amount) + for name, amount in model_mix_revenue_totals.items() + if name != "unknown" + ), + key=lambda item: item[1], + reverse=True, + )[:top_model_limit] + ] + top_models_tokens = [ + model_name + for model_name, _ in sorted( + ( + (name, token_count) + for name, token_count in model_mix_token_totals.items() + if name != "unknown" + ), + key=lambda item: item[1], + reverse=True, + )[:top_model_limit] + ] + selected_models: list[str] = [] + for model in top_models_requests + top_models_revenue + top_models_tokens: + if model not in selected_models: + selected_models.append(model) + top_model_set = set(selected_models) + + model_usage_mix_metrics: list[dict[str, Any]] = [] + mix_bucket_keys = sorted( + set(model_mix_buckets.keys()) + | set(model_mix_revenue_buckets.keys()) + | set(model_mix_token_buckets.keys()) + ) + for bucket_key in mix_bucket_keys: + counts = model_mix_buckets.get(bucket_key, {}) + revenue_counts = model_mix_revenue_buckets.get(bucket_key, {}) + token_counts = model_mix_token_buckets.get(bucket_key, {}) + others = 0 + others_revenue_msats = 0.0 + others_tokens = 0 + model_counts: dict[str, int] = {} + model_revenue_msats: dict[str, float] = {} + model_tokens: dict[str, int] = {} + for model_name, successful_count in counts.items(): + if model_name in top_model_set: + model_counts[model_name] = int(successful_count) + else: + others += int(successful_count) + for model_name, revenue_value in revenue_counts.items(): + if model_name in top_model_set: + model_revenue_msats[model_name] = float(revenue_value) + else: + others_revenue_msats += float(revenue_value) + for model_name, token_value in token_counts.items(): + if model_name in top_model_set: + model_tokens[model_name] = int(token_value) + else: + others_tokens += int(token_value) + + model_usage_mix_metrics.append( + { + "timestamp": bucket_key, + "total_successful": int(sum(counts.values())), + "total_revenue_msats": float(sum(revenue_counts.values())), + "total_tokens": int(sum(token_counts.values())), + "others": others, + "others_revenue_msats": others_revenue_msats, + "others_tokens": others_tokens, + "model_counts": model_counts, + "model_revenue_msats": model_revenue_msats, + "model_tokens": model_tokens, + } + ) + + return { + "metrics": { + "metrics": metrics_result, + "interval_minutes": interval_minutes, + "hours_back": hours_back, + "total_buckets": len(metrics_result), + }, + "summary": self._build_summary_response(summary_stats), + "error_details": { + "errors": latest_errors, + "total_count": total_error_count, + }, + "revenue_by_model": { + "models": models[:model_limit], + "total_revenue_sats": total_revenue, + "total_models": len(models), + }, + "model_usage_mix": { + "top_models": top_models_requests, + "top_models_by_metric": { + "requests": top_models_requests, + "revenue": top_models_revenue, + "tokens": top_models_tokens, + }, + "metrics": model_usage_mix_metrics, + "interval_minutes": interval_minutes, + "hours_back": hours_back, + "total_buckets": len(model_usage_mix_metrics), + }, + } + def _aggregate_metrics_by_time( self, entries: list[dict], interval_minutes: int, hours_back: int ) -> dict: diff --git a/routstr/core/logging.py b/routstr/core/logging.py index 00474949..14b1b1ff 100644 --- a/routstr/core/logging.py +++ b/routstr/core/logging.py @@ -3,36 +3,38 @@ Logging configuration for Routstr. CRITICAL LOG MESSAGES FOR USAGE STATISTICS: =========================================== -The following log messages are parsed by the usage tracking system (routstr/core/admin.py). +The following log messages are parsed by the usage tracking system +(routstr/core/usage_analytics_store.py and routstr/core/log_manager.py). DO NOT modify or remove these messages without updating the usage tracking logic: 1. "Received proxy request" (INFO) - routstr/proxy.py - Used to count total incoming requests - Includes model information in context - 2. "Payment adjustment completed for streaming" (INFO) - routstr/upstream/base.py - "Payment adjustment completed for non-streaming" (INFO) - routstr/upstream/base.py +2. "Calculated token-based cost" (INFO) - routstr/auth.py - Used to track successful completions and revenue - - The 'cost_data.total_msats' field is extracted for revenue calculation - - Must include 'cost_data' in extra dict + - The 'token_cost', 'model', 'input_tokens', and 'output_tokens' fields are extracted for dashboard metrics -3. "Payment processed successfully" (INFO) - routstr/auth.py +3. "Max cost payment finalized" (INFO) - routstr/auth.py + - Used as the successful completion fallback when token usage is unavailable + - The 'charged_amount', 'model', 'input_tokens', and 'output_tokens' fields are extracted for dashboard metrics + +4. "Payment processed successfully" (INFO) - routstr/auth.py - Used to count successful payment processing events - Tracks payment-related metrics -4. "Upstream request failed, revert payment" (WARNING) - routstr/proxy.py +5. "Upstream request failed, revert payment" (WARNING) - routstr/proxy.py - Used to track failed requests and refunds - The 'max_cost_for_model' field is extracted for refund calculation - Must include 'max_cost_for_model' in extra dict -5. Any ERROR level logs with "upstream" in the message +6. Any ERROR level logs with "upstream" in the message - Used to count upstream provider errors - Helps identify service reliability issues If you need to modify these messages, ensure you also update the parsing logic in: -- routstr/core/admin.py:_aggregate_metrics_by_time() -- routstr/core/admin.py:_get_summary_stats() -- routstr/core/admin.py:get_revenue_by_model() +- routstr/core/usage_analytics_store.py +- routstr/core/log_manager.py """ import logging.config diff --git a/routstr/core/main.py b/routstr/core/main.py index 27f46280..d0365dad 100644 --- a/routstr/core/main.py +++ b/routstr/core/main.py @@ -12,7 +12,11 @@ from starlette.exceptions import HTTPException from ..auth import periodic_key_reset from ..balance import balance_router, deprecated_wallet_router -from ..nostr import announce_provider, providers_cache_refresher +from ..nostr import ( + announce_provider, + providers_cache_refresher, + publish_usage_analytics, +) from ..nostr.discovery import providers_router from ..payment.models import models_router, update_sats_pricing from ..payment.price import update_prices_periodically @@ -45,6 +49,7 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: pricing_task = None payout_task = None nip91_task = None + analytics_task = None providers_task = None models_refresh_task = None model_maps_refresh_task = None @@ -104,6 +109,7 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: payout_task = asyncio.create_task(periodic_payout()) if global_settings.nsec: nip91_task = asyncio.create_task(announce_provider()) + analytics_task = asyncio.create_task(publish_usage_analytics()) if global_settings.providers_refresh_interval_seconds > 0: providers_task = asyncio.create_task(providers_cache_refresher()) key_reset_task = asyncio.create_task(periodic_key_reset()) @@ -132,6 +138,8 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: payout_task.cancel() if nip91_task is not None: nip91_task.cancel() + if analytics_task is not None: + analytics_task.cancel() if providers_task is not None: providers_task.cancel() if models_refresh_task is not None: @@ -155,6 +163,8 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: tasks_to_wait.append(payout_task) if nip91_task is not None: tasks_to_wait.append(nip91_task) + if analytics_task is not None: + tasks_to_wait.append(analytics_task) if providers_task is not None: tasks_to_wait.append(providers_task) if models_refresh_task is not None: diff --git a/routstr/core/settings.py b/routstr/core/settings.py index fba0383c..027cc0ba 100644 --- a/routstr/core/settings.py +++ b/routstr/core/settings.py @@ -93,6 +93,20 @@ class Settings(BaseSettings): # Discovery relays: list[str] = Field(default_factory=list, env="RELAYS") + enable_analytics_sharing: bool = Field( + default=True, env="ENABLE_ANALYTICS_SHARING" + ) + +def _normalize_settings_data(data: dict[str, Any]) -> dict[str, Any]: + """Discard unknown keys from persisted settings.""" + normalized: dict[str, Any] = {} + known_fields = Settings.__fields__ + + for key, value in data.items(): + if key in known_fields: + normalized[key] = value + + return normalized def _compute_primary_mint(cashu_mints: list[str]) -> str: @@ -232,17 +246,21 @@ class SettingsService: db_id, db_data, _updated_at = row try: - db_json = ( + db_json_raw = ( json.loads(db_data) if isinstance(db_data, str) else dict(db_data) ) + if not isinstance(db_json_raw, dict): + db_json_raw = {} except Exception: - db_json = {} + db_json_raw = {} + db_json = _normalize_settings_data(db_json_raw) valid_fields = set(env_resolved.dict().keys()) merged_dict: dict[str, Any] = dict(env_resolved.dict()) merged_dict.update( {k: v for k, v in db_json.items() if v not in (None, "", [], {}) and k in valid_fields} ) + merged_dict = Settings(**merged_dict).dict() # Ensure primary_mint is consistent with cashu_mints if not explicitly set if not merged_dict.get("primary_mint"): @@ -250,7 +268,7 @@ class SettingsService: merged_dict.get("cashu_mints", []) ) - if any(k not in db_json for k in merged_dict.keys()): + if db_json_raw != merged_dict: await db_session.exec( # type: ignore text( "UPDATE settings SET data = :data, updated_at = :updated_at WHERE id = 1" @@ -273,7 +291,7 @@ class SettingsService: ) -> Settings: async with cls._lock: current = cls.get() - candidate_dict = {**current.dict(), **partial} + candidate_dict = {**current.dict(), **_normalize_settings_data(partial)} candidate = Settings(**candidate_dict) from sqlmodel import text diff --git a/routstr/core/usage_analytics_store.py b/routstr/core/usage_analytics_store.py new file mode 100644 index 00000000..7ba90e24 --- /dev/null +++ b/routstr/core/usage_analytics_store.py @@ -0,0 +1,1380 @@ +import json +import sqlite3 +import time +from collections import defaultdict +from datetime import datetime, timedelta, timezone +from pathlib import Path +from threading import Lock +from typing import Any + +from .logging import get_logger + +logger = get_logger(__name__) + + +class UsageAnalyticsStore: + """ + Incremental usage analytics index backed by SQLite. + + Instead of rescanning raw JSON log files for every dashboard request, we keep + a rolling minute-level aggregate that is updated from only newly appended log + bytes. + """ + + SCHEMA_VERSION = "4" + + def __init__(self, logs_dir: Path, db_path: Path | None = None): + self.logs_dir = logs_dir + self.db_path = db_path or (logs_dir / "usage_analytics.db") + self._lock = Lock() + self._conn: sqlite3.Connection | None = None + + def get_dashboard( + self, + *, + interval_minutes: int, + hours_back: int, + error_limit: int, + model_limit: int, + ) -> dict[str, Any]: + with self._lock: + conn = self._get_connection_locked() + self._ensure_up_to_date_locked(conn) + cutoff_timestamp = self._cutoff_timestamp(hours_back) + + summary = self._query_summary_locked(conn, cutoff_timestamp) + metrics = self._query_metrics_locked( + conn, + cutoff_timestamp=cutoff_timestamp, + interval_minutes=interval_minutes, + hours_back=hours_back, + ) + error_details = self._query_error_details_locked( + conn, + cutoff_timestamp=cutoff_timestamp, + limit=error_limit, + total_error_count=summary["total_errors"], + ) + revenue_by_model = self._query_revenue_by_model_locked( + conn, + cutoff_timestamp=cutoff_timestamp, + limit=model_limit, + ) + model_usage_mix = self._query_model_usage_mix_locked( + conn, + cutoff_timestamp=cutoff_timestamp, + interval_minutes=interval_minutes, + hours_back=hours_back, + limit=model_limit, + ) + + return { + "metrics": metrics, + "summary": summary, + "error_details": error_details, + "revenue_by_model": revenue_by_model, + "model_usage_mix": model_usage_mix, + } + + def get_summary(self, *, hours_back: int) -> dict[str, Any]: + with self._lock: + conn = self._get_connection_locked() + self._ensure_up_to_date_locked(conn) + cutoff_timestamp = self._cutoff_timestamp(hours_back) + return self._query_summary_locked(conn, cutoff_timestamp) + + def get_metrics( + self, + *, + interval_minutes: int, + hours_back: int, + ) -> dict[str, Any]: + with self._lock: + conn = self._get_connection_locked() + self._ensure_up_to_date_locked(conn) + cutoff_timestamp = self._cutoff_timestamp(hours_back) + return self._query_metrics_locked( + conn, + cutoff_timestamp=cutoff_timestamp, + interval_minutes=interval_minutes, + hours_back=hours_back, + ) + + def get_error_details(self, *, hours_back: int, limit: int) -> dict[str, Any]: + with self._lock: + conn = self._get_connection_locked() + self._ensure_up_to_date_locked(conn) + cutoff_timestamp = self._cutoff_timestamp(hours_back) + return self._query_error_details_locked( + conn, + cutoff_timestamp=cutoff_timestamp, + limit=limit, + ) + + def get_revenue_by_model(self, *, hours_back: int, limit: int) -> dict[str, Any]: + with self._lock: + conn = self._get_connection_locked() + self._ensure_up_to_date_locked(conn) + cutoff_timestamp = self._cutoff_timestamp(hours_back) + return self._query_revenue_by_model_locked( + conn, + cutoff_timestamp=cutoff_timestamp, + limit=limit, + ) + + def _get_connection_locked(self) -> sqlite3.Connection: + if self._conn is not None: + return self._conn + + self.db_path.parent.mkdir(parents=True, exist_ok=True) + conn = sqlite3.connect( + self.db_path, + timeout=30.0, + check_same_thread=False, + ) + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA journal_mode=WAL") + conn.execute("PRAGMA synchronous=NORMAL") + conn.execute("PRAGMA temp_store=MEMORY") + conn.execute("PRAGMA cache_size=-20000") + self._initialize_schema_locked(conn) + self._conn = conn + return conn + + def _initialize_schema_locked(self, conn: sqlite3.Connection) -> None: + conn.execute( + """ + CREATE TABLE IF NOT EXISTS analytics_meta ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL + ) + """ + ) + + current_version_row = conn.execute( + "SELECT value FROM analytics_meta WHERE key = 'schema_version'" + ).fetchone() + current_version = current_version_row[0] if current_version_row else None + + conn.execute( + """ + CREATE TABLE IF NOT EXISTS analytics_file_state ( + path TEXT PRIMARY KEY, + inode INTEGER NOT NULL, + offset INTEGER NOT NULL, + size INTEGER NOT NULL, + updated_at REAL NOT NULL + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS analytics_minute ( + minute_ts TEXT PRIMARY KEY, + total_entries INTEGER NOT NULL DEFAULT 0, + total_requests INTEGER NOT NULL DEFAULT 0, + successful_chat_completions INTEGER NOT NULL DEFAULT 0, + failed_requests INTEGER NOT NULL DEFAULT 0, + errors INTEGER NOT NULL DEFAULT 0, + warnings INTEGER NOT NULL DEFAULT 0, + payment_processed INTEGER NOT NULL DEFAULT 0, + upstream_errors INTEGER NOT NULL DEFAULT 0, + revenue_msats REAL NOT NULL DEFAULT 0, + refunds_msats REAL NOT NULL DEFAULT 0, + input_tokens INTEGER NOT NULL DEFAULT 0, + output_tokens INTEGER NOT NULL DEFAULT 0, + total_tokens INTEGER NOT NULL DEFAULT 0 + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS analytics_model_minute ( + minute_ts TEXT NOT NULL, + model TEXT NOT NULL, + requests INTEGER NOT NULL DEFAULT 0, + successful INTEGER NOT NULL DEFAULT 0, + failed INTEGER NOT NULL DEFAULT 0, + revenue_msats REAL NOT NULL DEFAULT 0, + refunds_msats REAL NOT NULL DEFAULT 0, + input_tokens INTEGER NOT NULL DEFAULT 0, + output_tokens INTEGER NOT NULL DEFAULT 0, + total_tokens INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (minute_ts, model) + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS analytics_model_presence_minute ( + minute_ts TEXT NOT NULL, + model TEXT NOT NULL, + count INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (minute_ts, model) + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS analytics_error_type_minute ( + minute_ts TEXT NOT NULL, + error_type TEXT NOT NULL, + count INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (minute_ts, error_type) + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS analytics_error_events ( + timestamp TEXT NOT NULL, + message TEXT NOT NULL, + error_type TEXT NOT NULL, + pathname TEXT NOT NULL, + lineno INTEGER NOT NULL, + request_id TEXT NOT NULL + ) + """ + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_analytics_model_minute_ts ON analytics_model_minute (minute_ts)" + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_analytics_model_minute_model_ts ON analytics_model_minute (model, minute_ts)" + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_analytics_model_presence_ts ON analytics_model_presence_minute (minute_ts)" + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_analytics_error_type_minute_ts ON analytics_error_type_minute (minute_ts)" + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_analytics_error_events_ts ON analytics_error_events (timestamp DESC)" + ) + self._migrate_schema_locked(conn) + if current_version != self.SCHEMA_VERSION: + conn.execute( + """ + INSERT OR REPLACE INTO analytics_meta (key, value) + VALUES ('schema_version', ?) + """, + (self.SCHEMA_VERSION,), + ) + conn.commit() + + def _migrate_schema_locked(self, conn: sqlite3.Connection) -> None: + self._ensure_column_locked( + conn, + "analytics_minute", + "input_tokens", + "INTEGER NOT NULL DEFAULT 0", + ) + self._ensure_column_locked( + conn, + "analytics_minute", + "output_tokens", + "INTEGER NOT NULL DEFAULT 0", + ) + self._ensure_column_locked( + conn, + "analytics_minute", + "total_tokens", + "INTEGER NOT NULL DEFAULT 0", + ) + self._ensure_column_locked( + conn, + "analytics_model_minute", + "input_tokens", + "INTEGER NOT NULL DEFAULT 0", + ) + self._ensure_column_locked( + conn, + "analytics_model_minute", + "output_tokens", + "INTEGER NOT NULL DEFAULT 0", + ) + self._ensure_column_locked( + conn, + "analytics_model_minute", + "total_tokens", + "INTEGER NOT NULL DEFAULT 0", + ) + + def _ensure_column_locked( + self, + conn: sqlite3.Connection, + table: str, + column: str, + column_definition: str, + ) -> None: + existing_columns = { + str(row["name"]) + for row in conn.execute(f"PRAGMA table_info({table})").fetchall() + } + if column in existing_columns: + return + + conn.execute( + f"ALTER TABLE {table} ADD COLUMN {column} {column_definition}" + ) + logger.info(f"Migrated analytics schema: added {table}.{column}") + + def _drop_index_tables_locked(self, conn: sqlite3.Connection) -> None: + conn.execute("DROP TABLE IF EXISTS analytics_file_state") + conn.execute("DROP TABLE IF EXISTS analytics_minute") + conn.execute("DROP TABLE IF EXISTS analytics_model_minute") + conn.execute("DROP TABLE IF EXISTS analytics_model_presence_minute") + conn.execute("DROP TABLE IF EXISTS analytics_error_type_minute") + conn.execute("DROP TABLE IF EXISTS analytics_error_events") + + def _ensure_up_to_date_locked(self, conn: sqlite3.Connection) -> None: + if not self.logs_dir.exists(): + return + + log_files = sorted(self.logs_dir.glob("app_*.log")) + if not log_files: + return + + requires_rebuild = False + + for log_file in log_files: + try: + self._process_log_file_locked(conn, log_file) + except RuntimeError: + requires_rebuild = True + break + except Exception as exc: + logger.error(f"Failed indexing usage analytics for {log_file}: {exc}") + continue + + if requires_rebuild: + logger.warning( + "Usage analytics index out-of-sync, rebuilding from all log files" + ) + self._rebuild_locked(conn, log_files) + return + + # Commit even when only file-state metadata changed + # (for example when we intentionally keep offset at the last full line). + conn.commit() + + def _rebuild_locked( + self, conn: sqlite3.Connection, log_files: list[Path] | None = None + ) -> None: + self._drop_index_tables_locked(conn) + self._initialize_schema_locked(conn) + + files = log_files if log_files is not None else sorted(self.logs_dir.glob("app_*.log")) + for log_file in files: + try: + self._process_log_file_locked(conn, log_file, force_full_read=True) + except Exception as exc: + logger.error(f"Failed rebuilding usage analytics for {log_file}: {exc}") + conn.commit() + + def _process_log_file_locked( + self, + conn: sqlite3.Connection, + log_file: Path, + force_full_read: bool = False, + ) -> bool: + stat = log_file.stat() + inode = int(getattr(stat, "st_ino", 0)) + file_size = int(stat.st_size) + log_file_path = str(log_file.resolve()) + + previous_offset = 0 + if not force_full_read: + row = conn.execute( + """ + SELECT inode, offset + FROM analytics_file_state + WHERE path = ? + """, + (log_file_path,), + ).fetchone() + if row is not None: + previous_inode = int(row["inode"]) + previous_offset = int(row["offset"]) + if previous_inode and inode and previous_inode != inode: + raise RuntimeError("inode changed") + if previous_offset > file_size: + raise RuntimeError("file shrunk") + + if previous_offset >= file_size and not force_full_read: + self._upsert_file_state_locked( + conn, + path=log_file_path, + inode=inode, + offset=file_size, + size=file_size, + ) + return False + + ( + end_offset, + minute_updates, + model_updates, + model_presence_updates, + error_type_updates, + error_events, + ) = self._collect_updates_from_file(log_file, previous_offset) + + self._apply_updates_locked( + conn=conn, + minute_updates=minute_updates, + model_updates=model_updates, + model_presence_updates=model_presence_updates, + error_type_updates=error_type_updates, + error_events=error_events, + ) + + latest_size = int(log_file.stat().st_size) + self._upsert_file_state_locked( + conn, + path=log_file_path, + inode=inode, + offset=end_offset, + size=latest_size, + ) + return end_offset != previous_offset + + def _upsert_file_state_locked( + self, + conn: sqlite3.Connection, + *, + path: str, + inode: int, + offset: int, + size: int, + ) -> None: + conn.execute( + """ + INSERT INTO analytics_file_state (path, inode, offset, size, updated_at) + VALUES (?, ?, ?, ?, ?) + ON CONFLICT(path) DO UPDATE SET + inode = excluded.inode, + offset = excluded.offset, + size = excluded.size, + updated_at = excluded.updated_at + """, + (path, inode, offset, size, time.time()), + ) + + def _collect_updates_from_file( + self, log_file: Path, start_offset: int + ) -> tuple[ + int, + dict[str, dict[str, float]], + dict[tuple[str, str], dict[str, float]], + dict[tuple[str, str], int], + dict[tuple[str, str], int], + list[tuple[str, str, str, str, int, str]], + ]: + minute_updates: dict[str, dict[str, float]] = defaultdict( + self._new_minute_stats + ) + model_updates: dict[tuple[str, str], dict[str, float]] = defaultdict( + self._new_model_stats + ) + model_presence_updates: dict[tuple[str, str], int] = defaultdict(int) + error_type_updates: dict[tuple[str, str], int] = defaultdict(int) + error_events: list[tuple[str, str, str, str, int, str]] = [] + + end_offset = start_offset + with open(log_file, "rb") as f: + f.seek(start_offset) + + while True: + line_start = f.tell() + raw_line = f.readline() + if not raw_line: + break + + # If the writer is appending and we catch a partial line at EOF, + # do not advance beyond it. We'll parse it on the next refresh. + if not raw_line.endswith(b"\n"): + f.seek(line_start) + break + + end_offset = f.tell() + if not raw_line.strip(): + continue + + try: + entry = json.loads(raw_line) + except Exception: + continue + + if not isinstance(entry, dict): + continue + + minute_key = self._minute_key(entry.get("asctime")) + if minute_key is None: + continue + + bucket = minute_updates[minute_key] + bucket["total_entries"] += 1 + + message_value = entry.get("message", "") + message = str(message_value).lower() + level = str(entry.get("levelname", "")).upper() + + model_raw = entry.get("model", "unknown") + model = model_raw if isinstance(model_raw, str) else "unknown" + + if level == "ERROR": + bucket["errors"] += 1 + error_type = str(entry.get("error_type", "unknown")) + error_type_updates[(minute_key, error_type)] += 1 + + lineno_value = entry.get("lineno", 0) + try: + lineno = int(lineno_value) + except (TypeError, ValueError): + lineno = 0 + + error_events.append( + ( + str(entry.get("asctime", "")), + str(message_value), + error_type, + str(entry.get("pathname", "")), + lineno, + str(entry.get("request_id", "")), + ) + ) + elif level == "WARNING": + bucket["warnings"] += 1 + + completed, revenue_msats, input_tokens, output_tokens = ( + self._extract_success_metrics(entry, message) + ) + if completed: + bucket["total_requests"] += 1 + bucket["successful_chat_completions"] += 1 + model_bucket = model_updates[(minute_key, model)] + model_bucket["requests"] += 1 + model_bucket["successful"] += 1 + bucket["input_tokens"] += input_tokens + bucket["output_tokens"] += output_tokens + bucket["total_tokens"] += input_tokens + output_tokens + model_bucket["input_tokens"] += input_tokens + model_bucket["output_tokens"] += output_tokens + model_bucket["total_tokens"] += input_tokens + output_tokens + + if revenue_msats > 0: + bucket["revenue_msats"] += revenue_msats + model_bucket["revenue_msats"] += revenue_msats + + failed = ( + "upstream request failed" in message + or "revert payment" in message + ) + if failed: + bucket["total_requests"] += 1 + bucket["failed_requests"] += 1 + model_bucket = model_updates[(minute_key, model)] + model_bucket["requests"] += 1 + model_bucket["failed"] += 1 + + if "payment processed successfully" in message: + bucket["payment_processed"] += 1 + + if level == "ERROR" and "upstream" in message: + bucket["upstream_errors"] += 1 + + if model != "unknown": + model_presence_updates[(minute_key, model)] += 1 + + if "revert payment" in message: + max_cost = entry.get("max_cost_for_model", 0) + if isinstance(max_cost, (int, float)) and max_cost > 0: + max_cost_float = float(max_cost) + bucket["refunds_msats"] += max_cost_float + model_updates[(minute_key, model)][ + "refunds_msats" + ] += max_cost_float + + return ( + end_offset, + minute_updates, + model_updates, + model_presence_updates, + error_type_updates, + error_events, + ) + + def _apply_updates_locked( + self, + *, + conn: sqlite3.Connection, + minute_updates: dict[str, dict[str, float]], + model_updates: dict[tuple[str, str], dict[str, float]], + model_presence_updates: dict[tuple[str, str], int], + error_type_updates: dict[tuple[str, str], int], + error_events: list[tuple[str, str, str, str, int, str]], + ) -> None: + if minute_updates: + rows = [ + ( + minute_ts, + int(stats["total_entries"]), + int(stats["total_requests"]), + int(stats["successful_chat_completions"]), + int(stats["failed_requests"]), + int(stats["errors"]), + int(stats["warnings"]), + int(stats["payment_processed"]), + int(stats["upstream_errors"]), + float(stats["revenue_msats"]), + float(stats["refunds_msats"]), + int(stats["input_tokens"]), + int(stats["output_tokens"]), + int(stats["total_tokens"]), + ) + for minute_ts, stats in minute_updates.items() + ] + conn.executemany( + """ + INSERT INTO analytics_minute ( + minute_ts, + total_entries, + total_requests, + successful_chat_completions, + failed_requests, + errors, + warnings, + payment_processed, + upstream_errors, + revenue_msats, + refunds_msats, + input_tokens, + output_tokens, + total_tokens + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(minute_ts) DO UPDATE SET + total_entries = total_entries + excluded.total_entries, + total_requests = total_requests + excluded.total_requests, + successful_chat_completions = successful_chat_completions + excluded.successful_chat_completions, + failed_requests = failed_requests + excluded.failed_requests, + errors = errors + excluded.errors, + warnings = warnings + excluded.warnings, + payment_processed = payment_processed + excluded.payment_processed, + upstream_errors = upstream_errors + excluded.upstream_errors, + revenue_msats = revenue_msats + excluded.revenue_msats, + refunds_msats = refunds_msats + excluded.refunds_msats, + input_tokens = input_tokens + excluded.input_tokens, + output_tokens = output_tokens + excluded.output_tokens, + total_tokens = total_tokens + excluded.total_tokens + """, + rows, + ) + + if model_updates: + model_rows = [ + ( + minute_ts, + model, + int(stats["requests"]), + int(stats["successful"]), + int(stats["failed"]), + float(stats["revenue_msats"]), + float(stats["refunds_msats"]), + int(stats["input_tokens"]), + int(stats["output_tokens"]), + int(stats["total_tokens"]), + ) + for (minute_ts, model), stats in model_updates.items() + ] + conn.executemany( + """ + INSERT INTO analytics_model_minute ( + minute_ts, + model, + requests, + successful, + failed, + revenue_msats, + refunds_msats, + input_tokens, + output_tokens, + total_tokens + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(minute_ts, model) DO UPDATE SET + requests = requests + excluded.requests, + successful = successful + excluded.successful, + failed = failed + excluded.failed, + revenue_msats = revenue_msats + excluded.revenue_msats, + refunds_msats = refunds_msats + excluded.refunds_msats, + input_tokens = input_tokens + excluded.input_tokens, + output_tokens = output_tokens + excluded.output_tokens, + total_tokens = total_tokens + excluded.total_tokens + """, + model_rows, + ) + + if model_presence_updates: + presence_rows = [ + (minute_ts, model, count) + for (minute_ts, model), count in model_presence_updates.items() + ] + conn.executemany( + """ + INSERT INTO analytics_model_presence_minute ( + minute_ts, + model, + count + ) + VALUES (?, ?, ?) + ON CONFLICT(minute_ts, model) DO UPDATE SET + count = count + excluded.count + """, + presence_rows, + ) + + if error_type_updates: + error_type_rows = [ + (minute_ts, error_type, count) + for (minute_ts, error_type), count in error_type_updates.items() + ] + conn.executemany( + """ + INSERT INTO analytics_error_type_minute ( + minute_ts, + error_type, + count + ) + VALUES (?, ?, ?) + ON CONFLICT(minute_ts, error_type) DO UPDATE SET + count = count + excluded.count + """, + error_type_rows, + ) + + if error_events: + conn.executemany( + """ + INSERT INTO analytics_error_events ( + timestamp, + message, + error_type, + pathname, + lineno, + request_id + ) + VALUES (?, ?, ?, ?, ?, ?) + """, + error_events, + ) + + def _query_metrics_locked( + self, + conn: sqlite3.Connection, + *, + cutoff_timestamp: str, + interval_minutes: int, + hours_back: int, + ) -> dict[str, Any]: + bucket_seconds = max(60, int(interval_minutes) * 60) + rows = conn.execute( + """ + SELECT + datetime( + (CAST(strftime('%s', minute_ts) AS INTEGER) / ?) * ?, + 'unixepoch' + ) AS bucket_ts, + COALESCE(SUM(total_requests), 0) AS total_requests, + COALESCE(SUM(successful_chat_completions), 0) AS successful_chat_completions, + COALESCE(SUM(failed_requests), 0) AS failed_requests, + COALESCE(SUM(errors), 0) AS errors, + COALESCE(SUM(warnings), 0) AS warnings, + COALESCE(SUM(payment_processed), 0) AS payment_processed, + COALESCE(SUM(upstream_errors), 0) AS upstream_errors, + COALESCE(SUM(revenue_msats), 0) AS revenue_msats, + COALESCE(SUM(refunds_msats), 0) AS refunds_msats, + COALESCE(SUM(input_tokens), 0) AS input_tokens, + COALESCE(SUM(output_tokens), 0) AS output_tokens, + COALESCE(SUM(total_tokens), 0) AS total_tokens + FROM analytics_minute + WHERE minute_ts >= ? + GROUP BY bucket_ts + ORDER BY bucket_ts + """, + (bucket_seconds, bucket_seconds, cutoff_timestamp), + ).fetchall() + + totals: dict[str, float] = { + "total_requests": 0.0, + "successful_chat_completions": 0.0, + "failed_requests": 0.0, + "errors": 0.0, + "warnings": 0.0, + "payment_processed": 0.0, + "upstream_errors": 0.0, + "revenue_msats": 0.0, + "refunds_msats": 0.0, + "input_tokens": 0.0, + "output_tokens": 0.0, + "total_tokens": 0.0, + } + + points: list[dict[str, Any]] = [] + for row in rows: + total_requests = int(row["total_requests"]) + successful = int(row["successful_chat_completions"]) + failed = int(row["failed_requests"]) + errors = int(row["errors"]) + warnings = int(row["warnings"]) + payment_processed = int(row["payment_processed"]) + upstream_errors = int(row["upstream_errors"]) + revenue_msats = float(row["revenue_msats"]) + refunds_msats = float(row["refunds_msats"]) + input_tokens = int(row["input_tokens"]) + output_tokens = int(row["output_tokens"]) + total_tokens = int(row["total_tokens"]) + + totals["total_requests"] += total_requests + totals["successful_chat_completions"] += successful + totals["failed_requests"] += failed + totals["errors"] += errors + totals["warnings"] += warnings + totals["payment_processed"] += payment_processed + totals["upstream_errors"] += upstream_errors + totals["revenue_msats"] += revenue_msats + totals["refunds_msats"] += refunds_msats + totals["input_tokens"] += input_tokens + totals["output_tokens"] += output_tokens + totals["total_tokens"] += total_tokens + + points.append( + { + "timestamp": str(row["bucket_ts"]), + "total_requests": total_requests, + "successful_chat_completions": successful, + "failed_requests": failed, + "errors": errors, + "warnings": warnings, + "payment_processed": payment_processed, + "upstream_errors": upstream_errors, + "revenue_msats": revenue_msats, + "refunds_msats": refunds_msats, + "input_tokens": input_tokens, + "output_tokens": output_tokens, + "total_tokens": total_tokens, + "requests": total_requests, + } + ) + + normalized_totals: dict[str, int | float] = { + "total_requests": int(totals["total_requests"]), + "successful_chat_completions": int(totals["successful_chat_completions"]), + "failed_requests": int(totals["failed_requests"]), + "errors": int(totals["errors"]), + "warnings": int(totals["warnings"]), + "payment_processed": int(totals["payment_processed"]), + "upstream_errors": int(totals["upstream_errors"]), + "revenue_msats": float(totals["revenue_msats"]), + "refunds_msats": float(totals["refunds_msats"]), + "input_tokens": int(totals["input_tokens"]), + "output_tokens": int(totals["output_tokens"]), + "total_tokens": int(totals["total_tokens"]), + } + + return { + "metrics": points, + "interval_minutes": interval_minutes, + "hours_back": hours_back, + "total_buckets": len(points), + "totals": normalized_totals, + } + + def _query_summary_locked( + self, conn: sqlite3.Connection, cutoff_timestamp: str + ) -> dict[str, Any]: + totals = conn.execute( + """ + SELECT + COALESCE(SUM(total_entries), 0) AS total_entries, + COALESCE(SUM(total_requests), 0) AS total_requests, + COALESCE(SUM(successful_chat_completions), 0) AS successful_chat_completions, + COALESCE(SUM(failed_requests), 0) AS failed_requests, + COALESCE(SUM(errors), 0) AS total_errors, + COALESCE(SUM(warnings), 0) AS total_warnings, + COALESCE(SUM(payment_processed), 0) AS payment_processed, + COALESCE(SUM(upstream_errors), 0) AS upstream_errors, + COALESCE(SUM(revenue_msats), 0) AS revenue_msats, + COALESCE(SUM(refunds_msats), 0) AS refunds_msats, + COALESCE(SUM(input_tokens), 0) AS input_tokens, + COALESCE(SUM(output_tokens), 0) AS output_tokens, + COALESCE(SUM(total_tokens), 0) AS total_tokens + FROM analytics_minute + WHERE minute_ts >= ? + """, + (cutoff_timestamp,), + ).fetchone() + + unique_models = [ + str(row[0]) + for row in conn.execute( + """ + SELECT DISTINCT model + FROM analytics_model_presence_minute + WHERE minute_ts >= ? + ORDER BY model ASC + """, + (cutoff_timestamp,), + ).fetchall() + ] + + error_types = { + str(row[0]): int(row[1]) + for row in conn.execute( + """ + SELECT error_type, COALESCE(SUM(count), 0) AS total_count + FROM analytics_error_type_minute + WHERE minute_ts >= ? + GROUP BY error_type + """, + (cutoff_timestamp,), + ).fetchall() + } + + total_requests = int(totals["total_requests"]) + successful = int(totals["successful_chat_completions"]) + failed_requests = int(totals["failed_requests"]) + input_tokens = int(totals["input_tokens"]) + output_tokens = int(totals["output_tokens"]) + total_tokens = int(totals["total_tokens"]) + + revenue_msats = float(totals["revenue_msats"]) + refunds_msats = float(totals["refunds_msats"]) + net_revenue_msats = revenue_msats - refunds_msats + + revenue_sats = revenue_msats / 1000 + refunds_sats = refunds_msats / 1000 + net_revenue_sats = net_revenue_msats / 1000 + + return { + "total_entries": int(totals["total_entries"]), + "total_requests": total_requests, + "successful_chat_completions": successful, + "failed_requests": failed_requests, + "total_errors": int(totals["total_errors"]), + "total_warnings": int(totals["total_warnings"]), + "payment_processed": int(totals["payment_processed"]), + "upstream_errors": int(totals["upstream_errors"]), + "unique_models_count": len(unique_models), + "unique_models": unique_models, + "error_types": error_types, + "input_tokens": input_tokens, + "output_tokens": output_tokens, + "total_tokens": total_tokens, + "avg_input_tokens_per_completion": (input_tokens / successful) + if successful > 0 + else 0, + "avg_output_tokens_per_completion": (output_tokens / successful) + if successful > 0 + else 0, + "avg_total_tokens_per_completion": (total_tokens / successful) + if successful > 0 + else 0, + "success_rate": (successful / total_requests * 100) + if total_requests > 0 + else 0, + "revenue_msats": revenue_msats, + "refunds_msats": refunds_msats, + "revenue_sats": revenue_sats, + "refunds_sats": refunds_sats, + "net_revenue_msats": net_revenue_msats, + "net_revenue_sats": net_revenue_sats, + "avg_revenue_per_request_msats": (revenue_msats / successful) + if successful > 0 + else 0, + "refund_rate": (failed_requests / total_requests * 100) + if total_requests > 0 + else 0, + } + + def _query_error_details_locked( + self, + conn: sqlite3.Connection, + *, + cutoff_timestamp: str, + limit: int, + total_error_count: int | None = None, + ) -> dict[str, Any]: + rows = conn.execute( + """ + SELECT + timestamp, + message, + error_type, + pathname, + lineno, + request_id + FROM analytics_error_events + WHERE timestamp >= ? + ORDER BY timestamp DESC + LIMIT ? + """, + (cutoff_timestamp, limit), + ).fetchall() + + if total_error_count is None: + total_error_count_row = conn.execute( + """ + SELECT COALESCE(SUM(errors), 0) + FROM analytics_minute + WHERE minute_ts >= ? + """, + (cutoff_timestamp,), + ).fetchone() + total_error_count = int(total_error_count_row[0]) if total_error_count_row else 0 + + return { + "errors": [ + { + "timestamp": str(row["timestamp"]), + "message": str(row["message"]), + "error_type": str(row["error_type"]), + "pathname": str(row["pathname"]), + "lineno": int(row["lineno"]), + "request_id": str(row["request_id"]), + } + for row in rows + ], + "total_count": int(total_error_count), + } + + def _query_revenue_by_model_locked( + self, + conn: sqlite3.Connection, + *, + cutoff_timestamp: str, + limit: int, + ) -> dict[str, Any]: + rows = conn.execute( + """ + SELECT + model, + COALESCE(SUM(revenue_msats), 0) AS revenue_msats, + COALESCE(SUM(refunds_msats), 0) AS refunds_msats, + COALESCE(SUM(requests), 0) AS requests, + COALESCE(SUM(successful), 0) AS successful, + COALESCE(SUM(failed), 0) AS failed + FROM analytics_model_minute + WHERE minute_ts >= ? + GROUP BY model + ORDER BY (COALESCE(SUM(revenue_msats), 0) - COALESCE(SUM(refunds_msats), 0)) DESC + """, + (cutoff_timestamp,), + ).fetchall() + + models: list[dict[str, Any]] = [] + total_revenue_sats = 0.0 + + for row in rows: + revenue_msats = float(row["revenue_msats"]) + refunds_msats = float(row["refunds_msats"]) + revenue_sats = revenue_msats / 1000 + refunds_sats = refunds_msats / 1000 + net_revenue_sats = revenue_sats - refunds_sats + successful = int(row["successful"]) + + models.append( + { + "model": str(row["model"]), + "revenue_sats": revenue_sats, + "refunds_sats": refunds_sats, + "net_revenue_sats": net_revenue_sats, + "requests": int(row["requests"]), + "successful": successful, + "failed": int(row["failed"]), + "avg_revenue_per_request": (revenue_sats / successful) + if successful > 0 + else 0, + } + ) + total_revenue_sats += net_revenue_sats + + return { + "models": models[:limit], + "total_revenue_sats": total_revenue_sats, + "total_models": len(models), + } + + def _query_model_usage_mix_locked( + self, + conn: sqlite3.Connection, + *, + cutoff_timestamp: str, + interval_minutes: int, + hours_back: int, + limit: int, + ) -> dict[str, Any]: + top_limit = max(1, min(int(limit), 20)) + top_rows_requests = conn.execute( + """ + SELECT + model, + COALESCE(SUM(successful), 0) AS total_successful + FROM analytics_model_minute + WHERE minute_ts >= ? + AND model != 'unknown' + GROUP BY model + ORDER BY total_successful DESC + LIMIT ? + """, + (cutoff_timestamp, top_limit), + ).fetchall() + top_rows_revenue = conn.execute( + """ + SELECT + model, + COALESCE(SUM(revenue_msats), 0) AS total_revenue_msats + FROM analytics_model_minute + WHERE minute_ts >= ? + AND model != 'unknown' + GROUP BY model + ORDER BY total_revenue_msats DESC + LIMIT ? + """, + (cutoff_timestamp, top_limit), + ).fetchall() + top_rows_tokens = conn.execute( + """ + SELECT + model, + COALESCE(SUM(total_tokens), 0) AS total_tokens + FROM analytics_model_minute + WHERE minute_ts >= ? + AND model != 'unknown' + GROUP BY model + ORDER BY total_tokens DESC + LIMIT ? + """, + (cutoff_timestamp, top_limit), + ).fetchall() + + top_models_requests = [ + str(row["model"]) + for row in top_rows_requests + if int(row["total_successful"] or 0) > 0 + ] + top_models_revenue = [ + str(row["model"]) + for row in top_rows_revenue + if float(row["total_revenue_msats"] or 0.0) > 0 + ] + top_models_tokens = [ + str(row["model"]) + for row in top_rows_tokens + if int(row["total_tokens"] or 0) > 0 + ] + + selected_models: list[str] = [] + for model in top_models_requests + top_models_revenue + top_models_tokens: + if model not in selected_models: + selected_models.append(model) + + bucket_seconds = max(60, int(interval_minutes) * 60) + total_rows = conn.execute( + """ + SELECT + datetime( + (CAST(strftime('%s', minute_ts) AS INTEGER) / ?) * ?, + 'unixepoch' + ) AS bucket_ts, + COALESCE(SUM(successful), 0) AS total_successful, + COALESCE(SUM(revenue_msats), 0) AS total_revenue_msats, + COALESCE(SUM(total_tokens), 0) AS total_tokens + FROM analytics_model_minute + WHERE minute_ts >= ? + GROUP BY bucket_ts + ORDER BY bucket_ts + """, + (bucket_seconds, bucket_seconds, cutoff_timestamp), + ).fetchall() + + bucket_index: dict[str, dict[str, Any]] = {} + for row in total_rows: + total_successful = int(row["total_successful"]) + total_revenue_msats = float(row["total_revenue_msats"]) + total_tokens = int(row["total_tokens"]) + if ( + total_successful <= 0 + and total_revenue_msats <= 0 + and total_tokens <= 0 + ): + continue + + bucket_ts = str(row["bucket_ts"]) + bucket = bucket_index.setdefault( + bucket_ts, + { + "timestamp": bucket_ts, + "total_successful": 0, + "total_revenue_msats": 0.0, + "total_tokens": 0, + "others": 0, + "others_revenue_msats": 0.0, + "others_tokens": 0, + "model_counts": {}, + "model_revenue_msats": {}, + "model_tokens": {}, + }, + ) + bucket["total_successful"] = total_successful + bucket["total_revenue_msats"] = total_revenue_msats + bucket["total_tokens"] = total_tokens + bucket["others"] = total_successful + bucket["others_revenue_msats"] = total_revenue_msats + bucket["others_tokens"] = total_tokens + + if selected_models and bucket_index: + placeholders = ",".join("?" for _ in selected_models) + top_model_rows = conn.execute( + f""" + SELECT + datetime( + (CAST(strftime('%s', minute_ts) AS INTEGER) / ?) * ?, + 'unixepoch' + ) AS bucket_ts, + model, + COALESCE(SUM(successful), 0) AS successful, + COALESCE(SUM(revenue_msats), 0) AS revenue_msats, + COALESCE(SUM(total_tokens), 0) AS total_tokens + FROM analytics_model_minute + WHERE minute_ts >= ? + AND model IN ({placeholders}) + GROUP BY bucket_ts, model + ORDER BY bucket_ts + """, + (bucket_seconds, bucket_seconds, cutoff_timestamp, *selected_models), + ).fetchall() + + for row in top_model_rows: + bucket_ts = str(row["bucket_ts"]) + if bucket_ts not in bucket_index: + continue + bucket = bucket_index[bucket_ts] + + model = str(row["model"]) + successful = int(row["successful"]) + revenue_msats = float(row["revenue_msats"]) + total_tokens = int(row["total_tokens"]) + + model_counts = bucket["model_counts"] + model_counts[model] = successful + model_revenue_msats = bucket["model_revenue_msats"] + model_revenue_msats[model] = revenue_msats + model_tokens = bucket["model_tokens"] + model_tokens[model] = total_tokens + + bucket["others"] = max(0, int(bucket["others"]) - successful) + bucket["others_revenue_msats"] = max( + 0.0, + float(bucket["others_revenue_msats"]) - revenue_msats, + ) + bucket["others_tokens"] = max( + 0, + int(bucket["others_tokens"]) - total_tokens, + ) + + metrics = sorted(bucket_index.values(), key=lambda item: str(item["timestamp"])) + + return { + "top_models": top_models_requests, + "top_models_by_metric": { + "requests": top_models_requests, + "revenue": top_models_revenue, + "tokens": top_models_tokens, + }, + "metrics": metrics, + "hours_back": hours_back, + "interval_minutes": interval_minutes, + "total_buckets": len(metrics), + } + + def _cutoff_timestamp(self, hours_back: int) -> str: + cutoff = datetime.now(timezone.utc) - timedelta(hours=hours_back) + return cutoff.strftime("%Y-%m-%d %H:%M:%S") + + def _minute_key(self, timestamp: Any) -> str | None: + if not isinstance(timestamp, str) or len(timestamp) != 19: + return None + if timestamp[10] != " ": + return None + return f"{timestamp[:16]}:00" + + def _extract_success_metrics( + self, entry: dict[str, Any], message: str + ) -> tuple[bool, float, int, int]: + # These auth logs are emitted once per successful settlement across providers + # and avoid duplicate counting from provider-specific completion logs. + logger_name = str(entry.get("name", "")) + if not logger_name.startswith("routstr.auth"): + return False, 0.0, 0, 0 + + input_tokens = self._parse_token_count(entry.get("input_tokens", 0)) + output_tokens = self._parse_token_count(entry.get("output_tokens", 0)) + + if "calculated token-based cost" in message: + token_cost = entry.get("token_cost", 0) + if isinstance(token_cost, (int, float)) and token_cost > 0: + return True, float(token_cost), input_tokens, output_tokens + return True, 0.0, input_tokens, output_tokens + + if "max cost payment finalized" in message: + charged_amount = entry.get("charged_amount", 0) + if isinstance(charged_amount, (int, float)) and charged_amount > 0: + return True, float(charged_amount), input_tokens, output_tokens + return True, 0.0, input_tokens, output_tokens + + return False, 0.0, 0, 0 + + def _parse_token_count(self, value: Any) -> int: + if isinstance(value, bool): + return 0 + if isinstance(value, int): + return max(0, value) + if isinstance(value, float): + return max(0, int(value)) + if isinstance(value, str): + try: + return max(0, int(float(value))) + except ValueError: + return 0 + return 0 + + def _new_minute_stats(self) -> dict[str, float]: + return { + "total_entries": 0.0, + "total_requests": 0.0, + "successful_chat_completions": 0.0, + "failed_requests": 0.0, + "errors": 0.0, + "warnings": 0.0, + "payment_processed": 0.0, + "upstream_errors": 0.0, + "revenue_msats": 0.0, + "refunds_msats": 0.0, + "input_tokens": 0.0, + "output_tokens": 0.0, + "total_tokens": 0.0, + } + + def _new_model_stats(self) -> dict[str, float]: + return { + "requests": 0.0, + "successful": 0.0, + "failed": 0.0, + "revenue_msats": 0.0, + "refunds_msats": 0.0, + "input_tokens": 0.0, + "output_tokens": 0.0, + "total_tokens": 0.0, + } diff --git a/routstr/nostr/__init__.py b/routstr/nostr/__init__.py index b19039a7..afd165f5 100644 --- a/routstr/nostr/__init__.py +++ b/routstr/nostr/__init__.py @@ -1,4 +1,5 @@ +from .analytics import publish_usage_analytics from .discovery import providers_cache_refresher from .listing import announce_provider -__all__ = ["providers_cache_refresher", "announce_provider"] +__all__ = ["providers_cache_refresher", "announce_provider", "publish_usage_analytics"] diff --git a/routstr/nostr/analytics.py b/routstr/nostr/analytics.py new file mode 100644 index 00000000..e568b5e0 --- /dev/null +++ b/routstr/nostr/analytics.py @@ -0,0 +1,419 @@ +#!/usr/bin/env python3 +""" +Nostr usage analytics publisher. +Publishes a single replaceable analytics snapshot for each provider. +""" + +from __future__ import annotations + +import asyncio +import hashlib +import json +import time +from typing import Any + +from nostr.event import Event +from nostr.key import PrivateKey + +from ..core import get_logger +from ..core.log_manager import log_manager +from ..core.settings import settings +from .listing import nsec_to_keypair, publish_to_relay + +logger = get_logger(__name__) + +ANALYTICS_KIND = 38422 +ANALYTICS_SCHEMA = "routstr.analytics.snapshot.v1" +DEFAULT_RELAYS = [ + "wss://relay.nostr.band", + "wss://relay.damus.io", + "wss://relay.routstr.com", + "wss://nos.lol", +] +PUBLISH_INTERVAL_SECONDS = 15 * 60 +DISABLED_POLL_SECONDS = 60 +DASHBOARD_WINDOW_HOURS = 24 +DASHBOARD_INTERVAL_MINUTES = 60 +MODEL_LIMIT = 20 +WINDOW_DEFINITIONS: tuple[tuple[str, int, int], ...] = ( + ("24h", 24, 60), + ("7d", 7 * 24, 6 * 60), + ("30d", 30 * 24, 24 * 60), + ("3m", 90 * 24, 24 * 60), + ("1y", 365 * 24, 7 * 24 * 60), +) + + +def _event_to_dict(ev: Event) -> dict[str, Any]: + return { + "id": ev.id, + "pubkey": ev.public_key, + "created_at": ev.created_at, + "kind": int(ev.kind) if not isinstance(ev.kind, int) else ev.kind, + "tags": ev.tags, + "content": ev.content, + "sig": ev.signature, + } + + +def _resolve_provider_id(public_key_hex: str) -> str: + explicit_provider_id = (settings.provider_id or "").strip() + if explicit_provider_id: + return explicit_provider_id + return public_key_hex[:12] + + +def _resolve_endpoint_urls() -> list[str]: + urls: list[str] = [] + http_url = (settings.http_url or "").strip() + onion_url = (settings.onion_url or "").strip() + + if http_url and http_url != "http://localhost:8000": + urls.append(http_url) + + if onion_url: + if onion_url.endswith(".onion") and not ( + onion_url.startswith("http://") or onion_url.startswith("https://") + ): + onion_url = f"http://{onion_url}" + urls.append(onion_url) + + return urls + + +def _resolve_relays() -> list[str]: + configured = [url.strip() for url in settings.relays if url.strip()] + return configured if configured else list(DEFAULT_RELAYS) + + +def _to_int(value: Any) -> int: + if isinstance(value, bool): + return int(value) + if isinstance(value, int): + return value + if isinstance(value, float): + return int(value) + if isinstance(value, str): + try: + return int(float(value)) + except ValueError: + return 0 + return 0 + + +def _to_float(value: Any) -> float: + if isinstance(value, bool): + return float(int(value)) + if isinstance(value, (int, float)): + return float(value) + if isinstance(value, str): + try: + return float(value) + except ValueError: + return 0.0 + return 0.0 + + +def _aggregate_top_model_usage( + model_usage_mix: dict[str, Any], +) -> tuple[list[dict[str, Any]], dict[str, Any]]: + top_models_raw = model_usage_mix.get("top_models", []) + mix_metrics_raw = model_usage_mix.get("metrics", []) + + top_models = [model for model in top_models_raw if isinstance(model, str)] + metrics = [row for row in mix_metrics_raw if isinstance(row, dict)] + + model_totals: dict[str, dict[str, float | int]] = { + model: { + "successful_requests": 0, + "revenue_msats": 0.0, + "total_tokens": 0, + } + for model in top_models + } + others = { + "successful_requests": 0, + "revenue_msats": 0.0, + "total_tokens": 0, + } + + for metric in metrics: + model_counts = metric.get("model_counts", {}) + model_revenue = metric.get("model_revenue_msats", {}) + model_tokens = metric.get("model_tokens", {}) + + if isinstance(model_counts, dict): + for model, count in model_counts.items(): + if model in model_totals: + model_totals[model]["successful_requests"] += _to_int(count) + + if isinstance(model_revenue, dict): + for model, amount in model_revenue.items(): + if model in model_totals: + model_totals[model]["revenue_msats"] += _to_float(amount) + + if isinstance(model_tokens, dict): + for model, token_count in model_tokens.items(): + if model in model_totals: + model_totals[model]["total_tokens"] += _to_int(token_count) + + others["successful_requests"] += _to_int(metric.get("others", 0)) + others["revenue_msats"] += _to_float(metric.get("others_revenue_msats", 0.0)) + others["total_tokens"] += _to_int(metric.get("others_tokens", 0)) + + model_rows = [ + { + "model": model, + "successful_requests": int(values["successful_requests"]), + "revenue_msats": float(values["revenue_msats"]), + "total_tokens": int(values["total_tokens"]), + } + for model, values in model_totals.items() + ] + model_rows.sort( + key=lambda row: _to_int(row.get("successful_requests", 0)), + reverse=True, + ) + + return model_rows, others + + +def _build_summary_payload(summary: dict[str, Any]) -> dict[str, Any]: + return { + "total_requests": _to_int(summary.get("total_requests", 0)), + "successful_chat_completions": _to_int( + summary.get("successful_chat_completions", 0) + ), + "failed_requests": _to_int(summary.get("failed_requests", 0)), + "success_rate": _to_float(summary.get("success_rate", 0.0)), + "unique_models_count": _to_int(summary.get("unique_models_count", 0)), + "input_tokens": _to_int(summary.get("input_tokens", 0)), + "output_tokens": _to_int(summary.get("output_tokens", 0)), + "total_tokens": _to_int(summary.get("total_tokens", 0)), + "revenue_msats": _to_float(summary.get("revenue_msats", 0.0)), + "refunds_msats": _to_float(summary.get("refunds_msats", 0.0)), + "net_revenue_msats": _to_float(summary.get("net_revenue_msats", 0.0)), + "revenue_sats": _to_float(summary.get("revenue_sats", 0.0)), + "refunds_sats": _to_float(summary.get("refunds_sats", 0.0)), + "net_revenue_sats": _to_float(summary.get("net_revenue_sats", 0.0)), + } + + +def _build_window_payload( + *, + hours: int, + interval_minutes: int, + model_limit: int, +) -> dict[str, Any]: + dashboard = log_manager.get_usage_dashboard( + interval=interval_minutes, + hours=hours, + error_limit=1, + model_limit=model_limit, + ) + + summary = dashboard.get("summary", {}) + model_usage_mix = dashboard.get("model_usage_mix", {}) + + summary_payload = _build_summary_payload(summary if isinstance(summary, dict) else {}) + usage_mix_payload = model_usage_mix if isinstance(model_usage_mix, dict) else {} + top_model_usage, others_usage = _aggregate_top_model_usage(usage_mix_payload) + + return { + "window_hours": hours, + "interval_minutes": interval_minutes, + "summary": summary_payload, + "model_usage_mix": usage_mix_payload, + "top_model_usage": top_model_usage, + "others_usage": others_usage, + } + + +def build_stats_snapshot_payload( + provider_id: str, + *, + public_key_hex: str, + generated_at: int, + window_hours: int = DASHBOARD_WINDOW_HOURS, + interval_minutes: int = DASHBOARD_INTERVAL_MINUTES, + model_limit: int = MODEL_LIMIT, +) -> dict[str, Any]: + _ = (window_hours, interval_minutes) + windows: dict[str, dict[str, Any]] = {} + for key, hours, window_interval_minutes in WINDOW_DEFINITIONS: + windows[key] = _build_window_payload( + hours=hours, + interval_minutes=window_interval_minutes, + model_limit=model_limit, + ) + + primary_window = windows.get("24h", {}) + summary_payload = ( + primary_window.get("summary", {}) + if isinstance(primary_window.get("summary", {}), dict) + else {} + ) + usage_mix_payload = ( + primary_window.get("model_usage_mix", {}) + if isinstance(primary_window.get("model_usage_mix", {}), dict) + else {} + ) + top_model_usage = ( + primary_window.get("top_model_usage", []) + if isinstance(primary_window.get("top_model_usage", []), list) + else [] + ) + others_usage = ( + primary_window.get("others_usage", {}) + if isinstance(primary_window.get("others_usage", {}), dict) + else {} + ) + + return { + "schema": ANALYTICS_SCHEMA, + "generated_at": generated_at, + "provider_id": provider_id, + "pubkey": public_key_hex, + "npub": settings.npub or "", + "endpoint_urls": _resolve_endpoint_urls(), + "window_hours": DASHBOARD_WINDOW_HOURS, + "interval_minutes": DASHBOARD_INTERVAL_MINUTES, + "summary": summary_payload, + "model_usage_mix": usage_mix_payload, + "top_model_usage": top_model_usage, + "others_usage": others_usage, + "windows": windows, + } + + +def create_stats_snapshot_event( + private_key_hex: str, + provider_id: str, + payload_json: str, + *, + d_tag: str, +) -> dict[str, Any]: + private_key = PrivateKey(bytes.fromhex(private_key_hex)) + tags = [ + ["d", d_tag], + ["provider", provider_id], + ["schema", ANALYTICS_SCHEMA], + ] + + event = Event( + public_key=private_key.public_key.hex(), + content=payload_json, + kind=ANALYTICS_KIND, + tags=tags, + ) + private_key.sign_event(event) + return _event_to_dict(event) + + +def _fingerprint_payload(payload: dict[str, Any]) -> str: + normalized = dict(payload) + # Ignore generated timestamp for semantic dedupe. + normalized.pop("generated_at", None) + payload_json = json.dumps(normalized, separators=(",", ":"), sort_keys=True) + return hashlib.sha256(payload_json.encode("utf-8")).hexdigest() + + +async def publish_usage_analytics() -> None: + last_payload_hash: str | None = None + + parsed_nsec: str | None = None + private_key_hex: str | None = None + public_key_hex: str | None = None + provider_id: str | None = None + warned_missing_nsec = False + + logger.info("Usage analytics sharing task started") + + while True: + try: + if not settings.enable_analytics_sharing: + await asyncio.sleep(DISABLED_POLL_SECONDS) + continue + + nsec = (settings.nsec or "").strip() + if not nsec: + if not warned_missing_nsec: + logger.info("NSEC is not configured; skipping analytics sharing to Nostr") + warned_missing_nsec = True + await asyncio.sleep(DISABLED_POLL_SECONDS) + continue + + warned_missing_nsec = False + if nsec != parsed_nsec or private_key_hex is None or public_key_hex is None: + keypair = nsec_to_keypair(nsec) + if not keypair: + logger.error("Invalid NSEC; analytics sharing is paused") + await asyncio.sleep(DISABLED_POLL_SECONDS) + continue + private_key_hex, public_key_hex = keypair + parsed_nsec = nsec + provider_id = _resolve_provider_id(public_key_hex) + last_payload_hash = None + + if private_key_hex is None or public_key_hex is None: + await asyncio.sleep(DISABLED_POLL_SECONDS) + continue + + relay_urls = _resolve_relays() + if not relay_urls: + logger.warning("No Nostr relays configured; analytics sharing skipped") + await asyncio.sleep(DISABLED_POLL_SECONDS) + continue + + resolved_provider_id = provider_id or _resolve_provider_id(public_key_hex) + now_ts = int(time.time()) + payload = build_stats_snapshot_payload( + resolved_provider_id, + public_key_hex=public_key_hex, + generated_at=now_ts, + ) + + payload_hash = _fingerprint_payload(payload) + if last_payload_hash == payload_hash: + await asyncio.sleep(PUBLISH_INTERVAL_SECONDS) + continue + + payload_json = json.dumps(payload, separators=(",", ":"), sort_keys=True) + d_tag = f"{resolved_provider_id}:stats" + event = create_stats_snapshot_event( + private_key_hex, + resolved_provider_id, + payload_json, + d_tag=d_tag, + ) + + success_count = 0 + for relay_url in relay_urls: + if await publish_to_relay(relay_url, event): + success_count += 1 + + if success_count > 0: + last_payload_hash = payload_hash + + logger.info( + "Published analytics snapshot (success=%s/%s provider=%s)", + success_count, + len(relay_urls), + resolved_provider_id, + extra={ + "relay_success_count": success_count, + "relay_total": len(relay_urls), + "provider_id": resolved_provider_id, + }, + ) + await asyncio.sleep(PUBLISH_INTERVAL_SECONDS) + + except asyncio.CancelledError: + logger.info("Usage analytics sharing task cancelled") + break + except Exception as e: + logger.error( + "Usage analytics sharing error", + extra={"error": str(e), "error_type": type(e).__name__}, + ) + await asyncio.sleep(DISABLED_POLL_SECONDS) diff --git a/routstr/payment/cost_calculation.py b/routstr/payment/cost_calculation.py index 1df4caf8..2ed7e4b8 100644 --- a/routstr/payment/cost_calculation.py +++ b/routstr/payment/cost_calculation.py @@ -16,6 +16,8 @@ class CostData(BaseModel): output_msats: int total_msats: int total_usd: float = 0.0 + input_tokens: int = 0 + output_tokens: int = 0 class MaxCostData(CostData): @@ -63,10 +65,49 @@ async def calculate_cost( # todo: can be sync output_msats=0, total_msats=0, total_usd=0.0, + input_tokens=0, + output_tokens=0, ) usage_data = response_data["usage"] + def parse_token_count(value: object) -> int: + if isinstance(value, bool): + return 0 + if isinstance(value, int): + return max(0, value) + if isinstance(value, float): + return max(0, int(value)) + if isinstance(value, str): + try: + return max(0, int(float(value))) + except ValueError: + return 0 + return 0 + + input_tokens = parse_token_count(usage_data.get("prompt_tokens", 0)) + output_tokens = parse_token_count(usage_data.get("completion_tokens", 0)) + input_tokens = ( + input_tokens + if input_tokens != 0 + else parse_token_count(usage_data.get("input_tokens", 0)) + ) + output_tokens = ( + output_tokens + if output_tokens != 0 + else parse_token_count(usage_data.get("output_tokens", 0)) + ) + input_tokens = ( + input_tokens + if input_tokens != 0 + else parse_token_count(response_data.get("usage", {}).get("input_tokens", 0)) + ) + output_tokens = ( + output_tokens + if output_tokens != 0 + else parse_token_count(response_data.get("usage", {}).get("output_tokens", 0)) + ) + usd_cost = 0.0 # Prioritize cost_details.upstream_inference_cost @@ -104,6 +145,8 @@ async def calculate_cost( # todo: can be sync output_msats=-1, total_msats=cost_in_msats, total_usd=usd_cost, + input_tokens=input_tokens, + output_tokens=output_tokens, ) except Exception as e: logger.warning( @@ -184,31 +227,10 @@ async def calculate_cost( # todo: can be sync input_msats=0, output_msats=0, total_msats=max_cost, + input_tokens=input_tokens, + output_tokens=output_tokens, ) - input_tokens = usage_data.get("prompt_tokens", 0) - output_tokens = usage_data.get("completion_tokens", 0) - - # added for response api - input_tokens = ( - input_tokens if input_tokens != 0 else usage_data.get("input_tokens", 0) - ) - output_tokens = ( - output_tokens if output_tokens != 0 else usage_data.get("output_tokens", 0) - ) - - # added for response api - input_tokens = ( - input_tokens - if input_tokens != 0 - else response_data.get("usage", {}).get("input_tokens", 0) - ) - output_tokens = ( - output_tokens - if output_tokens != 0 - else response_data.get("usage", {}).get("output_tokens", 0) - ) - input_msats = round(input_tokens / 1000 * MSATS_PER_1K_INPUT_TOKENS, 3) output_msats = round(output_tokens / 1000 * MSATS_PER_1K_OUTPUT_TOKENS, 3) @@ -234,4 +256,6 @@ async def calculate_cost( # todo: can be sync output_msats=int(output_msats), total_msats=token_based_cost, total_usd=total_usd, + input_tokens=input_tokens, + output_tokens=output_tokens, ) diff --git a/tests/unit/test_nostr_analytics.py b/tests/unit/test_nostr_analytics.py new file mode 100644 index 00000000..9e159dd3 --- /dev/null +++ b/tests/unit/test_nostr_analytics.py @@ -0,0 +1,267 @@ +from __future__ import annotations + +import asyncio +from typing import Any + +import pytest + +from routstr.nostr import analytics + + +def test_aggregate_top_model_usage_sums_metrics() -> None: + model_usage_mix = { + "top_models": ["openai/gpt-4o", "anthropic/claude-3.5-sonnet"], + "metrics": [ + { + "model_counts": { + "openai/gpt-4o": 4, + "anthropic/claude-3.5-sonnet": 2, + }, + "model_revenue_msats": { + "openai/gpt-4o": 1500, + "anthropic/claude-3.5-sonnet": 700, + }, + "model_tokens": { + "openai/gpt-4o": 1200, + "anthropic/claude-3.5-sonnet": 600, + }, + "others": 1, + "others_revenue_msats": 300, + "others_tokens": 200, + }, + { + "model_counts": { + "openai/gpt-4o": 3, + "anthropic/claude-3.5-sonnet": 1, + }, + "model_revenue_msats": { + "openai/gpt-4o": 1000, + "anthropic/claude-3.5-sonnet": 500, + }, + "model_tokens": { + "openai/gpt-4o": 800, + "anthropic/claude-3.5-sonnet": 300, + }, + "others": 2, + "others_revenue_msats": 450, + "others_tokens": 350, + }, + ], + } + + rows, others = analytics._aggregate_top_model_usage(model_usage_mix) + assert rows == [ + { + "model": "openai/gpt-4o", + "successful_requests": 7, + "revenue_msats": 2500.0, + "total_tokens": 2000, + }, + { + "model": "anthropic/claude-3.5-sonnet", + "successful_requests": 3, + "revenue_msats": 1200.0, + "total_tokens": 900, + }, + ] + assert others == { + "successful_requests": 3, + "revenue_msats": 750.0, + "total_tokens": 550, + } + + +def test_build_stats_snapshot_payload_schema_and_shape(monkeypatch: Any) -> None: + seen_windows: set[tuple[int, int]] = set() + + def fake_usage_dashboard( + *, interval: int, hours: int, error_limit: int, model_limit: int + ) -> dict[str, Any]: + seen_windows.add((hours, interval)) + assert error_limit == 1 + assert model_limit == 20 + return { + "summary": { + "total_requests": hours, + "successful_chat_completions": max(1, hours - 1), + "failed_requests": 2, + "success_rate": 90.0, + "unique_models_count": 2, + "input_tokens": 2000, + "output_tokens": 1000, + "total_tokens": 3000, + "revenue_msats": 9000.0, + "refunds_msats": 1000.0, + "net_revenue_msats": 8000.0, + "revenue_sats": 9.0, + "refunds_sats": 1.0, + "net_revenue_sats": 8.0, + }, + "model_usage_mix": { + "top_models": ["openai/gpt-4o"], + "metrics": [ + { + "timestamp": "2026-03-02 10:00:00", + "model_counts": {"openai/gpt-4o": hours}, + "model_revenue_msats": {"openai/gpt-4o": float(hours * 100)}, + "model_tokens": {"openai/gpt-4o": hours * 10}, + "others": 4, + "others_revenue_msats": 1800.0, + "others_tokens": 400, + } + ], + }, + } + + monkeypatch.setattr( + analytics.log_manager, "get_usage_dashboard", fake_usage_dashboard + ) + monkeypatch.setattr(analytics.settings, "npub", "npub1example") + monkeypatch.setattr(analytics.settings, "http_url", "https://node.example.com") + monkeypatch.setattr(analytics.settings, "onion_url", "") + + payload = analytics.build_stats_snapshot_payload( + "provider123", + public_key_hex="ab" * 32, + generated_at=1772451600, + ) + + assert payload["schema"] == analytics.ANALYTICS_SCHEMA + assert payload["provider_id"] == "provider123" + assert payload["window_hours"] == 24 + assert payload["interval_minutes"] == 60 + assert payload["endpoint_urls"] == ["https://node.example.com"] + assert seen_windows == { + (24, 60), + (7 * 24, 6 * 60), + (30 * 24, 24 * 60), + (90 * 24, 24 * 60), + (365 * 24, 7 * 24 * 60), + } + assert set(payload["windows"].keys()) == {"24h", "7d", "30d", "3m", "1y"} + assert payload["windows"]["1y"]["interval_minutes"] == 7 * 24 * 60 + assert payload["summary"]["total_requests"] == 24 + assert payload["top_model_usage"] == [ + { + "model": "openai/gpt-4o", + "successful_requests": 24, + "revenue_msats": 2400.0, + "total_tokens": 240, + } + ] + assert payload["others_usage"] == { + "successful_requests": 4, + "revenue_msats": 1800.0, + "total_tokens": 400, + } + + +def test_create_stats_snapshot_event_tags() -> None: + private_key_hex = "11" * 32 + event = analytics.create_stats_snapshot_event( + private_key_hex, + "provider123", + payload_json='{"schema":"routstr.analytics.snapshot.v1"}', + d_tag="provider123:stats", + ) + + tags = event["tags"] + assert ["d", "provider123:stats"] in tags + assert ["provider", "provider123"] in tags + assert ["schema", analytics.ANALYTICS_SCHEMA] in tags + assert all(tag[0] != "period" for tag in tags) + + +def test_fingerprint_payload_ignores_generated_at() -> None: + a = {"schema": analytics.ANALYTICS_SCHEMA, "generated_at": 1000, "summary": {"x": 1}} + b = {"schema": analytics.ANALYTICS_SCHEMA, "generated_at": 2000, "summary": {"x": 1}} + + assert analytics._fingerprint_payload(a) == analytics._fingerprint_payload(b) + + +@pytest.mark.asyncio +async def test_publish_usage_analytics_skips_when_disabled(monkeypatch: Any) -> None: + delays: list[int] = [] + + async def fake_sleep(seconds: int) -> None: + delays.append(seconds) + raise asyncio.CancelledError() + + def fail_build(*args: Any, **kwargs: Any) -> dict[str, Any]: + raise AssertionError("build_stats_snapshot_payload should not be called") + + monkeypatch.setattr(analytics.settings, "enable_analytics_sharing", False) + monkeypatch.setattr(analytics, "build_stats_snapshot_payload", fail_build) + monkeypatch.setattr(analytics.asyncio, "sleep", fake_sleep) + + await analytics.publish_usage_analytics() + + assert delays == [analytics.DISABLED_POLL_SECONDS] + + +@pytest.mark.asyncio +async def test_publish_usage_analytics_skips_without_nsec(monkeypatch: Any) -> None: + delays: list[int] = [] + + async def fake_sleep(seconds: int) -> None: + delays.append(seconds) + raise asyncio.CancelledError() + + def fail_build(*args: Any, **kwargs: Any) -> dict[str, Any]: + raise AssertionError("build_stats_snapshot_payload should not be called") + + monkeypatch.setattr(analytics.settings, "enable_analytics_sharing", True) + monkeypatch.setattr(analytics.settings, "nsec", "") + monkeypatch.setattr(analytics, "build_stats_snapshot_payload", fail_build) + monkeypatch.setattr(analytics.asyncio, "sleep", fake_sleep) + + await analytics.publish_usage_analytics() + + assert delays == [analytics.DISABLED_POLL_SECONDS] + + +@pytest.mark.asyncio +async def test_publish_usage_analytics_dedupes_unchanged_payload(monkeypatch: Any) -> None: + published_events: list[dict[str, Any]] = [] + sleep_calls = 0 + + async def fake_sleep(seconds: int) -> None: + nonlocal sleep_calls + sleep_calls += 1 + if sleep_calls >= 2: + raise asyncio.CancelledError() + + def fake_build_payload( + provider_id: str, + *, + public_key_hex: str, + generated_at: int, + window_hours: int = 24, + interval_minutes: int = 60, + model_limit: int = 10, + ) -> dict[str, Any]: + _ = (public_key_hex, generated_at, window_hours, interval_minutes, model_limit) + return { + "schema": analytics.ANALYTICS_SCHEMA, + "generated_at": generated_at, + "provider_id": provider_id, + "summary": {"total_requests": 1}, + } + + async def fake_publish(relay_url: str, event: dict[str, Any]) -> bool: + _ = relay_url + published_events.append(event) + return True + + monkeypatch.setattr(analytics.settings, "enable_analytics_sharing", True) + monkeypatch.setattr(analytics.settings, "nsec", "11" * 32) + monkeypatch.setattr(analytics.settings, "relays", ["wss://relay.example.com"]) + monkeypatch.setattr(analytics.settings, "provider_id", "") + monkeypatch.setattr(analytics, "build_stats_snapshot_payload", fake_build_payload) + monkeypatch.setattr(analytics, "publish_to_relay", fake_publish) + monkeypatch.setattr(analytics.asyncio, "sleep", fake_sleep) + + await analytics.publish_usage_analytics() + + assert len(published_events) == 1 + assert ["schema", analytics.ANALYTICS_SCHEMA] in published_events[0].get("tags", []) diff --git a/tests/unit/test_settings.py b/tests/unit/test_settings.py index 5c5cb048..770af976 100644 --- a/tests/unit/test_settings.py +++ b/tests/unit/test_settings.py @@ -2,6 +2,7 @@ import os import pytest from sqlalchemy.ext.asyncio import create_async_engine +from sqlmodel import text from sqlmodel.ext.asyncio.session import AsyncSession from routstr.core.settings import SettingsService @@ -11,6 +12,7 @@ from routstr.core.settings import SettingsService async def test_settings_seed_from_env_and_persist() -> None: os.environ["UPSTREAM_BASE_URL"] = "https://api.test/v1" os.environ.pop("ONION_URL", None) + os.environ.pop("ENABLE_ANALYTICS_SHARING", None) engine = create_async_engine("sqlite+aiosqlite:///:memory:") async with AsyncSession(engine, expire_on_commit=False) as session: @@ -19,19 +21,53 @@ async def test_settings_seed_from_env_and_persist() -> None: assert settings.upstream_base_url == "https://api.test/v1" # ONION_URL may be empty if not discoverable assert isinstance(settings.onion_url, str) + assert settings.enable_analytics_sharing is True @pytest.mark.asyncio async def test_settings_db_precedence_over_env() -> None: os.environ["UPSTREAM_BASE_URL"] = "https://api.env/v1" + os.environ["ENABLE_ANALYTICS_SHARING"] = "true" engine = create_async_engine("sqlite+aiosqlite:///:memory:") async with AsyncSession(engine, expire_on_commit=False) as session: _ = await SettingsService.initialize(session) - updated = await SettingsService.update({"name": "DBName"}, session) + updated = await SettingsService.update( + {"name": "DBName", "enable_analytics_sharing": False}, session + ) assert updated.name == "DBName" + assert updated.enable_analytics_sharing is False # Change env and re-initialize; DB should still win os.environ["NAME"] = "EnvName" + os.environ["ENABLE_ANALYTICS_SHARING"] = "true" again = await SettingsService.initialize(session) assert again.name == "DBName" + assert again.enable_analytics_sharing is False + + +@pytest.mark.asyncio +async def test_settings_initialize_discards_unknown_keys() -> None: + engine = create_async_engine("sqlite+aiosqlite:///:memory:") + async with AsyncSession(engine, expire_on_commit=False) as session: + _ = await SettingsService.initialize(session) + + # 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}' + ) + ) + await session.commit() + + reloaded = await SettingsService.initialize(session) + assert reloaded.name == "LegacyNode" + assert reloaded.enable_analytics_sharing is True + + row = await session.exec(text("SELECT data FROM settings WHERE id = 1")) # type: ignore + stored_data = row.first()[0] + assert '"enable_analytics_sharing": true' in stored_data + assert "nostr_analytics_enabled" not in stored_data + assert "unknown_key" not in stored_data diff --git a/ui/app/page.tsx b/ui/app/page.tsx index 380f465e..76addfd7 100644 --- a/ui/app/page.tsx +++ b/ui/app/page.tsx @@ -9,6 +9,7 @@ import type { DateRange } from 'react-day-picker'; import { UsageMetricsChart } from '@/components/usage-metrics-chart'; import { UsageSummaryCards } from '@/components/usage-summary-cards'; import { ErrorDetailsTable } from '@/components/error-details-table'; +import { TopModelsUsageChart } from '@/components/top-models-usage-chart'; import { DashboardBalanceSummary } from '@/components/dashboard-balance-summary'; import { AdminService, @@ -79,6 +80,7 @@ const TIME_RANGE_PRESETS = [ { value: '3m', label: 'Last 3 Months', hours: 90 * 24 }, { value: '12m', label: 'Last 12 Months', hours: 365 * 24 }, ] as const; +const MAX_USAGE_RANGE_HOURS = 365 * 24; type TimeRangePresetValue = (typeof TIME_RANGE_PRESETS)[number]['value']; @@ -159,20 +161,6 @@ function getAutoIntervalMinutes(hours: number): number { ); } -function getQueryErrorMessage(error: unknown): string { - if ( - error && - typeof error === 'object' && - 'message' in error && - typeof error.message === 'string' && - error.message.trim().length > 0 - ) { - return error.message; - } - - return 'The analytics request failed. Refresh and try again.'; -} - function SectionLoading({ label }: { label: string }) { if (label === 'summary') { return ( @@ -554,19 +542,21 @@ export default function DashboardPage() { isCustomRangeActive && customRangeHours ? customRangeHours : activePreset.hours; - const autoInterval = getAutoIntervalMinutes(queryHours); + const safeQueryHours = Math.min(queryHours, MAX_USAGE_RANGE_HOURS); + const isUsageRangeCapped = safeQueryHours < queryHours; + const autoInterval = getAutoIntervalMinutes(safeQueryHours); const usageRefetchIntervalMs = useMemo(() => { - if (queryHours > 90 * 24) { + if (safeQueryHours > 90 * 24) { return 4 * 60 * 60_000; } - if (queryHours > 30 * 24) { + if (safeQueryHours > 30 * 24) { return 2 * 60 * 60_000; } - if (queryHours > 7 * 24) { + if (safeQueryHours > 7 * 24) { return 30 * 60_000; } return 60_000; - }, [queryHours]); + }, [safeQueryHours]); const revenueDisplayUnit: DisplayUnit = useMemo(() => { if (displayUnit === 'usd' && usdPerSat === null) { // Keep revenue charts meaningful while the USD rate is unavailable. @@ -584,43 +574,30 @@ export default function DashboardPage() { : revenueDisplayUnit; const { - data: metricsData, - isLoading: metricsLoading, - error: metricsError, - refetch: refetchMetrics, + data: usageDashboardData, + isLoading: usageDashboardLoading, + refetch: refetchUsageDashboard, } = useQuery({ - queryKey: ['usage-metrics', autoInterval, queryHours], - queryFn: () => AdminService.getUsageMetrics(autoInterval, queryHours), + queryKey: ['usage-dashboard', autoInterval, safeQueryHours], + queryFn: () => + AdminService.getUsageDashboard(safeQueryHours, autoInterval, 100, 20), enabled: isAuthenticated, refetchInterval: usageRefetchIntervalMs, staleTime: 30_000, }); - const { - data: summaryData, - isLoading: summaryLoading, - error: summaryError, - refetch: refetchSummary, - } = useQuery({ - queryKey: ['usage-summary', queryHours], - queryFn: () => AdminService.getUsageSummary(queryHours), - enabled: isAuthenticated, - refetchInterval: usageRefetchIntervalMs, - staleTime: 30_000, - }); + const metricsData = usageDashboardData?.metrics; + const summaryData = usageDashboardData?.summary; + const errorData = usageDashboardData?.error_details; + const modelUsageMixData = usageDashboardData?.model_usage_mix; + const hasModelUsageMixMetrics = + Array.isArray(modelUsageMixData?.metrics) && + modelUsageMixData.metrics.length > 0; - const { - data: errorData, - isLoading: errorLoading, - error: errorDetailsError, - refetch: refetchErrors, - } = useQuery({ - queryKey: ['usage-errors', queryHours], - queryFn: () => AdminService.getErrorDetails(queryHours, 100), - enabled: isAuthenticated, - refetchInterval: usageRefetchIntervalMs, - staleTime: 30_000, - }); + const metricsLoading = usageDashboardLoading; + const summaryLoading = usageDashboardLoading; + const errorLoading = usageDashboardLoading; + const metricsTotals = metricsData?.totals; const chartConfigs = useMemo(() => { if (!metricsData || metricsData.metrics.length === 0) { @@ -647,12 +624,6 @@ export default function DashboardPage() { }) ) as ChartDatum[]; - const hasTokenMetrics = metricPoints.some((metric) => - ['input_tokens', 'output_tokens', 'total_tokens'].some( - (key) => typeof metric[key] === 'number' - ) - ); - return [ { id: 'revenue', @@ -661,6 +632,11 @@ export default function DashboardPage() { description: 'Track collected revenue trends over time.', data: revenuePoints, metricType: 'currency', + totals: metricsTotals + ? { + revenue_display: convertRevenueMsats(metricsTotals.revenue_msats), + } + : undefined, dataKeys: [ { key: 'revenue_display', @@ -676,6 +652,14 @@ export default function DashboardPage() { description: 'Understand traffic and completion reliability over time.', data: metricPoints, metricType: 'count', + totals: metricsTotals + ? { + total_requests: metricsTotals.total_requests, + successful_chat_completions: + metricsTotals.successful_chat_completions, + failed_requests: metricsTotals.failed_requests, + } + : undefined, dataKeys: [ { key: 'total_requests', @@ -701,6 +685,13 @@ export default function DashboardPage() { description: 'Monitor warnings, handled errors, and upstream failures.', data: metricPoints, metricType: 'count', + totals: metricsTotals + ? { + errors: metricsTotals.errors, + warnings: metricsTotals.warnings, + upstream_errors: metricsTotals.upstream_errors, + } + : undefined, dataKeys: [ { key: 'errors', @@ -726,6 +717,11 @@ export default function DashboardPage() { description: 'Follow payment processing activity by interval.', data: metricPoints, metricType: 'count', + totals: metricsTotals + ? { + payment_processed: metricsTotals.payment_processed, + } + : undefined, dataKeys: [ { key: 'payment_processed', @@ -734,38 +730,41 @@ export default function DashboardPage() { }, ], }, - ...(hasTokenMetrics - ? [ - { - id: 'tokens', - title: 'Token Usage', - mobileTitle: 'Tokens', - description: - 'Track input, output, and total token throughput over time.', - data: metricPoints, - metricType: 'count' as const, - dataKeys: [ - { - key: 'total_tokens', - name: 'Total Tokens', - color: 'var(--chart-1)', - }, - { - key: 'input_tokens', - name: 'Input Tokens', - color: 'var(--chart-2)', - }, - { - key: 'output_tokens', - name: 'Output Tokens', - color: 'var(--chart-3)', - }, - ], - }, - ] - : []), + { + id: 'tokens', + title: 'Token Usage', + mobileTitle: 'Tokens', + description: + 'Track input, output, and total token throughput over time.', + data: metricPoints, + metricType: 'count', + totals: metricsTotals + ? { + input_tokens: metricsTotals.input_tokens, + output_tokens: metricsTotals.output_tokens, + total_tokens: metricsTotals.total_tokens, + } + : undefined, + dataKeys: [ + { + key: 'total_tokens', + name: 'Total Tokens', + color: 'var(--chart-1)', + }, + { + key: 'input_tokens', + name: 'Input Tokens', + color: 'var(--chart-2)', + }, + { + key: 'output_tokens', + name: 'Output Tokens', + color: 'var(--chart-3)', + }, + ], + }, ]; - }, [metricsData, revenueDisplayUnit, usdPerSat]); + }, [metricsData, metricsTotals, revenueDisplayUnit, usdPerSat]); useEffect(() => { if (chartConfigs.length === 0) { @@ -803,11 +802,7 @@ export default function DashboardPage() { setIsManualRefreshing(true); try { - await Promise.allSettled([ - refetchMetrics(), - refetchSummary(), - refetchErrors(), - ]); + await refetchUsageDashboard(); } finally { setIsManualRefreshing(false); } @@ -913,10 +908,16 @@ export default function DashboardPage() { All cards and charts in this section update from the selected range.

+ {isUsageRangeCapped ? ( +

+ Usage analytics are capped to the last{' '} + {MAX_USAGE_RANGE_HOURS / 24} days for server safety. +

+ ) : null} -
-
+
+
- - {isManualRefreshing ? 'Refreshing...' : 'Refresh'} - + {isManualRefreshing ? 'Refreshing...' : 'Refresh'}
{metricsLoading ? ( - ) : metricsError ? ( - - - - - - - - Unable to load analytics - - {getQueryErrorMessage(metricsError)} - - - - - ) : activeChartConfig ? ( )} + {!metricsLoading && modelUsageMixData && hasModelUsageMixMetrics ? ( + + ) : null} + {summaryLoading ? ( - ) : summaryError ? ( - - - - - Usage summary unavailable - - {getQueryErrorMessage(summaryError)} - - - - - ) : summaryData ? ( ) : null} @@ -1074,19 +1051,6 @@ export default function DashboardPage() { {errorLoading ? ( - ) : errorDetailsError ? ( - - - - - Error details unavailable - - {getQueryErrorMessage(errorDetailsError)} - - - - - ) : errorData ? ( ) : null} diff --git a/ui/components/settings/admin-settings.tsx b/ui/components/settings/admin-settings.tsx index ceab6128..30e0f79d 100644 --- a/ui/components/settings/admin-settings.tsx +++ b/ui/components/settings/admin-settings.tsx @@ -26,6 +26,7 @@ interface SettingsData { description?: string; npub?: string; nsec?: string; + enable_analytics_sharing?: boolean; upstream_api_key?: string; http_url?: string; onion_url?: string; @@ -43,6 +44,7 @@ const HANDLED_KEYS = [ 'nsec', 'cashu_mints', 'relays', + 'enable_analytics_sharing', 'admin_password', 'id', 'updated_at', @@ -366,6 +368,7 @@ export function AdminSettings() { const nostrChanged = ['npub', 'nsec'].some(hasFieldChanged); const cashuMintsChanged = hasFieldChanged('cashu_mints'); const relaysChanged = hasFieldChanged('relays'); + const analyticsSharingChanged = hasFieldChanged('enable_analytics_sharing'); const advancedKeys = Object.keys(settings).filter( (key) => !HANDLED_KEYS.includes(key) && !IGNORED_KEYS.includes(key) ); @@ -397,6 +400,7 @@ export function AdminSettings() { resetFields(['relays']); setNewRelay(''); }; + const resetAnalyticsSharing = () => resetFields(['enable_analytics_sharing']); const resetAdvanced = () => resetFields(advancedKeys); if (loading) { @@ -686,6 +690,52 @@ export function AdminSettings() { ) : null} + {/* Analytics Sharing */} + + + Analytics Sharing + + Publish aggregate usage stats to Nostr for external dashboards + + + +
+
+ +

+ When enabled, Routstr periodically publishes aggregate model + usage and revenue stats. +

+
+ + handleInputChange('enable_analytics_sharing', checked) + } + /> +
+
+ {analyticsSharingChanged ? ( + +
+ + +
+
+ ) : null} +
+ {/* Other Settings */} diff --git a/ui/components/top-models-usage-chart.tsx b/ui/components/top-models-usage-chart.tsx new file mode 100644 index 00000000..f6a64110 --- /dev/null +++ b/ui/components/top-models-usage-chart.tsx @@ -0,0 +1,872 @@ +'use client'; + +import { useEffect, useMemo, useRef, useState } from 'react'; +import { ExpandIcon, Minimize2Icon } from 'lucide-react'; +import { Bar, BarChart, CartesianGrid, XAxis, YAxis } from 'recharts'; +import { Card, CardContent, CardHeader, CardTitle } from '@/components/ui/card'; +import { Button } from '@/components/ui/button'; +import { + ChartConfig, + ChartContainer, + ChartTooltip, +} from '@/components/ui/chart'; +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'; + +interface TopModelsUsageChartProps { + mix: ModelUsageMix; + displayUnit: DisplayUnit; + usdPerSat: number | null; +} + +type ChartMode = 'requests' | 'revenue' | 'tokens'; + +interface TooltipRow { + color: string; + dataKey: string; + label: string; + value: number; +} + +type LeaderboardTrend = 'up' | 'down' | 'flat' | 'new'; + +interface LeaderboardRow { + chartDataKey: string | null; + displayName: string; + model: string; + provider: string; + rank: number; + totalRaw: number; + trend: LeaderboardTrend; + 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) { + hash = (hash << 5) - hash + input.charCodeAt(i); + hash |= 0; + } + return Math.abs(hash) % 360; +} + +function getSeriesColor(model: string, index: number): string { + const palette = [ + 'var(--chart-1)', + 'var(--chart-2)', + 'var(--chart-3)', + 'var(--chart-4)', + 'var(--chart-5)', + '#f59e0b', + '#06b6d4', + '#8b5cf6', + '#f97316', + '#34d399', + ]; + + if (index < palette.length) { + return palette[index]; + } + + const hue = (hueFromString(model) + index * 23) % 360; + return `hsl(${hue} 70% 56%)`; +} + +function formatTooltipTimestamp( + label: string, + intervalMinutes: number, + hoursBack: number +): string { + const date = parseBucketDate(label); + if (!date) { + return label; + } + const shouldShowTime = intervalMinutes <= 6 * 60 || hoursBack <= 48; + if (shouldShowTime) { + return date.toLocaleString([], { + month: 'long', + day: 'numeric', + year: 'numeric', + hour: '2-digit', + minute: '2-digit', + }); + } + return date.toLocaleString([], { + month: 'long', + day: 'numeric', + year: 'numeric', + }); +} + +function formatAxisTimestamp( + timestamp: string, + hasMultipleDays: boolean, + intervalMinutes: number, + hoursBack: number +): string { + const date = parseBucketDate(timestamp); + if (!date) { + return ''; + } + + const shouldShowTime = intervalMinutes <= 6 * 60 || hoursBack <= 48; + if (shouldShowTime && hasMultipleDays) { + return date.toLocaleString([], { + month: 'short', + day: 'numeric', + hour: '2-digit', + minute: '2-digit', + }); + } + + if (shouldShowTime) { + return date.toLocaleTimeString([], { + hour: '2-digit', + minute: '2-digit', + }); + } + + if (intervalMinutes >= 24 * 60 && hoursBack >= 24 * 180) { + return date.toLocaleDateString([], { + month: 'short', + year: '2-digit', + }); + } + + if (hasMultipleDays) { + return date.toLocaleDateString([], { + month: 'short', + day: 'numeric', + }); + } + + return date.toLocaleTimeString([], { + hour: '2-digit', + minute: '2-digit', + }); +} + +function convertRevenueMsats( + amountMsats: number, + displayUnit: DisplayUnit, + usdPerSat: number | null +): number { + if (displayUnit === 'msat') { + return amountMsats; + } + + const sats = amountMsats / 1000; + if (displayUnit === 'usd') { + return sats * (usdPerSat ?? 0); + } + + return sats; +} + +function prettifyProvider(provider: string): string { + const normalized = provider.trim().toLowerCase(); + const aliasMap: Record = { + 'x ai': 'x-ai', + xai: 'x-ai', + 'z ai': 'z-ai', + zai: 'z-ai', + open_ai: 'openai', + openai: 'openai', + }; + if (aliasMap[normalized]) { + return aliasMap[normalized]; + } + return normalized.replace(/[_-]+/g, ' '); +} + +function detectProviderFromModel(model: string): string { + const value = model.toLowerCase(); + if (value.includes('claude')) return 'anthropic'; + if (value.includes('gpt') || value.includes('openai')) return 'openai'; + if (value.includes('gemini')) return 'google'; + if ( + value.includes('grok') || + value.includes('x-ai') || + value.includes('xai') + ) { + return 'x-ai'; + } + if (value.includes('deepseek')) return 'deepseek'; + if (value.includes('minimax')) return 'minimax'; + if (value.includes('kimi') || value.includes('moonshot')) return 'moonshot'; + if (value.includes('mistral')) return 'mistral'; + if (value.includes('qwen') || value.includes('alibaba')) return 'alibaba'; + if ( + value.includes('glm') || + value.includes('z-ai') || + value.includes('z ai') + ) { + return 'z-ai'; + } + return 'unknown'; +} + +function getModelPresentation(model: string): { + displayName: string; + provider: string; +} { + const trimmed = model.trim(); + const slashIndex = trimmed.indexOf('/'); + if (slashIndex > 0 && slashIndex < trimmed.length - 1) { + const provider = prettifyProvider(trimmed.slice(0, slashIndex)); + const displayName = trimmed.slice(slashIndex + 1); + return { displayName, provider }; + } + + return { + displayName: trimmed, + provider: detectProviderFromModel(trimmed), + }; +} + +export function TopModelsUsageChart({ + mix, + displayUnit, + usdPerSat, +}: TopModelsUsageChartProps) { + const [mode, setMode] = useState('requests'); + const [hoveredSeriesKey, setHoveredSeriesKey] = useState(null); + const [isChartPointerInside, setIsChartPointerInside] = useState(false); + const [isFullscreen, setIsFullscreen] = useState(false); + const isMobile = useIsMobile(); + const containerRef = useRef(null); + const compactNumber = useMemo( + () => + new Intl.NumberFormat('en-US', { + notation: 'compact', + maximumFractionDigits: 2, + }), + [] + ); + const mixTopModels = useMemo( + () => (Array.isArray(mix.top_models) ? mix.top_models : []), + [mix.top_models] + ); + const mixMetrics = useMemo( + () => (Array.isArray(mix.metrics) ? mix.metrics : []), + [mix.metrics] + ); + + const chartModels = useMemo(() => mixTopModels.slice(0, 20), [mixTopModels]); + const leaderboardModels = useMemo( + () => mixTopModels.slice(0, 20), + [mixTopModels] + ); + const revenueDisplayUnit: DisplayUnit = useMemo(() => { + if (displayUnit === 'usd' && usdPerSat === null) { + return 'sat'; + } + return displayUnit; + }, [displayUnit, usdPerSat]); + const revenueUnitLabel = + revenueDisplayUnit === 'usd' + ? 'USD' + : revenueDisplayUnit === 'sat' + ? 'sats' + : revenueDisplayUnit === 'msat' + ? 'msats' + : revenueDisplayUnit; + + const series = useMemo( + () => + chartModels.map((model, index) => ({ + requestsKey: `model_req_${index}`, + revenueKey: `model_rev_${index}`, + tokensKey: `model_tok_${index}`, + label: model, + color: getSeriesColor(model, index), + })), + [chartModels] + ); + + const chartData = useMemo( + () => + mixMetrics.map((metric) => { + const modelCounts = metric.model_counts ?? {}; + const modelRevenue = metric.model_revenue_msats ?? {}; + const modelTokens = metric.model_tokens ?? {}; + const point: Record = { + timestamp: metric.timestamp, + total_successful: metric.total_successful, + total_revenue_msats: metric.total_revenue_msats, + total_tokens: metric.total_tokens, + others_requests: metric.others, + others_revenue_msats: metric.others_revenue_msats, + others_tokens: metric.others_tokens, + }; + + 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; + } + + return point; + }), + [mixMetrics, series] + ); + + const hasMultipleDays = useMemo(() => { + const daySet = new Set( + chartData.map((item) => + parseBucketDate(String(item.timestamp))?.toDateString() + ) + ); + return daySet.size > 1; + }, [chartData]); + + const chartConfig = useMemo(() => { + const config: ChartConfig = {}; + for (const item of series) { + config[item.requestsKey] = { + label: item.label, + color: item.color, + }; + config[item.revenueKey] = { + label: item.label, + color: item.color, + }; + config[item.tokensKey] = { + label: item.label, + color: item.color, + }; + } + config.others_requests = { + label: 'Others', + color: '#6b7280', + }; + config.others_revenue_msats = { + label: 'Others', + color: '#6b7280', + }; + config.others_tokens = { + label: 'Others', + color: '#6b7280', + }; + return config; + }, [series]); + + useEffect(() => { + setHoveredSeriesKey(null); + setIsChartPointerInside(false); + }, [mode]); + + useEffect(() => { + const handleFullscreenChange = () => { + setIsFullscreen(document.fullscreenElement === containerRef.current); + }; + + document.addEventListener('fullscreenchange', handleFullscreenChange); + + return () => { + document.removeEventListener('fullscreenchange', handleFullscreenChange); + }; + }, []); + + const toggleFullscreen = async () => { + if (!containerRef.current) { + return; + } + + try { + if (document.fullscreenElement === containerRef.current) { + await document.exitFullscreen(); + } else { + await containerRef.current.requestFullscreen(); + } + } catch (error) { + console.error('Failed to toggle top models chart fullscreen', error); + } + }; + + const formatValue = (rawValue: number): string => { + if (mode === 'requests') { + return compactNumber.format(rawValue); + } + + if (mode === 'tokens') { + return compactNumber.format(rawValue); + } + + const converted = convertRevenueMsats( + rawValue, + revenueDisplayUnit, + usdPerSat + ); + const compact = compactNumber.format(converted); + if (revenueDisplayUnit === 'usd') { + return `$${compact}`; + } + return `${compact} ${revenueUnitLabel}`; + }; + + const activeSeries = series.map((item) => ({ + dataKey: + mode === 'requests' + ? item.requestsKey + : mode === 'revenue' + ? item.revenueKey + : item.tokensKey, + name: item.label, + color: item.color, + })); + const othersKey = ( + mode === 'requests' + ? 'others_requests' + : mode === 'revenue' + ? 'others_revenue_msats' + : 'others_tokens' + ) as 'others_requests' | 'others_revenue_msats' | 'others_tokens'; + const activeSeriesKeys = [ + ...activeSeries.map((item) => item.dataKey), + othersKey, + ]; + const activeHoverSeriesKey = + hoveredSeriesKey && activeSeriesKeys.includes(hoveredSeriesKey) + ? hoveredSeriesKey + : null; + const getSeriesOpacity = (dataKey: string): number => + activeHoverSeriesKey && activeHoverSeriesKey !== dataKey ? 0.18 : 1; + const formatLeaderboardTotal = (rawValue: number): string => { + if (mode === 'requests') { + return `${compactNumber.format(rawValue)} requests`; + } + + if (mode === 'tokens') { + return `${compactNumber.format(rawValue)} tokens`; + } + + const converted = convertRevenueMsats( + rawValue, + revenueDisplayUnit, + usdPerSat + ); + const compact = compactNumber.format(converted); + if (revenueDisplayUnit === 'usd') { + return `$${compact}`; + } + return `${compact} ${revenueUnitLabel}`; + }; + const formatTrendPercent = (value: number): string => { + const abs = Math.abs(value); + const rounded = abs >= 10 ? abs.toFixed(0) : abs.toFixed(1); + return rounded.replace(/\.0$/, ''); + }; + const leaderboardRows = useMemo(() => { + if (leaderboardModels.length === 0 || mixMetrics.length === 0) { + return []; + } + + const windowSize = Math.floor(mixMetrics.length / 2); + const previousMetrics = + windowSize > 0 ? mixMetrics.slice(-windowSize * 2, -windowSize) : []; + const currentMetrics = + windowSize > 0 ? mixMetrics.slice(-windowSize) : mixMetrics; + + const rows = leaderboardModels + .map((model) => { + const readMetric = (metric: (typeof mixMetrics)[number]): number => + mode === 'requests' + ? ((metric.model_counts ?? {})[model] ?? 0) + : mode === 'revenue' + ? ((metric.model_revenue_msats ?? {})[model] ?? 0) + : ((metric.model_tokens ?? {})[model] ?? 0); + + const totalRaw = mixMetrics.reduce( + (sum, metric) => sum + readMetric(metric), + 0 + ); + const previousRaw = previousMetrics.reduce( + (sum, metric) => sum + readMetric(metric), + 0 + ); + const currentRaw = currentMetrics.reduce( + (sum, metric) => sum + readMetric(metric), + 0 + ); + const trendPercent = + previousRaw > 0 + ? ((currentRaw - previousRaw) / previousRaw) * 100 + : null; + + let trend: LeaderboardTrend = 'flat'; + if (previousRaw <= 0 && currentRaw > 0) { + trend = 'new'; + } else if (trendPercent !== null && trendPercent > 0.5) { + trend = 'up'; + } else if (trendPercent !== null && trendPercent < -0.5) { + trend = 'down'; + } + + const presentation = getModelPresentation(model); + const matchingSeries = series.find((item) => item.label === model); + const chartDataKey = matchingSeries + ? mode === 'requests' + ? matchingSeries.requestsKey + : mode === 'revenue' + ? matchingSeries.revenueKey + : matchingSeries.tokensKey + : null; + + return { + chartDataKey, + displayName: presentation.displayName, + model, + provider: presentation.provider, + rank: 0, + totalRaw, + trend, + trendPercent, + } satisfies LeaderboardRow; + }) + .filter((row) => row.totalRaw > 0) + .sort((a, b) => b.totalRaw - a.totalRaw) + .slice(0, 20) + .map((row, index) => ({ + ...row, + rank: index + 1, + })); + + return rows; + }, [leaderboardModels, mixMetrics, mode, series]); + + if (chartData.length === 0) { + return null; + } + + return ( +
+ + +
+
+ + Model Usage + +

+ Stacked requests, revenue, or tokens by model ( + {mix.interval_minutes}m buckets). +

+
+
+
+ + + +
+ +
+
+
+ + { + setHoveredSeriesKey(null); + setIsChartPointerInside(false); + }} + > + setIsChartPointerInside(true)} + onMouseMove={() => setIsChartPointerInside(true)} + onMouseLeave={() => { + setHoveredSeriesKey(null); + setIsChartPointerInside(false); + }} + margin={{ + top: 12, + right: isMobile ? 8 : 18, + left: isMobile ? 0 : 8, + bottom: 0, + }} + > + + + formatAxisTimestamp( + String(value), + hasMultipleDays, + mix.interval_minutes, + mix.hours_back + ) + } + /> + + formatValue( + typeof value === 'number' ? value : Number(value || 0) + ) + } + /> + { + if (!isChartPointerInside || !active || !payload?.length) { + return null; + } + + const rows = payload + .map((entry) => { + const value = + typeof entry.value === 'number' + ? entry.value + : Number(entry.value || 0); + + return { + color: String(entry.color || '#6b7280'), + dataKey: String(entry.dataKey || ''), + label: String(entry.name || ''), + value, + } satisfies TooltipRow; + }) + .filter( + (row) => Number.isFinite(row.value) && row.value > 0 + ) + .sort((a, b) => b.value - a.value); + + const total = rows.reduce((sum, row) => sum + row.value, 0); + if (rows.length === 0) { + return null; + } + + return ( +
+

+ {formatTooltipTimestamp( + String(label || ''), + mix.interval_minutes, + mix.hours_back + )} +

+
+ {rows.map((row) => ( +
+ + + {row.label} + + + {formatValue(row.value)} + +
+ ))} +
+
+
+ Total + + {formatValue(total)} + +
+
+
+ ); + }} + /> + {activeSeries.map((item) => ( + setHoveredSeriesKey(item.dataKey)} + onMouseLeave={() => setHoveredSeriesKey(null)} + /> + ))} + setHoveredSeriesKey(othersKey)} + onMouseLeave={() => setHoveredSeriesKey(null)} + /> +
+
+ +
+
+

+ Top models +

+

+ Change vs prior period +

+
+ + {leaderboardRows.length > 0 ? ( +
+ {leaderboardRows.map((row) => { + const rowIsLinked = Boolean(row.chartDataKey); + const rowIsActive = + row.chartDataKey !== null && + activeHoverSeriesKey === row.chartDataKey; + const rowIsDimmed = + Boolean(activeHoverSeriesKey) && + row.chartDataKey !== null && + row.chartDataKey !== activeHoverSeriesKey; + + let trendLabel = '0%'; + let trendClass = 'text-muted-foreground'; + if (row.trend === 'new') { + trendLabel = 'new'; + trendClass = 'text-blue-500'; + } else if (row.trend === 'up' && row.trendPercent !== null) { + trendLabel = `↑${formatTrendPercent(row.trendPercent)}%`; + trendClass = 'text-emerald-500'; + } else if ( + row.trend === 'down' && + row.trendPercent !== null + ) { + trendLabel = `↓${formatTrendPercent(row.trendPercent)}%`; + trendClass = 'text-red-500'; + } else if (row.trendPercent !== null) { + trendLabel = `${formatTrendPercent(row.trendPercent)}%`; + } + + return ( +
{ + if (row.chartDataKey) { + setHoveredSeriesKey(row.chartDataKey); + } + }} + onMouseLeave={() => { + if (row.chartDataKey) { + setHoveredSeriesKey(null); + } + }} + > + + {row.rank}. + +
+ + {row.displayName} + {' '} + + by {row.provider} + +
+ + {formatLeaderboardTotal(row.totalRaw)} + + + {trendLabel} + +
+ ); + })} +
+ ) : ( +

+ No model totals available for this range. +

+ )} +
+
+
+
+ ); +} diff --git a/ui/components/usage-summary-cards.tsx b/ui/components/usage-summary-cards.tsx index e578748a..7db46735 100644 --- a/ui/components/usage-summary-cards.tsx +++ b/ui/components/usage-summary-cards.tsx @@ -35,9 +35,10 @@ export function UsageSummaryCards({ summary }: UsageSummaryCardsProps) { const formatAmount = (msat: number) => formatFromMsat(msat, displayUnit, usdPerSat); - const hasTokenStats = - typeof summary.total_tokens === 'number' || - typeof summary.avg_total_tokens_per_completion === 'number'; + const totalTokens = Number(summary.total_tokens ?? 0); + const avgTotalTokensPerCompletion = Number( + summary.avg_total_tokens_per_completion ?? 0 + ); const cards = [ { @@ -52,26 +53,20 @@ export function UsageSummaryCards({ summary }: UsageSummaryCardsProps) { icon: CheckCircle2, iconClassName: 'text-emerald-600 dark:text-emerald-300', }, - ...(hasTokenStats - ? [ - { - title: 'Total Tokens', - value: Number(summary.total_tokens ?? 0).toLocaleString(), - icon: Database, - iconClassName: 'text-cyan-600 dark:text-cyan-300', - }, - { - title: 'Avg Tokens/Completion', - value: Number( - summary.avg_total_tokens_per_completion ?? 0 - ).toLocaleString(undefined, { - maximumFractionDigits: 1, - }), - icon: Activity, - iconClassName: 'text-indigo-600 dark:text-indigo-300', - }, - ] - : []), + { + title: 'Total Tokens', + value: totalTokens.toLocaleString(), + icon: Database, + iconClassName: 'text-cyan-600 dark:text-cyan-300', + }, + { + title: 'Avg Tokens/Completion', + value: avgTotalTokensPerCompletion.toLocaleString(undefined, { + maximumFractionDigits: 1, + }), + icon: Activity, + iconClassName: 'text-indigo-600 dark:text-indigo-300', + }, { title: 'Revenue', value: formatAmount(summary.revenue_msats), diff --git a/ui/lib/api/services/admin.ts b/ui/lib/api/services/admin.ts index 8ed8e3d7..e3a1676a 100644 --- a/ui/lib/api/services/admin.ts +++ b/ui/lib/api/services/admin.ts @@ -845,6 +845,23 @@ export class AdminService { ); } + static async getUsageDashboard( + hours: number = 24, + interval: number = 15, + errorLimit: number = 100, + modelLimit: number = 20 + ): 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)); + + return await apiClient.get( + `/admin/api/usage/dashboard?${params.toString()}` + ); + } + static async getUsageSummary(hours: number = 24): Promise { return await apiClient.get( `/admin/api/usage/summary?hours=${hours}` @@ -969,9 +986,9 @@ export interface UsageMetricData { upstream_errors: number; revenue_msats: number; refunds_msats: number; - input_tokens?: number; - output_tokens?: number; - total_tokens?: number; + input_tokens: number; + output_tokens: number; + total_tokens: number; [key: string]: unknown; } @@ -980,7 +997,20 @@ export interface UsageMetrics { interval_minutes: number; hours_back: number; total_buckets: number; - totals?: Partial>; + totals?: { + total_requests: number; + successful_chat_completions: number; + failed_requests: number; + errors: number; + warnings: number; + payment_processed: number; + upstream_errors: number; + revenue_msats: number; + refunds_msats: number; + input_tokens: number; + output_tokens: number; + total_tokens: number; + }; } export interface UsageSummary { @@ -995,6 +1025,12 @@ export interface UsageSummary { unique_models_count: number; unique_models: string[]; error_types: Record; + input_tokens: number; + output_tokens: number; + total_tokens: number; + avg_input_tokens_per_completion: number; + avg_output_tokens_per_completion: number; + avg_total_tokens_per_completion: number; success_rate: number; revenue_msats: number; refunds_msats: number; @@ -1004,8 +1040,6 @@ export interface UsageSummary { net_revenue_sats: number; avg_revenue_per_request_msats: number; refund_rate: number; - total_tokens?: number; - avg_total_tokens_per_completion?: number; } export interface ErrorDetail { @@ -1039,6 +1073,35 @@ export interface RevenueByModel { total_models: number; } +export interface ModelUsageMixMetric { + timestamp: string; + total_successful: number; + total_revenue_msats: number; + total_tokens: number; + others: number; + others_revenue_msats: number; + others_tokens: number; + model_counts: Record; + model_revenue_msats: Record; + model_tokens: Record; +} + +export interface ModelUsageMix { + top_models: string[]; + metrics: ModelUsageMixMetric[]; + interval_minutes: number; + hours_back: number; + total_buckets: number; +} + +export interface UsageDashboardResponse { + metrics: UsageMetrics; + summary: UsageSummary; + error_details: ErrorDetails; + revenue_by_model: RevenueByModel; + model_usage_mix?: ModelUsageMix; +} + export interface LogEntry { asctime: string; name: string;