diff --git a/routstr/auth.py b/routstr/auth.py index 07d83dff..62b6c5b8 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 ef363d2c..cdfb2e86 100644 --- a/routstr/core/admin.py +++ b/routstr/core/admin.py @@ -28,6 +28,8 @@ admin_router = APIRouter(prefix="/admin", include_in_schema=False) admin_sessions: dict[str, int] = {} ADMIN_SESSION_DURATION = 12 * 60 * 60 +# Usage analytics remain queryable up to 12 months. +MAX_USAGE_ANALYTICS_HOURS = 365 * 24 def _current_timestamp() -> int: @@ -971,6 +973,7 @@ async def get_usage_metrics( hours: int = Query( default=24, ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, description="Hours of history to analyze", ), ) -> dict: @@ -978,12 +981,44 @@ async def get_usage_metrics( 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, + le=MAX_USAGE_ANALYTICS_HOURS, description="Hours of history to analyze", ), ) -> dict: @@ -997,6 +1032,7 @@ async def get_error_details( hours: int = Query( default=24, ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, description="Hours of history to analyze", ), limit: int = Query( @@ -1015,6 +1051,7 @@ async def get_revenue_by_model( hours: int = Query( default=24, ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, description="Hours of history to analyze", ), limit: int = Query( diff --git a/routstr/core/log_manager.py b/routstr/core/log_manager.py index e3fd52fd..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, @@ -34,6 +90,7 @@ class LogManager: log_files = [] cutoff_date = None + cutoff_timestamp_str: str | None = None if specific_date: log_file = self.logs_dir / f"app_{specific_date}.log" @@ -47,6 +104,7 @@ class LogManager: # If we only care about hours back, we can optimize file selection if hours_back is not None: cutoff_date = datetime.now(timezone.utc) - timedelta(hours=hours_back) + cutoff_timestamp_str = cutoff_date.strftime("%Y-%m-%d %H:%M:%S") filtered_files = [] for log_path in log_files: try: @@ -69,27 +127,20 @@ class LogManager: for log_file in log_files: try: with open(log_file, "r") as f: - # For reverse search, we might want to read lines in reverse? - # But usually logs are append-only. - # If reverse_files is True, we iterate files newest to oldest. - # But lines within file are still oldest to newest unless we reverse them. - lines = f.readlines() - if reverse_files: - lines.reverse() + lines_iter = reversed(f.readlines()) if reverse_files else f - for line in lines: + for line in lines_iter: try: entry = json.loads(line.strip()) - if cutoff_date: + if cutoff_timestamp_str: timestamp_str = entry.get("asctime", "") - if not timestamp_str: + if ( + not isinstance(timestamp_str, str) + or len(timestamp_str) != 19 + ): continue - log_time = datetime.strptime( - timestamp_str, "%Y-%m-%d %H:%M:%S" - ) - log_time = log_time.replace(tzinfo=timezone.utc) - if log_time < cutoff_date: + if timestamp_str < cutoff_timestamp_str: continue yield entry @@ -166,7 +217,7 @@ class LogManager: methods: list[str] | None = None, endpoints: list[str] | None = None, ) -> bool: - if level and log_data.get("levelname", "").upper() != level.upper(): + if level and str(log_data.get("levelname", "")).upper() != level.upper(): return False if request_id and log_data.get("request_id") != request_id: @@ -216,207 +267,279 @@ class LogManager: return True + def _bucket_key_for_timestamp( + self, timestamp_str: str, interval_minutes: int + ) -> str | None: + if len(timestamp_str) != 19: + return None + if timestamp_str[10] != " ": + return None + + try: + hour = int(timestamp_str[11:13]) + minute = int(timestamp_str[14:16]) + except (TypeError, ValueError): + return None + + total_minutes = hour * 60 + minute + rounded_minutes = (total_minutes // interval_minutes) * interval_minutes + rounded_hour = rounded_minutes // 60 + rounded_minute = rounded_minutes % 60 + return f"{timestamp_str[:10]} {rounded_hour:02d}:{rounded_minute:02d}:00" + + def _extract_success_metrics( + self, entry: dict[str, Any], message: str + ) -> tuple[bool, float, int, int]: + # Use auth settlement logs as the canonical successful request signal. + 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 get_usage_summary(self, hours: int = 24) -> dict: - entries = list(self._yield_log_entries(hours_back=hours)) - 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)) - return self._aggregate_metrics_by_time(entries, interval, hours) - - def get_error_details(self, hours: int = 24, limit: int = 100) -> dict: - errors: list[dict] = [] - # Iterate newest to oldest for errors? - # yield_log_entries sorts files by name (date) ascending by default. - # usage stats logic usually expects ascending time for aggregation (though dictionaries don't care). - # For error details "last N errors", we probably want newest first. - - # Using list() loads everything into memory, which is what PR 229 did. - # For optimization, we could use reverse iterator. - - # Let's just stick to PR 229 logic which filters 'ERROR' level. - - entries = self._yield_log_entries(hours_back=hours) # oldest to newest - - for entry in entries: - if 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", ""), - } + 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)) - # Sort reverse time - errors.sort(key=lambda x: x["timestamp"], reverse=True) - return {"errors": errors[:limit], "total_count": len(errors)} - - def get_revenue_by_model(self, hours: int = 24, limit: int = 20) -> dict: - entries = list(self._yield_log_entries(hours_back=hours)) - - model_stats: dict[str, dict[str, int | float]] = defaultdict( - lambda: { - "revenue_msats": 0, - "refunds_msats": 0, - "requests": 0, - "successful": 0, - "failed": 0, - } + return self._cache_call( + ("usage_summary", hours), + compute, ) - for entry in entries: + def get_usage_metrics(self, interval: int = 15, hours: int = 24) -> dict: + def compute() -> dict: try: - model = entry.get("model", "unknown") - if not isinstance(model, str): - model = "unknown" + 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 + ) - message = entry.get("message", "").lower() + return self._cache_call( + ("usage_metrics", interval, hours), + compute, + ) - if "received proxy request" in message: - model_stats[model]["requests"] += 1 + 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 - if ( - "completed for streaming" in message - or "completed for non-streaming" in message - ): - model_stats[model]["successful"] += 1 - cost_data = entry.get("cost_data") - if isinstance(cost_data, dict): - actual_cost = cost_data.get("total_msats", 0) - if isinstance(actual_cost, (int, float)) and actual_cost > 0: - model_stats[model]["revenue_msats"] += actual_cost + 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, + ) - if "revert payment" in message or "upstream request failed" in message: - 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 + return self._cache_call( + ("usage_dashboard", interval, hours, error_limit, model_limit), + compute, + ttl_seconds=cache_ttl, + ) - except Exception: - continue + def get_error_details(self, hours: int = 24, limit: int = 100) -> dict: + 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}" + ) - models: list[dict[str, Any]] = [] - total_revenue = 0.0 + 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", ""), + } + ) - for model, stats in model_stats.items(): - revenue_msats = float(stats["revenue_msats"]) - refunds_msats = float(stats["refunds_msats"]) + errors.sort(key=lambda x: x["timestamp"], reverse=True) + return {"errors": errors[:limit], "total_count": len(errors)} - revenue_sats = revenue_msats / 1000 - refunds_sats = refunds_msats / 1000 - net_revenue_sats = revenue_sats - refunds_sats + return self._cache_call(("error_details", hours, limit), compute) - total_revenue += net_revenue_sats + def get_revenue_by_model(self, hours: int = 24, limit: int = 20) -> dict: + def compute() -> dict: + try: + return self._usage_store.get_revenue_by_model( + hours_back=hours, limit=limit + ) + except Exception as e: + logger.error( + f"Usage analytics index failed, falling back to log scan: {e}" + ) - requests = int(stats["requests"]) - successful = int(stats["successful"]) + entries = self._get_cached_entries(hours) - 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() - def _calculate_summary_stats(self, entries: list[dict]) -> dict: - 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, - } + 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 - for entry in entries: - try: - stats["total_entries"] += 1 + 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 - message = entry.get("message", "").lower() - level = entry.get("levelname", "").upper() + except Exception: + continue - if level == "ERROR": - stats["total_errors"] += 1 - if "error_type" in entry: - stats["error_types"][str(entry["error_type"])] += 1 - elif level == "WARNING": - stats["total_warnings"] += 1 + models: list[dict[str, Any]] = [] + total_revenue = 0.0 - if "received proxy request" in message: - stats["total_requests"] += 1 + for model, stats in model_stats.items(): + revenue_msats = float(stats["revenue_msats"]) + refunds_msats = float(stats["refunds_msats"]) - if ( - "completed for streaming" in message - or "completed for non-streaming" in message - ): - stats["successful_chat_completions"] += 1 + revenue_sats = revenue_msats / 1000 + refunds_sats = refunds_msats / 1000 + net_revenue_sats = revenue_sats - refunds_sats - if "upstream request failed" in message or "revert payment" in message: - stats["failed_requests"] += 1 + total_revenue += net_revenue_sats - if "payment processed successfully" in message: - stats["payment_processed"] += 1 + requests = int(stats["requests"]) + successful = int(stats["successful"]) - if "upstream" in message and level == "ERROR": - stats["upstream_errors"] += 1 + 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 + ), + } + ) - if "model" in entry: - model = entry["model"] - if isinstance(model, str) and model != "unknown": - stats["unique_models"].add(model) + models.sort(key=lambda x: float(x["net_revenue_sats"]), reverse=True) - if ( - "completed for streaming" in message - or "completed for non-streaming" in message - ): - cost_data = entry.get("cost_data") - if isinstance(cost_data, dict): - actual_cost = cost_data.get("total_msats", 0) - if isinstance(actual_cost, (int, float)) and actual_cost > 0: - stats["revenue_msats"] += float(actual_cost) + return { + "models": models[:limit], + "total_revenue_sats": total_revenue, + "total_models": len(models), + } - if "revert payment" in message: - max_cost = entry.get("max_cost_for_model", 0) - if isinstance(max_cost, (int, float)) and max_cost > 0: - stats["refunds_msats"] += float(max_cost) - - except Exception: - continue + 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 refunds_sats = stats["refunds_msats"] / 1000 net_revenue_sats = revenue_sats - refunds_sats total_requests = stats["total_requests"] successful = stats["successful_chat_completions"] + input_tokens = stats["input_tokens"] + output_tokens = stats["output_tokens"] + total_tokens = stats["total_tokens"] return { "total_entries": stats["total_entries"], @@ -430,6 +553,18 @@ class LogManager: "unique_models_count": len(stats["unique_models"]), "unique_models": sorted(list(stats["unique_models"])), "error_types": dict(stats["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, @@ -449,6 +584,419 @@ class LogManager: ), } + def _calculate_summary_stats(self, entries: list[dict]) -> dict: + 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, + } + + for entry in entries: + try: + stats["total_entries"] += 1 + + message = str(entry.get("message", "")).lower() + level = str(entry.get("levelname", "")).upper() + + if level == "ERROR": + stats["total_errors"] += 1 + if "error_type" in entry: + stats["error_types"][str(entry["error_type"])] += 1 + elif level == "WARNING": + stats["total_warnings"] += 1 + + completed, revenue_msats, input_tokens, output_tokens = ( + self._extract_success_metrics(entry, message) + ) + if completed: + stats["total_requests"] += 1 + stats["successful_chat_completions"] += 1 + stats["input_tokens"] += input_tokens + stats["output_tokens"] += output_tokens + stats["total_tokens"] += input_tokens + output_tokens + + failed = ( + "upstream request failed" in message + or "revert payment" in message + ) + if failed: + stats["total_requests"] += 1 + stats["failed_requests"] += 1 + + if "payment processed successfully" in message: + stats["payment_processed"] += 1 + + if "upstream" in message and level == "ERROR": + stats["upstream_errors"] += 1 + + if "model" in entry: + model = entry["model"] + if isinstance(model, str) and model != "unknown": + stats["unique_models"].add(model) + + if completed and revenue_msats > 0: + stats["revenue_msats"] += revenue_msats + + if "revert payment" in message: + max_cost = entry.get("max_cost_for_model", 0) + if isinstance(max_cost, (int, float)) and max_cost > 0: + stats["refunds_msats"] += float(max_cost) + + except Exception: + continue + + 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: @@ -463,39 +1011,37 @@ class LogManager: "upstream_errors": 0, "revenue_msats": 0.0, "refunds_msats": 0.0, + "input_tokens": 0, + "output_tokens": 0, + "total_tokens": 0, } ) for entry in entries: try: timestamp_str = entry.get("asctime", "") - if not timestamp_str: + if not isinstance(timestamp_str, str): continue - - log_time = datetime.strptime(timestamp_str, "%Y-%m-%d %H:%M:%S") - log_time = log_time.replace(tzinfo=timezone.utc) - - # Round down to nearest interval - minutes = log_time.minute - rounded_minutes = (minutes // interval_minutes) * interval_minutes - bucket_time = log_time.replace( - minute=rounded_minutes, second=0, microsecond=0 + bucket_key = self._bucket_key_for_timestamp( + timestamp_str, interval_minutes ) - bucket_key = bucket_time.strftime("%Y-%m-%d %H:%M:%S") + if not bucket_key: + continue bucket = time_buckets[bucket_key] - message = entry.get("message", "").lower() - level = entry.get("levelname", "").upper() + message = str(entry.get("message", "")).lower() + level = str(entry.get("levelname", "")).upper() - if "received proxy request" in message: + completed, revenue_msats, input_tokens, output_tokens = ( + self._extract_success_metrics(entry, message) + ) + if completed: bucket["total_requests"] += 1 - - if ( - "completed for streaming" in message - or "completed for non-streaming" in message - ): bucket["successful_chat_completions"] += 1 + bucket["input_tokens"] += input_tokens + bucket["output_tokens"] += output_tokens + bucket["total_tokens"] += input_tokens + output_tokens if level == "ERROR": bucket["errors"] += 1 @@ -504,21 +1050,19 @@ class LogManager: elif level == "WARNING": bucket["warnings"] += 1 - if "upstream request failed" in message or "revert payment" in message: + failed = ( + "upstream request failed" in message + or "revert payment" in message + ) + if failed: + bucket["total_requests"] += 1 bucket["failed_requests"] += 1 if "payment processed successfully" in message: bucket["payment_processed"] += 1 - if ( - "completed for streaming" in message - or "completed for non-streaming" in message - ): - cost_data = entry.get("cost_data") - if isinstance(cost_data, dict): - actual_cost = cost_data.get("total_msats", 0) - if isinstance(actual_cost, (int, float)) and actual_cost > 0: - bucket["revenue_msats"] += float(actual_cost) + if completed and revenue_msats > 0: + bucket["revenue_msats"] += revenue_msats if "revert payment" in message: max_cost = entry.get("max_cost_for_model", 0) @@ -534,11 +1078,42 @@ class LogManager: bucket["requests"] = bucket["total_requests"] result.append({"timestamp": bucket_key, **bucket}) + totals = { + "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, + } + for bucket in result: + totals["total_requests"] += int(bucket["total_requests"]) + totals["successful_chat_completions"] += int( + bucket["successful_chat_completions"] + ) + totals["failed_requests"] += int(bucket["failed_requests"]) + totals["errors"] += int(bucket["errors"]) + totals["warnings"] += int(bucket["warnings"]) + totals["payment_processed"] += int(bucket["payment_processed"]) + totals["upstream_errors"] += int(bucket["upstream_errors"]) + totals["revenue_msats"] += float(bucket["revenue_msats"]) + totals["refunds_msats"] += float(bucket["refunds_msats"]) + totals["input_tokens"] += int(bucket["input_tokens"]) + totals["output_tokens"] += int(bucket["output_tokens"]) + totals["total_tokens"] += int(bucket["total_tokens"]) + return { "metrics": result, "interval_minutes": interval_minutes, "hours_back": hours_back, "total_buckets": len(result), + "totals": totals, } 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/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/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/routstr/wallet.py b/routstr/wallet.py index 71ea18ac..f4817bb9 100644 --- a/routstr/wallet.py +++ b/routstr/wallet.py @@ -1,5 +1,6 @@ import asyncio import math +import time from typing import TypedDict from cashu.core.base import Proof, Token @@ -158,6 +159,14 @@ async def credit_balance( _wallets: dict[str, Wallet] = {} +_balances_cache_ttl_seconds = 300.0 +_balances_cache: dict[ + tuple[str, ...], tuple[float, tuple[list["BalanceDetail"], int, int, int]] +] = {} +_balances_refresh_tasks: dict[ + tuple[str, ...], asyncio.Task[tuple[list["BalanceDetail"], int, int, int]] +] = {} +_balances_cache_lock = asyncio.Lock() async def get_wallet(mint_url: str, unit: str = "sat", load: bool = True) -> Wallet: @@ -226,6 +235,12 @@ async def fetch_all_balances( """ if units is None: units = ["sat", "msat"] + units_key = tuple(units) + + now = time.time() + cached = _balances_cache.get(units_key) + if cached and cached[0] > now: + return cached[1] async def fetch_balance( session: db.AsyncSession, mint_url: str, unit: str @@ -261,47 +276,71 @@ async def fetch_all_balances( } return error_result - # Create tasks for all mint/unit combinations - async with db.create_session() as session: - tasks = [ - fetch_balance(session, mint_url, unit) - for mint_url in settings.cashu_mints - for unit in units - ] + async def compute_balances() -> tuple[list[BalanceDetail], int, int, int]: + # Create tasks for all mint/unit combinations + async with db.create_session() as session: + tasks = [ + fetch_balance(session, mint_url, unit) + for mint_url in settings.cashu_mints + for unit in units + ] - # Run all tasks concurrently - balance_details = list(await asyncio.gather(*tasks)) + # Run all tasks concurrently + balance_details = list(await asyncio.gather(*tasks)) - # Calculate totals - total_wallet_balance_sats = 0 - total_user_balance_sats = 0 + # Calculate totals + total_wallet_balance_sats = 0 + total_user_balance_sats = 0 - for detail in balance_details: - if not detail.get("error"): - # Convert to sats for total calculation - unit = detail["unit"] - proofs_balance_sats = ( - detail["wallet_balance"] - if unit == "sat" - else detail["wallet_balance"] // 1000 - ) - user_balance_sats = ( - detail["user_balance"] - if unit == "sat" - else detail["user_balance"] // 1000 - ) + for detail in balance_details: + if not detail.get("error"): + # Convert to sats for total calculation + unit = detail["unit"] + proofs_balance_sats = ( + detail["wallet_balance"] + if unit == "sat" + else detail["wallet_balance"] // 1000 + ) + user_balance_sats = ( + detail["user_balance"] + if unit == "sat" + else detail["user_balance"] // 1000 + ) - total_wallet_balance_sats += proofs_balance_sats - total_user_balance_sats += user_balance_sats + total_wallet_balance_sats += proofs_balance_sats + total_user_balance_sats += user_balance_sats - owner_balance = total_wallet_balance_sats - total_user_balance_sats + owner_balance = total_wallet_balance_sats - total_user_balance_sats + return ( + balance_details, + total_wallet_balance_sats, + total_user_balance_sats, + owner_balance, + ) - return ( - balance_details, - total_wallet_balance_sats, - total_user_balance_sats, - owner_balance, - ) + async with _balances_cache_lock: + now = time.time() + cached = _balances_cache.get(units_key) + if cached and cached[0] > now: + return cached[1] + + refresh_task = _balances_refresh_tasks.get(units_key) + if refresh_task is None or refresh_task.done(): + refresh_task = asyncio.create_task(compute_balances()) + _balances_refresh_tasks[units_key] = refresh_task + + result = await refresh_task + + async with _balances_cache_lock: + _balances_cache[units_key] = ( + time.time() + _balances_cache_ttl_seconds, + result, + ) + current_task = _balances_refresh_tasks.get(units_key) + if current_task is refresh_task and refresh_task.done(): + _balances_refresh_tasks.pop(units_key, None) + + return result async def periodic_payout() -> None: diff --git a/ui/app/page.tsx b/ui/app/page.tsx index 26534ed3..305daf1c 100644 --- a/ui/app/page.tsx +++ b/ui/app/page.tsx @@ -4,11 +4,12 @@ import { useEffect, useMemo, useState } from 'react'; import { format } from 'date-fns'; import { useQuery } from '@tanstack/react-query'; import { CalendarIcon, RefreshCw } from 'lucide-react'; +import { Bar, BarChart, CartesianGrid, XAxis, YAxis } from 'recharts'; 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 { RevenueByModelTable } from '@/components/revenue-by-model-table'; +import { TopModelsUsageChart } from '@/components/top-models-usage-chart'; import { DashboardBalanceSummary } from '@/components/dashboard-balance-summary'; import { AdminService, @@ -18,7 +19,12 @@ import { import { Button } from '@/components/ui/button'; import { Card, CardContent, CardHeader, CardTitle } from '@/components/ui/card'; import { Calendar } from '@/components/ui/calendar'; -import { Badge } from '@/components/ui/badge'; +import { + ChartContainer, + ChartTooltip, + ChartTooltipContent, + type ChartConfig as UiChartConfig, +} from '@/components/ui/chart'; import { Empty, EmptyDescription, @@ -45,6 +51,7 @@ import { CheatSheet } from '@/components/landing/cheat-sheet'; import { ConfigurationService } from '@/lib/api/services/configuration'; import { AppPageShell } from '@/components/app-page-shell'; import { useIsMobile } from '@/hooks/use-mobile'; +import type { DisplayUnit } from '@/lib/types/units'; import { cn } from '@/lib/utils'; type ChartDatum = Record & { timestamp: string }; @@ -62,6 +69,7 @@ type ChartConfig = { description: string; data: ChartDatum[]; dataKeys: ChartKeyConfig[]; + totals?: Partial>; metricType: 'currency' | 'count'; }; @@ -72,10 +80,13 @@ 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']; -const DEFAULT_TIME_RANGE_PRESET = TIME_RANGE_PRESETS[0]; +const DEFAULT_TIME_RANGE_PRESET = + TIME_RANGE_PRESETS.find((option) => option.value === '7d') ?? + TIME_RANGE_PRESETS[0]; function normalizeDateRange(range: DateRange): DateRange { if (!range.from || !range.to) { @@ -111,36 +122,6 @@ function getRangeHours(range?: DateRange): number | null { return Math.max(1, diffHours); } -function formatDateRangeLabel(range?: DateRange): string { - if (!range?.from && !range?.to) { - return 'Custom range'; - } - - if (range.from && !range.to) { - return `${format(range.from, 'MMM d, yyyy')} - ...`; - } - - if (!range.from || !range.to) { - return 'Custom range'; - } - - const normalized = normalizeDateRange(range); - const from = normalized.from; - const to = normalized.to; - - if (!from || !to) { - return 'Custom range'; - } - - const sameDay = format(from, 'yyyy-MM-dd') === format(to, 'yyyy-MM-dd'); - - if (sameDay) { - return format(from, 'MMM d, yyyy'); - } - - return `${format(from, 'MMM d')} - ${format(to, 'MMM d, yyyy')}`; -} - function formatCompactDateRangeLabel(range?: DateRange): string { if (!range?.from || !range.to) { return 'Custom range'; @@ -184,7 +165,7 @@ function SectionLoading({ label }: { label: string }) { if (label === 'summary') { return (
- {Array.from({ length: 12 }).map((_, index) => ( + {Array.from({ length: 14 }).map((_, index) => ( b - a - ); + const errorTypes = Object.entries(summary.error_types || {}) + .map(([type, count]) => ({ + type, + count: typeof count === 'number' ? count : Number(count || 0), + })) + .filter((entry) => Number.isFinite(entry.count) && entry.count > 0) + .sort((a, b) => b.count - a.count); - const hasModels = summary.unique_models.length > 0; const hasErrorTypes = errorTypes.length > 0; - if (!hasModels && !hasErrorTypes) { + if (!hasErrorTypes) { return null; } + const compactNumber = new Intl.NumberFormat('en-US', { + notation: 'compact', + maximumFractionDigits: 1, + }); + const truncateErrorType = (value: string): string => { + const maxLength = isMobile ? 16 : 26; + if (value.length <= maxLength) { + return value; + } + return `${value.slice(0, maxLength - 1)}…`; + }; + const chartData = errorTypes.slice(0, 10).map(({ type, count }) => ({ + type, + label: truncateErrorType(type), + count, + })); + const totalErrors = errorTypes.reduce((sum, entry) => sum + entry.count, 0); + const chartConfig: UiChartConfig = { + count: { + label: 'Errors', + color: 'var(--chart-4)', + }, + }; + return ( -
- {hasModels && ( - - - Active Models - - -
- {summary.unique_models.map((model) => ( - - {model} - - ))} -
-
-
- )} +
{hasErrorTypes && ( - + Error Types Distribution +

+ Top {chartData.length} categories in this range. +

- - {errorTypes.map(([type, count]) => ( -
- {type} - {count} -
- ))} + + + + + + compactNumber.format( + typeof value === 'number' ? value : Number(value || 0) + ) + } + /> + + + String(payload?.[0]?.payload?.type ?? '') + } + formatter={(value, name) => { + const numericValue = + typeof value === 'number' + ? value + : Number(value || 0); + return ( +
+ + {name} + + + {Number.isFinite(numericValue) + ? numericValue.toLocaleString() + : '-'} + +
+ ); + }} + /> + } + /> + +
+
+

+ Total errors: {totalErrors.toLocaleString()} +

)} @@ -428,8 +493,9 @@ function DashboardInsights({ summary }: { summary?: UsageSummary }) { } export default function DashboardPage() { - const [selectedPreset, setSelectedPreset] = - useState('24h'); + const [selectedPreset, setSelectedPreset] = useState( + DEFAULT_TIME_RANGE_PRESET.value + ); const [customRange, setCustomRange] = useState(); const [pendingCustomRange, setPendingCustomRange] = useState(); const [isCustomRangeActive, setIsCustomRangeActive] = useState(false); @@ -476,55 +542,62 @@ 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 (safeQueryHours > 90 * 24) { + return 4 * 60 * 60_000; + } + if (safeQueryHours > 30 * 24) { + return 2 * 60 * 60_000; + } + if (safeQueryHours > 7 * 24) { + return 30 * 60_000; + } + return 60_000; + }, [safeQueryHours]); + const revenueDisplayUnit: DisplayUnit = useMemo(() => { + if (displayUnit === 'usd' && usdPerSat === null) { + // Keep revenue charts meaningful while the USD rate is unavailable. + return 'sat'; + } + return displayUnit; + }, [displayUnit, usdPerSat]); + const revenueUnitLabel = + revenueDisplayUnit === 'usd' + ? 'USD' + : revenueDisplayUnit === 'sat' + ? 'sats' + : revenueDisplayUnit === 'msat' + ? 'msats' + : revenueDisplayUnit; const { - data: metricsData, - isLoading: metricsLoading, - 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: 60_000, + refetchInterval: usageRefetchIntervalMs, staleTime: 30_000, }); - const { - data: summaryData, - isLoading: summaryLoading, - refetch: refetchSummary, - } = useQuery({ - queryKey: ['usage-summary', queryHours], - queryFn: () => AdminService.getUsageSummary(queryHours), - enabled: isAuthenticated, - refetchInterval: 60_000, - 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, - refetch: refetchErrors, - } = useQuery({ - queryKey: ['usage-errors', queryHours], - queryFn: () => AdminService.getErrorDetails(queryHours, 100), - enabled: isAuthenticated, - refetchInterval: 60_000, - staleTime: 30_000, - }); - - const { - data: revenueByModelData, - isLoading: revenueByModelLoading, - refetch: refetchRevenueByModel, - } = useQuery({ - queryKey: ['revenue-by-model', queryHours], - queryFn: () => AdminService.getRevenueByModel(queryHours, 20), - enabled: isAuthenticated, - refetchInterval: 60_000, - 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) { @@ -532,39 +605,44 @@ export default function DashboardPage() { } const metricPoints = metricsData.metrics as ChartDatum[]; + const convertRevenueMsats = (amountMsats: number): number => { + if (revenueDisplayUnit === 'msat') { + return amountMsats; + } + + const sats = amountMsats / 1000; + if (revenueDisplayUnit === 'usd') { + return sats * (usdPerSat ?? 0); + } + + return sats; + }; const revenuePoints = metricsData.metrics.map( (metric: UsageMetricData) => ({ ...metric, - revenue_sats: metric.revenue_msats / 1000, - refunds_sats: metric.refunds_msats / 1000, - net_revenue_sats: (metric.revenue_msats - metric.refunds_msats) / 1000, + revenue_display: convertRevenueMsats(metric.revenue_msats), }) ) as ChartDatum[]; return [ { id: 'revenue', - title: 'Revenue Over Time (sats)', + title: 'Revenue Over Time', mobileTitle: 'Revenue', - description: 'Track gross revenue, refunds, and net revenue trends.', + description: 'Track collected revenue trends over time.', data: revenuePoints, metricType: 'currency', + totals: metricsTotals + ? { + revenue_display: convertRevenueMsats(metricsTotals.revenue_msats), + } + : undefined, dataKeys: [ { - key: 'revenue_sats', + key: 'revenue_display', name: 'Revenue', color: 'var(--chart-1)', }, - { - key: 'net_revenue_sats', - name: 'Net Revenue', - color: 'var(--chart-2)', - }, - { - key: 'refunds_sats', - name: 'Refunds', - color: 'var(--chart-5)', - }, ], }, { @@ -574,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', @@ -599,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', @@ -624,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', @@ -632,8 +730,41 @@ export default function DashboardPage() { }, ], }, + { + 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]); + }, [metricsData, metricsTotals, revenueDisplayUnit, usdPerSat]); useEffect(() => { if (chartConfigs.length === 0) { @@ -659,6 +790,7 @@ export default function DashboardPage() { ); } + if (!isAuthenticated) { return ; } @@ -669,13 +801,11 @@ export default function DashboardPage() { } setIsManualRefreshing(true); - await Promise.allSettled([ - refetchMetrics(), - refetchSummary(), - refetchErrors(), - refetchRevenueByModel(), - ]); - setIsManualRefreshing(false); + try { + await refetchUsageDashboard(); + } finally { + setIsManualRefreshing(false); + } }; const openRangePicker = () => { @@ -748,10 +878,6 @@ export default function DashboardPage() { isCustomRangeActive && customRange?.from && customRange?.to ? 'custom' : selectedPreset; - const activeRangeLabel = - selectedRangeValue === 'custom' - ? formatDateRangeLabel(customRange) - : activePreset.label; const compactCustomRangeLabel = formatCompactDateRangeLabel(customRange); return ( @@ -771,165 +897,167 @@ export default function DashboardPage() { usdPerSat={usdPerSat} /> -
-
-
+
+
+

Usage Analytics

-

- Select a preset or custom date range to analyze traffic and - revenue. -

-

- Showing {activeRangeLabel}. -

-
- -
-

- Range +

+ 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}
+ +
+
+
+ + + + + + + + + +
+ + +
+
+ + +
+ + {metricsLoading ? ( + + ) : activeChartConfig ? ( + ({ + id: config.id, + label: isMobile + ? (config.mobileTitle ?? config.title) + : config.title, + }))} + activeTabId={activeChartId} + onTabChange={setActiveChartId} + /> + ) : ( + + + + + + + + No data available + + No metrics data exists for this range yet. Try a broader + range. + + + + + + )} + + {!metricsLoading && modelUsageMixData && hasModelUsageMixMetrics ? ( + + ) : null} + + {summaryLoading ? ( + + ) : summaryData ? ( + + ) : null} + + + + {errorLoading ? ( + + ) : errorData ? ( + + ) : null}
- - {metricsLoading ? ( - - ) : activeChartConfig ? ( - ({ - id: config.id, - label: isMobile - ? (config.mobileTitle ?? config.title) - : config.title, - }))} - activeTabId={activeChartId} - onTabChange={setActiveChartId} - /> - ) : ( - - - - - - - - No data available - - No metrics data exists for this range yet. Try a broader - range. - - - - - - )} - - {summaryLoading ? ( - - ) : summaryData ? ( - - ) : null} - - - - {revenueByModelLoading ? ( - - ) : revenueByModelData && revenueByModelData.models.length > 0 ? ( - - ) : null} - - {errorLoading ? ( - - ) : errorData ? ( - - ) : null}
); diff --git a/ui/components/app-page-shell.tsx b/ui/components/app-page-shell.tsx index 73a7d8ac..ad2fdd31 100644 --- a/ui/components/app-page-shell.tsx +++ b/ui/components/app-page-shell.tsx @@ -8,7 +8,6 @@ import { FileTextIcon, LayoutDashboardIcon, LogOutIcon, - MoreHorizontalIcon, PanelLeftCloseIcon, PanelLeftOpenIcon, ServerIcon, @@ -22,12 +21,12 @@ import { Button } from '@/components/ui/button'; import { CurrencyToggle } from '@/components/currency-toggle'; import { ThemeToggle } from '@/components/theme-toggle'; import { - Drawer, - DrawerContent, - DrawerDescription, - DrawerHeader, - DrawerTitle, -} from '@/components/ui/drawer'; + Sheet, + SheetClose, + SheetContent, + SheetDescription, + SheetTitle, +} from '@/components/ui/sheet'; import { cn } from '@/lib/utils'; interface AppPageShellProps { @@ -45,9 +44,6 @@ const NAV_ITEMS = [ { title: 'Settings', url: '/settings', icon: SettingsIcon }, ] as const; -const MOBILE_NAV_TAB_WIDTH_CLASS = 'auto-cols-[22%]'; -const MOBILE_NAV_TAB_WIDTH_CLASS_XS = 'max-[359px]:auto-cols-[31%]'; - function isActivePath(pathname: string, itemUrl: string): boolean { if (itemUrl === '/') { return pathname === '/'; @@ -64,7 +60,7 @@ export function AppPageShell({ const pathname = usePathname(); const router = useRouter(); const [isSidebarCollapsed, setIsSidebarCollapsed] = useState(false); - const [isMobileMoreOpen, setIsMobileMoreOpen] = useState(false); + const [isMobileSidebarOpen, setIsMobileSidebarOpen] = useState(false); const handleLogout = async (): Promise => { try { @@ -82,36 +78,43 @@ export function AppPageShell({