算力中心新增消耗明细和趋势统计接口

- compute_client.py 新增 user_logs 方法,封装引擎日志接口(支持分页)
- rbac_opc.py 新增两个接口:
  - GET /opc/compute/usage/logs - 消耗明细(分页,统一字段格式为前端期望的 snake_case)
  - GET /opc/compute/usage/trend - 消耗趋势统计(近N天趋势、模型分布、时段分布)
- 趋势接口基于引擎日志数据聚合计算:
  - 近N天消耗趋势(按日期聚合 cost 和 tokens)
  - 模型分布(按模型名聚合调用次数,Top 5 + 其他)
  - 时段分布(按24小时聚合消耗)
- 引擎对接失败时返回空数据,不报错(避免前端白屏)
- 模型分布自动分配颜色,便于前端饼图展示
This commit is contained in:
Pine
2026-09-05 23:33:24 +08:00
parent 2cb93b7751
commit 1e9c4343d3
2 changed files with 141 additions and 0 deletions
+128
View File
@@ -366,6 +366,134 @@ async def opc_compute_usage(user: dict = Depends(require_roles("opc_member"))):
return data
@router.get("/compute/usage/logs", summary="我的消耗明细(分页)")
async def opc_compute_usage_logs(
page: int = 1,
page_size: int = 20,
user: dict = Depends(require_roles("opc_member")),
):
"""消耗流水明细,支持分页。"""
try:
data = await compute_client.user_logs(user.get("username"), page=page, page_size=page_size)
except compute_client.ComputeError as exc:
raise HTTPException(status_code=502, detail=f"算力引擎对接失败: {exc}") from exc
items = data if isinstance(data, list) else data.get("items", [])
total = data.get("total", len(items)) if isinstance(data, dict) else len(items)
# 统一字段格式(snake_case → 前端期望的格式)
out = []
for it in items or []:
it = dict(it)
out.append({
"id": it.get("id") or it.get("log_id") or "",
"time": it.get("created_at") or it.get("time") or it.get("created_time") or "",
"model": it.get("model_name") or it.get("model") or "",
"type": it.get("type") or it.get("mode") or "chat",
"input_tokens": it.get("input_tokens") or it.get("prompt_tokens") or 0,
"output_tokens": it.get("output_tokens") or it.get("completion_tokens") or 0,
"cost": it.get("cost") or it.get("cost_quota") or 0,
"status": it.get("status") or "success",
})
return {"items": out, "total": total, "page": page, "page_size": page_size}
@router.get("/compute/usage/trend", summary="消耗趋势统计(近7天/模型分布/时段分布)")
async def opc_compute_usage_trend(
days: int = 7,
user: dict = Depends(require_roles("opc_member")),
):
"""消耗趋势统计:近N天消耗趋势、模型分布、时段分布。"""
from collections import defaultdict
from datetime import datetime, timedelta
try:
# 获取最近的消耗日志(最多取 500 条用于统计)
data = await compute_client.user_logs(user.get("username"), page=1, page_size=500)
except compute_client.ComputeError as exc:
# 引擎对接失败时返回空数据,不报错
data = {"items": []}
items = data if isinstance(data, list) else data.get("items", [])
# 近N天消耗趋势
now = datetime.now()
trend_map = {}
for i in range(days):
d = (now - timedelta(days=days - 1 - i)).strftime("%m-%d")
trend_map[d] = {"date": d, "cost": 0, "tokens": 0}
# 模型分布
model_dist = defaultdict(lambda: {"name": "", "value": 0, "cost": 0})
# 时段分布(24小时)
hourly_dist = [{"hour": f"{h:02d}", "cost": 0} for h in range(24)]
total_cost = 0
total_tokens = 0
for it in items or []:
it = dict(it)
# 解析时间
time_str = it.get("created_at") or it.get("time") or it.get("created_time") or ""
cost = float(it.get("cost") or it.get("cost_quota") or 0)
input_tokens = int(it.get("input_tokens") or it.get("prompt_tokens") or 0)
output_tokens = int(it.get("output_tokens") or it.get("completion_tokens") or 0)
tokens = input_tokens + output_tokens
model_name = it.get("model_name") or it.get("model") or "unknown"
total_cost += cost
total_tokens += tokens
# 趋势统计
if time_str:
try:
dt = datetime.fromisoformat(time_str.replace("Z", "+00:00")) if "T" in time_str else datetime.strptime(time_str[:10], "%Y-%m-%d")
date_key = dt.strftime("%m-%d")
if date_key in trend_map:
trend_map[date_key]["cost"] += cost
trend_map[date_key]["tokens"] += tokens
# 时段统计
hour = dt.hour
hourly_dist[hour]["cost"] += cost
except (ValueError, TypeError):
pass
# 模型分布统计
model_dist[model_name]["name"] = model_name
model_dist[model_name]["value"] += 1
model_dist[model_name]["cost"] += cost
# 模型分布取 Top 5,其余归为"其他"
model_list = sorted(model_dist.values(), key=lambda x: x["value"], reverse=True)
if len(model_list) > 5:
top5 = model_list[:5]
others = model_list[5:]
other_value = sum(x["value"] for x in others)
other_cost = sum(x["cost"] for x in others)
top5.append({"name": "其他", "value": other_value, "cost": other_cost})
model_list = top5
# 模型分布颜色
colors = ["#6284ff", "#8f7bff", "#3fb68b", "#e8830c", "#e0526e", "#38a8e0"]
for i, m in enumerate(model_list):
m["color"] = colors[i % len(colors)]
return {
"trend": list(trend_map.values()),
"model_distribution": model_list,
"hourly_distribution": hourly_dist,
"summary": {
"total_cost": total_cost,
"total_tokens": total_tokens,
"total_requests": len(items or []),
"days": days,
},
}
@router.get("/compute/balance", summary="我的算力余额")
async def opc_compute_balance(user: dict = Depends(require_roles("opc_member"))):
try:
+13
View File
@@ -174,6 +174,19 @@ async def user_usage(username: str) -> dict:
return data.get("data") or {}
async def user_logs(username: str, page: int = 1, page_size: int = 50) -> dict:
"""用户消耗流水/日志(PAT 鉴权)。"""
key = await _user_key(username)
if not key:
return {"items": [], "total": 0}
data = await _request(
method="GET",
path=f"/api/log/self?page={page}&page_size={page_size}",
headers={"Authorization": f"Bearer {key}"},
)
return data.get("data") or data
async def delete_token(token_id: int) -> dict:
"""删除引擎令牌。"""
return await _request(