merge: 同步 origin/qyqy_develop(风控扫描、知识检索、客服 Agent、协商等)

冲突处理:均为双方各自新增,按并集保留——
- .gitignore:本地 data/logs 忽略项 + 远端 .dsh-drop/
- app/core/config.py:offsite/promotion 与 risk_scan 配置项并存
- app/service/agent/bootstrap.py:场外/推介/风控/NL2SQL/知识检索 工具与 Agent 全部注册
- app/worker/__main__.py:场外邮件 Worker 接线 + runtime 关系服务/投影清理注入并存

收尾:新增 20260911_merge_risk_heads 收敛迁移双 head;按 docs/21 生成 config/jwt/dev 开发密钥。
This commit is contained in:
2026-09-11 17:32:47 +08:00
153 changed files with 16846 additions and 267 deletions
+1 -34
View File
@@ -7,6 +7,7 @@ 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.negotiation import accepts_event_stream
from app.api.dependencies.rate_limit import enforce_rate_limit
from app.api.schemas.agent_runs import (
AgentRunAcceptedEnvelope,
@@ -25,40 +26,6 @@ from app.service.run_query_service import RunQueryService
router = APIRouter(prefix="/api/v1/agent-runs", tags=["agent-runs"],
dependencies=[Depends(enforce_rate_limit)])
SSE_MEDIA_TYPE = "text/event-stream"
def accepts_event_stream(accept: str | None) -> bool:
"""`Accept` 是否接受 `text/event-stream`(文档 §3.2:该头**非必填**)。
- 未携带(`None` 或空串)→ 放行:文档写明"默认 `application/json`;SSE 为
`text/event-stream`",即由接口自身决定响应类型,不是客户端错误;
- 携带 `text/event-stream`、`text/*` 或 `*/*` 且 `q != 0` → 放行;
- 显式携带但只接受其他类型(如 `application/json`)→ 拒绝,由调用方转
`406 SSE_NOT_ACCEPTABLE`(文档 §3.5/§6.4)。
只做"是否可接受"的判定,不参与内容协商排序:SSE 端点只有一种表示。
"""
if accept is None or not accept.strip():
return True
for entry in accept.split(","):
parts = entry.split(";")
media_type = parts[0].strip().lower()
if media_type not in {SSE_MEDIA_TYPE, "text/*", "*/*"}:
continue
quality = 1.0
for parameter in parts[1:]:
name, _, value = parameter.partition("=")
if name.strip().lower() == "q":
try:
quality = float(value.strip())
except ValueError:
quality = 0.0
if quality > 0:
return True
return False
@router.post(
"",
response_model=AgentRunAcceptedEnvelope,
+6 -2
View File
@@ -5,6 +5,7 @@ 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.conversations import FeedbackRequest
from app.api.views.envelope import envelope, list_envelope
from app.core.contracts import RequestContext
from app.core.cursor import parse_cursor
from app.service.conversation_service import ConversationService
@@ -28,9 +29,10 @@ async def list_messages(
`400 INVALID_CURSOR`,而不是被静默忽略后返回第一页。
"""
before = parse_cursor(cursor)
return await ConversationService(session).messages(
page = await ConversationService(session).messages(
session_id, context, limit, before=before
)
return list_envelope(page, context)
@router.post("/conversation-messages/{message_id}/feedback", status_code=status.HTTP_201_CREATED)
@@ -40,5 +42,7 @@ async def create_feedback(
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
return await ConversationService(session).feedback(
# 与消息列表同理:§3.3 的信封由 Controller 统一套,service 只负责业务数据。
data = await ConversationService(session).feedback(
message_id, context, payload.rating, payload.feedback_type, payload.feedback_content)
return envelope(data, context)
+5 -1
View File
@@ -2,6 +2,7 @@ from fastapi import APIRouter, Depends, Path
from app.api.dependencies.auth import build_request_context
from app.api.dependencies.rate_limit import enforce_rate_limit
from app.api.views.envelope import envelope
from app.core.contracts import RequestContext
from app.service.knowledge_service import KnowledgeReferenceService
@@ -14,4 +15,7 @@ async def resolve_reference(
reference_token: str = Path(min_length=20, max_length=300),
context: RequestContext = Depends(build_request_context), # noqa: B008
) -> dict[str, object]:
return await KnowledgeReferenceService().resolve(context, reference_token)
# §3.3:成功响应也要有 `meta.trace_id`。此前这里直接返回资源对象,客户端拿不到
# 本次请求的追踪标识,出问题时无法与服务端日志对上。
data = await KnowledgeReferenceService().resolve(context, reference_token)
return envelope(data, context)
+105 -19
View File
@@ -1,15 +1,17 @@
"""风控只读查询接口。"""
import json
from collections.abc import AsyncIterator
from collections.abc import AsyncIterator, Awaitable, Callable
from datetime import datetime, time
from typing import Any
from fastapi import APIRouter, Depends, File, Path, UploadFile
from fastapi import APIRouter, Depends, File, Header, Path, Request, 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.negotiation import accepts_event_stream
from app.api.dependencies.rate_limit import enforce_rate_limit
from app.api.schemas.risk import (
RiskAlertEscalationRequest,
@@ -22,14 +24,19 @@ from app.api.schemas.risk import (
RiskEvidenceSource,
RiskNotificationPageQuery,
)
from app.api.views.envelope import envelope as _envelope
from app.api.views.envelope import list_envelope as _list_envelope
from app.core.contracts import RequestContext
from app.core.errors import SseNotAcceptableError
from app.infrastructure.db import mysql_scan_lock
from app.service.api_transaction_service import ApiTransactionService
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
from app.service.risk_scan_service import RiskScanBusyError, RiskScanService
router = APIRouter(
prefix="/api/v1/risk",
@@ -38,6 +45,26 @@ router = APIRouter(
)
async def _idempotent_write(
session: AsyncSession,
context: RequestContext,
key: str | None,
scope: str,
body: Any,
action: Callable[[AsyncSession], Awaitable[dict[str, Any]]],
) -> dict[str, Any]:
"""风控写接口的统一幂等入口(`docs/05` §5.1、§5.2、§13.2)。
`scope` 用**实际**路径(含 `alert_no`):§5.1 的幂等范围是
`user_id + method + normalized_path + idempotency_key`,把路径参数折成模板会让
同一个键在不同预警之间互相回放 —— 那是把两次不同资源的操作当成一次。
幂等记录与业务写入同事务(见 `ApiTransactionService.execute_in`),重复请求直接
回放 `response_json`,不会二次驱动状态机。
"""
return await ApiTransactionService().execute_in(session, context, scope, key, body, action)
@router.get("/overview")
async def risk_overview(
context: RequestContext = Depends(build_request_context), # noqa: B008
@@ -54,15 +81,28 @@ async def list_risk_alerts(
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
data = await RiskQueryService(session).list_alerts(context, query)
return _envelope(data, context)
return _list_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
key: str | None = Header(default=None, alias="Idempotency-Key"),
) -> dict[str, object]:
data = await RiskScanService(session).scan(context)
async def run(inner: AsyncSession) -> dict[str, Any]:
# 手工触发的扫描必须与定时扫描互斥,否则两条路径会同时查不到重复、同时插入。
# 锁加在**入口层**而不是 `RiskScanService.scan()` 内部:`GET_LOCK` 是连接级的,
# 而调度器已在它自己的 session 上持锁 —— 被两个入口共用的服务方法若再取同一把锁,
# 取锁的连接不是持锁的那一个、必然失败,会**把定时扫描自己挡死**。
async with mysql_scan_lock() as acquired:
if not acquired:
raise RiskScanBusyError("规则扫描正在执行,请稍后重试")
return await RiskScanService(inner).scan(context)
data = await _idempotent_write(
session, context, key, "POST /api/v1/risk/alerts/scan", {}, run
)
return _envelope(data, context)
@@ -71,8 +111,16 @@ 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
key: str | None = Header(default=None, alias="Idempotency-Key"),
) -> dict[str, object]:
data = await RiskActionService(session).acknowledge(alert_no, context)
data = await _idempotent_write(
session,
context,
key,
f"POST /api/v1/risk/alerts/{alert_no}/acknowledgements",
{},
lambda inner: RiskActionService(inner).acknowledge(alert_no, context),
)
return _envelope(data, context)
@@ -81,8 +129,16 @@ 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
key: str | None = Header(default=None, alias="Idempotency-Key"),
) -> dict[str, object]:
data = await RiskActionService(session).investigate(alert_no, context)
data = await _idempotent_write(
session,
context,
key,
f"POST /api/v1/risk/alerts/{alert_no}/investigations",
{},
lambda inner: RiskActionService(inner).investigate(alert_no, context),
)
return _envelope(data, context)
@@ -92,8 +148,16 @@ async def exclude_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
key: str | None = Header(default=None, alias="Idempotency-Key"),
) -> dict[str, object]:
data = await RiskActionService(session).exclude(alert_no, payload.reason, context)
data = await _idempotent_write(
session,
context,
key,
f"POST /api/v1/risk/alerts/{alert_no}/exclusions",
{"reason": payload.reason},
lambda inner: RiskActionService(inner).exclude(alert_no, payload.reason, context),
)
return _envelope(data, context)
@@ -103,8 +167,16 @@ async def resolve_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
key: str | None = Header(default=None, alias="Idempotency-Key"),
) -> dict[str, object]:
data = await RiskActionService(session).resolve(alert_no, payload.resolution, context)
data = await _idempotent_write(
session,
context,
key,
f"POST /api/v1/risk/alerts/{alert_no}/resolutions",
{"resolution": payload.resolution},
lambda inner: RiskActionService(inner).resolve(alert_no, payload.resolution, context),
)
return _envelope(data, context)
@@ -114,8 +186,16 @@ async def escalate_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
key: str | None = Header(default=None, alias="Idempotency-Key"),
) -> dict[str, object]:
data = await RiskActionService(session).escalate(alert_no, payload.reason, context)
data = await _idempotent_write(
session,
context,
key,
f"POST /api/v1/risk/alerts/{alert_no}/escalations",
{"reason": payload.reason},
lambda inner: RiskActionService(inner).escalate(alert_no, payload.reason, context),
)
return _envelope(data, context)
@@ -151,7 +231,7 @@ async def list_risk_evidence(
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
data = await RiskQueryService(session).list_evidence(context, source, query)
return _envelope(data, context)
return _list_envelope(data, context)
@router.get("/notifications")
@@ -161,7 +241,7 @@ async def list_risk_notifications(
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
data = await RiskNotificationService(session).list_notifications(context, query)
return _envelope(data, context)
return _list_envelope(data, context)
@router.post("/daily-report")
@@ -182,6 +262,7 @@ async def generate_risk_daily_report(
@router.post("/daily-report/stream")
async def stream_risk_daily_report(
payload: RiskDailyReportGenerateRequest,
request: Request,
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
) -> StreamingResponse:
@@ -190,9 +271,16 @@ async def stream_risk_daily_report(
if payload.report_date is not None
else None
)
service = RiskDailyReportService(session)
# 鉴权与内容协商都必须在返回 StreamingResponse **之前**完成:`stream()` 是 async
# generator,函数体到第一次迭代才执行,而那时响应头已经发出去了 —— 403/406 只能
# 变成"200 + 半截流"(docs/25 P3 #24)。顺序与 §6.4 一致:先鉴权,后 Accept。
await service.authorize(context)
if not accepts_event_stream(request.headers.get("Accept")):
raise SseNotAcceptableError("Accept 必须接受 text/event-stream")
async def events() -> AsyncIterator[str]:
async for event in RiskDailyReportService(session).stream(context, report_time):
async for event in service.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"
@@ -209,16 +297,14 @@ async def send_risk_daily_report_mail(
payload: RiskDailyReportMailRequest,
context: RequestContext = Depends(build_request_context), # noqa: B008
) -> dict[str, object]:
data = RiskDailyReportMailService().send(
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},
}
# `_envelope` / `_list_envelope` 已抽到 `app/api/views/envelope.py`,与客服链路共用同一份
# §3.3 实现 —— 两处各写一份的结果就是其中一处漏了 `meta`。