diff --git a/routstr/core/admin.py b/routstr/core/admin.py index ebec8b00..a2986a5a 100644 --- a/routstr/core/admin.py +++ b/routstr/core/admin.py @@ -1,6 +1,7 @@ import json import secrets -from datetime import datetime, timezone +from collections import defaultdict +from datetime import datetime, timedelta, timezone from pathlib import Path from fastapi import APIRouter, Depends, HTTPException, Query, Request @@ -10,7 +11,6 @@ from sqlmodel import select from ..payment.models import _row_to_model, list_models from ..proxy import refresh_model_maps, reinitialize_upstreams -from ..search import search_logs from ..wallet import ( fetch_all_balances, get_proofs_per_mint_and_unit, @@ -20,8 +20,8 @@ from ..wallet import ( ) from .db import ApiKey, ModelRow, UpstreamProviderRow, create_session from .logging import get_logger +from .log_manager import log_manager from .settings import SettingsService, settings -from .usage_metrics import UsageMetricsService, list_metric_definitions logger = get_logger(__name__) @@ -170,38 +170,6 @@ async def get_balances_api(request: Request) -> list[dict[str, object]]: return [dict(d) for d in balance_details] -@admin_router.get( - "/api/usage-metrics/definitions", dependencies=[Depends(require_admin_api)] -) -async def get_usage_metric_definitions() -> list[dict[str, str]]: - return list_metric_definitions() - - -@admin_router.get("/api/usage-metrics", dependencies=[Depends(require_admin_api)]) -async def get_usage_metrics( - metrics: str | None = Query( - default=None, - description="Comma-separated list of usage metric identifiers to include", - ), - bucket_minutes: int = Query(default=15, ge=1, le=24 * 60), - hours: int = Query(default=24, ge=1, le=24 * 7), -) -> dict[str, object]: - metric_list = ( - [item.strip() for item in metrics.split(",") if item.strip()] - if metrics - else [] - ) - try: - result = await UsageMetricsService.collect( - metrics=metric_list, - bucket_minutes=bucket_minutes, - hours=hours, - ) - except ValueError as exc: - raise HTTPException(status_code=400, detail=str(exc)) from exc - return result.to_dict() - - @admin_router.get("/api/settings", dependencies=[Depends(require_admin_api)]) async def get_settings(request: Request) -> dict: data = settings.dict() @@ -860,8 +828,8 @@ async def view_logs(request: Request, request_id: str) -> str: except Exception as e: logger.error(f"Error reading log file {log_file}: {e}") - # Sort entries by timestamp if available (newest first) - log_entries.sort(key=lambda x: x.get("asctime", ""), reverse=True) + # Sort entries by timestamp if available + log_entries.sort(key=lambda x: x.get("asctime", ""), reverse=False) # Format log entries for display formatted_logs = [] @@ -942,72 +910,6 @@ async def view_logs(request: Request, request_id: str) -> str: ) -@admin_router.get("/api/logs", dependencies=[Depends(require_admin_api)]) -async def get_logs_api( - request: Request, - date: str | None = None, - level: str | None = None, - request_id: str | None = None, - search: str | None = None, - limit: int = 100, -) -> dict[str, object]: - """ - Get filtered log entries. - - Args: - date: Filter by specific date (YYYY-MM-DD) - level: Filter by log level - request_id: Filter by request ID - search: Search text in message and name fields (case-insensitive) - limit: Maximum number of entries to return - - Returns: - Dict containing logs and filter metadata - """ - logs_dir = Path("logs") - - # Use the search module for log filtering - log_entries = search_logs( - logs_dir=logs_dir, - date=date, - level=level, - request_id=request_id, - search_text=search, - limit=limit, - ) - - return { - "logs": log_entries, - "total": len(log_entries), - "date": date, - "level": level, - "request_id": request_id, - "search": search, - "limit": limit, - } - - -@admin_router.get("/api/logs/dates", dependencies=[Depends(require_admin_api)]) -async def get_log_dates_api(request: Request) -> dict[str, object]: - logs_dir = Path("logs") - dates = [] - - if logs_dir.exists(): - log_files = sorted( - logs_dir.glob("app_*.log"), key=lambda x: x.stat().st_mtime, reverse=True - ) - - for log_file in log_files[:30]: - try: - filename = log_file.name - date_str = filename.replace("app_", "").replace(".log", "") - dates.append(date_str) - except Exception: - continue - - return {"dates": dates} - - @admin_router.post("/withdraw", dependencies=[Depends(require_admin_api)]) async def withdraw( request: Request, withdraw_request: WithdrawRequest @@ -2504,32 +2406,6 @@ def upstream_providers_page() -> str: ) -def logs_page() -> str: - return """ - - -
- - -Redirecting to logs page...
- - - """ - - -@admin_router.get("/logs", response_class=HTMLResponse) -async def admin_logs(request: Request) -> str: - if is_admin_authenticated(request): - return logs_page() - return admin_auth() - - @admin_router.get("/upstream-providers", response_class=HTMLResponse) async def admin_upstream_providers(request: Request) -> str: if is_admin_authenticated(request): @@ -2991,3 +2867,123 @@ h1 { color: #333; } .no-logs { text-align: center; color: #666; padding: 40px; } .request-id-display { background-color: #e9ecef; padding: 10px; border-radius: 4px; margin-bottom: 20px; font-family: monospace; } """ + + + +@admin_router.get("/api/usage/metrics", dependencies=[Depends(require_admin_api)]) +async def get_usage_metrics( + request: Request, + interval: int = Query( + default=15, ge=1, le=1440, description="Time interval in minutes" + ), + hours: int = Query( + default=24, ge=1, le=168, 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/summary", dependencies=[Depends(require_admin_api)]) +async def get_usage_summary( + request: Request, + hours: int = Query( + default=24, ge=1, le=168, description="Hours of history to analyze" + ), +) -> dict: + """Get summary statistics for the specified time period.""" + return log_manager.get_usage_summary(hours=hours) + + +@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, le=168, description="Hours of history to analyze" + ), + limit: int = Query( + default=100, ge=1, le=1000, description="Maximum number of errors to return" + ), +) -> dict: + """Get detailed error information.""" + return log_manager.get_error_details(hours=hours, limit=limit) + + +@admin_router.get( + "/api/usage/revenue-by-model", dependencies=[Depends(require_admin_api)] +) +async def get_revenue_by_model( + request: Request, + hours: int = Query( + default=24, ge=1, le=168, description="Hours of history to analyze" + ), + limit: int = Query( + default=20, ge=1, le=100, description="Maximum number of models to return" + ), +) -> dict: + """ + Get revenue breakdown by model. + """ + return log_manager.get_revenue_by_model(hours=hours, limit=limit) + + +@admin_router.get("/api/logs", dependencies=[Depends(require_admin_api)]) +async def get_logs_api( + request: Request, + date: str | None = None, + level: str | None = None, + request_id: str | None = None, + search: str | None = None, + limit: int = 100, +) -> dict[str, object]: + """ + Get filtered log entries. + + Args: + date: Filter by specific date (YYYY-MM-DD) + level: Filter by log level + request_id: Filter by request ID + search: Search text in message and name fields (case-insensitive) + limit: Maximum number of entries to return + + Returns: + Dict containing logs and filter metadata + """ + log_entries = log_manager.search_logs( + date=date, + level=level, + request_id=request_id, + search_text=search, + limit=limit, + ) + + return { + "logs": log_entries, + "total": len(log_entries), + "date": date, + "level": level, + "request_id": request_id, + "search": search, + "limit": limit, + } + + +@admin_router.get("/api/logs/dates", dependencies=[Depends(require_admin_api)]) +async def get_log_dates_api(request: Request) -> dict[str, object]: + logs_dir = Path("logs") + dates = [] + + if logs_dir.exists(): + log_files = sorted( + logs_dir.glob("app_*.log"), key=lambda x: x.stat().st_mtime, reverse=True + ) + + for log_file in log_files[:30]: + try: + filename = log_file.name + date_str = filename.replace("app_", "").replace(".log", "") + dates.append(date_str) + except Exception: + continue + + return {"dates": dates} diff --git a/routstr/core/log_manager.py b/routstr/core/log_manager.py new file mode 100644 index 00000000..557a349a --- /dev/null +++ b/routstr/core/log_manager.py @@ -0,0 +1,434 @@ +import json +from collections import defaultdict +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Any, Generator, Iterator + +from .logging import get_logger + +logger = get_logger(__name__) + + +class LogManager: + def __init__(self, logs_dir: Path = Path("logs")): + self.logs_dir = logs_dir + + def _yield_log_entries( + self, + hours_back: int | None = None, + specific_date: str | None = None, + reverse_files: bool = False + ) -> Iterator[dict[str, Any]]: + """ + Yields log entries from files. + + Args: + hours_back: specific number of hours to look back. + specific_date: specific date string (YYYY-MM-DD) to look at. + reverse_files: if True, process files in reverse order (newest first). + """ + if not self.logs_dir.exists(): + return + + log_files = [] + cutoff_date = None + + if specific_date: + log_file = self.logs_dir / f"app_{specific_date}.log" + if log_file.exists(): + log_files.append(log_file) + else: + log_files = sorted(self.logs_dir.glob("app_*.log")) + if reverse_files: + log_files.reverse() + + # 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) + filtered_files = [] + for f in log_files: + try: + file_date_str = f.stem.split("_")[1] + file_date = datetime.strptime(file_date_str, "%Y-%m-%d").replace( + tzinfo=timezone.utc + ) + # Include file if it's from the same day or after the cutoff day + if file_date >= cutoff_date.replace(hour=0, minute=0, second=0, microsecond=0): + filtered_files.append(f) + except Exception: + continue + log_files = filtered_files + + 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() + + for line in lines: + try: + entry = json.loads(line.strip()) + + if cutoff_date: + timestamp_str = entry.get("asctime", "") + if not timestamp_str: + 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: + continue + + yield entry + except json.JSONDecodeError: + continue + except Exception as e: + logger.error(f"Error processing log file {log_file}: {e}") + continue + + def search_logs( + self, + date: str | None = None, + level: str | None = None, + request_id: str | None = None, + search_text: str | None = None, + limit: int = 100, + ) -> list[dict[str, Any]]: + """ + Search through log files and return matching entries. + """ + log_entries: list[dict[str, Any]] = [] + + # Use reverse=True to get newest logs first by default + # If date is specified, we only look at that file + + search_text_lower = search_text.lower() if search_text else None + + # We iterate efficiently + iterator = self._yield_log_entries(specific_date=date, reverse_files=True if not date else False) + + # If we are searching globally (no date), we might want to limit how far back we go? + # PR 228 did: "glob("app_*.log") sorted by mtime reverse [:7]" (last 7 files) + # My _yield_log_entries with reverse_files=True does all files. + # Let's rely on limit to stop us. + + files_processed = 0 + # Optimization: if we are not searching by date, maybe limit to last 7 files inside _yield? + # For now, let's just iterate. + + for log_data in iterator: + if not self._matches_filters(log_data, level, request_id, search_text_lower): + continue + + log_entries.append(log_data) + + if len(log_entries) >= limit: + break + + # Sort by time descending (newest first) + log_entries.sort(key=lambda x: x.get("asctime", ""), reverse=True) + return log_entries + + def _matches_filters( + self, + log_data: dict[str, Any], + level: str | None, + request_id: str | None, + search_text_lower: str | None, + ) -> bool: + if level and log_data.get("levelname", "").upper() != level.upper(): + return False + + if request_id and log_data.get("request_id") != request_id: + return False + + if search_text_lower: + message = str(log_data.get("message", "")).lower() + name = str(log_data.get("name", "")).lower() + pathname = str(log_data.get("pathname", "")).lower() + + if ( + search_text_lower not in message + and search_text_lower not in name + and search_text_lower not in pathname + ): + return False + + return True + + 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", ""), + }) + + # 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, + } + ) + + for entry in entries: + try: + model = entry.get("model", "unknown") + if not isinstance(model, str): + model = "unknown" + + message = entry.get("message", "").lower() + + if "received proxy request" in message: + model_stats[model]["requests"] += 1 + + if "token adjustment completed" 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 + + 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 + + except Exception: + continue + + models = [] + total_revenue = 0.0 + + for model, stats in model_stats.items(): + revenue_msats = float(stats["revenue_msats"]) + refunds_msats = float(stats["refunds_msats"]) + + revenue_sats = revenue_msats / 1000 + refunds_sats = refunds_msats / 1000 + net_revenue_sats = revenue_sats - refunds_sats + + total_revenue += net_revenue_sats + + requests = int(stats["requests"]) + successful = int(stats["successful"]) + + models.append({ + "model": model, + "revenue_sats": revenue_sats, + "refunds_sats": refunds_sats, + "net_revenue_sats": net_revenue_sats, + "requests": requests, + "successful": successful, + "failed": int(stats["failed"]), + "avg_revenue_per_request": ( + revenue_sats / successful if successful > 0 else 0 + ), + }) + + models.sort(key=lambda x: float(x["net_revenue_sats"]), reverse=True) + + return { + "models": models[:limit], + "total_revenue_sats": total_revenue, + "total_models": len(models), + } + + 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, + } + + for entry in entries: + try: + stats["total_entries"] += 1 + + message = entry.get("message", "").lower() + level = 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 + + if "received proxy request" in message: + stats["total_requests"] += 1 + + if "token adjustment completed" in message: + stats["successful_chat_completions"] += 1 + + if "upstream request failed" in message or "revert payment" in message: + 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 "token adjustment completed" 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) + + 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 + + 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"] + + return { + "total_entries": stats["total_entries"], + "total_requests": total_requests, + "successful_chat_completions": successful, + "failed_requests": stats["failed_requests"], + "total_errors": stats["total_errors"], + "total_warnings": stats["total_warnings"], + "payment_processed": stats["payment_processed"], + "upstream_errors": stats["upstream_errors"], + "unique_models_count": len(stats["unique_models"]), + "unique_models": sorted(list(stats["unique_models"])), + "error_types": dict(stats["error_types"]), + "success_rate": (successful / total_requests * 100) if total_requests > 0 else 0, + "revenue_msats": stats["revenue_msats"], + "refunds_msats": stats["refunds_msats"], + "revenue_sats": revenue_sats, + "refunds_sats": refunds_sats, + "net_revenue_msats": stats["revenue_msats"] - stats["refunds_msats"], + "net_revenue_sats": net_revenue_sats, + "avg_revenue_per_request_msats": ( + stats["revenue_msats"] / successful if successful > 0 else 0 + ), + "refund_rate": ( + (stats["failed_requests"] / total_requests * 100) if total_requests > 0 else 0 + ), + } + + def _aggregate_metrics_by_time( + self, entries: list[dict], interval_minutes: int, hours_back: int + ) -> dict: + time_buckets = defaultdict(lambda: { + "requests": 0, + "errors": 0, + "revenue_msats": 0.0 + }) + + for entry in entries: + try: + timestamp_str = entry.get("asctime", "") + if not timestamp_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 = bucket_time.strftime("%Y-%m-%d %H:%M:%S") + + bucket = time_buckets[bucket_key] + + message = entry.get("message", "").lower() + level = entry.get("levelname", "").upper() + + if "received proxy request" in message: + bucket["requests"] += 1 + + if level == "ERROR": + bucket["errors"] += 1 + + if "token adjustment completed" 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) + except Exception: + continue + + result = [] + for bucket_key in sorted(time_buckets.keys()): + result.append({"timestamp": bucket_key, **time_buckets[bucket_key]}) + + return { + "metrics": result, + "interval_minutes": interval_minutes, + "hours_back": hours_back, + "total_buckets": len(result), + } + +log_manager = LogManager() diff --git a/routstr/core/logging.py b/routstr/core/logging.py index d682b944..6dbd2cfc 100644 --- a/routstr/core/logging.py +++ b/routstr/core/logging.py @@ -1,3 +1,40 @@ +""" +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). +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. "Token adjustment completed for streaming" (INFO) - routstr/upstream/base.py + "Token adjustment completed for non-streaming" (INFO) - routstr/upstream/base.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 + +3. "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 + - 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 + - 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() +""" + import logging.config import logging.handlers import os diff --git a/routstr/core/main.py b/routstr/core/main.py index b88be3cc..320eb21a 100644 --- a/routstr/core/main.py +++ b/routstr/core/main.py @@ -273,6 +273,15 @@ if UI_DIST_PATH.exists() and UI_DIST_PATH.is_dir(): async def redirect_logs_index_txt() -> RedirectResponse: return RedirectResponse("/logs") + @app.get("/usage", include_in_schema=False) + async def serve_usage_ui() -> FileResponse: + return FileResponse(UI_DIST_PATH / "usage" / "index.html") + + # Add explicit route for /usage/index.txt to redirect to /usage + @app.get("/usage/index.txt", include_in_schema=False) + async def redirect_usage_index_txt() -> RedirectResponse: + return RedirectResponse("/usage") + @app.get("/unauthorized", include_in_schema=False) async def serve_unauthorized_ui() -> FileResponse: return FileResponse(UI_DIST_PATH / "unauthorized" / "index.html") diff --git a/routstr/core/usage_metrics.py b/routstr/core/usage_metrics.py deleted file mode 100644 index 6ab243b5..00000000 --- a/routstr/core/usage_metrics.py +++ /dev/null @@ -1,264 +0,0 @@ -from __future__ import annotations - -import asyncio -import json -import math -from dataclasses import asdict, dataclass -from datetime import datetime, timedelta -from pathlib import Path -from typing import Callable, Literal, cast - -UsageMetricName = Literal["errors", "chat_completions_success"] - - -@dataclass(frozen=True) -class MetricDefinition: - name: UsageMetricName - label: str - description: str - matcher: Callable[[dict[str, object]], bool] - - -@dataclass(frozen=True) -class UsageMetricPoint: - bucket_start: datetime - count: int - - -@dataclass(frozen=True) -class UsageMetricSeries: - name: UsageMetricName - label: str - description: str - total: int - points: list[UsageMetricPoint] - - -@dataclass(frozen=True) -class UsageMetricsComputation: - bucket_minutes: int - bucket_count: int - start: datetime - end: datetime - series: list[UsageMetricSeries] - - def to_dict(self) -> dict[str, object]: - return asdict(self) - - -def _is_error(record: dict[str, object]) -> bool: - level = record.get("levelname") - if not isinstance(level, str): - return False - return level.upper() in {"ERROR", "CRITICAL"} - - -def _coerce_int(value: object) -> int | None: - if isinstance(value, int): - return value - if isinstance(value, str): - try: - return int(value.strip()) - except ValueError: - return None - return None - - -def _is_successful_chat_completion(record: dict[str, object]) -> bool: - message = record.get("message") - if not isinstance(message, str) or message != "Received upstream response": - return False - status_code = _coerce_int(record.get("status_code")) - if status_code != 200: - return False - path = record.get("path") - return isinstance(path, str) and path.endswith("chat/completions") - - -METRIC_DEFINITIONS: dict[UsageMetricName, MetricDefinition] = { - "errors": MetricDefinition( - name="errors", - label="Errors", - description="Log entries emitted at ERROR or CRITICAL level", - matcher=_is_error, - ), - "chat_completions_success": MetricDefinition( - name="chat_completions_success", - label="200 chat/completions", - description="Successful upstream responses for chat/completions", - matcher=_is_successful_chat_completion, - ), -} - - -def list_metric_definitions() -> list[dict[str, str]]: - return [ - { - "name": definition.name, - "label": definition.label, - "description": definition.description, - } - for definition in METRIC_DEFINITIONS.values() - ] - - -def default_metric_names() -> list[UsageMetricName]: - return list(METRIC_DEFINITIONS.keys()) - - -class UsageMetricsService: - @classmethod - async def collect( - cls, - metrics: list[str], - bucket_minutes: int, - hours: int, - log_dir: Path | None = None, - now: datetime | None = None, - ) -> UsageMetricsComputation: - return await asyncio.to_thread( - cls._collect_sync, metrics, bucket_minutes, hours, log_dir, now - ) - - @classmethod - def _collect_sync( - cls, - metrics: list[str], - bucket_minutes: int, - hours: int, - log_dir: Path | None, - now: datetime | None, - ) -> UsageMetricsComputation: - metric_list = cls._sanitize_metrics(metrics) - bucket_minutes = cls._clamp(bucket_minutes, minimum=1, maximum=24 * 60) - hours = cls._clamp(hours, minimum=1, maximum=24 * 7) - now_dt = now or datetime.now() - start = now_dt - timedelta(hours=hours) - bucket_seconds = bucket_minutes * 60 - bucket_count = max(1, math.ceil((hours * 3600) / bucket_seconds)) - bucket_starts = [ - start + timedelta(seconds=bucket_seconds * index) - for index in range(bucket_count) - ] - counts: dict[UsageMetricName, list[int]] = { - name: [0] * bucket_count for name in metric_list - } - - for log_file in cls._iter_log_files(log_dir): - try: - with log_file.open("r", encoding="utf-8") as handle: - for line in handle: - record = cls._parse_record(line) - if not record: - continue - timestamp = cls._parse_timestamp(record.get("asctime")) - if timestamp is None or timestamp < start or timestamp > now_dt: - continue - bucket_index = cls._bucket_index( - timestamp, start, bucket_seconds, bucket_count - ) - if bucket_index is None: - continue - for metric_name in metric_list: - definition = METRIC_DEFINITIONS[metric_name] - if definition.matcher(record): - counts[metric_name][bucket_index] += 1 - except (OSError, UnicodeDecodeError): - continue - - series = [ - UsageMetricSeries( - name=metric_name, - label=METRIC_DEFINITIONS[metric_name].label, - description=METRIC_DEFINITIONS[metric_name].description, - total=sum(counts[metric_name]), - points=[ - UsageMetricPoint(bucket_start=bucket_starts[index], count=count) - for index, count in enumerate(counts[metric_name]) - ], - ) - for metric_name in metric_list - ] - - return UsageMetricsComputation( - bucket_minutes=bucket_minutes, - bucket_count=bucket_count, - start=start, - end=now_dt, - series=series, - ) - - @staticmethod - def _sanitize_metrics(metrics: list[str]) -> list[UsageMetricName]: - requested = [metric.strip() for metric in metrics if metric.strip()] - unique: list[UsageMetricName] = [] - for normalized in requested: - if normalized not in METRIC_DEFINITIONS: - continue - metric_name = cast(UsageMetricName, normalized) - if metric_name not in unique: - unique.append(metric_name) - if unique: - return unique - if requested: - raise ValueError("No valid usage metrics requested") - return default_metric_names() - - @staticmethod - def _clamp(value: int, *, minimum: int, maximum: int) -> int: - return max(minimum, min(maximum, value)) - - @staticmethod - def _iter_log_files(log_dir: Path | None) -> list[Path]: - directory = log_dir or Path("logs") - if not directory.exists(): - return [] - log_files = sorted( - directory.glob("*.log"), - key=UsageMetricsService._safe_mtime, - reverse=True, - ) - return log_files - - @staticmethod - def _safe_mtime(path: Path) -> float: - try: - return path.stat().st_mtime - except OSError: - return 0.0 - - @staticmethod - def _parse_record(line: str) -> dict[str, object] | None: - stripped = line.strip() - if not stripped: - return None - try: - data = json.loads(stripped) - except json.JSONDecodeError: - return None - return data if isinstance(data, dict) else None - - @staticmethod - def _parse_timestamp(value: object) -> datetime | None: - if not isinstance(value, str): - return None - try: - return datetime.strptime(value, "%Y-%m-%d %H:%M:%S") - except ValueError: - return None - - @staticmethod - def _bucket_index( - timestamp: datetime, - start: datetime, - bucket_seconds: int, - bucket_count: int, - ) -> int | None: - delta_seconds = (timestamp - start).total_seconds() - if delta_seconds < 0: - return None - index = int(delta_seconds // bucket_seconds) - if index >= bucket_count: - index = bucket_count - 1 - return index - diff --git a/routstr/search/__init__.py b/routstr/search/__init__.py deleted file mode 100644 index 963cd1b4..00000000 --- a/routstr/search/__init__.py +++ /dev/null @@ -1,3 +0,0 @@ -from .log_search import search_logs - -__all__ = ["search_logs"] diff --git a/routstr/search/log_search.py b/routstr/search/log_search.py deleted file mode 100644 index ff69ef3a..00000000 --- a/routstr/search/log_search.py +++ /dev/null @@ -1,108 +0,0 @@ -import json -from pathlib import Path -from typing import Any - - -def search_logs( - logs_dir: Path, - date: str | None = None, - level: str | None = None, - request_id: str | None = None, - search_text: str | None = None, - limit: int = 100, -) -> list[dict[str, Any]]: - """ - Search through log files and return matching entries. - Args: - logs_dir: Path to the logs directory - date: Filter by specific date (YYYY-MM-DD format) - level: Filter by log level (INFO, WARNING, ERROR, etc.) - request_id: Filter by exact request ID match - search_text: Search in message and name fields (case-insensitive) - limit: Maximum number of entries to return - - Returns: - List of log entries matching the criteria - """ - log_entries: list[dict[str, Any]] = [] - - if not logs_dir.exists(): - return log_entries - - log_files = [] - if date: - log_file = logs_dir / f"app_{date}.log" - if log_file.exists(): - log_files.append(log_file) - else: - log_files = sorted( - logs_dir.glob("app_*.log"), - key=lambda x: x.stat().st_mtime, - reverse=True, - )[:7] - - search_text_lower = search_text.lower() if search_text else None - - for log_file in log_files: - try: - with open(log_file, "r") as f: - for line in f: - try: - log_data = json.loads(line.strip()) - - if not _matches_filters( - log_data, level, request_id, search_text_lower - ): - continue - - log_entries.append(log_data) - - if len(log_entries) >= limit: - break - - except json.JSONDecodeError: - continue - - if len(log_entries) >= limit: - break - - except Exception: - continue - - log_entries.sort(key=lambda x: x.get("asctime", ""), reverse=True) - - return log_entries - - -def _matches_filters( - log_data: dict[str, Any], - level: str | None, - request_id: str | None, - search_text_lower: str | None, -) -> bool: - """ - Check if a log entry matches the given filters. - - Args: - log_data: The log entry to check - level: Log level filter (if any) - request_id: Request ID filter (if any) - search_text_lower: Lowercase search text (if any) - - Returns: - True if the log entry matches all filters, False otherwise - """ - if level and log_data.get("levelname", "").upper() != level.upper(): - return False - - if request_id and log_data.get("request_id") != request_id: - return False - - if search_text_lower: - message = str(log_data.get("message", "")).lower() - name = str(log_data.get("name", "")).lower() - - if search_text_lower not in message and search_text_lower not in name: - return False - - return True diff --git a/tests/unit/test_usage_metrics.py b/tests/unit/test_usage_metrics.py deleted file mode 100644 index 68f6afde..00000000 --- a/tests/unit/test_usage_metrics.py +++ /dev/null @@ -1,79 +0,0 @@ -from __future__ import annotations - -import json -from datetime import datetime, timedelta -from pathlib import Path - -import pytest - -from routstr.core.usage_metrics import UsageMetricsService - - -def _write_records(path: Path, records: list[dict[str, object]]) -> None: - with path.open("w", encoding="utf-8") as handle: - for record in records: - handle.write(json.dumps(record) + "\n") - - -@pytest.mark.asyncio -async def test_usage_metrics_counts(tmp_path: Path) -> None: - log_dir = tmp_path / "logs" - log_dir.mkdir() - now = datetime(2025, 1, 1, 12, 0, 0) - records = [ - { - "asctime": (now - timedelta(minutes=5)).strftime("%Y-%m-%d %H:%M:%S"), - "levelname": "ERROR", - "message": "Proxy failure", - }, - { - "asctime": (now - timedelta(minutes=10)).strftime("%Y-%m-%d %H:%M:%S"), - "levelname": "INFO", - "message": "Received upstream response", - "status_code": 200, - "path": "chat/completions", - }, - { - "asctime": (now - timedelta(minutes=20)).strftime("%Y-%m-%d %H:%M:%S"), - "levelname": "INFO", - "message": "Received upstream response", - "status_code": 500, - "path": "chat/completions", - }, - ] - _write_records(log_dir / "app_2025-01-01.log", records) - - result = await UsageMetricsService.collect( - metrics=["errors", "chat_completions_success"], - bucket_minutes=15, - hours=1, - log_dir=log_dir, - now=now, - ) - - errors_series = next(series for series in result.series if series.name == "errors") - completions_series = next( - series - for series in result.series - if series.name == "chat_completions_success" - ) - - assert errors_series.total == 1 - assert sum(point.count for point in errors_series.points) == 1 - assert completions_series.total == 1 - assert any(point.count == 1 for point in completions_series.points) - - -@pytest.mark.asyncio -async def test_usage_metrics_invalid_metric(tmp_path: Path) -> None: - log_dir = tmp_path / "logs" - log_dir.mkdir() - now = datetime(2025, 1, 1, 12, 0, 0) - with pytest.raises(ValueError): - await UsageMetricsService.collect( - metrics=["unknown"], - bucket_minutes=15, - hours=1, - log_dir=log_dir, - now=now, - ) diff --git a/ui/app/page.tsx b/ui/app/page.tsx index 5e7d6514..da8ffba3 100644 --- a/ui/app/page.tsx +++ b/ui/app/page.tsx @@ -10,7 +10,6 @@ import { TemporaryBalances } from '@/components/temporary-balances'; import { ToggleGroup, ToggleGroupItem } from '@/components/ui/toggle-group'; import type { DisplayUnit } from '@/lib/types/units'; import { fetchBtcUsdPrice, btcToSatsRate } from '@/lib/exchange-rate'; -import { UsageTracking } from '@/components/usage-tracking'; export default function Page() { const [displayUnit, setDisplayUnit] = useState+ Monitor system usage, requests, and errors over time +
++ No metrics data found for the selected time range. This + could be because no requests have been logged yet or the + log files are not available. +
++ No errors found in the selected time period +
++ Total Revenue: {totalRevenue.toLocaleString(undefined, { maximumFractionDigits: 2 })} sats +
+- total events in range -
-No data available for the selected window.
-