From 03f8b6378030f3a74cec5dcb59f7104e7534dd07 Mon Sep 17 00:00:00 2001 From: YUAN Date: Sun, 6 Sep 2026 17:40:13 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20FR-4=20payload=20=E4=B8=8A=E4=B8=8B?= =?UTF-8?q?=E6=96=87=E8=A1=A5=E7=BB=84=E8=A3=85/=E7=A7=8D=E5=AD=90?= =?UTF-8?q?=E5=90=8D=E5=8D=95=E5=94=AF=E4=B8=80=E5=91=BD=E4=B8=AD/?= =?UTF-8?q?=E8=A7=84=E5=88=99=E6=B5=81=E6=B0=B4=E5=8D=87=E5=BA=8F=E7=AD=89?= =?UTF-8?q?(B4=20=E8=AF=84=E5=AE=A1=20P1-1=E3=80=81P2-2~5=E3=80=81P3=C3=97?= =?UTF-8?q?7)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/repository/core_ro.py | 27 ++++++++++++++ app/service/risk/aml_service.py | 12 +++++- app/service/risk/engine.py | 51 +++++++++++++++++++++++--- app/utils/trace.py | 6 +++ docs/memory/TODO.md | 2 +- docs/项目框架设计/开发计划-风控模块.md | 4 +- docs/项目框架设计/架构设计-风控模块.md | 8 ++-- scripts/agent/seed-aml-list.sql | 8 ++-- tests/test_aml_service.py | 24 ++++++++++++ tests/test_risk_engine.py | 27 ++++++++++++++ 10 files changed, 153 insertions(+), 16 deletions(-) diff --git a/app/repository/core_ro.py b/app/repository/core_ro.py index f01f2e8..2736803 100644 --- a/app/repository/core_ro.py +++ b/app/repository/core_ro.py @@ -77,6 +77,33 @@ class CoreReadOnlyRepository: ).mappings() ] + def list_trades_range( + self, customer_id: str, start: datetime, end: datetime, limit: int = 10000 + ) -> list[dict[str, Any]]: + """[start, end) confirmed 申赎流水,**时间升序**,不 JOIN(规则统计专用)。 + + 升序保证 RISK-005「先小后大」的最早小额铺垫不被截断(B4 评审 P2-5); + 去掉 product_name JOIN 防脏产品数据丢行。规则演示规模远低于 limit。 + """ + sql = text( + """ + SELECT * FROM core_trade + WHERE customer_id = :cid + AND trade_type IN ('subscribe', 'redeem') + AND trade_status = 'confirmed' + AND traded_at >= :start AND traded_at < :end + ORDER BY traded_at ASC + LIMIT :lim + """ + ) + with self._engine.connect() as conn: + return [ + dict(r) + for r in conn.execute( + sql, {"cid": customer_id, "start": start, "end": end, "lim": limit} + ).mappings() + ] + def sum_trades_on_date(self, customer_id: str, day: date) -> Decimal: """当日申赎合计金额(RISK-002 累计口径:仅 confirmed 的 subscribe/redeem)。 diff --git a/app/service/risk/aml_service.py b/app/service/risk/aml_service.py index 6c64edc..6daa159 100644 --- a/app/service/risk/aml_service.py +++ b/app/service/risk/aml_service.py @@ -2,9 +2,14 @@ 一期降级口径:仅 display_name 归一化(去空白 + 大小写折叠)+ difflib 相似度; 证件/银行卡匹配待 Core 提供证件数据后启用(表 id_no/bank_card_no 已预留)。 -阈值:名单行 match_threshold 优先(表默认 0.85),缺省回落 settings.risk_aml_default_threshold。 +阈值:名单行 match_threshold 优先(表默认 0.85),仅缺 NULL 时回落 +settings.risk_aml_default_threshold(is None 判断,显式 0 不误回落——评审 P3-7)。 命中动作(独立 aml 预警单 + L3 high 标记)由调用方编排:engine(交易触发)/ scan_all(手动全量扫描);不冻结、不自动上报(附表 §2 行为边界)。 + +脱敏留痕(评审 P3-9):matched_name 进 payload/审计依赖种子「脱敏展示名口径」 +(名单表 full_name 注释);接入真实名单数据时须在出口接 utils/desensitize, +挂账见开发计划 B4 行。 """ from __future__ import annotations @@ -18,6 +23,7 @@ from app.repository.core_ro import CoreReadOnlyRepository from app.repository.risk_repository import RiskRepository from app.service.risk.alert_service import record_aml_alert from app.service.risk.profile_l3 import upsert_profile_l3 +from app.utils.trace import ensure_trace logger = logging.getLogger(__name__) @@ -42,7 +48,8 @@ def match_name( norm = normalize_name(customer_name) hits: list[dict[str, Any]] = [] for e in entries: - threshold = float(e.get("match_threshold") or settings.risk_aml_default_threshold) + raw = e.get("match_threshold") + threshold = float(raw) if raw is not None else float(settings.risk_aml_default_threshold) ratio = similarity(norm, normalize_name(e["full_name"])) if ratio >= threshold: hits.append( @@ -83,6 +90,7 @@ def scan_all( """ core = core_ro or CoreReadOnlyRepository() repo = risk_repo or RiskRepository() + ensure_trace() # 手动扫描入口兜底归因(B6 API 场景保留中间件 trace,评审 P3-8) entries = repo.list_active_aml_entries() customers = core.list_active_customers() alerts: list[str] = [] diff --git a/app/service/risk/engine.py b/app/service/risk/engine.py index 9afc749..de18fe4 100644 --- a/app/service/risk/engine.py +++ b/app/service/risk/engine.py @@ -14,7 +14,8 @@ RISK-004 窗口以 trade["traded_at"] 为事件时点(非墙钟 now):rebui from __future__ import annotations -from datetime import datetime +from datetime import datetime, timedelta +from decimal import Decimal from typing import Any from app.repository.core_ro import CoreReadOnlyRepository @@ -23,6 +24,7 @@ from app.service.risk.alert_service import record_aml_alert, record_trade_alerts from app.service.risk.aml_service import match_customer from app.service.risk.profile_l3 import upsert_profile_l3 from app.service.risk.rules import RiskThresholds, run_rules +from app.utils.trace import ensure_trace def _as_datetime(value: Any) -> datetime: @@ -41,6 +43,29 @@ def _normalize_trades(trades: list[dict[str, Any]]) -> list[dict[str, Any]]: return trades +def _build_customer_context( + core: CoreReadOnlyRepository, customer_id: str, day_start: datetime +) -> dict[str, Any]: + """预警 payload 客户上下文(PRD FR-4:L0 事实 + 近 30 天交易统计;B4 评审 P1-1 补组装)。 + + display_name 为脱敏展示名口径,可直存 payload;core_cash_flow 上下文一期 + 未接(core_ro 无对应查询,挂账见开发计划 B4 行)。 + """ + l0 = core.get_customer_l0(customer_id) or {} + month_ago = day_start - timedelta(days=30) + day_end = day_start + timedelta(days=1) # 近 30 天含当日(本笔在内) + trades_30d = core.list_trades_range(customer_id, month_ago, day_end) + total = sum((Decimal(str(t["amount"])) for t in trades_30d), Decimal(0)) + return { + "l0": { + k: l0[k] + for k in ("customer_id", "display_name", "age", "risk_code") + if l0.get(k) is not None + }, + "trades_30d": {"count": len(trades_30d), "total_amount": str(total)}, + } + + def process_trade_event( trade: dict[str, Any], core_ro: CoreReadOnlyRepository | None = None, @@ -55,17 +80,25 @@ def process_trade_event( core = core_ro or CoreReadOnlyRepository() repo = risk_repo or RiskRepository() th = thresholds or RiskThresholds.from_settings() + ensure_trace() # 脚本/重放入口兜底归因(中间件场景保留现有 trace,评审 P3-8) event_at = _as_datetime(trade["traded_at"]) day_start = event_at.replace(hour=0, minute=0, second=0, microsecond=0) - trades = _normalize_trades(core.list_trades(trade["customer_id"], since=day_start, limit=1000)) + day_end = day_start + timedelta(days=1) + trades = _normalize_trades( + core.list_trades_range(trade["customer_id"], day_start, day_end) + ) result: dict[str, Any] = {"triggered_rules": [], "alert_ids": [], "aml_hit": False} hits = run_rules(trades, th, now=event_at) - # 无条件走预警编排:空 hits 由 alert_service 落 pass 审计(架构 §3.1 ④ 未命中分支) - alert = record_trade_alerts(trade, hits, risk_repo=repo) if hits: + alert = record_trade_alerts( + trade, + hits, + risk_repo=repo, + customer_context=_build_customer_context(core, trade["customer_id"], day_start), + ) result["triggered_rules"] = sorted({h.rule_id for h in hits}) if alert: result["alert_ids"].append(alert["alert_id"]) @@ -76,13 +109,21 @@ def process_trade_event( last_alert_id=alert["alert_id"] if alert else None, risk_repo=repo, ) + else: + # 未命中分支:pass 审计由 alert_service 统一落库(架构 §3.1 ④) + record_trade_alerts(trade, [], risk_repo=repo) aml_hits = match_customer(trade["customer_id"], core_ro=core, risk_repo=repo) if aml_hits: result["aml_hit"] = True alert = record_aml_alert( trade["customer_id"], - {"trigger": "trade", "trade_id": trade.get("trade_id"), "matches": aml_hits}, + { + "trigger": "trade", + "trade_id": trade.get("trade_id"), + "product_id": trade.get("product_id"), # 评审 P3-6:payload/审计透传 + "matches": aml_hits, + }, risk_repo=repo, ) result["alert_ids"].append(alert["alert_id"]) diff --git a/app/utils/trace.py b/app/utils/trace.py index 3dd2569..bde6235 100644 --- a/app/utils/trace.py +++ b/app/utils/trace.py @@ -24,3 +24,9 @@ def new_trace(trace_id: str | None = None) -> str: def current_trace() -> str: """读取当前 trace_id;未初始化时返回空串(调用方应兜底生成)。""" return _trace_id.get() + + +def ensure_trace() -> None: + """无上下文时兜底归因(脚本/引擎入口),有值时保留(中间件场景不重新 set)。""" + if not current_trace(): + new_trace() diff --git a/docs/memory/TODO.md b/docs/memory/TODO.md index 75669e0..2bdc9be 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 +- [ ] 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+引擎编排 均已完成并经独立 AI 评审闭环;剩 B5 网关、B6 鉴权+API、B7 main 集成、B8 conftest+集成测试、B9a 脚本、B9b 演示走查** - [ ] T-31 `service/suitability.py` 公共校验(SUIT-001~008)+ `POST /api/risk/suitability/check` + 单测验收 A-8 —— **代码已完成(A1~A4 · 77 测试绿);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 diff --git a/docs/项目框架设计/开发计划-风控模块.md b/docs/项目框架设计/开发计划-风控模块.md index fbff5d3..a90ac8a 100644 --- a/docs/项目框架设计/开发计划-风控模块.md +++ b/docs/项目框架设计/开发计划-风控模块.md @@ -23,11 +23,11 @@ | B1 | `service/risk/rules.py`:RISK-001~005 纯函数 | 规则函数 | **单测:各规则命中/不命中 + RISK-004 窗口 + RISK-005 非连续** | A1 | | B2 | `alert_service.py`:聚合去重(进程内锁 + 锁内 check-insert)+ 审计落库 + `risk:pub:alert` PUBLISH(同步 Redis 单例)。**备注:merge 原语(append_alert_event)已下沉 repo(A3),本任务只做聚合决策与编排(评审 P2-5 口径)** | 预警服务 | 单测:聚合合并、risk_score 取 max、去重追加;**并发冒烟:两线程同客户同日首单 → 预警单数=1 且 events[] 含两笔** | A3 | | B3 | `profile_l3.py`:get_l3 → 最高档合并 → insert_l3/update_l3。**备注:非原子,并发首单需 catch IntegrityError 转更新(或复用 B2 锁)**(评审 P2-7③);**B3 评审后口径(P3-4 · 用户拍板 2026-09-06):L3 `risk_score` 一期不写(保持 NULL,归 R-05 评分模型首写),tier/tags/last_alert_id/computed_at 照 FR-7** | L3 写入 | 单测:normal→high 不降级、AML 后大额不回落 | A3 | -| B4 | `aml_service.py`(归一化+相似度匹配、scan_all)+ `engine.py`(process_trade_event 组装 + 预留客户事件钩子)+ **`scoring.py` 占位签名(FR-7 预留)** | 引擎完整 | 单测:AML 阈值边界;引擎集成冒烟 | B1、B2、B3 | +| B4 | `aml_service.py`(归一化+相似度匹配、scan_all)+ `engine.py`(process_trade_event 组装 + 预留客户事件钩子)+ **`scoring.py` 占位签名(FR-7 预留)** | 引擎完整 | 单测:AML 阈值边界;引擎集成冒烟 | B1、B2、B3。**B4 评审修复(2026-09-06):FR-4 payload 客户上下文(L0+近 30 天统计)补组装;core_ro 增 list_trades_range(升序)/list_active_customers;种子名单 3 条改名维持唯一演示命中;core_cash_flow 上下文一期未接(PRD「如有」),B8 前评估;真实名单数据接入时 payload/审计出口接 utils/desensitize(B9b 前核查项)** | | 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` 已留痕)** | 可运行应用 | `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 行。集成测试:A-1~A-5、A-7(状态机/compliance 403/GET 强制 aml)、A-9 越权、**trace 一致性断言** | 测试套件 + fixture | `pytest` 全绿 | B7 | +| 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 收敛)**。集成测试:A-1~A-5、A-7(状态机/compliance 403/GET 强制 aml)、A-9 越权、**trace 一致性断言** | 测试套件 + 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)** | 演示 SOP | 按 PRD §8 验收表逐条打勾 | B8、B9a | diff --git a/docs/项目框架设计/架构设计-风控模块.md b/docs/项目框架设计/架构设计-风控模块.md index 186278f..dc90acc 100644 --- a/docs/项目框架设计/架构设计-风控模块.md +++ b/docs/项目框架设计/架构设计-风控模块.md @@ -42,8 +42,10 @@ app/ ├── repository/ │ ├── core_ro.py # 【扩展(最小化)】SUIT 校验复用现有 get_customer_l0 │ │ # (已含 age+risk_code+evaluated_at)与 get_product; -│ │ # 当日流水明细复用 list_trades(since=当日0点); -│ │ # 仅新增 sum_trades_on_date(当日累计聚合) +│ │ # 仅新增(最小化,均纯 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 │ @@ -72,7 +74,7 @@ scripts/ tests/ ├── test_suitability.py # SUIT 矩阵全组合 + 70 岁/过期/NULL 年龄边界 ├── test_risk_rules.py # RISK-001~005 纯函数用例 -├── test_aml_match.py # 归一化/相似度边界 +├── test_aml_service.py # 归一化/相似度边界(B4 实际命名,评审 P3-10) └── test_gateway_flow.py # 集成:A-1~A-5 阻断与放行链路 ``` diff --git a/scripts/agent/seed-aml-list.sql b/scripts/agent/seed-aml-list.sql index 15bb0d5..b3644a2 100644 --- a/scripts/agent/seed-aml-list.sql +++ b/scripts/agent/seed-aml-list.sql @@ -5,6 +5,8 @@ -- 说明:full_name 与 core_customer.display_name 同为脱敏展示名口径; -- AML-0001 故意与种子客户 CUST-1002 的展示名「客户·赵**」一致(演示命中,PRD A-5 用例; -- CUST-1002 测评已被 scripts/demo/prepare_risk_demo.sql 刷新,交易可成功走通 AML 链路)。 +-- 其余名单名(名单样例·壹/贰/叁)与全部种子客户展示名相似度 <0.85,不构成命中(B4 评审 P2-2, +-- 维持 PRD §10.1「故意包含 1 条同名记录」口径)。 -- ============================================================================= USE jinrong_agent; @@ -14,11 +16,11 @@ INSERT INTO risk_aml_list ('AML-0001', 'sanction', '客户·赵**', NULL, NULL, 0.85, '外部名单镜像(模拟)', 'V2026.09', '2026-01-01'), ('AML-0002', 'terror', '客户·测试命中**', NULL, NULL, 0.85, '外部名单镜像(模拟)', 'V2026.09', '2026-01-01'), ('AML-0003', 'sanction', '客户·李**', NULL, NULL, 0.85, '外部名单镜像(模拟)', 'V2026.09', '2026-01-01'), -('AML-0004', 'pep', '客户·王**', NULL, NULL, 0.85, '外部名单镜像(模拟)', 'V2026.09', '2026-01-01'), +('AML-0004', 'pep', '名单样例·壹', NULL, NULL, 0.85, '外部名单镜像(模拟)', 'V2026.09', '2026-01-01'), ('AML-0005', 'sanction', '张某某', NULL, NULL, 0.85, '外部名单镜像(模拟)', 'V2026.09', '2026-01-01'), -('AML-0006', 'terror', '客户·陈**', NULL, NULL, 0.85, '外部名单镜像(模拟)', 'V2026.09', '2026-01-01'), +('AML-0006', 'terror', '名单样例·贰', NULL, NULL, 0.85, '外部名单镜像(模拟)', 'V2026.09', '2026-01-01'), ('AML-0007', 'pep', '客户·刘**', NULL, NULL, 0.85, '外部名单镜像(模拟)', 'V2026.09', '2026-01-01'), -('AML-0008', 'sanction', '客户·周**', NULL, NULL, 0.85, '外部名单镜像(模拟)', 'V2026.09', '2026-01-01'); +('AML-0008', 'sanction', '名单样例·叁', NULL, NULL, 0.85, '外部名单镜像(模拟)', 'V2026.09', '2026-01-01'); -- 校验 SELECT list_id, list_type, full_name, is_active FROM risk_aml_list; diff --git a/tests/test_aml_service.py b/tests/test_aml_service.py index 396f623..e222613 100644 --- a/tests/test_aml_service.py +++ b/tests/test_aml_service.py @@ -217,6 +217,30 @@ def test_scan_all_creates_alert_and_l3(env): assert len(pub.messages) == 2 # 每个命中客户一次紧急推送 +def test_multi_entry_match_merges_single_alert(env): + """B4 评审 P2-4:一客户命中多条名单 → 仅一张 aml 单,matches 合并全量进 payload。""" + core, repo, pub = env + with core._engine.begin() as conn: + conn.execute( + text( + "INSERT INTO risk_aml_list (list_id, list_type, full_name, match_threshold," + " source, list_version, is_active) VALUES" + " ('DUP-1', 'sanction', '赵六六', 0.85, 'mock', 'v1', 1)," + " ('DUP-2', 'pep', '赵六六', 0.85, 'mock', 'v1', 1)" + ) + ) + summary = scan_all(core_ro=core, risk_repo=repo) + assert summary["hit_customers"] == 3 # C1、C2(原种子命中)+ C3(双名单) + c3_ids = [ + aid for aid in summary["alerts"] if repo.get_alert(aid)["customer_id"] == "C3" + ] + assert len(c3_ids) == 1 # 单事件单张 + matches = repo.get_alert(c3_ids[0])["payload"]["events"][0]["matches"] + assert len(matches) == 2 + assert {m["list_type"] for m in matches} == {"sanction", "pep"} + assert len(pub.messages) == 3 # 每命中客户一次推送 + + def test_scan_all_no_hit_creates_nothing(env): core, repo, pub = env with core._engine.begin() as conn: diff --git a/tests/test_risk_engine.py b/tests/test_risk_engine.py index 0182cec..132e0c7 100644 --- a/tests/test_risk_engine.py +++ b/tests/test_risk_engine.py @@ -236,6 +236,33 @@ def test_aml_and_event_rule_both_fire(env): assert len(pub.messages) == 2 (_, aml_payload), = [m for m in pub.messages if m[1]["alert_type"] == "aml"] assert "compliance" in aml_payload["notify_role"] + # FR-4 payload 完整性(B4 评审 P1-1):客户上下文 = L0 摘要 + 近 30 天统计 + ev_alert = next( + a for a in (repo.get_alert(aid) for aid in result["alert_ids"]) + if a["alert_type"] == "large_amount" + ) + ctx = ev_alert["payload"]["customer_context"] + assert ctx["l0"]["display_name"] == "李四" and ctx["l0"]["age"] == 35 + assert ctx["trades_30d"]["count"] >= 1 + + +def test_second_trade_same_day_appends_to_same_alert(env): + """B4 评审 P2-3:同日第二笔 → 追加同一事件单,triggered_rules/risk_score/L3 联动。""" + core, repo, _, engine = env + with engine.begin() as conn: + _seed_trade(conn, "T1", "600000") + _seed_trade(conn, "T2", "200000", at=datetime(2026, 9, 6, 14, 1, 0)) + r1 = process_trade_event(_trade("T1", "600000"), core_ro=core, risk_repo=repo) + r2 = process_trade_event( + _trade("T2", "200000", at=datetime(2026, 9, 6, 14, 1, 0)), core_ro=core, risk_repo=repo + ) + assert r2["alert_ids"] == [r1["alert_ids"][0]] # 同日仅一张事件类单 + alert = repo.get_alert(r1["alert_ids"][0]) + assert len(alert["payload"]["events"]) == 2 + assert alert["triggered_rules"] == ["RISK-001", "RISK-002"] # T2 仅命中累计,追加合并 + assert alert["risk_score"] == 70 # 取 max + assert _counts(engine, "audit_log", "decision='alert_appended'") == 1 + assert repo.get_l3("C1")["last_alert_id"] == alert["alert_id"] # L3 联动最新 def test_probe_window_uses_trade_time_not_wall_clock(env):