#!/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 事件。**勿并行运行多个本脚本实例**(聚合锁为进程内 锁,跨进程并发重放同一 trade_id 无防护)。 **基金转换补偿(T-12 · PRD §7.1 / 架构 §5.4)**:`--convert-group CNV-xxx` 从 Core 侧(`core_trade` 两条流水 + `core_convert_lot_detail` 的 `nav`/`nav_date`) 补偿一次转换的**详情 + 预警**两件事,全部委托 `convert_service.compensate_convert` (锁在 service 内,锁键 `convert:rerun:{gid}`,与客户端带同键重试共用 → 人工补跑与客户端重试不会并发重复出单)。幂等锚点为**转出端 out_trade_id**: 一次转换有两条流水,只认转出端才不重复出单(评审 Q4)。 可复现性与局限: - 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 ...] python scripts/demo/rebuild_alerts.py --convert-group CNV-20260907-AB12CD34 退出码:0 = 全部成功/幂等跳过;1 = 有 trade_id 未找到,或 convert 组 Core 侧 流水不足(missing);3 = convert 组未抢到执行权(locked,稍后重试即可)。 (2 保留给 argparse 的参数用法错误。) """ 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.convert_repository import ConvertRepository # noqa: E402 from app.repository.core_ro import CoreReadOnlyRepository # noqa: E402 from app.repository.risk_repository import RiskRepository # noqa: E402 from app.service.convert.convert_service import compensate_convert # 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,当初 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: 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", "trade_id": trade_id, "alert_ids": result["alert_ids"], "triggered_rules": result["triggered_rules"], "aml_hit": result["aml_hit"], } def rebuild_convert_group( group_id: str, core: CoreReadOnlyRepository, repo: RiskRepository, crepo: ConvertRepository | None = None, ) -> dict[str, Any]: """按 convert_group_id 补偿一次转换(CLI 与单测共用入口)。 薄封装:全部逻辑在 `convert_service.compensate_convert`(锁 / 详情补写 / 预警 补偿都是那**一份**实现,本脚本不复制任何一条 —— 自检第 13 问)。 返回 `{"state": "rebuilt"|"skipped"|"missing"|"locked", ...}`。 """ return compensate_convert( group_id, core_ro=core, risk_repo=repo, convert_repo=crepo ) def main() -> None: parser = argparse.ArgumentParser( description="按 trade_id 幂等重放风控引擎;或按 --convert-group 补偿一次基金转换" ) parser.add_argument("trade_ids", nargs="*", help="core_trade 中的 trade_id,可一次多笔") parser.add_argument( "--convert-group", metavar="CNV-xxx", help="按 convert_group_id 从 Core 侧补偿一次转换(详情 + 预警),与 trade_id 互斥", ) args = parser.parse_args() if args.convert_group: if args.trade_ids: parser.error("--convert-group 与 trade_id 位置参数互斥,请择一使用") elif not args.trade_ids: parser.error("请给出 trade_id 位置参数,或使用 --convert-group CNV-xxx") 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) if args.convert_group: out = rebuild_convert_group(args.convert_group, core, repo, ConvertRepository()) gid = args.convert_group state = out["state"] if state == "rebuilt": rules = ",".join(out.get("triggered_rules") or []) or "-" print( f"[rebuilt] {gid}: detail={out['detail']} engine={out['engine']} " f"rules={rules} alerts={out['alert_ids']} " f"out={out['out_trade_id']} in={out['in_trade_id']}" ) elif state == "skipped": print( f"[skipped] {gid}: 详情已 completed、预警单已存在 " f"{out['alert_ids']}(幂等跳过)" ) else: print( f"[{state}] {gid}: " + ( f"Core 侧流水不足(实得 {out.get('trade_count')} 条,需 2 条)" if state == "missing" else "未抢到 convert:rerun 执行权,稍后重试" ), file=sys.stderr, ) if out.get("warning"): print(f"[warning] {gid}: {out['warning']}", file=sys.stderr) if state == "missing": raise SystemExit(1) if state == "locked": raise SystemExit(3) return 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']}(幂等跳过)") 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 if missing: raise SystemExit(f"{missing}/{len(args.trade_ids)} 笔 trade_id 未找到") if __name__ == "__main__": main()