"""T-12 真 MySQL 验证脚本:补偿 + SLA 清理在真库/真账号下实跑(T+1 受理/确认分离链路)。 **为什么 sqlite 单测全绿还不够(T-12 · T+1 视角)** 1. **`core_convert_request.status='expired'` 的真库 6 态 ENUM 值域**:写错值会**直接 报错**;而 sqlite 测试库该列是 `VARCHAR`,任何字符串都收。故「cleanup 能把 受理单写进 expired」这件事必须在真库证明一次。 2. **`input_summary` 是真 JSON 列**:`has_engine_error_audit` 用 `input_summary LIKE '%gid%'` 做幂等/警示判定。sqlite 里它是 TEXT、必定能 LIKE; MySQL 的 JSON 列能否直接 LIKE 属方言行为,必须实跑确认(E 组专测)。 3. **DATETIME(3) 的时间边界**:`list_inflight_before` 的 `requested_at < cutoff` 在 MySQL 毫秒精度列上是否按预期分流(超 SLA 进候选、当日恒不进),sqlite 的 字符串时间给不了这个保证。 4. **端到端**:受理 → T+1 确认(第⑦步引擎炸 + 第⑧步镜像炸)→ Core 三件套完整 → 补偿 `rebuilt`,走真事务、真账号、真引擎。 5. **真库交易日历**:cleanup 的 SLA 边界与 `confirm_batch` 捞单窗口(`_window_start`) 在**同一份真库日历**上数值相等(B 组对表);受理日 T / 确认日 T1 由生产 `trading_calendar` 实时推算,脚本对任意运行日稳、不手算日期、不猜自然日。 **断言分组** A 待补偿态(真实链路:`accept_convert` → `confirm_one` 引擎 + 镜像双失败)在真库如实落痕 B cleanup:超 SLA 在途受理单 → `expired`(ENUM 值域 / 窗口同口径 / dry-run / 幂等 / S2 不硬删 / 镜像不推) C `compensate_convert` 端到端:镜像 `completed` + 1 张单 + `confirmed` 审计(补偿来源) D 重复补偿 → `skipped` + 「人工核对」提示(engine_error 中间态) E 真 **JSON 列**上 `input_summary LIKE` 生效(`has_engine_error_audit`) F 清理后残留为零(自检) 用法: python scripts/dev/verify_convert_compensate.py 约定(与 T-6 ~ T-11 脚本一致): - 用 **T12C 前缀**隔离数据(客户 / 产品 / 批次 / group_id / client_request_id), 跑完**两个库**(core + agent)全清; - 建/清数据走 `role="admin"`(需 DELETE);业务链路走默认账号(D20 账号分离在 链路上生效:core_ro=ro / 仓储默认 rw); - 固定 `trace_id = "T12C-VERIFY-TRACE"`,便于精准清理 `audit_log`; - 期望值一律「与 Core 侧三件套对表」或由生产纯函数推出,禁手算(自检第 13 问); - 交易日历留痕:A 股连续休市历史最长为 **2020 年春节 10 个自然日**(国务院办公厅 延长春节假期,上交所《关于2020年春节延长休市相关业务衔接安排的通知》:休市延至 2 月 2 日、2 月 3 日恢复开市;联网核实 2026-09-12)。本脚本所有跨日推算都用 `previous_biz_day` / `next_biz_day` 链式走**真库交易日历**,不猜自然日 —— 假期 安排怎么变都不影响(`trading_calendar._MAX_SCAN_DAYS=256` 的量级依据即此); - 退出码 1 = 有断言不一致。 """ from __future__ import annotations import json import sys import uuid from datetime import date, datetime, time, timedelta from decimal import Decimal from pathlib import Path from unittest.mock import patch from sqlalchemy import text ROOT = Path(__file__).resolve().parents[2] sys.path.insert(0, str(ROOT)) sys.path.insert(0, str(ROOT / "scripts" / "agent")) # import cleanup_pending_convert from app.config.settings import settings # noqa: E402 from app.repository.convert_repository import ConvertRepository # noqa: E402 from app.repository.convert_request_repository import ( # noqa: E402 STATUS_NAV_PENDING, ConvertRequestRepository, ) from app.repository.core_ro import CoreReadOnlyRepository # noqa: E402 from app.repository.risk_repository import RiskRepository # noqa: E402 from app.service.convert.confirm_service import ( # noqa: E402 _window_start, confirm_one, ) from app.service.convert.convert_service import ( # noqa: E402 accept_convert, compensate_convert, ) from app.service.convert.trading_calendar import ( # noqa: E402 is_biz_day, next_biz_day, parse_cutoff, previous_biz_day, resolve_accept_date, ) from app.service.risk.rules import RiskThresholds # noqa: E402 from app.utils.db import dispose_engines, get_engine # noqa: E402 from app.utils.trace import new_trace # noqa: E402 from cleanup_pending_convert import cleanup as sla_cleanup # noqa: E402 # ── 隔离常量(全部 T12C 前缀)──────────────────────────────────────── CUSTOMER = "CUST-T12C" PROD_OUT = "PROD-T12CO" PROD_IN = "PROD-T12CI" COMPANY = "华夏模拟基金" TA = "TA-CN-001" OUT_RATE = Decimal("0.0030") IN_RATE = Decimal("0.0080") OUT_NAV = Decimal("1.3604") # T 日净值(T+1 未知价法:转出/转入折算都用 T 日净值) IN_NAV = Decimal("1.9194") FEE_TIERS = [ (0, 7, "0.0150"), (7, 30, "0.0100"), (30, 180, "0.0050"), (180, 365, "0.0025"), (365, None, "0.0000"), ] #: 两批:A1 持 100 天(0.0050 档)+ A2 持 3 天(0.0150 档);QTY=120000 → FIFO 两批都参与 LOT_SPEC = [("LOT-T12C-A1", "100000", 100), ("LOT-T12C-A2", "50000", 3)] TOTAL_QTY = sum(Decimal(q) for _, q, _ in LOT_SPEC) QTY = "120000" CID_REQ = "T12C-REQ-0001" TRACE = "T12C-VERIFY-TRACE" #: 大额阈值 100000 —— 转出金额 120000 × 1.3604 ≈ 163248 > 100000,确保补偿时 #: 引擎**真出一张单**(RISK-001,可断言);daily_total 调大避免 RISK-002 干扰。 TH = RiskThresholds( large_amount=Decimal("100000"), daily_total=Decimal("1000000"), freq_count=3, probe_window_minutes=5, probe_count=3, probe_amount=Decimal("400000"), small_amount=Decimal("10000"), small_count=3, concentration_threshold=1.01, # 持仓画像不参与(避免 RISK-006 干扰) ) _passed = 0 _failed = 0 #: 建/清数据用 admin 引擎(需 DELETE,R-e);在 main 中初始化 core_engine = None agent_engine = None def check(name: str, actual, expected) -> None: global _passed, _failed ok = actual == expected if ok: _passed += 1 else: _failed += 1 flag = "✅" if ok else "❌" print(f" {flag} {name}: 实际 {actual!r}" + ("" if ok else f" / 期望 {expected!r}")) def _t12c_id(prefix: str, now: datetime) -> str: """固定前缀 id 工厂:默认 `CNV-<日期>-` 与清理口径 `LIKE 'CNV-T12C%'` 不匹配, 重跑会留撞唯一键的残留(T-7 脚本踩过同款坑)。""" return f"{prefix}-T12C-{uuid.uuid4().hex[:8].upper()}" def q1(engine, sql: str, **params): with engine.connect() as conn: return conn.execute(text(sql), params).scalar_one_or_none() def rows(engine, sql: str, **params) -> list[dict]: with engine.connect() as conn: return [dict(r) for r in conn.execute(text(sql), params).mappings()] # ── 数据准备 / 清理 ───────────────────────────────────────────────── def cleanup_env(core_engine, agent_engine) -> None: """按 T12C 前缀清理两个库(含幂等重跑前置清理)。 清场顺序遵循 FK:core_convert_request(from/to_product_id → core_product, customer_id → core_customer)必须先于 product / customer 删除 (T-12 真库 fixture 踩过 FK 1451 坑)。 """ with core_engine.begin() as conn: for sql, params in [ # 受理单:按客户(链路产生的单)+ 按 client_request_id(B 组直插的单)双路清 ("DELETE FROM core_convert_request WHERE customer_id = :c", {"c": CUSTOMER}), ("DELETE FROM core_convert_request WHERE client_request_id LIKE 'T12C-%'", {}), ("DELETE FROM core_convert_lot_detail WHERE convert_group_id LIKE 'CNV-T12C%'", {}), ("DELETE FROM core_trade WHERE customer_id = :c", {"c": CUSTOMER}), ("DELETE FROM core_share_lot WHERE customer_id = :c", {"c": CUSTOMER}), ("DELETE FROM core_holding WHERE customer_id = :c", {"c": CUSTOMER}), ("DELETE FROM core_customer_risk WHERE customer_id = :c", {"c": CUSTOMER}), ("DELETE FROM core_fee_rule WHERE product_id LIKE 'PROD-T12C%'", {}), ("DELETE FROM core_product_nav WHERE product_id LIKE 'PROD-T12C%'", {}), ("DELETE FROM core_product WHERE product_id LIKE 'PROD-T12C%'", {}), # T-16 起确认事务写 core_cash_flow(fk_cf_customer 挡在客户删除前) ("DELETE FROM core_cash_flow WHERE remark LIKE 'convert:%'", {}), ("DELETE FROM core_customer WHERE customer_id = :c", {"c": CUSTOMER}), ]: conn.execute(text(sql), params) with agent_engine.begin() as conn: conn.execute( text( "DELETE FROM risk_convert_detail " "WHERE convert_group_id LIKE 'CNV-T12C%' OR client_request_id LIKE 'T12C-%'" ) ) conn.execute(text("DELETE FROM risk_alert WHERE customer_id = :c"), {"c": CUSTOMER}) conn.execute( text("DELETE FROM risk_suitability_log WHERE customer_id = :c"), {"c": CUSTOMER} ) conn.execute(text("DELETE FROM customer_profile_l3 WHERE customer_id = :c"), {"c": CUSTOMER}) conn.execute(text("DELETE FROM audit_log WHERE trace_id LIKE 'T12C-VERIFY%'")) def seed(core_engine, accept_date: date) -> None: """基础数据:两端净值都落在**受理日 T**(确认段 `get_nav_on(pid, accept_date)` 精确匹配)。""" base = datetime.combine(accept_date, time(10, 0)) with core_engine.begin() as conn: conn.execute( text( "INSERT INTO core_customer (customer_id, display_name, open_date) " "VALUES (:c, 'T12补偿真库验证', :d)" ), {"c": CUSTOMER, "d": accept_date - timedelta(days=400)}, ) conn.execute( text( "INSERT INTO core_customer_risk (customer_id, risk_code, evaluated_at, expires_at) " "VALUES (:c, 'C3', :t, :exp)" ), {"c": CUSTOMER, "t": base - timedelta(days=30), "exp": base + timedelta(days=300)}, ) for pid, name, ptype, rate in [ (PROD_OUT, "T12转出基金", "bond", OUT_RATE), (PROD_IN, "T12转入基金", "stock", IN_RATE), ]: conn.execute( text( "INSERT INTO core_product (product_id, product_name, min_risk_code, " "product_type, can_subscribe, can_redeem, subscribe_fee_rate, " "min_redeem_qty, min_hold_qty, min_hold_action, fund_company, ta_code) " "VALUES (:p, :n, 'R2', :t, 1, 1, :r, 0, 0, 'force_transfer', :co, :ta)" ), {"p": pid, "n": name, "t": ptype, "r": str(rate), "co": COMPANY, "ta": TA}, ) for mh, mh_max, rate in FEE_TIERS: conn.execute( text( "INSERT INTO core_fee_rule (product_id, fee_type, min_hold_days, " "max_hold_days, rate) VALUES (:p, 'redeem', :mh, :mh_max, :rate)" ), {"p": PROD_OUT, "mh": mh, "mh_max": mh_max, "rate": rate}, ) for pid, nav in ((PROD_OUT, OUT_NAV), (PROD_IN, IN_NAV)): conn.execute( text( "INSERT INTO core_product_nav (product_id, nav, daily_chg_pct, nav_date) " "VALUES (:p, :nav, 0, :d)" ), {"p": pid, "nav": str(nav), "d": accept_date}, ) for lot_id, qty, days in LOT_SPEC: conn.execute( text( "INSERT INTO core_share_lot (lot_id, customer_id, product_id, qty, " "remain_qty, nav, confirmed_at) VALUES (:l, :c, :p, :q, :q, :nav, :cat)" ), { "l": lot_id, "c": CUSTOMER, "p": PROD_OUT, "q": qty, "nav": str(OUT_NAV), # 持有天数按自然日差推算(hold_days 纯函数口径);档位边界宽裕, # 100/3 天无论交易日口径差异都稳定落在 0.0050 / 0.0150 档 "cat": base - timedelta(days=days), }, ) conn.execute( text( "INSERT INTO core_holding (customer_id, product_id, qty, cost_amount, " "market_value, pnl_pct, as_of) VALUES (:c, :p, :q, :cost, :mv, 0, :d)" ), { "c": CUSTOMER, "p": PROD_OUT, "q": str(TOTAL_QTY), "cost": str(TOTAL_QTY * OUT_NAV), "mv": str(TOTAL_QTY * OUT_NAV), "d": accept_date, }, ) def _boom_engine(out_trade, in_trade): raise RuntimeError("T12C 故意:确认段规则引擎异常") # ── 断言主体 ──────────────────────────────────────────────────────── def build_pending_state(core, repo, crepo, rrepo, *, now_submit, accept_date, t1) -> str: """[A] 真实链路构造待补偿态:受理成功 → T+1 确认(第⑦步引擎炸 + 第⑧步镜像炸)。 Core 三件套(受理单 confirmed + 两条流水 + 计费明细)**同事务真实落库**; agent 侧镜像停 pending、预警为 0 —— 补偿的两个补写对象都缺位。 """ print("\n[A] 待补偿态构造:accept_convert → confirm_one(引擎 + 镜像双失败)") req = { "customer_id": CUSTOMER, "from_product_id": PROD_OUT, "to_product_id": PROD_IN, "qty": QTY, "client_request_id": CID_REQ, } resp = accept_convert( req, core_ro=core, risk_repo=repo, convert_repo=crepo, request_repo=rrepo, now=now_submit, id_factory=_t12c_id, ) check("受理成功(accepted,未触发强制全转)", resp.get("accepted"), True) gid = resp["convert_group_id"] check( "受理单落库 status='accepted'", q1(core_engine, "SELECT status FROM core_convert_request WHERE convert_group_id = :g", g=gid), "accepted", ) check( "受理日 = 生产 resolve_accept_date 判定的 T", str(q1(core_engine, "SELECT DATE(requested_at) FROM core_convert_request WHERE convert_group_id = :g", g=gid)), str(accept_date), ) check( "受理镜像 pending(第⑦步受理版正常写入)", q1(agent_engine, "SELECT status FROM risk_convert_detail WHERE convert_group_id = :g", g=gid), "pending", ) # patch 窗口内确认:第⑧步镜像写入全部失败(Core 已 confirmed,不回滚); # 第⑦步引擎 hook 抛异常(engine_error 审计 + confirmed 审计标记 True)。 with patch.object( ConvertRepository, "sync_mirror", side_effect=RuntimeError("T12C 故意:确认镜像写入失败") ): out = confirm_one( gid, core_ro=core, risk_repo=repo, convert_repo=crepo, request_repo=rrepo, thresholds=TH, now=datetime.combine(t1, time(10, 0)), as_of=t1, engine_hook=_boom_engine, id_factory=_t12c_id, ) check("确认返回 status=confirmed(Core 不受 agent 故障影响)", out["status"], "confirmed") check("确认响应标记引擎异常(engine_error=True,T-12 三链路对齐)", out["engine_error"], True) trades = rows(core_engine, "SELECT * FROM core_trade WHERE convert_group_id = :g", g=gid) check("Core 流水两条(交易不可回滚)", len(trades), 2) check("流水类型 redeem + subscribe", sorted(t["trade_type"] for t in trades), ["redeem", "subscribe"]) lot_details = rows( core_engine, "SELECT * FROM core_convert_lot_detail WHERE convert_group_id = :g ORDER BY hold_days DESC", g=gid, ) check("计费明细两批(FIFO 跨 A1/A2)", len(lot_details), 2) check( "逐批净值 = T 日净值(未知价法)", {str(r["nav"]) for r in lot_details}, {str(OUT_NAV)}, ) check( "逐批净值日 = 受理日 T", {str(r["nav_date"]) for r in lot_details}, {str(accept_date)}, ) check( "第⑧步失败:镜像停 pending、流水号未回填", rows(agent_engine, "SELECT status, out_trade_id, in_trade_id FROM risk_convert_detail " "WHERE convert_group_id = :g", g=gid), [{"status": "pending", "out_trade_id": None, "in_trade_id": None}], ) check( "阶段⑦引擎异常审计已落(decision='engine_error',主链路留痕)", q1(agent_engine, "SELECT COUNT(*) FROM audit_log WHERE decision = 'engine_error' " "AND trace_id LIKE 'T12C-VERIFY%'"), 1, ) check( "此时本客户无预警单(引擎中断,预警缺位)", q1(agent_engine, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", c=CUSTOMER), 0, ) return gid def check_cleanup(core, crepo, rrepo, *, now: datetime) -> None: """[B] SLA 清理:超时在途受理单 → expired(ENUM 值域 / 窗口同口径 / 幂等 / 不硬删 / 镜像不推)。""" print("\n[B] cleanup:超 SLA 在途受理单 → expired(真库日历链式回推,不猜自然日)") # ── B0 窗口同口径:cleanup 的 SLA 边界 == confirm_batch 捞单窗口下界 ── # 当前业务日 = 今天开市取今天,否则回退最近交易日(与 _slab_cutoff 第一步一致)。 today = now.date() cur_biz = today if is_biz_day(today, core.is_open) else previous_biz_day(today, core.is_open) sla = settings.convert_confirm_sla_days # ── B1 造数:老 accepted / 老 nav_pending / 新 accepted;受理日按真库日历回推 ── old_date = today for _ in range(6): # 6 个交易日前(远超 SLA=2),跨周末仍稳 old_date = previous_biz_day(old_date, core.is_open) old_requested = datetime.combine(old_date, time(10, 0)) g_old, g_oldnav, g_new = (_t12c_id("CNV", now) for _ in range(3)) rrepo.insert(g_old, CUSTOMER, PROD_OUT, PROD_IN, Decimal("1000"), client_request_id="T12C-REQ-OLD", requested_at=old_requested) rrepo.insert(g_oldnav, CUSTOMER, PROD_OUT, PROD_IN, Decimal("2000"), client_request_id="T12C-REQ-OLDNAV", requested_at=old_requested) check( "条件 UPDATE 守卫:accepted → nav_pending 迁移成功", rrepo.transition_status(g_oldnav, "accepted", STATUS_NAV_PENDING), True, ) rrepo.insert(g_new, CUSTOMER, PROD_OUT, PROD_IN, Decimal("3000"), client_request_id="T12C-REQ-NEW", requested_at=now) crepo.insert_placeholder(g_old, "T12C-REQ-OLD") # 镜像行:cleanup 不得触碰 # ── B2 dry-run:只报告候选,不写库 ── dry = sla_cleanup(rrepo, core, sla, now=now, dry_run=True) cand = {c["convert_group_id"] for c in dry["candidates"]} check("dry-run 候选含两张老单", {g_old, g_oldnav} <= cand, True) check("dry-run 候选不含新单(requested_at 恒 ≥ cutoff)", g_new in cand, False) check( "dry-run 不写库(OLD 仍 accepted)", q1(core_engine, "SELECT status FROM core_convert_request WHERE convert_group_id = :g", g=g_old), "accepted", ) check( "SLA 边界与确认捞单窗口同口径(同一公式对表 _window_start)", dry["cutoff_requested_at"], str(_window_start(core, cur_biz)), ) # ── B3 真跑:老单 → expired ── summary = sla_cleanup(rrepo, core, sla, now=now) check("真跑 expired 恰两张老单", sorted(summary["expired"]), sorted([g_old, g_oldnav])) check( "真库 6 态 ENUM 接受 status='expired'", q1(core_engine, "SELECT status FROM core_convert_request WHERE convert_group_id = :g", g=g_old), "expired", ) check( "nav_pending 老单同收口", q1(core_engine, "SELECT status FROM core_convert_request WHERE convert_group_id = :g", g=g_oldnav), "expired", ) check( "新单原样(accepted,不动)", q1(core_engine, "SELECT status FROM core_convert_request WHERE convert_group_id = :g", g=g_new), "accepted", ) check( "S2 标记不硬删:受理单行仍在", q1(core_engine, "SELECT COUNT(*) FROM core_convert_request WHERE convert_group_id = :g", g=g_old), 1, ) check( "镜像不推 expired(T-4 契约:镜像不写 cancelled/expired)", q1(agent_engine, "SELECT status FROM risk_convert_detail WHERE convert_group_id = :g", g=g_old), "pending", ) # ── B4 复跑幂等:终态行不再进候选 ── again = sla_cleanup(rrepo, core, sla, now=now) again_cand = {c["convert_group_id"] for c in again["candidates"]} check("复跑老单不再进候选(终态)", g_old in again_cand, False) check("复跑 expired 计数为 0(本前缀)", [g for g in again["expired"] if g.startswith("CNV-T12C")], []) def check_compensate(gid: str, core, repo, crepo, *, accept_date: date, now: datetime) -> None: """C/D/E 组:补偿端到端 + 幂等 + 真 JSON 列 LIKE。""" print("\n[C] compensate_convert 端到端(详情 + 预警一次补齐)") out = compensate_convert( gid, core_ro=core, risk_repo=repo, convert_repo=crepo, thresholds=TH, now=now ) check("state", out["state"], "rebuilt") check("detail 分支", out["detail"], "rebuilt") check("engine 分支", out["engine"], "rebuilt") check("补偿后出一张单", len(out["alert_ids"]), 1) check("命中 RISK-001(大额)", "RISK-001" in out["triggered_rules"], True) detail = rows( agent_engine, "SELECT * FROM risk_convert_detail WHERE convert_group_id = :g", g=gid )[0] check("详情置 completed(T-4 契约:补偿 = 镜像唯一写入口)", detail["status"], "completed") check("回填 out_trade_id 与 Core 一致", detail["out_trade_id"], out["out_trade_id"]) check("回填 in_trade_id 与 Core 一致", detail["in_trade_id"], out["in_trade_id"]) check("回填 nav = T 日转出净值", str(detail["nav"]), str(OUT_NAV)) check("回填 nav_date = 受理日", str(detail["nav_date"]), str(accept_date)) check("回填 nav_stale = False(T+1 净值即 T 日,不存在过期)", detail["nav_stale"], 0) # 聚合口径与 Core 三件套对表(不手算 · 自检第 13 问) lot_details = rows( core_engine, "SELECT * FROM core_convert_lot_detail WHERE convert_group_id = :g", g=gid ) check( "fee_amount = Σ 逐批(Core 明细聚合)", str(detail["fee_amount"]), str(sum(Decimal(str(d["fee_amount"])) for d in lot_details)), ) check( "hold_days_min/max = Core 明细极值", (detail["hold_days_min"], detail["hold_days_max"]), (min(int(d["hold_days"]) for d in lot_details), max(int(d["hold_days"]) for d in lot_details)), ) confirmed_audits = rows( agent_engine, "SELECT input_summary FROM audit_log WHERE decision = 'confirmed' " "AND trace_id LIKE 'T12C-VERIFY%'", ) # T+1 真链路下恰 2 条:确认段第⑧步 1 条(phase='confirm')+ 补偿 1 条 # (phase='confirm-compensate',标注来源便于对账)—— v1.0 单测只看到补偿 # 那条(造数没跑 confirm_one),真库端到端才能验出两段各留各痕。 check("confirmed 审计恰 2 条(确认段 + 补偿各 1)", len(confirmed_audits), 2) phases = [] for a in confirmed_audits: raw = a["input_summary"] phases.append( json.loads(raw if isinstance(raw, str) else raw).get("phase") ) check("确认段审计 phase = confirm", phases.count("confirm"), 1) check("补偿审计 phase = confirm-compensate(来源可对账)", phases.count("confirm-compensate"), 1) confirm_summary = json.loads( next( a["input_summary"] for a in confirmed_audits if json.loads(a["input_summary"] if isinstance(a["input_summary"], str) else a["input_summary"]).get("phase") == "confirm" ) ) check("确认段审计 engine_error = True(第⑦步异常如实标记)", confirm_summary.get("engine_error"), True) alerts = rows(agent_engine, "SELECT * FROM risk_alert WHERE customer_id = :c", c=CUSTOMER) check("预警单 1 张", len(alerts), 1) payload = json.loads(alerts[0]["payload"]) if isinstance(alerts[0]["payload"], str) \ else alerts[0]["payload"] check("一张单两条事件(转出 + 转入,真 JSON 列反解)", len(payload["events"]), 2) print("\n[D] 重复补偿 → skipped(幂等,不产生第二张单)") again = compensate_convert( gid, core_ro=core, risk_repo=repo, convert_repo=crepo, thresholds=TH, now=now ) check("state", again["state"], "skipped") check("detail 分支", again["detail"], "already_completed") check("engine 分支", again["engine"], "skipped") check("alert_ids 与首次一致", again["alert_ids"], out["alert_ids"]) check( "预警单仍只 1 张", q1(agent_engine, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", c=CUSTOMER), 1, ) check( "补偿审计不重复(phase='confirm-compensate' 仍恰 1 条)", sum( 1 for a in rows( agent_engine, "SELECT input_summary FROM audit_log WHERE decision = 'confirmed' " "AND trace_id LIKE 'T12C-VERIFY%'", ) if json.loads( a["input_summary"] if isinstance(a["input_summary"], str) else a["input_summary"] ).get("phase") == "confirm-compensate" ), 1, ) check("幂等跳过时附人工核对提示(engine_error 中间态)", "人工核对" in (again.get("warning") or ""), True) print("\n[E] 真 JSON 列上 input_summary LIKE 生效(has_engine_error_audit)") col_type = q1( agent_engine, "SELECT DATA_TYPE FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = DATABASE() " "AND TABLE_NAME = 'audit_log' AND COLUMN_NAME = 'input_summary'", ) check("前提:input_summary 真库列类型为 json", str(col_type).lower(), "json") check( "按 group_id 命中(sqlite 是 TEXT,此断言只在真库有意义)", repo.has_engine_error_audit(gid, decision="engine_error"), True, ) check( "决策码不匹配则不命中(证明判定真的按 decision 过滤)", repo.has_engine_error_audit(gid, decision="risk_engine_error"), False, ) def main() -> int: global core_engine, agent_engine core_engine = get_engine(settings.mysql_core_database, "admin") agent_engine = get_engine(settings.mysql_database, "admin") print("T-12 真库验证:T+1 补偿(详情 + 预警)+ SLA 清理") print(f"客户={CUSTOMER} 产品={PROD_OUT}/{PROD_IN} SLA={settings.convert_confirm_sla_days} 交易日") cleanup_env(core_engine, agent_engine) new_trace(TRACE) # 固定 trace:审计行才能按 'T12C-VERIFY%' 精准清理与断言 try: # 业务路径用生产默认账号(ro 读 / 仓储默认 rw 写),仅建清数据用 admin core = CoreReadOnlyRepository(engine=get_engine(settings.mysql_core_database, "ro")) repo = RiskRepository(engine=get_engine(settings.mysql_database, "rw")) crepo = ConvertRepository(engine=get_engine(settings.mysql_database, "rw")) rrepo = ConvertRequestRepository(engine=get_engine(settings.mysql_core_database, "rw")) # 受理日 T / 确认日 T1:生产 trading_calendar 走真库日历实时推算(任意运行日稳) now_submit = datetime.now() cutoff = parse_cutoff(settings.convert_cutoff_time) accept_date = resolve_accept_date(now_submit, core.is_open, cutoff) t1 = next_biz_day(accept_date, core.is_open) print(f"受理日 T = {accept_date} / 确认日 T1 = {t1}(真实提交时刻 {now_submit:%H:%M:%S})") seed(core_engine, accept_date) gid = build_pending_state( core, repo, crepo, rrepo, now_submit=now_submit, accept_date=accept_date, t1=t1 ) check_cleanup(core, crepo, rrepo, now=datetime.now()) check_compensate( gid, core, repo, crepo, accept_date=accept_date, now=datetime.combine(t1, time(10, 0)) ) finally: cleanup_env(core_engine, agent_engine) left = q1( core_engine, "SELECT (SELECT COUNT(*) FROM core_trade WHERE customer_id LIKE 'CUST-T12C%') " "+ (SELECT COUNT(*) FROM core_convert_request WHERE customer_id = 'CUST-T12C') " "+ (SELECT COUNT(*) FROM core_convert_lot_detail WHERE convert_group_id LIKE 'CNV-T12C%') " "+ (SELECT COUNT(*) FROM core_share_lot WHERE customer_id LIKE 'CUST-T12C%') " "+ (SELECT COUNT(*) FROM core_holding WHERE customer_id LIKE 'CUST-T12C%') " "+ (SELECT COUNT(*) FROM core_product WHERE product_id LIKE 'PROD-T12C%') " "+ (SELECT COUNT(*) FROM core_customer WHERE customer_id = 'CUST-T12C')", ) + q1( agent_engine, "SELECT (SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id LIKE 'CNV-T12C%') " "+ (SELECT COUNT(*) FROM risk_alert WHERE customer_id = 'CUST-T12C') " "+ (SELECT COUNT(*) FROM audit_log WHERE trace_id LIKE 'T12C-VERIFY%')", ) print(f"\n清理后残留行数:{left}") dispose_engines() print(f"\n结果:{_passed} 通过 / {_failed} 失败") return 1 if _failed else 0 if __name__ == "__main__": raise SystemExit(main())