From d896a906cd9146d354f1f2e9024bf7c24e8ee472 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=8D=BF=E4=BA=91=E7=A7=8B=E6=9C=88?= <15273589815@163.com> Date: Fri, 11 Sep 2026 13:57:22 +0800 Subject: [PATCH] =?UTF-8?q?fix(risk):=20=E8=A7=84=E5=88=99=E6=89=AB?= =?UTF-8?q?=E6=8F=8F=E5=8A=A0=E8=B7=A8=E8=BF=9B=E7=A8=8B=E9=94=81=E2=80=94?= =?UTF-8?q?=E2=80=94=E6=89=8B=E5=B7=A5=E8=A7=A6=E5=8F=91=E4=B8=8E=E5=AE=9A?= =?UTF-8?q?=E6=97=B6=E6=89=AB=E6=8F=8F=E6=AD=A4=E5=89=8D=E5=8F=AF=E4=BB=A5?= =?UTF-8?q?=E5=90=8C=E6=97=B6=E8=B7=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit docs/25 P2 最后一项。核实后分清了两层,报告没区分: - **定时扫描是安全的**:risk_scan_scheduler.py 已有 MySQL 连接级咨询锁 (GET_LOCK,锁名 jr_risk_scan_schedule),跨进程互斥。 - **HTTP 端点不安全**:POST /api/v1/risk/alerts/scan → RiskScanService.scan() 只用了 **进程内** asyncio.Lock。多 Web worker、或 Worker 与 API 同时运行时形同虚设。 而扫描的幂等只有应用层的 _exists 查重 —— fin_risk_alert 的 trigger_rule_codes 是 JSON 数组,**无法建唯一索引兜底**(同一交易可命中多条规则,唯一键本应是"交易+规则",而规则 埋在 JSON 里)。所以两条路径并发时会同时查不到、同时插入,产生重复预警。 **改动**: 1. 把 mysql_scan_lock 与锁名移到 pp/infrastructure/db.py —— 端点与调度器**必须共用 同一把锁**,放在基础设施层两个入口才都能引用(service 不该反向依赖 worker)。 调度器改为从那里 import。 2. **端点层加锁**(controllers/risk.py 的 scan 端点):取不到锁就抛 RiskScanBusyError (与 service 内部那把进程内锁用同一错误类型与文案)。 **为什么不加在 RiskScanService.scan() 内部**:GET_LOCK 是**连接级**的,而调度器已经在 它自己的 session 上持锁;被两个入口共用的服务方法若再取同一把锁,取锁的连接并不是持锁的 那一个、必然返回 0 —— 会**把定时扫描自己挡死**。所以锁加在入口层,每个入口只取一次。 **实测**: - 无人持锁时扫描 → **200**「规则扫描完成」 - 本进程先取得跨进程锁后再调端点 → **409**「规则扫描正在执行,请稍后重试」 (同一进程内不同 session 也互斥,说明它是连接级的,正是跨进程所需) - 释放后再调 → **200**,恢复正常 顺带第 4 次遇到 409 复用错误码 RUN_NOT_CANCELLABLE,语义不符;属 P3 待处理项。 ruff / mypy(136 文件) / 639 unit+contract 全绿。 --- app/api/controllers/risk.py | 12 ++++++++-- app/infrastructure/db.py | 37 +++++++++++++++++++++++++++++++ app/worker/risk_scan_scheduler.py | 30 +++++-------------------- 3 files changed, 52 insertions(+), 27 deletions(-) diff --git a/app/api/controllers/risk.py b/app/api/controllers/risk.py index 18aec6a..025c97f 100644 --- a/app/api/controllers/risk.py +++ b/app/api/controllers/risk.py @@ -23,13 +23,14 @@ from app.api.schemas.risk import ( RiskNotificationPageQuery, ) from app.core.contracts import RequestContext +from app.infrastructure.db import mysql_scan_lock 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", @@ -62,7 +63,14 @@ 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) + # 手工触发的扫描必须与定时扫描互斥,否则两条路径会同时查不到重复、同时插入。 + # 锁加在**入口层**而不是 `RiskScanService.scan()` 内部:`GET_LOCK` 是连接级的, + # 而调度器已在它自己的 session 上持锁 —— 被两个入口共用的服务方法若再取同一把锁, + # 取锁的连接不是持锁的那一个、必然失败,会**把定时扫描自己挡死**。 + async with mysql_scan_lock() as acquired: + if not acquired: + raise RiskScanBusyError("规则扫描正在执行,请稍后重试") + data = await RiskScanService(session).scan(context) return _envelope(data, context) diff --git a/app/infrastructure/db.py b/app/infrastructure/db.py index bc60deb..4201a31 100644 --- a/app/infrastructure/db.py +++ b/app/infrastructure/db.py @@ -1,3 +1,7 @@ +from collections.abc import AsyncIterator +from contextlib import asynccontextmanager + +from sqlalchemy import text from sqlalchemy.engine import make_url from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine @@ -29,3 +33,36 @@ if "init_command" not in _url.query: engine = create_async_engine(_url, pool_pre_ping=True) SessionFactory = async_sessionmaker(engine, class_=AsyncSession, expire_on_commit=False) + +# 规则扫描的**跨进程**咨询锁名。用 MySQL 连接级 `GET_LOCK` 而不是进程内 `asyncio.Lock`: +# 后者只在单个进程内互斥,多 Web worker、或 Worker 与 API 同时运行时形同虚设。 +SCAN_LOCK_NAME = "jr_risk_scan_schedule" + + +@asynccontextmanager +async def mysql_scan_lock() -> AsyncIterator[bool]: + """约束规则扫描的跨进程并发;yield 是否取得锁。 + + **为什么两个入口都必须用它**:定时扫描(`worker/risk_scan_scheduler.py`)与手工触发的 + HTTP 端点(`POST /api/v1/risk/alerts/scan`)最终都调 `RiskScanService.scan()`, + 而扫描的幂等只有应用层的 `_exists` 查重 —— 两条路径并发时会同时查不到、同时插入, + 产生重复预警。锁要**共用同一个名字**才能互相排斥。 + + **为什么不放在 `scan()` 内部**:`GET_LOCK` 是**连接级**的。调度器已在自己的 session 上 + 持锁,再新建 session 去调 `scan()`;若 `scan()` 又去取同一把锁,取锁的连接并不是持锁的 + 那一个,`GET_LOCK` 返回 0 —— **调度器会把自己挡死**。所以锁加在**入口层** + (调度器 / 端点各取一次),而不是被两个入口共用的服务方法里。 + + `GET_LOCK(name, 0)` 的 0 表示不等待、立即返回:取不到就跳过本轮,而不是排队。 + """ + async with SessionFactory() as session: + acquired = bool( + await session.scalar(text("SELECT GET_LOCK(:name, 0)"), {"name": SCAN_LOCK_NAME}) + ) + try: + yield acquired + finally: + if acquired: + await session.scalar( + text("SELECT RELEASE_LOCK(:name)"), {"name": SCAN_LOCK_NAME} + ) diff --git a/app/worker/risk_scan_scheduler.py b/app/worker/risk_scan_scheduler.py index d21c6dd..9c5a4f1 100644 --- a/app/worker/risk_scan_scheduler.py +++ b/app/worker/risk_scan_scheduler.py @@ -9,17 +9,15 @@ from __future__ import annotations import argparse import asyncio import logging -from collections.abc import AsyncIterator, Awaitable, Callable -from contextlib import AbstractAsyncContextManager, asynccontextmanager +from collections.abc import Awaitable, Callable +from contextlib import AbstractAsyncContextManager from datetime import UTC, datetime, timedelta from typing import Any from uuid import uuid4 -from sqlalchemy import text - from app.core.config import get_settings from app.core.contracts import RequestContext -from app.infrastructure.db import SessionFactory, engine +from app.infrastructure.db import SessionFactory, engine, mysql_scan_lock from app.model.audit import InteractionAudit from app.service.risk_scan_schedule_config import ( RiskScanScheduleConfig, @@ -27,7 +25,6 @@ from app.service.risk_scan_schedule_config import ( ) from app.service.risk_scan_service import RiskScanService -SCAN_LOCK_NAME = "jr_risk_scan_schedule" logger = logging.getLogger(__name__) ConfigLoader = Callable[[], Awaitable[RiskScanScheduleConfig]] @@ -35,25 +32,8 @@ ScanExecutor = Callable[[], Awaitable[dict[str, int | str]]] AuditWriter = Callable[[str, dict[str, Any]], Awaitable[None]] LockFactory = Callable[[], AbstractAsyncContextManager[bool]] - -@asynccontextmanager -async def mysql_scan_lock() -> AsyncIterator[bool]: - """使用 MySQL 连接级咨询锁约束跨进程并发。""" - async with SessionFactory() as session: - acquired = bool( - await session.scalar( - text("SELECT GET_LOCK(:name, 0)"), - {"name": SCAN_LOCK_NAME}, - ) - ) - try: - yield acquired - finally: - if acquired: - await session.scalar( - text("SELECT RELEASE_LOCK(:name)"), - {"name": SCAN_LOCK_NAME}, - ) +# 跨进程扫描锁已移到 `app/infrastructure/db.py`:端点与调度器**必须共用同一把锁**, +# 放在基础设施层两个入口才都能引用(service 不该反向依赖 worker)。 async def default_scan_executor() -> dict[str, int | str]: