依据《实现方案-风控追加需求v1.1-C4C6.md》§2;不改表结构(alert_type/status 复用 payload 承载,audit_log.event_type 为 VARCHAR 可直接扩)。 1. settings.py + .env.example:一次性加齐风控追加 v1.1 共 11 项配置(C4~C6 共用)。 2. core_ro.concentration_profile(customer_id, limit=500):一次 SQL 取明细 (LIMIT limit+1 探测截断)+ Python 端按 min_risk_code in (R4,R5) 聚合; 收口挂账 #1(PRD 字面为 list_holdings,改聚合封装,docstring 注明偏离)。 3. rules.py:RULE_SCORES/RULE_ALERT_TYPES 加 RISK-006=60/pattern;RuleHit 加 alert_subtype;RiskThresholds 加 concentration_threshold 且 from_settings 必须补读(评审 P1-2:漏读会让 conftest monkeypatch 失效打穿现有断言); 新增纯函数 rule_concentration——空仓不触发、截断视同达标(保守告警)、 阈值边界 79.9% 不触发 / 80% 触发、R4+R5 为 0 不触发。 4. engine.process_trade_event:run_rules 之后、record_trade_alerts 之前并入 集中度命中(不动 run_rules 签名);命中后 L3 打 high_risk_concentration 标签 + 写 risk_concentration 审计(金额只落合计与前 5 条摘要)。 5. risk_repository:find_pending_event_alert 改候选 LIMIT 50 + Python 过滤掉 payload.alert_subtype 含 agent_behavior 的单(评审 P0-1:代理人维度行为链单 不得充当客户维度事件单的聚合锚点);append_alert_event 加 extra_subtypes 合并进 payload.alert_subtype(不传时行为与原先一致,向后兼容)。 6. alert_service:subtypes 集合维护(空集不注入 payload,评审 P2-3); 追加时 alert_type 按「老单规则 ∪ 本批规则」重算(评审 P1-3,修掉既有 large_amount 单被本批仅 RISK-006(60) 翻转为 pattern 的缺陷); _publish_alert 加 notify_role/extra 可选参数(C5/C6 复用)。 7. 对话线:chat_tools.customer_context 加 profile(concentration_ratio/ r45_value/total_value/holdings_truncated),tool_service.summarize 加 「高风险持仓占比 X%(仅供参考)」;不新增意图词。 8. 02-redis-keys.md 增补 alert_subtype / escalation_level 附加推送字段。 测试:conftest 加 autouse _disable_concentration_rule(阈值推 1.01 做回归隔离, 现有用例断言零改动);test_risk_rules 加 RISK-006 纯函数 6 例;新建 tests/test_concentration_c4.py 11 例(与 RISK-001 同单聚合、score max=70、 L3 tag、risk_concentration 审计、仅集中度也出单、subtype 合并、P0-1 回归、 alert_type 不翻转、对话线 ratio)。全量 453 绿(436 + 17)。
220 lines
8.6 KiB
Python
220 lines
8.6 KiB
Python
"""风控事件规则 RISK-001~005(纯函数 · 规则权威 docs/PRD/附-风控规则表.md §2)。
|
||
|
||
输入约定:trades 为**当日 confirmed 的 subscribe/redeem 明细**(引擎组装;函数内再做防御过滤),
|
||
amount 为 Decimal,traded_at 为 datetime;输出命中规则列表(RuleHit)。
|
||
阈值由 RiskThresholds 打包(默认来自 settings,可 .env 覆盖)。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
from dataclasses import dataclass, field
|
||
from datetime import datetime, timedelta
|
||
from decimal import Decimal
|
||
from typing import Any
|
||
|
||
from app.config.settings import settings
|
||
|
||
# 规则 → 预警类型/静态风险分(PRD FR-3 静态映射;R-05 时替换动态评分)
|
||
RULE_SCORES: dict[str, int] = {
|
||
"RISK-001": 70,
|
||
"RISK-002": 70,
|
||
"RISK-003": 50,
|
||
"RISK-004": 80,
|
||
"RISK-005": 80,
|
||
"RISK-006": 60, # FR-8 持仓集中度(PRD §4A)
|
||
}
|
||
RULE_ALERT_TYPES: dict[str, str] = {
|
||
"RISK-001": "large_amount",
|
||
"RISK-002": "large_amount",
|
||
"RISK-003": "freq_trade",
|
||
"RISK-004": "pattern",
|
||
"RISK-005": "pattern",
|
||
"RISK-006": "pattern",
|
||
}
|
||
|
||
|
||
@dataclass(frozen=True)
|
||
class RiskThresholds:
|
||
large_amount: Decimal = Decimal("500000")
|
||
daily_total: Decimal = Decimal("500000")
|
||
freq_count: int = 3
|
||
probe_window_minutes: int = 5
|
||
probe_count: int = 3
|
||
probe_amount: Decimal = Decimal("400000")
|
||
small_amount: Decimal = Decimal("10000")
|
||
small_count: int = 3
|
||
# FR-8 RISK-006:R4+R5 市值占比阈值(0~1)。
|
||
# 评审 P1-2:from_settings 必须补读本字段——引擎/网关测试全走 from_settings
|
||
# 默认路径,漏读会让 conftest 的 autouse monkeypatch 失效,现有断言被 RISK-006 打穿。
|
||
concentration_threshold: float = 0.80
|
||
|
||
@classmethod
|
||
def from_settings(cls) -> "RiskThresholds":
|
||
s = settings
|
||
return cls(
|
||
large_amount=s.risk_large_amount,
|
||
daily_total=s.risk_daily_total,
|
||
freq_count=s.risk_freq_count,
|
||
probe_window_minutes=s.risk_probe_window_minutes,
|
||
probe_count=s.risk_probe_count,
|
||
probe_amount=s.risk_probe_amount,
|
||
small_amount=s.risk_small_amount,
|
||
small_count=s.risk_small_count,
|
||
concentration_threshold=s.risk_concentration_threshold,
|
||
)
|
||
|
||
|
||
@dataclass(frozen=True)
|
||
class RuleHit:
|
||
rule_id: str
|
||
alert_type: str
|
||
risk_score: int
|
||
detail: str # 触发说明(进预警单 payload / audit)
|
||
# 出单子类型:RISK-006→"concentration"、RISK-008→"agent_behavior";
|
||
# RISK-001~005 恒 None(不写 payload.alert_subtype)。C4 起启用。
|
||
alert_subtype: str | None = None
|
||
|
||
|
||
def _eligible(trade: dict[str, Any]) -> bool:
|
||
"""防御过滤:仅 confirmed 的 subscribe/redeem 计入(PRD FR-1 convert 不进事件线)。"""
|
||
return (
|
||
trade.get("trade_status", "confirmed") == "confirmed"
|
||
and trade.get("trade_type") in ("subscribe", "redeem")
|
||
)
|
||
|
||
|
||
def _amount(trade: dict[str, Any]) -> Decimal:
|
||
v = trade["amount"]
|
||
return v if isinstance(v, Decimal) else Decimal(str(v))
|
||
|
||
|
||
def rule_large_amount(trades: list[dict[str, Any]], th: RiskThresholds) -> RuleHit | None:
|
||
"""RISK-001 单笔大额:amount ≥ large_amount。"""
|
||
for t in trades:
|
||
if _amount(t) >= th.large_amount:
|
||
return RuleHit(
|
||
"RISK-001", RULE_ALERT_TYPES["RISK-001"], RULE_SCORES["RISK-001"],
|
||
f"单笔交易 {_amount(t)} 元 ≥ 阈值 {th.large_amount} 元(trade_id={t['trade_id']})",
|
||
)
|
||
return None
|
||
|
||
|
||
def rule_daily_total(trades: list[dict[str, Any]], th: RiskThresholds) -> RuleHit | None:
|
||
"""RISK-002 单日累计大额:当日申赎合计 ≥ daily_total(含本笔)。"""
|
||
total = sum((_amount(t) for t in trades), Decimal(0))
|
||
if total >= th.daily_total:
|
||
return RuleHit(
|
||
"RISK-002", RULE_ALERT_TYPES["RISK-002"], RULE_SCORES["RISK-002"],
|
||
f"当日申赎累计 {total} 元 ≥ 阈值 {th.daily_total} 元(共 {len(trades)} 笔)",
|
||
)
|
||
return None
|
||
|
||
|
||
def rule_freq_trade(trades: list[dict[str, Any]], th: RiskThresholds) -> RuleHit | None:
|
||
"""RISK-003 频繁交易:同一客户同一产品当日申赎合计 ≥ freq_count 笔。"""
|
||
by_product: dict[str, list[dict[str, Any]]] = {}
|
||
for t in trades:
|
||
by_product.setdefault(t["product_id"], []).append(t)
|
||
for product_id, pts in by_product.items():
|
||
if len(pts) >= th.freq_count:
|
||
return RuleHit(
|
||
"RISK-003", RULE_ALERT_TYPES["RISK-003"], RULE_SCORES["RISK-003"],
|
||
f"产品 {product_id} 当日申赎 {len(pts)} 笔 ≥ {th.freq_count} 笔",
|
||
)
|
||
return None
|
||
|
||
|
||
def rule_probe_pattern(
|
||
trades: list[dict[str, Any]], th: RiskThresholds, now: datetime
|
||
) -> RuleHit | None:
|
||
"""RISK-004 接近阈值试探:now 前 probe_window_minutes 内 ≥probe_count 笔且每笔 ≥probe_amount。"""
|
||
window_start = now - timedelta(minutes=th.probe_window_minutes)
|
||
in_window = [
|
||
t for t in trades
|
||
if window_start <= t["traded_at"] <= now and _amount(t) >= th.probe_amount
|
||
]
|
||
if len(in_window) >= th.probe_count:
|
||
return RuleHit(
|
||
"RISK-004", RULE_ALERT_TYPES["RISK-004"], RULE_SCORES["RISK-004"],
|
||
f"近 {th.probe_window_minutes} 分钟内 {len(in_window)} 笔每笔 ≥ {th.probe_amount} 元,"
|
||
f"接近大额阈值试探模式",
|
||
)
|
||
return None
|
||
|
||
|
||
def rule_small_then_large(trades: list[dict[str, Any]], th: RiskThresholds) -> RuleHit | None:
|
||
"""RISK-005 先小后大:当日时间序上首次大额之前已存在 ≥small_count 笔 ≤small_amount
|
||
(不要求连续、可穿插其他金额,PRD FR-3)。"""
|
||
ordered = sorted(trades, key=lambda t: t["traded_at"])
|
||
small_count = 0
|
||
for t in ordered:
|
||
if _amount(t) >= th.large_amount:
|
||
if small_count >= th.small_count:
|
||
return RuleHit(
|
||
"RISK-005", RULE_ALERT_TYPES["RISK-005"], RULE_SCORES["RISK-005"],
|
||
f"当日首次大额前已存在 {small_count} 笔 ≤{th.small_amount} 元小额交易,"
|
||
f"先小后大模式(trade_id={t['trade_id']})",
|
||
)
|
||
return None # 首次大额即判定结束
|
||
if _amount(t) <= th.small_amount:
|
||
small_count += 1
|
||
return None
|
||
|
||
|
||
def rule_concentration(profile: dict[str, Any], th: RiskThresholds) -> RuleHit | None:
|
||
"""RISK-006 持仓集中度:R4+R5 市值占比 ≥ 阈值即命中(FR-8 纯函数)。
|
||
|
||
输入为 `core_ro.concentration_profile` 的返回值(持仓画像,非流水),
|
||
故不并入 run_rules(输入域不同,见实现方案 §2.3:避免改动现有 13 处调用)。
|
||
|
||
- 空仓(total_value ≤ 0)不触发——无持仓无从谈集中度;
|
||
- `holdings_truncated=True` 视同达标(PRD FR-8 截断防护:明细被 limit 截断后
|
||
占比可能失真,按保守口径告警,宁可多报不可漏报)。
|
||
"""
|
||
total = profile.get("total_value") or Decimal(0)
|
||
total = total if isinstance(total, Decimal) else Decimal(str(total))
|
||
if total <= 0:
|
||
return None
|
||
|
||
r45 = profile.get("r45_value") or Decimal(0)
|
||
r45 = r45 if isinstance(r45, Decimal) else Decimal(str(r45))
|
||
ratio = r45 / total # Decimal 除法,避免 float 精度引发边界误判
|
||
|
||
truncated = bool(profile.get("holdings_truncated"))
|
||
if not truncated and ratio < Decimal(str(th.concentration_threshold)):
|
||
return None
|
||
|
||
if truncated:
|
||
reason = f"持仓明细触及查询上限,按保守口径视同集中度达标(实际占比 {float(ratio):.1%})"
|
||
else:
|
||
reason = f"R4+R5 持仓占比 {float(ratio):.1%} ≥ 阈值 {float(th.concentration_threshold):.0%}"
|
||
|
||
return RuleHit(
|
||
"RISK-006",
|
||
RULE_ALERT_TYPES["RISK-006"],
|
||
RULE_SCORES["RISK-006"],
|
||
f"{reason}(R4+R5 市值 {r45} 元 / 总市值 {total} 元)",
|
||
alert_subtype="concentration",
|
||
)
|
||
|
||
|
||
def run_rules(
|
||
trades: list[dict[str, Any]],
|
||
thresholds: RiskThresholds | None = None,
|
||
now: datetime | None = None,
|
||
) -> list[RuleHit]:
|
||
"""引擎入口的规则编排:对单笔交易事件后的当日流水跑全部规则,返回全部命中。"""
|
||
th = thresholds or RiskThresholds.from_settings()
|
||
now = now or datetime.now()
|
||
eligible = [t for t in trades if _eligible(t)]
|
||
if not eligible:
|
||
return []
|
||
hits: list[RuleHit | None] = [
|
||
rule_large_amount(eligible, th),
|
||
rule_daily_total(eligible, th),
|
||
rule_freq_trade(eligible, th),
|
||
rule_probe_pattern(eligible, th, now),
|
||
rule_small_then_large(eligible, th),
|
||
]
|
||
return [h for h in hits if h is not None]
|