diff --git a/scripts/dev/verify_convert_engine.py b/scripts/dev/verify_convert_engine.py index d94f101..5ad0608 100644 --- a/scripts/dev/verify_convert_engine.py +++ b/scripts/dev/verify_convert_engine.py @@ -1,40 +1,45 @@ -"""T-8 真 MySQL 验证脚本:`process_convert_event` 引擎改造在真库/真账号下实跑 + DoD 断言。 +"""T-8 真 MySQL 验证脚本:**引擎时机**(受理不触引擎 / 确认事务提交后恰跑一次)实跑 + DoD 断言。 **为什么 sqlite 单测全绿还不够(T-8 视角)** -1. **「阶段 1.5 从跳过变生效」是主流程行为变化**:T-7 落地时 `process_convert_event` 尚不存在, - `convert_service._run_engine` 走 `ImportError` 分支**静默跳过**;T-8 落地后同一笔转换会 - **真实出单 + 写 L3 + 落审计**。真库跑一遍才能确认没有连带破坏(审计条数、L3 写入、单落库)。 +1. **「引擎时机」是主流程行为**:T+1 模型下引擎必须**只在确认事务提交后**跑一次(FR-C28)—— + 受理段触引擎会让「尚未确认的交易」提前进风控台账(验收 29 前半), + 确认段重复触引擎会让 RISK-002 当日累计**翻倍**(验收 5/7)。真库跑一遍才能证明两条都成立。 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。 +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 清理后残留为零(自检) +A0 受理**不触引擎**:受理后 `risk_alert` 0 新增、无 `pass` / `engine_error` 审计(验收 29 前半) +A 确认后引擎真跑出单:1 张单 / `payload.events` 两条 / 顺序 [转出, 转入] +B **RISK-002 不翻倍**(阈值夹逼:单条 < 阈值 < 两条之和) +C 去重不删行:`core_trade` 仍 2 条且同 `convert_group_id` +D 幂等重放(重复确认)**不产生第二张单** +E 无命中场景(独立小额客户)不建单、仅 pass 审计 +F 清理后残留为零(自检) 用法: python scripts/dev/verify_convert_engine.py 约定(与 T-6/T-7 脚本一致): - 用 **T8M 前缀**隔离数据(客户/产品/批次/group_id),跑完**两个库**(core + agent)全清; -- 建/清数据走 `role="admin"`(需 DELETE);业务本身走 `convert_fund` 默认账号; +- 建/清数据走 `role="admin"`(需 DELETE);业务链路走默认账号(确认段写 Core 由 + `ConvertCoreRepository` 内部固定用 `xh_core_rw`,D20); +- 基准日取真库 `core_trade_calendar` 的**中位开市日**作 T 日(避开数据边界), + T+1 = 下一交易日;净值插在 T 日(确认段 `get_nav_on` 精确匹配,缺则 `nav_pending`); +- 批次 `confirmed_at` 相对 T 日计算(T−400 / T−3),使持有期档位与费额确定; - 固定 `trace_id = "T8M-VERIFY-TRACE"`,便于精准清理 `audit_log`。 """ from __future__ import annotations +import itertools import json import sys -import uuid -from datetime import date, datetime, timedelta -from decimal import ROUND_HALF_UP, Decimal +from datetime import date, datetime, time, timedelta +from decimal import Decimal from pathlib import Path from sqlalchemy import text @@ -43,7 +48,13 @@ 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.calc import ( # noqa: E402 convert_amount, diff_fee, @@ -53,21 +64,23 @@ from app.service.convert.calc import ( # noqa: E402 lot_fee, plan_lots, ) -from app.service.convert.convert_service import convert_fund # noqa: E402 +from app.service.convert.confirm_service import confirm_one # noqa: E402 +from app.service.convert.convert_service import accept_convert # noqa: E402 from app.service.convert.fee import pick_fee_rate # noqa: E402 +from app.service.convert.trading_calendar import next_biz_day # 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" +CUSTOMER_SMALL = "CUST-T8MS" # E 组:无命中场景专用(独立客户避免当日累计互相干扰) 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 + +OUT_NAV = Decimal("1.0000") # T 日转出净值(未知价法 → 逐批金额按它计) IN_NAV = Decimal("0.9500") OUT_RATE = Decimal("0.0030") IN_RATE = Decimal("0.0080") @@ -75,10 +88,17 @@ FEE_TIERS = [ (0, 7, "0.0150"), (7, 30, "0.0100"), (30, 180, "0.0050"), (180, 365, "0.0025"), (365, None, "0.0000"), ] +QTY_ALL = "500000" CID_REQ = "T8M-IDEM-0001" TRACE = "T8M-VERIFY-TRACE" -# 夹逼阈值:单条转出 512000 < daily_total 800000 < 两条之和(约 1020000) +#: 基准受理日(T 日)—— 运行时由真库日历在 `main` 内赋值(模块级无法静态确定) +TA_DAY: date | None = None + +# 夹逼阈值:单条转出 500000 < daily_total 800000 < 两条之和(约 996035) +# ⚠️ **必须显式传给 `confirm_one(thresholds=TH)`**:缺省会退回 `RiskThresholds.from_settings()`, +# settings 的 daily_total 一旦低于单条金额,RISK-002 就会「命中」—— 症状看着像去重失效, +# 实为阈值口径没传(首轮实跑即栽在此,B 组断言直接捕获)。 TH = RiskThresholds( large_amount=Decimal("100000"), daily_total=Decimal("800000"), @@ -106,10 +126,13 @@ def check(name: str, actual, expected) -> None: print(f" {flag} {name}: 实际 {actual!r}" + ("" if ok else f" / 期望 {expected!r}")) +_SEQ = itertools.count(1) + + def _t8m_id(prefix: str, now: datetime) -> str: - """固定前缀的 id 工厂:`convert_fund` 默认生成 `CNV-<日期>-`,与清理口径 - `LIKE 'CNV-T8M%'` 不匹配 → 会在重跑时留下撞唯一键的残留行(T-7 脚本踩过同款坑)。""" - return f"{prefix}-T8M-{uuid.uuid4().hex[:8].upper()}" + """固定前缀的 id 工厂:生产 `_new_id` 生成 `CNV-<日期>-`, + 与清理口径 `LIKE 'CNV-T8M%'` 不匹配 → 重跑会留撞唯一键的残行(T-6/T-7 均踩过)。""" + return f"CNV-T8M-{next(_SEQ):03d}" def q1(engine, sql: str, **params): @@ -128,14 +151,21 @@ def cleanup(core_engine, agent_engine) -> None: 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_convert_request WHERE customer_id IN (:c1, :c2)", + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}), + ("DELETE FROM core_trade WHERE customer_id IN (:c1, :c2)", + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}), + ("DELETE FROM core_share_lot WHERE customer_id IN (:c1, :c2)", + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}), + ("DELETE FROM core_holding WHERE customer_id IN (:c1, :c2)", + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}), + ("DELETE FROM core_customer_risk WHERE customer_id IN (:c1, :c2)", + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}), ("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}), + ("DELETE FROM core_customer WHERE customer_id IN (:c1, :c2)", + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}), ]: conn.execute(text(sql), params) with agent_engine.begin() as conn: @@ -145,32 +175,37 @@ def cleanup(core_engine, agent_engine) -> None: "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} + text("DELETE FROM risk_alert WHERE customer_id IN (:c1, :c2)"), + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}, ) conn.execute( - text("DELETE FROM audit_log WHERE trace_id LIKE 'T8M-VERIFY%'") + text("DELETE FROM customer_profile_l3 WHERE customer_id IN (:c1, :c2)"), + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}, ) + conn.execute(text("DELETE FROM audit_log WHERE trace_id LIKE 'T8M-VERIFY%'")) -def seed(core_engine) -> None: +def seed(core_engine, base_day: date, submit_at: datetime) -> None: + """两个客户 + 两只产品 + 费率五档 + T 日净值 + 批次(T−400 / T−3)。""" 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 cid, name in ((CUSTOMER, "T8引擎时机"), (CUSTOMER_SMALL, "T8小额客户")): + conn.execute( + text( + "INSERT INTO core_customer (customer_id, display_name, open_date) " + "VALUES (:c, :n, :d)" + ), + {"c": cid, "n": name, "d": base_day - timedelta(days=500)}, + ) + conn.execute( + text( + "INSERT INTO core_customer_risk " + "(customer_id, risk_code, evaluated_at, expires_at) " + "VALUES (:c, 'C3', :t, :exp)" + ), + {"c": cid, "t": submit_at - timedelta(days=30), + "exp": submit_at + timedelta(days=300)}, + ) for pid, name, ptype, rate in [ (PROD_OUT, "T8转出基金", "bond", OUT_RATE), (PROD_IN, "T8转入基金", "stock", IN_RATE), @@ -179,8 +214,8 @@ def seed(core_engine) -> None: 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)" + "min_hold_qty, min_hold_action, fund_company, ta_code) " + "VALUES (:p, :n, 'R2', :t, 1, 1, :r, 0, 'force_transfer', :co, :ta)" ), {"p": pid, "n": name, "t": ptype, "r": str(rate), "co": COMPANY, "ta": TA}, ) @@ -192,45 +227,63 @@ def seed(core_engine) -> None: ), {"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)), - ]: + # 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, :nav, 0, :d)" + ), + {"p": pid, "nav": str(nav), "d": base_day}, + ) + # 主客户两批:40 万份(持 400 天,费率 0)+ 10 万份(持 3 天,费率 1.5%) + for lot_id, qty, confirmed_delta in ( + ("LOT-T8M-A1", "400000", 400), + ("LOT-T8M-A2", "100000", 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, :nav, :cat)" + "remain_qty, nav, confirmed_at) VALUES (:l, :c, :p, :q, :q, 1.0300, :cat)" ), - {"l": lot_id, "c": CUSTOMER, "p": PROD_OUT, "q": qty, "nav": nav, "cat": confirmed}, + {"l": lot_id, "c": CUSTOMER, "p": PROD_OUT, "q": qty, + "cat": datetime.combine(base_day - timedelta(days=confirmed_delta), time(10, 0))}, ) + # 小额客户一批(1000 份,持 400 天 → 费率 0) 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)" + "INSERT INTO core_share_lot (lot_id, customer_id, product_id, qty, " + "remain_qty, nav, confirmed_at) VALUES ('LOT-T8M-S1', :c, :p, 1000, 1000, " + "1.0300, :cat)" ), - {"c": CUSTOMER, "p": PROD_OUT, "d": TRADE_DATE}, + {"c": CUSTOMER_SMALL, "p": PROD_OUT, + "cat": datetime.combine(base_day - timedelta(days=400), time(10, 0))}, ) + for cid, qty in ((CUSTOMER, "500000"), (CUSTOMER_SMALL, "1000")): + 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, :q, 0, :d)" + ), + {"c": cid, "p": PROD_OUT, "q": qty, "d": base_day}, + ) def expected_quote(core_ro, requested: str) -> dict: - """期望折算(全部走生产纯函数,与 service 内部同口径)。""" + """期望折算(全部走生产纯函数,与 service 同口径)。 + + ⚠️ T+1 变更点:逐批金额按 **T 日净值**(`OUT_NAV`)计,**不是**批次买入净值 —— + 这正是「未知价法」与 v1.0 的实质差异,脚本不得沿用 `alloc.nav`。 + """ 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) + amount = lot_amount(alloc.qty, OUT_NAV) rate = pick_fee_rate( - fee_rules, hold_days(TRADE_DATE, alloc.confirmed_at), product_id=PROD_OUT + fee_rules, hold_days(TA_DAY, alloc.confirmed_at), product_id=PROD_OUT ) out_amount += amount redeem_fee += lot_fee(amount, rate) @@ -242,168 +295,212 @@ def expected_quote(core_ro, requested: str) -> dict: "convert_amount": conv, "diff_fee": gap, "in_amount": conv - gap, + "in_qty": in_qty(conv - gap, IN_NAV), } +def _accept_services() -> dict: + return dict( + core_ro=CoreReadOnlyRepository(), + risk_repo=RiskRepository(), + convert_repo=ConvertRepository(), + request_repo=ConvertRequestRepository(), + ) + + +def _confirm_services() -> dict: + return dict(**_accept_services(), core_writer=ConvertCoreRepository()) + + +def _pick_days(core_admin, is_open) -> tuple[date, date]: + """取日历**中位**开市日作 T 日,及其下一交易日(T+1 确认业务日)。 + + ⚠️ 不能用「最新开市日」:日历种子到 2027 年底为止,取末日会让 T+1 撞 + `next_biz_day` 的数据边界(T-6 实测暴露过该缺陷)。 + """ + rws = rows( + core_admin, + "SELECT cal_date FROM core_trade_calendar WHERE is_open = 1 ORDER BY cal_date ASC", + ) + if len(rws) < 10: + raise RuntimeError("core_trade_calendar 开市日不足,请先跑 10-seed-trade-calendar.sql") + day = rws[len(rws) // 2]["cal_date"] + if not isinstance(day, date): + day = date.fromisoformat(str(day)[:10]) + return day, next_biz_day(day, is_open) + + def main() -> int: + global TA_DAY core_engine = get_engine("jinrong_core", role="admin") agent_engine = get_engine("jinrong_agent", role="admin") + core_ro = CoreReadOnlyRepository() + print("=" * 60) - print("T-8 真库验证:规则引擎改造(amount_view 去重 + process_convert_event 出单)") + print("T-8 真库验证:引擎时机(受理不触引擎 / 确认后恰跑一次)") print("=" * 60) + try: + ta_day, t1_day = _pick_days(core_engine, core_ro.is_open) + except Exception as exc: # noqa: BLE001 + print(f"❌ 环境不可用:{exc}") + dispose_engines() + return 2 + TA_DAY = ta_day + submit_at = datetime.combine(ta_day, time(10, 0)) + confirm_at = datetime.combine(t1_day, time(9, 0)) + print(f"T 日={ta_day}(受理)→ T+1={t1_day}(确认业务日)") + 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 两条") + seed(core_engine, ta_day, submit_at) new_trace(TRACE) - resp = convert_fund( + exp = expected_quote(core_ro, QTY_ALL) + + # ── A0:受理**不触引擎**(FR-C28 · 验收 29 前半)───────────────── + print("\n【A0】受理不触引擎:risk_alert 0 新增、无 pass / engine_error 审计") + accepted = accept_convert( { "customer_id": CUSTOMER, "from_product_id": PROD_OUT, "to_product_id": PROD_IN, - "qty": req_qty, + "qty": Decimal(QTY_ALL), "client_request_id": CID_REQ, }, - now=NOW, - thresholds=TH, + now=submit_at, id_factory=_t8m_id, + **_accept_services(), ) - 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) + gid = accepted["convert_group_id"] + check("受理成功(accepted)", accepted["status"], "accepted") + check("受理后 risk_alert 0 行", + q1(agent_engine, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", + c=CUSTOMER), 0) + check("受理后无 pass 审计", + q1(agent_engine, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t " + "AND decision = 'pass'", t=TRACE), 0) + check("受理后无 engine_error 审计", + q1(agent_engine, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t " + "AND decision = 'engine_error'", t=TRACE), 0) + check("受理后 core_trade 0 行", + q1(core_engine, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g", + g=gid), 0) - alerts = rows( - agent_engine, "SELECT * FROM risk_alert WHERE customer_id = :c", c=CUSTOMER - ) + # ── A:确认后引擎真跑出单 ──────────────────────────────────────── + print("\n【A】确认后引擎真跑:一张单 + payload.events 两条") + res = confirm_one(gid, now=confirm_at, as_of=t1_day, thresholds=TH, **_confirm_services()) + check("确认状态 confirmed", res["status"], "confirmed") + check("折算 out_amount 与生产纯函数一致", Decimal(res["out_amount"]), exp["out_amount"]) + check("折算 in_amount 与生产纯函数一致", Decimal(res["in_amount"]), exp["in_amount"]) + check("折算 in_qty 与生产纯函数一致", Decimal(res["in_qty"]), exp["in_qty"]) + check("RISK-001 命中(转出端 500000 ≥ 100000)", "RISK-001" in res["triggered_rules"], True) + check("只出一张单", len(res["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"]) + check("events 顺序 [转出, 转入]", + [e["trade_type"] for e in payload["events"]], ["redeem", "subscribe"]) + check("events[0].trade_id 为转出端", payload["events"][0]["trade_id"], res["out_trade_id"]) + check("events[1].trade_id 为转入端", payload["events"][1]["trade_id"], res["in_trade_id"]) + check("主流水口径取转出端", alerts[0]["trade_id"], res["out_trade_id"]) - print("\n【B】RISK-002 不翻倍(单条 512000 < 800000 < 两条之和)") - check("triggered_rules 不含 RISK-002", "RISK-002" in resp["triggered_rules"], False) - # 同日两条流水之和确实超过阈值 → 证明「不命中」来自去重而非阈值太松 + # ── B:RISK-002 不翻倍(阈值夹逼)──────────────────────────────── + print("\n【B】RISK-002 不翻倍(单条 500000 < 800000 < 两条之和)") + check("triggered_rules 不含 RISK-002", "RISK-002" in res["triggered_rules"], False) total_two = exp["out_amount"] + exp["in_amount"] - check("两条之和确实 ≥ 阈值(阈值夹逼成立)", total_two >= TH.daily_total, True) + check("两条之和确实 ≥ 阈值(夹逼成立)", total_two >= TH.daily_total, True) print(f" └ 转出 {exp['out_amount']} + 转入 {exp['in_amount']} = {total_two}") + # ── C:去重不删行 ──────────────────────────────────────────────── 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, - ) + 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"], - ) + 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, + # ── D:幂等重放(重复确认)不产生第二张单 ──────────────────────── + print("\n【D】幂等重放:重复确认 skipped,不产生第二张单") + again = confirm_one( + gid, now=confirm_at, as_of=t1_day, thresholds=TH, **_confirm_services() ) - 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) + check("重放返回 skipped", again["status"], "skipped") + check("重放 core_trade 仍 2 条", + q1(core_engine, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g", + g=gid), 2) + check("重放后预警单仍 1 张", + q1(agent_engine, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", + c=CUSTOMER), 1) + check("重放后 lot 明细仍 2 条", + q1(core_engine, "SELECT COUNT(*) FROM core_convert_lot_detail " + "WHERE convert_group_id = :g", g=gid), 2) # ── 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( + print("\n【E】无命中场景(独立小额客户)不建单、仅 pass 审计") + small = accept_convert( { - "customer_id": CUSTOMER, + "customer_id": CUSTOMER_SMALL, "from_product_id": PROD_OUT, "to_product_id": PROD_IN, - "qty": "1000", + "qty": Decimal("1000"), "client_request_id": "T8M-IDEM-0002", }, - now=NOW + timedelta(days=1), - thresholds=TH, + now=submit_at, id_factory=_t8m_id, + **_accept_services(), ) - 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) + res3 = confirm_one( + small["convert_group_id"], + now=confirm_at, + as_of=t1_day, + thresholds=TH, + **_confirm_services(), + ) + check("小额转换确认成功", res3["status"], "confirmed") + check("未命中任何规则", res3["triggered_rules"], []) + check("未新建预警单", + q1(agent_engine, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", + c=CUSTOMER_SMALL), 0) + check("落 pass 审计 1 条", + q1(agent_engine, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t " + "AND decision = 'pass'", t=TRACE), 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) + for label, engine, sql, params in [ + ("core_trade", core_engine, + "SELECT COUNT(*) FROM core_trade WHERE customer_id IN (:c1, :c2)", + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}), + ("core_convert_request", core_engine, + "SELECT COUNT(*) FROM core_convert_request WHERE customer_id IN (:c1, :c2)", + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}), + ("core 客户", core_engine, + "SELECT COUNT(*) FROM core_customer WHERE customer_id IN (:c1, :c2)", + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}), + ("risk_alert", agent_engine, + "SELECT COUNT(*) FROM risk_alert WHERE customer_id IN (:c1, :c2)", + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}), + ("customer_profile_l3", agent_engine, + "SELECT COUNT(*) FROM customer_profile_l3 WHERE customer_id IN (:c1, :c2)", + {"c1": CUSTOMER, "c2": CUSTOMER_SMALL}), + ("audit_log", agent_engine, + "SELECT COUNT(*) FROM audit_log WHERE trace_id LIKE 'T8M-VERIFY%'", {}), + ("risk_convert_detail", agent_engine, + "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id LIKE 'CNV-T8M%' " + "OR client_request_id LIKE 'T8M-%'", {}), + ]: + check(f"{label} 零残留", q1(engine, sql, **params), 0) dispose_engines() print("\n" + "=" * 60) - print(f"真库验证(T-8 规则引擎改造):{_passed} 项一致 / {_failed} 项不一致") + print(f"真库验证(T-8 引擎时机):{_passed} 项一致 / {_failed} 项不一致") print("=" * 60) return 1 if _failed else 0 diff --git a/tests/test_convert_accept.py b/tests/test_convert_accept.py index b081f4e..98ceb8e 100644 --- a/tests/test_convert_accept.py +++ b/tests/test_convert_accept.py @@ -233,6 +233,14 @@ def test_accept_success_has_no_side_effects(sqlite_engine): ) assert [a["decision"] for a in audits] == ["accepted"] + # 引擎时机(T-8 / FR-C28 · 验收 29 前半):受理段**不得**触规则引擎 —— + # 引擎只在 T+1 确认事务提交后跑一次(`confirm_service` → `engine_call`)。 + assert _rows(sqlite_engine, "SELECT 1 FROM risk_alert LIMIT 1") == [] + assert _rows( + sqlite_engine, + "SELECT 1 FROM audit_log WHERE decision IN ('engine_error', 'pass')", + ) == [] + # ── 2/3/4. 幂等语义 ───────────────────────────────────────────────── def test_accept_same_client_request_id_returns_same_group(sqlite_engine): diff --git a/tests/test_convert_confirm.py b/tests/test_convert_confirm.py index d44c2cf..4d684e8 100644 --- a/tests/test_convert_confirm.py +++ b/tests/test_convert_confirm.py @@ -18,6 +18,7 @@ sqlite 内存库;真库侧(ENUM 严校验 / DATETIME(3) / 死锁重试)由 from __future__ import annotations +import json from datetime import date, datetime, timedelta from decimal import Decimal @@ -36,6 +37,7 @@ from app.service.convert.errors import ( ConfirmConflict, ConvertRequestNotFound, ) +from app.service.risk.rules import RiskThresholds CUST = "CUST-CC" # C3 客户 → 转入 R4 = 放行 CUST_LOW = "CUST-CCL" # C1 客户 → 转入 R4 = 阻断(T+1 复核用) @@ -599,6 +601,64 @@ def test_engine_exception_does_not_block_confirmation(sqlite_engine): assert res["engine_error"] is False # 异常被吞,未标记 engine_error 字段 +def test_real_engine_on_confirm_dedupes_daily_total(sqlite_engine): + """确认链路接**真引擎**(不注入 hook):一张单、两条事件、RISK-002 只计一次。 + + 验收 5/6/7 的端到端版本 —— 本文件其余用例用 `engine_hook` 注入计数(验证时机), + 这条走真实 `process_convert_event`,证明确认段确实在事务提交后把两条流水 + 投给了引擎,且金额聚合去重生效。 + + 折算值:转出端 68020.00 / 转入端 67074.46(和 135094.46)。阈值夹在 + 「单条」与「两条之和」之间 → 去重后不命中;若聚合未去重则 RISK-002 误报。 + """ + _seed(sqlite_engine) + base = dict( + large_amount=Decimal("1000000"), # 关掉 RISK-001,避免与 RISK-002 混淆 + freq_count=99, + probe_window_minutes=5, + probe_count=99, + probe_amount=Decimal("99999999"), + small_amount=Decimal("1"), + small_count=99, + concentration_threshold=1.01, + ) + + # ① 正证:68020 < 100000 < 135094.46 → 去重后不命中 + gid = _accept(sqlite_engine, cid="CC-ENG-1")["convert_group_id"] + res = _confirm( + sqlite_engine, gid, thresholds=RiskThresholds(daily_total=Decimal("100000"), **base) + ) + assert res["status"] == "confirmed" + assert res["triggered_rules"] == [], "金额聚合未去重时 RISK-002 会误报" + assert res["alert_ids"] == [] + assert _rows(sqlite_engine, "SELECT 1 FROM risk_alert") == [] + + # ② 反证:阈值降到 60000(< 单条 68020)→ 命中,证明引擎确实在确认段跑过 + _exec( + sqlite_engine, + "UPDATE core_share_lot SET remain_qty = qty " + "WHERE customer_id = :c AND product_id = :p", + c=CUST, p=PROD_OUT, + ) + gid2 = _accept(sqlite_engine, cid="CC-ENG-2")["convert_group_id"] + res2 = _confirm( + sqlite_engine, gid2, thresholds=RiskThresholds(daily_total=Decimal("60000"), **base) + ) + assert res2["status"] == "confirmed" + assert "RISK-002" in res2["triggered_rules"] + + alerts = _rows(sqlite_engine, "SELECT * FROM risk_alert") + assert len(alerts) == 1, "一次转换只出一张单" + payload = json.loads(alerts[0]["payload"]) + assert len(payload["events"]) == 2, "一张单承载两条事件" + assert [e["trade_type"] for e in payload["events"]] == ["redeem", "subscribe"] + assert alerts[0]["trade_id"] == res2["out_trade_id"], "主流水口径取转出端" + + # 去重不删行:两单各 2 条流水仍在(验收 6) + assert len(_trades(sqlite_engine, gid)) == 2 + assert len(_trades(sqlite_engine, gid2)) == 2 + + # ── 9. 资金流归属(T-7 明确不写 core_cash_flow,归 T-16)───────────── def test_confirm_writes_no_cash_flow(sqlite_engine): """确认事务**不写** `core_cash_flow` —— 资金流唯一归属 T-16(反双重归属)。 diff --git a/tests/test_convert_engine.py b/tests/test_convert_engine.py index 7f3524d..32d374e 100644 --- a/tests/test_convert_engine.py +++ b/tests/test_convert_engine.py @@ -13,8 +13,10 @@ sqlite 内存库(`conftest.sqlite_engine`),数据自建。阈值一律** from __future__ import annotations +import ast from datetime import datetime, timedelta from decimal import Decimal +from pathlib import Path import pytest from sqlalchemy import text @@ -308,3 +310,46 @@ def test_error_hook_none_keeps_original_exception(sqlite_engine, monkeypatch): out, inn, core_ro=core, risk_repo=RiskRepository(engine=sqlite_engine), thresholds=_th(), ) + + +# ── 4. 引擎调用点唯一性(T-8 DoD 1 / R-7)──────────────────────────── +def test_process_convert_event_has_single_call_site(): + """`process_convert_event(` 在 `app/` 下**恰 1 处调用**,且位于 `engine_call.run_convert_engine`。 + + T-8(FR-C28)把引擎时机从「受理后」挪到「确认事务提交后」,并要求**全仓唯一调用点** —— + 两处调用会让 RISK-002 当日累计翻倍(PRD 验收 5/7)。用 AST 统计 `Call` 节点, + **不靠字符串匹配**:docstring / 注释里提到函数名不算调用(`engine_call.py` 的说明段 + 就含该字样),定义点(`def`)与 `import` 也不是 `Call`。 + """ + root = Path(__file__).resolve().parents[1] + hits: list[str] = [] + for path in sorted((root / "app").rglob("*.py")): + tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path)) + for node in ast.walk(tree): + if not isinstance(node, ast.Call): + continue + func = node.func + name = getattr(func, "id", None) or getattr(func, "attr", None) + if name == "process_convert_event": + hits.append(f"{path.relative_to(root).as_posix()}:{node.lineno}") + + assert len(hits) == 1, ( + "`process_convert_event(` 应恰有 1 处调用点(engine_call.run_convert_engine);" + f"实际 {len(hits)} 处:{hits}" + ) + assert hits[0].startswith("app/service/convert/engine_call.py:"), ( + f"唯一调用点必须落在 engine_call.py,实际:{hits[0]}" + ) + # 调用所在函数体应是 `run_convert_engine`(防「文件对但函数被搬走」) + tree = ast.parse( + (root / "app/service/convert/engine_call.py").read_text(encoding="utf-8") + ) + owner = None + for node in ast.walk(tree): + if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef)): + for child in ast.walk(node): + if isinstance(child, ast.Call): + name = getattr(child.func, "id", None) or getattr(child.func, "attr", None) + if name == "process_convert_event": + owner = node.name + assert owner == "run_convert_engine", f"调用点宿主函数应为 run_convert_engine,实际 {owner!r}"