fix: FR-4 payload 上下文补组装/种子名单唯一命中/规则流水升序等(B4 评审 P1-1、P2-2~5、P3×7)
This commit is contained in:
@@ -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)。
|
||||
|
||||
|
||||
@@ -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] = []
|
||||
|
||||
@@ -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"])
|
||||
|
||||
@@ -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()
|
||||
|
||||
+1
-1
@@ -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
|
||||
|
||||
|
||||
@@ -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 |
|
||||
|
||||
|
||||
@@ -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 阻断与放行链路
|
||||
```
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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):
|
||||
|
||||
Reference in New Issue
Block a user