fix: B9a 复审闭环——P2-1 skipped 分支 has_engine_error_audit 中间态警示(input_summary LIKE)+docstring 局限声明; P2-2 非首笔 trade 幂等用例; P3 顺手: 双脚本连接自检/ImportError 分支/非 dict payload raw 兜底/get_trade_by_id confirmed 过滤/勿并行声明, 8 例 204 绿, warning 路径真库复验后清理; docs 同步复审结论
This commit is contained in:
@@ -6,8 +6,16 @@
|
||||
预警缺失。用本脚本按 trade_id 补放引擎(也可用于人工补放任何一笔已落库交易)。
|
||||
|
||||
幂等:trade_id 已存在于任一 risk_alert.payload(含 aml 单与已处置单)即跳过,
|
||||
不重复出单、不重复 append 事件。重放以 core_trade.traded_at 为引擎事件时点,
|
||||
RISK-004 窗口判定可复现(engine.py 模块注释口径),不受脚本执行时刻影响。
|
||||
不重复出单、不重复 append 事件。**勿并行运行多个本脚本实例**(聚合锁为进程内
|
||||
锁,跨进程并发重放同一 trade_id 无防护)。
|
||||
|
||||
可复现性与局限:
|
||||
- RISK-004/005 规则窗口以 core_trade.traded_at 为事件时点,不受脚本执行时刻影响;
|
||||
聚合锚点(同日 pending 单查找)按执行日——跨日补放历史交易时事件并入执行日
|
||||
活跃单(events 明细含原始 traded_at 可追溯),即时补偿(交易日=执行日)不受影响。
|
||||
- 引擎编排非原子:若当初 engine_error 发生在首张预警单落库之后(aml 单/审计/L3
|
||||
部分缺失的中间态),该笔会被幂等跳过,脚本检测到 risk_engine_error 审计时输出
|
||||
warning 提示人工核对,不做选择性补齐。
|
||||
|
||||
用法:
|
||||
python scripts/demo/rebuild_alerts.py TRD-20260907-AB12CD34 [更多 trade_id ...]
|
||||
@@ -36,19 +44,26 @@ def rebuild_trade(
|
||||
"""重放单笔已落库交易(CLI 与单测共用入口)。
|
||||
|
||||
返回 {"state": "rebuilt"|"skipped"|"missing", "trade_id", "alert_ids", ...}:
|
||||
rebuilt 附 triggered_rules/aml_hit;skipped 附已入的预警单 alert_ids;
|
||||
missing 表示 core_trade 无此流水(未做任何写操作)。
|
||||
rebuilt 附 triggered_rules/aml_hit;skipped 附已入的预警单 alert_ids,当初
|
||||
engine_error 中断过的笔附 warning(部分失败中间态,人工核对);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 {
|
||||
out = {
|
||||
"state": "skipped",
|
||||
"trade_id": trade_id,
|
||||
"alert_ids": [row["alert_id"] for row in existing],
|
||||
}
|
||||
if repo.has_engine_error_audit(trade_id):
|
||||
out["warning"] = (
|
||||
"该笔曾引擎异常中断(risk_engine_error):已入单可能不完整"
|
||||
"(aml 单/审计/L3 或缺失),请人工核对预警台账与 audit_log"
|
||||
)
|
||||
return out
|
||||
result = process_trade_event(trade, core_ro=core, risk_repo=repo)
|
||||
return {
|
||||
"state": "rebuilt",
|
||||
@@ -67,6 +82,15 @@ def main() -> None:
|
||||
core = CoreReadOnlyRepository()
|
||||
repo = RiskRepository()
|
||||
|
||||
# 连接自检(惰性引擎构造不连库;零写入探针,兼验证双库可达)
|
||||
try:
|
||||
core.get_trade_by_id("__connectivity_probe__")
|
||||
repo.find_alerts_by_trade("__connectivity_probe__")
|
||||
except Exception as exc:
|
||||
print(f"MySQL 不可达:{exc}", file=sys.stderr)
|
||||
print("请确认本机 MySQL 服务已启动(见 FLOW §0 本机状态)。", file=sys.stderr)
|
||||
raise SystemExit(1)
|
||||
|
||||
missing = 0
|
||||
for tid in args.trade_ids:
|
||||
out = rebuild_trade(tid, core, repo)
|
||||
@@ -75,6 +99,8 @@ def main() -> None:
|
||||
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']}(幂等跳过)")
|
||||
if out.get("warning"):
|
||||
print(f"[warning] {tid}: {out['warning']}", file=sys.stderr)
|
||||
else:
|
||||
print(f"[missing] {tid}: core_trade 无此流水", file=sys.stderr)
|
||||
missing += 1
|
||||
|
||||
@@ -32,6 +32,8 @@ CHANNEL = "risk:pub:alert"
|
||||
|
||||
def format_alert(payload: dict[str, Any]) -> str:
|
||||
"""推送 payload → 单行可读文本(演示展示用,字段缺失容错)。"""
|
||||
if "raw" in payload:
|
||||
return f"[raw] {payload['raw']}"
|
||||
roles = ",".join(payload.get("notify_role") or [])
|
||||
return (
|
||||
f"[{payload.get('alert_type', '?'):>12}] {payload.get('alert_id', '?')}"
|
||||
@@ -43,7 +45,11 @@ def format_alert(payload: dict[str, Any]) -> str:
|
||||
|
||||
|
||||
def run(duration: float) -> None:
|
||||
import redis
|
||||
try:
|
||||
import redis
|
||||
except ImportError:
|
||||
print("缺少 redis 包:请执行 python -m pip install -r requirements.txt", file=sys.stderr)
|
||||
raise SystemExit(1)
|
||||
|
||||
client = redis.Redis.from_url(settings.redis_url, decode_responses=True)
|
||||
try:
|
||||
@@ -68,6 +74,8 @@ def run(duration: float) -> None:
|
||||
received += 1
|
||||
try:
|
||||
payload = json.loads(msg["data"])
|
||||
if not isinstance(payload, dict):
|
||||
payload = {"raw": str(msg["data"])}
|
||||
except (TypeError, ValueError):
|
||||
payload = {"raw": str(msg["data"])}
|
||||
print(f"{time.strftime('%H:%M:%S')} #{received} {format_alert(payload)}")
|
||||
|
||||
Reference in New Issue
Block a user