Files
group_xinghuo_jinrong/scripts/dev/verify_convert_confirm.py
T

644 lines
31 KiB
Python
Raw Normal View History

"""T-7 真 MySQL 验证脚本:`confirm_service` 确认事务 + T+1 批处理实跑 + DoD 断言。
**为什么 sqlite 单测全绿还不够(T-7 视角)**
确认段是 T+1 模型里唯一**跨 6 张表多写**的段落,而这些写入的真实约束只有 InnoDB
才成立 —— sqlite 上「代码分支走过」不等于「真库落得对」:
1. **确认事务原子性**:`core_convert_request.status` 的条件 UPDATE
(`WHERE status IN ('accepted','nav_pending')`,`rowcount != 1` → 整事务回滚)
依赖 InnoDB 的**行锁 + 影响行数语义**;sqlite 的 `rowcount` 是 matched rows,
MySQL 是 **changed rows** —— 这正是 `nav_pending` 重试会被误判成冲突的根源。
2. **`core_trade` 恰 2 行 / 同 `convert_group_id`**:`uk_idem` 与索引在真库生效
(验收 2 / 24)。
3. **`core_convert_lot_detail.nav_date` / `nav` / `fee_rate` 为 NOT NULL + DECIMAL**:
T 日净值逐批落库零漂移(PRD §5.3.2 示例数字即由此重放)。
4. **`core_trade.traded_at` / `confirmed_at` 是 DATETIME(3)**:确认日时刻零漂移。
5. **批处理锁 + 捞单窗口**:`idx_status_asof` 上 `status IN (accepted, nav_pending)`
的 ENUM 比较 + 按受理日区间捞单,只有真库才验证得到。
6. 确认**不写** `core_cash_flow`(资金流在 T-16 落地前必须为 0 行)。
**核心锚点**:本脚本用 PRD §5.3.2 的真实示例输入(30000+20000 份、T 日净值
1.3604 / 1.9194、申购费率 0.30% / 0.80%、赎回费分档)**在真库上重放**,断言
`68020.00 / 612.18 / 67407.82 / 333.36 / 67074.46 / 34945.54` 逐字节一致 ——
真库与文档不漂移的最终关卡。
用法:
python scripts/dev/verify_convert_confirm.py # 建隔离数据 → 跑断言 → 清理
约定(与 `verify_convert_accept.py` / `verify_convert_service.py` 一致):
- 隔离前缀 **CVF**(客户 / 产品 / group_id / client_request_id),跑完**两个库**全清;
- 建/清数据走 `role="admin"`(需 DELETE);业务链路走默认 ro/rw 账号
(确认段写 Core 由 `ConvertCoreRepository` 内部固定用 `xh_core_rw`,D20);
- 交易日历沿用真库 `core_trade_calendar`(T-1 已灌 503 交易日),基准日取**日历中位
开市日**,期望值动态推算 —— 日历内容变了脚本不该误报;
- 批次 `confirmed_at` 相对基准日 D 计算(D−100 / D−3),使持有天数档位确定;
- 退出码:0=全一致 / 1=有断言不一致 / 2=环境不可用(连不上库或日历缺数据)。
"""
from __future__ import annotations
import itertools
import sys
from datetime import date, datetime, time, 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.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.confirm_service import confirm_batch, confirm_one # noqa: E402
from app.service.convert.convert_service import accept_convert # noqa: E402
from app.service.convert.trading_calendar import next_biz_day # noqa: E402
from app.utils.db import dispose_engines, get_engine # noqa: E402
CUSTOMER = "CUST-CVF" # C5 客户 → 任何产品放行
CUSTOMER_LOW = "CUST-CVFL" # 受理时 C5,确认前降级 C1 → 复核阻断
PROD_OUT = "PROD-CVFO"
PROD_IN = "PROD-CVFI"
COMPANY = "CVF模拟基金"
TA = "TA-CVF-001"
# PRD §5.3.2 示例输入(真实净值口径)
OUT_NAV = Decimal("1.3604")
IN_NAV = Decimal("1.9194")
OUT_RATE = Decimal("0.0030") # 转出端申购费率 → 价外法基准
IN_RATE = Decimal("0.0080") # 转入端申购费率(更高 → 收补差费)
QTY_LOT_1 = Decimal("30000")
QTY_LOT_2 = Decimal("20000")
QTY_TOTAL = QTY_LOT_1 + QTY_LOT_2
_passed = 0
_failed = 0
def check(name: str, actual, expected) -> None:
"""逐条断言并打印(与 `verify_convert_accept.py` 同款输出)。"""
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:
"""真库读回是 `Decimal`(MySQL 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) -> list[dict]:
with engine.connect() as conn:
return [dict(r) for r in conn.execute(text(sql), params).mappings()]
# ── 数据准备 / 清理 ─────────────────────────────────────────────────
def seed(engine, base_day: date) -> None:
"""灌入两个客户、两只产品、费率四档、两批份额、T 日净值。
批次 `confirmed_at` 相对 `base_day`(= 确认日 T+1 的前一交易日 T)计算:
`D−100` → 持有 100 天 → 档位 (30,180) = 0.50%;`D−3` → 持有 3 天 → (0,7) = 1.50%。
"""
with engine.begin() as conn:
for cid, name in ((CUSTOMER, "CVF真库验证"), (CUSTOMER_LOW, "CVF复核降级")):
conn.execute(
text(
"INSERT INTO core_customer (customer_id, display_name, open_date, age, is_active) "
"VALUES (:c, :n, :d, 40, 1)"
),
{"c": cid, "n": name, "d": base_day - timedelta(days=400)},
)
conn.execute(
text(
"INSERT INTO core_customer_risk "
"(customer_id, risk_code, evaluated_at, expires_at) "
"VALUES (:c, 'C5', :t, :exp)"
),
{
"c": cid,
"t": base_day - timedelta(days=30),
"exp": base_day + timedelta(days=300),
},
)
# 转出债基 R2(min_hold 0,不触发强制全转);转入股基 R4(C5 放行)
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, 'CVF转出债基', 'R2', 'bond', 1, 1, :r, 0, 0, "
"'force_transfer', :co, :ta)"
),
{"p": PROD_OUT, "r": float(OUT_RATE), "co": COMPANY, "ta": TA},
)
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, 'CVF转入股基', 'R4', 'stock', 1, 1, :r, :co, :ta)"
),
{"p": PROD_IN, "r": float(IN_RATE), "co": COMPANY, "ta": TA},
)
# 赎回费分档(22 号文 §10 下限口径)
for lo, hi, rate in [
(0, 7, "0.0150"), (7, 30, "0.0100"), (30, 180, "0.0050"),
(180, 365, "0.0025"), (365, None, "0.0000"),
]:
conn.execute(
text(
"INSERT INTO core_fee_rule "
"(product_id, fee_type, min_hold_days, max_hold_days, rate) "
"VALUES (:p, 'redeem', :a, :b, :r)"
),
{"p": PROD_OUT, "a": lo, "b": hi, "r": rate},
)
# 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, :n, 0, :d)"
),
{"p": pid, "n": float(nav), "d": base_day},
)
# 两批份额(转出端),**两个客户各一份**(D 组用 CUSTOMER_LOW),合计 50000
for cid, tag in ((CUSTOMER, "C"), (CUSTOMER_LOW, "L")):
_insert_lots_by_conn(conn, cid, tag, base_day)
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, 51000, 0, :d)"
),
{"c": cid, "p": PROD_OUT, "q": float(QTY_TOTAL), "d": base_day},
)
def _insert_lots_by_conn(conn, customer: str, tag: str, base_day: date) -> None:
"""插两批份额:`D−100` → 持有 100 天 → 档位 (30,180) = 0.50%;`D−3` → 3 天 → 1.50%。"""
for idx, (qty, confirmed_delta) in enumerate(((QTY_LOT_1, 100), (QTY_LOT_2, 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, 1.0300, :cat)"
),
{
"l": f"LOT-CVF-{tag}{idx + 1}",
"c": customer,
"p": PROD_OUT,
"q": float(qty),
"cat": datetime.combine(base_day - timedelta(days=confirmed_delta), time(10, 0)),
},
)
def reset_shares(engine, customer: str, base_day: date) -> None:
"""恢复某客户的转出端批次与持仓 —— 多组用例复用同一客户时的**隔离手段**。
每组用例都会把 50000 份额用光;不恢复会让后一组受理时「可转份额不足 0」。
同时清掉转入端持仓与占用(终态单本身已不占额度,此处只清持仓余额)。
"""
tag = "C" if customer == CUSTOMER else "L"
with engine.begin() as conn:
conn.execute(
text("DELETE FROM core_share_lot WHERE customer_id = :c"), {"c": customer}
)
conn.execute(
text("DELETE FROM core_holding WHERE customer_id = :c AND product_id = :p"),
{"c": customer, "p": PROD_IN},
)
conn.execute(
text(
"UPDATE core_holding SET qty = :q, market_value = 51000 "
"WHERE customer_id = :c AND product_id = :p"
),
{"c": customer, "p": PROD_OUT, "q": float(QTY_TOTAL)},
)
_insert_lots_by_conn(conn, customer, tag, base_day)
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-CVF%'", {}),
("DELETE FROM core_convert_request WHERE customer_id IN (:c1, :c2)",
{"c1": CUSTOMER, "c2": CUSTOMER_LOW}),
("DELETE FROM core_trade WHERE customer_id IN (:c1, :c2)",
{"c1": CUSTOMER, "c2": CUSTOMER_LOW}),
("DELETE FROM core_share_lot WHERE customer_id IN (:c1, :c2)",
{"c1": CUSTOMER, "c2": CUSTOMER_LOW}),
("DELETE FROM core_holding WHERE customer_id IN (:c1, :c2)",
{"c1": CUSTOMER, "c2": CUSTOMER_LOW}),
("DELETE FROM core_cash_flow WHERE customer_id IN (:c1, :c2)",
{"c1": CUSTOMER, "c2": CUSTOMER_LOW}),
("DELETE FROM core_customer_risk WHERE customer_id IN (:c1, :c2)",
{"c1": CUSTOMER, "c2": CUSTOMER_LOW}),
("DELETE FROM core_fee_rule WHERE product_id LIKE 'PROD-CVF%'", {}),
("DELETE FROM core_product_nav WHERE product_id LIKE 'PROD-CVF%'", {}),
("DELETE FROM core_product WHERE product_id LIKE 'PROD-CVF%'", {}),
("DELETE FROM core_customer WHERE customer_id IN (:c1, :c2)",
{"c1": CUSTOMER, "c2": CUSTOMER_LOW}),
]:
conn.execute(text(sql), params)
def cleanup_agent(engine) -> None:
with engine.begin() as conn:
# 双保险:gid 走注入 id_factory(CNV-CVF-*);client_request_id 前缀兜底 ——
# `risk_convert_detail.uk_idem` 一旦残留会让下一轮全红。
conn.execute(
text("DELETE FROM risk_convert_detail WHERE convert_group_id LIKE 'CNV-CVF%'")
)
conn.execute(
text("DELETE FROM risk_convert_detail WHERE client_request_id LIKE 'CVF-REQ-%'")
)
conn.execute(
text("DELETE FROM risk_alert WHERE customer_id IN (:c1, :c2)"),
{"c1": CUSTOMER, "c2": CUSTOMER_LOW},
)
conn.execute(text("DELETE FROM audit_log WHERE input_summary LIKE '%CNV-CVF%'"))
# ── 链路辅助 ────────────────────────────────────────────────────────
def _accept_services() -> dict:
"""受理段四仓储走默认账号(core ro/rw + agent rw)。"""
return dict(
core_ro=CoreReadOnlyRepository(),
risk_repo=RiskRepository(),
convert_repo=ConvertRepository(),
request_repo=ConvertRequestRepository(),
)
def _confirm_services() -> dict:
"""确认段多一个 `core_writer`(内部固定用 `xh_core_rw`,D20 账号分离)。"""
return dict(**_accept_services(), core_writer=ConvertCoreRepository())
_SEQ = itertools.count(1)
def _id_factory(prefix: str, now: datetime) -> str:
"""固定前缀 `CNV-CVF-nnn`(生产 `_new_id` 带随机 hex,脚本无法按前缀清理)。"""
return f"CNV-CVF-{next(_SEQ):03d}"
def _accept(req: dict, now: datetime) -> dict:
return accept_convert(req, now=now, id_factory=_id_factory, **_accept_services())
def _confirm(gid: str, now: datetime, as_of: date, **kw) -> dict:
return confirm_one(gid, now=now, as_of=as_of, **_confirm_services(), **kw)
def _base_req(**over) -> dict:
req = {
"customer_id": CUSTOMER,
"from_product_id": PROD_OUT,
"to_product_id": PROD_IN,
"qty": QTY_TOTAL,
}
req.update(over)
return req
def _pick_days(core_admin, is_open) -> tuple[date, date]:
"""取日历**中位**开市日作 T 日,及其下一交易日(T+1 确认日)。"""
rows = _rows(
core_admin,
"SELECT cal_date FROM core_trade_calendar WHERE is_open = 1 ORDER BY cal_date ASC",
)
if len(rows) < 10:
raise RuntimeError("core_trade_calendar 开市日不足,请先跑 10-seed-trade-calendar.sql")
day = rows[len(rows) // 2]["cal_date"]
if not isinstance(day, date):
day = date.fromisoformat(str(day)[:10])
return day, next_biz_day(day, is_open)
def _req_row(core_admin, gid: str) -> dict:
return _rows(
core_admin,
"SELECT * FROM core_convert_request WHERE convert_group_id = :g",
g=gid,
)[0]
def _trade_rows(core_admin, gid: str) -> list[dict]:
return _rows(
core_admin,
"SELECT trade_id, trade_type, amount, qty, trade_status, traded_at "
"FROM core_trade WHERE convert_group_id = :g ORDER BY trade_type",
g=gid,
)
def _remain(core_admin, lot_id: str) -> Decimal:
return dec(
q1(core_admin, "SELECT remain_qty FROM core_share_lot WHERE lot_id = :l", l=lot_id),
"0.0001",
)
def _inflight(core_ro, customer: str = CUSTOMER) -> Decimal:
return dec(core_ro.sum_inflight_qty(customer, PROD_OUT))
# ── 主流程 ──────────────────────────────────────────────────────────
def main() -> int:
core_admin = get_engine(settings.mysql_core_database, "admin")
agent_admin = get_engine(settings.mysql_database, "admin")
core_ro = CoreReadOnlyRepository()
try:
day, next_day = _pick_days(core_admin, core_ro.is_open)
except Exception as exc: # noqa: BLE001
print(f"❌ 环境不可用:{exc}")
dispose_engines()
return 2
submit_at = datetime.combine(day, time(10, 0))
confirm_at = datetime.combine(next_day, time(9, 0))
print(f"T 日={day}(受理)→ T+1={next_day}(确认业务日)")
try:
cleanup_core(core_admin)
cleanup_agent(agent_admin)
seed(core_admin, day)
# ══ A. 确认全流程 + PRD §5.3.2 逐字节对齐(验收 2 / 10 / 24)══
print("\n【A】确认全流程:扣批次 / 2 行流水 / 两端持仓 / 状态 confirmed")
gid = _accept(_base_req(client_request_id="CVF-REQ-A"), submit_at)["convert_group_id"]
check("受理后占用 = 50000", _inflight(core_ro), Decimal("50000.00"))
check("受理后批次未动(T 日不扣份额)", _remain(core_admin, "LOT-CVF-C1"), Decimal("30000.0000"))
res = _confirm(gid, confirm_at, next_day)
check("确认状态 confirmed", res["status"], "confirmed")
check("确认非部分成交", res["partial"], False)
check("无强制全转", res["forced_full_transfer"], False)
# ── PRD §5.3.2 示例数字(真库重放,逐字节)──
check("转出净值 = T 日净值 1.3604", res["out_nav"], "1.3604")
check("净值日期 = T 日", res["nav_date"], str(day))
check("转出金额 68020.00", res["out_amount"], "68020.00")
check("赎回费 612.18", res["redeem_fee"], "612.18")
check("净转出金额 67407.82", res["convert_amount"], "67407.82")
check("补差费 333.36", res["diff_fee"], "333.36")
check("转入金额 67074.46", res["in_amount"], "67074.46")
check("转入净值 1.9194", res["in_nav"], "1.9194")
check("转入份额 34945.54", res["in_qty"], "34945.54")
check("尾差 -0.0049", res["rounding_diff"], "-0.0049")
check("批次明细 2 条", res["lot_count"], 2)
# ── 批次扣减(FIFO)──
check("批次 1 remain 归零", _remain(core_admin, "LOT-CVF-C1"), Decimal("0.0000"))
check("批次 2 remain 归零", _remain(core_admin, "LOT-CVF-C2"), Decimal("0.0000"))
# ── 流水恰 2 行、同 gid、一赎一申 ──
trades = _trade_rows(core_admin, gid)
check("core_trade 恰 2 行", len(trades), 2)
check("两行类型", sorted(t["trade_type"] for t in trades), ["redeem", "subscribe"])
out_tr = next(t for t in trades if t["trade_type"] == "redeem")
in_tr = next(t for t in trades if t["trade_type"] == "subscribe")
check("转出流水金额 68020.00", dec(out_tr["amount"]), Decimal("68020.00"))
check("转出流水份额 50000.00", dec(out_tr["qty"]), Decimal("50000.00"))
check("转入流水金额 67074.46", dec(in_tr["amount"]), Decimal("67074.46"))
check("转入流水份额 34945.54", dec(in_tr["qty"]), Decimal("34945.54"))
check("流水状态 confirmed", {t["trade_status"] for t in trades}, {"confirmed"})
check("流水交易日 = 确认业务日", {t["traded_at"].date() for t in trades}, {next_day})
# ── 明细表(Core 侧自包含数据源)──
details = _rows(
core_admin,
"SELECT lot_id, qty, hold_days, fee_rate, fee_amount, nav, nav_date "
"FROM core_convert_lot_detail WHERE convert_group_id = :g ORDER BY lot_id",
g=gid,
)
check("明细 2 行", len(details), 2)
d1, d2 = details[0], details[1]
check("明细批 1 持有天数 100", d1["hold_days"], 100)
check("明细批 1 费率 0.0050", dec(d1["fee_rate"], "0.0001"), Decimal("0.0050"))
check("明细批 1 费额 204.06", dec(d1["fee_amount"]), Decimal("204.06"))
check("明细批 1 净值 1.3604", dec(d1["nav"], "0.0001"), Decimal("1.3604"))
check("明细批 2 持有天数 3", d2["hold_days"], 3)
check("明细批 2 费率 0.0150", dec(d2["fee_rate"], "0.0001"), Decimal("0.0150"))
check("明细批 2 费额 408.12", dec(d2["fee_amount"]), Decimal("408.12"))
check("明细净值日期 = T 日", d1["nav_date"], day)
# ── 两端持仓 ──
hold_out = dec(
q1(core_admin, "SELECT qty FROM core_holding WHERE customer_id = :c "
"AND product_id = :p", c=CUSTOMER, p=PROD_OUT)
)
hold_in_row = _rows(
core_admin,
"SELECT qty FROM core_holding WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_IN,
)
check("转出端持仓归零", hold_out, Decimal("0.00"))
check("转入端持仓 = 34945.54",
dec(hold_in_row[0]["qty"]) if hold_in_row else None, Decimal("34945.54"))
# ── 受理单回填 ──
row = _req_row(core_admin, gid)
check("受理单状态 confirmed", row["status"], "confirmed")
check("受理单 actual_qty 回填 50000", dec(row["actual_qty"]), Decimal("50000.00"))
check("受理单 coupon 回填 333.36", dec(row["coupon"]), Decimal("333.36"))
check("受理单 confirmed_at 非空", row["confirmed_at"] is not None, True)
check("确认后占用释放(终态不在途)", _inflight(core_ro), Decimal("0.00"))
# ── 确认不写资金流(T-16 前必须 0 行)──
check("无 core_cash_flow 写入",
q1(core_admin, "SELECT COUNT(*) FROM core_cash_flow WHERE customer_id = :c",
c=CUSTOMER), 0)
# ── agent 镜像 ──
mirror = _rows(
agent_admin,
"SELECT status, fee_amount FROM risk_convert_detail WHERE convert_group_id = :g",
g=gid,
)
check("镜像 completed", mirror[0]["status"] if mirror else None, "completed")
check("审计 confirmed 1 条",
q1(agent_admin,
"SELECT COUNT(*) FROM audit_log WHERE event_type = 'convert_request' "
"AND decision = 'confirmed' AND input_summary LIKE :p",
p=f"%{gid}%"), 1)
# ══ B. 幂等:同 as_of 重复批处理不重复确认(验收 24 后半 / 15)══
print("\n【B】幂等:重复 confirm_batch 同 as_of → 第二次 0 处理、流水不重复")
again = _confirm(gid, confirm_at, next_day)
check("单笔重放返回 skipped", again["status"], "skipped")
batch = confirm_batch(next_day, now=confirm_at, **_confirm_services())
check("批处理锁定成功", batch["locked"], True)
check("第二轮 scanned = 0(终态不再入窗)", batch["scanned"], 0)
check("第二条流水未产生",
q1(core_admin, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g", g=gid), 2)
# ══ C. 缺 T 日净值 → nav_pending → 补净值 → confirmed(验收 25)══
print("\n【C】缺 T 日净值:nav_pending → 补 nav → 重跑 confirmed")
reset_shares(core_admin, CUSTOMER, day) # 恢复份额(A 组已用光)
gid_c = _accept(_base_req(client_request_id="CVF-REQ-C"), submit_at)["convert_group_id"]
with core_admin.begin() as conn:
conn.execute(
text("DELETE FROM core_product_nav WHERE product_id = :p AND nav_date = :d"),
{"p": PROD_IN, "d": day},
)
res_c = _confirm(gid_c, confirm_at, next_day)
check("缺净值 → nav_pending", res_c["status"], "nav_pending")
check("nav_pending 标记 transitioned", res_c["transitioned"], True)
check("缺净值单状态落库 nav_pending", _req_row(core_admin, gid_c)["status"], "nav_pending")
check("缺净值不写流水",
q1(core_admin, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g",
g=gid_c), 0)
check("缺净值批次不动", _remain(core_admin, "LOT-CVF-C1"), Decimal("30000.0000"))
# 重试(仍缺净值)→ 不应被误判成并发冲突
res_c2 = _confirm(gid_c, confirm_at, next_day)
check("缺净值重试仍 nav_pending", res_c2["status"], "nav_pending")
check("重试 transitioned=False(未误判冲突)", res_c2["transitioned"], False)
with core_admin.begin() as conn:
conn.execute(
text(
"INSERT INTO core_product_nav (product_id, nav, daily_chg_pct, nav_date) "
"VALUES (:p, :n, 0, :d)"
),
{"p": PROD_IN, "n": float(IN_NAV), "d": day},
)
res_c3 = _confirm(gid_c, confirm_at, next_day)
check("补净值后 confirmed", res_c3["status"], "confirmed")
check("补净值后流水 2 行",
q1(core_admin, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g",
g=gid_c), 2)
# ══ D. T+1 复核不通过 → rejected + 占用释放(验收 26)══
print("\n【D】T+1 复核不通过:rejected + 占用释放 + 份额不变 + 无流水")
gid_d = _accept(
_base_req(customer_id=CUSTOMER_LOW, client_request_id="CVF-REQ-D"), submit_at
)["convert_group_id"]
check("受理时占用 = 50000", _inflight(core_ro, CUSTOMER_LOW), Decimal("50000.00"))
# 风评到期 + 降级(C1 不匹配 R4)
with core_admin.begin() as conn:
conn.execute(
text("UPDATE core_customer_risk SET risk_code = 'C1', expires_at = :e "
"WHERE customer_id = :c"),
{"c": CUSTOMER_LOW, "e": submit_at - timedelta(days=1)},
)
res_d = _confirm(gid_d, confirm_at, next_day)
check("复核不通过 → rejected", res_d["status"], "rejected")
check("拒绝原因 = 适当性", res_d["reject_reason"] not in (None, ""), True)
check("单状态落库 rejected", _req_row(core_admin, gid_d)["status"], "rejected")
check("不写流水",
q1(core_admin, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g",
g=gid_d), 0)
check("份额不变", _remain(core_admin, "LOT-CVF-L1"), Decimal("30000.0000"))
check("占用释放", _inflight(core_ro, CUSTOMER_LOW), Decimal("0.00"))
# ══ E. 部分成交(占用被抢)→ actual_qty < qty + remark(验收 30)══
print("\n【E】部分成交:占用被抢 → actual_qty < qty、remark=partial")
reset_shares(core_admin, CUSTOMER, day)
gid_e = _accept(_base_req(client_request_id="CVF-REQ-E"), submit_at)["convert_group_id"]
# 模拟 T→T+1 之间份额被别处用掉:批次 1 只剩 15000(可用 15000+20000=35000)
with core_admin.begin() as conn:
conn.execute(
text("UPDATE core_share_lot SET remain_qty = 15000 WHERE lot_id = 'LOT-CVF-C1'")
)
res_e = _confirm(gid_e, confirm_at, next_day)
check("部分成交标记", res_e["partial"], True)
check("实际成交 = 35000(申请 50000)", res_e["actual_qty"], "35000.00")
check("申请量回显 50000", res_e["requested_qty"], "50000.00")
row_e = _req_row(core_admin, gid_e)
check("受理单 remark = partial", row_e["remark"], "partial")
check("受理单 actual_qty 记 35000", dec(row_e["actual_qty"]), Decimal("35000.00"))
check("被抢走后批次 1 归零", _remain(core_admin, "LOT-CVF-C1"), Decimal("0.0000"))
check("批次 2 归零", _remain(core_admin, "LOT-CVF-C2"), Decimal("0.0000"))
check("未确认部分占用释放", _inflight(core_ro), Decimal("0.00"))
# ══ F. 引擎恰好一次(验收 5 / 7 / 29)══
print("\n【F】规则引擎:恰好跑一次(重复确认不再投递)")
reset_shares(core_admin, CUSTOMER, day)
gid_f = _accept(_base_req(client_request_id="CVF-REQ-F"), submit_at)["convert_group_id"]
calls: list[str] = []
def _hook(out_trade, in_trade):
calls.append(str(out_trade.get("convert_group_id")))
return {"triggered_rules": [], "alert_ids": [], "engine_error": False}
res_f = _confirm(gid_f, confirm_at, next_day, engine_hook=_hook)
check("引擎被调用 1 次", len(calls), 1)
check("引擎入参为同组 gid", calls[0], gid_f)
check("确认成功", res_f["status"], "confirmed")
_confirm(gid_f, confirm_at, next_day, engine_hook=_hook)
check("重放不再投递引擎(仍 1 次)", len(calls), 1)
# F2:真实引擎(不注入 hook)+ 重放不翻倍 —— 验收 5/7/29 的真库侧
reset_shares(core_admin, CUSTOMER, day)
gid_f2 = _accept(_base_req(client_request_id="CVF-REQ-F2"), submit_at)["convert_group_id"]
res_f2 = _confirm(gid_f2, confirm_at, next_day)
check("真实引擎路径确认成功", res_f2["status"], "confirmed")
alerts_before = q1(
agent_admin, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", c=CUSTOMER
)
_confirm(gid_f2, confirm_at, next_day)
alerts_after = q1(
agent_admin, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", c=CUSTOMER
)
check("重放不新增 risk_alert(RISK-002 不翻倍)", alerts_after, alerts_before)
# ══ G. 真库专属约束 ══
print("\n【G】真库专属:ENUM 严校验 / 精度 / 索引")
enum_err = None
try:
with core_admin.begin() as conn:
conn.execute(
text("UPDATE core_convert_request SET status = 'bogus' "
"WHERE convert_group_id = :g"),
{"g": gid},
)
except Exception as exc: # noqa: BLE001
enum_err = type(exc).__name__
check("非法 status 被 ENUM 拒绝", bool(enum_err), True)
check("拒绝后状态未被污染", _req_row(core_admin, gid)["status"], "confirmed")
check("confirmed_at 列精度为 DATETIME(3)",
q1(core_admin,
"SELECT datetime_precision FROM information_schema.columns "
"WHERE table_schema = :db AND table_name = 'core_convert_request' "
"AND column_name = 'confirmed_at'",
db=settings.mysql_core_database), 3)
check("idx_status_asof 索引存在",
q1(core_admin,
"SELECT COUNT(*) FROM information_schema.statistics "
"WHERE table_schema = :db AND table_name = 'core_convert_request' "
"AND index_name = 'idx_status_asof'",
db=settings.mysql_core_database), 2)
finally:
cleanup_core(core_admin)
cleanup_agent(agent_admin)
dispose_engines()
print(f"\n{'=' * 60}")
print(f"真库验证(T-7 confirm_service):{_passed} 项一致 / {_failed} 项不一致")
return 1 if _failed else 0
if __name__ == "__main__":
sys.exit(main())