diff --git a/routstr/core/admin.py b/routstr/core/admin.py index c3160c9e..cdfb2e86 100644 --- a/routstr/core/admin.py +++ b/routstr/core/admin.py @@ -2,6 +2,7 @@ import json import secrets from datetime import datetime, timezone from pathlib import Path +from typing import NoReturn from fastapi import APIRouter, Depends, HTTPException, Query, Request from pydantic import BaseModel @@ -26,18 +27,55 @@ logger = get_logger(__name__) admin_router = APIRouter(prefix="/admin", include_in_schema=False) admin_sessions: dict[str, int] = {} -ADMIN_SESSION_DURATION = 3600 +ADMIN_SESSION_DURATION = 12 * 60 * 60 +# Usage analytics remain queryable up to 12 months. +MAX_USAGE_ANALYTICS_HOURS = 365 * 24 + + +def _current_timestamp() -> int: + return int(datetime.now(timezone.utc).timestamp()) + + +def _cleanup_expired_admin_sessions(now_timestamp: int | None = None) -> None: + current_timestamp = ( + now_timestamp if now_timestamp is not None else _current_timestamp() + ) + expired_tokens = [ + token + for token, expiry_timestamp in admin_sessions.items() + if expiry_timestamp <= current_timestamp + ] + for token in expired_tokens: + admin_sessions.pop(token, None) + + +def _raise_unauthorized(detail: str) -> NoReturn: + raise HTTPException( + status_code=401, + detail=detail, + headers={"WWW-Authenticate": "Bearer"}, + ) def require_admin_api(request: Request) -> None: - auth_header = request.headers.get("Authorization") - if auth_header and auth_header.startswith("Bearer "): - token = auth_header.split(" ", 1)[1] - expiry = admin_sessions.get(token) - if expiry and expiry > int(datetime.now(timezone.utc).timestamp()): - return + auth_header = request.headers.get("Authorization", "") + if not auth_header.startswith("Bearer "): + _raise_unauthorized("Missing bearer token") - raise HTTPException(status_code=403, detail="Unauthorized") + token = auth_header.split(" ", 1)[1].strip() + if not token: + _raise_unauthorized("Missing bearer token") + + now_timestamp = _current_timestamp() + expiry_timestamp = admin_sessions.get(token) + if expiry_timestamp is None: + _raise_unauthorized("Invalid session token") + + if expiry_timestamp <= now_timestamp: + admin_sessions.pop(token, None) + _raise_unauthorized("Session expired") + + _cleanup_expired_admin_sessions(now_timestamp) @admin_router.get("/api/temporary-balances", dependencies=[Depends(require_admin_api)]) @@ -206,18 +244,10 @@ async def admin_login( raise HTTPException(status_code=401, detail="Invalid password") token = secrets.token_urlsafe(32) - expiry_timestamp = ( - int(datetime.now(timezone.utc).timestamp()) + ADMIN_SESSION_DURATION - ) + expiry_timestamp = _current_timestamp() + ADMIN_SESSION_DURATION admin_sessions[token] = expiry_timestamp - expired_tokens = [ - t - for t, exp in admin_sessions.items() - if exp <= int(datetime.now(timezone.utc).timestamp()) - ] - for t in expired_tokens: - del admin_sessions[t] + _cleanup_expired_admin_sessions() return {"ok": True, "token": token, "expires_in": ADMIN_SESSION_DURATION} @@ -544,7 +574,6 @@ class UpstreamProviderCreate(BaseModel): api_version: str | None = None enabled: bool = True provider_fee: float = 1.01 - provider_settings: dict | None = None class UpstreamProviderUpdate(BaseModel): @@ -554,7 +583,6 @@ class UpstreamProviderUpdate(BaseModel): api_version: str | None = None enabled: bool | None = None provider_fee: float | None = None - provider_settings: dict | None = None @admin_router.get("/api/upstream-providers", dependencies=[Depends(require_admin_api)]) @@ -571,9 +599,6 @@ async def get_upstream_providers() -> list[dict[str, object]]: "api_version": p.api_version, "enabled": p.enabled, "provider_fee": p.provider_fee, - "provider_settings": json.loads(p.provider_settings) - if p.provider_settings - else None, } for p in providers ] @@ -603,9 +628,6 @@ async def create_upstream_provider( api_version=payload.api_version, enabled=payload.enabled, provider_fee=payload.provider_fee, - provider_settings=json.dumps(payload.provider_settings) - if payload.provider_settings - else None, ) session.add(provider) await session.commit() @@ -621,7 +643,6 @@ async def create_upstream_provider( "api_version": provider.api_version, "enabled": provider.enabled, "provider_fee": provider.provider_fee, - "provider_settings": payload.provider_settings, } @@ -641,9 +662,6 @@ async def get_upstream_provider(provider_id: int) -> dict[str, object]: "api_version": provider.api_version, "enabled": provider.enabled, "provider_fee": provider.provider_fee, - "provider_settings": json.loads(provider.provider_settings) - if provider.provider_settings - else None, } @@ -670,8 +688,6 @@ async def update_upstream_provider( provider.enabled = payload.enabled if payload.provider_fee is not None: provider.provider_fee = payload.provider_fee - if payload.provider_settings is not None: - provider.provider_settings = json.dumps(payload.provider_settings) session.add(provider) await session.commit() @@ -687,9 +703,6 @@ async def update_upstream_provider( "api_version": provider.api_version, "enabled": provider.enabled, "provider_fee": provider.provider_fee, - "provider_settings": json.loads(provider.provider_settings) - if provider.provider_settings - else None, } @@ -809,47 +822,6 @@ class TopupRequest(BaseModel): amount: int -class TopupTokenRequest(BaseModel): - token: str - - -@admin_router.post( - "/api/upstream-providers/{provider_id}/topup-token", - dependencies=[Depends(require_admin_api)], -) -async def topup_provider_with_token( - provider_id: int, payload: TopupTokenRequest -) -> dict: - """Redeem a Cashu token for an upstream provider.""" - async with create_session() as session: - provider = await session.get(UpstreamProviderRow, provider_id) - if not provider: - raise HTTPException(status_code=404, detail="Provider not found") - - import httpx - - async with httpx.AsyncClient() as client: - clean_url = provider.base_url.rstrip("/") - headers = {} - if provider.api_key: - headers["Authorization"] = f"Bearer {provider.api_key}" - resp = await client.post( - f"{clean_url}/v1/balance/topup", - json={"cashu_token": payload.token}, - headers=headers, - ) - - if resp.status_code == 200: - return {"ok": True, "message": "Token redeemed successfully"} - else: - logger.error(f"Upstream token topup failed: {resp.text}") - try: - error_detail = resp.json() - except Exception: - error_detail = resp.text - return {"ok": False, "message": f"Upstream error: {error_detail}"} - - @admin_router.post( "/api/upstream-providers/{provider_id}/topup", dependencies=[Depends(require_admin_api)], @@ -876,49 +848,7 @@ async def initiate_provider_topup( f"Initiating top-up for provider {provider_id}", extra={"amount": payload.amount}, ) - - # For Routstr providers, we might be doing a Lightning top-up or a direct token transfer - if provider.provider_type == "routstr": - # UI sends sats for Routstr topup - import httpx - - async with httpx.AsyncClient() as client: - clean_url = provider.base_url.rstrip("/") - # Proxy the request to upstream Routstr - # Use the actual API key from the database - resp = await client.post( - f"{clean_url}/v1/balance/lightning/invoice", - json={ - "amount_sats": int(payload.amount), - "purpose": "topup", - "api_key": provider.api_key, - }, - headers={"Authorization": f"Bearer {provider.api_key}"} if provider.api_key else {}, - ) - - if resp.status_code == 200: - data = resp.json() - return { - "ok": True, - "topup_data": { - "payment_request": data.get("bolt11"), - "invoice_id": data.get("invoice_id"), - "status": "pending", - }, - } - else: - logger.error(f"Upstream topup request failed: {resp.text}") - # Check if it's JSON error - try: - error_detail = resp.json() - except Exception: - error_detail = resp.text - raise HTTPException( - status_code=resp.status_code, detail=error_detail - ) - topup_data = await upstream_instance.initiate_topup(payload.amount) - logger.info( "Top-up initiated successfully", extra={ @@ -969,23 +899,6 @@ async def check_topup_status(provider_id: int, invoice_id: str) -> dict[str, obj if not provider: raise HTTPException(status_code=404, detail="Provider not found") - # For Routstr providers, proxy the status check - if provider.provider_type == "routstr": - import httpx - - async with httpx.AsyncClient() as client: - clean_url = provider.base_url.rstrip("/") - resp = await client.get( - f"{clean_url}/v1/balance/lightning/invoice/{invoice_id}/status", - headers={"Authorization": f"Bearer {provider.api_key}"} if provider.api_key else {}, - ) - if resp.status_code == 200: - status_data = resp.json() - return {"ok": True, "paid": status_data.get("status") == "paid"} - else: - logger.error(f"Upstream status check failed: {resp.text}") - return {"ok": False, "paid": False} - upstream_instance = _instantiate_provider(provider) if not upstream_instance: raise HTTPException( @@ -1013,7 +926,7 @@ async def check_topup_status(provider_id: int, invoice_id: str) -> dict[str, obj dependencies=[Depends(require_admin_api)], ) async def get_provider_balance(provider_id: int) -> dict[str, object]: - """Get the current balance for an upstream provider account.""" + """Get the current account balance for the upstream provider.""" from ..upstream.helpers import _instantiate_provider async with create_session() as session: @@ -1021,30 +934,6 @@ async def get_provider_balance(provider_id: int) -> dict[str, object]: if not provider: raise HTTPException(status_code=404, detail="Provider not found") - # For Routstr providers, proxy the balance check - if provider.provider_type == "routstr": - import httpx - - async with httpx.AsyncClient() as client: - clean_url = provider.base_url.rstrip("/") - headers = {} - if provider.api_key: - headers["Authorization"] = f"Bearer {provider.api_key}" - resp = await client.get( - f"{clean_url}/v1/balance/info", - headers=headers, - ) - if resp.status_code == 200: - data = resp.json() - # Return balance in sats - balance = data.get("balance", 0) - if isinstance(balance, (int, float)): - return {"ok": True, "balance_data": balance // 1000} - return {"ok": True, "balance_data": balance} - else: - logger.error(f"Failed to fetch Routstr balance: {resp.text}") - return {"ok": False, "balance_data": None} - upstream_instance = _instantiate_provider(provider) if not upstream_instance: raise HTTPException( @@ -1081,16 +970,57 @@ async def get_usage_metrics( interval: int = Query( default=15, ge=1, le=1440, description="Time interval in minutes" ), - hours: int = Query(default=24, ge=1, description="Hours of history to analyze"), + hours: int = Query( + default=24, + ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, + description="Hours of history to analyze", + ), ) -> dict: """Get usage metrics aggregated by time interval.""" return log_manager.get_usage_metrics(interval=interval, hours=hours) +@admin_router.get("/api/usage/dashboard", dependencies=[Depends(require_admin_api)]) +async def get_usage_dashboard( + request: Request, + interval: int = Query( + default=15, ge=1, le=1440, description="Time interval in minutes" + ), + hours: int = Query( + default=24, + ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, + description="Hours of history to analyze", + ), + error_limit: int = Query( + default=100, ge=1, le=1000, description="Maximum number of errors to return" + ), + model_limit: int = Query( + default=20, ge=1, le=100, description="Maximum number of models to return" + ), +) -> dict: + """ + Get all dashboard analytics in one request. + This runs one combined aggregation pass and avoids repeated scans. + """ + return log_manager.get_usage_dashboard( + interval=interval, + hours=hours, + error_limit=error_limit, + model_limit=model_limit, + ) + + @admin_router.get("/api/usage/summary", dependencies=[Depends(require_admin_api)]) async def get_usage_summary( request: Request, - hours: int = Query(default=24, ge=1, description="Hours of history to analyze"), + hours: int = Query( + default=24, + ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, + description="Hours of history to analyze", + ), ) -> dict: """Get summary statistics for the specified time period.""" return log_manager.get_usage_summary(hours=hours) @@ -1099,7 +1029,12 @@ async def get_usage_summary( @admin_router.get("/api/usage/error-details", dependencies=[Depends(require_admin_api)]) async def get_error_details( request: Request, - hours: int = Query(default=24, ge=1, description="Hours of history to analyze"), + hours: int = Query( + default=24, + ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, + description="Hours of history to analyze", + ), limit: int = Query( default=100, ge=1, le=1000, description="Maximum number of errors to return" ), @@ -1113,7 +1048,12 @@ async def get_error_details( ) async def get_revenue_by_model( request: Request, - hours: int = Query(default=24, ge=1, description="Hours of history to analyze"), + hours: int = Query( + default=24, + ge=1, + le=MAX_USAGE_ANALYTICS_HOURS, + description="Hours of history to analyze", + ), limit: int = Query( default=20, ge=1, le=100, description="Maximum number of models to return" ), @@ -1206,71 +1146,3 @@ async def get_log_dates_api(request: Request) -> dict[str, object]: continue return {"dates": dates} - - -@admin_router.post( - "/api/upstream-providers/{provider_id}/routstr/refund", - dependencies=[Depends(require_admin_api)], -) -async def refund_routstr_provider_balance(provider_id: int) -> dict[str, object]: - """Refund balance from an upstream Routstr provider back to the local wallet.""" - from ..upstream.helpers import _instantiate_provider - from ..upstream.routstr import RoutstrUpstreamProvider - - async with create_session() as session: - provider_row = await session.get(UpstreamProviderRow, provider_id) - if not provider_row: - raise HTTPException(status_code=404, detail="Provider not found") - - if provider_row.provider_type != "routstr": - raise HTTPException( - status_code=400, detail="Refund only supported for Routstr providers" - ) - - provider = _instantiate_provider(provider_row) - if not isinstance(provider, RoutstrUpstreamProvider): - raise HTTPException(status_code=400, detail="Invalid provider instance") - - try: - # Request refund from upstream - data = await provider.refund_balance() - if "error" in data: - # If the upstream returned an OpenAI-style error (like the model unknown error) - # it means the request likely didn't even reach the refund endpoint handler - # but was intercepted by the proxy layer. - error_info = data.get("error", {}) - message = ( - error_info.get("message") - if isinstance(error_info, dict) - else str(error_info) - ) - return { - "ok": False, - "message": f"Upstream refund failed: {message}", - } - - token = data.get("token") - if not token: - return {"ok": False, "message": "Upstream did not return a token"} - - # Receive token into local wallet - from ..wallet import recieve_token - - try: - # Use current wallet to receive - await recieve_token(token) - return { - "ok": True, - "message": "Successfully received refund from upstream provider", - } - except Exception as e: - logger.error(f"Failed to receive refund token: {e}") - return { - "ok": False, - "message": f"Failed to receive refund token: {str(e)}", - "token": token, - } - - except Exception as e: - logger.exception(f"Refund failed for provider {provider_id}") - raise HTTPException(status_code=500, detail=str(e)) diff --git a/routstr/core/log_manager.py b/routstr/core/log_manager.py index 61dcceb9..0444dcbf 100644 --- a/routstr/core/log_manager.py +++ b/routstr/core/log_manager.py @@ -1,17 +1,73 @@ import json +import time from collections import defaultdict from datetime import datetime, timedelta, timezone +from heapq import heappush, heapreplace from pathlib import Path -from typing import Any, Iterator +from threading import Lock +from typing import Any, Callable, Iterator, TypeVar from .logging import get_logger +from .usage_analytics_store import UsageAnalyticsStore logger = get_logger(__name__) +T = TypeVar("T") class LogManager: def __init__(self, logs_dir: Path = Path("logs")): self.logs_dir = logs_dir + self._usage_store = UsageAnalyticsStore(logs_dir=logs_dir) + self._analytics_cache_ttl_seconds = 30.0 + self._analytics_cache: dict[tuple[Any, ...], tuple[float, Any]] = {} + self._analytics_cache_lock = Lock() + self._cache_miss = object() + + def _get_cached(self, key: tuple[Any, ...]) -> Any: + now = time.time() + with self._analytics_cache_lock: + cached = self._analytics_cache.get(key) + if cached is None: + return self._cache_miss + + expires_at, value = cached + if expires_at <= now: + self._analytics_cache.pop(key, None) + return self._cache_miss + + return value + + def _set_cached( + self, key: tuple[Any, ...], value: Any, ttl_seconds: float | None = None + ) -> None: + ttl = ( + self._analytics_cache_ttl_seconds + if ttl_seconds is None + else max(1.0, ttl_seconds) + ) + expires_at = time.time() + ttl + with self._analytics_cache_lock: + self._analytics_cache[key] = (expires_at, value) + + def _cache_call( + self, + key: tuple[Any, ...], + compute: Callable[[], T], + ttl_seconds: float | None = None, + ) -> T: + cached = self._get_cached(key) + if cached is not self._cache_miss: + return cached + + value = compute() + self._set_cached(key, value, ttl_seconds=ttl_seconds) + return value + + def _get_cached_entries(self, hours: int) -> list[dict[str, Any]]: + return self._cache_call( + ("usage_entries", hours), + lambda: list(self._yield_log_entries(hours_back=hours)), + ) def _yield_log_entries( self, @@ -19,7 +75,6 @@ class LogManager: specific_date: str | None = None, reverse_files: bool = False, max_files: int | None = None, - window_center: datetime | None = None, ) -> Iterator[dict[str, Any]]: """ Yields log entries from files. @@ -29,7 +84,6 @@ class LogManager: specific_date: specific date string (YYYY-MM-DD) to look at. reverse_files: if True, process files in reverse order (newest first). max_files: maximum number of log files to process (most recent if reverse_files is True). - window_center: datetime object to center a 5-month window around. """ if not self.logs_dir.exists(): return @@ -44,36 +98,6 @@ class LogManager: log_files.append(log_file) else: log_files = sorted(self.logs_dir.glob("app_*.log")) - - if window_center: - # Calculate the 5 months: [center-2, center-1, center, center+1, center+2] - allowed_month_years = [] - cur_m = window_center.month - cur_y = window_center.year - - for offset in range(-2, 3): - m = cur_m + offset - y = cur_y - while m <= 0: - m += 12 - y -= 1 - while m > 12: - m -= 12 - y += 1 - allowed_month_years.append(f"{y}-{m:02d}") - - filtered_files = [] - for log_path in log_files: - try: - # Stem is "app_YYYY-MM-DD" - file_date_str = log_path.stem.split("_")[1] - file_month_year = file_date_str[:7] # YYYY-MM - if file_month_year in allowed_month_years: - filtered_files.append(log_path) - except Exception: - continue - log_files = filtered_files - if reverse_files: log_files.reverse() @@ -303,45 +327,124 @@ class LogManager: return 0 def get_usage_summary(self, hours: int = 24) -> dict: + def compute() -> dict: + try: + return self._usage_store.get_summary(hours_back=hours) + except Exception as e: + logger.error( + f"Usage analytics index failed, falling back to log scan: {e}" + ) + return self._calculate_summary_stats(self._get_cached_entries(hours)) + return self._cache_call( ("usage_summary", hours), - lambda: self._calculate_summary_stats(self._get_cached_entries(hours)), + compute, ) def get_usage_metrics(self, interval: int = 15, hours: int = 24) -> dict: + def compute() -> dict: + try: + return self._usage_store.get_metrics( + interval_minutes=interval, + hours_back=hours, + ) + except Exception as e: + logger.error( + f"Usage analytics index failed, falling back to log scan: {e}" + ) + return self._aggregate_metrics_by_time( + self._get_cached_entries(hours), interval, hours + ) + return self._cache_call( ("usage_metrics", interval, hours), - lambda: self._aggregate_metrics_by_time( - self._get_cached_entries(hours), interval, hours - ), + compute, + ) + + def get_usage_dashboard( + self, + interval: int = 15, + hours: int = 24, + error_limit: int = 100, + model_limit: int = 20, + ) -> dict: + # Large ranges are expensive to scan; keep cached longer. + if hours <= 24: + cache_ttl = 60.0 + elif hours <= 7 * 24: + cache_ttl = 300.0 + elif hours <= 30 * 24: + cache_ttl = 1800.0 + elif hours <= 90 * 24: + cache_ttl = 7200.0 + else: + cache_ttl = 21600.0 + + def compute() -> dict: + try: + return self._usage_store.get_dashboard( + interval_minutes=interval, + hours_back=hours, + error_limit=error_limit, + model_limit=model_limit, + ) + except Exception as e: + logger.error( + f"Usage analytics index failed, falling back to log scan: {e}" + ) + return self._aggregate_dashboard( + interval_minutes=interval, + hours_back=hours, + error_limit=error_limit, + model_limit=model_limit, + ) + + return self._cache_call( + ("usage_dashboard", interval, hours, error_limit, model_limit), + compute, + ttl_seconds=cache_ttl, ) def get_error_details(self, hours: int = 24, limit: int = 100) -> dict: def compute() -> dict: - errors: list[dict[str, Any]] = [] - - for entry in self._get_cached_entries(hours): - if str(entry.get("levelname", "")).upper() != "ERROR": - continue - - errors.append( - { - "timestamp": entry.get("asctime", ""), - "message": entry.get("message", ""), - "error_type": entry.get("error_type", "unknown"), - "pathname": entry.get("pathname", ""), - "lineno": entry.get("lineno", 0), - "request_id": entry.get("request_id", ""), - } + 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}" ) - errors.sort(key=lambda x: str(x["timestamp"]), reverse=True) + errors: list[dict] = [] + for entry in self._get_cached_entries(hours): + if str(entry.get("levelname", "")).upper() == "ERROR": + timestamp_str = entry.get("asctime", "") + errors.append( + { + "timestamp": timestamp_str, + "message": entry.get("message", ""), + "error_type": entry.get("error_type", "unknown"), + "pathname": entry.get("pathname", ""), + "lineno": entry.get("lineno", 0), + "request_id": entry.get("request_id", ""), + } + ) + + errors.sort(key=lambda x: x["timestamp"], reverse=True) return {"errors": errors[:limit], "total_count": len(errors)} return self._cache_call(("error_details", hours, limit), compute) def get_revenue_by_model(self, hours: int = 24, limit: int = 20) -> dict: 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}" + ) + entries = self._get_cached_entries(hours) model_stats: dict[str, dict[str, int | float]] = defaultdict( @@ -380,10 +483,7 @@ class LogManager: 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 - ): + if isinstance(max_cost, (int, float)) and max_cost > 0: model_stats[model]["refunds_msats"] += max_cost except Exception: @@ -559,6 +659,344 @@ class LogManager: return self._build_summary_response(stats) + def _aggregate_dashboard( + self, + interval_minutes: int, + hours_back: int, + error_limit: int, + model_limit: int, + ) -> dict[str, Any]: + time_buckets: dict[str, dict[str, Any]] = defaultdict( + lambda: { + "total_requests": 0, + "successful_chat_completions": 0, + "failed_requests": 0, + "errors": 0, + "warnings": 0, + "payment_processed": 0, + "upstream_errors": 0, + "revenue_msats": 0.0, + "refunds_msats": 0.0, + "input_tokens": 0, + "output_tokens": 0, + "total_tokens": 0, + } + ) + summary_stats: dict[str, Any] = { + "total_entries": 0, + "total_requests": 0, + "successful_chat_completions": 0, + "failed_requests": 0, + "total_errors": 0, + "total_warnings": 0, + "payment_processed": 0, + "upstream_errors": 0, + "unique_models": set(), + "error_types": defaultdict(int), + "revenue_msats": 0.0, + "refunds_msats": 0.0, + "input_tokens": 0, + "output_tokens": 0, + "total_tokens": 0, + } + model_stats: dict[str, dict[str, int | float]] = defaultdict( + lambda: { + "revenue_msats": 0, + "refunds_msats": 0, + "requests": 0, + "successful": 0, + "failed": 0, + } + ) + model_mix_buckets: dict[str, dict[str, int]] = defaultdict( + lambda: defaultdict(int) + ) + model_mix_revenue_buckets: dict[str, dict[str, float]] = defaultdict( + lambda: defaultdict(float) + ) + model_mix_token_buckets: dict[str, dict[str, int]] = defaultdict( + lambda: defaultdict(int) + ) + model_mix_totals: dict[str, int] = defaultdict(int) + model_mix_revenue_totals: dict[str, float] = defaultdict(float) + model_mix_token_totals: dict[str, int] = defaultdict(int) + latest_errors_heap: list[tuple[str, dict[str, Any]]] = [] + total_error_count = 0 + + for entry in self._yield_log_entries(hours_back=hours_back): + try: + summary_stats["total_entries"] += 1 + + timestamp_str = entry.get("asctime", "") + message = str(entry.get("message", "")).lower() + level = str(entry.get("levelname", "")).upper() + model = entry.get("model", "unknown") + if not isinstance(model, str): + model = "unknown" + + bucket_key = ( + self._bucket_key_for_timestamp(timestamp_str, interval_minutes) + if isinstance(timestamp_str, str) + else None + ) + bucket = time_buckets[bucket_key] if bucket_key else None + + if level == "ERROR": + summary_stats["total_errors"] += 1 + if bucket: + bucket["errors"] += 1 + if "error_type" in entry: + summary_stats["error_types"][str(entry["error_type"])] += 1 + + total_error_count += 1 + error_item = { + "timestamp": timestamp_str, + "message": entry.get("message", ""), + "error_type": entry.get("error_type", "unknown"), + "pathname": entry.get("pathname", ""), + "lineno": entry.get("lineno", 0), + "request_id": entry.get("request_id", ""), + } + if len(latest_errors_heap) < error_limit: + heappush(latest_errors_heap, (timestamp_str, error_item)) + elif timestamp_str > latest_errors_heap[0][0]: + heapreplace(latest_errors_heap, (timestamp_str, error_item)) + elif level == "WARNING": + summary_stats["total_warnings"] += 1 + if bucket: + bucket["warnings"] += 1 + + completed, revenue_msats, input_tokens, output_tokens = ( + self._extract_success_metrics(entry, message) + ) + if completed: + summary_stats["total_requests"] += 1 + summary_stats["successful_chat_completions"] += 1 + summary_stats["input_tokens"] += input_tokens + summary_stats["output_tokens"] += output_tokens + summary_stats["total_tokens"] += input_tokens + output_tokens + model_stats[model]["requests"] += 1 + model_stats[model]["successful"] += 1 + model_mix_totals[model] += 1 + if bucket: + bucket["total_requests"] += 1 + bucket["successful_chat_completions"] += 1 + bucket["input_tokens"] += input_tokens + bucket["output_tokens"] += output_tokens + bucket["total_tokens"] += input_tokens + output_tokens + if bucket_key: + model_mix_buckets[bucket_key][model] += 1 + if revenue_msats > 0: + model_mix_revenue_buckets[bucket_key][model] += revenue_msats + model_mix_revenue_totals[model] += revenue_msats + if input_tokens > 0 or output_tokens > 0: + token_total = input_tokens + output_tokens + model_mix_token_buckets[bucket_key][model] += token_total + model_mix_token_totals[model] += token_total + + if revenue_msats > 0: + summary_stats["revenue_msats"] += revenue_msats + model_stats[model]["revenue_msats"] += revenue_msats + if bucket: + bucket["revenue_msats"] += revenue_msats + + failed = ( + "upstream request failed" in message + or "revert payment" in message + ) + if failed: + summary_stats["total_requests"] += 1 + summary_stats["failed_requests"] += 1 + model_stats[model]["requests"] += 1 + model_stats[model]["failed"] += 1 + if bucket: + bucket["total_requests"] += 1 + bucket["failed_requests"] += 1 + + if "payment processed successfully" in message: + summary_stats["payment_processed"] += 1 + if bucket: + bucket["payment_processed"] += 1 + + if "upstream" in message and level == "ERROR": + summary_stats["upstream_errors"] += 1 + if bucket: + bucket["upstream_errors"] += 1 + + if model != "unknown": + summary_stats["unique_models"].add(model) + + if "revert payment" in message: + max_cost = entry.get("max_cost_for_model", 0) + if isinstance(max_cost, (int, float)) and max_cost > 0: + max_cost_float = float(max_cost) + summary_stats["refunds_msats"] += max_cost_float + model_stats[model]["refunds_msats"] += max_cost_float + if bucket: + bucket["refunds_msats"] += max_cost_float + except Exception: + continue + + metrics_result = [] + for bucket_key in sorted(time_buckets.keys()): + bucket = dict(time_buckets[bucket_key]) + bucket["requests"] = bucket["total_requests"] + metrics_result.append({"timestamp": bucket_key, **bucket}) + + models: list[dict[str, Any]] = [] + total_revenue = 0.0 + for model_name, stats in model_stats.items(): + revenue_msats = float(stats["revenue_msats"]) + refunds_msats = float(stats["refunds_msats"]) + revenue_sats = revenue_msats / 1000 + refunds_sats = refunds_msats / 1000 + net_revenue_sats = revenue_sats - refunds_sats + total_revenue += net_revenue_sats + + successful = int(stats["successful"]) + models.append( + { + "model": model_name, + "revenue_sats": revenue_sats, + "refunds_sats": refunds_sats, + "net_revenue_sats": net_revenue_sats, + "requests": int(stats["requests"]), + "successful": successful, + "failed": int(stats["failed"]), + "avg_revenue_per_request": ( + revenue_sats / successful if successful > 0 else 0 + ), + } + ) + + models.sort(key=lambda x: float(x["net_revenue_sats"]), reverse=True) + latest_errors = [ + item + for _, item in sorted( + latest_errors_heap, key=lambda x: x[0], reverse=True + ) + ] + top_model_limit = max(1, min(model_limit, 20)) + top_models_requests = [ + model_name + for model_name, _ in sorted( + ( + (name, count) + for name, count in model_mix_totals.items() + if name != "unknown" + ), + key=lambda item: item[1], + reverse=True, + )[:top_model_limit] + ] + top_models_revenue = [ + model_name + for model_name, _ in sorted( + ( + (name, amount) + for name, amount in model_mix_revenue_totals.items() + if name != "unknown" + ), + key=lambda item: item[1], + reverse=True, + )[:top_model_limit] + ] + top_models_tokens = [ + model_name + for model_name, _ in sorted( + ( + (name, token_count) + for name, token_count in model_mix_token_totals.items() + if name != "unknown" + ), + key=lambda item: item[1], + reverse=True, + )[:top_model_limit] + ] + selected_models: list[str] = [] + for model in top_models_requests + top_models_revenue + top_models_tokens: + if model not in selected_models: + selected_models.append(model) + top_model_set = set(selected_models) + + model_usage_mix_metrics: list[dict[str, Any]] = [] + mix_bucket_keys = sorted( + set(model_mix_buckets.keys()) + | set(model_mix_revenue_buckets.keys()) + | set(model_mix_token_buckets.keys()) + ) + for bucket_key in mix_bucket_keys: + counts = model_mix_buckets.get(bucket_key, {}) + revenue_counts = model_mix_revenue_buckets.get(bucket_key, {}) + token_counts = model_mix_token_buckets.get(bucket_key, {}) + others = 0 + others_revenue_msats = 0.0 + others_tokens = 0 + model_counts: dict[str, int] = {} + model_revenue_msats: dict[str, float] = {} + model_tokens: dict[str, int] = {} + for model_name, successful_count in counts.items(): + if model_name in top_model_set: + model_counts[model_name] = int(successful_count) + else: + others += int(successful_count) + for model_name, revenue_value in revenue_counts.items(): + if model_name in top_model_set: + model_revenue_msats[model_name] = float(revenue_value) + else: + others_revenue_msats += float(revenue_value) + for model_name, token_value in token_counts.items(): + if model_name in top_model_set: + model_tokens[model_name] = int(token_value) + else: + others_tokens += int(token_value) + + model_usage_mix_metrics.append( + { + "timestamp": bucket_key, + "total_successful": int(sum(counts.values())), + "total_revenue_msats": float(sum(revenue_counts.values())), + "total_tokens": int(sum(token_counts.values())), + "others": others, + "others_revenue_msats": others_revenue_msats, + "others_tokens": others_tokens, + "model_counts": model_counts, + "model_revenue_msats": model_revenue_msats, + "model_tokens": model_tokens, + } + ) + + return { + "metrics": { + "metrics": metrics_result, + "interval_minutes": interval_minutes, + "hours_back": hours_back, + "total_buckets": len(metrics_result), + }, + "summary": self._build_summary_response(summary_stats), + "error_details": { + "errors": latest_errors, + "total_count": total_error_count, + }, + "revenue_by_model": { + "models": models[:model_limit], + "total_revenue_sats": total_revenue, + "total_models": len(models), + }, + "model_usage_mix": { + "top_models": top_models_requests, + "top_models_by_metric": { + "requests": top_models_requests, + "revenue": top_models_revenue, + "tokens": top_models_tokens, + }, + "metrics": model_usage_mix_metrics, + "interval_minutes": interval_minutes, + "hours_back": hours_back, + "total_buckets": len(model_usage_mix_metrics), + }, + } + def _aggregate_metrics_by_time( self, entries: list[dict], interval_minutes: int, hours_back: int ) -> dict: diff --git a/routstr/core/usage_analytics_store.py b/routstr/core/usage_analytics_store.py new file mode 100644 index 00000000..809604e3 --- /dev/null +++ b/routstr/core/usage_analytics_store.py @@ -0,0 +1,1043 @@ +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 = "2" + + 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, + max_points: int | None, + ) -> 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, + max_points=max_points, + ) + 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, + ) + + return { + "metrics": metrics, + "summary": summary, + "error_details": error_details, + "revenue_by_model": revenue_by_model, + } + + 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, + max_points=None, + ) + + 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 + if current_version != self.SCHEMA_VERSION: + self._drop_index_tables_locked(conn) + conn.execute( + """ + INSERT OR REPLACE INTO analytics_meta (key, value) + VALUES ('schema_version', ?) + """, + (self.SCHEMA_VERSION,), + ) + + 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 + ) + """ + ) + 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, + 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_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)" + ) + conn.commit() + + 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 = ( + "completed for streaming" in message + or "completed for non-streaming" in 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 + + 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: + cost_float = float(actual_cost) + bucket["revenue_msats"] += cost_float + model_bucket["revenue_msats"] += cost_float + + 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"]), + ) + 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 + ) + 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 + """, + rows, + ) + + if model_updates: + rows = [ + ( + minute_ts, + model, + int(stats["requests"]), + int(stats["successful"]), + int(stats["failed"]), + float(stats["revenue_msats"]), + float(stats["refunds_msats"]), + ) + 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 + ) + 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 + """, + rows, + ) + + if model_presence_updates: + 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 + """, + rows, + ) + + if error_type_updates: + 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 + """, + 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, + max_points: int | None, + ) -> 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 + 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, + } + + 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"]) + + 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 + + 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, + "requests": total_requests, + } + ) + + points = self._downsample_metric_points(points, max_points) + + 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"]), + } + + 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 + 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"]) + + 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, + "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 _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 _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, + } + + 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, + } + + def _downsample_metric_points( + self, points: list[dict[str, Any]], max_points: int | None + ) -> list[dict[str, Any]]: + if max_points is None or max_points <= 0: + return points + if len(points) <= max_points: + return points + if max_points == 1: + return [points[-1]] + + step = (len(points) - 1) / (max_points - 1) + sampled: list[dict[str, Any]] = [] + last_index = -1 + + for i in range(max_points): + index = int(round(i * step)) + if index <= last_index: + index = min(last_index + 1, len(points) - 1) + sampled.append(points[index]) + last_index = index + + sampled[0] = points[0] + sampled[-1] = points[-1] + return sampled 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/layout.tsx b/ui/app/layout.tsx index 92c51b56..0ff625d7 100644 --- a/ui/app/layout.tsx +++ b/ui/app/layout.tsx @@ -3,7 +3,6 @@ import { GeistMono } from 'geist/font/mono'; import { GeistSans } from 'geist/font/sans'; import './globals.css'; import { Providers } from './providers'; -import { SuppressHydrationWarning } from '@/components/suppress-hydration-warning'; export const metadata: Metadata = { title: 'Routstr', @@ -23,9 +22,7 @@ export default function RootLayout({ - - {children} - + {children} ); diff --git a/ui/app/page.tsx b/ui/app/page.tsx index 380f465e..76addfd7 100644 --- a/ui/app/page.tsx +++ b/ui/app/page.tsx @@ -9,6 +9,7 @@ import type { DateRange } from 'react-day-picker'; import { UsageMetricsChart } from '@/components/usage-metrics-chart'; import { UsageSummaryCards } from '@/components/usage-summary-cards'; import { ErrorDetailsTable } from '@/components/error-details-table'; +import { TopModelsUsageChart } from '@/components/top-models-usage-chart'; import { DashboardBalanceSummary } from '@/components/dashboard-balance-summary'; import { AdminService, @@ -79,6 +80,7 @@ const TIME_RANGE_PRESETS = [ { value: '3m', label: 'Last 3 Months', hours: 90 * 24 }, { value: '12m', label: 'Last 12 Months', hours: 365 * 24 }, ] as const; +const MAX_USAGE_RANGE_HOURS = 365 * 24; type TimeRangePresetValue = (typeof TIME_RANGE_PRESETS)[number]['value']; @@ -159,20 +161,6 @@ function getAutoIntervalMinutes(hours: number): number { ); } -function getQueryErrorMessage(error: unknown): string { - if ( - error && - typeof error === 'object' && - 'message' in error && - typeof error.message === 'string' && - error.message.trim().length > 0 - ) { - return error.message; - } - - return 'The analytics request failed. Refresh and try again.'; -} - function SectionLoading({ label }: { label: string }) { if (label === 'summary') { return ( @@ -554,19 +542,21 @@ export default function DashboardPage() { isCustomRangeActive && customRangeHours ? customRangeHours : activePreset.hours; - const autoInterval = getAutoIntervalMinutes(queryHours); + const safeQueryHours = Math.min(queryHours, MAX_USAGE_RANGE_HOURS); + const isUsageRangeCapped = safeQueryHours < queryHours; + const autoInterval = getAutoIntervalMinutes(safeQueryHours); const usageRefetchIntervalMs = useMemo(() => { - if (queryHours > 90 * 24) { + if (safeQueryHours > 90 * 24) { return 4 * 60 * 60_000; } - if (queryHours > 30 * 24) { + if (safeQueryHours > 30 * 24) { return 2 * 60 * 60_000; } - if (queryHours > 7 * 24) { + if (safeQueryHours > 7 * 24) { return 30 * 60_000; } return 60_000; - }, [queryHours]); + }, [safeQueryHours]); const revenueDisplayUnit: DisplayUnit = useMemo(() => { if (displayUnit === 'usd' && usdPerSat === null) { // Keep revenue charts meaningful while the USD rate is unavailable. @@ -584,43 +574,30 @@ export default function DashboardPage() { : revenueDisplayUnit; const { - data: metricsData, - isLoading: metricsLoading, - error: metricsError, - refetch: refetchMetrics, + data: usageDashboardData, + isLoading: usageDashboardLoading, + refetch: refetchUsageDashboard, } = useQuery({ - queryKey: ['usage-metrics', autoInterval, queryHours], - queryFn: () => AdminService.getUsageMetrics(autoInterval, queryHours), + queryKey: ['usage-dashboard', autoInterval, safeQueryHours], + queryFn: () => + AdminService.getUsageDashboard(safeQueryHours, autoInterval, 100, 20), enabled: isAuthenticated, refetchInterval: usageRefetchIntervalMs, staleTime: 30_000, }); - const { - data: summaryData, - isLoading: summaryLoading, - error: summaryError, - refetch: refetchSummary, - } = useQuery({ - queryKey: ['usage-summary', queryHours], - queryFn: () => AdminService.getUsageSummary(queryHours), - enabled: isAuthenticated, - refetchInterval: usageRefetchIntervalMs, - staleTime: 30_000, - }); + const metricsData = usageDashboardData?.metrics; + const summaryData = usageDashboardData?.summary; + const errorData = usageDashboardData?.error_details; + const modelUsageMixData = usageDashboardData?.model_usage_mix; + const hasModelUsageMixMetrics = + Array.isArray(modelUsageMixData?.metrics) && + modelUsageMixData.metrics.length > 0; - const { - data: errorData, - isLoading: errorLoading, - error: errorDetailsError, - refetch: refetchErrors, - } = useQuery({ - queryKey: ['usage-errors', queryHours], - queryFn: () => AdminService.getErrorDetails(queryHours, 100), - enabled: isAuthenticated, - refetchInterval: usageRefetchIntervalMs, - staleTime: 30_000, - }); + const metricsLoading = usageDashboardLoading; + const summaryLoading = usageDashboardLoading; + const errorLoading = usageDashboardLoading; + const metricsTotals = metricsData?.totals; const chartConfigs = useMemo(() => { if (!metricsData || metricsData.metrics.length === 0) { @@ -647,12 +624,6 @@ export default function DashboardPage() { }) ) as ChartDatum[]; - const hasTokenMetrics = metricPoints.some((metric) => - ['input_tokens', 'output_tokens', 'total_tokens'].some( - (key) => typeof metric[key] === 'number' - ) - ); - return [ { id: 'revenue', @@ -661,6 +632,11 @@ export default function DashboardPage() { description: 'Track collected revenue trends over time.', data: revenuePoints, metricType: 'currency', + totals: metricsTotals + ? { + revenue_display: convertRevenueMsats(metricsTotals.revenue_msats), + } + : undefined, dataKeys: [ { key: 'revenue_display', @@ -676,6 +652,14 @@ export default function DashboardPage() { description: 'Understand traffic and completion reliability over time.', data: metricPoints, metricType: 'count', + totals: metricsTotals + ? { + total_requests: metricsTotals.total_requests, + successful_chat_completions: + metricsTotals.successful_chat_completions, + failed_requests: metricsTotals.failed_requests, + } + : undefined, dataKeys: [ { key: 'total_requests', @@ -701,6 +685,13 @@ export default function DashboardPage() { description: 'Monitor warnings, handled errors, and upstream failures.', data: metricPoints, metricType: 'count', + totals: metricsTotals + ? { + errors: metricsTotals.errors, + warnings: metricsTotals.warnings, + upstream_errors: metricsTotals.upstream_errors, + } + : undefined, dataKeys: [ { key: 'errors', @@ -726,6 +717,11 @@ export default function DashboardPage() { description: 'Follow payment processing activity by interval.', data: metricPoints, metricType: 'count', + totals: metricsTotals + ? { + payment_processed: metricsTotals.payment_processed, + } + : undefined, dataKeys: [ { key: 'payment_processed', @@ -734,38 +730,41 @@ export default function DashboardPage() { }, ], }, - ...(hasTokenMetrics - ? [ - { - id: 'tokens', - title: 'Token Usage', - mobileTitle: 'Tokens', - description: - 'Track input, output, and total token throughput over time.', - data: metricPoints, - metricType: 'count' as const, - dataKeys: [ - { - key: 'total_tokens', - name: 'Total Tokens', - color: 'var(--chart-1)', - }, - { - key: 'input_tokens', - name: 'Input Tokens', - color: 'var(--chart-2)', - }, - { - key: 'output_tokens', - name: 'Output Tokens', - color: 'var(--chart-3)', - }, - ], - }, - ] - : []), + { + id: 'tokens', + title: 'Token Usage', + mobileTitle: 'Tokens', + description: + 'Track input, output, and total token throughput over time.', + data: metricPoints, + metricType: 'count', + totals: metricsTotals + ? { + input_tokens: metricsTotals.input_tokens, + output_tokens: metricsTotals.output_tokens, + total_tokens: metricsTotals.total_tokens, + } + : undefined, + dataKeys: [ + { + key: 'total_tokens', + name: 'Total Tokens', + color: 'var(--chart-1)', + }, + { + key: 'input_tokens', + name: 'Input Tokens', + color: 'var(--chart-2)', + }, + { + key: 'output_tokens', + name: 'Output Tokens', + color: 'var(--chart-3)', + }, + ], + }, ]; - }, [metricsData, revenueDisplayUnit, usdPerSat]); + }, [metricsData, metricsTotals, revenueDisplayUnit, usdPerSat]); useEffect(() => { if (chartConfigs.length === 0) { @@ -803,11 +802,7 @@ export default function DashboardPage() { setIsManualRefreshing(true); try { - await Promise.allSettled([ - refetchMetrics(), - refetchSummary(), - refetchErrors(), - ]); + await refetchUsageDashboard(); } finally { setIsManualRefreshing(false); } @@ -913,10 +908,16 @@ export default function DashboardPage() { All cards and charts in this section update from the selected range.

+ {isUsageRangeCapped ? ( +

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

+ ) : null} -
-
+
+
- - {isManualRefreshing ? 'Refreshing...' : 'Refresh'} - + {isManualRefreshing ? 'Refreshing...' : 'Refresh'}
{metricsLoading ? ( - ) : metricsError ? ( - - - - - - - - Unable to load analytics - - {getQueryErrorMessage(metricsError)} - - - - - ) : activeChartConfig ? ( )} + {!metricsLoading && modelUsageMixData && hasModelUsageMixMetrics ? ( + + ) : null} + {summaryLoading ? ( - ) : summaryError ? ( - - - - - Usage summary unavailable - - {getQueryErrorMessage(summaryError)} - - - - - ) : summaryData ? ( ) : null} @@ -1074,19 +1051,6 @@ export default function DashboardPage() { {errorLoading ? ( - ) : errorDetailsError ? ( - - - - - Error details unavailable - - {getQueryErrorMessage(errorDetailsError)} - - - - - ) : errorData ? ( ) : null} diff --git a/ui/components/revenue-by-model-table.tsx b/ui/components/revenue-by-model-table.tsx new file mode 100644 index 00000000..00183bab --- /dev/null +++ b/ui/components/revenue-by-model-table.tsx @@ -0,0 +1,202 @@ +'use client'; + +import { useCallback, useMemo } from 'react'; +import { Bar, BarChart, CartesianGrid, XAxis, YAxis } from 'recharts'; +import { Card, CardContent, CardHeader, CardTitle } from '@/components/ui/card'; +import { + ChartConfig, + ChartContainer, + ChartTooltip, + ChartTooltipContent, +} from '@/components/ui/chart'; +import { ModelRevenueData } from '@/lib/api/services/admin'; +import { convertToMsat, formatFromMsat } from '@/lib/currency'; +import { useIsMobile } from '@/hooks/use-mobile'; +import type { DisplayUnit } from '@/lib/types/units'; + +interface RevenueByModelTableProps { + models: ModelRevenueData[]; + displayUnit: DisplayUnit; + usdPerSat: number | null; +} + +function truncateModelName(value: string, maxLength: number): string { + if (value.length <= maxLength) { + return value; + } + return `${value.slice(0, maxLength - 1)}…`; +} + +export function RevenueByModelTable({ + models, + displayUnit, + usdPerSat, +}: RevenueByModelTableProps) { + const isMobile = useIsMobile(); + + const revenueDisplayUnit: DisplayUnit = useMemo(() => { + if (displayUnit === 'usd' && usdPerSat === null) { + return 'sat'; + } + return displayUnit; + }, [displayUnit, usdPerSat]); + const unitLabel = revenueDisplayUnit === 'usd' ? 'USD' : revenueDisplayUnit; + + const compactNumber = useMemo( + () => + new Intl.NumberFormat('en-US', { + notation: 'compact', + maximumFractionDigits: 1, + }), + [] + ); + + const convertSatsToDisplay = useCallback( + (sats: number): number => { + if (revenueDisplayUnit === 'msat') { + return sats * 1000; + } + if (revenueDisplayUnit === 'usd') { + return sats * (usdPerSat ?? 0); + } + return sats; + }, + [revenueDisplayUnit, usdPerSat] + ); + + const formatAmount = (sats: number) => + formatFromMsat(convertToMsat(sats, 'sat'), revenueDisplayUnit, usdPerSat); + + const formatCompactAmount = (value: number): string => { + const compact = compactNumber.format(value); + if (revenueDisplayUnit === 'usd') { + return `$${compact}`; + } + return `${compact} ${unitLabel}`; + }; + + const totalCollectedRevenue = models.reduce( + (sum, model) => sum + model.revenue_sats, + 0 + ); + const totalOperationalNet = models.reduce( + (sum, model) => sum + model.net_revenue_sats, + 0 + ); + + const chartData = useMemo( + () => + [...models] + .sort((a, b) => b.revenue_sats - a.revenue_sats) + .slice(0, 12) + .map((model) => ({ + model: model.model, + modelLabel: truncateModelName(model.model, isMobile ? 16 : 28), + revenueDisplay: convertSatsToDisplay(model.revenue_sats), + })), + [models, isMobile, convertSatsToDisplay] + ); + + const chartConfig: ChartConfig = { + revenueDisplay: { + label: 'Revenue', + color: 'var(--chart-1)', + }, + }; + + if (chartData.length === 0) { + return ( + + + Revenue by Model + + + No model data available + + + ); + } + + return ( + + + Revenue by Model +

+ Total Collected Revenue:{' '} + + {formatAmount(totalCollectedRevenue)} + +

+

+ Operational Net:{' '} + + {formatAmount(totalOperationalNet)} + +

+
+ + + + + + compactNumber.format( + typeof value === 'number' ? value : Number(value || 0) + ) + } + /> + + String(label)} + formatter={(value, name) => { + const numericValue = + typeof value === 'number' ? value : Number(value || 0); + return ( +
+ + {name} + + + {Number.isFinite(numericValue) + ? formatCompactAmount(numericValue) + : '-'} + +
+ ); + }} + /> + } + /> + +
+
+
+
+ ); +} diff --git a/ui/lib/api/services/admin.ts b/ui/lib/api/services/admin.ts index d91e24cb..b3f7677f 100644 --- a/ui/lib/api/services/admin.ts +++ b/ui/lib/api/services/admin.ts @@ -20,7 +20,6 @@ export const UpstreamProviderSchema = z.object({ api_version: z.string().nullable().optional(), enabled: z.boolean(), provider_fee: z.number().optional(), - provider_settings: z.record(z.string(), z.any()).nullable().optional(), }); export const CreateUpstreamProviderSchema = z.object({ @@ -30,7 +29,6 @@ export const CreateUpstreamProviderSchema = z.object({ api_version: z.string().nullable().optional(), enabled: z.boolean().default(true), provider_fee: z.number().optional(), - provider_settings: z.record(z.string(), z.any()).nullable().optional(), }); export const UpdateUpstreamProviderSchema = z.object({ @@ -40,7 +38,6 @@ export const UpdateUpstreamProviderSchema = z.object({ api_version: z.string().nullable().optional(), enabled: z.boolean().optional(), provider_fee: z.number().optional(), - provider_settings: z.record(z.string(), z.any()).nullable().optional(), }); export const AdminModelPricingSchema = z.object({ @@ -845,6 +842,23 @@ export class AdminService { ); } + static async getUsageDashboard( + hours: number = 24, + interval: number = 15, + errorLimit: number = 100, + modelLimit: number = 20 + ): Promise { + const params = new URLSearchParams(); + params.set('interval', String(interval)); + params.set('hours', String(hours)); + params.set('error_limit', String(errorLimit)); + params.set('model_limit', String(modelLimit)); + + return await apiClient.get( + `/admin/api/usage/dashboard?${params.toString()}` + ); + } + static async getUsageSummary(hours: number = 24): Promise { return await apiClient.get( `/admin/api/usage/summary?hours=${hours}` @@ -895,17 +909,9 @@ export class AdminService { ok: boolean; topup_data: Record; message: string; - }>(`/admin/api/upstream-providers/${providerId}/topup`, { amount }); - } - - static async topupProviderWithToken( - providerId: number, - token: string - ): Promise<{ ok: boolean; message?: string }> { - return await apiClient.post<{ ok: boolean; message?: string }>( - `/admin/api/upstream-providers/${providerId}/topup-token`, - { token } - ); + }>(`/admin/api/upstream-providers/${providerId}/topup`, { + amount: amount, + }); } static async checkTopupStatus( @@ -952,9 +958,9 @@ export interface UsageMetricData { upstream_errors: number; revenue_msats: number; refunds_msats: number; - input_tokens?: number; - output_tokens?: number; - total_tokens?: number; + input_tokens: number; + output_tokens: number; + total_tokens: number; [key: string]: unknown; } @@ -963,7 +969,20 @@ export interface UsageMetrics { interval_minutes: number; hours_back: number; total_buckets: number; - totals?: Partial>; + totals?: { + total_requests: number; + successful_chat_completions: number; + failed_requests: number; + errors: number; + warnings: number; + payment_processed: number; + upstream_errors: number; + revenue_msats: number; + refunds_msats: number; + input_tokens: number; + output_tokens: number; + total_tokens: number; + }; } export interface UsageSummary { @@ -978,6 +997,12 @@ export interface UsageSummary { unique_models_count: number; unique_models: string[]; error_types: Record; + input_tokens: number; + output_tokens: number; + total_tokens: number; + avg_input_tokens_per_completion: number; + avg_output_tokens_per_completion: number; + avg_total_tokens_per_completion: number; success_rate: number; revenue_msats: number; refunds_msats: number; @@ -987,8 +1012,6 @@ export interface UsageSummary { net_revenue_sats: number; avg_revenue_per_request_msats: number; refund_rate: number; - total_tokens?: number; - avg_total_tokens_per_completion?: number; } export interface ErrorDetail { @@ -1022,6 +1045,35 @@ export interface RevenueByModel { total_models: number; } +export interface ModelUsageMixMetric { + timestamp: string; + total_successful: number; + total_revenue_msats: number; + total_tokens: number; + others: number; + others_revenue_msats: number; + others_tokens: number; + model_counts: Record; + model_revenue_msats: Record; + model_tokens: Record; +} + +export interface ModelUsageMix { + top_models: string[]; + metrics: ModelUsageMixMetric[]; + interval_minutes: number; + hours_back: number; + total_buckets: number; +} + +export interface UsageDashboardResponse { + metrics: UsageMetrics; + summary: UsageSummary; + error_details: ErrorDetails; + revenue_by_model: RevenueByModel; + model_usage_mix?: ModelUsageMix; +} + export interface LogEntry { asctime: string; name: string;