Files
group_xinghuo_jinrong/docs/项目框架设计/架构设计-风控模块.md
GaoYiYuan_0626 037ce7edca docs(架构改进): 补齐 PRD/开发计划/TODO/交接文档,落地无密钥告警与 Redis 分布式锁
一、流程文档(按 AIcoding 六步落地,供新会话从交接文档开工)
- 新增 docs/PRD/PRD-架构改进与稳定性加固.md:6 条 FR(文档勘误、非缺陷说明、
  无密钥启动告警、审计失败告警、Redis 分布式锁、中间件顺序测试)
- 新增 docs/项目框架设计/改进方案评审-问题清单与对比.md:24 项问题分档 A~G,
  经两轮独立 AI 评审,无阻断级错误
- 新增 docs/项目框架设计/开发计划-架构改进.md:HOW 层设计,含合并前只做低风险
  11 项的批次策略
- 新增 docs/项目框架设计/TODO-架构改进.md:T-101~T-109、T-201~T-202 可勾选项
- 新增 docs/交接文档-架构改进.md:自包含交接入口,hy3 新会话可直接开工
- 新增 docs/项目框架设计/架构设计说明书.md:按模块/分层逐一讲解的全量架构说明

二、代码改动(T-107/108/109、T-201.1、T-201.2)
- app/main.py:启动时 DEEPSEEK_API_KEY 缺失告警,明确告知将走降级回复
- app/utils/authz.py:越权审计失败日志补 trace_id,便于串联全链路
- app/api/audit_middleware.py:审计失败日志补 status/path/request_id
- app/service/risk/redis_gateway.py:新增 acquire_lock(SET NX EX)与
  release_lock(Lua 原子释放,只删自己的锁)
- app/service/risk/locks.py:run_locked 改为双层锁,Redis 为主、进程内锁为备;
  Redis 超时沿用 fn(locked=False) 降级语义,Redis 不可用(含测试 Fake 缺方法的
  AttributeError)安全退回进程内锁,绝不抛异常

三、文档勘误(A1/A2/A3)
- MEMORY.md:文件数 42→45、Tools 4→5
- 02-mysql-agent专用.sql:会话表 5→6
- 架构设计-风控模块.md:同步更正

四、测试
- 新增 tests/test_locks_redis.py:覆盖抢锁成功、占用超时、Redis 故障降级、
  Fake 缺方法降级、只删自己锁、三处调用点 key 前缀
- tests/test_audit_middleware.py:补充告警字段断言
- 全量 pytest 510 passed(原基线 503)
2026-09-09 18:10:03 +08:00

269 lines
16 KiB
Markdown
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.
# 架构设计说明书 · 风控监测 Agent 模块
> 版本:v1.1(AI 评审通过:修复 P1×5——Redis 同步选型、并发首单锁方案、contextvars 实现约束、测试基建、rebuild 降级为脚本;P2×7 顺手收敛)· 2026-09-06
> 分支:`risk-control-agent`(旧称 feature/risk 已过时)· 负责人:E
> 上游:`docs/PRD/PRD-风控监测Agent.md`、`docs/PRD/附-风控规则表.md`、`docs/memory/FRAMEWORK.md`
> 技术选型不变(FastAPI + SQLAlchemy + MySQL 双库 + Redis + LangGraph + DeepSeek),本文只定义**目录划分、模块职责、核心时序、关键技术决策**。
---
## 1. 设计原则
1. **分层遵守 FRAMEWORK §3**:api(薄)→ service(业务)→ repository/tool;**唯一例外** `app/gateway/`(模拟外部交易系统,PRD 授权,仅 `gateway_repository` 可 INSERT `core_trade`)
2. **事件驱动线不用消息队列**:网关落库后进程内同步调规则引擎(阻断演示可靠性优先);Redis Pub/Sub 只做**预警通知广播**,不做事件总线
3. **规则即纯函数**:RISK-001~005 与 SUIT-001~008 全部实现为无副作用纯函数,输入事实、输出判定,便于单测与阈值调整
4. **双库无跨库事务**:`core_trade`(core 库)与预警单(agent 库)分开提交,取舍见 §5.3
5. **审计与 trace_id 用 contextvars 贯通**,业务代码不手工传递
---
## 2. 目录与文件划分(新增)
```text
app/
├── gateway/ # 【PRD 例外层】模拟交易网关 = 外部 Core 交易系统替身
│ ├── __init__.py
│ ├── trade_gateway.py # 网关编排:生成 trade_id/trace_id → suitability → 落库 → 触发引擎
│ └── gateway_repository.py # 仅 INSERT jinrong_core.core_trade(不碰其他 core 表)
│
├── service/
│ ├── suitability.py # 公共校验:SUIT-001~008 纯函数 + 落 suitability_log
│ └── risk/
│ ├── __init__.py
│ ├── engine.py # 引擎入口 process_trade_event(trade) → 跑规则 → 聚合出单
│ │ # → L3 → 推送;预留空钩子 on_customer_created/
│ │ # on_customer_updated(AML 开户触发,本期 no-op)
│ ├── rules.py # RISK-001~005 纯函数(输入当日流水上下文,输出命中列表)
│ ├── alert_service.py # 预警单聚合/去重/落库/更新 + Pub/Sub 推送
│ ├── aml_service.py # AML:姓名归一化 + difflib 相似度匹配 + 全量扫描
│ ├── profile_l3.py # L3 UPSERT(最高档合并 + tags 追加)后 DEL profile:l3:{cid} 缓存
│ └── scoring.py # R-05 预留桩(非实际评分器):recompute_customer_score(customer_id) 直接 raise NotImplementedError(签名已冻结,待 R-05 实现);L3 的 risk_score 一期恒 NULL,由 R-05 首写,避免静态分污染语义
│
├── repository/
│ ├── core_ro.py # 【扩展(最小化)】SUIT 校验复用现有 get_customer_l0
│ │ # (已含 age+risk_code+evaluated_at)与 get_product;
│ │ # 仅新增(最小化,均纯 SELECT):sum_trades_on_date
│ │ # (当日累计)/ list_trades_range(规则流水,升序防
│ │ # RISK-005 截断,B4 评审 P2-5)/ list_active_customers
│ │ # (AML scan_all,B4)
│ └── risk_repository.py # agent 库:risk_alert / risk_suitability_log /
│ # customer_profile_l3 / risk_aml_list
│
├── api/
│ ├── simulate.py # POST /api/simulate/trade(薄路由 → gateway)
│ └── risk.py # GET /api/risk/alerts、POST /api/risk/alerts/{id}/handle、
│ # POST /api/risk/suitability/check、POST /api/risk/aml/scan
│
├── config/
│ └── settings.py # 【扩展】risk_* 阈值配置(§6)
├── utils/
│ ├── trace.py # 【新增】contextvars trace_id 生成/读取
│ └── desensitize.py # 【新增】DESENS-001~005 脱敏纯函数
└── main.py # 【扩展】挂载 simulate/risk 路由;lifespan:共享 Engine 单例
# (core+agent 两池,注入各 repository)+ 同步 Redis 单例
# + trace 中间件注册(T-01 后再加审计中间件)
chat 对话线(T-32,依赖 T-01/T-03/T-07 就绪后实施):
├── service/agent_service.py # StateGraph(risk) 分支注册(见 §4)
└── service/risk/chat_tools.py # 4 个 Tool:alert_query / customer_context / suitability_check / aml_lookup
scripts/
├── agent/seed-aml-list.sql # 已建:AML 名单种子(含演示命中记录)
└── demo/prepare_risk_demo.sql # 已建:演示测评日期刷新(reset 后重跑)
tests/
├── test_suitability.py # SUIT 矩阵全组合 + 70 岁/过期/NULL 年龄边界
├── test_risk_rules.py # RISK-001~005 纯函数用例
├── test_aml_service.py # 归一化/相似度边界(B4 实际命名,评审 P3-10)
└── test_gateway_flow.py # 集成:A-1~A-5 阻断与放行链路
```
**不新增**:tool/(风控不用 Milvus)、model/entities 中无新 ORM(沿用 SQL 文本 + Pydantic schema,与 core_ro 风格一致)。
---
## 3. 核心时序
### 3.1 交易放行链路(正常)
```text
POST /api/simulate/trade
→ api/simulate.py(JWT:risk_demo 或本人)生成 trace_id → utils/trace.set()
→ gateway/trade_gateway.py
① 参数校验(trade_type ∈ {subscribe,redeem},convert→400)
② suitability.py :: check(customer_id, product_id)
读 core_ro:customer_risk(C 级+evaluated_at) / customer.age / product.min_risk_code
→ SUIT-001~008 纯函数判定 → 落 risk_suitability_log
③a 不匹配 → alert_service.suitability_alert(去重追加) → 返回 blocked=true+阻断文案
③b 匹配 → gateway_repository.insert_trade(core_trade) [core 库提交]
④ risk/engine.py :: process_trade_event(trade)
rules 跑 RISK-001~005 + aml_service.match_name
→ 命中 → alert_service.agg_upsert(聚合/去重)→ profile_l3.upsert
→ audit_log INSERT → PUBLISH risk:pub:alert
→ 未命中 → audit_log INSERT(decision=pass)
⑤ 返回 blocked=false + trade_id + 触发规则列表(如有)
```
### 3.2 AML 全量扫描
```text
POST /api/risk/aml/scan(risk_officer)
→ aml_service.scan_all():core_ro 全客户 ↔ risk_aml_list(is_active)
归一化全等 → 命中;否则 difflib.SequenceMatcher 相似度 ≥ threshold → 命中
→ 每命中客户:独立 aml 预警单(score=95)+ L3 high + audit + PUBLISH(notify_role 含 compliance)
```
### 3.3 人工处置
```text
POST /api/risk/alerts/{id}/handle(risk_officer)
→ 状态机校验 pending_review → 目标态(其余 409)→ UPDATE handler_* 字段
→ audit_log INSERT(agent_type='risk')
```
---
## 4. 对话线 LangGraph 图(T-32)
```text
StateGraph: RiskAgentState(messages, intent, tool_results, auth_context)
entry → intent_node # LLM(DeepSeek) 分类:alert_query / customer_context /
# suitability_check / aml_lookup / chitchat
intent_node →(条件边)→ tool_node[alert_query | customer_context | suitability_check | aml_lookup]
# chitchat 直达 respond
tool_node → respond_node # DeepSeek 汇总;规范:引用必带 alert_id、
# 评分/分层标注"仅供参考,不自动决策"、不输出处置动作
respond_node → END
边界:对话线除 suitability_check(按 FR-2 落日志与预警,属其定义的副作用)外无写 Tool;
处置引导文案指向 POST /alerts/{id}/handle,图中无其他写 Tool
```
Tool 数据源:`risk_repository`(预警台账只读)、`core_ro`(L0)、`customer_profile_l1/l2/l3`(只读);入参均过 AuthContext 归属校验(risk_officer 全量)。
---
## 5. 关键技术决策
### 5.1 trace_id 贯通(contextvars · 含实现约束)
```python
# utils/trace.py
_trace_id: ContextVar[str] = ContextVar("trace_id", default="")
def new_trace() -> str: tid = f"trc-{uuid4().hex[:16]}"; _trace_id.set(tid); return tid
def current() -> str: return _trace_id.get()
```
- api 层中间件:优先透传请求头 `X-Trace-Id`,否则生成
- suitability_log / risk_alert / audit_log / Pub/Sub payload 全部 `trace.current()`,业务函数不传参
**实现约束(防踩坑,编码前必读)**:
1. 中间件必须**在 `call_next` 之前** `set()`——若用 `BaseHTTPMiddleware`,endpoint 内已读不到后续设置;推荐纯 ASGI 中间件或放在鉴权依赖最前
2. 事件线为同步 def 路由(线程池执行),context 经 `anyio.to_thread` 传播需在 §7 集成测试用断言实测确认
3. **禁止**在 endpoint/service 内重新 `set()`;后续若引入 `create_task` 后台化必须显式 `contextvars.copy_context()`
### 5.2 聚合去重实现(PRD FR-4 · 首单并发正确性)
- **先取锁再 check-insert**(`risk_alert` 无 `(customer_id, date)` 唯一键且冻结不改表,裸 SELECT FOR UPDATE 在并发首单时空结果集锁不住,会产生双单):
- 进程内锁:`locks.py` 用 `threading.Lock` 字典,key=`agg:{customer_id}:{alert_class}:{date}`(alert_class ∈ event/suitability);**这是进程内锁,仅活在当前进程内存,多实例部署时两实例各有各的锁、同一客户并发会同时进临界区(原设计为 TODO,非已完成)**。本期(架构改进 T-201)已将其改造为「Redis 分布式锁为主 + 进程内锁为备」的双层结构:Redis 不可用时退回进程内锁,单进程行为不变;多实例场景由 Redis 锁提供跨进程互斥。
- **锁内**普通 `SELECT ... WHERE customer_id=? AND status='pending_review' AND alert_type IN (事件类) AND created_at >= 当日 00:00` → 存在则读改写 `payload.events[]` + `triggered_rules` 合并 + `risk_score=max` + `alert_type=最高分规则类型`;不存在 INSERT
- **降级策略**:锁获取超时(>2s)不阻塞交易——放行并照常独立出单(宁多勿漏),记 logger.warning
- suitability 单:同上,键 `customer_id + alert_class='suitability' + payload.product_id`
- Redis `risk:dedup:{cid}:{rule_id}:{date}` 仅作前置短路(规则级防重入),DB 查询为准
- `date` 口径:服务器本地时区自然日(`datetime.now()`,`traded_at` 同源)
### 5.3 双库一致性取舍
`core_trade` 与预警单分属两库、无跨库事务。顺序固定为 **先 core 后 agent**(交易成立是事实,预警可补偿);引擎异常时 `except` 记 `audit_log(decision='risk_engine_error')` + logger.error。
**补偿入口为运维脚本**(不进 HTTP 面,避免动冻结 PRD 的 API 清单):`scripts/demo/rebuild_alerts.py`——按指定日期重放 `core_trade` 给引擎。**重放幂等**:处理每笔前按 `trade_id` 检查该客户当日预警单 `payload.events[]` 是否已含此 `trade_id`,已存在则跳过(防 events/monitor_tags 重复追加)。**不做**分布式事务。
### 5.4 金额与时间
- 全链路 `Decimal`(SQLAlchemy DECIMAL 原生映射);阈值比较 `Decimal` vs `Decimal(str(env))`
- 时间统一服务器本地时间 `datetime.now()`;`traded_at` 由网关写入;"当日/5 分钟窗口"均基于 `traded_at`
### 5.5 AML 匹配算法
```python
norm = lambda s: re.sub(r"\s+", "", s).lower() # 归一化
hit = (norm(name) == norm(list_name)) or \
(SequenceMatcher(None, norm(name), norm(list_name)).ratio() >= threshold)
```
- 一期仅 `full_name`(`id_no/bank_card_no` 字段在表中预留,Core 有证件数据后启用精确匹配)
- 阈值默认 0.85,取 `risk_aml_list.match_threshold`(每条名单可独立配置)
### 5.6 Redis 客户端(同步/异步分线)
- **事件驱动线:同步 `redis.Redis` 单例**(与同步 Engine、同步调用链一致;lifespan 创建、shutdown 关闭);Pub/Sub PUBLISH 为同步调用
- **对话线(T-32):LangGraph 异步栈自行引入 `redis.asyncio`**,由 agent_service 管理
- **两套客户端不共用连接**,各自独立单例;订阅端(演示)见 `scripts/demo/subscribe_alerts.py`
### 5.7 归属校验落点(全模块统一)
- 归属校验统一放 **FastAPI 依赖层**(`Depends(get_auth_context)` 内:JWT 解析 + 角色判定 + 归属断言——customer 需 `customer_id == subject`、advisor 走 `core_ro.is_advisor_assigned`(已有方法)、risk_officer 放行、compliance 强制 `alert_type=aml`)
- service 层不重复校验,函数签名接收 `AuthContext`;越权 403 + audit 由依赖层统一产出
---
### 5.8 设计标注(非缺陷项 · 评审记录)
下列三项在评审中被列为「非问题」,单独列出避免读者误判为缺陷:
- **G1 · `X-Trace-Id` 允许前端伪造**:`main.py` trace 中间件对传入头做白名单 `^[A-Za-z0-9._-]{1,64}$` 校验,不合规一律丢弃并新生成,防响应头注入。**身份只认 JWT,trace_id 绝不参与任何权限判定**,伪造它最多改个日志追踪号,动不了鉴权。
- **G2 · `trace_id` 概率唯一**:生成用 `uuid4().hex[:16]`(64 bit 熵),为请求级标识,约 2^32 次请求才有 50% 碰撞概率,无需数据库唯一性校验;碰撞最坏后果是日志追踪归并错误,非数据问题。
- **G3 · `admin.py` / `knowledge.py` 空壳**:T-21 拍板一期只做脚本入库(知识库已建集合并 upsert,能力在 `service/rag_service.py`),这两文件未挂载路由。**空 ≠ RAG 缺失**,RAG 检索能力已就绪,缺的只是前端上传/重建端点(产品排期后置)。
---
## 6. 配置新增(`.env.example` 同步)
```ini
# ===== Risk(默认值=冻结规则,均可覆盖)=====
RISK_ASSESSMENT_VALID_DAYS=365
RISK_LARGE_AMOUNT=500000
RISK_DAILY_TOTAL=500000
RISK_FREQ_COUNT=3
RISK_PROBE_WINDOW_MINUTES=5
RISK_PROBE_COUNT=3
RISK_PROBE_AMOUNT=400000
RISK_SMALL_AMOUNT=10000
RISK_SMALL_COUNT=3
RISK_AML_DEFAULT_THRESHOLD=0.85
```
---
## 7. 测试与验收策略
**测试基建**:`requirements.txt` 增补 `pytest`(+`pytest-asyncio` 供 T-32);**不引入第三套测试库**——复用本地 `jinrong_core`/`jinrong_agent` 两库,集成测试前置跑 `reset.ps1` + `prepare_risk_demo.sql`,测试产生的交易统一用 `trade_id` 前缀 `TRD-TEST-`,teardown 按 `trade_id LIKE 'TRD-TEST-%'` + 对应预警/日志清理,不污染演示数据。
| 层 | 内容 | 对照 |
| --- | --- | --- |
| 单测 | SUIT 矩阵 25 组合 + SUIT-006(69/70 岁边界、NULL 年龄)+ SUIT-008(364/365 天)| A-8 |
| 单测 | RISK-001~005 纯函数(含 RISK-005 非连续、RISK-004 窗口边界) | 规则表 §2 |
| 单测 | AML 归一化/相似度边界(阈值±) | 规则表 §3 |
| 集成 | 网关链路 A-1~A-5、A-9 越权 | PRD §8 |
| 集成 | 处置状态机:pending→三态、重复处置 409、compliance 调 handle 403 且 GET 仅返 aml 单 | A-7 |
| 集成 | **trace 一致性断言**:响应体 trace_id == suitability_log == risk_alert == audit_log | PRD §7.5 |
| 对话线 | A-6 台账问答(T-32 实施,依赖 T-01/T-03/T-07) | A-6 |
| 演示 | `scripts/demo/prepare_risk_demo.sql` → Swagger 逐条过 A-1~A-9 | PRD §8 |
---
## 8. 与 Wave 0 的集成点(依赖提醒)
| 集成点 | 说明 |
| --- | --- |
| T-01 JWT | `api/simulate.py`、`api/risk.py` 全部走统一鉴权依赖;`risk_demo` 角色判定在鉴权依赖内 |
| T-02 audit 中间件 | 风控复用统一 audit 写入;网关 `agent_type='platform'`、处置 `'risk'` |
| T-03 输入防护 | 仅对话线(专员输入);事件线无用户输入不受影响 |
| T-04 core_ro | 本模块扩展方法见 §2;扩展后仍是纯 SELECT |
| T-07 LangGraph | 仅 §4 对话线;事件驱动线完全独立可先行 |
---
*开发计划与任务拆分见下一步《开发计划-风控模块》。*