**问题**:POST /api/v1/risk/daily-report/mail 原先只取 context 做 401 判定,
**没有任何授权校验**;RiskDailyReportMailService.send 既拿不到 context、也不调用
AuthorizationService。收件人、标题、正文**全部由客户端决定** —— 一旦运维开启 SMTP
(RISK_DAILY_REPORT_MAIL_ENABLED),它就是一个未授权的邮件发送器。
默认关闭(ENABLED 默认 false + DRY_RUN 默认 true)让它至今没出事,但那不是可依赖的保护。
**改动**:
- send 改为 async 并接收 context,入口处 wait AuthorizationService.require(
context, "risk:report:mail")。校验放在 **service 层**而不是 controller —— 本项目风控
端点的授权一律落在 service(risk_query / risk_action / risk_scan 等都是这样),
controller 只负责取 context;这个端点是唯一的例外,现在补齐。
- controller 相应改为 wait ...send(..., context=context)。
- 权限
isk:report:mail 已在上一轮随另外三个一起创建并授予 risk_operator 与 admin。
**实测**:
- 9002(risk_operator) → **200** + {"status":"disabled","recipient_count":1}(默认关闭)
- 9001(customer) → **403** AGENT_PERMISSION_DENIED「缺少操作权限」
**测试**:3 个既有用例改为 async 并传入带权限的 context;**新增**
est_mail_service_requires_the_permission,断言无权限身份必须被拒 —— 钉住本次修复。
ruff / mypy(136 文件) / 623 unit+contract / 29 integration 全绿。
226 lines
8.5 KiB
Python
226 lines
8.5 KiB
Python
"""风控只读查询接口。"""
|
|
|
|
import json
|
|
from collections.abc import AsyncIterator
|
|
from datetime import datetime, time
|
|
|
|
from fastapi import APIRouter, Depends, File, Path, UploadFile
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
from starlette.responses import StreamingResponse
|
|
|
|
from app.api.dependencies.auth import build_request_context
|
|
from app.api.dependencies.database import get_session
|
|
from app.api.dependencies.rate_limit import enforce_rate_limit
|
|
from app.api.schemas.risk import (
|
|
RiskAlertEscalationRequest,
|
|
RiskAlertExclusionRequest,
|
|
RiskAlertPageQuery,
|
|
RiskAlertResolutionRequest,
|
|
RiskDailyReportGenerateRequest,
|
|
RiskDailyReportMailRequest,
|
|
RiskEvidencePageQuery,
|
|
RiskEvidenceSource,
|
|
RiskNotificationPageQuery,
|
|
)
|
|
from app.core.contracts import RequestContext
|
|
from app.service.risk_action_service import RiskActionService
|
|
from app.service.risk_daily_report_mail_service import RiskDailyReportMailService
|
|
from app.service.risk_daily_report_service import RiskDailyReportService
|
|
from app.service.risk_evidence_archive_service import RiskEvidenceArchiveService
|
|
from app.service.risk_notification_service import RiskNotificationService
|
|
from app.service.risk_query_service import RiskQueryService
|
|
from app.service.risk_scan_service import RiskScanService
|
|
|
|
router = APIRouter(
|
|
prefix="/api/v1/risk",
|
|
tags=["risk"],
|
|
dependencies=[Depends(enforce_rate_limit)],
|
|
)
|
|
|
|
|
|
@router.get("/overview")
|
|
async def risk_overview(
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskQueryService(session).overview(context)
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.get("/alerts")
|
|
async def list_risk_alerts(
|
|
query: RiskAlertPageQuery = Depends(), # noqa: B008
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskQueryService(session).list_alerts(context, query)
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.post("/alerts/scan")
|
|
async def scan_risk_alerts(
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskScanService(session).scan(context)
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.post("/alerts/{alert_no}/acknowledgements")
|
|
async def acknowledge_risk_alert(
|
|
alert_no: str = Path(min_length=1, max_length=64, pattern=r"^[A-Za-z0-9_-]+$"),
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskActionService(session).acknowledge(alert_no, context)
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.post("/alerts/{alert_no}/investigations")
|
|
async def investigate_risk_alert(
|
|
alert_no: str = Path(min_length=1, max_length=64, pattern=r"^[A-Za-z0-9_-]+$"),
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskActionService(session).investigate(alert_no, context)
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.post("/alerts/{alert_no}/exclusions")
|
|
async def exclude_risk_alert(
|
|
payload: RiskAlertExclusionRequest,
|
|
alert_no: str = Path(min_length=1, max_length=64, pattern=r"^[A-Za-z0-9_-]+$"),
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskActionService(session).exclude(alert_no, payload.reason, context)
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.post("/alerts/{alert_no}/resolutions")
|
|
async def resolve_risk_alert(
|
|
payload: RiskAlertResolutionRequest,
|
|
alert_no: str = Path(min_length=1, max_length=64, pattern=r"^[A-Za-z0-9_-]+$"),
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskActionService(session).resolve(alert_no, payload.resolution, context)
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.post("/alerts/{alert_no}/escalations")
|
|
async def escalate_risk_alert(
|
|
payload: RiskAlertEscalationRequest,
|
|
alert_no: str = Path(min_length=1, max_length=64, pattern=r"^[A-Za-z0-9_-]+$"),
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskActionService(session).escalate(alert_no, payload.reason, context)
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.post("/alerts/{alert_no}/evidence")
|
|
async def archive_risk_evidence(
|
|
evidence_file: UploadFile = File(...), # noqa: B008
|
|
alert_no: str = Path(min_length=1, max_length=64, pattern=r"^[A-Za-z0-9_-]+$"),
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
try:
|
|
data = await RiskEvidenceArchiveService(session).archive(alert_no, evidence_file, context)
|
|
return _envelope(data, context)
|
|
finally:
|
|
await evidence_file.close()
|
|
|
|
|
|
@router.get("/alerts/{alert_no}")
|
|
async def get_risk_alert(
|
|
alert_no: str = Path(min_length=1, max_length=64, pattern=r"^[A-Za-z0-9_-]+$"),
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskQueryService(session).get_alert_detail(context, alert_no.strip())
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.get("/evidence/{source}")
|
|
async def list_risk_evidence(
|
|
source: RiskEvidenceSource = Path(), # noqa: B008
|
|
query: RiskEvidencePageQuery = Depends(), # noqa: B008
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskQueryService(session).list_evidence(context, source, query)
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.get("/notifications")
|
|
async def list_risk_notifications(
|
|
query: RiskNotificationPageQuery = Depends(), # noqa: B008
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskNotificationService(session).list_notifications(context, query)
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.post("/daily-report")
|
|
async def generate_risk_daily_report(
|
|
payload: RiskDailyReportGenerateRequest,
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> dict[str, object]:
|
|
report_time = (
|
|
datetime.combine(payload.report_date, time.min)
|
|
if payload.report_date is not None
|
|
else None
|
|
)
|
|
data = await RiskDailyReportService(session).generate(context, report_time)
|
|
return _envelope(data, context)
|
|
|
|
|
|
@router.post("/daily-report/stream")
|
|
async def stream_risk_daily_report(
|
|
payload: RiskDailyReportGenerateRequest,
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
session: AsyncSession = Depends(get_session), # noqa: B008
|
|
) -> StreamingResponse:
|
|
report_time = (
|
|
datetime.combine(payload.report_date, time.min)
|
|
if payload.report_date is not None
|
|
else None
|
|
)
|
|
|
|
async def events() -> AsyncIterator[str]:
|
|
async for event in RiskDailyReportService(session).stream(context, report_time):
|
|
event_type = str(event.get("type", "message"))
|
|
payload = json.dumps(event, ensure_ascii=False, default=str)
|
|
yield f"event: {event_type}\ndata: {payload}\n\n"
|
|
|
|
return StreamingResponse(
|
|
events(),
|
|
media_type="text/event-stream",
|
|
headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
|
|
)
|
|
|
|
|
|
@router.post("/daily-report/mail")
|
|
async def send_risk_daily_report_mail(
|
|
payload: RiskDailyReportMailRequest,
|
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
|
) -> dict[str, object]:
|
|
data = await RiskDailyReportMailService().send(
|
|
payload.recipients,
|
|
payload.subject,
|
|
payload.content,
|
|
context=context,
|
|
)
|
|
return _envelope(data, context)
|
|
|
|
|
|
def _envelope(data: object, context: RequestContext) -> dict[str, object]:
|
|
return {
|
|
"data": data,
|
|
"meta": {"trace_id": context.trace_id},
|
|
}
|