"""T-6 真 MySQL 验证脚本:`accept_convert` 受理事务实跑 + DoD 断言。 **为什么 sqlite 单测全绿还不够(T-6 视角)** 受理段本身「不扣份额、不折算、不写流水」,所以「有没有漏写」在 sqlite 上看不出来; 但 T+1 模型新增的持久化面里,以下事项**只有真 MySQL(InnoDB)才成立**: 1. `core_convert_request.status` 是 **ENUM 6 态**:非法状态必须被库拒绝 (sqlite 无 ENUM,写什么都能进)。 2. `uk_idem`(`client_request_id`)**唯一索引真库生效**:幂等「同键只出一单」 最终靠它兜底,sqlite 只能证明代码分支走过。 3. `requested_at` / `cancel_before` 是 **DATETIME(3)**:受理日 00:00 与截点 15:00 必须零漂移 —— 受理日顺延(R-5)判定完全依赖这两列。 4. `qty` 是 **DECIMAL(18,4)**:受理份额落库零漂移。 5. 受理**不产生** `core_trade` / `core_convert_lot_detail`,且 `core_share_lot` 一分不动(验收 21)—— 必须在真库上重放确认「T 日只落单」。 6. 在途占用(R-3)跨行聚合在真库上正确(`status IN ('accepted','nav_pending')` 的 ENUM 比较语义)。 用法: python scripts/dev/verify_convert_accept.py # 建隔离数据 → 跑断言 → 清理 约定(与 `verify_convert_service.py` 一致): - 隔离前缀 **ACV**(客户 / 产品 / group_id / client_request_id),跑完**两个库**全清; - 建/清数据走 `role="admin"`(需 DELETE);业务链路走默认 ro/rw 账号; - 交易日历沿用真库 `core_trade_calendar`(T-1 已灌 503 交易日),**期望值全部动态推算** —— 日历内容变了脚本不该误报; - 退出码: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.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.convert_service import accept_convert # noqa: E402 from app.service.convert.errors import InsufficientShares # noqa: E402 from app.service.convert.trading_calendar import ( # noqa: E402 next_biz_day, parse_cutoff, ) from app.utils.db import dispose_engines, get_engine # noqa: E402 CUSTOMER = "CUST-ACV" # C5 客户(最高承受力)→ 任何产品均放行 CUSTOMER_LOW = "CUST-ACVL" # C1 客户 → 转入 R4 必被适当性阻断 PROD_OUT = "PROD-ACVO" PROD_IN = "PROD-ACVI" COMPANY = "ACV模拟基金" TA = "TA-ACV-001" QTY_TOTAL = Decimal("150") # 100 + 50 两批 _passed = 0 _failed = 0 def check(name: str, actual, expected) -> None: """逐条断言并打印(与 `verify_convert_service.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, submit_day: date) -> None: with engine.begin() as conn: conn.execute( text( "INSERT INTO core_customer (customer_id, display_name, open_date, age, is_active) " "VALUES (:c, 'ACV真库验证', :d, 40, 1)" ), {"c": CUSTOMER, "d": submit_day - timedelta(days=400)}, ) conn.execute( text( "INSERT INTO core_customer (customer_id, display_name, open_date, age, is_active) " "VALUES (:c, 'ACV低风险客户', :d, 40, 1)" ), {"c": CUSTOMER_LOW, "d": submit_day - timedelta(days=400)}, ) for cid, code in ((CUSTOMER, "C5"), (CUSTOMER_LOW, "C1")): conn.execute( text( "INSERT INTO core_customer_risk " "(customer_id, risk_code, evaluated_at, expires_at) " "VALUES (:c, :r, :t, :exp)" ), { "c": cid, "r": code, "t": submit_day - timedelta(days=30), "exp": submit_day + timedelta(days=300), }, ) # 转出 R2 / 转入 R4(C5 放行、C1 阻断 —— 同一对产品覆盖两条路径) 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, 'ACV转出债基', 'R2', 'bond', 1, 1, 0.0030, 0, 0, " "'force_transfer', :co, :ta)" ), {"p": PROD_OUT, "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, 'ACV转入股基', 'R4', 'stock', 1, 1, 0.0080, :co, :ta)" ), {"p": PROD_IN, "co": COMPANY, "ta": TA}, ) for pid in (PROD_OUT,): conn.execute( text( "INSERT INTO core_product_nav (product_id, nav, daily_chg_pct, nav_date) " "VALUES (:p, 1.0300, 0, :d)" ), {"p": pid, "d": submit_day}, ) for cid, tag in ((CUSTOMER, "ACV"), (CUSTOMER_LOW, "ACVL")): for idx, (qty, nav, confirmed) in enumerate( [ ("100", "1.0300", datetime(2026, 1, 5, 10, 0, 0)), ("50", "1.0000", datetime(2026, 2, 5, 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": f"LOT-{tag}-{idx}", "c": cid, "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": cid, "p": PROD_OUT, "d": submit_day}, ) def cleanup_core(engine) -> None: with engine.begin() as conn: for sql, params in [ ("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_customer_risk WHERE customer_id IN (:c1, :c2)", {"c1": CUSTOMER, "c2": CUSTOMER_LOW}), ("DELETE FROM core_product_nav WHERE product_id LIKE 'PROD-ACV%'", {}), ("DELETE FROM core_product WHERE product_id LIKE 'PROD-ACV%'", {}), ("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-ACV-*);client_request_id 前缀兜底 —— # 早期版本 gid 用生产 `_new_id`(CNV-<日期>-),前缀匹配不到会留残行, # 而 `risk_convert_detail.uk_idem`(client_request_id)一残留就让下一轮全红。 conn.execute( text("DELETE FROM risk_convert_detail WHERE convert_group_id LIKE 'CNV-ACV%'") ) conn.execute( text("DELETE FROM risk_convert_detail WHERE client_request_id LIKE 'ACV-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-ACV%'") ) # ── 链路辅助 ──────────────────────────────────────────────────────── def _services() -> dict: """受理段四仓储**全部走默认账号**(core ro/rw + agent rw)。 ⚠️ 不能像 sqlite 单测那样把同一个引擎塞给四个仓储:真库下 `core_convert_request` 在 **core 库**、`risk_convert_detail` 在 **agent 库**, 且 Core 写侧须用 `xh_core_rw`(D20 账号分离)。传同一引擎会写错库。 """ return dict( core_ro=CoreReadOnlyRepository(), risk_repo=RiskRepository(), convert_repo=ConvertRepository(), request_repo=ConvertRequestRepository(), ) _SEQ = itertools.count(1) def _id_factory(prefix: str, now: datetime) -> str: """固定前缀 `CNV-ACV-nnn`(生产 `_new_id` 会带随机 hex,脚本无法按前缀清理)。""" return f"CNV-ACV-{next(_SEQ):03d}" def _accept(req: dict, now: datetime) -> dict: """受理链路(注入固定前缀 id 工厂,让脚本能按 `CNV-ACV%` 精准清理)。""" return accept_convert(req, now=now, id_factory=_id_factory, **_services()) def _base_req(**over) -> dict: req = { "customer_id": CUSTOMER, "from_product_id": PROD_OUT, "to_product_id": PROD_IN, "qty": Decimal("100"), } req.update(over) return req def _pick_days(core_admin, is_open) -> tuple[date, date]: """取**日历中位**的开市日作基准受理日 D,及其下一交易日。 ⚠️ 不能用「最新开市日」:日历种子到 2027 年底为止,取最后一天会让「过截点顺延」 用例撞上 `next_biz_day` 的数据边界(首轮实跑即因此暴露了 `confirm_eta` 的缺陷)。 取中位日可保证前后都有交易日,脚本断言不受日历边界干扰。 """ 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 _find_closed_day(start: date, is_open) -> date: """从 `start` 往前找一个休市日(含日历缺行 —— `is_open` 对缺行返回 False)。""" cursor = start for _ in range(15): cursor -= timedelta(days=1) if not is_open(cursor): return cursor raise RuntimeError("前 15 天内未找到休市日,日历数据异常") # ── 主流程 ────────────────────────────────────────────────────────── def main() -> int: core_admin = get_engine(settings.mysql_core_database, "admin") agent_admin = get_engine(settings.mysql_database, "admin") core_ro = CoreReadOnlyRepository() cutoff = parse_cutoff(settings.convert_cutoff_time) 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)) print(f"基准受理日 D={day}(下一交易日 {next_day}),截点 {cutoff}") try: cleanup_core(core_admin) cleanup_agent(agent_admin) seed(core_admin, day) # ── A. 受理成功 + 无副作用(验收 21)───────────────────────── print("\n【A】受理成功:只落受理单,不扣份额 / 不写流水") before_lots = dec( q1(core_admin, "SELECT COALESCE(SUM(remain_qty), 0) FROM core_share_lot " "WHERE customer_id = :c AND product_id = :p", c=CUSTOMER, p=PROD_OUT) ) res = _accept(_base_req(client_request_id="ACV-REQ-1"), submit_at) check("受理返回 accepted", res["status"], "accepted") check("受理返回 accepted 标记", res["accepted"], True) check("受理返回非 blocked", res["blocked"], False) check("受理份额", dec(res["qty"]), Decimal("100.00")) check("预计确认日 = 下一交易日", res["confirm_eta"], str(next_day)) row = _rows( core_admin, "SELECT * FROM core_convert_request WHERE convert_group_id = :g", g=res["convert_group_id"], )[0] check("受理单状态 accepted", row["status"], "accepted") check("actual_qty 为空(确认后才有)", row["actual_qty"], None) check("qty 落 DECIMAL(18,4) 零漂移", dec(row["qty"], "0.0001"), Decimal("100.0000")) check("requested_at = 受理日 00:00", row["requested_at"], datetime.combine(day, time.min)) check("cancel_before = 受理日截点", row["cancel_before"], datetime.combine(day, cutoff)) after_lots = dec( q1(core_admin, "SELECT COALESCE(SUM(remain_qty), 0) FROM core_share_lot " "WHERE customer_id = :c AND product_id = :p", c=CUSTOMER, p=PROD_OUT) ) check("份额一分未动(验收 21)", after_lots, before_lots) check("无 core_trade 流水", q1(core_admin, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g", g=res["convert_group_id"]), 0) check("无 core_convert_lot_detail", q1(core_admin, "SELECT COUNT(*) FROM core_convert_lot_detail " "WHERE convert_group_id = :g", g=res["convert_group_id"]), 0) mirror = _rows( agent_admin, "SELECT status, estimated FROM risk_convert_detail WHERE convert_group_id = :g", g=res["convert_group_id"], ) check("镜像落 pending", mirror[0]["status"] if mirror else None, "pending") check("镜像 estimated=0(真实请求)", int(mirror[0]["estimated"]) if mirror else None, 0) check("审计 accepted 1 条", q1(agent_admin, "SELECT COUNT(*) FROM audit_log WHERE event_type = 'convert_request' " "AND decision = 'accepted' AND input_summary LIKE :p", p=f"%{res['convert_group_id']}%"), 1) # ── B. 幂等(验收 12)──────────────────────────────────────── print("\n【B】幂等:同 client_request_id 只出一单(uk_idem 真库生效)") dup = _accept(_base_req(client_request_id="ACV-REQ-1"), submit_at) check("重复提交返回同一 group_id", dup["convert_group_id"], res["convert_group_id"]) check("重复提交标记 idempotent", dup["idempotent"], True) check("受理单仍 1 行", q1(core_admin, "SELECT COUNT(*) FROM core_convert_request WHERE customer_id = :c", c=CUSTOMER), 1) check("uk_idem 唯一索引存在", q1(core_admin, "SELECT COUNT(*) FROM information_schema.statistics " "WHERE table_schema = :db AND table_name = 'core_convert_request' " "AND index_name = 'uk_idem'", db=settings.mysql_core_database), 1) dup_err = None try: with core_admin.begin() as conn: conn.execute( text( "INSERT INTO core_convert_request (convert_group_id, client_request_id, " "customer_id, from_product_id, to_product_id, qty, status, requested_at) " "VALUES ('CNV-ACV-DUP', 'ACV-REQ-1', :c, :o, :i, 1, 'accepted', NOW(3))" ), {"c": CUSTOMER, "o": PROD_OUT, "i": PROD_IN}, ) except Exception as exc: # noqa: BLE001 dup_err = type(exc).__name__ check("真库重复键报 IntegrityError", "IntegrityError" in (dup_err or ""), True) # ── C. 在途占用(验收 22,R-3)──────────────────────────────── print("\n【C】在途占用:可用份额 = 物理余量 − 在途") check("在途占用 = 100", dec(core_ro.sum_inflight_qty(CUSTOMER, PROD_OUT)), Decimal("100.00")) check("可用份额 = 50(150 − 100)", dec(core_ro.sum_remain_qty(CUSTOMER, PROD_OUT) - core_ro.sum_inflight_qty( CUSTOMER, PROD_OUT)), Decimal("50.00")) over_err = None try: _accept(_base_req(qty=Decimal("60")), submit_at) except InsufficientShares as exc: over_err = str(exc) check("超量申请被拒(InsufficientShares)", bool(over_err), True) check("被拒后未新增受理单", q1(core_admin, "SELECT COUNT(*) FROM core_convert_request WHERE customer_id = :c", c=CUSTOMER), 1) # ── D. 受理日顺延(验收 28,R-5)───────────────────────────── print("\n【D】受理日顺延:非交易日 / 过截点 → 下一交易日") closed = _find_closed_day(day, core_ro.is_open) expect_from_closed = next_biz_day(closed, core_ro.is_open) r_closed = _accept( _base_req(qty=Decimal("1"), client_request_id="ACV-REQ-CLOSED"), datetime.combine(closed, time(10, 0)), ) row_closed = _rows( core_admin, "SELECT requested_at FROM core_convert_request WHERE convert_group_id = :g", g=r_closed["convert_group_id"], )[0] check(f"非交易日 {closed} → 顺延 {expect_from_closed}", row_closed["requested_at"], datetime.combine(expect_from_closed, time.min)) r_late = _accept( _base_req(qty=Decimal("1"), client_request_id="ACV-REQ-LATE"), datetime.combine(day, cutoff), ) row_late = _rows( core_admin, "SELECT requested_at FROM core_convert_request WHERE convert_group_id = :g", g=r_late["convert_group_id"], )[0] check(f"恰好 {cutoff} 视为超时 → 顺延 {next_day}", row_late["requested_at"], datetime.combine(next_day, time.min)) # ── E. 适当性阻断(验收 4)─────────────────────────────────── print("\n【E】适当性阻断:C1 客户买 R4 → blocked 且受理单 0 行") low_cnt_before = q1( core_admin, "SELECT COUNT(*) FROM core_convert_request WHERE customer_id = :c", c=CUSTOMER_LOW, ) blocked = _accept( _base_req(customer_id=CUSTOMER_LOW, qty=Decimal("100")), submit_at ) check("返回 blocked", blocked["blocked"], True) check("blocked 无受理单", q1(core_admin, "SELECT COUNT(*) FROM core_convert_request WHERE customer_id = :c", c=CUSTOMER_LOW), low_cnt_before) check("blocked 落预警单", q1(agent_admin, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", c=CUSTOMER_LOW) >= 1, True) # ── F. 强制全转(验收 13 前半)─────────────────────────────── print("\n【F】强制全转:min_hold_qty=100,申请 148 → 记实际全转 150") with core_admin.begin() as conn: conn.execute( text("UPDATE core_product SET min_hold_qty = 100 WHERE product_id = :p"), {"p": PROD_OUT}, ) # 先把在途单撤掉(否则占用导致可用不足) with core_admin.begin() as conn: conn.execute( text("UPDATE core_convert_request SET status = 'cancelled' " "WHERE customer_id = :c AND status IN ('accepted', 'nav_pending')"), {"c": CUSTOMER}, ) forced = _accept( _base_req(qty=Decimal("148"), client_request_id="ACV-REQ-FORCED"), submit_at ) check("强制全转标记", forced["forced_full_transfer"], True) check("申请量回显 148", dec(forced["requested_qty"]), Decimal("148.00")) check("受理 qty 记实际全转 150", dec(forced["qty"]), Decimal("150.00")) # ── 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": forced["convert_group_id"]}, ) except Exception as exc: # noqa: BLE001 enum_err = type(exc).__name__ check("非法 status 被 ENUM 拒绝", bool(enum_err), True) check("拒绝后状态未被污染", q1(core_admin, "SELECT status FROM core_convert_request WHERE convert_group_id = :g", g=forced["convert_group_id"]), "accepted") check("requested_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 = 'requested_at'", db=settings.mysql_core_database), 3) finally: cleanup_core(core_admin) cleanup_agent(agent_admin) dispose_engines() print(f"\n{'=' * 60}") print(f"真库验证(T-6 accept_convert):{_passed} 项一致 / {_failed} 项不一致") return 1 if _failed else 0 if __name__ == "__main__": sys.exit(main())