diff --git a/app/gateway/trade_gateway.py b/app/gateway/trade_gateway.py index 7c99660..b4a7503 100644 --- a/app/gateway/trade_gateway.py +++ b/app/gateway/trade_gateway.py @@ -131,7 +131,7 @@ def _submit_convert( `confirm_service.confirm_batch` 完成(PRD §2.6 未知价法)。 ⚠️ **T-9 变更**:此前调 v1.0 `convert_fund`(实时八步编排,含阶段 1.5 引擎)。 - 改调受理后,`convert_fund` 的**生产调用点已摘除**(该函数随之退役, + 改调受理后,`convert_fund` 的生产调用点已摘除并**随 T-13 退役删除**( 见开发计划 R-回归 2/7 —— 其测试面重写归后续任务,本任务不删函数体)。 `ConvertRepository` / `ConvertRequestRepository` **按模块级符号引用** diff --git a/app/repository/convert_repository.py b/app/repository/convert_repository.py index 6316a69..8eaf911 100644 --- a/app/repository/convert_repository.py +++ b/app/repository/convert_repository.py @@ -198,7 +198,7 @@ class ConvertRepository: 置回 `pending` 是正确语义:同一 `group_id` 的这一次尝试正在进行中。 到达此处时既有行的状态只可能是 `pending`(上一轮中途崩溃)/ `failed`(阶段一或阶段二失败); - `completed` 已在 `convert_fund` 幂等前置分支返回,不会走到这里。 + `completed` 已在幂等前置分支返回,不会走到这里。 **为什么还要 catch IntegrityError**:`convert:idem:{cid}` 锁只包住幂等判定 (出块即释放),两笔同键请求可能**都判定为"无占位"**、各自生成了不同的 `group_id` diff --git a/app/service/convert/convert_service.py b/app/service/convert/convert_service.py index 5bc5ad8..5faca61 100644 --- a/app/service/convert/convert_service.py +++ b/app/service/convert/convert_service.py @@ -1,53 +1,38 @@ -"""基金转换编排(`convert_service`)· T-7 · 关键路径。 +"""基金转换编排(`convert_service`)· T+1 受理/确认分离模型。 -**八步顺序(PRD §7.0 固定顺序,前四步不落库)** +**本模块承载(T-13 起 v1.0 两阶段实时入口 `convert_fund` 已退役删除)** -``` -① 参数与产品校验 同产品 / can_redeem / can_subscribe / 同管理人+同 TA → 4xx,不占位 -② 份额校验 Σ core_share_lot.remain_qty(权威源,非 holding.qty) - + 最低转出份额(全额豁免)+ 批次数上限 → 4xx,不占位 -③ 净值取数与折算 纯函数 calc.py;无净值 → 503 NAV_NOT_READY → 不占位 -④ 适当性校验 仅**转入端**(FR-C5「转换即销售」)→ blocked → R-02 预警+审计 → return,不占位 -⑤ 阶段零 try_lock("convert:idem:{cid_req}") + agent 库占位 -⑥ 阶段一 apply_convert(core 库单事务) -⑦ 阶段 1.5 同步跑规则引擎(D17):异常不阻断已成立的交易 -⑧ 阶段二 回写 completed + 主审计(+ nav_stale 副审计) -``` +- `accept_convert`:T 日**受理**(八步,PRD §7.0:校验 + 落受理单 + 镜像 + 审计; + 不扣份额、不折算、不写流水 —— 扣减与折算全部推迟到 T+1 `confirm_service`); +- `cancel_convert`:T 日撤单(两道闸门:状态 + 撤单窗口 `cancel_before`); +- `compensate_convert`:补偿单点(凭 Core 三件套补 agent 镜像 + 预警,T-12); +- `rebuild_convert_response`:确认后查询的读路径(由 Core 侧重建折算结果, + `convert_admin.py` 查询接口在役消费); +- `t1_t2_dates` / `PROCESSING`:受理回执辅助与 202 语义常量。 -**为什么 blocked 与 4xx 都必须在占位之前(PRD §7.0)**:阶段零一旦占位, -失败就会留下 `pending` 孤儿;把纯校验前置后,这些路径**根本不产生持久化**, -不需要任何清理。 +**幂等锚点(受理段)**:`client_request_id` → Core 库 `uk_idem` 唯一键, +命中既有单直接回执(`_accept_idempotent`);同键并发由 +`convert_req_lock_key(customer)` 客户级锁 + `uk_idem` 兜底收敛(R-2 双守卫)。 -**幂等的两个锚点** -- `client_request_id` → `uk_idem`(agent 库唯一键)兜底重复提交; -- `convert_group_id`(**阶段零预生成**)→ 判定「阶段一是否已成」, - 杜绝阶段二失败后重试产生**第二组流水**(PRD §7.4 v0.3 缺陷)。 - -**响应体**:PRD §5.3 字段,全部 `Decimal → str`(架构 §1 原则 11)。 -未抢到执行权时返回 `{"status": "processing", "convert_group_id": ...}`, -由 T-9(`api/simulate.py`)映射为 **HTTP 202**。 +**响应体**:PRD §5.3.1 受理回执字段,全部 `Decimal → str`(架构 §1 原则 11)。 +`PROCESSING` 由 T-9(`api/simulate.py`)映射为 **HTTP 202**。 """ from __future__ import annotations import logging from dataclasses import dataclass -from datetime import date, datetime, time, timedelta -from decimal import ROUND_HALF_UP, Decimal +from datetime import date, datetime, time +from decimal import Decimal from typing import Any, Callable from uuid import uuid4 from sqlalchemy.exc import IntegrityError from app.config.settings import settings -from app.gateway.convert_core_repository import ( - ConvertApplyInput, - ConvertCoreRepository, - LotCharge, -) +from app.gateway.convert_core_repository import LotCharge from app.repository.convert_repository import ConvertRepository from app.repository.convert_request_repository import ( - FINAL_STATUSES, REMARK_FULL_TRANSFER, STATUS_ACCEPTED, STATUS_CANCELLED, @@ -63,8 +48,6 @@ from app.service.convert.calc import ( ensure_batch_limit, hold_days, in_qty, - lot_amount, - lot_fee, plan_lots, round2, rounding_diff, @@ -79,14 +62,11 @@ from app.service.convert.errors import ( CancelNotAllowed, ConvertRequestNotFound, CrossEntityNotSupported, - IdempotencyUnavailable, InsufficientShares, - NavNotReady, ProductNotRedeemable, ProductNotSubscribable, SameProduct, ) -from app.service.convert.fee import pick_fee_rate from app.service.convert.format import D2 as _D2, D4 as _D4, q as _fmt_q from app.service.convert.trading_calendar import ( next_biz_day, @@ -98,7 +78,7 @@ from app.service.risk.alert_service import record_suitability_alert from app.service.risk.locks import convert_req_lock_key, run_locked, try_lock from app.service.risk.rules import RiskThresholds from app.service.suitability import suitability_check -from app.utils.trace import current_trace, ensure_trace, new_trace +from app.utils.trace import ensure_trace logger = logging.getLogger(__name__) @@ -255,85 +235,6 @@ def _plan_shares( return plan, available -# ── ②③ 份额校验 + 净值折算 ────────────────────────────────────────── -def _plan_and_quote( - core: CoreReadOnlyRepository, - req: dict[str, Any], - out_product: dict[str, Any], - in_product: dict[str, Any], - now: datetime, -) -> _Quote: - from_pid = str(req["from_product_id"]) - to_pid = str(req["to_product_id"]) - trade_date = now.date() - plan, _physical = _plan_shares(core, req, out_product) - - # ── 净值:转出端用**各批次自身成交净值**,转入端取 T 日净值(未知价法)── - rules = [FeeRule.from_row(r) for r in core.get_redeem_fee_rules(from_pid)] - charges: list[LotCharge] = [] - for alloc in plan.allocations: - days = hold_days(trade_date, alloc.confirmed_at) - rate = pick_fee_rate(rules, days, product_id=from_pid) - amount = lot_amount(alloc.qty, alloc.nav) - charges.append( - LotCharge( - lot_id=alloc.lot_id, - qty=alloc.qty, - hold_days=days, - amount=amount, - fee_rate=rate, - fee_amount=lot_fee(amount, rate), - nav=alloc.nav, - nav_date=trade_date, - ) - ) - - in_nav_row = core.get_nav_as_of(to_pid, trade_date) - if in_nav_row is None: - raise NavNotReady(f"转入基金 {to_pid} 尚无 {trade_date} 当日或之前的净值") - in_nav = to_decimal(in_nav_row["nav"]) - nav_date = in_nav_row["nav_date"] - if not isinstance(nav_date, date): # sqlite 读回为字符串 - nav_date = date.fromisoformat(str(nav_date)[:10]) - - out_amount = sum((c.amount for c in charges), Decimal("0")) - redeem_fee = sum((c.fee_amount for c in charges), Decimal("0")) - conv = convert_amount(out_amount, redeem_fee) - out_rate = to_decimal(out_product.get("subscribe_fee_rate") or 0) - in_rate = to_decimal(in_product.get("subscribe_fee_rate") or 0) - gap = diff_fee(conv, out_rate, in_rate, settings.convert_diff_fee_mode) - in_amount = conv - gap - shares = in_qty(in_amount, in_nav) - - # 转出端展示净值 = 金额 ÷ 份额(加权平均;计费仍逐批用各自 nav) - out_nav = round2(out_amount / plan.actual_qty) if plan.actual_qty else Decimal("0") - - return _Quote( - plan=plan, - charges=tuple(charges), - out_nav=out_nav, - out_amount=out_amount, - redeem_fee=redeem_fee, - convert_amount=conv, - diff_fee=gap, - in_amount=in_amount, - in_qty=shares, - rounding_diff=rounding_diff(in_amount, in_nav, shares), - in_nav=in_nav, - nav_date=nav_date, - nav_stale=(trade_date - nav_date).days > settings.convert_nav_stale_days, - out_subscribe_fee_rate=out_rate, - in_subscribe_fee_rate=in_rate, - ) - - -def _pick_rate(rules: list[FeeRule], days: int, product_id: str) -> Decimal: - """持有天数 → 赎回费率(委托 `fee.pick_fee_rate`,无命中即 500 FeeRuleMissing)。""" - from app.service.convert.fee import pick_fee_rate - - return pick_fee_rate(rules, days, product_id=product_id) - - def _build_response( req: dict[str, Any], group_id: str, @@ -503,7 +404,7 @@ def accept_convert( ====== ================================================================== `req`:`{customer_id, from_product_id, to_product_id, qty, client_request_id?}` - (与 `convert_fund` 同形,便于 T-9 双入口并存)。 + (字段面与 v1.0 请求同形,`trade_gateway` 分派与测试共用)。 """ core = core_ro or CoreReadOnlyRepository() repo = risk_repo or RiskRepository() @@ -783,297 +684,6 @@ def cancel_convert( } -# ── v1.0 旧三阶段主入口(⛔ Deprecated · T-9 切换 API 后删除)────────── -def convert_fund( - req: dict[str, Any], - core_ro: CoreReadOnlyRepository | None = None, - risk_repo: RiskRepository | None = None, - convert_repo: ConvertRepository | None = None, - core_writer: ConvertCoreRepository | None = None, - thresholds: RiskThresholds | None = None, - now: datetime | None = None, - actor_id: str | None = None, - *, - engine_hook: Callable[[dict[str, Any], dict[str, Any]], dict[str, Any]] | None = None, - id_factory: Callable[[str, datetime], str] | None = None, -) -> dict[str, Any]: - """执行一次基金转换(PRD §7.0 八步)。 - - `req`:`{customer_id, from_product_id, to_product_id, qty, client_request_id?}`。 - `engine_hook`:阶段 1.5 的注入点(T-8 未落地时传假函数;缺省自动尝试 - `engine.process_convert_event`,不存在则跳过并记录 warning)。 - `id_factory`:`(前缀, now) -> id`,测试注入点(与 `trade_gateway` 同款)。 - """ - core = core_ro or CoreReadOnlyRepository() - repo = risk_repo or RiskRepository() - crepo = convert_repo or ConvertRepository() - writer = core_writer or ConvertCoreRepository() - th = thresholds or RiskThresholds.from_settings() - ensure_trace() - - now = now or datetime.now() - new_id = id_factory or _new_id - customer_id = str(req["customer_id"]) - cid_req = req.get("client_request_id") or None - - # ── 幂等前置(必须先于 ①②③④,实施期修正)── - # 同 client_request_id 重试时,首次已扣减 core_share_lot 份额,若先跑 ② plan_lots - # 会误报 InsufficientShares;故先判定「已完成 / 阶段一成」并直接返回首次结果。 - if cid_req is None: - # 未带幂等键:免占位直跑(PRD §7.3 / Q7,无幂等语义) - group_id = new_id("CNV", now) - placeholder = False - else: - with try_lock( - f"convert:idem:{cid_req}", settings.convert_lock_ttl_seconds - ) as acquired: - if not acquired: - # 有并发请求正在执行 → 202,不查占位、不进阶段一 - return {"status": PROCESSING, "convert_group_id": None} - existing = crepo.get_by_client_request_id(cid_req) - if existing is not None: - hit_gid = str(existing["convert_group_id"]) - if str(existing["status"]) == "completed": - rebuilt = rebuild_convert_response(hit_gid, core_ro=core) - if rebuilt is not None: - return rebuilt - # 阶段一已成、阶段二未成 → 只补跑阶段二(加锁防并发重试审计双写) - if core.has_convert_trades(hit_gid): - with try_lock( - f"convert:rerun:{hit_gid}", settings.convert_lock_ttl_seconds - ): - rebuilt = _finalize_from_core( - req, hit_gid, core, repo, crepo, now, actor_id - ) - if rebuilt is not None: - return rebuilt - # 阶段一未成:必须区分「确定没跑成」与「在飞/未知」,否则同键并发会双跑 - # —— `convert:idem:` 锁只包住本段幂等判定(出块即释放),窗口内重入会 - # 与在飞的那笔**同时进入阶段一**(T-13 真库确定性交错实测: - # 未映射的 IntegrityError 直穿 → 500;若日后 `in_lot_id` 的派生规则被 - # 改成随机值,同一窗口会升级为**真·双扣**)。 - # · `failed`:阶段一定性失败(无流水)→ 复用同一 group_id 重跑, - # 这正是架构 §8.3「LOT_CONFLICT → 调用方退避重试」的服务端契约; - # · `expired`:占位已被 SLA 巡检判死(`cleanup_pending_convert.py`,S2) - # → 同样可安全复用; - # · `pending`:在飞与崩溃**不可区分** → 按并发处理,返回 202 交客户端 - # 稍后重试(架构 §9「同键并发 → 202」);崩溃残留由 SLA 巡检置 expired - # 后自动放行(急用可 `cleanup_pending_convert.py --hours 0` 立即判死)。 - if str(existing["status"]) in ("failed", "expired"): - group_id = hit_gid - else: - return {"status": PROCESSING, "convert_group_id": hit_gid} - else: - group_id = new_id("CNV", now) - placeholder = True - - # ① 参数与产品校验(不落库) - out_product, in_product = _validate_products(core, req) - # ② 份额校验 + ③ 净值取数与折算(不落库) - quote = _plan_and_quote(core, req, out_product, in_product, now) - - # ④ 适当性校验(转入端 · 唯一的业务阻断点)→ blocked 时**不占位** - suit = suitability_check( - customer_id, - str(req["to_product_id"]), - core_ro=core, - risk_repo=repo, - check_source="r02_trade", - actor_id=actor_id, - request_ref=group_id, - ) - if suit.blocked: - record_suitability_alert( - { - "trade_id": group_id, - "customer_id": customer_id, - "product_id": str(req["to_product_id"]), - "trade_type": "convert", - "amount": str(quote.in_amount), - "traded_at": str(now), - }, - rule_id=suit.rule_id, - block_reason=suit.block_reason, - risk_repo=repo, - ) - _audit( - repo, - decision="suitability_blocked", - group_id=group_id, - customer_id=customer_id, - rule_id=suit.rule_id, - actor_id=actor_id, - summary={ - "from_product_id": req["from_product_id"], - "to_product_id": req["to_product_id"], - "requested_qty": _q(quote.plan.requested_qty, _D2), - "block_reason": suit.block_reason, - "block_response_code": suit.block_response_code, - "reasons": list(suit.reasons), - }, - ) - return { - "blocked": True, - "convert_group_id": group_id, - "match_result": suit.match_result, - "mismatch_type": suit.mismatch_type, - "requires_disclosure": suit.requires_disclosure, - "needs_branch_confirm": suit.needs_branch_confirm, - "block_reason": suit.block_reason, - "block_response_code": suit.block_response_code, - "rule_refs": suit.rule_refs, - "reasons": list(suit.reasons), - "advice": "请联系持证投资顾问", - "notice": "本次请求已记录", - } - - # ⑤ 阶段零:占位(仅带幂等键时;在全部 4xx 之后,故 4xx 不留占位) - if placeholder: - try: - owns = crepo.insert_placeholder(group_id, cid_req) - except Exception as exc: # noqa: BLE001 - # 占位失败 = 无法保证幂等 → 不放行(PRD §7.3) - logger.exception("convert 占位失败:%s", group_id) - raise IdempotencyUnavailable(f"幂等占位失败:{exc}") from exc - if not owns: - # 同一 client_request_id 的并发请求已占住 uk_idem(各自生成过不同 group_id) - # → 本笔让路,回 202(架构 §9「同键并发 → 202」),由客户端稍后重试并发起幂等命中。 - # 不这样收敛的话,此处会直穿未映射的 IntegrityError(503)。(T-13 真库实测) - logger.info("同键并发占位让路:cid=%s gid=%s", cid_req, group_id) - return {"status": PROCESSING, "convert_group_id": None} - - out_trade_id = new_id("TRD", now) - in_trade_id = new_id("TRD", now) - apply_input = ConvertApplyInput( - convert_group_id=group_id, - out_trade_id=out_trade_id, - in_trade_id=in_trade_id, - customer_id=customer_id, - from_product_id=str(req["from_product_id"]), - to_product_id=str(req["to_product_id"]), - traded_at=now, - out_qty=quote.plan.actual_qty, - out_amount=quote.out_amount, - in_qty=quote.in_qty, - in_amount=quote.in_amount, - in_nav=quote.in_nav, - in_nav_date=quote.nav_date, - in_lot_id=f"LOT-{group_id}-IN", - in_confirmed_at=now + timedelta(days=settings.convert_confirm_offset_days), - charges=quote.charges, - ) - - # ⑥ 阶段一:core 库单事务(失败 → 占位置 failed,供人工补偿) - try: - writer.apply_convert(apply_input) - except Exception: - logger.exception("convert 阶段一失败:%s", group_id) - if cid_req is not None: - try: - crepo.mark_failed(group_id) - except Exception: # noqa: BLE001 - logger.exception("占位标记 failed 失败(不影响原始异常):%s", group_id) - raise - - # ⑦ 阶段 1.5:同步跑规则引擎(D17:异常不阻断已成立的交易) - engine_result: dict[str, Any] | None = None - try: - engine_result = _run_engine( - { - "trade_id": out_trade_id, - "customer_id": customer_id, - "product_id": str(req["from_product_id"]), - "trade_type": "redeem", - "amount": quote.out_amount, - "qty": quote.plan.actual_qty, - "trade_status": "confirmed", - "traded_at": now, - "convert_group_id": group_id, - }, - { - "trade_id": in_trade_id, - "customer_id": customer_id, - "product_id": str(req["to_product_id"]), - "trade_type": "subscribe", - "amount": quote.in_amount, - "qty": quote.in_qty, - "trade_status": "confirmed", - "traded_at": now, - "convert_group_id": group_id, - }, - core_ro=core, - risk_repo=repo, - thresholds=th, - engine_hook=engine_hook, - ) - except Exception: # noqa: BLE001 - logger.exception("convert 阶段 1.5 引擎失败(不阻断交易):%s", group_id) - _audit( - repo, - decision="engine_error", - group_id=group_id, - customer_id=customer_id, - actor_id=actor_id, - summary={"error_stage": "process_convert_event"}, - ) - engine_result = {"triggered_rules": [], "alert_ids": [], "aml_hit": False, - "engine_error": True} - - response = _build_response( - req, group_id, quote, out_trade_id=out_trade_id, in_trade_id=in_trade_id, - engine_result=engine_result, - ) - - # ⑧ 阶段二:回写 completed + 主审计(失败**不回滚 Core**) - lo, hi = quote.hold_days_range - try: - crepo.complete_convert( - group_id, - out_trade_id=out_trade_id, - in_trade_id=in_trade_id, - related_trade_id=out_trade_id, - nav=quote.in_nav, - nav_date=quote.nav_date, - fee_amount=quote.redeem_fee, - hold_days_min=lo, - hold_days_max=hi, - nav_stale=quote.nav_stale, - ) - _write_main_audit( - req, group_id, quote, out_trade_id, in_trade_id, repo, now, actor_id, engine_result - ) - except Exception: # noqa: BLE001 - # 交易已成立:只能留痕 + 本地日志兜底(PRD §7.1 第三轮第 8 条) - logger.exception( - "convert 阶段二失败(交易已成立,待补偿)group_id=%s quote=%s", - group_id, - { - "out_amount": _q(quote.out_amount, _D2), - "in_amount": _q(quote.in_amount, _D2), - "in_qty": _q(quote.in_qty, _D2), - "out_trade_id": out_trade_id, - "in_trade_id": in_trade_id, - }, - ) - try: - _audit( - repo, - decision="convert_detail_write_failed", - group_id=group_id, - customer_id=customer_id, - actor_id=actor_id, - summary={"out_trade_id": out_trade_id, "in_trade_id": in_trade_id}, - ) - except Exception: # noqa: BLE001 - logger.exception("阶段二失败审计亦写入失败:%s", group_id) - if cid_req is not None: - try: - crepo.mark_failed(group_id) - except Exception: # noqa: BLE001 - logger.exception("占位标记 failed 失败:%s", group_id) - return response - - # ── 阶段 1.5 的引擎调用(D17)────────────────────────────────────── def _run_engine( out_trade: dict[str, Any], @@ -1099,104 +709,6 @@ def _run_engine( ) -def _finalize_from_core( - req: dict[str, Any], - group_id: str, - core: CoreReadOnlyRepository, - repo: RiskRepository, - crepo: ConvertRepository, - now: datetime, - actor_id: str | None, -) -> dict[str, Any] | None: - """补跑阶段二:仅凭 Core 侧数据回填详情 + 审计(阶段一已成、阶段二未成)。 - - **幂等窗口闭合(验收 15)**:阶段二失败后带同键重试 → 不重跑阶段一, - 不产生第二组流水,RISK-002 当日累计也不翻倍。 - """ - rebuilt = _rebuild_quote(group_id, core) - if rebuilt is None: - return None - quote, out_trade_id, in_trade_id = rebuilt - lo, hi = quote.hold_days_range - try: - crepo.complete_convert( - group_id, - out_trade_id=out_trade_id, - in_trade_id=in_trade_id, - related_trade_id=out_trade_id, - nav=quote.in_nav, - nav_date=quote.nav_date, - fee_amount=quote.redeem_fee, - hold_days_min=lo, - hold_days_max=hi, - nav_stale=quote.nav_stale, - ) - _write_main_audit( - req, group_id, quote, out_trade_id, in_trade_id, repo, now, actor_id, None - ) - except Exception: # noqa: BLE001 - logger.exception("补跑阶段二失败:%s", group_id) - return _build_response( - req, group_id, quote, out_trade_id=out_trade_id, in_trade_id=in_trade_id - ) - - -def _write_main_audit( - req: dict[str, Any], - group_id: str, - quote: _Quote, - out_trade_id: str, - in_trade_id: str, - repo: RiskRepository, - now: datetime, - actor_id: str | None, - engine_result: dict[str, Any] | None, -) -> None: - """主审计 + `nav_stale` 副审计(PRD §7.3:实际为 1~2 条)。""" - _audit( - repo, - decision="convert_accepted", - group_id=group_id, - customer_id=str(req["customer_id"]), - actor_id=actor_id, - summary={ - "from_product_id": req.get("from_product_id"), - "to_product_id": req.get("to_product_id"), - "requested_qty": _q(quote.plan.requested_qty, _D2), - "actual_qty": _q(quote.plan.actual_qty, _D2), - "forced_full_transfer": quote.plan.forced_full_transfer, - "out_amount": _q(quote.out_amount, _D2), - "redeem_fee": _q(quote.redeem_fee, _D2), - "convert_amount": _q(quote.convert_amount, _D2), - "diff_fee": _q(quote.diff_fee, _D2), - "in_amount": _q(quote.in_amount, _D2), - "in_qty": _q(quote.in_qty, _D2), - "rounding_diff": _q(quote.rounding_diff, _D4), - "nav": _q(quote.in_nav, _D4), - "nav_date": str(quote.nav_date), - "nav_stale": quote.nav_stale, - "estimated": True, - "lot_count": quote.plan.batch_count, - "out_trade_id": out_trade_id, - "in_trade_id": in_trade_id, - **(dict(engine_result or {})), - }, - ) - if quote.nav_stale: - _audit( - repo, - decision="nav_stale", - group_id=group_id, - customer_id=str(req["customer_id"]), - actor_id=actor_id, - summary={ - "nav_date": str(quote.nav_date), - "trade_date": str(now.date()), - "stale_days": (now.date() - quote.nav_date).days, - }, - ) - - def _out_nav_rebuild( details: list[Any], out_amount: Decimal, out_qty_val: Decimal ) -> Decimal: @@ -1538,17 +1050,14 @@ def _as_date(value: Any) -> date: __all__ = [ - # T-6/T-7 新增的 T+1 主入口(T-9 补登记 —— 此前漏进 __all__,通配导入取不到) + # T+1 主入口(T-6/T-7 新增;v1.0 实时入口 convert_fund 已随 T-13 退役删除) "accept_convert", "cancel_convert", "t1_t2_dates", - # T-12 补偿单点(T+1 改写后仍是「凭 Core 数据补 agent 写入」的唯一入口; - # v0.x 曾随 v1.0 标 Deprecated —— T+1 下其详情侧已切 `sync_mirror` + confirmed - # 审计,与确认段第 ⑧ 步同口径,恢复为在役函数) + # T-12 补偿单点(T+1 改写后仍是「凭 Core 数据补 agent 写入」的唯一入口) "compensate_convert", - # v1.0 实时链路(⛔ Deprecated,待退役;生产调用点已由 T-9 摘除; - # _finalize_from_core 仅剩 convert_fund 幂等重试一个调用点) - "convert_fund", + # 确认后查询读路径(convert_admin 查询接口在役消费) "rebuild_convert_response", + # 202 语义常量(api/simulate.py 在役消费) "PROCESSING", ] diff --git a/scripts/dev/verify_convert_lots.py b/scripts/dev/verify_convert_lots.py index 56f34c0..700bb51 100644 --- a/scripts/dev/verify_convert_lots.py +++ b/scripts/dev/verify_convert_lots.py @@ -189,11 +189,12 @@ def assert_subscribe(admin, writer, core) -> None: def assert_redeem_fifo(admin, writer, core) -> None: print("【B/C】赎回 FIFO 扣减(★ 真库验证 `(:q + 0.0)` 方言语义)") - # 120 元 ÷ 1.0 = 120 份 → C1 扣满 100 归零,C2 再扣 20 + # D26/R-6(T-9 起):赎回**份额申报**——直接申报 120 份 → C1 扣满 100 归零,C2 再扣 20 + # (T-13 全量复跑修正:v1.0 的 amount= 入参会被 redeem 分支静默跳过) _maintain_lots( writer=writer, core=core, trade_id="TRD-LOTT-RED", customer_id=CUSTOMER, product_id=PROD_REDEEM, - trade_type="redeem", amount=Decimal("120"), traded_at=TRADE_AT, + trade_type="redeem", qty=Decimal("120"), traded_at=TRADE_AT, ) check( "C1(最老批)remain_qty 归零 —— 扣减真库生效", @@ -220,10 +221,11 @@ def assert_redeem_fifo(admin, writer, core) -> None: def assert_bootstrap(admin, writer, core) -> None: print("【D】D8 兜底补建:无批次有持仓 → 按 as_of 补建后再扣") + # D26/R-6:赎回份额申报(v1.0 的 amount= 入参会被 redeem 分支静默跳过) _maintain_lots( writer=writer, core=core, trade_id="TRD-LOTT-BOOT", customer_id=CUSTOMER, product_id=PROD_BOOT, - trade_type="redeem", amount=Decimal("30"), traded_at=TRADE_AT, + trade_type="redeem", qty=Decimal("30"), traded_at=TRADE_AT, ) check( "补建批次行数", diff --git a/scripts/dev/verify_convert_service.py b/scripts/dev/verify_convert_service.py deleted file mode 100644 index a86b1b3..0000000 --- a/scripts/dev/verify_convert_service.py +++ /dev/null @@ -1,402 +0,0 @@ -"""T-7 真 MySQL 验证脚本:`convert_service.convert_fund` 八步编排实跑 + DoD 断言。 - -**为什么 sqlite 单测全绿还不够(T-7 视角)** - -sqlite 单测验证了八步顺序与折算值,但**证明不了**以下只有真 MySQL(InnoDB)才成立、 -且正是本次实施期修正点的事项: - -1. **`complete_convert` 三步法 upsert(R-a)在真库成立**:无 `client_request_id` 时 - 不走占位、`complete_convert` 必须**先查再 INSERT** 完成行。sqlite 单测因 happy path - 不带键、从未触达「无占位即 INSERT」分支,看不出 MySQL 下是否真的能落行。 - 本次修复前该分支缺失 → `risk_convert_detail` 0 行(已被单测暴露)。 -2. **`uk_idem` 唯一约束真实存在并生效**:幂等「同键只出一组流水」除了代码判定, - 还依赖 `risk_convert_detail.client_request_id` 的 UNIQUE——真库迁移是否真的建了这条 - 约束、重复插入是否真报 IntegrityError,sqlite 单测无法证明。 -3. **`idx_convert_group` 真实存在**:④ 重试判定 `has_convert_trades` 依赖它; - 约束缺失会让重试扫描全表或判错。 -4. **DECIMAL(18,2/18,4) 精度**:`risk_convert_detail.fee_amount` / `nav` 在 MySQL 下 - 落库零漂移(sqlite 用 REAL 无此保证)。 -5. **阶段二失败兜底不回滚 Core**:本脚本覆盖 happy path 与幂等,4xx 不落库亦验证。 - -用法: - python scripts/dev/verify_convert_service.py # 建隔离数据 → 跑断言 → 清理 - -约定(与 T-6 `verify_convert_apply.py` 一致): -- 用 **T7M 前缀**的隔离数据(客户/产品/批次/group_id),跑完**两个库**(core + agent)全清, - 不碰既有种子; -- 建/清数据走 `role="admin"`(需 DELETE);事务本身走 `convert_fund` 默认的 ro/rw 账号; -- 固定 `trace_id = "T7M-VERIFY-TRACE"`,便于精准清理 `audit_log`(agent 库)。 -""" - -from __future__ import annotations - -import sys -import threading -import uuid -from datetime import date, datetime, timedelta -from decimal import ROUND_HALF_UP, Decimal -from pathlib import Path - -from sqlalchemy import text - -ROOT = Path(__file__).resolve().parents[2] -sys.path.insert(0, str(ROOT)) - -from app.config.settings import settings # noqa: E402 -from app.repository.convert_repository import ConvertRepository # noqa: E402 -from app.service.convert.calc import ( # noqa: E402 - convert_amount, - diff_fee, - hold_days, - in_qty, - lot_amount, - lot_fee, - plan_lots, -) -from app.service.convert.convert_service import convert_fund, PROCESSING # noqa: E402 -from app.service.convert.errors import InsufficientShares # noqa: E402 -from app.service.convert.fee import pick_fee_rate # noqa: E402 -from app.service.convert.types import FeeRule, Lot # noqa: E402 -from app.utils.db import dispose_engines, get_engine # noqa: E402 -from app.utils.trace import new_trace # noqa: E402 - -CUSTOMER = "CUST-T7M" -PROD_OUT = "PROD-T7MO" -PROD_IN = "PROD-T7MI" -COMPANY = "华夏模拟基金" -TA = "TA-CN-001" -TRADE_AT = datetime(2026, 9, 4, 10, 0, 0) -TRADE_DATE = date(2026, 9, 4) -NOW = TRADE_AT -OUT_RATE = Decimal("0.0030") -IN_RATE = Decimal("0.0080") -IN_NAV = Decimal("0.9500") -FEE_TIERS = [ - (0, 7, "0.0150"), (7, 30, "0.0100"), (30, 180, "0.0050"), - (180, 365, "0.0025"), (365, None, "0.0000"), -] -TRACE = "T7M-VERIFY-TRACE" -TRACE_E = "T7M-VERIFY-E" # 隔离 E 组审计,便于按 trace 精确计数 - -_passed = 0 -_failed = 0 - - -def check(name: str, actual, expected) -> None: - global _passed, _failed - ok = actual == expected - if ok: - _passed += 1 - else: - _failed += 1 - flag = "✅" if ok else "❌" - print(f" {flag} {name}: 实际 {actual!r}" + ("" if ok else f" / 期望 {expected!r}")) - - -def dec(value, places: str = "0.01") -> Decimal: - return Decimal(str(value)).quantize(Decimal(places), rounding=ROUND_HALF_UP) - - -def q1(engine, sql: str, **params): - with engine.connect() as conn: - return conn.execute(text(sql), params).scalar() - - -def rows(engine, sql: str, **params): - with engine.connect() as conn: - return [dict(r) for r in conn.execute(text(sql), params).mappings()] - - -# ── 数据准备 / 清理(两个库) ──────────────────────────────────────── -def seed_core(engine) -> None: - with engine.begin() as conn: - conn.execute( - text( - "INSERT INTO core_customer (customer_id, display_name, open_date) " - "VALUES (:c, 'T7真库验证', :d)" - ), - {"c": CUSTOMER, "d": TRADE_DATE}, - ) - conn.execute( - text( - "INSERT INTO core_customer_risk (customer_id, risk_code, evaluated_at, expires_at) " - "VALUES (:c, 'C3', :t, :exp)" - ), - {"c": CUSTOMER, "t": TRADE_AT - timedelta(days=30), - "exp": TRADE_AT + timedelta(days=300)}, - ) - for pid, name, ptype, rate in [ - (PROD_OUT, "T7转出基金", "bond", OUT_RATE), - (PROD_IN, "T7转入基金", "stock", IN_RATE), - ]: - conn.execute( - text( - "INSERT INTO core_product (product_id, product_name, min_risk_code, " - "product_type, can_subscribe, can_redeem, subscribe_fee_rate, " - "fund_company, ta_code) " - "VALUES (:p, :n, 'R2', :t, 1, 1, :r, :co, :ta)" - ), - {"p": pid, "n": name, "t": ptype, "r": str(rate), "co": COMPANY, "ta": TA}, - ) - for mh, mh_max, rate in FEE_TIERS: - conn.execute( - text( - "INSERT INTO core_fee_rule (product_id, fee_type, min_hold_days, " - "max_hold_days, rate) VALUES (:p, 'redeem', :mh, :mh_max, :rate)" - ), - {"p": PROD_OUT, "mh": mh, "mh_max": mh_max, "rate": rate}, - ) - conn.execute( - text( - "INSERT INTO core_product_nav (product_id, nav, daily_chg_pct, nav_date) " - "VALUES (:p, :nav, 0, :d)" - ), - {"p": PROD_IN, "nav": str(IN_NAV), "d": TRADE_DATE}, - ) - for lot_id, qty, nav, confirmed in [ - ("LOT-T7M-A1", "100", "1.0300", datetime(2026, 8, 1, 10, 0, 0)), - ("LOT-T7M-A2", "50", "1.0000", datetime(2026, 9, 1, 10, 0, 0)), - ]: - conn.execute( - text( - "INSERT INTO core_share_lot (lot_id, customer_id, product_id, qty, " - "remain_qty, nav, confirmed_at) VALUES (:l, :c, :p, :q, :q, :nav, :cat)" - ), - {"l": lot_id, "c": CUSTOMER, "p": PROD_OUT, "q": qty, "nav": nav, "cat": confirmed}, - ) - conn.execute( - text( - "INSERT INTO core_holding (customer_id, product_id, qty, cost_amount, " - "market_value, pnl_pct, as_of) VALUES (:c, :p, 150, 150.00, 154.50, 0, :d)" - ), - {"c": CUSTOMER, "p": PROD_OUT, "d": TRADE_DATE}, - ) - - -def seed_nav_stale(engine) -> None: - """把转入端净值改成过期(距交易日 10 天 > 阈值 3)以触发 nav_stale 副审计。""" - with engine.begin() as conn: - conn.execute( - text("UPDATE core_product_nav SET nav_date = :d WHERE product_id = :p"), - {"d": date(2026, 8, 25), "p": PROD_IN}, - ) - - -def cleanup_core(engine) -> None: - with engine.begin() as conn: - for sql, params in [ - ("DELETE FROM core_convert_lot_detail WHERE convert_group_id LIKE 'CNV-T7M%'", {}), - ("DELETE FROM core_trade WHERE customer_id = :c", {"c": CUSTOMER}), - ("DELETE FROM core_share_lot WHERE customer_id = :c", {"c": CUSTOMER}), - ("DELETE FROM core_holding WHERE customer_id = :c", {"c": CUSTOMER}), - ("DELETE FROM core_customer_risk WHERE customer_id = :c", {"c": CUSTOMER}), - ("DELETE FROM core_fee_rule WHERE product_id LIKE 'PROD-T7M%'", {}), - ("DELETE FROM core_product_nav WHERE product_id LIKE 'PROD-T7M%'", {}), - ("DELETE FROM core_product WHERE product_id LIKE 'PROD-T7M%'", {}), - ("DELETE FROM core_customer WHERE customer_id = :c", {"c": CUSTOMER}), - ]: - conn.execute(text(sql), params) - - -def cleanup_agent(engine) -> None: - with engine.begin() as conn: - conn.execute( - text("DELETE FROM risk_convert_detail WHERE convert_group_id LIKE 'CNV-T7M%'") - ) - conn.execute( - text("DELETE FROM audit_log WHERE trace_id IN (:t1, :t2)"), - {"t1": TRACE, "t2": TRACE_E}, - ) - - -# ── 期望折算(全部走生产纯函数,与 service 内部同口径) ───────────────── -def compute_expected(engine, requested: str) -> dict: - with engine.connect() as conn: - lot_rows = [ - Lot.from_row(r) - for r in conn.execute( - text( - "SELECT * FROM core_share_lot WHERE customer_id = :c AND product_id = :p " - "ORDER BY confirmed_at ASC, lot_id ASC" - ), - {"c": CUSTOMER, "p": PROD_OUT}, - ).mappings() - ] - rules = [ - FeeRule.from_row(r) - for r in conn.execute( - text( - "SELECT * FROM core_fee_rule WHERE product_id = :p AND fee_type = 'redeem' " - "ORDER BY min_hold_days ASC" - ), - {"p": PROD_OUT}, - ).mappings() - ] - plan = plan_lots(lot_rows, Decimal(requested)) - out_amount = Decimal("0") - redeem_fee = Decimal("0") - for alloc in plan.allocations: - amount = lot_amount(alloc.qty, alloc.nav) - rate = pick_fee_rate(rules, hold_days(TRADE_DATE, alloc.confirmed_at), product_id=PROD_OUT) - out_amount += amount - redeem_fee += lot_fee(amount, rate) - conv = convert_amount(out_amount, redeem_fee) - in_amount = conv - diff_fee(conv, OUT_RATE, IN_RATE, "amount_diff") - return { - "out_amount": out_amount, - "redeem_fee": redeem_fee, - "convert_amount": conv, - "in_amount": in_amount, - "in_qty": in_qty(in_amount, IN_NAV), - "actual_qty": plan.actual_qty, - } - - -def _req(customer: str = CUSTOMER, qty: str = "120", cid_req: str | None = None) -> dict: - return { - "customer_id": customer, - "from_product_id": PROD_OUT, - "to_product_id": PROD_IN, - "qty": qty, - "client_request_id": cid_req, - } - - -def _t7m_id(prefix: str, now: datetime) -> str: - """注入 `convert_fund` 的 id_factory:让 group_id 带 T7M 前缀, - 与 `cleanup_agent` 的 `LIKE 'CNV-T7M%'` 对齐,跑完可精准清理(避免遗留行 - 撞 `uk_group` 唯一约束)。默认 `_new_id` 生成 `CNV-<日期>-`,不含 T7M。""" - return f"CNV-T7M-{uuid.uuid4().hex[:8].upper()}" - - -def main() -> int: - new_trace(TRACE) # 固定 trace,便于清理 audit_log - core_admin = get_engine(settings.mysql_core_database, "admin") - agent_admin = get_engine(settings.mysql_database, "admin") - - run_ts = TRADE_AT.strftime("%H%M%S") # 同一次运行内 cid_req 唯一,避免跨运行锁冲突 - try: - # 初清 + 建隔离数据 - cleanup_core(core_admin) - cleanup_agent(agent_admin) - seed_core(core_admin) - - # ── A. 八步贯通(带幂等键):占位 pending → 2 流水 → completed + 审计 ── - print("\n【A】八步贯通(cid_req 带键):2 流水 + risk_convert_detail 完成 + 审计") - cid_a = f"T7M-REQ-A-{run_ts}" - # 期望折算必须在转换前算(转换会扣减份额,转换后读库会得到不足份额) - exp = compute_expected(core_admin, "120") - first = convert_fund(_req(cid_req=cid_a), now=NOW, actor_id="verify-t7m", id_factory=_t7m_id) - - check("未拦截", first.get("blocked"), False) - check("estimated=True", first.get("estimated"), True) - check("group_id 前缀", first["convert_group_id"].startswith("CNV-"), True) - check("流水条数", q1(core_admin, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g", g=first["convert_group_id"]), 2) - check("R-b trade_type", sorted(r["trade_type"] for r in rows(core_admin, "SELECT trade_type FROM core_trade WHERE convert_group_id = :g", g=first["convert_group_id"])), ["redeem", "subscribe"]) - check("risk_convert_detail 条数", q1(agent_admin, "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id = :g", g=first["convert_group_id"]), 1) - check("risk_convert_detail 状态", q1(agent_admin, "SELECT status FROM risk_convert_detail WHERE convert_group_id = :g", g=first["convert_group_id"]), "completed") - # DECIMAL 精度:真库落库零漂移 - detail = rows(agent_admin, "SELECT fee_amount, nav FROM risk_convert_detail WHERE convert_group_id = :g", g=first["convert_group_id"])[0] - check("fee_amount 精度(2位)", dec(detail["fee_amount"]), dec(exp["redeem_fee"])) - check("nav 精度(4位)", dec(detail["nav"], "0.0001"), dec(exp["in_amount"] / exp["in_qty"], "0.0001")) - # 审计:主审计 1 条;净值新鲜 → 无 nav_stale 副审计 - check("convert_accepted 审计", q1(agent_admin, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'convert_accepted'", t=TRACE), 1) - check("无 nav_stale 审计", q1(agent_admin, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'nav_stale'", t=TRACE), 0) - # 引擎未落地 → 不阻断、无 engine_error 审计 - check("engine_error=False", first.get("engine_error"), False) - check("无 engine_error 审计", q1(agent_admin, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'engine_error'", t=TRACE), 0) - # 折算值对齐生产 calc - check("out_amount", dec(first["out_amount"]), dec(exp["out_amount"])) - check("redeem_fee", dec(first["redeem_fee"]), dec(exp["redeem_fee"])) - check("convert_amount", dec(first["convert_amount"]), dec(exp["convert_amount"])) - check("in_amount", dec(first["in_amount"]), dec(exp["in_amount"])) - check("in_qty", dec(first["in_qty"]), dec(exp["in_qty"])) - check("actual_qty", dec(first["actual_qty"]), dec(exp["actual_qty"])) - - # ── B. 幂等:同键重复 → 返回首次结果、只一组流水 ──────────────── - print("\n【B】幂等同键重试:返回首次结果、不产生第二组流水") - again = convert_fund(_req(cid_req=cid_a), now=NOW, actor_id="verify-t7m", id_factory=_t7m_id) - check("group_id 一致", again["convert_group_id"], first["convert_group_id"]) - check("in_qty 一致", dec(again["in_qty"]), dec(first["in_qty"])) - check("out_trade_id 一致", again["out_trade_id"], first["out_trade_id"]) - check("core_trade 仅 2 条(无第二组)", q1(core_admin, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g", g=first["convert_group_id"]), 2) - check("risk_convert_detail 仅 1 行", q1(agent_admin, "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id = :g", g=first["convert_group_id"]), 1) - - # ── C. 无幂等键:complete_convert 三步法 upsert 必须 INSERT 完成行 ── - print("\n【C】无 cid_req:complete_convert upsert 直接 INSERT 完成行(修复点)") - resp_c = convert_fund(_req(qty="30"), now=NOW, actor_id="verify-t7m", id_factory=_t7m_id) - check("未拦截", resp_c.get("blocked"), False) - check("risk_convert_detail 落 1 完成行", q1(agent_admin, "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id = :g AND status = 'completed'", g=resp_c["convert_group_id"]), 1) - check("core_trade 2 条", q1(core_admin, "SELECT COUNT(*) FROM core_trade WHERE convert_group_id = :g", g=resp_c["convert_group_id"]), 2) - - # ── D. 4xx(份额不足):不占位、不落流水、不审计 ───────────────── - print("\n【D】4xx 分支(InsufficientShares):零残留") - before_detail = q1(agent_admin, "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id LIKE 'CNV-T7M%'") - raised = None - try: - convert_fund(_req(qty="500"), now=NOW, actor_id="verify-t7m", id_factory=_t7m_id) - except InsufficientShares: - raised = "InsufficientShares" - check("抛 InsufficientShares", raised, "InsufficientShares") - check("无新增占位", q1(agent_admin, "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id LIKE 'CNV-T7M%'"), before_detail) - - # ── E. nav_stale 副审计(真库,净值过期) ─────────────────────── - print("\n【E】nav_stale 副审计:净值过期额外落 1 条审计") - cleanup_core(core_admin) - seed_core(core_admin) - seed_nav_stale(core_admin) - cid_e = f"T7M-REQ-E-{run_ts}" - new_trace(TRACE_E) # 切独立 trace,按 trace 精确计数 E 组审计 - resp_e = convert_fund(_req(cid_req=cid_e, qty="120"), now=NOW, actor_id="verify-t7m", id_factory=_t7m_id) - check("nav_stale=True", resp_e.get("nav_stale"), True) - check("nav_stale 审计 1 条", q1(agent_admin, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'nav_stale'", t=TRACE_E), 1) - check("convert_accepted 审计 1 条", q1(agent_admin, "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'convert_accepted'", t=TRACE_E), 1) - - # ── F. 重试判定依赖的索引真实存在 ────────────────────────────── - print("\n【F】约束真实存在(重试判定 / 幂等兜底依赖)") - check( - "core_trade.idx_convert_group 存在", - q1(core_admin, "SELECT COUNT(*) FROM information_schema.statistics " - "WHERE table_schema = :db AND table_name = 'core_trade' " - "AND index_name = 'idx_convert_group'", db=settings.mysql_core_database), - 1, - ) - check( - "risk_convert_detail.uk_idem 存在", - q1(agent_admin, "SELECT COUNT(*) FROM information_schema.statistics " - "WHERE table_schema = :db AND table_name = 'risk_convert_detail' " - "AND index_name = 'uk_idem'", db=settings.mysql_database), - 1, - ) - - # ── G. uk_idem 真库生效:重复 client_request_id 必须 IntegrityError ── - print("\n【G】uk_idem 真库生效:重复键报 IntegrityError") - with agent_admin.begin() as conn: - conn.execute( - text("INSERT INTO risk_convert_detail (convert_group_id, client_request_id, status, estimated) " - "VALUES ('CNV-T7M-DUP', 'T7M-DUP-KEY', 'pending', 0)"), - ) - dup_err = None - try: - with agent_admin.begin() as conn: - conn.execute( - text("INSERT INTO risk_convert_detail (convert_group_id, client_request_id, status, estimated) " - "VALUES ('CNV-T7M-DUP2', 'T7M-DUP-KEY', 'pending', 0)"), - ) - except Exception as exc: # noqa: BLE001 - dup_err = type(exc).__name__ - check("重复键报 IntegrityError", "IntegrityError" in (dup_err or ""), True) - with agent_admin.begin() as conn: - conn.execute(text("DELETE FROM risk_convert_detail WHERE client_request_id = 'T7M-DUP-KEY'")) - finally: - cleanup_core(core_admin) - cleanup_agent(agent_admin) - dispose_engines() - - print(f"\n{'=' * 60}") - print(f"真库验证(T-7 convert_service):{_passed} 项一致 / {_failed} 项不一致") - return 1 if _failed else 0 - - -if __name__ == "__main__": - sys.exit(main()) diff --git a/tests/test_convert_concurrency.py b/tests/test_convert_concurrency.py index 6f237c1..fd01eaa 100644 --- a/tests/test_convert_concurrency.py +++ b/tests/test_convert_concurrency.py @@ -1,4 +1,4 @@ -"""T-13 真 MySQL 并发测试(架构 §9 测试清单 · §10「50 并发压测口径(评审 Q6)」)。 +"""T-13 真 MySQL 并发测试 · **T+1 受理/确认分离模型**(架构 §9 · PRD 验收 18/31)。 **为什么必须是真 MySQL**:本文件验的是 InnoDB **行锁 + 条件 UPDATE rowcount** 这套 数据库级并发原语。sqlite 内存库(`conftest.sqlite_engine`)共用单连接,多线程跑不出 @@ -9,6 +9,21 @@ 未设置时重载用例逐条 skip(给出可操作的原因),常规 `pytest -q` 不受拖累。 T-13 实测:`CONVERT_STRESS=1 pytest tests/test_convert_concurrency.py -q -s`。 +**T+1 语义下的并发面(v1.0 用例的映射,convert_fund 退役后整体重写)** + +- **份额争抢从受理挪到确认**:受理段只做 R-3 在途占用校验(不扣份额); + 真正的 FIFO 扣减争抢在确认段 `apply_convert`(批次哨兵 `remain_qty >= :q` + + 1213 死锁整事务重试)。硬不变量「不超卖」钉在**确认后**的流水/余量上。 +- **受理超授竞态**:并发受理下 R-3 的占用读数会滞后(TOCTOU)→ 可能出现 + Σ占用 > 池的**超授**;T+1 的兜底是确认段逐笔 FIFO 严格扣减(份额不足 → + 部分成交 R-10 / `rejected`),绝不超卖 —— 这是设计语义,不是缺陷。 +- **紧池 80% 病根消失**:v1.0「需求=供给并发扣减只到 80% 成交」的根因是实时 + 扣减的批次碎片争抢;T+1 确认**串行**(FR-C23)后同场景 100% 成交 + (验收 31 设计目标,S1 用例钉住)。 +- **v1.0 占位竞态用例(3b/3c)退役**:T+1 受理单是 Core 6 态状态机(无 agent + 占位),同键幂等由 `uk_idem` + 客户级锁承接(N1 + 单测 + `test_accept_uk_idem_race_falls_back_to_idempotent`)。 + **隔离策略**(与 `test_convert_integration.py` 同款,三条同时成立,缺一即污染种子): 1. id 前缀:`CNV-CONC-` / `TRD-CONC-`(monkeypatch `convert_service._new_id`); 2. 数据自建:客户 `CUST-CONCTEST` / 产品 `PROD-CONCTEST?` 全自建,**绝不碰种子**—— @@ -41,11 +56,20 @@ from conftest import ensure_risk_demo_ready 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 ConvertRequestRepository # noqa: E402 from app.repository.core_ro import CoreReadOnlyRepository # noqa: E402 from app.repository.risk_repository import RiskRepository # noqa: E402 from app.service.convert import convert_service as cs # noqa: E402 -from app.service.convert.convert_service import PROCESSING, convert_fund # noqa: E402 +from app.service.convert.confirm_service import confirm_batch, confirm_one # noqa: E402 +from app.service.convert.convert_service import ( # noqa: E402 + accept_convert, + compensate_convert, +) from app.service.convert.errors import LotConflict # noqa: E402 +from app.service.convert.trading_calendar import ( # noqa: E402 + parse_cutoff, + resolve_accept_date, +) from app.service.risk.rules import RiskThresholds # noqa: E402 from app.utils.db import _resolve_credentials # noqa: E402 @@ -72,10 +96,13 @@ _LOT_DAYS = 100 # 全部批次同一费率档(0.0050)→ 断言只看份额 _GROUP_LIKE = "CNV-CONC-%" -#: 架构 §8.3 建议的退避间隔(`LOT_CONFLICT` → 调用方重试 ≤3 次) +#: 架构 §8.3 建议的退避间隔(瞬时失败 → 调用方重试 ≤3 次) BACKOFF_MS = (100, 200, 400) -#: 可重试异常:409 是明确设计为可重试的;死锁/锁等待超时同属瞬时状态 +#: 可重试异常(架构 §8.3 外部退避兜底的对象):批次争抢 `LotConflict` + +#: 确认事务 1213 死锁重试耗尽后上抛的 `OperationalError`。串行批处理(FR-C23) +# 下二者不应出现;绕过批处理锁并发 confirm_one 时由调用方退避兜底。 _RETRYABLE = (LotConflict, OperationalError) +_RETRYABLE_NAMES = {t.__name__ for t in _RETRYABLE} def _stress_enabled() -> bool: @@ -89,7 +116,7 @@ requires_stress = pytest.mark.skipif( def _noop_engine(out_trade: dict, in_trade: dict) -> dict: - """阶段 1.5 的空引擎:并发用例只验**份额争抢**,不掺引擎/出单噪声。""" + """空引擎 hook:并发用例只验**份额争抢**,不掺引擎/出单噪声。""" return {"triggered_rules": [], "alert_ids": [], "aml_hit": False} @@ -111,10 +138,23 @@ def _big_pool_engine(database: str, role: str, size: int = 60): # ── 种子与清理 ────────────────────────────────────────────────────── +def _accept_date(core) -> date: + """受理日 T(生产 `resolve_accept_date` 走真库日历实时判定,脚本不手算)。""" + return resolve_accept_date( + _NOW, + CoreReadOnlyRepository(engine=core).is_open, + parse_cutoff(settings.convert_cutoff_time), + ) + + def _seed(core, lots: list[str]) -> Decimal: - """自建隔离种子;`lots` = 各批次份额。返回份额池总量。""" + """自建隔离种子;`lots` = 各批次份额。返回份额池总量。 + + 两端净值都锚在**受理日 T**(确认段 `get_nav_on(pid, T)` 精确匹配)。 + """ today = date.today() base = datetime.combine(today, dtime(10, 0)) + accept_day = _accept_date(core) total = sum(Decimal(q) for q in lots) with core.begin() as conn: conn.execute( @@ -156,13 +196,14 @@ def _seed(core, lots: list[str]) -> Decimal: ), {"p": PROD_OUT, "mh": mh, "mm": mh_max, "r": rate}, ) - 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": str(IN_NAV), "d": today}, - ) + 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": str(nav), "d": accept_day}, + ) for i, qty in enumerate(lots, start=1): conn.execute( text( @@ -202,6 +243,7 @@ _PROD_LIKE = "PROD-CONCTEST%" #: `core_holding`/`core_share_lot`/`core_trade` 同时引用 `core_customer` 与 #: `core_product`,故必须先删它们,再删 `core_product`,最后删 `core_customer`。 _CORE_CLEANUP_BY_CUSTOMER = [ + "DELETE FROM core_convert_lot_detail WHERE convert_group_id LIKE 'CNV-CONC-%'", "DELETE FROM core_trade WHERE customer_id LIKE :cl", "DELETE FROM core_share_lot WHERE customer_id LIKE :cl", "DELETE FROM core_holding WHERE customer_id LIKE :cl", @@ -270,15 +312,15 @@ def conc_env(risk_demo_env, monkeypatch): eng.dispose() -# ── 调用与度量 ────────────────────────────────────────────────────── +# ── T+1 两段调用(受理 / 确认;生产账号 + 大池引擎)──────────────────── _NOW = datetime.combine(date.today(), dtime(10, 0)) -def _convert(env, qty: str, cid: str | None, hook=_noop_engine): - """一次转换(生产账号、生产口径;仅 id 工厂与引擎 hook 被替换)。""" - return convert_fund( +def _accept(env, qty: str, cid: str | None = None, customer: str = CUSTOMER) -> dict: + """T 日受理(不扣份额;生产口径,仅 id 工厂被替换)。""" + return accept_convert( { - "customer_id": CUSTOMER, + "customer_id": customer, "from_product_id": PROD_OUT, "to_product_id": PROD_IN, "qty": qty, @@ -287,32 +329,50 @@ def _convert(env, qty: str, cid: str | None, hook=_noop_engine): core_ro=CoreReadOnlyRepository(engine=env["core_ro"]), risk_repo=RiskRepository(engine=env["agent_rw"]), convert_repo=ConvertRepository(engine=env["agent_rw"]), + request_repo=ConvertRequestRepository(engine=env["core_rw"]), + now=_NOW, + ) + + +def _confirm(env, gid: str, hook=_noop_engine) -> dict: + """T+1 确认一笔(as_of=受理日,补跑语义;确认事务含 FIFO 扣减哨兵)。 + + `id_factory` 显式注入:`confirm_one` 的 trade_id 走 `format.new_id` + (与 `convert_service._new_id` 是两个符号),monkeypatch 罩不到 → 必须传参。 + """ + return confirm_one( + gid, + core_ro=CoreReadOnlyRepository(engine=env["core_ro"]), + risk_repo=RiskRepository(engine=env["agent_rw"]), + convert_repo=ConvertRepository(engine=env["agent_rw"]), + request_repo=ConvertRequestRepository(engine=env["core_rw"]), core_writer=ConvertCoreRepository(engine=env["core_rw"]), thresholds=RiskThresholds.from_settings(), now=_NOW, + as_of=_accept_date(env["core"]), engine_hook=hook, + id_factory=_new_id, ) class Outcome: """一个并发请求的终态(是否成功、重试了几次、最终异常)。""" - __slots__ = ("ok", "retries", "error", "processing", "group_id") + __slots__ = ("ok", "retries", "error", "group_id") - def __init__(self, ok, retries, error=None, processing=False, group_id=None): + def __init__(self, ok, retries, error=None, group_id=None): self.ok = ok self.retries = retries self.error = error - self.processing = processing self.group_id = group_id -def _run_with_backoff(env, qty: str, cid: str | None) -> Outcome: - """按架构 §8.3 的建议间隔退避重试(≤3 次);`client_request_id` 全程不变。""" +def _confirm_with_backoff(env, gid: str) -> Outcome: + """按架构 §8.3 的建议间隔退避重试确认(≤3 次;1213 类瞬时失败)。""" attempt = 0 while True: try: - resp = _convert(env, qty, cid) + resp = _confirm(env, gid) except _RETRYABLE as exc: if attempt >= len(BACKOFF_MS): return Outcome(False, attempt, type(exc).__name__) @@ -321,9 +381,12 @@ def _run_with_backoff(env, qty: str, cid: str | None) -> Outcome: continue except Exception as exc: # noqa: BLE001 return Outcome(False, attempt, type(exc).__name__) - if resp.get("status") == PROCESSING: - return Outcome(False, attempt, "PROCESSING", processing=True) - return Outcome(True, attempt, group_id=resp["convert_group_id"]) + return Outcome( + resp["status"] == "confirmed", + attempt, + None if resp["status"] == "confirmed" else resp["status"], + group_id=gid, + ) def _concurrently(fn, n: int, *, stagger_ms: float = 0.0) -> list: @@ -395,343 +458,95 @@ def _group_counts(core) -> dict[str, int]: return counts -def _assert_no_oversell(core, pool: Decimal, ok_count: int) -> None: - """三条硬不变量:不超卖 / 守恒 / 每组恰好两条流水(无重复组)。""" +def _request_statuses(core) -> dict[str, int]: + """受理单 6 态分布(T+1 权威状态机)。""" + with core.connect() as conn: + rows = conn.execute( + text( + "SELECT status, COUNT(*) AS n FROM core_convert_request" + " WHERE customer_id = :c GROUP BY status" + ), + {"c": CUSTOMER}, + ).mappings() + return {r["status"]: int(r["n"]) for r in rows} + + +def _assert_no_oversell(core, pool: Decimal) -> None: + """硬不变量:不超卖 / 守恒 / 每组恰好两条流水(无重复组)。""" deducted = _deducted_total(core) remain = _remain_total(core) assert deducted <= pool, f"超卖:扣减 {deducted} > 池 {pool}" assert remain == pool - deducted, f"不守恒:余 {remain} ≠ 池 {pool} − 扣 {deducted}" counts = _group_counts(core) assert set(counts.values()) <= {2}, f"存在流水条数 ≠2 的组:{counts}" - assert len(counts) == ok_count, f"成功数 {ok_count} 与组数 {len(counts)} 不一致" - assert len(_trades(core)) == ok_count * 2, "流水总数应为成功的转换数 ×2" -# ── 1. 50 并发争抢同一批份额:不许超卖(评审 Q6 断言 ①③)────────────── -@requires_stress -def test_50_concurrent_requests_never_oversell(conc_env): - """50 线程 × 各申请 2000、份额池 50000(5 批 × 10000)→ **不许超卖**。 +# ── N1. 同键并发受理:只允许落一张受理单(uk_idem + 客户级锁)────────── +def test_same_client_request_id_concurrent_produces_single_acceptance(conc_env): + """同 `client_request_id` 并发**受理** → core_convert_request 恰 1 行。 - 典型失效形态:若 `_deduct_lots` 的条件 UPDATE 丢掉 `remain_qty >= :q` - (并发哨兵),50 笔会**全部**成交 → 扣减 100000 > 池 50000 → 本用例红。 - - ⚠️ **为什么成交笔数不是恒等于 25(理论上限)**——实测在 **20~25** 间浮动: - `plan_lots` 会把一笔需求跨批次拆成多笔扣减(如 `400 + 1600`),而 - `_deduct_lots` 要求**本次全部扣减都成功**,否则整个阶段一事务回滚 → - 只要其中任一批次被别人清零,这一笔就整体 `LotConflict` 并进入重试; - 重试窗口重叠时彼此让路,落点因此不定。**不超卖是硬不变量,成交笔数是观测量。** + T+1 语义(取代 v1.0「同键并发只扣一次」):受理不扣份额,幂等收敛目标是 + **受理单唯一** —— 客户级锁串行化 + `uk_idem` 兜底让路,终态允许 + 「幂等命中(同 gid)」或「让路 202」,绝不允许第二张受理单、也绝无 5xx。 + (uk_idem 竞态兜底的单测锚点:`test_accept_uk_idem_race_falls_back_to_idempotent`) """ core = conc_env["core"] - pool = _seed(core, ["10000"] * 5) - - outcomes = _concurrently(lambda i: _run_with_backoff(conc_env, "2000", None), 50) - ok = [o for o in outcomes if o.ok] - - errs = {o.error for o in outcomes if not o.ok} - print( - f"\n[50 并发·不超卖] 成交 {len(ok)}/50(理论上限 25)· 失败类型 {errs} · " - f"余 {_remain_total(core)} · 扣减 {_deducted_total(core)}" - ) - # ① 硬不变量:不超卖 / 守恒 / 每组恰好两条流水(无重复组) - _assert_no_oversell(core, pool, len(ok)) - # ② 失败面:只能是"份额被抢走"这类可重试失败,不得是引擎/映射/锁异常 - assert errs <= {"LotConflict", "InsufficientShares", "OperationalError"}, ( - f"出现了非预期失败类型:{errs}" - ) - # ③ 反"假绿"地板:成交数塌穿 20 说明并发路径退化(哨兵过严/锁失效),必须红 - assert len(ok) >= 20, f"成交仅 {len(ok)} 笔,并发路径疑似退化:{errs}" - - -# ── 2. 退避重试的最终成功率(评审 Q6 断言 ②)────────────────────────── -@requires_stress -def test_backoff_retry_success_rate_on_exactly_sufficient_pool(conc_env): - """50 线程 × 各申请 1000、池恰好 50000 → 需求与供给**恰好相等**(最紧张)。 - - 这是「验证 §8.3 建议间隔(100/200/400ms)够不够」的场景:并发读-改-写必然 - 产生一批 `LotConflict`,退避重试后能否全部成交,直接量化建议间隔的充分性。 - - **实测结论(2026-09-10 · 本机 MySQL 8.0.46)**:3 次 × 100/200/400ms 只到 - **80%**(40/50,10 笔重试耗尽后仍 `LotConflict`)。根因不是"间隔太短",而是 - **近空批次的反复争抢**:`plan_lots` 会把一笔需求跨批次拆成 - `400 + 600` 这类碎片,碎片所在批次随时被别人清零 → 该笔重试仍可能抢不到。 - 要收敛到 100% 需要**更多次重试**或**冲突后换批次重规划**,而非单纯拉长间隔。 - 该结论已写入 `交接文档.md` §B;本用例把数字**钉住**(防日后倒退)。 - """ - core = conc_env["core"] - pool = _seed(core, ["10000"] * 5) - - outcomes = _concurrently(lambda i: _run_with_backoff(conc_env, "1000", None), 50) - ok = [o for o in outcomes if o.ok] - retried = [o for o in outcomes if o.retries] - rate = len(ok) / 50 - - _assert_no_oversell(core, pool, len(ok)) - assert _remain_total(core) == pool - _deducted_total(core) - # 失败者必须是**可重试类型**(份额被抢走),不得出现引擎/锁/映射异常 - errs = {o.error for o in outcomes if not o.ok} - assert errs <= {"LotConflict", "InsufficientShares"}, f"非预期失败类型:{errs}" - print( - f"\n[退避重试·紧池] 成功率 {len(ok)}/50 = {rate:.0%} · " - f"发生重试的请求 {len(retried)}/50 · 重试次数 {sorted(o.retries for o in outcomes)} · " - f"失败 {sorted(o.error for o in outcomes if not o.ok)}" - ) - # 护栏(**非指标**):<60% 说明退避链路整体失效(例如锁/哨兵被改坏),必须红 - assert rate >= 0.6, f"退避链路疑似失效:成功率仅 {rate:.0%}" - - -@requires_stress -def test_backoff_retry_all_succeed_when_pool_is_not_tight(conc_env): - """50 线程 × 各申请 500、池 50000(需求 25000 < 供给)→ **50/50 全成功**。 - - 与紧池用例互为对照:证明条件 UPDATE 哨兵 + 退避链路在**非耗尽**场景下 - 完全可靠——紧池的 80% 是"供给恰好用尽"的固有争抢,不是链路缺陷。 - """ - core = conc_env["core"] - pool = _seed(core, ["10000"] * 5) - - outcomes = _concurrently(lambda i: _run_with_backoff(conc_env, "500", None), 50) - ok = [o for o in outcomes if o.ok] - - assert len(ok) == 50, f"非紧池应全部成交,实际 {len(ok)}:{[o.error for o in outcomes if not o.ok]}" - _assert_no_oversell(core, pool, len(ok)) - print( - f"\n[退避重试·松池] 成功率 {len(ok)}/50 · " - f"发生重试的请求 {sum(1 for o in outcomes if o.retries)}/50 · 余 {_remain_total(core)}" - ) - - -# ── 3. 同键并发:只允许产生一次转换(架构 §9「同键并发 → 202」)──────── -def test_same_client_request_id_concurrent_produces_single_conversion(conc_env): - """同 `client_request_id` 并发提交 → 只允许一组流水、只扣一次份额。 - - 期望(架构 §9):未抢到执行权的那些返回 **202 `processing`**。 - 允许的终态:`processing`(202)或**幂等命中**(200,拿到首次结果); - 绝不允许出现**第二组流水**(双扣)。 - """ - core = conc_env["core"] - pool = _seed(core, ["50000"]) + _seed(core, ["50000"]) n = 8 cid = f"CONC-SAME-{uuid4().hex[:8]}" - def _one(_i: int) -> Outcome: + def _one(_i: int) -> dict: try: - resp = _convert(conc_env, "120", cid) + return _accept(conc_env, "120", cid) except Exception as exc: # noqa: BLE001 - return Outcome(False, 0, type(exc).__name__) - if resp.get("status") == PROCESSING: - return Outcome(False, 0, None, processing=True) - return Outcome(True, 0, group_id=resp["convert_group_id"]) + return {"_error": type(exc).__name__} outcomes = _concurrently(_one, n, stagger_ms=1.5) - kinds: dict[str, int] = {} - for o in outcomes: - key = "processing" if o.processing else ("ok" if o.ok else f"error:{o.error}") - kinds[key] = kinds.get(key, 0) + 1 - print(f"\n[同键并发 {n} 路] 终态 {kinds} · 组数 {len(_group_counts(core))}") - - # 唯一硬约束:**只扣一次**(不得双组流水) - assert _deducted_total(core) == Decimal("120.0000"), ( - f"同键并发了 {_deducted_total(core)} 份,应为 120 份(不得双扣)" - ) - assert len(_group_counts(core)) == 1 - assert len(_trades(core)) == 2 - # 错误面约束:不得出现 5xx(幂等语义应表现为 202 或幂等命中) - assert not [o for o in outcomes if o.error], ( - f"同键并发出现了 5xx:{[o.error for o in outcomes if o.error]}" - ) - assert pool == Decimal("50000") - - -# ── 3b. 同键并发的**竞态窗口**(确定性交错 · 缺陷已修复后的回归闸门)──── -@requires_stress -def test_same_client_request_id_interleaving_must_not_leak_5xx(conc_env, monkeypatch): - """把 T1 **确定性**停在「占位已落、幂等锁已释放、阶段一未提交」那一点,再放 T2 进来。 - - 这是并发缺陷的标准证法(确定性交错,不靠碰运气):`convert_fund` 的 - `with try_lock("convert:idem:{cid}")` 只包住**幂等判定**(出块即释放), - 紧随其后的 ①②③④(校验与折算,多次 DB 往返)+ ⑤(占位)与 T1 的阶段一之间 - 形成一个数毫秒的窗口。窗口内进入的第二笔请求读到 `existing.status='pending'` - 且 `has_convert_trades=False` → 判定为「阶段一未成 → 复用 group_id 重跑」, - 于是**两笔带着同一个 `group_id` 同时跑阶段一**。 - - 已实测的失效形态(**这个子窗口不是双扣**,是 500): - `in_lot_id = f"LOT-{group_id}-IN"` 派生自 group_id,后跑的那笔 INSERT 撞 - `core_share_lot` 主键 → 整个阶段一事务回滚 → **只扣一次**(120 份)。代价是 - **未映射的 `IntegrityError` 直穿到 API = 500**,而架构 §9 的契约是 - 「同键并发 → 202」、§7.3 的契约是幂等命中。 - - ⚠️ **但"不是双扣"只在本子窗口成立** —— 见 3c:若两笔各自生成了**不同**的 - group_id,`in_lot_id` 不会相撞,那边是**真·双扣**(实测扣 240 份)。 - 也就是说 3b 的 120 份是**顺带**被派生主键挡下的,不是被并发设计挡下的。 - - 本用例断言正确契约。**修复**(T-13 落地):幂等判定区分「确定没跑成」与 - 「在飞/未知」——`failed`/`expired` 才复用 group_id 重跑,`pending` 一律回 202。 - 本用例即修复后的回归闸门:把 `pending` 改回"也重跑"即变红。 - """ - core = conc_env["core"] - _seed(core, ["50000"]) - cid = f"CONC-RACE-{uuid4().hex[:8]}" - - original = ConvertCoreRepository.apply_convert - in_stage1 = threading.Event() - release = threading.Event() - gate_used = {"v": False} - - def gated(self, req): # noqa: ANN001 - if not gate_used["v"]: # 只拦第一笔(T1) - gate_used["v"] = True - in_stage1.set() - release.wait(10) - return original(self, req) - - monkeypatch.setattr(ConvertCoreRepository, "apply_convert", gated) - - got: dict = {} - - def _t1() -> None: - try: - got["t1"] = _convert(conc_env, "120", cid) - except Exception as exc: # noqa: BLE001 - got["t1_err"] = type(exc).__name__ - - th = threading.Thread(target=_t1, name="T1") - th.start() - assert in_stage1.wait(10), "T1 未进入阶段一,交错构造失败" - # 此刻 T1:占位 pending 已落库、`convert:idem:` 锁已释放、阶段一未提交 - try: - got["t2"] = _convert(conc_env, "120", cid) - except Exception as exc: # noqa: BLE001 - got["t2_err"] = type(exc).__name__ - release.set() - th.join(20) - - deducted = _deducted_total(core) - counts = _group_counts(core) + errs = [o["_error"] for o in outcomes if "_error" in o] + gids = {o["convert_group_id"] for o in outcomes if "_error" not in o} print( - f"\n[同键竞态·确定性交错] T1={got.get('t1_err') or 'ok'} " - f"T2={got.get('t2_err') or 'ok'} · 扣减 {deducted} · 组→流水数 {counts}" + f"\n[同键并发受理 {n} 路] 受理单 {len(_request_statuses(core))} 张 · " + f"gid 去重 {len(gids)} · 错误 {errs}" ) - # ① 资金安全:只扣一次(当前由 in_lot_id 主键**顺带**保证,非设计保证) - assert deducted == Decimal("120.0000"), f"同键被扣了 {deducted} 份(应为 120)" - assert len(counts) == 1 and len(_trades(core)) == 2, f"应一组两条流水:{counts}" - # ② 错误面:不得有未映射异常直穿(契约是 202 / 幂等命中,不是 5xx) - leaked = {k: v for k, v in got.items() if k.endswith("_err")} - assert not leaked, f"同键并发出现未映射异常直穿(应为 202 或幂等命中):{leaked}" + + # 唯一硬约束:恰一张受理单(uk_idem),且无 5xx + statuses = _request_statuses(core) + assert sum(statuses.values()) == 1, f"同键并发落了 {sum(statuses.values())} 张受理单" + assert len(gids) <= 1, f"同键应幂等命中同一 gid:{gids}" + assert not errs, f"同键并发出现了异常:{errs}" -# ── 3c. 同键并发的**第二个子窗口**:两笔都判定"无占位"(确定性交错)──── -@requires_stress -def test_same_client_request_id_placeholder_race_must_not_leak_5xx(conc_env, monkeypatch): - """把 T1 停在「幂等判定已过、占位尚未插入」那一点,让 T2 先占位成功。 - - 这是与 3b **不同的**子窗口: - - 3b:T2 在 T1 占位**之后**进入 → 看到 `pending` → 旧代码当"未成"重跑 → 双跑阶段一 - (同 group_id,被 `in_lot_id` 主键顺带挡下,表现为 500); - - 3c:T2 在 T1 占位**之前**进入 → 也判定"无占位" → **各自生成不同的 group_id** - → 后插入的那笔撞 `uk_idem`。 - - **旧代码表现**:撞键直穿未映射的 `IntegrityError`(503)。 - - ⚠️ **这条路径比 503 更危险**:若只修仓储(让 `insert_placeholder` 撞键返回 False) - 而**不接住这个让路信号**,两笔会各带**不同 group_id** 一路跑到阶段一 —— - `in_lot_id` 不再相撞 → **真·双扣**。突变验证实测:扣 **240** 份、**两组流水** - (T1=ok T2=ok,`{'CNV-…': 2, 'CNV-…': 2}`)。这是本次 T-13 压测最严重的一处发现, - 也是"必须由调用方 `return 202`"而非"仓储静默吞掉"的原因。 - - 修复后:`insert_placeholder` 返回"本笔是否持有占位",撞键即让路 → 调用方回 202。 - 断言:① 无未映射异常直穿;② 只扣一次;③ 只有一组流水。 - """ - core = conc_env["core"] - _seed(core, ["50000"]) - cid = f"CONC-UKIDEM-{uuid4().hex[:8]}" - - original = ConvertRepository.insert_placeholder - in_zero = threading.Event() - release = threading.Event() - gate_used = {"v": False} - - def gated(self, group_id, client_request_id): # noqa: ANN001 - if not gate_used["v"]: # 只拦第一笔(T1),且**拦在插入之前** - gate_used["v"] = True - in_zero.set() - release.wait(10) - return original(self, group_id, client_request_id) - - monkeypatch.setattr(ConvertRepository, "insert_placeholder", gated) - - got: dict = {} - - def _t1() -> None: - try: - got["t1"] = _convert(conc_env, "120", cid) - except Exception as exc: # noqa: BLE001 - got["t1_err"] = type(exc).__name__ - - th = threading.Thread(target=_t1, name="T1") - th.start() - assert in_zero.wait(10), "T1 未到达阶段零,交错构造失败" - # 此刻 T1:幂等判定已过(无占位)、`convert:idem:` 锁已释放、占位未插 - try: - got["t2"] = _convert(conc_env, "120", cid) - except Exception as exc: # noqa: BLE001 - got["t2_err"] = type(exc).__name__ - release.set() - th.join(20) - - deducted = _deducted_total(core) - counts = _group_counts(core) - print( - f"\n[占位竞态·确定性交错] T1={got.get('t1_err') or 'ok'} " - f"T2={got.get('t2_err') or 'ok'} · 扣减 {deducted} · 组→流水数 {counts}" - ) - leaked = {k: v for k, v in got.items() if k.endswith("_err")} - assert not leaked, f"uk_idem 竞态直穿未映射异常(应为 202):{leaked}" - assert deducted == Decimal("120.0000"), f"同键被扣了 {deducted} 份(应为 120)" - assert len(counts) == 1 and len(_trades(core)) == 2, f"应一组两条流水:{counts}" - - -# ── 4. 补偿加锁:并发补偿只落一次(T-12 的并发面 · T+1 前置)──────────── +# ── N2. 并发补偿只落一次(T-12 的并发面 · T+1 造数)──────────────────── def test_concurrent_compensation_writes_detail_once(conc_env, monkeypatch): """待补偿现场并发补偿(`convert:rerun:` 锁)→ 详情只补一次、只有一组流水。 - T+1 造数:走 v1.0 `_convert` 真实跑出 Core 侧两条流水 + 计费明细(受理单是 - T+1 才有的表,SQL 补插 confirmed 行凑齐补偿前置)→ 镜像 UPDATE 回 `failed` - 制造「确认段第 ⑧ 步失败」现场。并发文件整体切 T+1 底座属 T-13 - (convert_fund 退役时一并重写)。 + T+1 造数:真受理 → 真确认(Core 三件套 + 镜像 completed + 引擎 hook 静默) + → 镜像 UPDATE 回 `failed` 制造「确认段第 ⑧ 步失败」现场(受理单已 confirmed, + 补偿前置满足)。 """ - from app.service.convert.convert_service import compensate_convert - core = conc_env["core"] agent = conc_env["agent"] _seed(core, ["50000"]) - gid = _convert(conc_env, "120", f"CONC-CMP-{uuid4().hex[:8]}")["convert_group_id"] - with core.begin() as conn: + resp = _accept(conc_env, "120", f"CONC-CMP-{uuid4().hex[:8]}") + gid = resp["convert_group_id"] + out = _confirm(conc_env, gid) + assert out["status"] == "confirmed" + # 制造待补偿现场:镜像回退 failed(Core 权威不动) + with agent.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, actual_qty, status," - " requested_at, confirmed_at)" - " VALUES (:g, NULL, :c, :fp, :tp, 120.0, 120.0, 'confirmed', :t, :t)" - ), - {"g": gid, "c": CUSTOMER, "fp": PROD_OUT, "tp": PROD_IN, "t": _NOW}, - ) - with agent.begin() as conn: # 镜像在 agent 库 - conn.execute( - text( - "UPDATE risk_convert_detail SET status = 'failed'" - " WHERE convert_group_id = :g" - ), + text("UPDATE risk_convert_detail SET status = 'failed' WHERE convert_group_id = :g"), {"g": gid}, ) + def _comp(_i: int) -> str: svc = dict( core_ro=CoreReadOnlyRepository(engine=conc_env["core_ro"]), risk_repo=RiskRepository(engine=conc_env["agent_rw"]), convert_repo=ConvertRepository(engine=conc_env["agent_rw"]), ) - out = compensate_convert(gid, thresholds=RiskThresholds.from_settings(), now=_NOW, **svc) - return out["state"] + result = compensate_convert(gid, thresholds=RiskThresholds.from_settings(), now=_NOW, **svc) + return result["state"] states = _concurrently(_comp, 6, stagger_ms=1.0) print(f"\n[并发补偿 6 路] 状态分布 {states}") @@ -755,20 +570,135 @@ def test_concurrent_compensation_writes_detail_once(conc_env, monkeypatch): assert status == "confirmed" -# ── 5. 多客户并发:互不阻塞,但会出 1213 死锁(须靠退避重试收敛)───────── +# ── S1. 紧池(需求=供给):受理 100% + 串行确认 100%(验收 31 载体)──── +@requires_stress +def test_tight_pool_accept_all_then_confirm_100pct(conc_env): + """50 并发受理各 1000、池恰好 50000(需求=供给)→ **受理 100% + 确认 100%**。 + + 这是 PRD 验收 31 的 T+1 载体:v1.0 同场景实时扣减只到 **80%**(40/50, + 批次碎片争抢,见交接文档 §B.6.6)——根因是「确认前就扣份额」。T+1 受理 + 不扣份额(R-3 占用校验:50×1000 恰好 ≤ 池,**全部 accepted 是确定性结果**), + 确认**串行**(FR-C23)无争抢 → 逐笔 FIFO 全额成交 → 100%。 + + 硬门禁:① 受理恰 50 张 accepted;② 确认恰 50 张 confirmed(100%); + ③ Σ扣减 = 池 = 50000(需求=供给全额消化);④ 守恒 + 无超卖。 + """ + core = conc_env["core"] + pool = _seed(core, ["10000"] * 5) + n = 50 + + outcomes = _concurrently(lambda i: _accept(conc_env, "1000"), n, stagger_ms=0.5) + accepted = [o for o in outcomes if o.get("accepted")] + print( + f"\n[紧池·受理] accepted {len(accepted)}/{n} · 其余 " + f"{[o.get('block_reason') or o.get('_error') for o in outcomes if not o.get('accepted')]}" + ) + assert len(accepted) == n, f"受理应 100%(设计目标),实际 {len(accepted)}" + statuses = _request_statuses(core) + assert statuses.get("accepted") == n, f"受理单应全为 accepted:{statuses}" + + # T+1 确认:批处理串行(as_of = 受理日,补跑语义) + accept_day = _accept_date(core) + batch = confirm_batch( + as_of=accept_day, + core_ro=CoreReadOnlyRepository(engine=conc_env["core_ro"]), + risk_repo=RiskRepository(engine=conc_env["agent_rw"]), + convert_repo=ConvertRepository(engine=conc_env["agent_rw"]), + request_repo=ConvertRequestRepository(engine=conc_env["core_rw"]), + core_writer=ConvertCoreRepository(engine=conc_env["core_rw"]), + thresholds=RiskThresholds.from_settings(), + now=_NOW, + engine_hook=_noop_engine, + id_factory=_new_id, + ) + print( + f"\n[紧池·串行确认] scanned={batch['scanned']} confirmed={batch['confirmed']} " + f"rejected={batch.get('rejected', 0)} nav_pending={batch.get('nav_pending', 0)}" + ) + assert batch["scanned"] == n, f"批处理应捞到全部 {n} 张" + assert batch["confirmed"] == n, f"串行确认应 100%(验收 31 目标),实际 {batch}" + statuses = _request_statuses(core) + assert statuses.get("confirmed") == n, f"受理单应全为 confirmed:{statuses}" + # 需求=供给 → 池被全额消化且守恒(无超卖硬门禁) + assert _deducted_total(core) == pool, f"扣减 {_deducted_total(core)} 应恰等于池 {pool}" + _assert_no_oversell(core, pool) + assert len({g for g in _group_counts(core)}) == n, "50 张单应有 50 个独立组" + + +# ── S2. 并发确认争抢同一池:不许超卖(评审 Q6 断言的 T+1 承接)────────── +@requires_stress +def test_concurrent_confirms_never_oversell(conc_env): + """并发受理允许**超授**(R-3 占用读数滞后的 TOCTOU),并发确认兜底 → 不许超卖。 + + 50 线程 × 各受理 2000(池 50000):受理段占用校验读数滞后时可能出现 + Σ占用 > 池的超授 —— T+1 的资金安全不依赖受理段,而在**确认段扣减哨兵** + (`remain_qty >= :q` 条件 UPDATE + 死锁重试):份额不足时逐笔部分成交 + (R-10)或 `rejected`,Σ实际扣减绝不越过池。 + + 典型失效形态:若确认事务丢掉扣减哨兵,50 笔会**全部**全额成交 → + 扣减 100000 > 池 50000 → 本用例红。 + """ + core = conc_env["core"] + pool = _seed(core, ["10000"] * 5) + n = 50 + + # 受理段:并发窗口内 R-3 占用读数滞后 → 允许超授(Σ占用 > 池); + # 校验命中时直接 `InsufficientShares`(受理段 4xx,异常而非 blocked)。 + def _try_accept(_i: int) -> dict: + try: + return _accept(conc_env, "2000") + except Exception as exc: # noqa: BLE001 - 受理失败本身是合法观测面 + return {"_error": type(exc).__name__} + + accepted = _concurrently(_try_accept, n, stagger_ms=0.5) + gids = [o["convert_group_id"] for o in accepted if o.get("accepted")] + accept_errors = {o["_error"] for o in accepted if "_error" in o} + print( + f"\n[并发确认·不超卖] 受理 {len(gids)}/{n}(允许超授)· 受理失败类型 {accept_errors} · " + f"受理单分布 {_request_statuses(core)}" + ) + assert gids, "至少应有部分受理成功" + + # 确认段:并发 confirm_one 争抢扣减 → 兜底不允许超卖 + confirmed = _concurrently(lambda i: _confirm_with_backoff(conc_env, gids[i]), len(gids)) + ok = [o for o in confirmed if o.ok] + # 合法失败面(都是**零扣减**的让路终态,不碰资金安全): + # - `rejected`:份额被抢空 → 复核拒绝(占用释放); + # - `LotConflict`:退避耗尽仍撞批次 → 让路。并发 confirm_one 本身是**超纲场景** + # (生产由批处理锁保证串行,FR-C23),让路单受理单保持 accepted, + # 占用保持、下轮批处理再确认 —— 与 v1.0 用例 1 允许 LotConflict 同口径。 + legit_fail = {"rejected", "LotConflict"} + unexpected = {o.error for o in confirmed if not o.ok and o.error not in legit_fail} + print( + f"\n[并发确认·不超卖] confirmed {len(ok)}/{len(gids)} · 让路/拒绝 " + f"{sorted({o.error for o in confirmed if not o.ok})} · " + f"未映射异常 {unexpected} · 扣减 {_deducted_total(core)} / 池 {pool}" + ) + # ① 硬不变量:不超卖 / 守恒 / 每组恰好两条流水 + _assert_no_oversell(core, pool) + # ② 失败面:终态只能是 confirmed/rejected/LotConflict(部分成交含在 confirmed 内), + # 不得有未映射异常直穿 + assert not unexpected, f"确认出现未映射异常:{unexpected}" + # ③ 让路不吞单:没确认成的单,受理单必须保持 accepted(占用保持,下轮可确认) + statuses = _request_statuses(core) + assert statuses.get("accepted", 0) == len(gids) - len(ok), ( + f"让路单受理单状态异常:{statuses}(成功 {len(ok)}/{len(gids)})" + ) + # ④ 反"假绿"地板:confirmed 塌穿 20 说明扣减路径退化(哨兵过严/锁失效) + assert len(ok) >= 20, f"确认仅 {len(ok)} 笔,扣减路径疑似退化" + + +# ── S3. 多客户并发确认:互不阻塞,但会出 1213 死锁(双层重试收敛)─────── @requires_stress def test_concurrent_different_customers_need_deadlock_retry(conc_env): - """8 个客户并发转换同一转出产品 → 全部成功,但**必须靠重试死锁**。 + """8 个客户**并发确认**同一转出产品 → 全部成功,但**必须靠重试死锁**。 - **实测发现(2026-09-10)**:跨客户并发在 `core_holding`/`core_share_lot` 的 - 唯一索引上会触发 InnoDB **1213 死锁**(RR 隔离级下 INSERT/UPDATE 的 - 插入意向锁与间隙锁互斥)。**T-13 后补(同日用户拍板)**:`apply_convert` - 已在数据访问层对 1213 自动重试整事务(settings 可配);本用例保留 §8.3 - 的**外部退避兜底**,双层重试叠加下 8 客户并发仍须全部成功、份额守恒、 - group_id 两两不同。 - - 注意与「50 并发不超卖」的区别:那批是**同一** (客户, 产品),走的是行锁阻塞; - 这批是**不同**客户,走的是间隙锁 → 死锁而非阻塞。 + **实测发现(2026-09-10,v1.0 阶段一)**:跨客户并发在 `core_holding`/ + `core_share_lot` 的唯一索引上会触发 InnoDB **1213 死锁**(RR 隔离级别下 + INSERT/UPDATE 的插入意向锁与间隙锁互斥)。T+1 下同一争抢面在**确认事务** + (`apply_convert`,1213 自动重试整事务,settings 可配);本用例保留 §8.3 + 的**外部退避兜底**,双层重试叠加下 8 客户并发确认仍须全部成功、份额守恒、 + group_id 两两不同。受理段先串行完成(不扣份额、无争抢面)。 """ core = conc_env["core"] _seed(core, ["5000"] * 2) @@ -800,42 +730,21 @@ def test_concurrent_different_customers_need_deadlock_retry(conc_env): ) customers = [CUSTOMER, *others] + # 受理先串行完成(不扣份额;其他客户的受理走同款服务路径) + gids = {} + for c in customers: + resp = _accept(conc_env, "1000", customer=c) + assert resp.get("accepted"), f"{c} 受理失败:{resp}" + gids[c] = resp["convert_group_id"] def _one(i: int) -> Outcome: - """走 `_run_with_backoff` 的等价逻辑,但客户可变(死锁重试在此体现)。""" - attempt = 0 - while True: - try: - resp = convert_fund( - { - "customer_id": customers[i], - "from_product_id": PROD_OUT, - "to_product_id": PROD_IN, - "qty": "1000", - }, - core_ro=CoreReadOnlyRepository(engine=conc_env["core_ro"]), - risk_repo=RiskRepository(engine=conc_env["agent_rw"]), - convert_repo=ConvertRepository(engine=conc_env["agent_rw"]), - core_writer=ConvertCoreRepository(engine=conc_env["core_rw"]), - thresholds=RiskThresholds.from_settings(), - now=_NOW, - engine_hook=_noop_engine, - ) - except _RETRYABLE as exc: - if attempt >= len(BACKOFF_MS): - return Outcome(False, attempt, type(exc).__name__) - time.sleep(BACKOFF_MS[attempt] / 1000.0) - attempt += 1 - continue - except Exception as exc: # noqa: BLE001 - return Outcome(False, attempt, type(exc).__name__) - return Outcome(True, attempt, group_id=resp["convert_group_id"]) + return _confirm_with_backoff(conc_env, gids[customers[i]]) outcomes = _concurrently(_one, len(customers)) ok = [o for o in outcomes if o.ok] deadlock_retries = sum(o.retries for o in outcomes) print( - f"\n[8 客户并发] 成功 {len(ok)}/8 · 死锁重试累计 {deadlock_retries} 次 · " + f"\n[8 客户并发确认] 成功 {len(ok)}/8 · 死锁重试累计 {deadlock_retries} 次 · " f"失败 {[o.error for o in outcomes if not o.ok]}" ) assert len(ok) == 8, f"退避重试后应全部成功:{[o.error for o in outcomes if not o.ok]}" @@ -858,9 +767,9 @@ def test_concurrent_different_customers_need_deadlock_retry(conc_env): assert len({o.group_id for o in ok}) == 8, "8 个客户的 group_id 必须两两不同" -# ── 6. 性能实测(PRD §9 第 18 条)──────────────────────────────────── +# ── S4. 性能实测(PRD §9 第 18 条 · T+1 两段)────────────────────────── class _TimingWriter: - """`ConvertCoreRepository` 代理:只测**阶段一单库事务**耗时,不改任何行为。""" + """`ConvertCoreRepository` 代理:只测**确认事务**耗时,不改任何行为。""" def __init__(self, inner): self._inner = inner @@ -881,48 +790,58 @@ def _pct(values: list[float], p: float) -> float: @requires_stress -def test_performance_probe_stage1_and_end_to_end(conc_env): - """实测:阶段一单库事务 + 端到端一次转换(PRD §9 第 18 条补录来源)。 +def test_performance_probe_accept_confirm_end_to_end(conc_env): + """实测:**受理事务 + 确认事务 + 端到端两段**耗时(PRD §9 第 18 条补录来源)。 测量条件:本机 MySQL 8.0.46(127.0.0.1)、模拟库、**单线程顺序**、 - 份额池充足的稳态(每笔 1 份,避免 `plan_lots` 跨批波动)。 - 端到端含阶段 1.5 空引擎(真实引擎的额外开销在 `test_convert_integration` 另行体现)。 + 份额池充足的稳态(每笔 1 份,避免 FIFO 跨批波动);端到端含空引擎 hook + (真实引擎的额外开销在 `test_convert_integration` 另行体现)。 + 端到端 = 受理(校验+落单+镜像+审计)+ 确认(折算+扣减+流水+镜像+审计)两段之和, + PRD 验收 18 的硬指标是 **< 2s**。 """ core = conc_env["core"] n = 60 - _seed(core, [str(n)]) + _seed(core, [str(n)] * 1) writer = _TimingWriter(ConvertCoreRepository(engine=conc_env["core_rw"])) - totals: list[float] = [] + accept_ms: list[float] = [] + confirm_ms: list[float] = [] for i in range(n): - req = { - "customer_id": CUSTOMER, - "from_product_id": PROD_OUT, - "to_product_id": PROD_IN, - "qty": "1", - "client_request_id": f"CONC-PERF-{i}", - } t0 = time.perf_counter() - convert_fund( - req, + resp = _accept(conc_env, "1", f"CONC-PERF-{i}") + accept_ms.append((time.perf_counter() - t0) * 1000.0) + gid = resp["convert_group_id"] + + t1 = time.perf_counter() + out = confirm_one( + gid, core_ro=CoreReadOnlyRepository(engine=conc_env["core_ro"]), risk_repo=RiskRepository(engine=conc_env["agent_rw"]), convert_repo=ConvertRepository(engine=conc_env["agent_rw"]), + request_repo=ConvertRequestRepository(engine=conc_env["core_rw"]), core_writer=writer, thresholds=RiskThresholds.from_settings(), now=_NOW, + as_of=_accept_date(core), engine_hook=_noop_engine, + id_factory=_new_id, ) - totals.append((time.perf_counter() - t0) * 1000.0) + confirm_ms.append((time.perf_counter() - t1) * 1000.0) + assert out["status"] == "confirmed", out # 首个样本含连接池冷启动,剔除(否则把"建连"算进"业务耗时") - stage1 = writer.samples[1:] - e2e = totals[1:] + stage_confirm = writer.samples[1:] + acc = accept_ms[1:] + cfm = confirm_ms[1:] + e2e = [a + c for a, c in zip(acc, cfm)] print( - f"\n[性能实测 n={len(e2e)}] 阶段一 ms: P50={_pct(stage1, 50):.1f} " - f"P95={_pct(stage1, 95):.1f} max={max(stage1):.1f} mean={mean(stage1):.1f}" - f"\n 端到端 ms: P50={_pct(e2e, 50):.1f} P95={_pct(e2e, 95):.1f} " + f"\n[性能实测 n={len(e2e)}] 受理 ms: P50={_pct(acc, 50):.1f} P95={_pct(acc, 95):.1f} " + f"max={max(acc):.1f} mean={mean(acc):.1f}" + f"\n 确认事务 ms: P50={_pct(stage_confirm, 50):.1f} " + f"P95={_pct(stage_confirm, 95):.1f} max={max(stage_confirm):.1f} mean={mean(stage_confirm):.1f}" + f"\n 确认全程 ms: P50={_pct(cfm, 50):.1f} P95={_pct(cfm, 95):.1f} max={max(cfm):.1f}" + f"\n 端到端(受理+确认) ms: P50={_pct(e2e, 50):.1f} P95={_pct(e2e, 95):.1f} " f"max={max(e2e):.1f} mean={mean(e2e):.1f}" ) - # PRD §9 第 18 条的验收硬指标是端到端 < 2s(阶段一 <100ms 是预估、非硬指标) + # PRD §9 第 18 条的验收硬指标是端到端 < 2s assert max(e2e) < 2000, f"端到端最大 {max(e2e):.1f}ms 超 PRD §9 第 18 条的 2s" diff --git a/tests/test_convert_service.py b/tests/test_convert_service.py index 24afb81..f51aa53 100644 --- a/tests/test_convert_service.py +++ b/tests/test_convert_service.py @@ -1,19 +1,19 @@ -"""T-7 `convert_service` 八步编排单测(开发计划 §6.2 DoD)。 +"""`convert_service` 补偿单点单测(T-12/T-13 · 开发计划 §6.2 DoD)。 sqlite 内存库(`conftest.sqlite_engine`,含 C×R 矩阵种子),数据自建。 -折算期望值一律由生产 `calc.py` 实算(不自造公式副本,自检第 13 问)。 -覆盖: -1. 三阶段贯通(占位 pending→completed / 2 流水 / 响应字段与 PRD §5.3 对齐) -2. **blocked 不占位**(PRD §7.0:前四步不落库) -3. 幂等命中返回首次结果、**不产生第二组流水** -4. 未抢到执行权 → `status=processing`(T-9 映射 202) -5. 全额转出豁免 `min_redeem_qty`(验收 16)/ 强制全转留痕 -6. `nav_stale` 额外落副审计(1~2 条) -7. 阶段 1.5 引擎异常**不阻断**已成立交易 -8. 阶段二失败 → `convert_detail_write_failed` 审计 + 占位 `failed` -9. 各 4xx/503 分支:SameProduct / 不可赎回 / 跨主体 / 份额不足 / 低于最低份额 / - 无净值 / 批次数超限 +**T-13 退役说明**:本文件原承载 v1.0 两阶段实时入口 `convert_fund` 的八步编排 +用例(三阶段贯通 / 占位幂等 / nav_stale 副审计 / 阶段一失败重试等)。`convert_fund` +随 T+1 模型退役删除后,其验证点的 **T+1 替代断言**已由以下用例承接: + +- 受理段(blocked 不占位 / 同键幂等 / 最低份额 / 强制全转 / 产品校验 / 在途占用 / + uk_idem 竞态)→ `tests/test_convert_accept.py`; +- 确认段(happy path 与 PRD §5.3 对表 / 缺净值 nav_pending / 引擎恰好一次 / + 引擎异常不阻断但留痕 / 真引擎出单去重 / 强制全转继承)→ `tests/test_convert_confirm.py`。 + +本文件保留 **补偿单点 `compensate_convert`** 的全部用例,并补 2 条 T+1 缺口: +确认第 ⑧ 步镜像失败不回滚 Core(v1.0「阶段二失败」语义的承接)、 +确认段批次数超限 rejected(v1.0 `TooManyLots` 503 语义的 T+1 承接)。 """ from __future__ import annotations @@ -24,42 +24,20 @@ from decimal import Decimal import pytest from sqlalchemy import text +from unittest.mock import patch +from app.config.settings import settings from app.gateway.convert_core_repository import ConvertCoreRepository from app.repository.convert_repository import ConvertRepository +from app.repository.convert_request_repository import ConvertRequestRepository from app.repository.core_ro import CoreReadOnlyRepository from app.repository.risk_repository import RiskRepository -from app.service.convert.calc import ( - convert_amount, - diff_fee, - hold_days, - in_qty, - lot_amount, - lot_fee, - plan_lots, -) -from app.service.convert.convert_service import ( - PROCESSING, - compensate_convert, - convert_fund, -) -from app.service.convert.errors import ( - BelowMinQty, - CrossEntityNotSupported, - InsufficientShares, - LotConflict, - NavNotReady, - ProductNotRedeemable, - SameProduct, - TooManyLots, -) -from app.service.convert.fee import pick_fee_rate -from app.service.convert.types import FeeRule, Lot +from app.service.convert.confirm_service import confirm_one +from app.service.convert.convert_service import compensate_convert from app.service.risk.locks import try_lock from app.service.risk.rules import RiskThresholds CUST = "CUST-T7" # C3 客户 → R4 产品 allowed_with_disclosure(放行) -CUST_LOW = "CUST-T7L" # C1 客户 → R4 forbidden(用于 blocked 分支) PROD_OUT = "PROD-T7A" PROD_IN = "PROD-T7B" COMPANY = "华夏模拟基金" @@ -92,24 +70,12 @@ def _seed(engine, *, nav_date: date = TODAY, min_redeem: str = "0", min_hold: st "VALUES (:c, 'T7客户', 40, 1)", c=CUST, ) - _exec( - engine, - "INSERT INTO core_customer (customer_id, display_name, age, is_active) " - "VALUES (:c, 'T7低风险客户', 40, 1)", - c=CUST_LOW, - ) _exec( engine, "INSERT INTO core_customer_risk (customer_id, risk_code, evaluated_at, expires_at) " "VALUES (:c, 'C3', :t, :exp)", c=CUST, t=NOW - timedelta(days=30), exp=NOW + timedelta(days=300), ) - _exec( - engine, - "INSERT INTO core_customer_risk (customer_id, risk_code, evaluated_at, expires_at) " - "VALUES (:c, 'C1', :t, :exp)", - c=CUST_LOW, t=NOW - timedelta(days=30), exp=NOW + timedelta(days=300), - ) _exec( engine, "INSERT INTO core_product (product_id, product_name, min_risk_code, product_type, " @@ -146,17 +112,6 @@ def _seed(engine, *, nav_date: date = TODAY, min_redeem: str = "0", min_hold: st "pnl_pct, as_of) VALUES (:c, :p, 150, 150, 154.5, 0, :d)", c=CUST, p=PROD_OUT, d=TODAY, ) - # CUST_LOW(C1,R4 应被适当性拦截)也需有份额,否则 ② 份额校验会先于 ④ 触发 - _seed_lot(engine, "LOT-T7L-1", "100", "1.0300", datetime(2026, 8, 1, 10, 0, 0), - customer=CUST_LOW) - _seed_lot(engine, "LOT-T7L-2", "50", "1.0000", datetime(2026, 9, 1, 10, 0, 0), - customer=CUST_LOW) - _exec( - engine, - "INSERT INTO core_holding (customer_id, product_id, qty, cost_amount, market_value, " - "pnl_pct, as_of) VALUES (:c, :p, 150, 150, 154.5, 0, :d)", - c=CUST_LOW, p=PROD_OUT, d=TODAY, - ) def _seed_lot(engine, lot_id: str, qty: str, nav: str, confirmed_at: datetime, @@ -174,415 +129,14 @@ def _services(engine): core_ro=CoreReadOnlyRepository(engine=engine), risk_repo=RiskRepository(engine=engine), convert_repo=ConvertRepository(engine=engine), + # 确认事务的 core 写侧也必须钉在同一 sqlite 引擎(缺省构造会连真 MySQL, + # CUST-T7 在真库无 customer 行 → FK 1452) core_writer=ConvertCoreRepository(engine=engine), + request_repo=ConvertRequestRepository(engine=engine), ) -def _req(customer: str = CUST, qty: str = "120", cid_req: str | None = None) -> dict: - return { - "customer_id": customer, - "from_product_id": PROD_OUT, - "to_product_id": PROD_IN, - "qty": qty, - "client_request_id": cid_req, - } - - -def _expected(engine, requested: str) -> dict: - """用生产纯函数算一遍期望值(与 service 内部同一套口径)。""" - core = CoreReadOnlyRepository(engine=engine) - lots = [Lot.from_row(r) for r in core.list_share_lots(CUST, PROD_OUT)] - rules = [FeeRule.from_row(r) for r in core.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) - rate = pick_fee_rate(rules, hold_days(TODAY, alloc.confirmed_at), product_id=PROD_OUT) - out_amount += amount - redeem_fee += lot_fee(amount, rate) - conv = convert_amount(out_amount, redeem_fee) - gap = diff_fee(conv, OUT_RATE, IN_RATE, "amount_diff") - in_amount = conv - gap - return { - "out_amount": out_amount, - "redeem_fee": redeem_fee, - "convert_amount": conv, - "diff_fee": gap, - "in_amount": in_amount, - "in_qty": in_qty(in_amount, IN_NAV), - "actual_qty": plan.actual_qty, - } - - -# ── 1. 三阶段贯通 ─────────────────────────────────────────────────── -def test_happy_path_writes_two_trades_and_completes_placeholder(sqlite_engine): - _seed(sqlite_engine) - exp = _expected(sqlite_engine, "120") - resp = convert_fund(_req(), now=NOW, **_services(sqlite_engine)) - - assert resp["blocked"] is False - assert resp["estimated"] is True - assert resp["convert_group_id"].startswith("CNV-") - assert Decimal(resp["out_amount"]) == exp["out_amount"] - assert Decimal(resp["redeem_fee"]) == exp["redeem_fee"] - assert Decimal(resp["convert_amount"]) == exp["convert_amount"] - assert Decimal(resp["diff_fee"]) == exp["diff_fee"] - assert Decimal(resp["in_amount"]) == exp["in_amount"] - assert Decimal(resp["in_qty"]) == exp["in_qty"] - assert Decimal(resp["actual_qty"]) == exp["actual_qty"] - assert resp["lot_count"] == 2 - assert [b["hold_days"] for b in resp["lot_breakdown"]] == [34, 3] - assert resp["confirm_basis"] == "natural_day_approx" - - # 阶段二:占位 completed - rows = _rows( - sqlite_engine, - "SELECT * FROM risk_convert_detail WHERE convert_group_id = :g", - g=resp["convert_group_id"], - ) - assert len(rows) == 1 and rows[0]["status"] == "completed" - # 两条流水同组、redeem/subscribe(R-b) - trades = _rows( - sqlite_engine, - "SELECT * FROM core_trade WHERE convert_group_id = :g ORDER BY trade_type", - g=resp["convert_group_id"], - ) - assert [t["trade_type"] for t in trades] == ["redeem", "subscribe"] - # 主审计 1 条(净值新鲜 → 无 nav_stale 副审计) - assert len( - _rows(sqlite_engine, "SELECT * FROM audit_log WHERE decision = 'convert_accepted'") - ) == 1 - assert _rows(sqlite_engine, "SELECT * FROM audit_log WHERE decision = 'nav_stale'") == [] - - -# ── 2. blocked 不占位(PRD §7.0)───────────────────────────────────── -def test_suitability_blocked_does_not_placeholder(sqlite_engine): - _seed(sqlite_engine) - resp = convert_fund( - _req(customer=CUST_LOW), now=NOW, **_services(sqlite_engine) - ) - assert resp["blocked"] is True - assert resp["block_response_code"] - # 关键:前四步不落库 —— 无占位、无流水 - assert _rows(sqlite_engine, "SELECT * FROM risk_convert_detail") == [] - assert _rows(sqlite_engine, "SELECT * FROM core_trade") == [] - assert len( - _rows(sqlite_engine, "SELECT * FROM audit_log WHERE decision = 'suitability_blocked'") - ) == 1 - - -# ── 3. 幂等:同键重复提交不产生第二组流水 ──────────────────────────── -def test_idempotent_repeat_returns_first_result(sqlite_engine): - _seed(sqlite_engine) - svc = _services(sqlite_engine) - first = convert_fund(_req(cid_req="CLI-T7-001"), now=NOW, **svc) - again = convert_fund(_req(cid_req="CLI-T7-001"), now=NOW, **svc) - - assert again["convert_group_id"] == first["convert_group_id"] - assert again["in_qty"] == first["in_qty"] - assert again["out_trade_id"] == first["out_trade_id"] - # 只有一组流水(2 条),没有第二组 - assert len(_rows(sqlite_engine, "SELECT * FROM core_trade")) == 2 - assert len(_rows(sqlite_engine, "SELECT * FROM risk_convert_detail")) == 1 - - -# ── 4. 未抢到执行权 → processing(T-9 映射 202)─────────────────────── -def test_lock_not_acquired_returns_processing(sqlite_engine, monkeypatch): - import app.service.convert.convert_service as cs - from app.service.risk.locks import _NoLock - - _seed(sqlite_engine) - monkeypatch.setattr(cs, "try_lock", lambda *a, **k: _NoLock()) - resp = convert_fund(_req(cid_req="CLI-T7-002"), now=NOW, **_services(sqlite_engine)) - - assert resp["status"] == PROCESSING - assert "convert_group_id" in resp - assert _rows(sqlite_engine, "SELECT * FROM core_trade") == [] - assert _rows(sqlite_engine, "SELECT * FROM risk_convert_detail") == [] - - -# ── 5. 全额转出豁免最低份额(验收 16)──────────────────────────────── -def test_full_transfer_waives_min_redeem_qty(sqlite_engine): - _seed(sqlite_engine, min_redeem="10000") - # 持 150(< min_redeem 10000)但申请全额 → 豁免,成功 - resp = convert_fund(_req(qty="150"), now=NOW, **_services(sqlite_engine)) - assert resp["blocked"] is False - assert Decimal(resp["actual_qty"]) == Decimal("150") - - -def test_below_min_qty_rejected_when_not_full(sqlite_engine): - _seed(sqlite_engine, min_redeem="10000") - # 持 150、申请 100(< 最低 10000,且非全额)→ 400 BELOW_MIN_QTY - with pytest.raises(BelowMinQty): - convert_fund(_req(qty="100"), now=NOW, **_services(sqlite_engine)) - - -# ── 6. 强制全转留痕 ───────────────────────────────────────────────── -def test_forced_full_transfer_flag(sqlite_engine): - _seed(sqlite_engine, min_hold="100") - # 持 150、申请 140 → 余额 10 < 100 → 强制全转 150 - resp = convert_fund(_req(qty="140"), now=NOW, **_services(sqlite_engine)) - assert resp["forced_full_transfer"] is True - assert resp["min_hold_action"] == "force_transfer" - assert Decimal(resp["actual_qty"]) == Decimal("150") - - -# ── 7. nav_stale 额外落副审计 ──────────────────────────────────────── -def test_nav_stale_adds_second_audit(sqlite_engine): - _seed(sqlite_engine, nav_date=date(2026, 8, 25)) # 距今 10 天 > 3 - resp = convert_fund(_req(), now=NOW, **_services(sqlite_engine)) - assert resp["nav_stale"] is True - assert len(_rows(sqlite_engine, "SELECT * FROM audit_log WHERE decision = 'nav_stale'")) == 1 - assert len( - _rows(sqlite_engine, "SELECT * FROM audit_log WHERE decision = 'convert_accepted'") - ) == 1 - - -# ── 8. 阶段 1.5:引擎异常不阻断已成立的交易 ─────────────────────────── -def test_engine_exception_does_not_block_trade(sqlite_engine): - _seed(sqlite_engine) - - def boom(out_trade, in_trade): # noqa: ANN001 - raise RuntimeError("引擎炸了") - - resp = convert_fund(_req(), now=NOW, engine_hook=boom, **_services(sqlite_engine)) - assert resp["blocked"] is False - assert resp["engine_error"] is True - assert len(_rows(sqlite_engine, "SELECT * FROM core_trade")) == 2 # 交易仍成立 - assert len(_rows(sqlite_engine, "SELECT * FROM audit_log WHERE decision = 'engine_error'")) == 1 - - -def test_engine_result_merged_into_response(sqlite_engine): - _seed(sqlite_engine) - hook = lambda o, i: { # noqa: E731 - "triggered_rules": ["RISK-002"], - "alert_ids": ["ALT-1"], - "aml_hit": False, - } - resp = convert_fund(_req(), now=NOW, engine_hook=hook, **_services(sqlite_engine)) - assert resp["triggered_rules"] == ["RISK-002"] - assert resp["alert_ids"] == ["ALT-1"] - - -# ── 9. 阶段二失败:留痕 + 占位 failed(不回滚 Core)─────────────────── -def test_phase_two_failure_keeps_trade_and_marks_failed(sqlite_engine, monkeypatch): - _seed(sqlite_engine) - - def boom(self, group_id, **kwargs): # noqa: ANN001 - raise RuntimeError("阶段二写失败") - - monkeypatch.setattr(ConvertRepository, "complete_convert", boom) - resp = convert_fund( - _req(cid_req="CLI-T7-003"), now=NOW, **_services(sqlite_engine) - ) - # 交易已成立:core 有两条流水,响应照常返回 - assert len(_rows(sqlite_engine, "SELECT * FROM core_trade")) == 2 - assert resp["blocked"] is False - assert len( - _rows( - sqlite_engine, - "SELECT * FROM audit_log WHERE decision = 'convert_detail_write_failed'", - ) - ) == 1 - rows = _rows(sqlite_engine, "SELECT * FROM risk_convert_detail") - assert rows[0]["status"] == "failed" - - -# ── 9b. 阶段一失败 → 同键重试必须可成功(T-13 前置修复的回归闸门)────── -def test_retry_after_phase_one_failure_succeeds_with_same_key( - sqlite_engine, monkeypatch, -): - """`LotConflict`(409) 后带**同一** `client_request_id` 重试 → 成功、只有一组流水。 - - 这是架构 §8.3「`LOT_CONFLICT` → 调用方重试,建议 ≤3 次、间隔 100/200/400ms」的 - 服务端契约。修复前:占位被 `mark_failed` 置 failed,重试复用同一 `group_id` - 再次进入阶段零,`insert_placeholder` 的朴素 INSERT 撞 `uk_group`/`uk_idem` - → `IdempotencyUnavailable`(503) —— **确定性失败**,重试永远不会成功。 - """ - _seed(sqlite_engine) - original = ConvertCoreRepository.apply_convert - - def conflict(self, req): # noqa: ANN001 - raise LotConflict("并发争抢:条件 UPDATE rowcount=0") - - monkeypatch.setattr(ConvertCoreRepository, "apply_convert", conflict) - with pytest.raises(LotConflict): - convert_fund(_req(cid_req="CLI-T13-RETRY"), now=NOW, **_services(sqlite_engine)) - monkeypatch.setattr(ConvertCoreRepository, "apply_convert", original) - - # 前置态:占位 failed、两条流水都没落 - pre = _rows(sqlite_engine, "SELECT * FROM risk_convert_detail") - assert len(pre) == 1 and pre[0]["status"] == "failed" - assert _rows(sqlite_engine, "SELECT * FROM core_trade") == [] - first_gid = pre[0]["convert_group_id"] - - # 带同键重试 → 必须成功,且复用原 group_id(杜绝第二组流水) - resp = convert_fund(_req(cid_req="CLI-T13-RETRY"), now=NOW, **_services(sqlite_engine)) - assert resp["convert_group_id"] == first_gid - assert len(_rows(sqlite_engine, "SELECT * FROM core_trade")) == 2 - assert _rows(sqlite_engine, "SELECT * FROM risk_convert_detail")[0]["status"] == "completed" - assert len(_rows(sqlite_engine, "SELECT * FROM risk_convert_detail")) == 1 # 不新增占位行 - - -def test_retry_after_phase_one_failure_can_fail_again_and_still_retry( - sqlite_engine, monkeypatch, -): - """连续两次 409 后再重试仍能成功 —— 证明「failed → pending」可反复回置。""" - _seed(sqlite_engine) - original = ConvertCoreRepository.apply_convert - state = {"n": 0} - - def conflict_twice(self, req): # noqa: ANN001 - state["n"] += 1 - if state["n"] <= 2: - raise LotConflict("第 %d 次争抢失败" % state["n"]) - original(self, req) - - monkeypatch.setattr(ConvertCoreRepository, "apply_convert", conflict_twice) - for _ in range(2): - with pytest.raises(LotConflict): - convert_fund(_req(cid_req="CLI-T13-RETRY2"), now=NOW, **_services(sqlite_engine)) - - resp = convert_fund(_req(cid_req="CLI-T13-RETRY2"), now=NOW, **_services(sqlite_engine)) - assert resp["blocked"] is False and resp["in_qty"] - assert len(_rows(sqlite_engine, "SELECT * FROM core_trade")) == 2 - assert [r["status"] for r in _rows(sqlite_engine, "SELECT * FROM risk_convert_detail")] == [ - "completed" - ] - - -def test_placeholder_uk_idem_race_returns_processing_without_writing( - sqlite_engine, monkeypatch, -): - """同键并发子窗口(两笔都判定"无占位")→ 撞 `uk_idem` 的那笔必须回 202,且零写入。 - - 构造方式:**只把"按 cid 查"这一读打桩成 None**(等价于并发下"那一行还没插进来"), - 占位表里预置另一笔同 cid 的占位。于是: - ① 幂等判定读不到 → 本笔自生成新 group_id; - ② 阶段零 INSERT 撞 `uk_idem` → `insert_placeholder` 返回 False → 本笔让路。 - - 这是**资金安全**的闸门:若不接住这个信号(仓储吞掉 False、调用方照跑), - 两笔会各带不同 group_id 跑完阶段一 —— `in_lot_id` 派生自 group_id 不再相撞 - → **双扣**(真库突变验证实测:120 份变 240 份、两组流水,见 - `tests/test_convert_concurrency.py::test_same_client_request_id_placeholder_race_must_not_leak_5xx`)。 - """ - _seed(sqlite_engine) - _exec( - sqlite_engine, - "INSERT INTO risk_convert_detail" - " (convert_group_id, client_request_id, status, estimated)" - " VALUES ('CNV-OTHER-CONC', 'CID-RACE', 'pending', 0)", - ) - monkeypatch.setattr( - ConvertRepository, "get_by_client_request_id", lambda self, cid: None - ) - - resp = convert_fund(_req(cid_req="CID-RACE"), now=NOW, **_services(sqlite_engine)) - - assert resp["status"] == PROCESSING # 让路 → 202(T-9 映射) - assert resp["convert_group_id"] is None - # 零写入:没有第二组流水、没有第二行占位、没有扣份额 - assert _rows(sqlite_engine, "SELECT * FROM core_trade") == [] - placeholders = _rows(sqlite_engine, "SELECT * FROM risk_convert_detail") - assert len(placeholders) == 1 and placeholders[0]["convert_group_id"] == "CNV-OTHER-CONC" - assert Decimal( - _rows( - sqlite_engine, - "SELECT SUM(remain_qty) AS s FROM core_share_lot WHERE customer_id = :c", - c=CUST, - )[0]["s"] - ) == Decimal("150") - - -# ── 10. 各 4xx / 503 分支 ─────────────────────────────────────────── -def test_same_product_rejected(sqlite_engine): - _seed(sqlite_engine) - req = _req() - req["to_product_id"] = PROD_OUT - with pytest.raises(SameProduct): - convert_fund(req, now=NOW, **_services(sqlite_engine)) - - -def test_cross_entity_rejected(sqlite_engine): - _seed(sqlite_engine) - _exec( - sqlite_engine, - "UPDATE core_product SET fund_company = '易方达模拟基金' WHERE product_id = :p", - p=PROD_IN, - ) - with pytest.raises(CrossEntityNotSupported): - convert_fund(_req(), now=NOW, **_services(sqlite_engine)) - - -def test_out_product_not_redeemable_rejected(sqlite_engine): - _seed(sqlite_engine) - _exec( - sqlite_engine, - "UPDATE core_product SET can_redeem = 0 WHERE product_id = :p", - p=PROD_OUT, - ) - with pytest.raises(ProductNotRedeemable): - convert_fund(_req(), now=NOW, **_services(sqlite_engine)) - - -def test_insufficient_shares_rejected(sqlite_engine): - _seed(sqlite_engine) - with pytest.raises(InsufficientShares): - convert_fund(_req(qty="500"), now=NOW, **_services(sqlite_engine)) - - -def test_nav_not_ready_503(sqlite_engine): - _seed(sqlite_engine) - _exec(sqlite_engine, "DELETE FROM core_product_nav") - with pytest.raises(NavNotReady) as exc: - convert_fund(_req(), now=NOW, **_services(sqlite_engine)) - assert exc.value.status_code == 503 - - -def test_too_many_lots_rejected(sqlite_engine, monkeypatch): - from app.config import settings as settings_module - - _seed(sqlite_engine) - monkeypatch.setattr(settings_module.settings, "convert_batch_max_lots", 1) - with pytest.raises(TooManyLots) as exc: - convert_fund(_req(qty="120"), now=NOW, **_services(sqlite_engine)) - assert exc.value.extra == {"batch_count": 2, "max_lots": 1} - - -# ── 10. 阶段 1.5 接线(T-8):引擎真跑并出单 ───────────────────────── -def test_engine_wired_produces_single_alert_with_two_events(sqlite_engine): - """T-8 落地后阶段 1.5 不再跳过:大额转换**真出一张单**、`payload.events` 两条(验收 7)。 - - 本用例是「接线回归」:若 `_run_engine` 又被改回静默跳过(或签名对不上被 - ImportError 吞掉),这里会因 `triggered_rules` 为空而变红。 - """ - _seed(sqlite_engine) - th = RiskThresholds( - large_amount=Decimal("100"), # 调低以让 120 份的折算额命中 RISK-001 - daily_total=Decimal("1000000"), - freq_count=3, - probe_window_minutes=5, - probe_count=3, - probe_amount=Decimal("400000"), - small_amount=Decimal("10000"), - small_count=3, - concentration_threshold=1.01, - ) - resp = convert_fund(_req(), now=NOW, thresholds=th, **_services(sqlite_engine)) - - assert resp["engine_error"] is False, "引擎真的跑了且没炸(未走 ImportError 跳过分支)" - assert "RISK-001" in resp["triggered_rules"] - assert len(resp["alert_ids"]) == 1, "一次转换只出一张单" - - 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"] - - -# ── 11. T-12 补偿:阶段二失败 → 按 group 补写详情 + 预警(FR-C17 / 架构 §5.4)── +# ── 11. T-12 补偿:确认段第 ⑦/⑧ 步失败 → 按 group 补写详情 + 预警 ───────── def _low_thresholds() -> RiskThresholds: """调低大额阈值,让 120 份的折算额足以命中 RISK-001(与 §10 接线用例同口径)。""" return RiskThresholds( @@ -717,9 +271,7 @@ def test_compensate_rebuilds_detail_and_alerts(sqlite_engine): ) assert len(audits) == 1 # 补偿审计可追溯:phase 标注补偿来源 + 引擎异常事实入 summary - import json as _json - - summary = _json.loads( + summary = json.loads( _rows( sqlite_engine, "SELECT input_summary FROM audit_log WHERE decision = 'confirmed'", @@ -802,4 +354,102 @@ def test_compensate_returns_locked_when_execution_right_taken(sqlite_engine): out = _compensate(sqlite_engine, gid) assert out["state"] == "locked" and out["alert_ids"] == [] assert _detail_rows(sqlite_engine, gid)[0]["status"] == "failed" # 仍待补偿 - assert _rows(sqlite_engine, "SELECT * FROM risk_alert") == [] + + +# ── 12. T+1 缺口补用(v1.0 用例退役后的承接,T-13)───────────────────── +def _seed_accepted_request(engine, cid_req: str, qty: str = "30000") -> str: + """直插一张 accepted 受理单(确认段输入态;日历缺失由 `_safe_next_biz_day` 容错)。 + + 两端净值都锚在受理日 T(`get_nav_on` 精确匹配);转出端净值取批次 nav 同值, + 使逐批折算与明细口径自洽。 + """ + _exec( + engine, + "INSERT INTO core_product_nav (product_id, nav, daily_chg_pct, nav_date) " + "VALUES (:p, 1.0300, 0, :d)", + p=PROD_OUT, d=TODAY, + ) + _exec( + engine, + "INSERT INTO core_convert_request (convert_group_id, client_request_id," + " customer_id, from_product_id, to_product_id, qty, status," + " cancel_before, requested_at)" + " VALUES (:gid, :cid_req, :c, :fp, :tp, :q, 'accepted', NULL, :rat)", + gid="CNV-T13-0001", cid_req=cid_req, c=CUST, fp=PROD_OUT, tp=PROD_IN, + q=float(qty), rat=NOW, + ) + return "CNV-T13-0001" + + +def test_confirm_mirror_write_failure_keeps_core_confirmed(sqlite_engine): + """确认第 ⑧ 步镜像写入失败 → Core 已 confirmed 不回滚、镜像停 pending(v1.0 + 「阶段二失败不回滚」语义的 T+1 承接;失败现场由补偿 T-12 收口,真库证据 + `verify_convert_compensate.py` A 组)。""" + _seed(sqlite_engine) + gid = _seed_accepted_request(sqlite_engine, "CLI-T13-MIRROR") + svc = _services(sqlite_engine) + + with patch.object( + ConvertRepository, "sync_mirror", side_effect=RuntimeError("T13 故意:镜像失败") + ): + out = confirm_one( + gid, + core_ro=svc["core_ro"], + risk_repo=svc["risk_repo"], + convert_repo=svc["convert_repo"], + core_writer=svc["core_writer"], + request_repo=svc["request_repo"], + thresholds=_low_thresholds(), + now=NOW, + as_of=TODAY, + ) + + assert out["status"] == "confirmed" # Core 事务不受 agent 故障影响 + assert len(_rows(sqlite_engine, "SELECT * FROM core_trade")) == 2 + req_status = _rows( + sqlite_engine, + "SELECT status FROM core_convert_request WHERE convert_group_id = :g", g=gid, + )[0]["status"] + assert req_status == "confirmed" + mirror = _detail_rows(sqlite_engine, gid) + assert mirror == [] or mirror[0]["status"] == "pending" # 第 ⑧ 步未推进(绝不写 completed) + + +def test_confirm_rejects_too_many_lots(sqlite_engine, monkeypatch): + """参与批次数超上限 → `rejected` + `TOO_MANY_LOTS`(占用释放,份额不变; + v1.0 `TooManyLots` 4xx 语义在确认段的承接)。""" + _seed(sqlite_engine) + # 再补 2 个小批次 → 共 4 批;把确认批上限钳到 2 + _seed_lot(sqlite_engine, "LOT-T13-3", "30", "1.0100", datetime(2026, 7, 1, 10, 0, 0)) + _seed_lot(sqlite_engine, "LOT-T13-4", "20", "1.0200", datetime(2026, 6, 1, 10, 0, 0)) + gid = _seed_accepted_request(sqlite_engine, "CLI-T13-LOTS") + monkeypatch.setattr(settings, "convert_batch_max_lots", 2) + svc = _services(sqlite_engine) + + out = confirm_one( + gid, + core_ro=svc["core_ro"], + risk_repo=svc["risk_repo"], + core_writer=svc["core_writer"], + request_repo=svc["request_repo"], + convert_repo=svc["convert_repo"], + thresholds=_low_thresholds(), + now=NOW, + as_of=TODAY, + ) + + assert out["status"] == "rejected" + assert out["reject_reason"] == "TOO_MANY_LOTS" + assert len(_rows(sqlite_engine, "SELECT * FROM core_trade")) == 0 # 无流水 + remain = _rows( + sqlite_engine, + "SELECT COALESCE(SUM(remain_qty), 0) AS s FROM core_share_lot " + "WHERE customer_id = :c AND product_id = :p", + c=CUST, p=PROD_OUT, + )[0]["s"] + assert Decimal(str(remain)) == Decimal("200") # 份额原封不动(占用自然释放) + req_status = _rows( + sqlite_engine, + "SELECT status FROM core_convert_request WHERE convert_group_id = :g", g=gid, + )[0]["status"] + assert req_status == "rejected"