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

279 lines
9.2 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""圆桌会商诊断日志查询 API(管理员)。
会商全流程(第一/二/三步 + 后台执行器 + 进程内网关)各处的报错与关键事件落在
``roundtable_diagnostics`` 表,本路由暴露**管理员**筛选查询 + facets + 清理。
Routes(prefix ``/api/roundtable-diagnostics``,全部 admin-gated)::
GET "" 筛选 + 分页列出(scope/stage/level/event/job_id/draft_id/task_id/user_id/q/since/until)
GET "/facets" 各筛选维度的 distinct 取值(下拉用)
POST "/admin/cleanup" 按保留天数清理旧日志
"""
from __future__ import annotations
import time
from datetime import UTC, datetime, timedelta
from typing import Any
from fastapi import APIRouter, HTTPException, Query, Request
from pydantic import BaseModel, Field
from deerflow.runtime.user_context import get_effective_user_id
router = APIRouter(prefix="/api/roundtable-diagnostics", tags=["roundtable-diagnostics"])
def _get_store(request: Request):
store = getattr(request.app.state, "roundtable_diagnostic_store", None)
if store is None:
raise HTTPException(status_code=503, detail="Roundtable diagnostic store not available")
return store
def _resolve_user_id(request: Request) -> str | None:
user = getattr(request.state, "user", None)
if user is not None:
return str(user.id)
try:
return get_effective_user_id()
except Exception: # noqa: BLE001
return None
async def _require_admin(request: Request) -> None:
"""放行 admin / auth-disabled,其余 403(与 roundtable_jobs / task_buttons 同套语义)。"""
from app.gateway.deps import get_optional_user_from_request
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="会商诊断日志仅限管理员查看")
class DiagnosticItem(BaseModel):
id: str
scope: str
stage: str
level: str
event: str
message: str | None = None
detail: Any = None
jobId: str | None = None
draftId: str | None = None
taskId: str | None = None
userId: str | None = None
agentId: str | None = None
agentName: str | None = None
cycle: int | None = None
createdAt: str | None = None
userName: str | None = None
class DiagnosticListResponse(BaseModel):
items: list[DiagnosticItem]
total: int
limit: int
offset: int
class FacetsResponse(BaseModel):
scope: list[str] = Field(default_factory=list)
stage: list[str] = Field(default_factory=list)
level: list[str] = Field(default_factory=list)
event: list[str] = Field(default_factory=list)
class CleanupRequest(BaseModel):
retention_days: int = Field(default=30, ge=0)
batch_size: int = Field(default=1000, ge=1, le=10000)
class CleanupResponse(BaseModel):
deleted: int
elapsed_seconds: float
class ReportRequest(BaseModel):
"""前端上报一条会商调用报错(任意登录用户可调,**非 admin**)。"""
stage: str = "client"
event: str = ""
level: str = "error"
message: str | None = None
detail: Any = None
jobId: str | None = None
draftId: str | None = None
taskId: str | None = None
agentId: str | None = None
agentName: str | None = None
cycle: int | None = None
def _to_item(row: dict[str, Any]) -> DiagnosticItem:
return DiagnosticItem(
id=row["id"],
scope=row.get("scope") or "",
stage=row.get("stage") or "",
level=row.get("level") or "",
event=row.get("event") or "",
message=row.get("message"),
detail=row.get("detail"),
jobId=row.get("job_id"),
draftId=row.get("draft_id"),
taskId=row.get("task_id"),
userId=row.get("user_id"),
agentId=row.get("agent_id"),
agentName=row.get("agent_name"),
cycle=row.get("cycle"),
createdAt=row.get("created_at"),
userName=row.get("_user_name"),
)
def _username_from_email(email: str | None, uid: str | None) -> str | None:
"""展示用用户名:邮箱 ``@`` 前缀(与个人页一致);无邮箱回退到短 uid。"""
if email:
return email.split("@", 1)[0] if "@" in email else email
if uid:
return uid[:8]
return None
async def _enrich_user_names(request: Request, items: list[dict[str, Any]]) -> None:
"""把每条诊断的 user_id 解析成用户名(邮箱 @ 前缀),写进 dict 的 ``_user_name``。
Best-effort + 批量去重缓存:解析不到就回退短 uid,绝不因此让列表查询失败。
"""
try:
from app.gateway.deps import get_local_provider
provider = get_local_provider()
except Exception: # noqa: BLE001
provider = None
cache: dict[str, str | None] = {}
for it in items:
uid = it.get("user_id")
if not uid:
it["_user_name"] = None
continue
if uid not in cache:
email: str | None = None
if provider is not None:
try:
user = await provider.get_user(uid)
email = getattr(user, "email", None) if user else None
except Exception: # noqa: BLE001
email = None
cache[uid] = _username_from_email(email, uid)
it["_user_name"] = cache[uid]
def _parse_dt(raw: str | None) -> datetime | None:
if not raw:
return None
try:
# 容忍前端传的 ``...Z`` 后缀(datetime.fromisoformat 在 3.11+ 接受,旧版不接受)。
return datetime.fromisoformat(raw.replace("Z", "+00:00"))
except Exception:
return None
@router.get("", response_model=DiagnosticListResponse)
async def list_diagnostics(
request: Request,
scope: str | None = None,
stage: str | None = None,
level: str | None = None,
event: str | None = None,
job_id: str | None = None,
draft_id: str | None = None,
task_id: str | None = None,
user_id: str | None = None,
q: str | None = None,
since: str | None = None,
until: str | None = None,
limit: int = Query(default=100, ge=1, le=500),
offset: int = Query(default=0, ge=0),
) -> DiagnosticListResponse:
await _require_admin(request)
store = _get_store(request)
# scope 支持逗号分隔多值(前端「前台调用」= foreground,client 两种 scope 合并查)。
scope_param: str | list[str] | None = None
if scope:
parts = [s.strip() for s in scope.split(",") if s.strip()]
scope_param = parts if len(parts) > 1 else (parts[0] if parts else None)
items, total = await store.list_logs(
scope=scope_param,
stage=stage or None,
level=level or None,
event=event or None,
job_id=job_id or None,
draft_id=draft_id or None,
task_id=task_id or None,
user_id=user_id or None,
q=q or None,
since=_parse_dt(since),
until=_parse_dt(until),
limit=limit,
offset=offset,
)
await _enrich_user_names(request, items)
return DiagnosticListResponse(items=[_to_item(r) for r in items], total=total, limit=limit, offset=offset)
@router.post("/report")
async def report_client_error(request: Request, body: ReportRequest) -> dict[str, bool]:
"""前端把一次会商调用的报错上报上来(HTTP 502/网络失败/SSE error 帧等用户实际遇到的)。
**不 admin-gated**——任何登录用户在遇到报错时都应能上报;写入 scope="client" 以区分
「服务端观测到的(foreground)」与「前端上报的(client)」。best-effort,写失败也返回 ok。
"""
store = getattr(request.app.state, "roundtable_diagnostic_store", None)
if store is None:
return {"ok": False}
try:
await store.record(
scope="client",
stage=body.stage or "client",
level=body.level or "error",
event=body.event or "",
message=body.message,
detail=body.detail,
job_id=body.jobId,
draft_id=body.draftId,
task_id=body.taskId,
user_id=_resolve_user_id(request),
agent_id=body.agentId,
agent_name=body.agentName,
cycle=body.cycle,
)
except Exception: # noqa: BLE001 — 上报失败不影响前端
return {"ok": False}
return {"ok": True}
@router.get("/facets", response_model=FacetsResponse)
async def get_facets(request: Request) -> FacetsResponse:
await _require_admin(request)
store = _get_store(request)
facets = await store.facets()
return FacetsResponse(**facets)
@router.post("/admin/cleanup", response_model=CleanupResponse)
async def cleanup(request: Request, body: CleanupRequest) -> CleanupResponse:
await _require_admin(request)
store = _get_store(request)
cutoff = datetime.now(UTC) - timedelta(days=body.retention_days)
t0 = time.perf_counter()
deleted = 0
while True:
ids = await store.list_older_than(cutoff, limit=body.batch_size)
if not ids:
break
deleted += await store.delete_by_ids(ids)
if len(ids) < body.batch_size:
break
return CleanupResponse(deleted=deleted, elapsed_seconds=round(time.perf_counter() - t0, 3))