Files
group_xinghuo_jinrong/scripts/dev/verify_convert_engine.py
T
GaoYiYuan_0626 7ad1c8204d 基金转换 T-8:规则引擎改造(_amount_view 去重视图 + process_convert_event)
一次转换落两条流水(转出 redeem + 转入 subscribe,同 convert_group_id),
金额聚合类规则若两条都算会翻倍报假预警,故按「逐笔 / 聚合」拆成两个视图。

- rules.py:新增 _amount_view —— 同组内只保留 redeem 那条(无 convert_group_id
  的交易恒等通过、组内无 redeem 保首条、绝不删行);run_rules 双视图分流:
  RISK-001/003/004 用全量 eligible(逐笔判定),RISK-002/005 走金额视图去重。
- alert_service.py:抽出公开 build_trade_event(事件体结构唯一定义),
  record_trade_alerts 新增可选 events 参数;缺省 None 退化为单条,
  既有调用零改动。新建单落 payload.events 全部,聚合追加只追首条(只认转出端)。
- engine.py:抽 _run 共用实现;process_trade_event 变薄封装(签名与行为不变);
  新增 process_convert_event(out_trade, in_trade, ...) 与 _notify_error_hook
  (hook 自身异常吞掉,原始异常照常上抛)。

实质影响:阶段 1.5 从「ImportError 静默跳过」变为「真跑」——T-7 部署时
process_convert_event 不存在,T-8 落地后同一笔转换会真实出单 + 写 L3 + 落审计;
T-7 真库脚本复跑仍 35/35,无连带破坏。

验证:新增 tests/test_convert_engine.py(15 用例)+ test_convert_service.py
接线回归 1 条;pytest 672 passed / 3 skipped(基线 656 +16,零回归);
新增 scripts/dev/verify_convert_engine.py 真 MySQL 验证 31/31;
突变验证(关掉去重视图 → 4 条变红)确认用例非假绿。
2026-09-10 17:34:35 +08:00

413 lines
18 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""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())