【新增】
- app/service/convert/confirm_service.py:confirm_one(8 步确认)+ confirm_batch
(按业务日捞单 / 整批锁 / 串行 / 单笔容错)。关键裁定:
· T 日净值用 get_nav_on 精确匹配,缺则 nav_pending(绝不回退旧净值);
· 扣批次 + 2 条流水 + 转入批次 + 两端持仓 + 明细 + 受理单置 confirmed
同处一个 Core 单库事务,状态被抢(rowcount!=1)→ 整事务回滚;
· 引擎在事务 commit 后跑,异常不阻断已成立的交易(FR-C28);
· 部分成交 actual=min(申请,可用),被抢部分占用自然释放(R-10)。
- app/service/convert/format.py / audit.py / engine_call.py:展示规格、审计、
引擎调用三处共用出口抽出(受理/确认两段不再各写一份,避免口径漂移)。
- scripts/dev/verify_convert_confirm.py:真库验证 81 项断言,含 PRD §5.3.2
示例在真库上逐字节重放(68020.00/612.18/67407.82/333.36/67074.46/34945.54/-0.0049)。
【修复 · 受理-确认接口契约缺口】
强制全转是**受理段决策**(受理时 qty 已收敛为实际全转量),确认段拿不到原始
申请量、无法复现该判定。修法:
- 受理段把 forced_full_transfer 落受理单 remark(新增 REMARK_FULL_TRANSFER);
- 确认段改为**继承受理决策、不再重判最低持有**(plan_lots 不传 min_hold_qty)
—— 重判会因 T→T+1 可用份额变化得出与受理承诺不一致的结论(擅自扩大客户指令);
- remark 支持多标记 `;` 连接(full_transfer;partial)。
- MySQL rowcount=changed rows 陷阱:nav_pending 重试不得复用 transition_status
的冲突判定,改为 status 未变时不迁移、返回 transitioned=False。
【其他】
- convert_service:新增 cancel_convert(T 日撤单,两道闸门)、_t1_t2_dates
(日历边界 None 容错);受理响应改为 PRD §5.3.1 字段;convert_fund 标 Deprecated。
- convert_core_repository:ConvertApplyInput 加 convert_request_id/diff_fee/
request_remark;_apply_convert_once 末步 _confirm_request 事务内置状态守卫。
- convert_request_repository:_end_of_day 统一闭区间语义;新增 reject()。
- 测试:新增 tests/test_convert_confirm.py(24 例);全量 827 passed / 10 skipped。
630 lines
30 KiB
Python
630 lines
30 KiB
Python
"""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)
|
||
|
||
# ══ 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())
|