本次提交对多个外部系统路由进行了多维度优化: 1. 统一分页参数:将所有路由的`page = offset//limit +1`、`page_size=limit`替换为标准的`limit`+`offset`分页格式 2. 完善接口文档:补充多个端点的功能说明、参数含义与返回字段解释 3. 增强参数校验:新增字段长度限制、正则校验、枚举类型约束与业务逻辑校验 4. 优化代码复用:提取重复逻辑为辅助函数,减少样板代码 5. 修复接口问题:修正工具健康检查端点路径参数类型,优化导出接口响应格式 6. 补充异常处理:为批量操作添加异常捕获与日志记录,避免流程中断
359 lines
14 KiB
Python
359 lines
14 KiB
Python
"""Alert 子域 Router。
|
||
|
||
外部系统限界上下文的告警事件管理 API,覆盖告警列表 / 详情 / 触发 /
|
||
确认 / 恢复 / 统计聚合 / 按链路查询 / 批量确认 / 批量恢复 / 告警上下文 /
|
||
通知状态。所有端点通过 ``create_use_cases_from_db`` 装配 use_cases,
|
||
经 ``alert_service`` 端口调用用例。
|
||
|
||
设计要点:
|
||
|
||
- ``list_alerts`` 支持 keyword / resource_type / resource_id / related_trace_id /
|
||
dedup_key / 时间范围 多维过滤,结果携带 ``system_name`` / ``system_slug``(LEFT JOIN 填充)。
|
||
- ``get_alert_stats`` 支持 system_id / env_key / 时间范围 过滤,聚合按 status / severity /
|
||
alert_type 分组与 Top 告警排行。
|
||
- ``list_alerts_by_trace`` 路径参数 ``trace_id`` 通过 ``Path(min_length=1, max_length=64)``
|
||
约束,避免空串与超长输入。
|
||
- ``fire_alert`` 显式字段映射构造 ``FireAlertInput``,``triggered_by`` 由当前管理员 uid 填充;
|
||
若 ``system_id`` 非空但系统不存在,Service 层抛 404;同一 ``dedup_key`` 在 firing 状态下去重更新。
|
||
- ``acknowledge_alert`` / ``resolve_alert`` 在告警不存在时返回 404,状态冲突时返回 409
|
||
(由 Service 层区分并通过 unified_error_handler 自动映射)。
|
||
- Request Schema 与 Input DTO 不共享类,Router 内显式构造 DTO,确保 API 边界与用例边界解耦。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
from datetime import datetime
|
||
from typing import Any, Literal
|
||
|
||
from fastapi import APIRouter, Body, Depends, HTTPException, Path, Query
|
||
from pydantic import BaseModel, ConfigDict, Field
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
from yuxi.external_systems.infrastructure.container import create_use_cases_from_db
|
||
from yuxi.external_systems.use_cases.dto.alert import (
|
||
AcknowledgeAlertInput,
|
||
AlertStatsInput,
|
||
BatchAcknowledgeAlertsInput,
|
||
BatchResolveAlertsInput,
|
||
FireAlertInput,
|
||
GetAlertContextInput,
|
||
GetAlertInput,
|
||
GetAlertNotificationsInput,
|
||
ListAlertsByTraceInput,
|
||
ListAlertsInput,
|
||
ResolveAlertInput,
|
||
)
|
||
from yuxi.storage.postgres.models_business import User
|
||
|
||
from server.utils.auth_middleware import get_admin_user, get_db, get_required_user
|
||
|
||
alert_router = APIRouter(prefix="/alerts", tags=["external-systems-alert"])
|
||
|
||
|
||
# ---------------- Request Schemas ----------------
|
||
|
||
|
||
class FireAlertRequest(BaseModel):
|
||
"""触发告警请求体。字段对齐 ``FireAlertInput``。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
system_id: int | None = Field(default=None, ge=1)
|
||
alert_type: str = Field(..., min_length=1, max_length=64)
|
||
title: str = Field(..., min_length=1, max_length=256)
|
||
triggered_at: datetime
|
||
env_key: str | None = Field(default=None, max_length=32)
|
||
severity: Literal["info", "warning", "critical", "fatal"] = "warning"
|
||
description: str | None = None
|
||
detail: dict[str, Any] = Field(default_factory=dict)
|
||
resource_type: str | None = Field(default=None, max_length=32)
|
||
resource_id: str | None = Field(default=None, max_length=128)
|
||
related_execution_id: str | None = Field(default=None, max_length=64)
|
||
related_trace_id: str | None = Field(default=None, max_length=64)
|
||
metric_snapshot: dict[str, Any] | None = None
|
||
dedup_key: str | None = Field(default=None, max_length=256)
|
||
|
||
|
||
class AcknowledgeAlertRequest(BaseModel):
|
||
"""确认告警请求体。``acknowledged_by`` 由当前管理员填充,仅暴露 ``note``。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
note: str | None = None
|
||
|
||
|
||
class ResolveAlertRequest(BaseModel):
|
||
"""恢复告警请求体。字段对齐 ``ResolveAlertInput``(不含 id / user)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
resolution_note: str | None = None
|
||
|
||
|
||
class BatchAcknowledgeAlertsRequest(BaseModel):
|
||
"""批量确认告警请求体。字段对齐 ``BatchAcknowledgeAlertsInput``(不含 acknowledged_by)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
alert_ids: list[int] = Field(..., min_length=1, max_length=100)
|
||
note: str | None = None
|
||
|
||
|
||
class BatchResolveAlertsRequest(BaseModel):
|
||
"""批量恢复告警请求体。字段对齐 ``BatchResolveAlertsInput``(不含 user)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
alert_ids: list[int] = Field(..., min_length=1, max_length=100)
|
||
resolution_note: str | None = None
|
||
|
||
|
||
# ---------------- Endpoints ----------------
|
||
|
||
|
||
@alert_router.get("", response_model=dict)
|
||
async def list_alerts(
|
||
limit: int = Query(100, ge=1, le=500),
|
||
offset: int = Query(0, ge=0),
|
||
system_id: int | None = Query(None, ge=1),
|
||
env_key: str | None = Query(None, max_length=32),
|
||
alert_type: str | None = Query(None, max_length=64),
|
||
severity: Literal["info", "warning", "critical", "fatal"] | None = Query(
|
||
None, description="严重度:info/warning/critical/fatal"
|
||
),
|
||
status: Literal["firing", "acknowledged", "resolved", "suppressed"] | None = Query(
|
||
None, description="告警状态:firing/acknowledged/resolved/suppressed"
|
||
),
|
||
start_time: datetime | None = Query(None, description="ISO8601 开始时间"),
|
||
end_time: datetime | None = Query(None, description="ISO8601 结束时间"),
|
||
keyword: str | None = Query(
|
||
None, max_length=128, description="关键词搜索(title/description/alert_type/system name/slug)"
|
||
),
|
||
# 新增查询参数
|
||
resource_type: str | None = Query(None, max_length=32),
|
||
resource_id: str | None = Query(None, max_length=128),
|
||
related_trace_id: str | None = Query(None, max_length=64),
|
||
dedup_key: str | None = Query(None, max_length=256),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""分页列出告警事件(支持 keyword / resource_type / resource_id / related_trace_id / dedup_key 过滤)。
|
||
|
||
时间范围校验:若 ``start_time`` 晚于 ``end_time``,返回 422 校验错误。
|
||
列表结果包含 ``system_name`` / ``system_slug``(由 LEFT JOIN 填充)。
|
||
"""
|
||
if start_time is not None and end_time is not None and start_time > end_time:
|
||
raise HTTPException(
|
||
status_code=422,
|
||
detail="start_time 不能晚于 end_time",
|
||
)
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = ListAlertsInput(
|
||
limit=limit,
|
||
offset=offset,
|
||
system_id=system_id,
|
||
env_key=env_key,
|
||
alert_type=alert_type,
|
||
severity=severity,
|
||
status=status,
|
||
start_at=start_time,
|
||
end_at=end_time,
|
||
keyword=keyword,
|
||
resource_type=resource_type,
|
||
resource_id=resource_id,
|
||
related_trace_id=related_trace_id,
|
||
dedup_key=dedup_key,
|
||
)
|
||
output = await use_cases.alert_service.list_alerts(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@alert_router.get("/stats", response_model=dict)
|
||
async def get_alert_stats(
|
||
system_id: int | None = Query(None, ge=1),
|
||
env_key: str | None = Query(None, max_length=32),
|
||
start_time: datetime | None = Query(None, description="ISO8601 开始时间"),
|
||
end_time: datetime | None = Query(None, description="ISO8601 结束时间"),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""查询告警统计聚合(支持 system_id / env_key / 时间范围过滤)。
|
||
|
||
时间范围基于 ``triggered_at`` 字段,若 ``start_time`` 晚于 ``end_time`` 返回 422。
|
||
"""
|
||
if start_time is not None and end_time is not None and start_time > end_time:
|
||
raise HTTPException(
|
||
status_code=422,
|
||
detail="start_time 不能晚于 end_time",
|
||
)
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = AlertStatsInput(
|
||
system_id=system_id,
|
||
env_key=env_key,
|
||
start_at=start_time,
|
||
end_at=end_time,
|
||
)
|
||
output = await use_cases.alert_service.get_alert_stats(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@alert_router.get("/by-trace/{trace_id}", response_model=dict)
|
||
async def list_alerts_by_trace(
|
||
trace_id: str = Path(..., min_length=1, max_length=64, description="链路 ID"),
|
||
limit: int = Query(100, ge=1, le=500),
|
||
offset: int = Query(0, ge=0),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""按链路 ID 查询告警(通过 related_trace_id 关联,支持分页)。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = ListAlertsByTraceInput(
|
||
trace_id=trace_id,
|
||
limit=limit,
|
||
offset=offset,
|
||
)
|
||
output = await use_cases.alert_service.list_alerts_by_trace(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@alert_router.post("/batch-acknowledge", response_model=dict)
|
||
async def batch_acknowledge_alerts(
|
||
payload: BatchAcknowledgeAlertsRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""批量确认告警。``acknowledged_by`` 由当前管理员填充。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = BatchAcknowledgeAlertsInput(
|
||
alert_ids=payload.alert_ids,
|
||
note=payload.note,
|
||
acknowledged_by=current_user.uid,
|
||
)
|
||
output = await use_cases.alert_service.batch_acknowledge_alerts(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@alert_router.post("/batch-resolve", response_model=dict)
|
||
async def batch_resolve_alerts(
|
||
payload: BatchResolveAlertsRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""批量恢复告警。``user`` 由当前管理员填充,确保审计可追溯。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = BatchResolveAlertsInput(
|
||
alert_ids=payload.alert_ids,
|
||
user=current_user.uid,
|
||
resolution_note=payload.resolution_note,
|
||
)
|
||
output = await use_cases.alert_service.batch_resolve_alerts(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@alert_router.get("/{alert_id}", response_model=dict)
|
||
async def get_alert(
|
||
alert_id: int = Path(ge=1),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""获取告警详情。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetAlertInput(id=alert_id)
|
||
output = await use_cases.alert_service.get_alert(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@alert_router.post("", response_model=dict)
|
||
async def fire_alert(
|
||
payload: FireAlertRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""触发告警事件(含去重与自动通知派发)。
|
||
|
||
``triggered_by`` 由当前管理员 uid 填充,确保审计可追溯。
|
||
若 ``system_id`` 非空但对应系统不存在,Service 层抛 404。
|
||
同一 ``dedup_key`` 在 ``firing`` 状态下重复触发时更新已有记录。
|
||
"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = FireAlertInput(
|
||
system_id=payload.system_id,
|
||
alert_type=payload.alert_type,
|
||
title=payload.title,
|
||
triggered_at=payload.triggered_at,
|
||
env_key=payload.env_key,
|
||
severity=payload.severity,
|
||
description=payload.description,
|
||
detail=payload.detail,
|
||
resource_type=payload.resource_type,
|
||
resource_id=payload.resource_id,
|
||
related_execution_id=payload.related_execution_id,
|
||
related_trace_id=payload.related_trace_id,
|
||
metric_snapshot=payload.metric_snapshot,
|
||
dedup_key=payload.dedup_key,
|
||
triggered_by=current_user.uid,
|
||
)
|
||
output = await use_cases.alert_service.fire_alert(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@alert_router.post("/{alert_id}/acknowledge", response_model=dict)
|
||
async def acknowledge_alert(
|
||
alert_id: int = Path(ge=1),
|
||
payload: AcknowledgeAlertRequest = Body(...),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""确认告警。``acknowledged_by`` 由当前管理员 uid 填充,``note`` 可选。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = AcknowledgeAlertInput(
|
||
id=alert_id,
|
||
acknowledged_by=current_user.uid,
|
||
note=payload.note,
|
||
)
|
||
output = await use_cases.alert_service.acknowledge_alert(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@alert_router.post("/{alert_id}/resolve", response_model=dict)
|
||
async def resolve_alert(
|
||
alert_id: int = Path(ge=1),
|
||
payload: ResolveAlertRequest = Body(...),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""恢复告警。``user`` 由当前管理员 uid 填充,确保审计可追溯。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = ResolveAlertInput(
|
||
id=alert_id,
|
||
user=current_user.uid,
|
||
resolution_note=payload.resolution_note,
|
||
)
|
||
output = await use_cases.alert_service.resolve_alert(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@alert_router.get("/{alert_id}/context", response_model=dict)
|
||
async def get_alert_context(
|
||
alert_id: int = Path(ge=1),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""查询告警上下文(含执行记录 + 指标快照)。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetAlertContextInput(alert_id=alert_id)
|
||
output = await use_cases.alert_service.get_alert_context(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@alert_router.get("/{alert_id}/notifications", response_model=dict)
|
||
async def get_alert_notifications(
|
||
alert_id: int = Path(ge=1),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""查询告警通知状态。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetAlertNotificationsInput(alert_id=alert_id)
|
||
output = await use_cases.alert_service.get_alert_notifications(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|