diff --git a/app/utils/db.py b/app/utils/db.py index 2468886..3dfab89 100644 --- a/app/utils/db.py +++ b/app/utils/db.py @@ -2,32 +2,48 @@ 按库名缓存单例 Engine:deps / risk / simulate 路由每请求实例化 Repository 时 复用同一连接池,不再每次 create_engine(B6 复审 P3:实例化点泄漏);应用 -shutdown 经 dispose_engines 统一释放。测试直传 engine= 参数的用法不受影响; +shutdown 经 dispose_engines 统一释放(B7 评审 P1-1:显式 Engine.dispose(), +不依赖 GC 兜底)。测试直传 engine= 参数的用法不受影响; monkeypatch settings 后须先 dispose_engines() 清缓存。 """ from __future__ import annotations -from functools import lru_cache +import threading from sqlalchemy import create_engine from sqlalchemy.engine import Engine from app.config.settings import settings +_engines: dict[str, Engine] = {} +_engines_lock = threading.Lock() + -@lru_cache(maxsize=None) def get_engine(database: str) -> Engine: """取指定库的单例 Engine(进程内共享连接池)。""" - pwd = settings.mysql_password - auth = f"{settings.mysql_user}:{pwd}" if pwd else settings.mysql_user - url = ( - f"mysql+pymysql://{auth}@{settings.mysql_host}:{settings.mysql_port}" - f"/{database}?charset=utf8mb4" - ) - return create_engine(url, pool_pre_ping=True) + with _engines_lock: + engine = _engines.get(database) + if engine is None: + pwd = settings.mysql_password + auth = f"{settings.mysql_user}:{pwd}" if pwd else settings.mysql_user + url = ( + f"mysql+pymysql://{auth}@{settings.mysql_host}:{settings.mysql_port}" + f"/{database}?charset=utf8mb4" + ) + engine = create_engine(url, pool_pre_ping=True) + _engines[database] = engine + return engine def dispose_engines() -> None: - """释放全部缓存 Engine(lifespan shutdown / 测试隔离)。""" - get_engine.cache_clear() + """释放全部缓存 Engine(lifespan shutdown / 测试隔离)。 + + 逐个 Engine.dispose() 关闭连接池空闲连接后清空缓存;调用方仍持有的 + checked-out 连接不受影响,归还时由旧池关闭。dispose 后再 get_engine + 按当前 settings 重建(测试 monkeypatch 配置依赖此语义)。 + """ + with _engines_lock: + for engine in _engines.values(): + engine.dispose() + _engines.clear() diff --git a/docs/memory/MEMORY.md b/docs/memory/MEMORY.md index 6e17e15..518973a 100644 --- a/docs/memory/MEMORY.md +++ b/docs/memory/MEMORY.md @@ -9,7 +9,7 @@ **项目是什么:** 金融四 Agent(客户财富 / 代理人 / 数据分析 / 风控)共用数据层与合规底座;**不**互调 LLM,跨 Agent 走 L1/L2/L3 画像与预警表。 -**当前进度:** 需求与表设计已定 · **风控模块已落地 B1~B7**(规则/预警聚合/L3/AML/引擎/交易网关/鉴权+4 API/适当性校验/**main 集成+挂账①~⑦**,**180 测试绿**)· 剩 **B8 conftest+集成测试 → B9a 脚本 → B9b 演示走查** · **JWT(T-01) / 审计中间件(T-02) / LangGraph 对话线(T-07) 未做**(chat/knowledge/admin 仍空壳)。**开发在分支 `feature/risk`(未合入 main)。** +**当前进度:** 需求与表设计已定 · **风控模块已落地 B1~B7**(规则/预警聚合/L3/AML/引擎/交易网关/鉴权+4 API/适当性校验/**main 集成+挂账①~⑦**,B7 复审闭环,**183 测试绿**)· 剩 **B8 conftest+集成测试 → B9a 脚本 → B9b 演示走查** · **JWT(T-01) / 审计中间件(T-02) / LangGraph 对话线(T-07) 未做**(chat/knowledge/admin 仍空壳)。**开发在分支 `feature/risk`(未合入 main)。** **仓库地图:** @@ -28,7 +28,7 @@ | `scripts/core/*.sql` + `reset.ps1` | **已实现** | Core 模拟库 DDL + 种子 | | `scripts/agent/` `scripts/demo/` | **已实现** | AML 名单种子 + 风控演示数据 | | `scripts/sync/*.py` | **已实现** | 归属同步 + Neo4j 全图 | -| `tests/` | **已实现** | 13 个测试文件 173 用例(sqlite 隔离) | +| `tests/` | **已实现** | 15 个测试文件 183 用例(sqlite 隔离) | | `docs/需求拆解/` | 已定 | 场景 P0、矩阵、合规原文 | | `docs/PRD/PRD-风控监测Agent.md` | **已冻结** | 风控 PRD v1.0 + 规则表附录 | | `docs/项目框架设计/表设计/` | 已定 | Agent 共用 11 表 + agent 专用 SQL | @@ -122,7 +122,7 @@ Core 模拟:scripts/core/reset.ps1 · 文档 docs/项目框架设计/Core模 种子:scripts/agent/seed-aml-list.sql(AML 名单)· scripts/demo/prepare_risk_demo.sql(reset 后重跑) 依赖:requirements.txt(LangGraph + langchain-core/openai + FastAPI + SQLAlchemy) 启动:uvicorn app.main:app --reload → GET /health -测试:python -m pytest(173 用例,sqlite 隔离) +测试:python -m pytest(183 用例,sqlite 隔离) 配置:.env(见 .env.example) RBAC 联调账号:scripts/dev/rbac-seed-reference.md ``` @@ -163,6 +163,6 @@ RBAC 联调账号:scripts/dev/rbac-seed-reference.md 2. 改动属于 api / service / tool / repository 哪一层? 3. 是否需 customer_id 归属与 JWT RBAC? 4. Core 是模拟库只读还是 agent 库读写? -5. 如何验证?(`python -m pytest` 全量(当前 180 绿)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) +5. 如何验证?(`python -m pytest` 全量(当前 183 绿)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) 大任务:FRAMEWORK/FLOW 与实现状态不符时先更新 memory 再编码(用户确认跳过除外)。 diff --git a/docs/memory/TODO.md b/docs/memory/TODO.md index 48b9c5b..f3831a2 100644 --- a/docs/memory/TODO.md +++ b/docs/memory/TODO.md @@ -20,7 +20,7 @@ ### 风控模块(PRD v1.0 已冻结 · `docs/PRD/PRD-风控监测Agent.md`,事件驱动线不依赖 T-07 可先行) -- [ ] T-30 风控事件线:`app/gateway/` 交易网关 + `service/risk/` 规则引擎(RISK-001~005)+ 预警单聚合 + L3 最小写入 + `risk:pub:alert` 推送 + AML(含 `risk_aml_list` 种子)+ 演示数据脚本验收 A-1~A-5/A-9 —— **进度(2026-09-06):B1 规则纯函数 / B2 预警服务 / B3 L3 写入(risk_score 一期不写)/ B4 AML+引擎编排 / B5 交易网关 / B6 鉴权+4 API+处置编排 / B7 main 集成(路由挂载 + trace 中间件 + lifespan 启动期 debug 校验;挂账①锁公共化 ②L3 缓存 DEL ③死代码 ④统一错误体 ⑤启动校验 ⑥引擎工厂 ⑦处置原子事务均落地;⑧ input_guard_log 双写仍随 T-02)均完成并经独立 AI 评审闭环(180 测试绿);剩 B8 conftest+集成测试、B9a 脚本、B9b 演示走查** +- [ ] T-30 风控事件线:`app/gateway/` 交易网关 + `service/risk/` 规则引擎(RISK-001~005)+ 预警单聚合 + L3 最小写入 + `risk:pub:alert` 推送 + AML(含 `risk_aml_list` 种子)+ 演示数据脚本验收 A-1~A-5/A-9 —— **进度(2026-09-06):B1 规则纯函数 / B2 预警服务 / B3 L3 写入(risk_score 一期不写)/ B4 AML+引擎编排 / B5 交易网关 / B6 鉴权+4 API+处置编排 / B7 main 集成(路由挂载 + trace 中间件 + lifespan 启动期 debug 校验;挂账①锁公共化 ②L3 缓存 DEL ③死代码 ④统一错误体 ⑤启动校验 ⑥引擎工厂 ⑦处置原子事务均落地;⑧ input_guard_log 双写仍随 T-02)均完成并经独立 AI 评审闭环(B7 复审有条件通过→P1-1 dispose 已修闭环;183 测试绿;P2/P3 已登记开发计划 B7/B8/B9b/T-02);剩 B8 conftest+集成测试、B9a 脚本、B9b 演示走查** - [ ] T-31 `service/suitability.py` 公共校验(SUIT-001~008)+ `POST /api/risk/suitability/check` + 单测验收 A-8 —— **代码已完成(A1~A4 · 77 测试绿;B6 起该 API 带鉴权+直调审计);MySQL 手工 SQL 对照挂账至 B9b 执行(阶段 A 评审 P2-8)** - [ ] T-32 预警台账与人工处置 API(`GET /alerts`、`POST /handle`,risk_officer/compliance 权限)+ 对话线(依赖 T-01/T-03/T-07)验收 A-6/A-7 —— **API 部分已由 B6 覆盖并复审通过(A-7/A-9 用例绿);对话线仍依赖 T-01/T-03/T-07** diff --git a/docs/项目框架设计/开发计划-风控模块.md b/docs/项目框架设计/开发计划-风控模块.md index 17bea34..e24f453 100644 --- a/docs/项目框架设计/开发计划-风控模块.md +++ b/docs/项目框架设计/开发计划-风控模块.md @@ -27,9 +27,11 @@ | B5 | `app/gateway/`(trade_gateway + gateway_repository 仅 INSERT core_trade)+ `api/simulate.py` 薄路由 | 网关 | 集成:convert 400、阻断不落 trade | A4、B4 | | B6 | **`app/api/deps.py`:`AuthContext`(actor_id/roles/customer_id,字段按 JWT 手册冻结)+ `get_auth_context()` 工厂**——dev 模式从 `X-Debug-Role`/`X-Debug-Actor` 请求头构造、`app_env != development` 启动时检测 debug 头直接拒绝;T-01 就绪后仅替换工厂内部为 JWT 解析,签名不变。另:`api/risk.py` 4 个 API(GET alerts / POST handle / POST suitability/check / POST aml/scan)+ 归属校验(含 compliance 强制 aml 过滤)。**备注:依赖层须校验 handler_result 枚举(repo 不校验);本阶段顺手统一 `NotFoundError` 异常(utils/exceptions.py 现为占位)**(评审 P2-7①②) | 鉴权依赖 + 4 个 API | Swagger 手测 + **权限矩阵(按 debug 头切换角色/身份执行 A-7/A-9 用例)** | A4、B2、**B4**(aml/scan 依赖 scan_all) | | B7 | `main.py` 集成:路由挂载 + lifespan(双 Engine 单例注入 + Redis 单例 + trace 中间件)。**备注:顺手提取 `utils/db.py` 引擎工厂收敛 core_ro/risk_repository 双份 _default_engine**(评审 P2-6);**B3 挂账(B3 评审 P2-5):① `_run_locked` 锁原语公共化(alert_service/profile_l3 现复用私有实现)② L3 写侧 Redis 缓存 DEL 钩子(PRD §5.1 `profile:l3:{customer_id}` 更新时 DEL,`profile_l3.upsert_profile_l3` 已留痕)③ 删除 core_ro.list_trades 死代码(B4 改用 list_trades_range 后无调用方,复审 N3)④ 统一响应外壳落地(utils/response.py 现占位,simulate/risk 路由届时一并包裹,B5 评审 P2-2;错误体对齐手册 §10 error_code/message/trace_id,B6 评审 P3-2)**;**B6 挂账(B6 复审):⑤ lifespan 启动期检测 `app_env != development` + debug 头依赖直接拒绝启动(现为请求时 RuntimeError,B6 评审遗漏①)⑥ 引擎工厂须覆盖 deps/simulate/risk 三处每请求 `RiskRepository()`/`CoreReadOnlyRepository()` 实例化点并 dispose(现复用 _default_engine 不释放,B6 复审 P3)⑦ `handle_alert` 的 update_alert_status 与 insert_audit_log 两事务非原子——统一事务或补偿记录(B6 复审 P3)⑧ input_guard_log 双写缺口(手册 P-05 要求 audit_log+input_guard_log,表归 T-03 底座,B7 接 T-02 审计中间件时统一补)** | 可运行应用 | `uvicorn` 启动 + `/health` + 全路由可达 | B5、B6 | -| B8 | **`tests/conftest.py`**:a) session fixture 启动校验演示数据就位(CUST-4001 测评 <365 天、risk_aml_list ≥8),缺失则中止并提示先跑 FLOW §0 ③④;b) fixture 幂等代跑 `prepare_risk_demo.sql`;c) teardown 按 `TRD-TEST-` 清 core_trade + 关联 risk_alert/risk_suitability_log/audit_log + 还原 L3 行。**顺手集中 sqlite 测试 DDL 为单一事实源(B4 评审 P3-12,各测试文件手写 DDL 收敛;CURRENT_TIMESTAMP 改 localtime 或 fixture 固定时间,防 UTC/本地日界错位——B5 评审 P3-4)**;集成测试交易统一走 `trade_id_factory` 注入 `TRD-TEST-` 前缀(trade_gateway 已留参数,B5 评审 P3-2)。集成测试:A-1~A-5、A-7(状态机/compliance 403/GET 强制 aml)、A-9 越权、**trace 一致性断言**;补 platform 审计 input_summary 的 JSON 解析断言(含引擎输出/阻断 reasons,B5 复审 L1) | 测试套件 + fixture | `pytest` 全绿 | B7 | + +**B7 复审(独立 AI 评审 · 2026-09-06):有条件通过 → 已闭环。** 挂账①②③④⑤⑦证实落地、⑧确认随 T-02 无烂尾、红线全守住。P1-1 `dispose_engines` 仅 cache_clear 未真 dispose——已修(`utils/db.py` 手写单例字典 + 显式 `Engine.dispose()`,新增 `tests/test_db.py` 3 例,183 绿)。挂账:P2-1 L3 DEL 钩子无行为断言随 **B8**;P2-2 未捕获异常 500 无 trace 头回写 + 422/404/405 错误码补齐归 **T-02**;P3-1 `update_alert_status` 建议标注 deprecated(防绕过同事务审计)、P3-2 locks 降级文案中性化(可随 B8 顺手);P3-3 生产误配 `app_env=development` 检查归 **B9b SOP**;P3-4 独立 request_id(现沿用 trace_id)归 **T-02**。偏差留痕:成功响应不包裹外壳 = B5 P2-2「返回体不变」演进口径,M4 复核时对齐开发计划文本。 +| B8 | **`tests/conftest.py`**:a) session fixture 启动校验演示数据就位(CUST-4001 测评 <365 天、risk_aml_list ≥8),缺失则中止并提示先跑 FLOW §0 ③④;b) fixture 幂等代跑 `prepare_risk_demo.sql`;c) teardown 按 `TRD-TEST-` 清 core_trade + 关联 risk_alert/risk_suitability_log/audit_log + 还原 L3 行。**顺手集中 sqlite 测试 DDL 为单一事实源(B4 评审 P3-12,各测试文件手写 DDL 收敛;CURRENT_TIMESTAMP 改 localtime 或 fixture 固定时间,防 UTC/本地日界错位——B5 评审 P3-4)**;集成测试交易统一走 `trade_id_factory` 注入 `TRD-TEST-` 前缀(trade_gateway 已留参数,B5 评审 P3-2)。集成测试:A-1~A-5、A-7(状态机/compliance 403/GET 强制 aml)、A-9 越权、**trace 一致性断言**;补 platform 审计 input_summary 的 JSON 解析断言(含引擎输出/阻断 reasons,B5 复审 L1);**B7 复审 P2-1:补 L3 DEL 钩子行为断言(fake 记录 deletes,断言 key=`profile:l3:{customer_id}` 与降级路径);P3-2 locks 降级文案中性化顺手改** | 测试套件 + fixture | `pytest` 全绿 | B7 | | B9a | 演示/运维脚本开发:`scripts/demo/subscribe_alerts.py`(订阅演示)+ `scripts/demo/rebuild_alerts.py`(按 trade_id 幂等重放补偿) | 2 个脚本 | 手工执行验证 | B2、B4(可与 B5~B8 并行) | -| B9b | 演示链路走查:`reset.ps1` → `prepare_risk_demo.sql` → agent 库建表 → `seed-aml-list.sql` → Swagger 逐条过 **A-1~A-5、A-7~A-9(A-6 归 M3)**。**B6 挂账核查单:① 预警类 API 响应体含固定 disclaimer「本预警由系统自动生成,最终判定需经风控专员人工审核」(PRD §6/规则表 §5,B6 复审 P3-7)② aml/scan 幂等防护(重复扫描同命中客户重复出单,演示点击即复现,B6 评审 P3-6)③ B6 时代码以 TestClient 独立挂 router 等价验证,本次补一次真 Swagger 手测(B6 复审遗漏⑤)④ analyst 台账只读权限扩展待 Wave 3 分析 Agent 接入时定权限矩阵(手册 §5.3 有 risk:alert:read,现 fail-closed 拒绝,B6 复审观察③)** | 演示 SOP | 按 PRD §8 验收表逐条打勾 | B8、B9a | +| B9b | 演示链路走查:`reset.ps1` → `prepare_risk_demo.sql` → agent 库建表 → `seed-aml-list.sql` → Swagger 逐条过 **A-1~A-5、A-7~A-9(A-6 归 M3)**。**B6 挂账核查单:① 预警类 API 响应体含固定 disclaimer「本预警由系统自动生成,最终判定需经风控专员人工审核」(PRD §6/规则表 §5,B6 复审 P3-7)② aml/scan 幂等防护(重复扫描同命中客户重复出单,演示点击即复现,B6 评审 P3-6)③ B6 时代码以 TestClient 独立挂 router 等价验证,本次补一次真 Swagger 手测(B6 复审遗漏⑤)④ analyst 台账只读权限扩展待 Wave 3 分析 Agent 接入时定权限矩阵(手册 §5.3 有 risk:alert:read,现 fail-closed 拒绝,B6 复审观察③)⑤ 生产/演示机 `app_env=development` 误配检查(B7 复审 P3-3:误配时 debug 头可达且启动校验放行,列入演示 SOP)** | 演示 SOP | 按 PRD §8 验收表逐条打勾 | B8、B9a | ## 阶段 C · 对话线(依赖 Wave 0 的 T-01 JWT / T-03 输入防护 / T-07 LangGraph) diff --git a/tests/test_db.py b/tests/test_db.py new file mode 100644 index 0000000..1fb84b0 --- /dev/null +++ b/tests/test_db.py @@ -0,0 +1,54 @@ +"""引擎工厂(utils/db.py)单测(B7 评审 P1-1):单例复用 + dispose 真实释放。 + +create_engine 与缓存字典均 monkeypatch 替换,不触网;验证 dispose_engines +对每个缓存 Engine 显式调用 dispose() 并清空缓存(不再是仅 cache_clear)。 +""" + +from app.utils import db + + +class SpyEngine: + def __init__(self): + self.dispose_calls = 0 + + def dispose(self, close: bool = True) -> None: + self.dispose_calls += 1 + + +def _patch(monkeypatch): + created_urls = [] + + def fake_create_engine(url, **kwargs): + created_urls.append(url) + return SpyEngine() + + monkeypatch.setattr(db, "create_engine", fake_create_engine) + monkeypatch.setattr(db, "_engines", {}) + return created_urls + + +def test_get_engine_caches_one_engine_per_database(monkeypatch): + created = _patch(monkeypatch) + e1 = db.get_engine("db_a") + e2 = db.get_engine("db_a") + e3 = db.get_engine("db_b") + assert e1 is e2 + assert e1 is not e3 + assert len(created) == 2 + + +def test_dispose_engines_calls_dispose_and_clears_cache(monkeypatch): + _patch(monkeypatch) + e = db.get_engine("db_a") + db.dispose_engines() + assert e.dispose_calls == 1 + assert db._engines == {} + # 缓存已清:重建新实例而非复用已 dispose 的旧实例 + fresh = db.get_engine("db_a") + assert fresh is not e + assert fresh.dispose_calls == 0 + + +def test_dispose_engines_on_empty_cache_is_noop(monkeypatch): + _patch(monkeypatch) + db.dispose_engines() # 不抛异常即可