2026-09-11 13:57:22 +08:00
|
|
|
|
from collections.abc import AsyncIterator
|
|
|
|
|
|
from contextlib import asynccontextmanager
|
|
|
|
|
|
|
|
|
|
|
|
from sqlalchemy import text
|
2026-09-10 15:55:54 +08:00
|
|
|
|
from sqlalchemy.engine import make_url
|
2026-09-09 21:55:37 +08:00
|
|
|
|
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine
|
|
|
|
|
|
|
|
|
|
|
|
from app.core.config import get_settings
|
|
|
|
|
|
|
2026-09-10 15:55:54 +08:00
|
|
|
|
# 存储层统一 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)
|
2026-09-09 21:55:37 +08:00
|
|
|
|
SessionFactory = async_sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
|
2026-09-11 13:57:22 +08:00
|
|
|
|
|
|
|
|
|
|
# 规则扫描的**跨进程**咨询锁名。用 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}
|
|
|
|
|
|
)
|