Files
lzf_0626 d896a906cd fix(risk): 规则扫描加跨进程锁——手工触发与定时扫描此前可以同时跑
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 全绿。
2026-09-11 13:57:22 +08:00

69 lines
3.7 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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
from app.core.config import get_settings
# 存储层统一 UTC:`init_command` 是 asyncmy(driver 层)在**每次建立连接**时执行的
# 初始化 SQL,SQLAlchemy 的 asyncmy 方言会把 DSN 上的 query 参数原样透传给
# `asyncmy.connect`,因此写在 DSN 上即可生效,无需事件监听。
#
# 为什么不挂 `connect` 事件:`connect` 事件在 asyncmy 这套 asyncio 方言上**不会**
# 等待 async 监听器——实测注册后 `@@session.time_zone` 仍是 `SYSTEM`,同时抛
# `RuntimeWarning: coroutine ... was never awaited`。也就是说那种写法不报错、不生效,
# 是"看着修好了、其实照旧差 8 小时"的典型陷阱。
#
# 背景:建库时未显式指定时区,MySQL 会话继承系统时区(本机 Asia/Shanghai),
# 表上 `created_at DATETIME(6) DEFAULT CURRENT_TIMESTAMP(6)` 因此写入**北京时间**,
# 而应用代码统一写 `datetime.now(UTC)`:同一行里 DB 默认值与应用写入值相差 8 小时。
# 更严重的真实故障是给 `sys_user_role.assigned_at` 用 MySQL `NOW()` 会超前 UTC 8 小时,
# 被 RBAC 的 `assigned_at <= now` 判为"尚未生效",接口直接 403。
#
# `.env` 的 `TIMEZONE=Asia/Shanghai` 只作用于展示/日志,改它**不会**影响 MySQL 会话时区,
# 所以必须在连接层设置。改这里即可:新增的连接一律 UTC,DB 默认值与应用 UTC 语义一致。
_SESSION_UTC_INIT_COMMAND = "SET time_zone = '+00:00'"
settings = get_settings()
_url = make_url(settings.mysql_dsn)
if "init_command" not in _url.query:
_url = _url.update_query_dict({"init_command": _SESSION_UTC_INIT_COMMAND})
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}
)