1. 新增统一分页常量定义,替换各路由零散的Query参数声明 2. 批量替换多个列表接口的分页参数为预设常量 3. 优化插件目录、消息搜索等接口的参数与文档 4. 修复死信导出接口的参数命名不一致问题
325 lines
15 KiB
Python
325 lines
15 KiB
Python
"""审计日志域 Router(AUD-QUERY-01 / AUD-05 / AUD-06 / AUD-07 / AUD-RETENTION + 同步导出)。
|
||
|
||
本 router 实现审计日志域的全部 HTTP 端点,覆盖审计日志查询、统计、操作类型
|
||
枚举、单条详情、保留策略管理与同步导出共 7 个操作。所有端点统一采用模板 A
|
||
(控制面端口路由),通过 ``get_channel_use_cases`` 装配 ``ChannelUseCases``,
|
||
经 ``audit_query`` 端口调用类型化方法,由 ``ChannelControlService`` 内部走
|
||
控制面管道(auth → permission → rate_limit → dispatch → audit)。
|
||
|
||
鉴权策略:全部端点使用 ``get_admin_user`` 依赖(``PUT /audit/retention-policy``
|
||
升级为 ``get_superadmin_user``),要求管理员或超级管理员角色。角色校验由依赖
|
||
函数完成,router 内不做角色判断(规范 §4)。
|
||
|
||
模板选型:模板 A(控制面端口路由)——全部端点需要审计留痕(记录"谁查了 /
|
||
导了审计日志"),数据面端口不经控制面管道无法在 audit 阶段留痕。
|
||
|
||
路径设计:静态路径(``/audit/logs`` / ``/audit/logs/stats`` /
|
||
``/audit/operations`` / ``/audit/retention-policy`` / ``/audit/export``)先于
|
||
动态路径(``/audit/logs/{log_id}``)声明,避免被动态路径捕获(规范 §6.5)。
|
||
子 router 不自行设置 prefix,根前缀 ``/channels`` 由 ``channels_router`` 聚合
|
||
router 统一追加。完整 HTTP 路径为 ``/channels/audit/*``(不含 ``{channel_type}``
|
||
段,审计日志为跨渠道全局能力)。
|
||
|
||
特殊处理:``GET /audit/export`` 返回 ``StreamingResponse`` 输出 JSON 数组
|
||
``[entry1, entry2, ...]``,不遵循 ``{"success": True, "data": ...}`` 标准响应
|
||
结构(文件下载语义特例)。控制面管道返回 ``{"entries": [...], "total": N}``,
|
||
Router 层提取 ``entries`` 序列化为 JSON 数组。响应头 ``X-Total-Count`` 携带
|
||
匹配总数、``X-Truncated`` 标识是否因上限截断,使调用方可感知数据完整性。
|
||
审计留痕使用独立 operation ``audit/export``(audit_type=``audit_exported``),
|
||
与 ``audit/query``(audit_type=``admin_query``)区分,满足"谁导了审计日志"
|
||
的合规要求。
|
||
|
||
tags 命名:``audit_router = APIRouter(tags=["channels-audit"])``,与现有
|
||
``channels-config`` / ``channels-account`` 等子 router 命名规范一致。
|
||
|
||
端点清单(对应《审计日志域设计方案》§2.1):
|
||
- GET /audit/logs AUD-QUERY-01 query_audit_logs
|
||
- GET /audit/logs/stats AUD-05 get_audit_log_stats
|
||
- GET /audit/operations AUD-06 list_audit_operations
|
||
- GET /audit/export 同步导出审计日志(audit/export)
|
||
- GET /audit/retention-policy AUD-RETENTION-GET get_retention_policy
|
||
- PUT /audit/retention-policy AUD-RETENTION-PUT update_retention_policy
|
||
- GET /audit/logs/{log_id} AUD-07 get_audit_log
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
from datetime import datetime
|
||
from typing import Any
|
||
from zoneinfo import ZoneInfo
|
||
|
||
from fastapi import APIRouter, Depends, Query, Request
|
||
from fastapi.responses import StreamingResponse
|
||
from pydantic import BaseModel, ConfigDict, Field
|
||
from yuxi.channels.contract.dtos.audit import (
|
||
AuditOperationType,
|
||
AuditQuery,
|
||
RetentionPolicyUpdateCmd,
|
||
)
|
||
from yuxi.channels.contract.dtos.channel import ChannelType
|
||
from yuxi.storage.postgres.models_business import User
|
||
|
||
from server.routers.channels import (
|
||
LARGE_LIMIT,
|
||
OFFSET,
|
||
build_operator,
|
||
get_channel_use_cases,
|
||
parse_datetime,
|
||
raiseOnControlFailure,
|
||
serialize_control_data,
|
||
)
|
||
from server.utils.auth_middleware import get_admin_user, get_superadmin_user
|
||
|
||
audit_router = APIRouter(tags=["channels-audit"])
|
||
|
||
# 导出上限常量(与 adapter 5s 查询超时匹配,避免大查询超时与内存峰值)
|
||
AUDIT_EXPORT_LIMIT = 10000
|
||
|
||
|
||
class UpdateRetentionPolicyRequest(BaseModel):
|
||
"""更新审计保留策略请求体(AUD-RETENTION-PUT)。
|
||
|
||
字段约束与 ``RetentionPolicyUpdateCmd`` 一致,DTO ``__post_init__``
|
||
二次校验 ``auto_archive_before_days < default_retention_days`` 与
|
||
``by_operation_type`` key/value 合法性。
|
||
"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
default_retention_days: int = Field(default=90, ge=1, description="默认保留天数")
|
||
by_operation_type: dict[str, int] | None = Field(
|
||
default=None, description="按操作类型分组的保留天数(key 为操作类型,value 为正整数)"
|
||
)
|
||
auto_archive_enabled: bool = Field(default=True, description="是否启用自动归档(变更需重启生效)")
|
||
auto_archive_before_days: int = Field(default=80, ge=1, description="归档提前天数(须小于 default_retention_days)")
|
||
|
||
|
||
@audit_router.get("/audit/logs", response_model=dict)
|
||
async def query_audit_logs(
|
||
request: Request,
|
||
operation_type: AuditOperationType | None = Query(default=None, description="操作类型"),
|
||
operator: str | None = Query(default=None, description="操作人用户 ID"),
|
||
target_channel: ChannelType | None = Query(default=None, description="目标渠道类型"),
|
||
target_account: str | None = Query(default=None, description="目标账户 ID"),
|
||
start_time: str | None = Query(default=None, description="起始时间(ISO 8601)"),
|
||
end_time: str | None = Query(default=None, description="结束时间(ISO 8601)"),
|
||
trace_id: str | None = Query(default=None, description="按链路追踪 ID 过滤"),
|
||
limit: int = LARGE_LIMIT,
|
||
offset: int = OFFSET,
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""查询审计日志列表(AUD-QUERY-01)。对应控制面操作 audit/query。
|
||
|
||
``operation_type`` 查询参数由 FastAPI 按 ``AuditOperationType`` 枚举校验,
|
||
非法值返回 422。
|
||
"""
|
||
operator_vo = build_operator(current_user, request)
|
||
start_time_dt = parse_datetime("start_time", start_time)
|
||
end_time_dt = parse_datetime("end_time", end_time)
|
||
query = AuditQuery(
|
||
operation_type=operation_type,
|
||
operator=operator,
|
||
target_channel=target_channel,
|
||
target_account=target_account,
|
||
start_time=start_time_dt,
|
||
end_time=end_time_dt,
|
||
trace_id=trace_id,
|
||
limit=limit,
|
||
offset=offset,
|
||
)
|
||
result = await use_cases.audit_query.queryAuditLogs(query=query, operator=operator_vo)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@audit_router.get("/audit/logs/stats", response_model=dict)
|
||
async def get_audit_log_stats(
|
||
request: Request,
|
||
operation_type: AuditOperationType | None = Query(default=None, description="操作类型"),
|
||
target_channel: ChannelType | None = Query(default=None, description="目标渠道类型"),
|
||
target_account: str | None = Query(default=None, description="目标账户 ID"),
|
||
start_time: str | None = Query(default=None, description="起始时间(ISO 8601)"),
|
||
end_time: str | None = Query(default=None, description="结束时间(ISO 8601)"),
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""查询审计日志统计聚合(AUD-05)。对应控制面操作 audit/stats。
|
||
|
||
统计查询不含 ``limit`` / ``offset`` 分页参数,对全量匹配数据聚合。
|
||
"""
|
||
operator_vo = build_operator(current_user, request)
|
||
start_time_dt = parse_datetime("start_time", start_time)
|
||
end_time_dt = parse_datetime("end_time", end_time)
|
||
query = AuditQuery(
|
||
operation_type=operation_type,
|
||
target_channel=target_channel,
|
||
target_account=target_account,
|
||
start_time=start_time_dt,
|
||
end_time=end_time_dt,
|
||
)
|
||
result = await use_cases.audit_query.getAuditLogStats(query=query, operator=operator_vo)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@audit_router.get("/audit/operations", response_model=dict)
|
||
async def list_audit_operations(
|
||
request: Request,
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""查询审计操作类型枚举(AUD-06)。对应控制面操作 audit/operations。
|
||
|
||
返回 ``AuditOperationType`` 全部枚举值列表,dispatch 阶段直接返回静态枚举,
|
||
无 DB 操作。
|
||
"""
|
||
operator_vo = build_operator(current_user, request)
|
||
result = await use_cases.audit_query.listAuditOperations(operator=operator_vo)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@audit_router.get("/audit/export")
|
||
async def export_audit_logs(
|
||
request: Request,
|
||
operation_type: AuditOperationType | None = Query(default=None, description="操作类型"),
|
||
operator: str | None = Query(default=None, description="操作人用户 ID"),
|
||
target_channel: ChannelType | None = Query(default=None, description="目标渠道类型"),
|
||
target_account: str | None = Query(default=None, description="目标账户 ID"),
|
||
start_time: str | None = Query(default=None, description="起始时间(ISO 8601)"),
|
||
end_time: str | None = Query(default=None, description="结束时间(ISO 8601)"),
|
||
trace_id: str | None = Query(default=None, description="按链路追踪 ID 过滤"),
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> StreamingResponse:
|
||
"""同步导出审计日志。对应控制面操作 audit/export。
|
||
|
||
走控制面管道(模板 A)调用 ``exportAuditLogs``,按查询条件过滤后以
|
||
``StreamingResponse`` 输出 JSON 数组 ``[entry1, entry2, ...]``。
|
||
响应头 ``Content-Disposition: attachment;
|
||
filename="audit_export_{YYYYMMDD_HHMMSS}.json"`` 触发浏览器下载,
|
||
文件名时间使用 Asia/Shanghai 本地时区。
|
||
|
||
导出上限 ``AUDIT_EXPORT_LIMIT`` 条(10000),与 adapter 5s 查询超时
|
||
匹配,避免大查询超时与内存峰值;无匹配数据时返回空数组 ``[]``。
|
||
响应头 ``X-Total-Count`` 携带匹配总数、``X-Truncated`` 标识是否因
|
||
上限截断,调用方可据此判断数据完整性。如需导出更大数据量,应改用
|
||
异步导出任务(本期未实现)。
|
||
|
||
审计留痕使用独立 operation ``audit/export``(audit_type=
|
||
``audit_exported``),与查询操作(``admin_query``)区分,满足"谁导了
|
||
审计日志"的合规要求。文件下载语义特例不遵循
|
||
``{"success": True, "data": ...}`` 标准响应结构。
|
||
"""
|
||
operator_vo = build_operator(current_user, request)
|
||
start_time_dt = parse_datetime("start_time", start_time)
|
||
end_time_dt = parse_datetime("end_time", end_time)
|
||
query = AuditQuery(
|
||
operation_type=operation_type,
|
||
operator=operator,
|
||
target_channel=target_channel,
|
||
target_account=target_account,
|
||
start_time=start_time_dt,
|
||
end_time=end_time_dt,
|
||
trace_id=trace_id,
|
||
limit=AUDIT_EXPORT_LIMIT,
|
||
offset=0,
|
||
)
|
||
result = await use_cases.audit_query.exportAuditLogs(query=query, operator=operator_vo)
|
||
raiseOnControlFailure(result)
|
||
data = serialize_control_data(result.data)
|
||
entries = data["entries"]
|
||
total = data["total"]
|
||
|
||
async def _stream_json_array():
|
||
yield "["
|
||
first = True
|
||
for entry in entries:
|
||
if not first:
|
||
yield ","
|
||
first = False
|
||
yield json.dumps(entry, ensure_ascii=False, default=str)
|
||
yield "]"
|
||
|
||
filename = f"audit_export_{datetime.now(ZoneInfo('Asia/Shanghai')).strftime('%Y%m%d_%H%M%S')}.json"
|
||
truncated = total > AUDIT_EXPORT_LIMIT
|
||
return StreamingResponse(
|
||
_stream_json_array(),
|
||
media_type="application/json",
|
||
headers={
|
||
"Content-Disposition": f'attachment; filename="{filename}"',
|
||
"X-Total-Count": str(total),
|
||
"X-Truncated": "true" if truncated else "false",
|
||
},
|
||
)
|
||
|
||
|
||
@audit_router.get("/audit/retention-policy", response_model=dict)
|
||
async def get_retention_policy(
|
||
request: Request,
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""查询审计保留策略(AUD-RETENTION-GET)。
|
||
|
||
对应控制面操作 ``audit/retention_policy_get``。复用 ConfigManager 读取
|
||
``audit_retention_policy`` 配置键(GLOBAL 作用域),配置未初始化时返回
|
||
默认值(default_retention_days=90、auto_archive_enabled=True、
|
||
auto_archive_before_days=80、updated_at=None)。
|
||
"""
|
||
operator_vo = build_operator(current_user, request)
|
||
result = await use_cases.audit_query.getRetentionPolicy(operator=operator_vo)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@audit_router.put("/audit/retention-policy", response_model=dict)
|
||
async def update_retention_policy(
|
||
payload: UpdateRetentionPolicyRequest,
|
||
request: Request,
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_superadmin_user),
|
||
) -> dict[str, Any]:
|
||
"""更新审计保留策略(AUD-RETENTION-PUT)。
|
||
|
||
对应控制面操作 ``audit/retention_policy_update``。复用 ConfigManager 更新
|
||
``audit_retention_policy`` 配置键(GLOBAL 作用域),dispatch 阶段校验
|
||
``auto_archive_before_days < default_retention_days``,``updated_at``
|
||
刷新为当前时间并随配置值持久化。
|
||
|
||
``auto_archive_enabled`` 为不可热更新字段:配置正常持久化(HTTP 200),
|
||
但变更不会热生效。响应 data 中 ``requires_restart`` 字段为 ``True`` 时,
|
||
客户端应提示用户需重启服务才能使该字段变更生效。
|
||
"""
|
||
operator_vo = build_operator(current_user, request)
|
||
cmd = RetentionPolicyUpdateCmd(
|
||
operator=operator_vo,
|
||
default_retention_days=payload.default_retention_days,
|
||
by_operation_type=payload.by_operation_type,
|
||
auto_archive_enabled=payload.auto_archive_enabled,
|
||
auto_archive_before_days=payload.auto_archive_before_days,
|
||
)
|
||
result = await use_cases.audit_query.updateRetentionPolicy(cmd=cmd)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|
||
|
||
|
||
@audit_router.get("/audit/logs/{log_id}", response_model=dict)
|
||
async def get_audit_log(
|
||
log_id: str,
|
||
request: Request,
|
||
use_cases=Depends(get_channel_use_cases),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""查询单条审计日志详情(AUD-07)。对应控制面操作 audit/get。
|
||
|
||
``log_id`` 不存在时返回 ``NOT_FOUND``,由 ``raiseOnControlFailure`` 映射为
|
||
``NotFoundError``(HTTP 404)。
|
||
"""
|
||
operator_vo = build_operator(current_user, request)
|
||
result = await use_cases.audit_query.getAuditLog(log_id=log_id, operator=operator_vo)
|
||
raiseOnControlFailure(result)
|
||
return {"success": True, "data": serialize_control_data(result.data)}
|