From ca2e442656af716e712f82dcec61e1b30f44bf53 Mon Sep 17 00:00:00 2001 From: Shroominic Date: Sun, 16 Nov 2025 11:45:53 -0800 Subject: [PATCH] ruff fmt+lint --- routstr/core/admin.py | 269 ++++++++++++++++++++++++------------------ 1 file changed, 152 insertions(+), 117 deletions(-) diff --git a/routstr/core/admin.py b/routstr/core/admin.py index 79bd17de..a3690142 100644 --- a/routstr/core/admin.py +++ b/routstr/core/admin.py @@ -2903,7 +2903,7 @@ def _aggregate_metrics_by_time( """ now = datetime.now(timezone.utc) cutoff = now - timedelta(hours=hours_back) - + time_buckets: dict[str, dict[str, int | float]] = defaultdict( lambda: { "total_requests": 0, @@ -2917,70 +2917,70 @@ def _aggregate_metrics_by_time( "refunds_msats": 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) - + if log_time < cutoff: continue - + bucket_time = log_time.replace( minute=(log_time.minute // interval_minutes) * interval_minutes, second=0, microsecond=0, ) bucket_key = bucket_time.isoformat() - + message = entry.get("message", "").lower() level = entry.get("levelname", "").upper() - + if level == "ERROR": time_buckets[bucket_key]["errors"] += 1 elif level == "WARNING": time_buckets[bucket_key]["warnings"] += 1 - + if "received proxy request" in message: time_buckets[bucket_key]["total_requests"] += 1 - + if "token adjustment completed for non-streaming" in message: time_buckets[bucket_key]["successful_chat_completions"] += 1 elif "token adjustment completed for streaming" in message: time_buckets[bucket_key]["successful_chat_completions"] += 1 - + if "upstream request failed" in message or "revert payment" in message: time_buckets[bucket_key]["failed_requests"] += 1 - + if "payment processed successfully" in message: time_buckets[bucket_key]["payment_processed"] += 1 - + if "upstream" in message and level == "ERROR": time_buckets[bucket_key]["upstream_errors"] += 1 - + if "token adjustment completed" in message: cost_data = entry.get("cost_data") if isinstance(cost_data, dict): actual_cost = cost_data.get("actual_cost", 0) if isinstance(actual_cost, (int, float)) and actual_cost > 0: time_buckets[bucket_key]["revenue_msats"] += 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: time_buckets[bucket_key]["refunds_msats"] += max_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, @@ -2999,7 +2999,7 @@ def _get_summary_stats(entries: list[dict], hours_back: int = 24) -> dict: """ now = datetime.now(timezone.utc) cutoff = now - timedelta(hours=hours_back) - + stats = { "total_entries": 0, "total_requests": 0, @@ -3018,24 +3018,24 @@ def _get_summary_stats(entries: list[dict], hours_back: int = 24) -> dict: "net_revenue_msats": 0, "net_revenue_sats": 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) - + if log_time < cutoff: continue - + 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: @@ -3043,47 +3043,47 @@ def _get_summary_stats(entries: list[dict], hours_back: int = 24) -> dict: stats["error_types"][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("actual_cost", 0) if isinstance(actual_cost, (int, float)) and actual_cost > 0: stats["revenue_msats"] += 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"] += max_cost - + except Exception: continue - + stats["revenue_sats"] = stats["revenue_msats"] / 1000 stats["refunds_sats"] = stats["refunds_msats"] / 1000 stats["net_revenue_msats"] = stats["revenue_msats"] - stats["refunds_msats"] stats["net_revenue_sats"] = stats["net_revenue_msats"] / 1000 - + return { "total_entries": stats["total_entries"], "total_requests": stats["total_requests"], @@ -3123,13 +3123,17 @@ def _get_summary_stats(entries: list[dict], hours_back: int = 24) -> dict: @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"), + 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.""" logs_dir = Path("logs") all_entries: list[dict] = [] - + if not logs_dir.exists(): return { "metrics": [], @@ -3137,35 +3141,41 @@ async def get_usage_metrics( "hours_back": hours, "total_buckets": 0, } - + cutoff_date = datetime.now(timezone.utc) - timedelta(hours=hours) - + for log_file in sorted(logs_dir.glob("app_*.log")): try: file_date_str = log_file.stem.split("_")[1] - file_date = datetime.strptime(file_date_str, "%Y-%m-%d").replace(tzinfo=timezone.utc) - - if file_date < cutoff_date.replace(hour=0, minute=0, second=0, microsecond=0): + file_date = datetime.strptime(file_date_str, "%Y-%m-%d").replace( + tzinfo=timezone.utc + ) + + if file_date < cutoff_date.replace( + hour=0, minute=0, second=0, microsecond=0 + ): continue - + entries = _parse_log_file(log_file) all_entries.extend(entries) except Exception as e: logger.error(f"Error processing log file {log_file}: {e}") continue - + return _aggregate_metrics_by_time(all_entries, interval, 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"), + hours: int = Query( + default=24, ge=1, le=168, description="Hours of history to analyze" + ), ) -> dict: """Get summary statistics for the specified time period.""" logs_dir = Path("logs") all_entries: list[dict] = [] - + if not logs_dir.exists(): return { "total_entries": 0, @@ -3181,88 +3191,108 @@ async def get_usage_summary( "error_types": {}, "success_rate": 0, } - + cutoff_date = datetime.now(timezone.utc) - timedelta(hours=hours) - + for log_file in sorted(logs_dir.glob("app_*.log")): try: file_date_str = log_file.stem.split("_")[1] - file_date = datetime.strptime(file_date_str, "%Y-%m-%d").replace(tzinfo=timezone.utc) - - if file_date < cutoff_date.replace(hour=0, minute=0, second=0, microsecond=0): + file_date = datetime.strptime(file_date_str, "%Y-%m-%d").replace( + tzinfo=timezone.utc + ) + + if file_date < cutoff_date.replace( + hour=0, minute=0, second=0, microsecond=0 + ): continue - + entries = _parse_log_file(log_file) all_entries.extend(entries) except Exception as e: logger.error(f"Error processing log file {log_file}: {e}") continue - + return _get_summary_stats(all_entries, 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"), + 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.""" logs_dir = Path("logs") errors: list[dict] = [] - + if not logs_dir.exists(): return {"errors": [], "total_count": 0} - + cutoff_date = datetime.now(timezone.utc) - timedelta(hours=hours) - + for log_file in sorted(logs_dir.glob("app_*.log"), reverse=True): try: file_date_str = log_file.stem.split("_")[1] - file_date = datetime.strptime(file_date_str, "%Y-%m-%d").replace(tzinfo=timezone.utc) - - if file_date < cutoff_date.replace(hour=0, minute=0, second=0, microsecond=0): + file_date = datetime.strptime(file_date_str, "%Y-%m-%d").replace( + tzinfo=timezone.utc + ) + + if file_date < cutoff_date.replace( + hour=0, minute=0, second=0, microsecond=0 + ): continue - + entries = _parse_log_file(log_file) - + for entry in entries: if entry.get("levelname", "").upper() == "ERROR": timestamp_str = entry.get("asctime", "") if timestamp_str: 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: - errors.append({ - "timestamp": timestamp_str, - "message": entry.get("message", ""), - "error_type": entry.get("error_type", "unknown"), - "pathname": entry.get("pathname", ""), - "lineno": entry.get("lineno", 0), - "request_id": entry.get("request_id", ""), - }) - + errors.append( + { + "timestamp": 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(errors) >= limit: break - + except Exception as e: logger.error(f"Error processing log file {log_file}: {e}") continue - + if len(errors) >= limit: break - + errors.sort(key=lambda x: x["timestamp"], reverse=True) - + return {"errors": errors[:limit], "total_count": len(errors)} -@admin_router.get("/api/usage/revenue-by-model", dependencies=[Depends(require_admin_api)]) +@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"), + 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. @@ -3273,27 +3303,30 @@ async def get_revenue_by_model( """ logs_dir = Path("logs") all_entries: list[dict] = [] - + if not logs_dir.exists(): return {"models": [], "total_revenue_sats": 0} - + cutoff_date = datetime.now(timezone.utc) - timedelta(hours=hours) - + for log_file in sorted(logs_dir.glob("app_*.log")): try: file_date_str = log_file.stem.split("_")[1] - file_date = datetime.strptime(file_date_str, "%Y-%m-%d").replace(tzinfo=timezone.utc) - - if file_date < cutoff_date.replace(hour=0, minute=0, second=0, microsecond=0): + file_date = datetime.strptime(file_date_str, "%Y-%m-%d").replace( + tzinfo=timezone.utc + ) + + if file_date < cutoff_date.replace( + hour=0, minute=0, second=0, microsecond=0 + ): continue - + entries = _parse_log_file(log_file) all_entries.extend(entries) except Exception as e: logger.error(f"Error processing log file {log_file}: {e}") continue - - now = datetime.now(timezone.utc) + model_stats: dict[str, dict[str, int | float]] = defaultdict( lambda: { "revenue_msats": 0, @@ -3303,28 +3336,28 @@ async def get_revenue_by_model( "failed": 0, } ) - + for entry in all_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) - + if log_time < cutoff_date: continue - + 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") @@ -3332,42 +3365,44 @@ async def get_revenue_by_model( actual_cost = cost_data.get("actual_cost", 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 - + for model, stats in model_stats.items(): revenue_sats = stats["revenue_msats"] / 1000 refunds_sats = stats["refunds_msats"] / 1000 net_revenue_sats = revenue_sats - refunds_sats - + total_revenue += net_revenue_sats - - models.append({ - "model": model, - "revenue_sats": revenue_sats, - "refunds_sats": refunds_sats, - "net_revenue_sats": net_revenue_sats, - "requests": stats["requests"], - "successful": stats["successful"], - "failed": stats["failed"], - "avg_revenue_per_request": ( - revenue_sats / stats["successful"] if stats["successful"] > 0 else 0 - ), - }) - + + models.append( + { + "model": model, + "revenue_sats": revenue_sats, + "refunds_sats": refunds_sats, + "net_revenue_sats": net_revenue_sats, + "requests": stats["requests"], + "successful": stats["successful"], + "failed": stats["failed"], + "avg_revenue_per_request": ( + revenue_sats / stats["successful"] if stats["successful"] > 0 else 0 + ), + } + ) + models.sort(key=lambda x: x["net_revenue_sats"], reverse=True) - + return { "models": models[:limit], "total_revenue_sats": total_revenue,