Files
group_xinghuo_jinrong/scripts/dev/verify_convert_service.py
T
GaoYiYuan_0626 318cb39a1f 基金转换 T-7:convert_service 八步编排(关键路径 · 幂等前置 + 三阶段)
- 新增 app/service/convert/convert_service.py:八步编排(①-④ 校验不落库 / ⑤ 占位 /
  ⑥ apply_convert / ⑦ 阶段 1.5 引擎不阻断 / ⑧ 回写 + 审计),幂等判定前置到校验之前
  以修复同键重试误报 InsufficientShares
- 新增 tests/test_convert_service.py(17 用例)
- 新增 scripts/dev/verify_convert_service.py(真 MySQL 验证 35/35)
- 改 app/repository/convert_repository.py:complete_convert 由纯 UPDATE 改三步法
  upsert(无占位直跑也能落完成行,R-a 口径)
- 改 app/repository/core_ro.py:新增 has_convert_trades / list_convert_trades /
  list_convert_lot_details 三个只读方法
- 改 app/config/settings.py:新增 6 个 convert 配置项
- 修 tests/test_convert_calc.py:TestPurity 排除编排层 convert_service.py

基线 639 → 656 passed / 3 skipped,零回归
2026-09-10 17:12:01 +08:00

403 lines
20 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-7 真 MySQL 验证脚本:`convert_service.convert_fund` 八步编排实跑 + DoD 断言。
**为什么 sqlite 单测全绿还不够(T-7 视角)**
sqlite 单测验证了八步顺序与折算值,但**证明不了**以下只有真 MySQL(InnoDB)才成立、
且正是本次实施期修正点的事项:
1. **`complete_convert` 三步法 upsert(R-a)在真库成立**:无 `client_request_id` 时
不走占位、`complete_convert` 必须**先查再 INSERT** 完成行。sqlite 单测因 happy path
不带键、从未触达「无占位即 INSERT」分支,看不出 MySQL 下是否真的能落行。
本次修复前该分支缺失 → `risk_convert_detail` 0 行(已被单测暴露)。
2. **`uk_idem` 唯一约束真实存在并生效**:幂等「同键只出一组流水」除了代码判定,
还依赖 `risk_convert_detail.client_request_id` 的 UNIQUE——真库迁移是否真的建了这条
约束、重复插入是否真报 IntegrityError,sqlite 单测无法证明。
3. **`idx_convert_group` 真实存在**:④ 重试判定 `has_convert_trades` 依赖它;
约束缺失会让重试扫描全表或判错。
4. **DECIMAL(18,2/18,4) 精度**:`risk_convert_detail.fee_amount` / `nav` 在 MySQL 下
落库零漂移(sqlite 用 REAL 无此保证)。
5. **阶段二失败兜底不回滚 Core**:本脚本覆盖 happy path 与幂等,4xx 不落库亦验证。
用法:
python scripts/dev/verify_convert_service.py # 建隔离数据 → 跑断言 → 清理
约定(与 T-6 `verify_convert_apply.py` 一致):
- 用 **T7M 前缀**的隔离数据(客户/产品/批次/group_id),跑完**两个库**(core + agent)全清,
不碰既有种子;
- 建/清数据走 `role="admin"`(需 DELETE);事务本身走 `convert_fund` 默认的 ro/rw 账号;
- 固定 `trace_id = "T7M-VERIFY-TRACE"`,便于精准清理 `audit_log`(agent 库)。
"""
from __future__ import annotations
import sys
import threading
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.convert_repository import ConvertRepository # 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, PROCESSING # noqa: E402
from app.service.convert.errors import InsufficientShares # 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.utils.db import dispose_engines, get_engine # noqa: E402
from app.utils.trace import new_trace # noqa: E402
CUSTOMER = "CUST-T7M"
PROD_OUT = "PROD-T7MO"
PROD_IN = "PROD-T7MI"
COMPANY = "华夏模拟基金"
TA = "TA-CN-001"
TRADE_AT = datetime(2026, 9, 4, 10, 0, 0)
TRADE_DATE = date(2026, 9, 4)
NOW = TRADE_AT
OUT_RATE = Decimal("0.0030")
IN_RATE = Decimal("0.0080")
IN_NAV = Decimal("0.9500")
FEE_TIERS = [
(0, 7, "0.0150"), (7, 30, "0.0100"), (30, 180, "0.0050"),
(180, 365, "0.0025"), (365, None, "0.0000"),
]
TRACE = "T7M-VERIFY-TRACE"
TRACE_E = "T7M-VERIFY-E" # 隔离 E 组审计,便于按 trace 精确计数
_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 dec(value, places: str = "0.01") -> Decimal:
return Decimal(str(value)).quantize(Decimal(places), rounding=ROUND_HALF_UP)
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 seed_core(engine) -> None:
with engine.begin() as conn:
conn.execute(
text(
"INSERT INTO core_customer (customer_id, display_name, open_date) "
"VALUES (:c, 'T7真库验证', :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, "T7转出基金", "bond", OUT_RATE),
(PROD_IN, "T7转入基金", "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},
)
for lot_id, qty, nav, confirmed in [
("LOT-T7M-A1", "100", "1.0300", datetime(2026, 8, 1, 10, 0, 0)),
("LOT-T7M-A2", "50", "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, 150, 150.00, 154.50, 0, :d)"
),
{"c": CUSTOMER, "p": PROD_OUT, "d": TRADE_DATE},
)
def seed_nav_stale(engine) -> None:
"""把转入端净值改成过期(距交易日 10 天 > 阈值 3)以触发 nav_stale 副审计。"""
with engine.begin() as conn:
conn.execute(
text("UPDATE core_product_nav SET nav_date = :d WHERE product_id = :p"),
{"d": date(2026, 8, 25), "p": PROD_IN},
)
def cleanup_core(engine) -> None:
with engine.begin() as conn:
for sql, params in [
("DELETE FROM core_convert_lot_detail WHERE convert_group_id LIKE 'CNV-T7M%'", {}),
("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-T7M%'", {}),
("DELETE FROM core_product_nav WHERE product_id LIKE 'PROD-T7M%'", {}),
("DELETE FROM core_product WHERE product_id LIKE 'PROD-T7M%'", {}),
("DELETE FROM core_customer WHERE customer_id = :c", {"c": CUSTOMER}),
]:
conn.execute(text(sql), params)
def cleanup_agent(engine) -> None:
with engine.begin() as conn:
conn.execute(
text("DELETE FROM risk_convert_detail WHERE convert_group_id LIKE 'CNV-T7M%'")
)
conn.execute(
text("DELETE FROM audit_log WHERE trace_id IN (:t1, :t2)"),
{"t1": TRACE, "t2": TRACE_E},
)
# ── 期望折算(全部走生产纯函数,与 service 内部同口径) ─────────────────
def compute_expected(engine, requested: str) -> dict:
with engine.connect() as conn:
lot_rows = [
Lot.from_row(r)
for r in conn.execute(
text(
"SELECT * FROM core_share_lot WHERE customer_id = :c AND product_id = :p "
"ORDER BY confirmed_at ASC, lot_id ASC"
),
{"c": CUSTOMER, "p": PROD_OUT},
).mappings()
]
rules = [
FeeRule.from_row(r)
for r in conn.execute(
text(
"SELECT * FROM core_fee_rule WHERE product_id = :p AND fee_type = 'redeem' "
"ORDER BY min_hold_days ASC"
),
{"p": PROD_OUT},
).mappings()
]
plan = plan_lots(lot_rows, 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(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)
in_amount = conv - diff_fee(conv, OUT_RATE, IN_RATE, "amount_diff")
return {
"out_amount": out_amount,
"redeem_fee": redeem_fee,
"convert_amount": conv,
"in_amount": in_amount,
"in_qty": in_qty(in_amount, IN_NAV),
"actual_qty": plan.actual_qty,
}
def _req(customer: str = CUSTOMER, qty: str = "120", cid_req: str | None = None) -> dict:
return {
"customer_id": customer,
"from_product_id": PROD_OUT,
"to_product_id": PROD_IN,
"qty": qty,
"client_request_id": cid_req,
}
def _t7m_id(prefix: str, now: datetime) -> str:
"""注入 `convert_fund` 的 id_factory:让 group_id 带 T7M 前缀,
与 `cleanup_agent` 的 `LIKE 'CNV-T7M%'` 对齐,跑完可精准清理(避免遗留行
撞 `uk_group` 唯一约束)。默认 `_new_id` 生成 `CNV-<日期>-<uuid>`,不含 T7M。"""
return f"CNV-T7M-{uuid.uuid4().hex[:8].upper()}"
def main() -> int:
new_trace(TRACE) # 固定 trace,便于清理 audit_log
core_admin = get_engine(settings.mysql_core_database, "admin")
agent_admin = get_engine(settings.mysql_database, "admin")
run_ts = TRADE_AT.strftime("%H%M%S") # 同一次运行内 cid_req 唯一,避免跨运行锁冲突
try:
# 初清 + 建隔离数据
cleanup_core(core_admin)
cleanup_agent(agent_admin)
seed_core(core_admin)
# ── A. 八步贯通(带幂等键):占位 pending → 2 流水 → completed + 审计 ──
print("\n【A】八步贯通(cid_req 带键):2 流水 + risk_convert_detail 完成 + 审计")
cid_a = f"T7M-REQ-A-{run_ts}"
# 期望折算必须在转换前算(转换会扣减份额,转换后读库会得到不足份额)
exp = compute_expected(core_admin, "120")
first = convert_fund(_req(cid_req=cid_a), now=NOW, actor_id="verify-t7m", id_factory=_t7m_id)
check("未拦截", first.get("blocked"), False)
check("estimated=True", first.get("estimated"), True)
check("group_id 前缀", first["convert_group_id"].startswith("CNV-"), True)
check("流水条数", q1(core_admin, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g", g=first["convert_group_id"]), 2)
check("R-b trade_type", sorted(r["trade_type"] for r in rows(core_admin, "SELECT trade_type FROM core_trade WHERE convert_group_id = :g", g=first["convert_group_id"])), ["redeem", "subscribe"])
check("risk_convert_detail 条数", q1(agent_admin, "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id = :g", g=first["convert_group_id"]), 1)
check("risk_convert_detail 状态", q1(agent_admin, "SELECT status FROM risk_convert_detail WHERE convert_group_id = :g", g=first["convert_group_id"]), "completed")
# DECIMAL 精度:真库落库零漂移
detail = rows(agent_admin, "SELECT fee_amount, nav FROM risk_convert_detail WHERE convert_group_id = :g", g=first["convert_group_id"])[0]
check("fee_amount 精度(2位)", dec(detail["fee_amount"]), dec(exp["redeem_fee"]))
check("nav 精度(4位)", dec(detail["nav"], "0.0001"), dec(exp["in_amount"] / exp["in_qty"], "0.0001"))
# 审计:主审计 1 条;净值新鲜 → 无 nav_stale 副审计
check("convert_accepted 审计", q1(agent_admin, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'convert_accepted'", t=TRACE), 1)
check("无 nav_stale 审计", q1(agent_admin, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'nav_stale'", t=TRACE), 0)
# 引擎未落地 → 不阻断、无 engine_error 审计
check("engine_error=False", first.get("engine_error"), False)
check("无 engine_error 审计", q1(agent_admin, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'engine_error'", t=TRACE), 0)
# 折算值对齐生产 calc
check("out_amount", dec(first["out_amount"]), dec(exp["out_amount"]))
check("redeem_fee", dec(first["redeem_fee"]), dec(exp["redeem_fee"]))
check("convert_amount", dec(first["convert_amount"]), dec(exp["convert_amount"]))
check("in_amount", dec(first["in_amount"]), dec(exp["in_amount"]))
check("in_qty", dec(first["in_qty"]), dec(exp["in_qty"]))
check("actual_qty", dec(first["actual_qty"]), dec(exp["actual_qty"]))
# ── B. 幂等:同键重复 → 返回首次结果、只一组流水 ────────────────
print("\n【B】幂等同键重试:返回首次结果、不产生第二组流水")
again = convert_fund(_req(cid_req=cid_a), now=NOW, actor_id="verify-t7m", id_factory=_t7m_id)
check("group_id 一致", again["convert_group_id"], first["convert_group_id"])
check("in_qty 一致", dec(again["in_qty"]), dec(first["in_qty"]))
check("out_trade_id 一致", again["out_trade_id"], first["out_trade_id"])
check("core_trade 仅 2 条(无第二组)", q1(core_admin, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g", g=first["convert_group_id"]), 2)
check("risk_convert_detail 仅 1 行", q1(agent_admin, "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id = :g", g=first["convert_group_id"]), 1)
# ── C. 无幂等键:complete_convert 三步法 upsert 必须 INSERT 完成行 ──
print("\n【C】无 cid_req:complete_convert upsert 直接 INSERT 完成行(修复点)")
resp_c = convert_fund(_req(qty="30"), now=NOW, actor_id="verify-t7m", id_factory=_t7m_id)
check("未拦截", resp_c.get("blocked"), False)
check("risk_convert_detail 落 1 完成行", q1(agent_admin, "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id = :g AND status = 'completed'", g=resp_c["convert_group_id"]), 1)
check("core_trade 2 条", q1(core_admin, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g", g=resp_c["convert_group_id"]), 2)
# ── D. 4xx(份额不足):不占位、不落流水、不审计 ─────────────────
print("\n【D】4xx 分支(InsufficientShares):零残留")
before_detail = q1(agent_admin, "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id LIKE 'CNV-T7M%'")
raised = None
try:
convert_fund(_req(qty="500"), now=NOW, actor_id="verify-t7m", id_factory=_t7m_id)
except InsufficientShares:
raised = "InsufficientShares"
check("抛 InsufficientShares", raised, "InsufficientShares")
check("无新增占位", q1(agent_admin, "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id LIKE 'CNV-T7M%'"), before_detail)
# ── E. nav_stale 副审计(真库,净值过期) ───────────────────────
print("\n【E】nav_stale 副审计:净值过期额外落 1 条审计")
cleanup_core(core_admin)
seed_core(core_admin)
seed_nav_stale(core_admin)
cid_e = f"T7M-REQ-E-{run_ts}"
new_trace(TRACE_E) # 切独立 trace,按 trace 精确计数 E 组审计
resp_e = convert_fund(_req(cid_req=cid_e, qty="120"), now=NOW, actor_id="verify-t7m", id_factory=_t7m_id)
check("nav_stale=True", resp_e.get("nav_stale"), True)
check("nav_stale 审计 1 条", q1(agent_admin, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'nav_stale'", t=TRACE_E), 1)
check("convert_accepted 审计 1 条", q1(agent_admin, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'convert_accepted'", t=TRACE_E), 1)
# ── F. 重试判定依赖的索引真实存在 ──────────────────────────────
print("\n【F】约束真实存在(重试判定 / 幂等兜底依赖)")
check(
"core_trade.idx_convert_group 存在",
q1(core_admin, "SELECT COUNT(*) FROM information_schema.statistics "
"WHERE table_schema = :db AND table_name = 'core_trade' "
"AND index_name = 'idx_convert_group'", db=settings.mysql_core_database),
1,
)
check(
"risk_convert_detail.uk_idem 存在",
q1(agent_admin, "SELECT COUNT(*) FROM information_schema.statistics "
"WHERE table_schema = :db AND table_name = 'risk_convert_detail' "
"AND index_name = 'uk_idem'", db=settings.mysql_database),
1,
)
# ── G. uk_idem 真库生效:重复 client_request_id 必须 IntegrityError ──
print("\n【G】uk_idem 真库生效:重复键报 IntegrityError")
with agent_admin.begin() as conn:
conn.execute(
text("INSERT INTO risk_convert_detail (convert_group_id, client_request_id, status, estimated) "
"VALUES ('CNV-T7M-DUP', 'T7M-DUP-KEY', 'pending', 0)"),
)
dup_err = None
try:
with agent_admin.begin() as conn:
conn.execute(
text("INSERT INTO risk_convert_detail (convert_group_id, client_request_id, status, estimated) "
"VALUES ('CNV-T7M-DUP2', 'T7M-DUP-KEY', 'pending', 0)"),
)
except Exception as exc: # noqa: BLE001
dup_err = type(exc).__name__
check("重复键报 IntegrityError", "IntegrityError" in (dup_err or ""), True)
with agent_admin.begin() as conn:
conn.execute(text("DELETE FROM risk_convert_detail WHERE client_request_id = 'T7M-DUP-KEY'"))
finally:
cleanup_core(core_admin)
cleanup_agent(agent_admin)
dispose_engines()
print(f"\n{'=' * 60}")
print(f"真库验证(T-7 convert_service):{_passed} 项一致 / {_failed} 项不一致")
return 1 if _failed else 0
if __name__ == "__main__":
sys.exit(main())