deerflow-code/offline-backend-20260512/backend/app/gateway/routers/llm_metrics.py
2026-09-07 18:24:55 +08:00

347 lines
11 KiB
Python

"""Admin-only LLM call metrics: list + Excel export.
Lets administrators audit every provider invocation made by the runtime —
including failures — and download the full table as ``.xlsx``. The recording
is unconditional on the backend, so this view always reflects every call.
"""
from __future__ import annotations
import logging
from datetime import datetime
from io import BytesIO
from typing import Any
from urllib.parse import quote
from fastapi import APIRouter, HTTPException, Query, Request
from fastapi.responses import StreamingResponse
from pydantic import BaseModel, Field
from app.gateway.deps import get_llm_metrics_store, get_local_provider, get_optional_user_from_request
from deerflow.persistence.llm_metrics.base import LlmMetricsQuery, LlmMetricsStore
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/admin/llm-metrics", tags=["llm-metrics"])
class LlmMetricResponse(BaseModel):
id: str
created_at: str | None = None
user_id: str | None = None
user_name: str | None = Field(
default=None,
description="Local part (before @) of the user's email, resolved from user_id",
)
thread_id: str | None = None
run_id: str | None = None
agent_name: str | None = None
model_name: str | None = None
duration_ms: int = 0
input_tokens: int | None = None
output_tokens: int | None = None
total_tokens: int | None = None
tokens_per_sec: float | None = None
status: str = "success"
error_type: str | None = None
error_message: str | None = None
class LlmMetricsListResponse(BaseModel):
items: list[LlmMetricResponse]
total: int = Field(..., description="Total rows matching the filter, ignoring pagination")
limit: int
offset: int
async def _require_admin(request: Request) -> None:
user = await get_optional_user_from_request(request)
if user is None:
return
if getattr(user, "system_role", None) != "admin":
raise HTTPException(status_code=403, detail="LLM metrics admin view is admin-only")
def _require_store(request: Request) -> LlmMetricsStore:
store = get_llm_metrics_store(request)
if store is None:
raise HTTPException(status_code=503, detail="LLM metrics store is not configured")
return store
def _parse_dt(value: str | None) -> datetime | None:
if not value:
return None
try:
# Accept both ``Z`` suffix and ``+00:00``.
return datetime.fromisoformat(value.replace("Z", "+00:00"))
except ValueError as exc:
raise HTTPException(status_code=422, detail=f"Invalid datetime: {value!r} ({exc})") from exc
def _email_local_part(email: str | None) -> str | None:
if not email:
return None
at = email.find("@")
return email[:at] if at > 0 else email
def _derive_tokens_per_sec(row: dict[str, Any]) -> None:
"""Fill in ``tokens_per_sec`` when the row was recorded before we knew how
to extract ``output_tokens`` from response_metadata.token_usage.
Uses ``output_tokens / elapsed`` when available, otherwise falls back to
``total_tokens / elapsed`` as an approximation. Newer rows already have
the field set during recording, so this only kicks in for legacy data.
"""
if not isinstance(row, dict):
return
if row.get("tokens_per_sec") not in (None, 0):
return
duration_ms = row.get("duration_ms") or 0
if duration_ms <= 0:
return
numerator = row.get("output_tokens") or row.get("total_tokens")
if not numerator:
return
row["tokens_per_sec"] = round((numerator / duration_ms) * 1000, 3)
async def _resolve_user_name_to_id(name: str | None) -> str | None:
"""Look up a UUID for the given username (email local-part)."""
if not name:
return None
try:
provider = get_local_provider()
except Exception:
return None
try:
user = await provider.find_user_by_email_local_part(name)
except Exception:
return None
return str(user.id) if user is not None else None
async def _resolve_user_names_to_ids(name: str | None) -> list[str]:
"""Fuzzy-resolve a username fragment to every matching account's UUID.
Powers the 按用户名模糊搜索 filter: typing part of a username matches all
users whose email local part contains it. Returns an empty list when the
fragment matches nobody, so the caller can short-circuit to zero rows.
"""
if not name or not name.strip():
return []
try:
provider = get_local_provider()
except Exception:
return []
try:
users = await provider.find_users_by_email_local_part_like(name.strip())
except Exception:
return []
return [str(u.id) for u in users]
async def _attach_usernames(rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
"""Resolve each row's ``user_id`` UUID to a friendly ``user_name``.
Batches lookups so one query per distinct user_id, not per row. Tolerates
missing accounts (deleted users) by leaving ``user_name`` blank.
"""
if not rows:
return rows
user_ids = {row["user_id"] for row in rows if isinstance(row, dict) and row.get("user_id")}
if not user_ids:
return rows
try:
provider = get_local_provider()
except Exception:
provider = None
email_by_id: dict[str, str] = {}
if provider is not None:
for uid in user_ids:
try:
user = await provider.get_user(uid)
except Exception:
user = None
if user is not None and user.email:
email_by_id[uid] = user.email
for row in rows:
if isinstance(row, dict) and row.get("user_id"):
email = email_by_id.get(row["user_id"])
row["user_name"] = _email_local_part(email) or row["user_id"]
return rows
def _build_query(
*,
user_id: str | None,
model_name: str | None,
status: str | None,
since: str | None,
until: str | None,
limit: int,
offset: int,
user_ids: list[str] | None = None,
) -> LlmMetricsQuery:
return LlmMetricsQuery(
user_id=user_id or None,
user_ids=user_ids or None,
model_name=model_name or None,
status=status or None,
since=_parse_dt(since),
until=_parse_dt(until),
limit=limit,
offset=offset,
)
@router.get("", response_model=LlmMetricsListResponse)
async def list_llm_metrics(
request: Request,
user_id: str | None = Query(default=None),
user_name: str | None = Query(default=None, description="Username (email local part) to filter by"),
model_name: str | None = Query(default=None),
status: str | None = Query(default=None, pattern="^(success|error)$"),
since: str | None = Query(default=None, description="ISO 8601 datetime, inclusive lower bound"),
until: str | None = Query(default=None, description="ISO 8601 datetime, inclusive upper bound"),
limit: int = Query(default=50, ge=1, le=500),
offset: int = Query(default=0, ge=0),
) -> LlmMetricsListResponse:
await _require_admin(request)
store = _require_store(request)
user_ids: list[str] | None = None
if not user_id and user_name:
user_ids = await _resolve_user_names_to_ids(user_name)
if not user_ids:
# No username matches → no rows; short-circuit the global scan.
return LlmMetricsListResponse(items=[], total=0, limit=limit, offset=offset)
query = _build_query(
user_id=user_id,
user_ids=user_ids,
model_name=model_name,
status=status,
since=since,
until=until,
limit=limit,
offset=offset,
)
rows = await store.list(query)
total = await store.count(query)
await _attach_usernames(rows)
for row in rows:
_derive_tokens_per_sec(row)
return LlmMetricsListResponse(
items=[LlmMetricResponse(**row) for row in rows],
total=total,
limit=limit,
offset=offset,
)
def _format_cell(value: Any) -> Any:
if value is None:
return ""
if isinstance(value, datetime):
return value.replace(tzinfo=None) if value.tzinfo else value
return value
@router.get("/export")
async def export_llm_metrics(
request: Request,
user_id: str | None = Query(default=None),
user_name: str | None = Query(default=None),
model_name: str | None = Query(default=None),
status: str | None = Query(default=None, pattern="^(success|error)$"),
since: str | None = Query(default=None),
until: str | None = Query(default=None),
):
"""Stream the full filtered metric set as an ``.xlsx`` file."""
await _require_admin(request)
store = _require_store(request)
user_ids: list[str] | None = None
if not user_id and user_name:
user_ids = await _resolve_user_names_to_ids(user_name)
if not user_ids:
# No matching user — sentinel id that matches nothing, so the export
# still produces an empty xlsx (matching the list view's empty state).
user_ids = ["__no_such_user__"]
try:
from openpyxl import Workbook
except ImportError as exc:
raise HTTPException(status_code=500, detail=f"openpyxl is not installed: {exc}") from exc
query = _build_query(
user_id=user_id,
user_ids=user_ids,
model_name=model_name,
status=status,
since=since,
until=until,
limit=1, # unused: iter_all paginates internally
offset=0,
)
workbook = Workbook(write_only=True)
sheet = workbook.create_sheet("llm_call_metrics")
columns = [
"id",
"created_at",
"status",
"user_name",
"user_id",
"thread_id",
"run_id",
"agent_name",
"model_name",
"duration_ms",
"input_tokens",
"output_tokens",
"total_tokens",
"tokens_per_sec",
"error_type",
"error_message",
]
sheet.append(columns)
# Pre-resolve all distinct user_ids in chunks so the export doesn't pay one
# round trip per row.
buffered: list[dict[str, Any]] = []
FLUSH_AT = 200
async def flush() -> int:
if not buffered:
return 0
await _attach_usernames(buffered)
for row in buffered:
_derive_tokens_per_sec(row)
sheet.append([_format_cell(row.get(col)) for col in columns])
n = len(buffered)
buffered.clear()
return n
count = 0
async for row in store.iter_all(query):
buffered.append(row)
if len(buffered) >= FLUSH_AT:
count += await flush()
count += await flush()
buffer = BytesIO()
workbook.save(buffer)
buffer.seek(0)
logger.info("Exported %d llm_call_metrics rows", count)
timestamp = datetime.utcnow().strftime("%Y%m%d-%H%M%S")
filename = f"llm_call_metrics-{timestamp}.xlsx"
headers = {
"Content-Disposition": f"attachment; filename=\"{filename}\"; filename*=UTF-8''{quote(filename)}",
}
return StreamingResponse(
buffer,
media_type="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
headers=headers,
)