@@ -1,40 +1,45 @@
""" T-8 真 MySQL 验证脚本:`process_convert_event` 引擎改造在真库/真账号下 实跑 + DoD 断言。
""" T-8 真 MySQL 验证脚本:**引擎时机**(受理不触引擎 / 确认事务提交后恰跑一次) 实跑 + DoD 断言。
**为什么 sqlite 单测全绿还不够(T-8 视角)**
1. **「阶段 1.5 从跳过变生效 」是主流程行为变化 **: T-7 落地时 `process_convert_event` 尚不存在,
`convert_service._run_engine` 走 `ImportError` 分支**静默跳过**;T-8 落地后同一笔转换会
**真实出单 + 写 L3 + 落审计**。真库跑一遍才能确认没有连带破坏(审计条数、L3 写入、单落库) 。
1. **「引擎时机 」是主流程行为**: T+1 模型下引擎必须**只在确认事务提交后**跑一次(FR-C28)——
受理段触引擎会让「尚未确认的交易」提前进风控台账(验收 29 前半),
确认段重复触引擎会让 RISK-002 当日累计**翻倍**(验收 5/7)。真库跑一遍才能证明两条都成立 。
2. **`convert_group_id` 的 NULL 值域**: MySQL 里普通交易该列为 NULL、convert 两条为同一串值;
`amount_view` 的 `if not gid` 必须在真库值域下成立(NULL / 空串都不得误聚合)。
3. **DECIMAL 精度**: `core_trade.amount` 在 MySQL 是 DECIMAL,进 Python 后与生产纯函数
逐项比对(sqlite 用 REAL 无此保证,需 `float()` 绑定)。
4. **JSON 列反解**: `risk_alert.payload` 真库落库后反解出 的 `events` 长度必须为 2。
4. **JSON 列反解**: `risk_alert.payload` 真库落库后反解的 `events` 长度必须为 2。
**断言分组**
A 引擎真跑出单:1 张单 / `payload.events` 两条 / 顺序 [转出, 转入]
B **RISK-002 不翻倍**(阈值夹逼:单条 < 阈值 < 两条之和)
C 去重不删行:`core_trade` 仍 2 条且同 `convert_group_id`
D 幂等重试(同 `client_request_id`) **不产生第二张单**
E 无命中场景(小额)不建单、仅 pass 审计
F 清理后残留为零(自检)
A0 受理**不触引擎**:受理后 `risk_alert` 0 新增、无 `pass` / `engine_error` 审计(验收 29 前半)
A 确认后引擎真跑出单:1 张单 / `payload.events` 两条 / 顺序 [转出, 转入]
B **RISK-002 不翻倍**(阈值夹逼:单条 < 阈值 < 两条之和)
C 去重不删行:`core_trade` 仍 2 条且同 `convert_group_id`
D 幂等重放(重复确认)**不产生第二张单**
E 无命中场景(独立小额客户)不建单、仅 pass 审计
F 清理后残留为零(自检)
用法:
python scripts/dev/verify_convert_engine.py
约定(与 T-6/T-7 脚本一致):
- 用 **T8M 前缀**隔离数据(客户/产品/批次/group_id),跑完**两个库**( core + agent)全清;
- 建/清数据走 `role= " admin " `(需 DELETE);业务本身走 `convert_fund` 默认账号;
- 建/清数据走 `role= " admin " `(需 DELETE);业务链路走默认账号(确认段写 Core 由
`ConvertCoreRepository` 内部固定用 `xh_core_rw`, D20);
- 基准日取真库 `core_trade_calendar` 的**中位开市日**作 T 日(避开数据边界),
T+1 = 下一交易日;净值插在 T 日(确认段 `get_nav_on` 精确匹配,缺则 `nav_pending`);
- 批次 `confirmed_at` 相对 T 日计算(T−400 / T−3),使持有期档位与费额确定;
- 固定 `trace_id = " T8M-VERIFY-TRACE " `,便于精准清理 `audit_log`。
"""
from __future__ import annotations
import itertools
import json
import sys
import uuid
from datetime import date , datetime , timedelta
from decimal import ROUND_HALF_UP , Decimal
from datetime import date , datetime , time , timedelta
from decimal import Decimal
from pathlib import Path
from sqlalchemy import text
@@ -43,7 +48,13 @@ ROOT = Path(__file__).resolve().parents[2]
sys . path . insert ( 0 , str ( ROOT ) )
from app . config . settings import settings # noqa: E402
from app . gateway . convert_core_repository import ConvertCoreRepository # noqa: E402
from app . repository . convert_repository import ConvertRepository # noqa: E402
from app . repository . convert_request_repository import ( # noqa: E402
ConvertRequestRepository ,
)
from app . repository . core_ro import CoreReadOnlyRepository # noqa: E402
from app . repository . risk_repository import RiskRepository # noqa: E402
from app . service . convert . calc import ( # noqa: E402
convert_amount ,
diff_fee ,
@@ -53,21 +64,23 @@ from app.service.convert.calc import ( # noqa: E402
lot_fee ,
plan_lots ,
)
from app . service . convert . convert_service import convert_fund # noqa: E402
from app . service . convert . confirm_service import confirm_one # noqa: E402
from app . service . convert . convert_service import accept_convert # noqa: E402
from app . service . convert . fee import pick_fee_rate # noqa: E402
from app . service . convert . trading_calendar import next_biz_day # noqa: E402
from app . service . convert . types import FeeRule , Lot # noqa: E402
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
CUSTOMER = " CUST-T8M "
CUSTOMER_SMALL = " CUST-T8MS " # E 组:无命中场景专用(独立客户避免当日累计互相干扰)
PROD_OUT = " PROD-T8MO "
PROD_IN = " PROD-T8MI "
COMPANY = " 华夏模拟基金 "
TA = " TA-CN-001 "
TRADE_AT = datetime ( 2026 , 9 , 4 , 10 , 0 , 0 )
TRADE_DATE = date ( 2026 , 9 , 4 )
NOW = TRADE_AT
OUT_NAV = Decimal ( " 1.0000 " ) # T 日转出净值(未知价法 → 逐批金额按它计)
IN_NAV = Decimal ( " 0.9500 " )
OUT_RATE = Decimal ( " 0.0030 " )
IN_RATE = Decimal ( " 0.0080 " )
@@ -75,10 +88,17 @@ FEE_TIERS = [
( 0 , 7 , " 0.0150 " ) , ( 7 , 30 , " 0.0100 " ) , ( 30 , 180 , " 0.0050 " ) ,
( 180 , 365 , " 0.0025 " ) , ( 365 , None , " 0.0000 " ) ,
]
QTY_ALL = " 500000 "
CID_REQ = " T8M-IDEM-0001 "
TRACE = " T8M-VERIFY-TRACE "
# 夹逼阈值:单条转出 512000 < daily_total 800000 < 两条之和(约 1020000 )
#: 基准受理日(T 日)—— 运行时由真库日历在 `main` 内赋值(模块级无法静态确定 )
TA_DAY : date | None = None
# 夹逼阈值:单条转出 500000 < daily_total 800000 < 两条之和(约 996035)
# ⚠️ **必须显式传给 `confirm_one(thresholds=TH)`**:缺省会退回 `RiskThresholds.from_settings()`,
# settings 的 daily_total 一旦低于单条金额,RISK-002 就会「命中」—— 症状看着像去重失效,
# 实为阈值口径没传(首轮实跑即栽在此,B 组断言直接捕获)。
TH = RiskThresholds (
large_amount = Decimal ( " 100000 " ) ,
daily_total = Decimal ( " 800000 " ) ,
@@ -106,10 +126,13 @@ def check(name: str, actual, expected) -> None:
print ( f " { flag } { name } : 实际 { actual !r} " + ( " " if ok else f " / 期望 { expected !r} " ) )
_SEQ = itertools . count ( 1 )
def _t8m_id ( prefix : str , now : datetime ) - > str :
""" 固定前缀的 id 工厂:`convert_fun d` 默认 生成 `CNV-<日期>-<uuid>`, 与清理口径
`LIKE ' CNV-T8M % ' ` 不匹配 → 会在重跑时留下 撞唯一键的残留 行(T-7 脚本踩过同款坑 )。 """
return f " { prefix } -T8M- { uuid . uuid4 ( ) . hex [ : 8 ] . upper ( ) } "
""" 固定前缀的 id 工厂:生产 `_new_i d` 生成 `CNV-<日期>-<uuid>`,
与清理口径 `LIKE ' CNV-T8M % ' ` 不匹配 → 重跑会留 撞唯一键的残行(T-6/T-7 均踩过 )。 """
return f " CNV-T8M- { next ( _SEQ ) : 03d } "
def q1 ( engine , sql : str , * * params ) :
@@ -128,14 +151,21 @@ def cleanup(core_engine, agent_engine) -> None:
with core_engine . begin ( ) as conn :
for sql , params in [
( " DELETE FROM core_convert_lot_detail WHERE convert_group_id LIKE ' CNV-T8M % ' " , { } ) ,
( " 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_convert_request WHERE customer_id IN (:c1, :c2) " ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ) ,
( " DELETE FROM core_trade WHERE customer_id IN (:c1, :c2) " ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ) ,
( " DELETE FROM core_share_lot WHERE customer_id IN (:c1, :c2) " ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ) ,
( " DELETE FROM core_holding WHERE customer_id IN (:c1, :c2) " ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ) ,
( " DELETE FROM core_customer_risk WHERE customer_id IN (:c1, :c2) " ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ) ,
( " DELETE FROM core_fee_rule WHERE product_id LIKE ' PROD-T8M % ' " , { } ) ,
( " DELETE FROM core_product_nav WHERE product_id LIKE ' PROD-T8M % ' " , { } ) ,
( " DELETE FROM core_product WHERE product_id LIKE ' PROD-T8M % ' " , { } ) ,
( " DELETE FROM core_customer WHERE customer_id = :c " , { " c " : CUSTOMER } ) ,
( " DELETE FROM core_customer WHERE customer_id IN (:c1, :c2) " ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ) ,
] :
conn . execute ( text ( sql ) , params )
with agent_engine . begin ( ) as conn :
@@ -145,32 +175,37 @@ def cleanup(core_engine, agent_engine) -> None:
" WHERE convert_group_id LIKE ' CNV-T8M % ' OR client_request_id LIKE ' T8M- % ' "
)
)
conn . execute ( text ( " DELETE FROM risk_alert WHERE customer_id = :c " ) , { " c " : CUSTOMER } )
conn . execute (
text ( " DELETE FROM customer_profile_l3 WHERE customer_id = :c " ) , { " c " : CUSTOMER }
text ( " DELETE FROM risk_alert WHERE customer_id IN (:c1, :c2) " ) ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ,
)
conn . execute (
text ( " DELETE FROM audit_log WHERE trace_id LIKE ' T8M-VERIFY % ' " )
text ( " DELETE FROM customer_profile_l3 WHERE customer_id IN (:c1, :c2) " ) ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ,
)
conn . execute ( text ( " DELETE FROM audit_log WHERE trace_id LIKE ' T8M-VERIFY % ' " ) )
def seed ( core_engine ) - > None :
def seed ( core_engine , base_day : date , submit_at : datetime ) - > None :
""" 两个客户 + 两只产品 + 费率五档 + T 日净值 + 批次(T−400 / T−3)。 """
with core_engine . begin ( ) as conn :
conn . execute (
text (
" INSERT INTO core_customer (customer_id, display_name, open_date) "
" VALUES (:c, ' T8真库验证 ' , :d) "
) ,
{ " c " : CUSTOMER , " d " : TRADE_DATE } ,
)
conn . execute (
text (
" INSERT INTO core_customer_risk (customer_id, risk_code, evaluated_at, expires_at) "
" VALUES (:c, ' C3 ' , :t, :exp) "
) ,
{ " c " : CUSTOMER , " t " : TRADE_AT - timedelta ( days = 30 ) ,
" exp " : TRADE_AT + timedelta ( days = 300 ) } ,
)
for cid , name in ( ( CUSTOMER , " T8引擎时机 " ) , ( CUSTOMER_SMALL , " T8小额客户 " ) ) :
conn . execute (
text (
" INSERT INTO core_customer (customer_id, display_name, open_date) "
" VALUES (:c, :n, :d) "
) ,
{ " c " : cid , " n " : name , " d " : base_day - timedelta ( days = 500 ) } ,
)
conn . execute (
text (
" INSERT INTO core_customer_risk "
" (customer_id, risk_code, evaluated_at, expires_at) "
" VALUES (:c, ' C3 ' , :t, :exp) "
) ,
{ " c " : cid , " t " : submit_at - timedelta ( days = 30 ) ,
" exp " : submit_at + timedelta ( days = 300 ) } ,
)
for pid , name , ptype , rate in [
( PROD_OUT , " T8转出基金 " , " bond " , OUT_RATE ) ,
( PROD_IN , " T8转入基金 " , " stock " , IN_RATE ) ,
@@ -179,8 +214,8 @@ def seed(core_engine) -> None:
text (
" INSERT INTO core_product (product_id, product_name, min_risk_code, "
" product_type, can_subscribe, can_redeem, subscribe_fee_rate, "
" fund_company, ta_code) "
" VALUES (:p, :n, ' R2 ' , :t, 1, 1, :r, :co, :ta) "
" min_hold_qty, min_hold_action, fund_company, ta_code) "
" VALUES (:p, :n, ' R2 ' , :t, 1, 1, :r, 0, ' force_transfer ' , :co, :ta)"
) ,
{ " p " : pid , " n " : name , " t " : ptype , " r " : str ( rate ) , " co " : COMPANY , " ta " : TA } ,
)
@@ -192,45 +227,63 @@ def seed(core_engine) -> None:
) ,
{ " p " : PROD_OUT , " mh " : mh , " mh_max " : mh_max , " rate " : rate } ,
)
conn . execute (
text (
" INSERT INTO core_product_nav (product_id, nav, daily_chg_pct, nav_date) "
" VALUES (:p, :nav, 0, :d) "
) ,
{ " p " : PROD_IN , " nav " : str ( IN_NAV ) , " d " : TRADE_DATE } ,
)
# 两批:400000 份(持有 ~400 天,费率 0)+ 100000 份(持有 3 天,费率 1.5%)
for lot_id , qty , nav , confirmed in [
( " LOT-T8M-A1 " , " 400000 " , " 1.0300 " , datetime ( 2025 , 8 , 1 , 10 , 0 , 0 ) ) ,
( " LOT-T8M-A2 " , " 100000 " , " 1.0000 " , datetime ( 2026 , 9 , 1 , 10 , 0 , 0 ) ) ,
] :
# T 日净值:两端都必须有(确认段 `get_nav_on` 精确匹配)
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 " : base_day } ,
)
# 主客户两批:40 万份(持 400 天,费率 0)+ 10 万份(持 3 天,费率 1.5%)
for lot_id , qty , confirmed_delta in (
( " LOT-T8M-A1 " , " 400000 " , 400 ) ,
( " LOT-T8M-A2 " , " 100000 " , 3 ) ,
) :
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) "
" remain_qty, nav, confirmed_at) VALUES (:l, :c, :p, :q, :q, 1.0300 , :cat) "
) ,
{ " l " : lot_id , " c " : CUSTOMER , " p " : PROD_OUT , " q " : qty , " nav " : nav , " cat " : confirmed } ,
{ " l " : lot_id , " c " : CUSTOMER , " p " : PROD_OUT , " q " : qty ,
" cat " : datetime . combine ( base_day - timedelta ( days = confirmed_delta ) , time ( 10 , 0 ) ) } ,
)
# 小额客户一批(1000 份,持 400 天 → 费率 0)
conn . execute (
text (
" INSERT INTO core_holding ( customer_id, product_id, qty, cost_amount, "
" market_value, pnl_pct, as_of) VALUES (:c, :p, 500000, 512000.00, 512000.00, 0, :d) "
" INSERT INTO core_share_lot (lot_id, customer_id, product_id, qty, "
" remain_qty, nav, confirmed_at) VALUES ( ' LOT-T8M-S1 ' , :c, :p, 1000, 1000, "
" 1.0300, :cat) "
) ,
{ " c " : CUSTOMER , " p " : PROD_OUT , " d " : TRADE_DATE } ,
{ " c " : CUSTOMER_SMALL , " p " : PROD_OUT ,
" cat " : datetime . combine ( base_day - timedelta ( days = 400 ) , time ( 10 , 0 ) ) } ,
)
for cid , qty in ( ( CUSTOMER , " 500000 " ) , ( CUSTOMER_SMALL , " 1000 " ) ) :
conn . execute (
text (
" INSERT INTO core_holding (customer_id, product_id, qty, cost_amount, "
" market_value, pnl_pct, as_of) VALUES (:c, :p, :q, :q, :q, 0, :d) "
) ,
{ " c " : cid , " p " : PROD_OUT , " q " : qty , " d " : base_day } ,
)
def expected_quote ( core_ro , requested : str ) - > dict :
""" 期望折算(全部走生产纯函数,与 service 内部 同口径)。 """
""" 期望折算(全部走生产纯函数,与 service 同口径)。
⚠️ T+1 变更点:逐批金额按 **T 日净值**(`OUT_NAV`)计,**不是**批次买入净值 ——
这正是「未知价法」与 v1.0 的实质差异,脚本不得沿用 `alloc.nav`。
"""
lots = [ Lot . from_row ( r ) for r in core_ro . list_share_lots ( CUSTOMER , PROD_OUT ) ]
fee_rules = [ FeeRule . from_row ( r ) for r in core_ro . get_redeem_fee_rules ( PROD_OUT ) ]
plan = plan_lots ( lots , Decimal ( requested ) )
out_amount = Decimal ( " 0 " )
redeem_fee = Decimal ( " 0 " )
for alloc in plan . allocations :
amount = lot_amount ( alloc . qty , alloc . nav )
amount = lot_amount ( alloc . qty , OUT_NAV )
rate = pick_fee_rate (
fee_rules , hold_days ( TRADE_DATE , alloc . confirmed_at ) , product_id = PROD_OUT
fee_rules , hold_days ( TA_DAY , alloc . confirmed_at ) , product_id = PROD_OUT
)
out_amount + = amount
redeem_fee + = lot_fee ( amount , rate )
@@ -242,168 +295,212 @@ def expected_quote(core_ro, requested: str) -> dict:
" convert_amount " : conv ,
" diff_fee " : gap ,
" in_amount " : conv - gap ,
" in_qty " : in_qty ( conv - gap , IN_NAV ) ,
}
def _accept_services ( ) - > dict :
return dict (
core_ro = CoreReadOnlyRepository ( ) ,
risk_repo = RiskRepository ( ) ,
convert_repo = ConvertRepository ( ) ,
request_repo = ConvertRequestRepository ( ) ,
)
def _confirm_services ( ) - > dict :
return dict ( * * _accept_services ( ) , core_writer = ConvertCoreRepository ( ) )
def _pick_days ( core_admin , is_open ) - > tuple [ date , date ] :
""" 取日历**中位**开市日作 T 日,及其下一交易日(T+1 确认业务日)。
⚠️ 不能用「最新开市日」:日历种子到 2027 年底为止,取末日会让 T+1 撞
`next_biz_day` 的数据边界(T-6 实测暴露过该缺陷)。
"""
rws = rows (
core_admin ,
" SELECT cal_date FROM core_trade_calendar WHERE is_open = 1 ORDER BY cal_date ASC " ,
)
if len ( rws ) < 10 :
raise RuntimeError ( " core_trade_calendar 开市日不足,请先跑 10-seed-trade-calendar.sql " )
day = rws [ len ( rws ) / / 2 ] [ " cal_date " ]
if not isinstance ( day , date ) :
day = date . fromisoformat ( str ( day ) [ : 10 ] )
return day , next_biz_day ( day , is_open )
def main ( ) - > int :
global TA_DAY
core_engine = get_engine ( " jinrong_core " , role = " admin " )
agent_engine = get_engine ( " jinrong_agent " , role = " admin " )
core_ro = CoreReadOnlyRepository ( )
print ( " = " * 60 )
print ( " T-8 真库验证:规则引擎改造(amount_view 去重 + process_convert_event 出单 ) " )
print ( " T-8 真库验证:引擎时机(受理不触引擎 / 确认后恰跑一次 ) " )
print ( " = " * 60 )
try :
ta_day , t1_day = _pick_days ( core_engine , core_ro . is_open )
except Exception as exc : # noqa: BLE001
print ( f " ❌ 环境不可用: { exc } " )
dispose_engines ( )
return 2
TA_DAY = ta_day
submit_at = datetime . combine ( ta_day , time ( 10 , 0 ) )
confirm_at = datetime . combine ( t1_day , time ( 9 , 0 ) )
print ( f " T 日= { ta_day } (受理)→ T+1= { t1_day } (确认业务日) " )
cleanup ( core_engine , agent_engine )
seed ( core_engine )
core_ro = CoreReadOnlyRepository ( )
req_qty = " 500000 " # 全部转出
exp = expected_quote ( core_ro , req_qty )
# ── A/B/C/D:带幂等键的大额转换 ─────────────────────────────────
print ( " \n 【A】引擎真跑(阶段 1.5 不再跳过):一张单 + payload.events 两条 " )
seed ( core_engine , ta_day , submit_at )
new_trace ( TRACE )
resp = convert_fund (
exp = expected_quote ( core_ro , QTY_ALL )
# ── A0:受理**不触引擎**(FR-C28 · 验收 29 前半)─────────────────
print ( " \n 【A0】受理不触引擎:risk_alert 0 新增、无 pass / engine_error 审计 " )
accepted = accept_convert (
{
" customer_id " : CUSTOMER ,
" from_product_id " : PROD_OUT ,
" to_product_id " : PROD_IN ,
" qty " : req_qty ,
" qty " : Decimal ( QTY_ALL ) ,
" client_request_id " : CID_REQ ,
} ,
now = NOW ,
thresholds = TH ,
now = submit_at ,
id_factory = _t8m_id ,
* * _accept_services ( ) ,
)
check ( " 转换未被阻断 " , resp [ " blocked " ] , False )
check ( " 引擎异常标记为 False(说明真的跑了且没炸 ) " , resp [ " engine_error " ] , False )
check ( " 折算 out_amount 与生产纯函数一致 " , Decimal ( resp [ " out_amount " ] ) , exp [ " out_amount " ] )
check ( " 折算 in_amount 与生产纯函数一致 " , Decimal ( resp [ " in_amount " ] ) , exp [ " in_amount " ] )
check ( " RISK-001 命中(转出端 512000 ≥ 100000) " , " RISK-001 " in resp [ " triggered_rules " ] , True )
check ( " 只有一张预警单 " , len ( resp [ " alert_ids " ] ) , 1 )
gid = accepted [ " convert_group_id " ]
check ( " 受理成功(accepted ) " , accepted [ " status " ] , " accepted " )
check ( " 受理后 risk_alert 0 行 " ,
q1 ( agent_engine , " SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c " ,
c = CUSTOMER ) , 0 )
check ( " 受理后无 pass 审计 " ,
q1 ( agent_engine , " SELECT COUNT(*) FROM audit_log WHERE trace_id = :t "
" AND decision = ' pass ' " , t = TRACE ) , 0 )
check ( " 受理后无 engine_error 审计 " ,
q1 ( agent_engine , " SELECT COUNT(*) FROM audit_log WHERE trace_id = :t "
" AND decision = ' engine_error ' " , t = TRACE ) , 0 )
check ( " 受理后 core_trade 0 行 " ,
q1 ( core_engine , " SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g " ,
g = gid ) , 0 )
alerts = rows (
agent_engine , " SELECT * FROM risk_alert WHERE customer_id = :c " , c = CUSTOMER
)
# ── A:确认后引擎真跑出单 ────────────────────────────────────────
print ( " \n 【A】确认后引擎真跑:一张单 + payload.events 两条 " )
res = confirm_one ( gid , now = confirm_at , as_of = t1_day , thresholds = TH , * * _confirm_services ( ) )
check ( " 确认状态 confirmed " , res [ " status " ] , " confirmed " )
check ( " 折算 out_amount 与生产纯函数一致 " , Decimal ( res [ " out_amount " ] ) , exp [ " out_amount " ] )
check ( " 折算 in_amount 与生产纯函数一致 " , Decimal ( res [ " in_amount " ] ) , exp [ " in_amount " ] )
check ( " 折算 in_qty 与生产纯函数一致 " , Decimal ( res [ " in_qty " ] ) , exp [ " in_qty " ] )
check ( " RISK-001 命中(转出端 500000 ≥ 100000) " , " RISK-001 " in res [ " triggered_rules " ] , True )
check ( " 只出一张单 " , len ( res [ " alert_ids " ] ) , 1 )
alerts = rows ( agent_engine , " SELECT * FROM risk_alert WHERE customer_id = :c " , c = CUSTOMER )
check ( " risk_alert 落库 1 行 " , len ( alerts ) , 1 )
payload = alerts [ 0 ] [ " payload " ]
payload = json . loads ( payload ) if isinstance ( payload , str ) else payload
check ( " payload.events 长度 2 " , len ( payload [ " events " ] ) , 2 )
check (
" events 顺序 [转出, 转入] " ,
[ e [ " trade_type " ] for e in payload [ " events " ] ] ,
[ " redeem " , " subscribe " ] ,
)
check ( " events[0].trade_id 为转出端 " , payload [ " events " ] [ 0 ] [ " trade_id " ] , resp [ " out_trade_id " ] )
check ( " events[1].trade_id 为转入端 " , payload [ " events " ] [ 1 ] [ " trade_id " ] , resp [ " in_trade_id " ] )
check ( " events 顺序 [转出, 转入] " ,
[ e [ " trade_type " ] for e in payload [ " events " ] ] , [ " redeem " , " subscribe " ] )
check ( " events[0].trade_id 为转出端 " , payload [ " events " ] [ 0 ] [ " trade_id " ] , res [ " out_trade_id " ] )
check ( " events[1].trade_id 为转入端 " , payload [ " events " ] [ 1 ] [ " trade_id " ] , res [ " in_trade_id " ] )
check ( " 主流水口径取转出端 " , alerts [ 0 ] [ " trade_id " ] , res [ " out_trade_id " ] )
print ( " \n 【B】RISK-002 不翻倍(单条 512000 < 800000 < 两条之和) " )
check ( " triggered_rules 不含 RISK-002 " , " RISK-002 " in resp [ " triggered_rules " ] , False )
# 同日两条流水之和确实超过阈值 → 证明「不命中」来自去重而非阈值太松
# ── B:RISK-002 不翻倍(阈值夹逼)────────────────────────────────
print ( " \n 【B】RISK-002 不翻倍(单条 500000 < 800000 < 两条之和) " )
check ( " triggered_rules 不含 RISK-002 " , " RISK-002 " in res [ " triggered_rules " ] , False )
total_two = exp [ " out_amount " ] + exp [ " in_amount " ]
check ( " 两条之和确实 ≥ 阈值(阈值 夹逼成立) " , total_two > = TH . daily_total , True )
check ( " 两条之和确实 ≥ 阈值(夹逼成立) " , total_two > = TH . daily_total , True )
print ( f " └ 转出 { exp [ ' out_amount ' ] } + 转入 { exp [ ' in_amount ' ] } = { total_two } " )
# ── C:去重不删行 ────────────────────────────────────────────────
print ( " \n 【C】去重不删行:core_trade 仍 2 条且同组 " )
gid = resp [ " convert_group_id " ]
trades = rows (
core_engine ,
" SELECT * FROM core_trade WHERE convert_group_id = :g " ,
g = gid ,
)
trades = rows ( core_engine , " SELECT * FROM core_trade WHERE convert_group_id = :g " , g = gid )
by_type = { t [ " trade_type " ] : t for t in trades }
check ( " core_trade 2 条 " , len ( trades ) , 2 )
check ( " 类型覆盖 [redeem, subscribe] " , sorted ( by_type ) , [ " redeem " , " subscribe " ] )
check ( " 两条同 convert_group_id " , { t [ " convert_group_id " ] for t in trades } , { gid } )
check (
" 转出端金额精度零漂移(DECIMAL 18,2) " ,
Decimal ( str ( by_type [ " redeem " ] [ " amount " ] ) ) ,
exp [ " out_amount " ] ,
)
check (
" 转入端金额精度零漂移(DECIMAL 18,2) " ,
Decimal ( str ( by_type [ " subscribe " ] [ " amount " ] ) ) ,
exp [ " in_amount " ] ,
)
check ( " 转出端金额精度零漂移(DECIMAL 18,2) " ,
Decimal ( str ( by_type [ " redeem " ] [ " amount " ] ) ) , exp [ " out_amount " ] )
check ( " 转入端金额精度零漂移(DECIMAL 18,2) " ,
Decimal ( str ( by_type [ " subscribe " ] [ " amount " ] ) ) , exp [ " in_amount " ] )
print ( " \n 【D】幂等重试(同 client_request_id)不产生第二张单 " )
resp2 = convert_fund (
{
" customer_id " : CUSTOMER ,
" from_product_id " : PROD_OUT ,
" to_product_id " : PROD_IN ,
" qty " : req_qty ,
" client_request_id " : CID_REQ ,
} ,
now = NOW + timedelta ( minutes = 1 ) ,
thresholds = TH ,
id_factory = _t8m_id ,
# ── D:幂等重放(重复确认)不产生第二张单 ────────────────────────
print ( " \n 【D】幂等重放:重复确认 skipped,不产生第二张单 " )
again = confirm_one (
gid , now = confirm_at , as_of = t1_day , thresholds = TH , * * _confirm_services ( )
)
check ( " 重试 返回同一 group_i d " , resp2 [ " convert_group_id " ] , gid )
check ( " 重试 core_trade 仍 2 条 " , q1 (
core_engine , " SELECT COUNT(*) FROM core_trade WHERE customer _id = :c " , c = CUSTOMER ) , 2 )
check ( " 重试后预警单仍 1 张 " , q1 (
agent_engine , " SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c " , c = CUSTOMER ) , 1 )
check ( " 重放 返回 skippe d " , again [ " status " ] , " skipped " )
check ( " 重放 core_trade 仍 2 条 " ,
q1 ( core_engine , " SELECT COUNT(*) FROM core_trade WHERE convert_group _id = :g " ,
g = gid ) , 2 )
check ( " 重放后预警单仍 1 张 " ,
q1 ( agent_engine , " SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c " ,
c = CUSTOMER ) , 1 )
check ( " 重放后 lot 明细仍 2 条 " ,
q1 ( core_engine , " SELECT COUNT(*) FROM core_convert_lot_detail "
" WHERE convert_group_id = :g " , g = gid ) , 2 )
# ── E:无命中场景 ────────────────────────────────────────────────
print ( " \n 【E】无命中场景(换日 → 当日无其他流水 )不建单、仅 pass 审计 " )
before_alerts = q1 (
agent_engine , " SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c " , c = CUSTOMER )
before_pass = q1 (
agent_engine ,
" SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = ' pass ' " , t = TRACE )
# 该客户份额已全转出 → 4xx 不落库;故先补一批小额份额再造一笔
with core_engine . begin ( ) as conn :
conn . execute (
text (
" INSERT INTO core_share_lot (lot_id, customer_id, product_id, qty, remain_qty, "
" nav, confirmed_at) VALUES ( ' LOT-T8M-B1 ' , :c, :p, 1000, 1000, 1.0000, :cat) "
) ,
{ " c " : CUSTOMER , " p " : PROD_OUT , " cat " : datetime ( 2025 , 8 , 1 , 10 , 0 , 0 ) } ,
)
# 换到次日:`list_trades_range` 按 [日初, 次日) 取数,前一日的大额流水不再进本批规则输入
resp3 = convert_fund (
print ( " \n 【E】无命中场景(独立小额客户 )不建单、仅 pass 审计 " )
small = accept_convert (
{
" customer_id " : CUSTOMER ,
" customer_id " : CUSTOMER_SMALL ,
" from_product_id " : PROD_OUT ,
" to_product_id " : PROD_IN ,
" qty " : " 1000 " ,
" qty " : Decimal ( " 1000 " ) ,
" client_request_id " : " T8M-IDEM-0002 " ,
} ,
now = NOW + timedelta ( days = 1 ) ,
thresholds = TH ,
now = submit_at ,
id_factory = _t8m_id ,
* * _accept_services ( ) ,
)
check ( " 小额转换成功(非 blocked) " , resp3 [ " blocked " ] , False )
check ( " 未命中任何规则 " , resp3 [ " triggered_rules " ] , [ ] )
check ( " 未新建预警单 " , q1 (
agent_engine , " SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c " , c = CUSTOMER ) ,
before_alerts )
check ( " 落 pass 审计 1 条 " , q1 (
agent_engine ,
" SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = ' pass ' " , t = TRACE )
- before_pass , 1 )
res3 = confirm_one (
small [ " convert_group_id " ] ,
now = confirm_at ,
as_of = t1_day ,
thresholds = TH ,
* * _confirm_services ( ) ,
)
check ( " 小额转换确认成功 " , res3 [ " status " ] , " confirmed " )
check ( " 未命中任何规则 " , res3 [ " triggered_rules " ] , [ ] )
check ( " 未新建预警单 " ,
q1 ( agent_engine , " SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c " ,
c = CUSTOMER_SMALL ) , 0 )
check ( " 落 pass 审计 1 条 " ,
q1 ( agent_engine , " SELECT COUNT(*) FROM audit_log WHERE trace_id = :t "
" AND decision = ' pass ' " , t = TRACE ) , 1 )
# ── F:清理自检 ─────────────────────────────────────────────────
print ( " \n 【F】清理 T8M 前缀数据并自检无残留 " )
cleanup ( core_engine , agent_engine )
check ( " core_trade 零残留 " , q1 (
core_engine , " SELECT COUNT(*) FROM core_trade WHERE customer_id = :c " , c = CUSTOMER ) , 0 )
check ( " core 客户零残留 " , q1 (
core_engine , " SELECT COUNT(*) FROM core_customer WHERE customer_id = :c " , c = CUSTOMER ) , 0 )
check ( " risk_alert 零残留 " , q1 (
agent_engine , " SELECT COUNT(*) FROM risk_aler t WHERE customer_id = :c " , c = CUSTOMER ) , 0 )
check ( " customer_profile_l3 零残留 " , q1 (
agent_engine ,
" SELECT COUNT(*) FROM customer_profile_l3 WHERE customer_id = :c " , c = CUSTOMER ) , 0 )
check ( " audit_log 零残留 " , q1 (
agent_engine , " SELECT COUNT(*) FROM audit_log WHERE trace_id LIKE ' T8M-VERIFY % ' " ) , 0 )
check ( " risk_convert_detail 零残留 " , q1 (
agent_engine ,
" SELECT COUNT(*) FROM risk_convert_detail "
" WHERE convert_group_id LIKE ' CNV-T8M % ' OR client_request_id LIKE ' T8M- % ' " ) , 0 )
for label , engine , sql , params in [
( " core_trade " , core_engine ,
" SELECT COUNT(*) FROM core_trade WHERE customer_id IN (:c1, :c2) " ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ) ,
( " core_convert_request " , core_engine ,
" SELECT COUNT(*) FROM core_convert_reques t WHERE customer_id IN (:c1, :c2) " ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ) ,
( " core 客户 " , core_engine ,
" SELECT COUNT(*) FROM core_ customer WHERE customer_id IN (:c1, :c2) " ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ) ,
( " risk_alert " , agent_engine ,
" SELECT COUNT(*) FROM risk_alert WHERE customer_id IN (:c1, :c2) " ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ) ,
( " customer_profile_l3 " , agent_engine ,
" SELECT COUNT(*) FROM customer_profile_l3 WHERE customer_id IN (:c1, :c2) " ,
{ " c1 " : CUSTOMER , " c2 " : CUSTOMER_SMALL } ) ,
( " audit_log " , agent_engine ,
" SELECT COUNT(*) FROM audit_log WHERE trace_id LIKE ' T8M-VERIFY % ' " , { } ) ,
( " risk_convert_detail " , agent_engine ,
" SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id LIKE ' CNV-T8M % ' "
" OR client_request_id LIKE ' T8M- % ' " , { } ) ,
] :
check ( f " { label } 零残留 " , q1 ( engine , sql , * * params ) , 0 )
dispose_engines ( )
print ( " \n " + " = " * 60 )
print ( f " 真库验证(T-8 规则引擎改造 ): { _passed } 项一致 / { _failed } 项不一致 " )
print ( f " 真库验证(T-8 引擎时机 ): { _passed } 项一致 / { _failed } 项不一致 " )
print ( " = " * 60 )
return 1 if _failed else 0