347 lines
11 KiB
Python
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,
|
|
)
|