Files
group_xinghuo_jinrong/scripts/dev/verify_convert_engine.py
T

413 lines
18 KiB
Python
Raw Normal View History

"""T-8 真 MySQL 验证脚本:`process_convert_event` 引擎改造在真库/真账号下实跑 + DoD 断言。
**为什么 sqlite 单测全绿还不够(T-8 视角)**
1. **「阶段 1.5 从跳过变生效」是主流程行为变化**:T-7 落地时 `process_convert_event` 尚不存在,
`convert_service._run_engine` 走 `ImportError` 分支**静默跳过**;T-8 落地后同一笔转换会
**真实出单 + 写 L3 + 落审计**。真库跑一遍才能确认没有连带破坏(审计条数、L3 写入、单落库)。
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。
**断言分组**
A 引擎真跑出单:1 张单 / `payload.events` 两条 / 顺序 [转出, 转入]
B **RISK-002 不翻倍**(阈值夹逼:单条 < 阈值 < 两条之和)
C 去重不删行:`core_trade` 仍 2 条且同 `convert_group_id`
D 幂等重试(同 `client_request_id`)**不产生第二张单**
E 无命中场景(小额)不建单、仅 pass 审计
F 清理后残留为零(自检)
用法:
python scripts/dev/verify_convert_engine.py
约定(与 T-6/T-7 脚本一致):
- 用 **T8M 前缀**隔离数据(客户/产品/批次/group_id),跑完**两个库**(core + agent)全清;
- 建/清数据走 `role="admin"`(需 DELETE);业务本身走 `convert_fund` 默认账号;
- 固定 `trace_id = "T8M-VERIFY-TRACE"`,便于精准清理 `audit_log`。
"""
from __future__ import annotations
import json
import sys
import uuid
from datetime import date, datetime, timedelta
from decimal import ROUND_HALF_UP, Decimal
from pathlib import Path
from sqlalchemy import text
ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(ROOT))
from app.config.settings import settings # noqa: E402
from app.repository.core_ro import CoreReadOnlyRepository # noqa: E402
from app.service.convert.calc import ( # noqa: E402
convert_amount,
diff_fee,
hold_days,
in_qty,
lot_amount,
lot_fee,
plan_lots,
)
from app.service.convert.convert_service import convert_fund # noqa: E402
from app.service.convert.fee import pick_fee_rate # 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"
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
IN_NAV = Decimal("0.9500")
OUT_RATE = Decimal("0.0030")
IN_RATE = Decimal("0.0080")
FEE_TIERS = [
(0, 7, "0.0150"), (7, 30, "0.0100"), (30, 180, "0.0050"),
(180, 365, "0.0025"), (365, None, "0.0000"),
]
CID_REQ = "T8M-IDEM-0001"
TRACE = "T8M-VERIFY-TRACE"
# 夹逼阈值:单条转出 512000 < daily_total 800000 < 两条之和(约 1020000)
TH = RiskThresholds(
large_amount=Decimal("100000"),
daily_total=Decimal("800000"),
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
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 _t8m_id(prefix: str, now: datetime) -> str:
"""固定前缀的 id 工厂:`convert_fund` 默认生成 `CNV-<日期>-<uuid>`,与清理口径
`LIKE 'CNV-T8M%'` 不匹配 → 会在重跑时留下撞唯一键的残留行(T-7 脚本踩过同款坑)。"""
return f"{prefix}-T8M-{uuid.uuid4().hex[:8].upper()}"
def q1(engine, sql: str, **params):
with engine.connect() as conn:
return conn.execute(text(sql), params).scalar()
def rows(engine, sql: str, **params):
with engine.connect() as conn:
return [dict(r) for r in conn.execute(text(sql), params).mappings()]
# ── 数据准备 / 清理 ─────────────────────────────────────────────────
def cleanup(core_engine, agent_engine) -> None:
"""按 T8M 前缀清理两个库(含幂等重跑前置清理)。"""
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_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}),
]:
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-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}
)
conn.execute(
text("DELETE FROM audit_log WHERE trace_id LIKE 'T8M-VERIFY%'")
)
def seed(core_engine) -> None:
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 pid, name, ptype, rate in [
(PROD_OUT, "T8转出基金", "bond", OUT_RATE),
(PROD_IN, "T8转入基金", "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, "
"fund_company, ta_code) "
"VALUES (:p, :n, 'R2', :t, 1, 1, :r, :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},
)
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)),
]:
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": nav, "cat": confirmed},
)
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)"
),
{"c": CUSTOMER, "p": PROD_OUT, "d": TRADE_DATE},
)
def expected_quote(core_ro, requested: str) -> dict:
"""期望折算(全部走生产纯函数,与 service 内部同口径)。"""
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)
rate = pick_fee_rate(
fee_rules, hold_days(TRADE_DATE, alloc.confirmed_at), product_id=PROD_OUT
)
out_amount += amount
redeem_fee += lot_fee(amount, rate)
conv = convert_amount(out_amount, redeem_fee)
gap = diff_fee(conv, OUT_RATE, IN_RATE, settings.convert_diff_fee_mode)
return {
"out_amount": out_amount,
"redeem_fee": redeem_fee,
"convert_amount": conv,
"diff_fee": gap,
"in_amount": conv - gap,
}
def main() -> int:
core_engine = get_engine("jinrong_core", role="admin")
agent_engine = get_engine("jinrong_agent", role="admin")
print("=" * 60)
print("T-8 真库验证:规则引擎改造(amount_view 去重 + process_convert_event 出单)")
print("=" * 60)
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 两条")
new_trace(TRACE)
resp = 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,
thresholds=TH,
id_factory=_t8m_id,
)
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)
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"])
print("\n【B】RISK-002 不翻倍(单条 512000 < 800000 < 两条之和)")
check("triggered_rules 不含 RISK-002", "RISK-002" in resp["triggered_rules"], False)
# 同日两条流水之和确实超过阈值 → 证明「不命中」来自去重而非阈值太松
total_two = exp["out_amount"] + exp["in_amount"]
check("两条之和确实 ≥ 阈值(阈值夹逼成立)", total_two >= TH.daily_total, True)
print(f" └ 转出 {exp['out_amount']} + 转入 {exp['in_amount']} = {total_two}")
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,
)
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"],
)
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,
)
check("重试返回同一 group_id", 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)
# ── 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(
{
"customer_id": CUSTOMER,
"from_product_id": PROD_OUT,
"to_product_id": PROD_IN,
"qty": "1000",
"client_request_id": "T8M-IDEM-0002",
},
now=NOW + timedelta(days=1),
thresholds=TH,
id_factory=_t8m_id,
)
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)
# ── 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_alert 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)
dispose_engines()
print("\n" + "=" * 60)
print(f"真库验证(T-8 规则引擎改造):{_passed} 项一致 / {_failed} 项不一致")
print("=" * 60)
return 1 if _failed else 0
if __name__ == "__main__":
raise SystemExit(main())