feat: B9a 演示/运维脚本——subscribe_alerts(risk:pub:alert 订阅演示/连接自检/--duration) + rebuild_alerts(按 trade_id 幂等重放: find_alerts_by_trade payload LIKE 查已入单防重复出单与重复 append, 含 aml/已处置单) + core_ro.get_trade_by_id 只读扩展 + test_demo_scripts 5 例, 201 绿; 真库手工验证(出单→幂等 skip→missing exit1; 订阅端到端 trace 贯通)后现场清理; docs 同步(TODO/MEMORY/FLOW/开发计划)
This commit is contained in:
@@ -0,0 +1,87 @@
|
||||
#!/usr/bin/env python3
|
||||
"""按 trade_id 幂等重放风控引擎(B9a 补偿 · trade_gateway engine_error 场景)。
|
||||
|
||||
背景(PRD FR-1 / 架构 §5.3):网关 INSERT core_trade 成功后 process_trade_event
|
||||
异常时交易已成立——响应带 engine_error=true、审计 decision='risk_engine_error',
|
||||
预警缺失。用本脚本按 trade_id 补放引擎(也可用于人工补放任何一笔已落库交易)。
|
||||
|
||||
幂等:trade_id 已存在于任一 risk_alert.payload(含 aml 单与已处置单)即跳过,
|
||||
不重复出单、不重复 append 事件。重放以 core_trade.traded_at 为引擎事件时点,
|
||||
RISK-004 窗口判定可复现(engine.py 模块注释口径),不受脚本执行时刻影响。
|
||||
|
||||
用法:
|
||||
python scripts/demo/rebuild_alerts.py TRD-20260907-AB12CD34 [更多 trade_id ...]
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[2]
|
||||
sys.path.insert(0, str(ROOT))
|
||||
|
||||
from app.repository.core_ro import CoreReadOnlyRepository # noqa: E402
|
||||
from app.repository.risk_repository import RiskRepository # noqa: E402
|
||||
from app.service.risk.engine import process_trade_event # noqa: E402
|
||||
|
||||
|
||||
def rebuild_trade(
|
||||
trade_id: str,
|
||||
core: CoreReadOnlyRepository,
|
||||
repo: RiskRepository,
|
||||
) -> dict[str, Any]:
|
||||
"""重放单笔已落库交易(CLI 与单测共用入口)。
|
||||
|
||||
返回 {"state": "rebuilt"|"skipped"|"missing", "trade_id", "alert_ids", ...}:
|
||||
rebuilt 附 triggered_rules/aml_hit;skipped 附已入的预警单 alert_ids;
|
||||
missing 表示 core_trade 无此流水(未做任何写操作)。
|
||||
"""
|
||||
trade = core.get_trade_by_id(trade_id)
|
||||
if trade is None:
|
||||
return {"state": "missing", "trade_id": trade_id, "alert_ids": []}
|
||||
existing = repo.find_alerts_by_trade(trade_id)
|
||||
if existing:
|
||||
return {
|
||||
"state": "skipped",
|
||||
"trade_id": trade_id,
|
||||
"alert_ids": [row["alert_id"] for row in existing],
|
||||
}
|
||||
result = process_trade_event(trade, core_ro=core, risk_repo=repo)
|
||||
return {
|
||||
"state": "rebuilt",
|
||||
"trade_id": trade_id,
|
||||
"alert_ids": result["alert_ids"],
|
||||
"triggered_rules": result["triggered_rules"],
|
||||
"aml_hit": result["aml_hit"],
|
||||
}
|
||||
|
||||
|
||||
def main() -> None:
|
||||
parser = argparse.ArgumentParser(description="按 trade_id 幂等重放风控引擎(引擎异常补偿)")
|
||||
parser.add_argument("trade_ids", nargs="+", help="core_trade 中的 trade_id,可一次多笔")
|
||||
args = parser.parse_args()
|
||||
|
||||
core = CoreReadOnlyRepository()
|
||||
repo = RiskRepository()
|
||||
|
||||
missing = 0
|
||||
for tid in args.trade_ids:
|
||||
out = rebuild_trade(tid, core, repo)
|
||||
if out["state"] == "rebuilt":
|
||||
rules = ",".join(out["triggered_rules"]) or "-"
|
||||
print(f"[rebuilt] {tid}: rules={rules} alerts={out['alert_ids']} aml={out['aml_hit']}")
|
||||
elif out["state"] == "skipped":
|
||||
print(f"[skipped] {tid}: 已入预警单 {out['alert_ids']}(幂等跳过)")
|
||||
else:
|
||||
print(f"[missing] {tid}: core_trade 无此流水", file=sys.stderr)
|
||||
missing += 1
|
||||
|
||||
if missing:
|
||||
raise SystemExit(f"{missing}/{len(args.trade_ids)} 笔 trade_id 未找到")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,95 @@
|
||||
#!/usr/bin/env python3
|
||||
"""risk:pub:alert 预警推送订阅演示(B9a · PRD FR-4 通知 / A-3 验收第三条)。
|
||||
|
||||
演示 SOP:
|
||||
终端 1 python scripts/demo/subscribe_alerts.py
|
||||
终端 2 uvicorn app.main:app --reload → Swagger 发演示交易
|
||||
(A-3:CUST-3001 申购 50 万 PROD-510300,RISK-001/002 命中)
|
||||
终端 1 应实时打印预警推送行;无推送时每 30s 打印心跳提示仍在监听。
|
||||
|
||||
payload 口径(PRD §4 FR-4 · 02-redis-keys.md §2.4):
|
||||
{alert_id, alert_type, customer_id_mask, risk_score, trace_id, notify_role}
|
||||
|
||||
--duration N:监听 N 秒后自动退出(默认 0 = 不限,Ctrl+C 退出)。
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import sys
|
||||
import time
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[2]
|
||||
sys.path.insert(0, str(ROOT))
|
||||
|
||||
from app.config.settings import settings # noqa: E402
|
||||
|
||||
CHANNEL = "risk:pub:alert"
|
||||
|
||||
|
||||
def format_alert(payload: dict[str, Any]) -> str:
|
||||
"""推送 payload → 单行可读文本(演示展示用,字段缺失容错)。"""
|
||||
roles = ",".join(payload.get("notify_role") or [])
|
||||
return (
|
||||
f"[{payload.get('alert_type', '?'):>12}] {payload.get('alert_id', '?')}"
|
||||
f" score={payload.get('risk_score', '?')}"
|
||||
f" customer={payload.get('customer_id_mask', '?')}"
|
||||
f" trace={payload.get('trace_id', '?')}"
|
||||
f" notify=[{roles}]"
|
||||
)
|
||||
|
||||
|
||||
def run(duration: float) -> None:
|
||||
import redis
|
||||
|
||||
client = redis.Redis.from_url(settings.redis_url, decode_responses=True)
|
||||
try:
|
||||
client.ping()
|
||||
except Exception as exc:
|
||||
print(f"Redis 不可达({settings.redis_url}):{exc}", file=sys.stderr)
|
||||
print("请先启动本机 Redis 服务(见 FLOW §0 本机状态)。", file=sys.stderr)
|
||||
raise SystemExit(1)
|
||||
|
||||
pubsub = client.pubsub(ignore_subscribe_messages=True)
|
||||
pubsub.subscribe(CHANNEL)
|
||||
print(f"已订阅 {CHANNEL} @ {settings.redis_url}(Ctrl+C 退出)")
|
||||
print("等待预警推送…(另开终端发演示交易,SOP 见脚本头注释)")
|
||||
|
||||
deadline = time.monotonic() + duration if duration > 0 else None
|
||||
received = 0
|
||||
last_heartbeat = time.monotonic()
|
||||
try:
|
||||
while True:
|
||||
msg = pubsub.get_message(timeout=1.0)
|
||||
if msg and msg.get("type") == "message":
|
||||
received += 1
|
||||
try:
|
||||
payload = json.loads(msg["data"])
|
||||
except (TypeError, ValueError):
|
||||
payload = {"raw": str(msg["data"])}
|
||||
print(f"{time.strftime('%H:%M:%S')} #{received} {format_alert(payload)}")
|
||||
elif deadline is not None and time.monotonic() >= deadline:
|
||||
print(f"已监听 {duration:g}s,共收到 {received} 条推送,退出。")
|
||||
return
|
||||
elif deadline is None and time.monotonic() - last_heartbeat >= 30:
|
||||
last_heartbeat = time.monotonic()
|
||||
print(f"…监听中(已收到 {received} 条)")
|
||||
except KeyboardInterrupt:
|
||||
print(f"\n退出(共收到 {received} 条推送)。")
|
||||
finally:
|
||||
pubsub.close()
|
||||
client.close()
|
||||
|
||||
|
||||
def main() -> None:
|
||||
parser = argparse.ArgumentParser(description="订阅 risk:pub:alert(风控预警推送演示)")
|
||||
parser.add_argument("--duration", type=float, default=0.0, help="监听秒数;0=不限(Ctrl+C 退出)")
|
||||
args = parser.parse_args()
|
||||
run(args.duration)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user