refactor: L3 risk_score 一期不写保持 NULL,归 R-05 首写(B3 评审 P3-4 用户拍板)
This commit is contained in:
@@ -1,10 +1,13 @@
|
||||
"""L3 监测画像写入(B3 · PRD FR-7 / R-05 最小写入)。
|
||||
|
||||
合并规则(防降级,PRD FR-7):monitor_tier 取最高档(normal < watch < high),
|
||||
monitor_tags 追加合并不覆盖,risk_score 取 max(FR-7 未定义的口径外推,见
|
||||
RISK-00n 静态分;R-05 动态评分接入时复核),last_alert_id 传入时联动最新预警
|
||||
(未传保留旧值,防调用方漏传抹掉),computed_at 每次写当前时间(列 NOT NULL,
|
||||
毫秒截断对齐 DATETIME(3),保证乐观锁读写比对一致)。
|
||||
monitor_tags 追加合并不覆盖,last_alert_id 传入时联动最新预警(未传保留旧值,
|
||||
防调用方漏传抹掉),computed_at 每次写当前时间(列 NOT NULL,毫秒截断对齐
|
||||
DATETIME(3),保证乐观锁读写比对一致)。
|
||||
|
||||
risk_score 口径(评审 P3-4 · 用户拍板 2026-09-06):一期**不写**(保持 NULL)。
|
||||
该列是 R-05 动态评分模型的客户风险分,与 risk_alert.risk_score(预警单严重度)
|
||||
语义不同;由 R-05 首写,避免静态分造成 tier/score 错位与二期语义污染。
|
||||
|
||||
并发(B3 评审 P1-1 修复):进程内锁按 customer_id 串行减少冲突(多进程部署换
|
||||
Redis SET NX,接口不变);锁超时降级与跨进程竞态由乐观锁兜底——update 比对
|
||||
@@ -59,7 +62,6 @@ def merge_l3(
|
||||
existing: dict[str, Any] | None,
|
||||
*,
|
||||
mapped_tier: str,
|
||||
risk_score: int | None,
|
||||
monitor_tags: list[str],
|
||||
last_alert_id: str | None,
|
||||
score_dimensions: dict[str, Any] | None = None,
|
||||
@@ -67,21 +69,19 @@ def merge_l3(
|
||||
) -> dict[str, Any]:
|
||||
"""纯函数:existing 行(repo.get_l3 输出)与新事件合并后的 L3 行。
|
||||
|
||||
existing 为 None 表示新客户首写;risk_score 双方均空时保持 NULL。
|
||||
existing 为 None 表示新客户首写;risk_score 一期恒 None(不读不写,归 R-05)。
|
||||
"""
|
||||
if existing is None:
|
||||
tier, old_score, old_tags, old_dims, old_alert = "normal", None, [], {}, None
|
||||
tier, old_tags, old_dims, old_alert = "normal", [], {}, None
|
||||
else:
|
||||
tier = existing.get("monitor_tier") or "normal"
|
||||
old_score = existing.get("risk_score")
|
||||
old_tags = existing.get("monitor_tags") or []
|
||||
old_dims = existing.get("score_dimensions") or {}
|
||||
old_alert = existing.get("last_alert_id")
|
||||
|
||||
scores = [int(s) for s in (old_score, risk_score) if s is not None]
|
||||
return {
|
||||
"monitor_tier": highest_tier(tier, mapped_tier),
|
||||
"risk_score": max(scores) if scores else None,
|
||||
"risk_score": None,
|
||||
"monitor_tags": sorted(set(old_tags) | set(monitor_tags)),
|
||||
"score_dimensions": score_dimensions if score_dimensions is not None else old_dims,
|
||||
"last_alert_id": last_alert_id if last_alert_id is not None else old_alert,
|
||||
@@ -92,7 +92,6 @@ def merge_l3(
|
||||
def upsert_profile_l3(
|
||||
customer_id: str,
|
||||
alert_type: str,
|
||||
risk_score: int | None = None,
|
||||
monitor_tags: list[str] | None = None,
|
||||
last_alert_id: str | None = None,
|
||||
score_dimensions: dict[str, Any] | None = None,
|
||||
@@ -116,7 +115,6 @@ def upsert_profile_l3(
|
||||
merged = merge_l3(
|
||||
existing,
|
||||
mapped_tier=mapped_tier,
|
||||
risk_score=risk_score,
|
||||
monitor_tags=tags,
|
||||
last_alert_id=last_alert_id,
|
||||
score_dimensions=score_dimensions,
|
||||
|
||||
@@ -22,7 +22,7 @@
|
||||
| --- | --- | --- | --- | --- |
|
||||
| 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③) | L3 写入 | 单测:normal→high 不降级、AML 后大额不回落 | 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 |
|
||||
| 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) |
|
||||
|
||||
+45
-47
@@ -1,7 +1,8 @@
|
||||
"""profile_l3 单测(B3 · 最高档合并防降级 / AML 后大额不回落 / 并发首单)。
|
||||
|
||||
sqlite StaticPool 单连接共享内存库(同 test_alert_service 模式);
|
||||
IntegrityError 兜底用 monkeypatch 模拟跨进程首写竞态。
|
||||
IntegrityError / 乐观锁丢竞态用 monkeypatch 模拟跨进程交错。
|
||||
risk_score 口径:一期不写(恒 NULL,归 R-05 评分模型首写,评审 P3-4 用户拍板)。
|
||||
"""
|
||||
|
||||
from datetime import datetime
|
||||
@@ -53,14 +54,14 @@ def env():
|
||||
engine.dispose()
|
||||
|
||||
|
||||
def _upsert(repo, cid, alert_type, score=None, alert_id=None, tags=None, dims=None):
|
||||
def _upsert(repo, cid, alert_type, alert_id=None, tags=None, dims=None, computed_at=None):
|
||||
return upsert_profile_l3(
|
||||
cid,
|
||||
alert_type,
|
||||
risk_score=score,
|
||||
monitor_tags=tags,
|
||||
last_alert_id=alert_id,
|
||||
score_dimensions=dims,
|
||||
computed_at=computed_at,
|
||||
risk_repo=repo,
|
||||
)
|
||||
|
||||
@@ -87,7 +88,7 @@ def test_alert_type_tier_mapping(alert_type, tier):
|
||||
def test_unknown_alert_type_rejected(env):
|
||||
repo, _ = env
|
||||
with pytest.raises(ValueError):
|
||||
_upsert(repo, "C1", "unknown_type", 70)
|
||||
_upsert(repo, "C1", "unknown_type")
|
||||
|
||||
|
||||
def test_highest_tier_order():
|
||||
@@ -98,20 +99,19 @@ def test_highest_tier_order():
|
||||
|
||||
|
||||
def test_merge_new_customer():
|
||||
merged = merge_l3(None, mapped_tier="watch", risk_score=70, monitor_tags=["t1"],
|
||||
last_alert_id="ALT-1")
|
||||
merged = merge_l3(None, mapped_tier="watch", monitor_tags=["t1"], last_alert_id="ALT-1")
|
||||
assert merged["monitor_tier"] == "watch"
|
||||
assert merged["risk_score"] == 70
|
||||
assert merged["risk_score"] is None # 一期不写(P3-4 口径)
|
||||
assert merged["monitor_tags"] == ["t1"]
|
||||
assert merged["last_alert_id"] == "ALT-1"
|
||||
assert merged["computed_at"] is not None
|
||||
|
||||
|
||||
def test_merge_both_scores_none_stays_none():
|
||||
existing = {"monitor_tier": "watch", "risk_score": None, "monitor_tags": [],
|
||||
def test_merge_risk_score_stays_null_even_if_existing_has_value():
|
||||
"""历史行若已有 score(异常数据),合并时也不维护/不传播。"""
|
||||
existing = {"monitor_tier": "watch", "risk_score": 90, "monitor_tags": [],
|
||||
"score_dimensions": {}}
|
||||
merged = merge_l3(existing, mapped_tier="watch", risk_score=None, monitor_tags=[],
|
||||
last_alert_id="ALT-2")
|
||||
merged = merge_l3(existing, mapped_tier="watch", monitor_tags=[], last_alert_id="ALT-2")
|
||||
assert merged["risk_score"] is None
|
||||
|
||||
|
||||
@@ -120,33 +120,33 @@ def test_merge_both_scores_none_stays_none():
|
||||
|
||||
def test_first_event_inserts(env):
|
||||
repo, engine = env
|
||||
merged = _upsert(repo, "C1", "large_amount", 70, alert_id="ALT-1")
|
||||
merged = _upsert(repo, "C1", "large_amount", alert_id="ALT-1")
|
||||
assert merged["monitor_tier"] == "watch" and merged["customer_id"] == "C1"
|
||||
row = _row(engine, "C1")
|
||||
assert row["monitor_tier"] == "watch" and row["risk_score"] == 70
|
||||
assert row["monitor_tier"] == "watch" and row["risk_score"] is None
|
||||
assert row["last_alert_id"] == "ALT-1" and row["computed_at"] is not None
|
||||
assert get_profile_l3("C1", risk_repo=repo)["monitor_tier"] == "watch"
|
||||
|
||||
|
||||
def test_normal_upgrades_to_watch(env):
|
||||
repo, _ = env
|
||||
_upsert(repo, "C1", "suitability", 90, alert_id="ALT-0")
|
||||
merged = _upsert(repo, "C1", "pattern", 80, alert_id="ALT-1")
|
||||
_upsert(repo, "C1", "suitability", alert_id="ALT-0")
|
||||
merged = _upsert(repo, "C1", "pattern", alert_id="ALT-1")
|
||||
assert merged["monitor_tier"] == "watch" # normal → watch 升档
|
||||
|
||||
|
||||
def test_watch_does_not_degrade_to_normal(env):
|
||||
"""suitability 映射 normal:已 watch 的客户不被 suitability 事件拉低。"""
|
||||
repo, engine = env
|
||||
_upsert(repo, "C1", "pattern", 80, alert_id="ALT-1")
|
||||
merged = _upsert(repo, "C1", "suitability", 90, alert_id="ALT-2")
|
||||
_upsert(repo, "C1", "pattern", alert_id="ALT-1")
|
||||
merged = _upsert(repo, "C1", "suitability", alert_id="ALT-2")
|
||||
assert merged["monitor_tier"] == "watch"
|
||||
assert _row(engine, "C1")["monitor_tier"] == "watch"
|
||||
|
||||
|
||||
def test_aml_marks_high_with_pending_review_tag(env):
|
||||
repo, engine = env
|
||||
merged = _upsert(repo, "C1", "aml", 95, alert_id="ALT-1")
|
||||
merged = _upsert(repo, "C1", "aml", alert_id="ALT-1")
|
||||
assert merged["monitor_tier"] == "high"
|
||||
assert AML_PENDING_TAG in merged["monitor_tags"]
|
||||
assert AML_PENDING_TAG in _row(engine, "C1")["monitor_tags"]
|
||||
@@ -155,37 +155,37 @@ def test_aml_marks_high_with_pending_review_tag(env):
|
||||
def test_high_does_not_degrade_to_normal(env):
|
||||
"""B3 验收(评审 P2-4):aml 后 suitability(映射 normal)不回落。"""
|
||||
repo, engine = env
|
||||
_upsert(repo, "C1", "aml", 95, alert_id="ALT-1")
|
||||
merged = _upsert(repo, "C1", "suitability", 90, alert_id="ALT-2")
|
||||
_upsert(repo, "C1", "aml", alert_id="ALT-1")
|
||||
merged = _upsert(repo, "C1", "suitability", alert_id="ALT-2")
|
||||
assert merged["monitor_tier"] == "high"
|
||||
assert _row(engine, "C1")["monitor_tier"] == "high"
|
||||
|
||||
|
||||
def test_large_amount_after_aml_does_not_fall_back(env):
|
||||
"""B3 验收:AML 后大额 → tier 仍 high、tags 并集、score 取 max。"""
|
||||
"""B3 验收:AML 后大额 → tier 仍 high、tags 并集、risk_score 保持 NULL。"""
|
||||
repo, engine = env
|
||||
_upsert(repo, "C1", "aml", 95, alert_id="ALT-1")
|
||||
merged = _upsert(repo, "C1", "large_amount", 70, alert_id="ALT-2", tags=["manual_review"])
|
||||
_upsert(repo, "C1", "aml", alert_id="ALT-1")
|
||||
merged = _upsert(repo, "C1", "large_amount", alert_id="ALT-2", tags=["manual_review"])
|
||||
assert merged["monitor_tier"] == "high"
|
||||
assert merged["risk_score"] == 95 # 取 max,不降
|
||||
assert merged["risk_score"] is None
|
||||
assert set(merged["monitor_tags"]) == {AML_PENDING_TAG, "manual_review"}
|
||||
assert merged["last_alert_id"] == "ALT-2" # 联动最新
|
||||
row = _row(engine, "C1")
|
||||
assert row["monitor_tier"] == "high" and row["risk_score"] == 95
|
||||
assert row["monitor_tier"] == "high" and row["risk_score"] is None
|
||||
|
||||
|
||||
def test_tags_accumulate_not_overwrite(env):
|
||||
repo, _ = env
|
||||
_upsert(repo, "C1", "pattern", 80, alert_id="ALT-1", tags=["freq"])
|
||||
merged = _upsert(repo, "C1", "pattern", 80, alert_id="ALT-2", tags=["manual_review"])
|
||||
_upsert(repo, "C1", "pattern", alert_id="ALT-1", tags=["freq"])
|
||||
merged = _upsert(repo, "C1", "pattern", alert_id="ALT-2", tags=["manual_review"])
|
||||
assert merged["monitor_tags"] == ["freq", "manual_review"] # 不同 tag 跨事件追加
|
||||
|
||||
|
||||
def test_last_alert_id_kept_when_not_passed(env):
|
||||
"""评审 P2-2:漏传 last_alert_id 不抹掉旧值(FR-7 联动语义防御)。"""
|
||||
repo, engine = env
|
||||
_upsert(repo, "C1", "aml", 95, alert_id="ALT-1")
|
||||
merged = upsert_profile_l3("C1", "large_amount", 70, risk_repo=repo)
|
||||
_upsert(repo, "C1", "aml", alert_id="ALT-1")
|
||||
merged = upsert_profile_l3("C1", "large_amount", risk_repo=repo)
|
||||
assert merged["last_alert_id"] == "ALT-1"
|
||||
assert _row(engine, "C1")["last_alert_id"] == "ALT-1"
|
||||
|
||||
@@ -193,8 +193,8 @@ def test_last_alert_id_kept_when_not_passed(env):
|
||||
def test_computed_at_refreshed_on_each_write(env):
|
||||
"""评审 P3-2:每次写 computed_at 均刷新(FR-7)。"""
|
||||
repo, engine = env
|
||||
_upsert(repo, "C1", "pattern", 80, alert_id="ALT-1")
|
||||
upsert_profile_l3("C1", "large_amount", 70, last_alert_id="ALT-2",
|
||||
_upsert(repo, "C1", "pattern", alert_id="ALT-1")
|
||||
upsert_profile_l3("C1", "large_amount", last_alert_id="ALT-2",
|
||||
computed_at=datetime(2027, 1, 1, 8, 0, 0), risk_repo=repo)
|
||||
# sqlite 读回为字符串,格式无关断言(核心是值已从首写的 now 刷新为传入时间)
|
||||
assert str(_row(engine, "C1")["computed_at"]).startswith("2027-01-01 08:00")
|
||||
@@ -202,10 +202,10 @@ def test_computed_at_refreshed_on_each_write(env):
|
||||
|
||||
def test_score_dimensions_replaced_only_when_passed(env):
|
||||
repo, _ = env
|
||||
_upsert(repo, "C1", "aml", 95, alert_id="ALT-1", dims={"amount": 1})
|
||||
merged = _upsert(repo, "C1", "large_amount", 70, alert_id="ALT-2")
|
||||
_upsert(repo, "C1", "aml", alert_id="ALT-1", dims={"amount": 1})
|
||||
merged = _upsert(repo, "C1", "large_amount", alert_id="ALT-2")
|
||||
assert merged["score_dimensions"] == {"amount": 1} # 未传保留旧值
|
||||
merged = _upsert(repo, "C1", "large_amount", 70, alert_id="ALT-3", dims={"amount": 2})
|
||||
merged = _upsert(repo, "C1", "large_amount", alert_id="ALT-3", dims={"amount": 2})
|
||||
assert merged["score_dimensions"] == {"amount": 2} # 传入则替换
|
||||
|
||||
|
||||
@@ -222,7 +222,7 @@ def test_concurrent_first_upsert_single_row(env):
|
||||
|
||||
def worker(alert_type, alert_id):
|
||||
try:
|
||||
_upsert(repo, "C1", alert_type, 70, alert_id=alert_id)
|
||||
_upsert(repo, "C1", alert_type, alert_id=alert_id)
|
||||
except Exception as exc: # pragma: no cover
|
||||
errors.append(exc)
|
||||
|
||||
@@ -250,10 +250,9 @@ def test_concurrent_first_upsert_single_row(env):
|
||||
def test_integrity_error_falls_back_to_remerge(env, monkeypatch):
|
||||
"""跨进程竞态(评审 P2-7③):首读 None、insert 撞主键 → 重读合并转更新,不丢对方写入。"""
|
||||
repo, engine = env
|
||||
from sqlalchemy.exc import IntegrityError
|
||||
|
||||
# 对方进程已写入 high(模拟 AML 先落库)
|
||||
repo.insert_l3("C1", "high", 95, {}, [AML_PENDING_TAG], "ALT-AML", datetime.now())
|
||||
repo.insert_l3("C1", "high", None, {}, [AML_PENDING_TAG], "ALT-AML", datetime.now())
|
||||
|
||||
calls = {"n": 0}
|
||||
orig_get = type(repo).get_l3
|
||||
@@ -266,12 +265,11 @@ def test_integrity_error_falls_back_to_remerge(env, monkeypatch):
|
||||
|
||||
monkeypatch.setattr(type(repo), "get_l3", racing_get)
|
||||
try:
|
||||
merged = _upsert(repo, "C1", "large_amount", 70, alert_id="ALT-2")
|
||||
merged = _upsert(repo, "C1", "large_amount", alert_id="ALT-2")
|
||||
finally:
|
||||
monkeypatch.undo()
|
||||
|
||||
assert merged["monitor_tier"] == "high" # 重读后合并,不降级
|
||||
assert merged["risk_score"] == 95
|
||||
assert merged["last_alert_id"] == "ALT-2"
|
||||
row = _row(engine, "C1")
|
||||
assert row["monitor_tier"] == "high" and row["last_alert_id"] == "ALT-2"
|
||||
@@ -279,9 +277,9 @@ def test_integrity_error_falls_back_to_remerge(env, monkeypatch):
|
||||
|
||||
def test_update_lost_race_retries_and_converges(env, monkeypatch):
|
||||
"""评审 P1-1:update 路径丢更新——本进程读旧值后对方先写 high,乐观锁未命中
|
||||
触发重读重试,最终收敛 high/95(修复前会被覆盖回退 watch/70)。"""
|
||||
触发重读重试,最终收敛 high(修复前会被覆盖回退 watch)。"""
|
||||
repo, engine = env
|
||||
repo.insert_l3("C1", "watch", 70, {}, [], "ALT-1", datetime.now())
|
||||
repo.insert_l3("C1", "watch", None, {}, [], "ALT-1", datetime.now())
|
||||
|
||||
calls = {"n": 0}
|
||||
orig_update = type(repo).update_l3
|
||||
@@ -289,27 +287,27 @@ def test_update_lost_race_retries_and_converges(env, monkeypatch):
|
||||
def racing_update(self, *args, **kwargs):
|
||||
calls["n"] += 1
|
||||
if calls["n"] == 1:
|
||||
orig_update(self, "C1", "high", 95, {}, [AML_PENDING_TAG], "ALT-AML", datetime.now())
|
||||
orig_update(self, "C1", "high", None, {}, [AML_PENDING_TAG], "ALT-AML", datetime.now())
|
||||
return False # 对方抢先提交,本进程乐观锁未命中
|
||||
return orig_update(self, *args, **kwargs)
|
||||
|
||||
monkeypatch.setattr(type(repo), "update_l3", racing_update)
|
||||
merged = _upsert(repo, "C1", "large_amount", 70, alert_id="ALT-2")
|
||||
merged = _upsert(repo, "C1", "large_amount", alert_id="ALT-2")
|
||||
|
||||
assert calls["n"] >= 2 # 确实走了重试
|
||||
assert merged["monitor_tier"] == "high" and merged["risk_score"] == 95
|
||||
assert merged["monitor_tier"] == "high"
|
||||
assert AML_PENDING_TAG in merged["monitor_tags"]
|
||||
row = _row(engine, "C1")
|
||||
assert row["monitor_tier"] == "high" and row["risk_score"] == 95
|
||||
assert row["monitor_tier"] == "high"
|
||||
assert row["last_alert_id"] == "ALT-2"
|
||||
|
||||
|
||||
def test_persistent_conflict_raises_not_silent(env, monkeypatch):
|
||||
"""重试耗尽抛错(宁失败不静默覆盖),不留半写状态。"""
|
||||
repo, engine = env
|
||||
repo.insert_l3("C1", "watch", 70, {}, [], "ALT-1", datetime.now())
|
||||
repo.insert_l3("C1", "watch", None, {}, [], "ALT-1", datetime.now())
|
||||
monkeypatch.setattr(type(repo), "update_l3", lambda self, *a, **k: False)
|
||||
with pytest.raises(RuntimeError, match="conflicted"):
|
||||
_upsert(repo, "C1", "large_amount", 70, alert_id="ALT-2")
|
||||
_upsert(repo, "C1", "large_amount", alert_id="ALT-2")
|
||||
row = _row(engine, "C1")
|
||||
assert row["monitor_tier"] == "watch" # 未被静默改写
|
||||
|
||||
Reference in New Issue
Block a user