2026-09-06 17:19:54 +08:00
|
|
|
|
"""风控事件引擎(B4 · 架构 §3.1 ④ / PRD FR-3)。
|
|
|
|
|
|
|
|
|
|
|
|
编排(网关 B5 在 core_trade 落库后**同步调用**,不用消息队列):
|
|
|
|
|
|
当日流水上下文(core_ro)→ RISK-001~005 纯函数 → AML 姓名匹配 → 预警落库
|
|
|
|
|
|
(alert_service:聚合/去重/审计/推送)→ L3 upsert(profile_l3)。
|
|
|
|
|
|
审计 pass(未命中)与命中审计均由 alert_service 完成,本层不重复落审计。
|
|
|
|
|
|
|
|
|
|
|
|
RISK-004 窗口以 trade["traded_at"] 为事件时点(非墙钟 now):rebuild_alerts
|
|
|
|
|
|
幂等重放可复现窗口判定,演示脚本不受执行时刻影响。
|
|
|
|
|
|
|
|
|
|
|
|
客户事件钩子 on_customer_created/on_customer_updated 为 FR-5 预留(本期 no-op,
|
|
|
|
|
|
模拟环境无开户流程)。
|
|
|
|
|
|
"""
|
|
|
|
|
|
|
|
|
|
|
|
from __future__ import annotations
|
|
|
|
|
|
|
2026-09-06 17:40:13 +08:00
|
|
|
|
from datetime import datetime, timedelta
|
|
|
|
|
|
from decimal import Decimal
|
2026-09-06 17:19:54 +08:00
|
|
|
|
from typing import Any
|
|
|
|
|
|
|
|
|
|
|
|
from app.repository.core_ro import CoreReadOnlyRepository
|
|
|
|
|
|
from app.repository.risk_repository import RiskRepository
|
|
|
|
|
|
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
|
2026-09-06 17:40:13 +08:00
|
|
|
|
from app.utils.trace import ensure_trace
|
2026-09-06 17:19:54 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _as_datetime(value: Any) -> datetime:
|
|
|
|
|
|
if isinstance(value, datetime):
|
|
|
|
|
|
return value
|
|
|
|
|
|
if isinstance(value, str):
|
|
|
|
|
|
return datetime.fromisoformat(value)
|
|
|
|
|
|
raise TypeError(f"traded_at must be datetime/str, got {type(value)!r}")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _normalize_trades(trades: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
|
|
|
|
|
"""驱动差异防御:sqlite text 查询返回 str 时间,统一转 datetime(MySQL 驱动本就返回 datetime)。"""
|
|
|
|
|
|
for t in trades:
|
|
|
|
|
|
if isinstance(t.get("traded_at"), str):
|
|
|
|
|
|
t["traded_at"] = datetime.fromisoformat(t["traded_at"])
|
|
|
|
|
|
return trades
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-06 17:40:13 +08:00
|
|
|
|
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)},
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-06 17:19:54 +08:00
|
|
|
|
def process_trade_event(
|
|
|
|
|
|
trade: dict[str, Any],
|
|
|
|
|
|
core_ro: CoreReadOnlyRepository | None = None,
|
|
|
|
|
|
risk_repo: RiskRepository | None = None,
|
|
|
|
|
|
thresholds: RiskThresholds | None = None,
|
|
|
|
|
|
) -> dict[str, Any]:
|
|
|
|
|
|
"""处理一笔已落库交易(PRD FR-1 ②③b 之后)。
|
|
|
|
|
|
|
|
|
|
|
|
返回 {"triggered_rules": [...], "alert_ids": [...], "aml_hit": bool},
|
|
|
|
|
|
网关据此拼装响应(FR-1 ⑤:blocked=false + trade_id + 触发规则列表)。
|
|
|
|
|
|
"""
|
|
|
|
|
|
core = core_ro or CoreReadOnlyRepository()
|
|
|
|
|
|
repo = risk_repo or RiskRepository()
|
|
|
|
|
|
th = thresholds or RiskThresholds.from_settings()
|
2026-09-06 17:40:13 +08:00
|
|
|
|
ensure_trace() # 脚本/重放入口兜底归因(中间件场景保留现有 trace,评审 P3-8)
|
2026-09-06 17:19:54 +08:00
|
|
|
|
|
|
|
|
|
|
event_at = _as_datetime(trade["traded_at"])
|
|
|
|
|
|
day_start = event_at.replace(hour=0, minute=0, second=0, microsecond=0)
|
2026-09-06 17:40:13 +08:00
|
|
|
|
day_end = day_start + timedelta(days=1)
|
|
|
|
|
|
trades = _normalize_trades(
|
|
|
|
|
|
core.list_trades_range(trade["customer_id"], day_start, day_end)
|
|
|
|
|
|
)
|
2026-09-06 17:19:54 +08:00
|
|
|
|
|
|
|
|
|
|
result: dict[str, Any] = {"triggered_rules": [], "alert_ids": [], "aml_hit": False}
|
|
|
|
|
|
|
|
|
|
|
|
hits = run_rules(trades, th, now=event_at)
|
|
|
|
|
|
if hits:
|
2026-09-06 17:40:13 +08:00
|
|
|
|
alert = record_trade_alerts(
|
|
|
|
|
|
trade,
|
|
|
|
|
|
hits,
|
|
|
|
|
|
risk_repo=repo,
|
|
|
|
|
|
customer_context=_build_customer_context(core, trade["customer_id"], day_start),
|
|
|
|
|
|
)
|
2026-09-06 17:19:54 +08:00
|
|
|
|
result["triggered_rules"] = sorted({h.rule_id for h in hits})
|
|
|
|
|
|
if alert:
|
|
|
|
|
|
result["alert_ids"].append(alert["alert_id"])
|
|
|
|
|
|
best = max(hits, key=lambda h: h.risk_score)
|
|
|
|
|
|
upsert_profile_l3(
|
|
|
|
|
|
trade["customer_id"],
|
|
|
|
|
|
best.alert_type,
|
|
|
|
|
|
last_alert_id=alert["alert_id"] if alert else None,
|
|
|
|
|
|
risk_repo=repo,
|
|
|
|
|
|
)
|
2026-09-06 17:40:13 +08:00
|
|
|
|
else:
|
|
|
|
|
|
# 未命中分支:pass 审计由 alert_service 统一落库(架构 §3.1 ④)
|
|
|
|
|
|
record_trade_alerts(trade, [], risk_repo=repo)
|
2026-09-06 17:19:54 +08:00
|
|
|
|
|
|
|
|
|
|
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"],
|
2026-09-06 17:40:13 +08:00
|
|
|
|
{
|
|
|
|
|
|
"trigger": "trade",
|
|
|
|
|
|
"trade_id": trade.get("trade_id"),
|
|
|
|
|
|
"product_id": trade.get("product_id"), # 评审 P3-6:payload/审计透传
|
|
|
|
|
|
"matches": aml_hits,
|
|
|
|
|
|
},
|
2026-09-06 17:19:54 +08:00
|
|
|
|
risk_repo=repo,
|
|
|
|
|
|
)
|
|
|
|
|
|
result["alert_ids"].append(alert["alert_id"])
|
|
|
|
|
|
upsert_profile_l3(
|
|
|
|
|
|
trade["customer_id"], "aml", last_alert_id=alert["alert_id"], risk_repo=repo
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
return result
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def on_customer_created(customer_id: str) -> None:
|
|
|
|
|
|
"""AML 开户触发预留(本期 no-op;模拟环境无开户流程,PRD FR-5)。"""
|
|
|
|
|
|
|
|
|
|
|
|
def on_customer_updated(customer_id: str) -> None:
|
|
|
|
|
|
"""客户信息变更触发预留(本期 no-op;PRD FR-5)。"""
|