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

381 lines
12 KiB
Python

"""Admin-only tool / skill call metrics: list + Excel export.
Companion to ``llm_metrics`` — where that audits every LLM provider call, this
audits every *tool* call (built-in tools, skills, MCP tools, subagent calls),
including the failures: retrieval errors, blocked/rate-limited IPs, timeouts,
auth problems. Failures are bucketed into a coarse ``error_category`` so an
admin can scan for systemic issues.
"""
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_local_provider, get_optional_user_from_request, get_tool_metrics_store
from deerflow.persistence.tool_metrics.base import ToolMetricsQuery, ToolMetricsStore
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/admin/tool-metrics", tags=["tool-metrics"])
_VALID_CATEGORIES = (
"rate_limited",
"ip_blocked",
"auth_error",
"timeout",
"network_error",
"not_found",
"empty_result",
"tool_error",
)
def _kind_is_skill(kind: str | None) -> bool:
"""``kind=skill`` keeps only rows that target a specific Agent Skill."""
return kind == "skill"
class ToolMetricResponse(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
tool_name: str | None = None
skill_name: str | None = Field(
default=None,
description="The Agent Skill this call targets (skill tools only)",
)
duration_ms: int = 0
status: str = "success"
error_category: str | None = None
error_message: str | None = None
class ToolMetricsListResponse(BaseModel):
items: list[ToolMetricResponse]
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="Tool metrics admin view is admin-only")
def _require_store(request: Request) -> ToolMetricsStore:
store = get_tool_metrics_store(request)
if store is None:
raise HTTPException(status_code=503, detail="Tool metrics store is not configured")
return store
def _parse_dt(value: str | None) -> datetime | None:
if not value:
return None
try:
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
async def _resolve_user_name_to_id(name: str | None) -> str | None:
if not name:
return None
try:
provider = get_local_provider()
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 _attach_usernames(rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
"""Resolve each row's ``user_id`` UUID to a friendly ``user_name``."""
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,
tool_name: str | None,
skill_only: bool = False,
skill_name: str | None = None,
status: str | None,
error_category: str | None,
since: str | None,
until: str | None,
limit: int,
offset: int,
) -> ToolMetricsQuery:
return ToolMetricsQuery(
user_id=user_id or None,
tool_name=tool_name or None,
skill_only=skill_only,
skill_name=skill_name or None,
status=status or None,
error_category=error_category or None,
since=_parse_dt(since),
until=_parse_dt(until),
limit=limit,
offset=offset,
)
@router.get("", response_model=ToolMetricsListResponse)
async def list_tool_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"),
tool_name: str | None = Query(default=None),
kind: str | None = Query(default=None, description="Coarse call kind: 'skill' keeps only Agent-Skill invocations"),
skill_name: str | None = Query(default=None, description="Filter to one specific Agent Skill by name"),
status: str | None = Query(default=None, pattern="^(success|error)$"),
error_category: str | None = Query(default=None, description="Coarse failure bucket"),
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),
) -> ToolMetricsListResponse:
await _require_admin(request)
store = _require_store(request)
if error_category and error_category not in _VALID_CATEGORIES:
raise HTTPException(status_code=422, detail=f"Unknown error_category: {error_category!r}")
if not user_id and user_name:
resolved = await _resolve_user_name_to_id(user_name)
if resolved is None:
return ToolMetricsListResponse(items=[], total=0, limit=limit, offset=offset)
user_id = resolved
query = _build_query(
user_id=user_id,
tool_name=tool_name,
skill_only=_kind_is_skill(kind),
skill_name=skill_name,
status=status,
error_category=error_category,
since=since,
until=until,
limit=limit,
offset=offset,
)
rows = await store.list(query)
total = await store.count(query)
await _attach_usernames(rows)
return ToolMetricsListResponse(
items=[ToolMetricResponse(**row) for row in rows],
total=total,
limit=limit,
offset=offset,
)
class SkillStat(BaseModel):
skill_name: str
total: int
success: int
error: int
class ToolMetricsSummaryResponse(BaseModel):
total: int
error_total: int
by_category: dict[str, int] = Field(default_factory=dict)
by_skill: list[SkillStat] = Field(
default_factory=list,
description="Per-skill invocation counts — answers 'was skill X called, and did it succeed'",
)
@router.get("/summary", response_model=ToolMetricsSummaryResponse)
async def tool_metrics_summary(
request: Request,
user_id: str | None = Query(default=None),
user_name: str | None = Query(default=None),
tool_name: str | None = Query(default=None),
kind: str | None = Query(default=None),
skill_name: str | None = Query(default=None),
since: str | None = Query(default=None),
until: str | None = Query(default=None),
) -> ToolMetricsSummaryResponse:
"""Aggregate counts: total, errors, per-category breakdown, and per-skill stats."""
await _require_admin(request)
store = _require_store(request)
if not user_id and user_name:
resolved = await _resolve_user_name_to_id(user_name)
if resolved is None:
return ToolMetricsSummaryResponse(total=0, error_total=0, by_category={}, by_skill=[])
user_id = resolved
skill_only = _kind_is_skill(kind)
def _q(*, status: str | None = None, error_category: str | None = None) -> ToolMetricsQuery:
return _build_query(
user_id=user_id,
tool_name=tool_name,
skill_only=skill_only,
skill_name=skill_name,
status=status,
error_category=error_category,
since=since,
until=until,
limit=1,
offset=0,
)
total = await store.count(_q())
error_total = await store.count(_q(status="error"))
by_category: dict[str, int] = {}
for category in _VALID_CATEGORIES:
count = await store.count(_q(status="error", error_category=category))
if count:
by_category[category] = count
by_skill = [SkillStat(**row) for row in await store.skill_breakdown(_q())]
return ToolMetricsSummaryResponse(
total=total, error_total=error_total, by_category=by_category, by_skill=by_skill
)
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_tool_metrics(
request: Request,
user_id: str | None = Query(default=None),
user_name: str | None = Query(default=None),
tool_name: str | None = Query(default=None),
kind: str | None = Query(default=None),
skill_name: str | None = Query(default=None),
status: str | None = Query(default=None, pattern="^(success|error)$"),
error_category: str | None = Query(default=None),
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)
if not user_id and user_name:
resolved = await _resolve_user_name_to_id(user_name)
user_id = resolved if resolved is not None else "__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,
tool_name=tool_name,
skill_only=_kind_is_skill(kind),
skill_name=skill_name,
status=status,
error_category=error_category,
since=since,
until=until,
limit=1,
offset=0,
)
workbook = Workbook(write_only=True)
sheet = workbook.create_sheet("tool_call_metrics")
columns = [
"id",
"created_at",
"status",
"error_category",
"user_name",
"user_id",
"thread_id",
"run_id",
"agent_name",
"tool_name",
"skill_name",
"duration_ms",
"error_message",
]
sheet.append(columns)
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:
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 tool_call_metrics rows", count)
timestamp = datetime.utcnow().strftime("%Y%m%d-%H%M%S")
filename = f"tool_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,
)