279 lines
11 KiB
Python
279 lines
11 KiB
Python
"""预警单聚合与通知(B2 · 架构 §5.2 / PRD FR-4)。
|
||
|
||
职责:单事件命中规则的**聚合决策与编排**(merge 原语 append_alert_event 在 repo,A3 评审口径);
|
||
聚合锁原语公共化至 locks(B7 挂账①);审计落库;Pub/Sub 通知广播经
|
||
redis_gateway 单例(B7:连接由 lifespan 管理,publish 失败降级不阻塞落库)。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
from datetime import date, datetime, time
|
||
from typing import Any
|
||
from uuid import uuid4
|
||
|
||
from app.repository.risk_repository import EVENT_ALERT_TYPES, RiskRepository
|
||
from app.service.risk import redis_gateway
|
||
from app.service.risk.locks import run_locked
|
||
from app.service.risk.rules import RuleHit
|
||
from app.utils.exceptions import NotFoundError, StateConflict
|
||
from app.utils.trace import current_trace, new_trace
|
||
|
||
TIER_SCORE = {"large_amount": 70, "freq_trade": 50, "pattern": 80, "suitability": 90, "aml": 95}
|
||
|
||
|
||
def set_publisher(publisher: Any | None) -> None:
|
||
"""测试/集成注入点(旧名兼容,B7 起转发 redis_gateway.set_gateway)。"""
|
||
redis_gateway.set_gateway(publisher)
|
||
|
||
|
||
# ---------- 进程内聚合锁:公共原语见 locks.py(B7 挂账①) ----------
|
||
|
||
|
||
def _day_start(day: date | None = None) -> datetime:
|
||
return datetime.combine(day or date.today(), time.min)
|
||
|
||
|
||
def _new_alert_id() -> str:
|
||
return f"ALT-{datetime.now():%Y%m%d}-{uuid4().hex[:8].upper()}"
|
||
|
||
|
||
def _audit(
|
||
risk_repo: RiskRepository,
|
||
*,
|
||
customer_id: str | None,
|
||
rule_id: str | None,
|
||
decision: str,
|
||
risk_score: int | None,
|
||
input_summary: dict[str, Any],
|
||
event_type: str = "risk_judgement",
|
||
actor_id: str = "SYSTEM",
|
||
) -> None:
|
||
risk_repo.insert_audit_log(
|
||
{
|
||
"trace_id": current_trace() or new_trace(),
|
||
"event_type": event_type,
|
||
"agent_type": "risk",
|
||
"actor_id": actor_id,
|
||
"customer_id": customer_id,
|
||
"rule_id": rule_id,
|
||
"input_summary": input_summary,
|
||
"decision": decision,
|
||
"risk_score": risk_score,
|
||
"handler_id": None,
|
||
"handler_result": None,
|
||
"handler_comment": None,
|
||
}
|
||
)
|
||
|
||
|
||
def _publish_alert(alert: dict[str, Any]) -> None:
|
||
redis_gateway.publish(
|
||
"risk:pub:alert",
|
||
{
|
||
"alert_id": alert["alert_id"],
|
||
"alert_type": alert["alert_type"],
|
||
"customer_id_mask": alert["customer_id"][:6] + "**",
|
||
"risk_score": alert["risk_score"],
|
||
"trace_id": alert["trace_id"],
|
||
"notify_role": ["risk_officer"] + (["compliance"] if alert["alert_type"] == "aml" else []),
|
||
},
|
||
)
|
||
|
||
|
||
def record_trade_alerts(
|
||
trade: dict[str, Any],
|
||
hits: list[RuleHit],
|
||
risk_repo: RiskRepository | None = None,
|
||
customer_context: dict[str, Any] | None = None,
|
||
) -> dict[str, Any] | None:
|
||
"""交易事件预警入口(引擎 B4 调用)。
|
||
|
||
hits 为空 → 审计 pass;否则事件类命中聚合为一张单(PRD FR-4:同客户同日仅一张事件类
|
||
pending 单,alert_type 随最高分规则动态更新)。
|
||
"""
|
||
risk_repo = risk_repo or RiskRepository()
|
||
if not hits:
|
||
_audit(risk_repo, customer_id=trade["customer_id"], rule_id=None, decision="pass",
|
||
risk_score=None, input_summary={"trade_id": trade["trade_id"]})
|
||
return None
|
||
|
||
best = max(hits, key=lambda h: h.risk_score)
|
||
triggered_rules = sorted({h.rule_id for h in hits})
|
||
risk_score = max(h.risk_score for h in hits)
|
||
event = {
|
||
"trade_id": trade["trade_id"],
|
||
"product_id": trade["product_id"],
|
||
"trade_type": trade["trade_type"],
|
||
"amount": str(trade["amount"]),
|
||
"traded_at": str(trade["traded_at"]),
|
||
"matched_rules": [h.rule_id for h in hits],
|
||
"details": {h.rule_id: h.detail for h in hits},
|
||
}
|
||
|
||
def _agg(locked: bool) -> dict[str, Any]:
|
||
if locked:
|
||
pending = risk_repo.find_pending_event_alert(trade["customer_id"], _day_start())
|
||
if pending:
|
||
risk_repo.append_alert_event(
|
||
pending["alert_id"], event, triggered_rules, risk_score, best.alert_type
|
||
)
|
||
updated = risk_repo.get_alert(pending["alert_id"])
|
||
_audit(risk_repo, customer_id=trade["customer_id"], rule_id=",".join(triggered_rules),
|
||
decision="alert_appended", risk_score=updated["risk_score"],
|
||
input_summary=event)
|
||
_publish_alert(updated)
|
||
return updated
|
||
alert = {
|
||
"alert_id": _new_alert_id(),
|
||
"trace_id": current_trace() or new_trace(),
|
||
"customer_id": trade["customer_id"],
|
||
"trade_id": trade["trade_id"],
|
||
"alert_type": best.alert_type,
|
||
"triggered_rules": triggered_rules,
|
||
"risk_score": risk_score,
|
||
"status": "pending_review",
|
||
"payload": {
|
||
"product_id": trade["product_id"],
|
||
"events": [event],
|
||
"customer_context": customer_context or {},
|
||
},
|
||
}
|
||
risk_repo.insert_alert(alert)
|
||
_audit(risk_repo, customer_id=trade["customer_id"], rule_id=",".join(triggered_rules),
|
||
decision="alert_created", risk_score=risk_score, input_summary=event)
|
||
_publish_alert(alert)
|
||
return alert
|
||
|
||
return run_locked(f"agg:event:{trade['customer_id']}:{date.today()}", _agg)
|
||
|
||
|
||
def record_suitability_alert(
|
||
trade_request: dict[str, Any],
|
||
rule_id: str,
|
||
block_reason: str,
|
||
risk_repo: RiskRepository | None = None,
|
||
) -> dict[str, Any]:
|
||
"""R-02 阻断预警(网关 FR-1 调用):同客户+产品+日仅一张 pending 单。"""
|
||
risk_repo = risk_repo or RiskRepository()
|
||
event = {
|
||
"trade_id": trade_request["trade_id"],
|
||
"product_id": trade_request["product_id"],
|
||
"trade_type": trade_request["trade_type"],
|
||
"amount": str(trade_request["amount"]),
|
||
"traded_at": str(trade_request["traded_at"]),
|
||
"matched_rules": [rule_id],
|
||
"details": {rule_id: block_reason},
|
||
}
|
||
best_type, score = "suitability", TIER_SCORE["suitability"]
|
||
|
||
def _agg(locked: bool) -> dict[str, Any]:
|
||
if locked:
|
||
pending = risk_repo.find_pending_suitability_alert(
|
||
trade_request["customer_id"], trade_request["product_id"], _day_start()
|
||
)
|
||
if pending:
|
||
risk_repo.append_alert_event(
|
||
pending["alert_id"], event, [rule_id], score, best_type
|
||
)
|
||
updated = risk_repo.get_alert(pending["alert_id"])
|
||
_audit(risk_repo, customer_id=trade_request["customer_id"], rule_id=rule_id,
|
||
decision="suitability_alert_appended", risk_score=score,
|
||
input_summary=event, event_type="suitability_block")
|
||
_publish_alert(updated)
|
||
return updated
|
||
alert = {
|
||
"alert_id": _new_alert_id(),
|
||
"trace_id": current_trace() or new_trace(),
|
||
"customer_id": trade_request["customer_id"],
|
||
"trade_id": trade_request["trade_id"],
|
||
"alert_type": best_type,
|
||
"triggered_rules": [rule_id],
|
||
"risk_score": score,
|
||
"status": "pending_review",
|
||
"payload": {
|
||
"product_id": trade_request["product_id"],
|
||
"events": [event],
|
||
"customer_context": {},
|
||
},
|
||
}
|
||
risk_repo.insert_alert(alert)
|
||
_audit(risk_repo, customer_id=trade_request["customer_id"], rule_id=rule_id,
|
||
decision="suitability_alert_created", risk_score=score,
|
||
input_summary=event, event_type="suitability_block")
|
||
_publish_alert(alert)
|
||
return alert
|
||
|
||
return run_locked(
|
||
f"agg:suitability:{trade_request['customer_id']}:{trade_request['product_id']}:{date.today()}",
|
||
_agg,
|
||
)
|
||
|
||
|
||
def handle_alert(
|
||
alert_id: str,
|
||
handler_result: str,
|
||
handler_id: str,
|
||
handler_comment: str | None = None,
|
||
risk_repo: RiskRepository | None = None,
|
||
) -> dict[str, Any]:
|
||
"""人工处置(FR-4;角色鉴权在依赖层,本层管状态机与审计)。
|
||
|
||
仅 pending_review 可处置(禁止跳改已处置单,A-7);状态变更与处置审计
|
||
**同事务提交**(B7 挂账⑦:原两事务非原子,审计失败会留下无痕的状态变更,
|
||
违反 P-05 全量留痕);处置动作全量写 audit_log(agent_type='risk')。
|
||
"""
|
||
repo = risk_repo or RiskRepository()
|
||
alert = repo.get_alert(alert_id)
|
||
if alert is None:
|
||
raise NotFoundError(f"alert not found: {alert_id}")
|
||
if alert["status"] != "pending_review":
|
||
raise StateConflict(f"alert {alert_id} already handled (status={alert['status']})")
|
||
ok = repo.handle_alert_with_audit(
|
||
alert_id,
|
||
handler_result,
|
||
handler_id,
|
||
handler_comment,
|
||
audit_entry={
|
||
"trace_id": current_trace() or new_trace(),
|
||
"event_type": "alert_handle",
|
||
"agent_type": "risk",
|
||
"actor_id": handler_id,
|
||
"customer_id": alert["customer_id"],
|
||
"rule_id": None,
|
||
"input_summary": {"alert_id": alert_id, "prev_status": alert["status"]},
|
||
"decision": "alert_handled",
|
||
"risk_score": None,
|
||
"handler_id": handler_id,
|
||
"handler_result": handler_result,
|
||
"handler_comment": handler_comment,
|
||
},
|
||
)
|
||
if not ok: # 读后并发窗口:他人已抢先处置
|
||
raise StateConflict(f"alert {alert_id} state changed concurrently")
|
||
return repo.get_alert(alert_id)
|
||
|
||
|
||
def record_aml_alert(
|
||
customer_id: str,
|
||
aml_detail: dict[str, Any],
|
||
risk_repo: RiskRepository | None = None,
|
||
) -> dict[str, Any]:
|
||
"""R-03 AML 命中:独立出单不聚合(PRD FR-4),单事件单张。"""
|
||
risk_repo = risk_repo or RiskRepository()
|
||
alert = {
|
||
"alert_id": _new_alert_id(),
|
||
"trace_id": current_trace() or new_trace(),
|
||
"customer_id": customer_id,
|
||
"trade_id": aml_detail.get("trade_id"),
|
||
"alert_type": "aml",
|
||
"triggered_rules": ["AML-001"],
|
||
"risk_score": TIER_SCORE["aml"],
|
||
"status": "pending_review",
|
||
"payload": {"product_id": aml_detail.get("product_id"), "events": [aml_detail], "customer_context": {}},
|
||
}
|
||
risk_repo.insert_alert(alert)
|
||
_audit(risk_repo, customer_id=customer_id, rule_id="AML-001", decision="aml_alert_created",
|
||
risk_score=TIER_SCORE["aml"], input_summary=aml_detail, event_type="aml_hit")
|
||
_publish_alert(alert)
|
||
return alert
|