基金转换 T-10:普通申赎批次维护(FR-C16)

让 core_trade 的普通申赎同步维护 core_share_lot 批次,为 convert 提供权威份额口径。

落地(改码 3 处 + 新增 2 个脚本)
- gateway_repository:新增 insert_share_lots / deduct_share_lots,各自单事务、走 rw 账号;
  只落已分配结果、不含分配逻辑,避免出现第二份 FIFO 口径
- share_lot_repository:__init__ 增可选 core_ro 入参,让 trade_gateway 复用同一份 FIFO 选批口径(D18)
- trade_gateway:新增 _nav_as_of / _maintain_lots,在落流水后、规则引擎前调用
  · subscribe:取 T 日(含)前最新净值 → qty = amount ÷ 净值(2 位 HALF_UP)→ 建批次(confirmed_at = T)
  · redeem:FIFO 扣减;无批次但有持仓 → D8 兜底补建后再扣;两者皆无 → 降级跳过
  · 整段 try 包住,任何异常只 warning,绝不阻断交易(R-c(1))
- 新增 scripts/core/rebuild_lots.py:按 core_holding 快照重建批次(L-7),与网关 D8 同调 bootstrap_lots;
  默认只补建无批次的持仓行(幂等),--force 才先删后建,--dry-run 只报告
- 新增 scripts/dev/verify_convert_lots.py:真库验证脚本(MySQL 8.0.46)

口径订正(联网查证 4 家管理人业务规则后)
- redeem 的 amount ÷ 净值 折算属本项目简化建模,不是行业标准
- 行业铁律是「金额申购、份额赎回」:投资者以份额申报,登记机构按 T 日净值算金额
  (睿远业务规则 §65 / 华泰保兴 §69 / 东方基金 §57 / 国投瑞银 §33;无一家公募支持按金额赎回)
- 派生风险:未知价法下 T 日净值当日不可得(T+1 公告),_nav_as_of 实取 T−1 净值,
  故此处算出的份额只是估算值
- 已记入 PRD §10.1 已知差异清单;trade_gateway 注释同步订正

验证
- pytest -q → 714 passed / 3 skipped(基线 697 加 17,零回归)
- R16 test_redeem_accepted_without_alert 零改动通过(由 R-c(1) 降级保住)
- 真库 20/20;T-6 24/24、T-7 35/35 复跑零回归;真库隔离数据零残留
- 突变验证 2 组:切断接线 → 精准 2 条红;关掉 D8 兜底 → 精准 2 条红(含 D18 同源断言)
This commit is contained in:
2026-09-10 18:28:48 +08:00
parent e80098b32f
commit 5e6fa06d03
12 changed files with 1383 additions and 22 deletions
+74 -1
View File
@@ -2,20 +2,52 @@
仅本类可 INSERT jinrong_core.core_trade;core_ro 仍只读。生产环境由真实 仅本类可 INSERT jinrong_core.core_trade;core_ro 仍只读。生产环境由真实
交易系统回调替代,本类随 app/gateway/ 退役。 交易系统回调替代,本类随 app/gateway/ 退役。
**T-10 起并承担 `core_share_lot` 的普通申赎批次维护写侧**(架构 §12 Q1:
`core_share_lot` 的写入方 = `convert_core_repository`(convert 转换)+
本类(普通申赎,含 D8 兜底补建)。分工刻意如此:
- **FIFO 选批/分配口径仍只有一处** —— `share_lot_repository`(D18);
本类**不含任何分配逻辑**,只把已分配好的结果落库,避免出现第二份口径。
- 批次维护属**尽力而为**(FR-C16 / R-c(1)):数据不全时由调用方
(`trade_gateway`)warning 跳过,**不抛异常、不阻断交易**。
""" """
from __future__ import annotations from __future__ import annotations
from datetime import datetime from datetime import datetime
from decimal import Decimal from decimal import Decimal
from typing import Any from typing import Any, Sequence
from sqlalchemy import text from sqlalchemy import text
from sqlalchemy.engine import Engine from sqlalchemy.engine import Engine
from app.config.settings import settings from app.config.settings import settings
from app.service.convert.types import Lot
from app.utils.db import get_engine from app.utils.db import get_engine
# ── 批次写入 SQL(模块级常量:口径集中,便于单测逐条核对)──────────────
# ⚠️ `(:q + 0.0)` **不能省**:数值参数统一以字符串绑定(`str()`,跨库无损),
# 而 sqlite 在 **UPDATE 的算术表达式**里不会把 TEXT 绑定参数转数值 ——
# `remain_qty = remain_qty - :q` 传 '50' 会扣 0(见开发计划 §B.8 第 11 条)。
# 加 `+ 0.0` 强制数值化后两库行为一致(MySQL 本就隐式转换)。
_SHARE_LOT_INSERT = text(
"""
INSERT INTO core_share_lot
(lot_id, customer_id, product_id, qty, remain_qty, nav,
confirmed_at, source_trade_id)
VALUES (:lot, :cid, :pid, :q, :q, :nav, :cat, :src)
"""
)
_SHARE_LOT_DEDUCT = text(
"""
UPDATE core_share_lot
SET remain_qty = remain_qty - (:q + 0.0)
WHERE lot_id = :lot
"""
)
class GatewayRepository: class GatewayRepository:
"""core_trade 唯一写入口;仅 INSERT,不改不删(模拟网关语义)。 """core_trade 唯一写入口;仅 INSERT,不改不删(模拟网关语义)。
@@ -61,3 +93,44 @@ class GatewayRepository:
"at": traded_at, "at": traded_at,
}, },
) )
# ── 批次维护写侧(T-10 · FR-C16 · 架构 §12 Q1)────────────────────
def insert_share_lots(self, lots: Sequence[Lot]) -> None:
"""批量写入 `core_share_lot`(**单事务**)—— 申购建批次 / D8 兜底补建共用。
`qty` 与 `remain_qty` 写**同一个值**(新批次未发生任何扣减);
取值统一用 `remain_qty`(`bootstrap_lots` 保证两者相等)。
`nav` 以字符串绑定:sqlite 列有 NUMERIC affinity 会转换,MySQL 精确。
"""
if not lots:
return
with self._engine.begin() as conn:
for lot in lots:
conn.execute(
_SHARE_LOT_INSERT,
{
"lot": lot.lot_id,
"cid": lot.customer_id,
"pid": lot.product_id,
"q": str(lot.remain_qty),
"nav": str(lot.nav),
"cat": lot.confirmed_at,
"src": lot.source_trade_id,
},
)
def deduct_share_lots(self, deductions: Sequence[tuple[str, Decimal]]) -> None:
"""按 FIFO 分配结果逐批扣减 `remain_qty`(**单事务**)。
⚠️ 与 `convert_core_repository._deduct_lots` 的**关键差别**:本方法
**不带 `remain_qty >= :q` 哨兵**。传入的份额已由 `share_lot_repository`
按「不超过可用余额」裁剪,普通赎回**不因份额不足阻断交易**(R-c(1),
与既有 redeem 语义一致);而 convert 的份额不足是**严格校验**(409
`LotConflict`)。两处口径不可互相套用。
"""
if not deductions:
return
with self._engine.begin() as conn:
for lot_id, qty in deductions:
conn.execute(_SHARE_LOT_DEDUCT, {"lot": lot_id, "q": str(qty)})
+170 -3
View File
@@ -3,8 +3,9 @@
`submit_trade` 为唯一入口:**按 `trade_type` 分派**—— `submit_trade` 为唯一入口:**按 `trade_type` 分派**——
- `subscribe` / `redeem`:参数校验 → suitability_check(FR-2,落校验日志; - `subscribe` / `redeem`:参数校验 → suitability_check(FR-2,落校验日志;
不匹配→ suitability 预警单 + 阻断响应,交易不落 core_trade)→ 匹配 → 不匹配→ suitability 预警单 + 阻断响应,交易不落 core_trade)→ 匹配 →
INSERT core_trade → 同步调规则引擎(FR-3)→ 返回 blocked + trade_id + INSERT core_trade → **份额批次维护(FR-C16 · T-10)** → 同步调规则引擎(FR-3)
触发规则。阻断/放行全量审计(agent_type='platform',FR-1 §6); → 返回 blocked + trade_id + 触发规则。阻断/放行全量审计
(agent_type='platform',FR-1 §6);
- `convert`(T-9 起走通):交 `convert_service.convert_fund` 八步编排 - `convert`(T-9 起走通):交 `convert_service.convert_fund` 八步编排
(含转入端适当性、幂等、三阶段、阶段 1.5)。**审计由该服务落 (含转入端适当性、幂等、三阶段、阶段 1.5)。**审计由该服务落
`convert_request`**,本层不重复写 `trade_request`;未抢到执行权时 `convert_request`**,本层不重复写 `trade_request`;未抢到执行权时
@@ -33,7 +34,11 @@ from app.gateway.gateway_repository import GatewayRepository
from app.repository.convert_repository import ConvertRepository from app.repository.convert_repository import ConvertRepository
from app.repository.core_ro import CoreReadOnlyRepository from app.repository.core_ro import CoreReadOnlyRepository
from app.repository.risk_repository import RiskRepository from app.repository.risk_repository import RiskRepository
from app.repository.share_lot_repository import ShareLotRepository
from app.service.convert.calc import round2
from app.service.convert.convert_service import convert_fund from app.service.convert.convert_service import convert_fund
from app.service.convert.lot_bootstrap import bootstrap_lots
from app.service.convert.types import Lot, to_decimal
from app.service.risk.engine import process_trade_event from app.service.risk.engine import process_trade_event
from app.service.risk.rules import RiskThresholds from app.service.risk.rules import RiskThresholds
from app.service.risk.alert_service import record_suitability_alert from app.service.risk.alert_service import record_suitability_alert
@@ -42,8 +47,13 @@ from app.utils.trace import ensure_trace
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
SUPPORTED_TRADE_TYPES = ("subscribe", "redeem")
CONVERT = "convert" CONVERT = "convert"
SUBSCRIBE = "subscribe"
REDEEM = "redeem"
SUPPORTED_TRADE_TYPES = (SUBSCRIBE, REDEEM)
#: 普通申购批次的 `lot_id` 前缀 —— 与转换转入批次(`LOT-CNV-…`)、
#: 兜底补建批次(`LOT-BOOT-…`)在库里可一眼区分来源。
SUBSCRIBE_LOT_PREFIX = "LOT-SUB"
ADVICE = "请联系持证投资顾问" ADVICE = "请联系持证投资顾问"
RECORDED_NOTICE = "本次请求已记录" RECORDED_NOTICE = "本次请求已记录"
@@ -133,6 +143,152 @@ def _submit_convert(
) )
def _nav_as_of(
core: CoreReadOnlyRepository, product_id: str, traded_at: datetime
) -> Decimal | None:
"""取 T 日(含)之前最新一期净值(D10,未知价法);缺失或非正数返回 None(触发降级跳过)。
⚠️ 未知价法下 **T 日净值当日不可得**(T+1 公告),真实场景一般落回 **T−1**,
故调用方须把据此算出的份额视为**估算值**(见 `_maintain_lots` 说明)。
"""
row = core.get_nav_as_of(product_id, traded_at.date())
if not row:
return None
nav = to_decimal(row["nav"])
return nav if nav > 0 else None
def _maintain_lots(
*,
writer: GatewayRepository,
core: CoreReadOnlyRepository,
trade_id: str,
customer_id: str,
product_id: str,
trade_type: str,
amount: Decimal,
traded_at: datetime,
) -> None:
"""普通申赎的份额批次维护(FR-C16 · T-10 · 架构 §12 Q1)。
**尽力而为,绝不阻断交易**(R-c(1)):整段被 `try` 包住,源数据缺失只
`logger.warning`。原因是批次表虽为 convert 正确性前提,但既有测试库
(sqlite `env` fixture)**无净值、无持仓**,若此处「缺数据即报错」会直接
打穿既有用例,违反开发计划 §1.2 硬约束。降级 ≠ 不覆盖:两条**真实路径**
由 `tests/test_share_lot.py` 自建完整种子单独覆盖(R-c(2))。
两条真实路径:
- `subscribe`:取 T 日(含)前最新净值 → `qty = amount ÷ nav`(2 位 HALF_UP)→ 建批次
(`confirmed_at = T`)。
- `redeem`:FIFO 扣减;**无批次但有持仓 → D8 兜底补建后再扣**(主场景,
必须执行);既无批次也无持仓 → 降级跳过。
⚠️ **redeem 的份额折算是「本项目简化」,不是行业标准**(2026-09-10 联网查证后订正):
行业铁律是「**金额申购、份额赎回**」—— 投资者以**份额**申报,登记机构按 T 日净值
反算金额(`赎回金额 = 赎回份额 × T日净值 − 费用`)。依据:睿远业务规则 §65
「以份额方式提出赎回申请」/ 华泰保兴 §69 / 东方基金 §57 / 永赢招募书;
**无一家公募支持「按金额赎回」**(货基面值 1 元、份额数值 = 金额,易被误读)。
本项目 `TradeRequest.amount` 是 `subscribe` / `redeem` **共用入参**
(`app/api/simulate.py:60`),故此处必须 `amount ÷ nav` 反算份额才能扣批次 ——
**方向与行业标准相反,属项目侧简化建模**。详见 PRD §10.1 已知差异。
⚠️ **由此派生的真实偏差(必须知晓)**:行业里份额由投资者申报、无歧义;本项目份额是
网关**反算**的,而未知价法下 **T 日净值当日不可得**(T+1 才公告),`_nav_as_of` 实取
**T−1 净值** → **此处算出的份额只是估算值**,与最终确认份额存在偏差。真实 Core 接入后
份额应由 TA 按 T 日净值确认,本折算逻辑随之退役。
取不到净值 → 降级跳过(不猜、不阻断)。
⚠️ **FIFO 分配只有一处实现**:补建批次**先落库、再读回**是刻意的 ——
这样选批仍由 `share_lot_repository` 独占(D18 / 自检第 13 问),
本模块不自带第二份贪心逻辑。
"""
try:
nav = _nav_as_of(core, product_id, traded_at)
if trade_type == SUBSCRIBE:
if nav is None:
logger.warning(
"批次维护跳过(申购取不到 T 日净值): trade_id=%s product=%s",
trade_id,
product_id,
)
return
qty = round2(amount / nav)
if qty <= 0:
return
writer.insert_share_lots(
[
Lot(
lot_id=f"{SUBSCRIBE_LOT_PREFIX}-{trade_id}",
customer_id=customer_id,
product_id=product_id,
qty=qty,
remain_qty=qty,
nav=nav,
confirmed_at=traded_at,
source_trade_id=trade_id,
)
]
)
return
# ── redeem ────────────────────────────────────────────────────
if nav is None:
logger.warning(
"批次维护跳过(赎回取不到 T 日净值,无法折算赎回份额): "
"trade_id=%s product=%s",
trade_id,
product_id,
)
return
qty = round2(amount / nav)
if qty <= 0:
return
lots_repo = ShareLotRepository(core_ro=core)
if lots_repo.available_qty(customer_id, product_id) <= 0:
# D8 兜底补建(lot_bootstrap 单点,D18)
holding = core.get_holding(customer_id, product_id)
if holding is None:
logger.warning(
"批次维护跳过(赎回既无批次也无持仓): trade_id=%s cid=%s pid=%s",
trade_id,
customer_id,
product_id,
)
return
bootstrap = bootstrap_lots(holding)
if not bootstrap:
logger.warning(
"批次维护跳过(赎回持仓份额为 0,无可补建批次): "
"trade_id=%s cid=%s pid=%s",
trade_id,
customer_id,
product_id,
)
return
writer.insert_share_lots(bootstrap)
selected = lots_repo.select_for_convert(customer_id, product_id, qty)
if not selected:
logger.warning(
"批次维护跳过(赎回无可用批次可扣): trade_id=%s cid=%s pid=%s",
trade_id,
customer_id,
product_id,
)
return
writer.deduct_share_lots(
[(item["lot_id"], to_decimal(item["qty"])) for item in selected]
)
except Exception: # 批次维护不得反噬交易主流程(R-c(1))
logger.warning("批次维护失败(交易仍成立): trade_id=%s", trade_id, exc_info=True)
def submit_trade( def submit_trade(
req: dict[str, Any], req: dict[str, Any],
core_ro: CoreReadOnlyRepository | None = None, core_ro: CoreReadOnlyRepository | None = None,
@@ -225,6 +381,17 @@ def submit_trade(
writer.insert_trade( writer.insert_trade(
trade_id, req["customer_id"], req["product_id"], trade_type, amount, traded_at trade_id, req["customer_id"], req["product_id"], trade_type, amount, traded_at
) )
# FR-C16:普通申赎的批次维护(含 D8 兜底补建)—— 在任何情况下都不阻断交易
_maintain_lots(
writer=writer,
core=core,
trade_id=trade_id,
customer_id=req["customer_id"],
product_id=req["product_id"],
trade_type=trade_type,
amount=amount,
traded_at=traded_at,
)
try: try:
engine_result = process_trade_event( engine_result = process_trade_event(
{ {
+19 -5
View File
@@ -21,12 +21,26 @@ from app.utils.db import get_engine
class ShareLotRepository: class ShareLotRepository:
"""core_share_lot 读侧:FIFO 选批 + 汇总(D2:写侧在 convert_core_repository)。""" """core_share_lot 读侧:FIFO 选批 + 汇总(D2:写侧在 convert_core_repository)。
def __init__(self, engine: Engine | None = None) -> None: `core_ro` 入参(T-10 补入):普通申赎的批次维护(`trade_gateway`)已经持有
self._engine = engine or get_engine(settings.mysql_core_database, "ro") 一个 `CoreReadOnlyRepository`(可能绑 sqlite 测试引擎),直接传入即可复用
# 复用 core_ro 的 ORDER BY(D18 单一副本),避免排序规则漂移 **同一份 FIFO 分配口径** —— 否则 `trade_gateway` 只能另写一份贪心逻辑,
self._core = CoreReadOnlyRepository(engine=self._engine) 正是自检第 13 问「同一规则只留一个副本」要防的漂移。传 `core_ro` 时
`engine` 参数被忽略(读侧一律委托该实例)。
"""
def __init__(
self,
engine: Engine | None = None,
core_ro: CoreReadOnlyRepository | None = None,
) -> None:
if core_ro is not None:
self._core = core_ro
else:
self._core = CoreReadOnlyRepository(
engine=engine or get_engine(settings.mysql_core_database, "ro")
)
def select_for_convert( def select_for_convert(
self, customer_id: str, product_id: str, qty: Decimal self, customer_id: str, product_id: str, qty: Decimal
+15
View File
@@ -938,6 +938,21 @@ SELECT 1 FROM core_trade WHERE convert_group_id = :gid LIMIT 1;
| Q11 | T+1 确认时序怎么建模? | **不模拟 T+1 权益登记延迟(转入份额立即可用),但转入新批次 `confirmed_at` 取 T+1** | 持有期自确认日起算,比 T 起算少 1 天、费率档更严(B-5);存量/普通申购批次仍按 T | | Q11 | T+1 确认时序怎么建模? | **不模拟 T+1 权益登记延迟(转入份额立即可用),但转入新批次 `confirmed_at` 取 T+1** | 持有期自确认日起算,比 T 起算少 1 天、费率档更严(B-5);存量/普通申购批次仍按 T |
| Q12 | 最低持有余额触发后怎么处置? | **两种都做**:`core_product.min_hold_action` `ENUM('force_transfer','force_redeem')` | 阈值大(100/500/1000 份)→ 强制赎回剩余(财通/汇添富);阈值小(0.1/1 份)→ 强制全转(中银/东海),由种子按产品设定(I-5) | | Q12 | 最低持有余额触发后怎么处置? | **两种都做**:`core_product.min_hold_action` `ENUM('force_transfer','force_redeem')` | 阈值大(100/500/1000 份)→ 强制赎回剩余(财通/汇添富);阈值小(0.1/1 份)→ 强制全转(中银/东海),由种子按产品设定(I-5) |
### 10.1 已知差异清单(本项目简化 vs 真实业务)
> **性质**:以下均为**有意简化**,非缺陷。单列此表,是为防止后续把「项目做法」误当「行业标准」。
> **§2.5.2 引用的「已知差异」即本表**(原引用悬空,本次补齐)。接真实 Core / TA 时须逐条评估。
| # | 项 | 真实业务怎么做 | 本项目怎么做 | 依据 / 出处 |
| --- | --- | --- | --- | --- |
| ① | **赎回申报方向** | **份额赎回**——投资者以**份额**申报,登记机构按 T 日净值**算金额**:`赎回金额 = 赎回份额 × T日净值 − 赎回费用` | **金额入参**——`TradeRequest.amount` 为申赎共用入参,网关按 `amount ÷ 净值` **反算份额**再扣批次(方向相反) | 睿远业务规则 §65 / 华泰保兴 §69 / 东方基金 §57 / 国投瑞银(货基)§33。实现:`app/gateway/trade_gateway.py::_maintain_lots` |
| ② | **折算份额的可信度** | 份额由投资者申报,**无歧义** | 份额由网关**反算**;未知价法下 **T 日净值当日不可得**(T+1 公告),实取 **T−1 净值** → **份额为估算值** | 派生自 ①;`_nav_as_of` 取 `nav_date <= T` 降序 1 条 |
| ③ | **赎回费用舍入** | 个别管理人规定**费用 2 位后直接舍去**(截断),仅赎回**金额**四舍五入 | 费用与金额**统一 `ROUND_HALF_UP`** | 东方基金业务规则 §59。误差归属一致:均计入基金财产 |
| ④ | **转入份额位数** | 不统一:易方达 ETF 场外取**整数位**、南方基金**两位后截断** | 取主流口径 **2 位四舍五入** | 原见 §2.5.2,此处汇总 |
> **① 的连带影响**:份额是反算的,则 `core_share_lot.remain_qty` 的精度链比真实 TA 多一环
> (净值 → 份额)。接真实 Core 后,份额应由 TA 按 T 日净值确认,`amount ÷ nav` 逻辑随之退役。
--- ---
## 11. 并发控制(FR-C18) ## 11. 并发控制(FR-C18)
+49
View File
@@ -430,3 +430,52 @@ skill `design-doc-selfcheck` → **v1.2.0**:新增 **「二补 · 格式契约
`.workbuddy/memory/MEMORY.md`(远程铁律条**重写**,含 7 条证据 + 绕过方案)· `.workbuddy/memory/MEMORY.md`(远程铁律条**重写**,含 7 条证据 + 绕过方案)·
`docs/memory/MEMORY.md`(交接清单**第 7 问重写** + 顺手订正验证条目里过期的「516 passed」→ **697** + 第 73 行旧表述订正)· 本条。 `docs/memory/MEMORY.md`(交接清单**第 7 问重写** + 顺手订正验证条目里过期的「516 passed」→ **697** + 第 73 行旧表述订正)· 本条。
---
## T-10(普通申赎批次维护 FR-C16)完成(2026-09-10 · 用户「进入 T10 阶段」)
**开工前基线快照**(计划硬要求):`pytest -q` → **697 passed / 3 skipped**
(开发计划 §8 写的「当前 516」是撰写期快照,已过期 → 已在文档中订正)。
**改码 4 处**:`gateway_repository.py`(+`insert_share_lots`/`deduct_share_lots`,各单事务、走 `rw`,
**不含分配逻辑**)· `share_lot_repository.py`(`__init__` 增可选 `core_ro` → 复用**同一份** FIFO 口径)·
`trade_gateway.py`(+`_nav_as_of`/`_maintain_lots`,在落流水后、引擎前调用,整段 try 包住)·
**新增 `scripts/core/rebuild_lots.py`**(快照重建批次 L-7,与网关 D8 同调 `bootstrap_lots`,`role="admin"`)。
**两条执行期裁定(计划未点明,已留痕)**
1. **`redeem` 份额折算 = `amount ÷ 净值`(2 位 HALF_UP)** —— 计划只写「FIFO 扣减」,
未说扣减的**份额**从哪来(赎回请求体给的是金额)。取不到净值 → 降级跳过。
> **⚠️ 2026-09-10 18:10 订正(联网查证 4 家管理人业务规则后)**:本条初稿写「按未知价法」,
> **定性错误**。真实行业铁律 = 「**金额申购、份额赎回**」—— 投资者以**份额**申报,
> 登记机构按 T 日净值**算金额**(睿远 §65 / 华泰保兴 §69 / 东方基金 §57 / 国投瑞银(货基)§33,
> **无一家公募支持「按金额赎回」**)。本项目是**简化建模**(`TradeRequest.amount` 为申赎共用
> 入参,`app/api/simulate.py:60`),方向与行业相反 → 已记入 **PRD §10.1 已知差异 ①②**。
> 另发现:**T 日净值当日不可得**(T+1 公告),`_nav_as_of` 实取 **T−1 净值** →
> **折算份额只是估算值**,与最终确认份额有偏差。折算逻辑**保留**,定性已同步订正
> (`trade_gateway.py` 注释 / 开发计划 §8 / 交接文档 §B / PRD §10.1)。
2. **`rebuild_lots.py` 默认只补建无批次行**(幂等、不误删转换转入的真实批次),`--force` 才先删后建。
**R-c 两条落地**:R-c(1) 降级(整段 try,源数据缺失只 warning,**绝不阻断交易**)→
既有 `env`(无净值无持仓)自然走降级,**§12 R16 `test_redeem_accepted_without_alert` 零改动通过**;
R-c(2) 覆盖率补偿(真实路径自建完整种子,**不得让降级充当覆盖**)→ `test_share_lot.py` 新增 14 条。
**验证**:`pytest -q` → **714 passed / 3 skipped**(**+17**,零回归);
**新增 `scripts/dev/verify_convert_lots.py` 真库 20/20**(隔离 `CUST-LOTT`/`PROD-LOTT*`,跑完零残留);
T-6 `24/24` · T-7 `35/35` 复跑零回归;**突变验证 2 组**(切断接线 → 精准 2 条红;关 D8 兜底 → 精准 2 条红含 D18),
均已恢复、无残留。
**★ 真库脚本的核心价值**:扣减 SQL 的 `(:q + 0.0)` 是**为绕开 sqlite 的 TEXT 绑定坑**而加(§B.8 第 11 条);
它在 **MySQL 上是否同样生效,sqlite 侧无从证明** —— 若 MySQL 未隐式转数值,扣减会静默失效
(批次永不减少 → convert 超扣)。**这正是用户 2026-09-10 定的「涉及方言语义的任务必须上真库」的意义**。
**D18 双保险**:除「同一持仓两侧算出的批次逐字段相等」外,另加**机制断言**
(`rebuild.bootstrap_lots is bootstrap_lots`)—— 防「两侧各写一份但碰巧同值」的假绿。
**未顺手改(已上报)**:`core_holding` 不随普通申赎更新(T-10 只维护 `core_share_lot`,计划范围即如此)。
**文档回写**:开发计划 §8(DoD 全勾 + 执行记录 + 2 条裁定 + 基线订正)· `交接文档.md`(§0 导航 / §B **v2.1** /
B.1 状态 / B.6 拓扑与基线 / **新增 B.6.3**)· `docs/memory/{TODO,MEMORY}` · `.workbuddy/memory/2026-09-10.md` · 本条。
**下一步 = T-11**(`core_tools.query_recent_trades` 汇总走 `_amount_view` FR-C15 + 持仓过滤 `qty <= 0` +
`core_ro.sum_trades_on_date` 加 convert 去重 **R-d**;DoD 见开发计划 §7.3)。
+2 -2
View File
@@ -199,8 +199,8 @@ RBAC 联调账号:scripts/dev/rbac-seed-reference.md
2. 改动属于 api / service / tool / repository 哪一层? 2. 改动属于 api / service / tool / repository 哪一层?
3. 是否需 customer_id 归属与 JWT RBAC? 3. 是否需 customer_id 归属与 JWT RBAC?
4. Core 是模拟库只读还是 agent 库读写? 4. Core 是模拟库只读还是 agent 库读写?
5. 如何验证?(`python -m pytest` 全量(当前 **697 passed / 3 skipped**,基线 510)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) 5. 如何验证?(`python -m pytest` 全量(当前 **714 passed / 3 skipped**,基线 510)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列)
6. **当前有哪两条并行线?**(① 风控/架构改进线:**已结项**(510 基线绿、§7.2 七项冒烟 7/7 PASS、`037ce7e` 已核实早已推送);② **基金转换线**:设计闭环,**第 5 步进行中 —— T-0~T-9 已完成(697 passed),下一步 = T-10 / T-11**)——动代码前先确认自己属于哪条线,别混淆前置条件。 6. **当前有哪两条并行线?**(① 风控/架构改进线:**已结项**(510 基线绿、§7.2 七项冒烟 7/7 PASS、`037ce7e` 已核实早已推送);② **基金转换线**:设计闭环,**第 5 步进行中 —— T-0~T-10 已完成(714 passed),下一步 = T-11**)——动代码前先确认自己属于哪条线,别混淆前置条件。
7. **远程分支到底还在不在?**(**在**。`git ls-remote --heads origin` 实测 `refs/heads/risk-control-agent` = `fffb78a`。⚠️ **判断远程存亡只能用 `git ls-remote`** —— 2026-09-10 **同日误报两次**,勿再踩。) 7. **远程分支到底还在不在?**(**在**。`git ls-remote --heads origin` 实测 `refs/heads/risk-control-agent` = `fffb78a`。⚠️ **判断远程存亡只能用 `git ls-remote`** —— 2026-09-10 **同日误报两次**,勿再踩。)
**本仓特有异常(2026-09-10 深挖确认,别误判为 stale ref)**:`.git/refs/remotes/**` **写入不落盘** —— **本仓特有异常(2026-09-10 深挖确认,别误判为 stale ref)**:`.git/refs/remotes/**` **写入不落盘** ——
`git update-ref refs/remotes/origin/X <sha>` 返回 0,但松散引用消失,**且整个 `refs/remotes/origin/` 目录被删** `git update-ref refs/remotes/origin/X <sha>` 返回 0,但松散引用消失,**且整个 `refs/remotes/origin/` 目录被删**
+2 -2
View File
@@ -7,7 +7,7 @@
**阶段一 AL-01~AL-08 与阶段二 C4~C6 均已完成(2026-09-07)**:全量 pytest **482 passed 0 failed 0 skipped**(真库集成)✓ · uvicorn 冒烟三端点 ✓ · 接口实调验收 ✓(2026-09-07:suitability/check 阻断+放行、simulate/trade 阻断、三条鉴权边界 401/403/403,契约零偏差)· risk-m1~m4 tag 齐。**合并 main 已移交合并执行人**(操作手册:《docs/项目框架设计/合并注意事项-风控模块并入main.md》,随分支上传),后续模块侧待办见下方。 **阶段一 AL-01~AL-08 与阶段二 C4~C6 均已完成(2026-09-07)**:全量 pytest **482 passed 0 failed 0 skipped**(真库集成)✓ · uvicorn 冒烟三端点 ✓ · 接口实调验收 ✓(2026-09-07:suitability/check 阻断+放行、simulate/trade 阻断、三条鉴权边界 401/403/403,契约零偏差)· risk-m1~m4 tag 齐。**合并 main 已移交合并执行人**(操作手册:《docs/项目框架设计/合并注意事项-风控模块并入main.md》,随分支上传),后续模块侧待办见下方。
**⚡ 并行新线 · 基金转换(convert)**(2026-09-10):**设计 + 开发计划均已闭环** —— PRD **v0.9.2** + 架构 **v1.0.1** + 独立评审 13 条 **0 悬空**(接受 10 / 修正性接受 3 / 驳回 0),门控 **M-7 已满足**;**第 4 步开发计划 v1.0 已产出并经独立审核**(4 条意见全接受、**驳回 0**,含新增 2 条回归面 R15/R16 + R-c 双条修订);**第 5 步进行中**:**T-0 ~ T-9 均已于 2026-09-10 完成**(**697 passed / 3 skipped**;T-1 断言 **8/8 PASS**、T-2 纯函数 **93 用例**、T-2b 实算 **15/15 一致**、T-6 真库 **24/24**、T-7 真库 **35/35**、**T-8 真库 31/31**、**T-9 真 MySQL 集成 8 条 + 4 组突变验证 + 展示位数修复**),**下一步 = T-10(普通申赎批次维护 · 回归风险最大)/ T-11(`core_tools` 与 `sum_trades_on_date` 汇总去重 · 依赖 T-8 已解锁)**。两个阻断前置(**T-0** sqlite/MySQL 列名统一 + 建库自校验 · **T-0b** DB 账号分离 D20:`xh_core_ro`/`xh_core_rw`/`xh_agent_rw`)**均已落地**。**入口:项目根 `交接文档.md` §B(三线合并版唯一入口,读这一节即可开工)**;**开工前必读开发计划 §1.4(15 条代码事实)+ §1.5(8 条实现级裁定 R-a~R-h)+ §12(16 条回归面)**。 **⚡ 并行新线 · 基金转换(convert)**(2026-09-10):**设计 + 开发计划均已闭环** —— PRD **v0.9.2** + 架构 **v1.0.1** + 独立评审 13 条 **0 悬空**(接受 10 / 修正性接受 3 / 驳回 0),门控 **M-7 已满足**;**第 4 步开发计划 v1.0 已产出并经独立审核**(4 条意见全接受、**驳回 0**,含新增 2 条回归面 R15/R16 + R-c 双条修订);**第 5 步进行中**:**T-0 ~ T-10 均已于 2026-09-10 完成**(**714 passed / 3 skipped**;T-1 断言 **8/8 PASS**、T-2 纯函数 **93 用例**、T-2b 实算 **15/15 一致**、T-6 真库 **24/24**、T-7 真库 **35/35**、**T-8 真库 31/31**、**T-9 真 MySQL 集成 8 条 + 4 组突变验证 + 展示位数修复**、**T-10 新增 17 用例 + 真库 20/20 + 2 组突变验证**),**下一步 = T-11(`core_tools` 与 `sum_trades_on_date` 汇总去重 · 依赖 T-8 已解锁)**。两个阻断前置(**T-0** sqlite/MySQL 列名统一 + 建库自校验 · **T-0b** DB 账号分离 D20:`xh_core_ro`/`xh_core_rw`/`xh_agent_rw`)**均已落地**。**入口:项目根 `交接文档.md` §B(三线合并版唯一入口,读这一节即可开工)**;**开工前必读开发计划 §1.4(15 条代码事实)+ §1.5(8 条实现级裁定 R-a~R-h)+ §12(16 条回归面)**。
### 基金转换线待办(推荐顺序) ### 基金转换线待办(推荐顺序)
@@ -21,7 +21,7 @@
- [ ] **【待用户裁定】幂等重放响应的数值位数偏差**(T-9 执行期发现,**未顺手改**):首次响应 2 位(`calc` 量化)vs 重放响应 4 位(`core_trade` `DECIMAL(18,4)` 直读)→ `"53456.95"` vs `"53456.9500"`,**数值相等**,违反 PRD「对外一律 2 位」展示契约,属 T-7 `_rebuild_quote` 范畴。集成测试已「比数值不比字符串」并钉住偏差 - [ ] **【待用户裁定】幂等重放响应的数值位数偏差**(T-9 执行期发现,**未顺手改**):首次响应 2 位(`calc` 量化)vs 重放响应 4 位(`core_trade` `DECIMAL(18,4)` 直读)→ `"53456.95"` vs `"53456.9500"`,**数值相等**,违反 PRD「对外一律 2 位」展示契约,属 T-7 `_rebuild_quote` 范畴。集成测试已「比数值不比字符串」并钉住偏差
> **第 0~2 批结果(2026-09-10)**:基线 **510 passed** → 批 0 后 **516 passed / 3 skipped**(+3 T-0 用例 +3 T-0b 引擎用例)→ 批 1(T-1)后**仍 516 passed / 3 skipped**(只加表与种子,未加用例 → **零回归**);**T-1 数据层断言 8/8 PASS**(含新增断言 ⑧:费率档 ↔ `product_type` 匹配,越档即 FAIL)→ 批 2(T-2 + T-2b)后 **609 passed / 3 skipped**(**+93 纯函数用例**,零回归);T-2b 实算脚本 15/15 与 PRD §5.3 一致(退出码 0)。 > **第 0~2 批结果(2026-09-10)**:基线 **510 passed** → 批 0 后 **516 passed / 3 skipped**(+3 T-0 用例 +3 T-0b 引擎用例)→ 批 1(T-1)后**仍 516 passed / 3 skipped**(只加表与种子,未加用例 → **零回归**);**T-1 数据层断言 8/8 PASS**(含新增断言 ⑧:费率档 ↔ `product_type` 匹配,越档即 FAIL)→ 批 2(T-2 + T-2b)后 **609 passed / 3 skipped**(**+93 纯函数用例**,零回归);T-2b 实算脚本 15/15 与 PRD §5.3 一致(退出码 0)。
> **下一步 = T-10**(普通申赎批次维护 · FR-C16 + D8 兜底补建 + `rebuild_lots.py`;**改 `trade_gateway` 主流程,回归风险最大 → 先跑基线再动**;DoD 见开发计划 §8)。**可并行**:T-11(`core_tools` / `sum_trades_on_date` 汇总去重 · 依赖 T-8 已解锁)。 > **下一步 = T-11**(`core_tools.query_recent_trades` 汇总走 `_amount_view` 去重 FR-C15 + 持仓查询过滤 `qty <= 0` + `core_ro.sum_trades_on_date` 加 convert 去重条件 **R-d**;DoD 见开发计划 §7.3)。基线 **714 passed / 3 skipped**。
> ⚠️ **基金转换的 T-0b 与下方「架构改进第 3/4 批」的 `core_ro` 只读账号是同一件事** —— 已由本线定案为 D20,**不再挂在架构改进线**(该线原「不要做」清单已更新)。 > ⚠️ **基金转换的 T-0b 与下方「架构改进第 3/4 批」的 `core_ro` 只读账号是同一件事** —— 已由本线定案为 D20,**不再挂在架构改进线**(该线原「不要做」清单已更新)。
@@ -1222,16 +1222,77 @@ sqlite 无 gap lock,故该分支由 `tests/test_convert_core.py` 用注入点
- 无批次时**补建再扣**,**不跳过、不阻断** - 无批次时**补建再扣**,**不跳过、不阻断**
**DoD** **DoD**
- [ ] 开工前:全量全绿基线已记录(当前 **516**) - [x] 开工前:全量全绿基线已记录(实测 **697 passed / 3 skipped**,非计划写作时的 516 —— 计划 §8 的基线数字是撰写期快照)
- [ ] 正常赎回后再 convert,**批次与持仓一致、不超扣**(验收 9) - [x] 正常赎回后再 convert,**批次与持仓一致、不超扣**(验收 9)
- [ ] **真实路径全覆盖**(R-c(2)):`test_share_lot.py` 自建种子 → `subscribe` 建批次(`qty = amount ÷ nav`,2 位 HALF_UP,`confirmed_at = T`)· `redeem` FIFO 扣减 · 无批次有持仓 → **D8 兜底补建后再扣** · `rebuild_lots.py` 入口 - [x] **真实路径全覆盖**(R-c(2)):`test_share_lot.py` 自建种子 → `subscribe` 建批次(`qty = amount ÷ nav`,2 位 HALF_UP,`confirmed_at = T`)· `redeem` FIFO 扣减 · 无批次有持仓 → **D8 兜底补建后再扣** · `rebuild_lots.py` 入口
- [ ] 降级路径单测(R-c(1)):`subscribe` 无净值 → warning + 不建批次 + 不抛异常;`redeem` 无持仓无批次 → warning + 跳过扣减 - [x] 降级路径单测(R-c(1)):`subscribe` 无净值 → warning + 不建批次 + 不抛异常;`redeem` 无持仓无批次 → warning + 跳过扣减
- [ ] D18 同源断言绿(两侧 `confirmed_at` 逐一相等) - [x] D18 同源断言绿(两侧 `confirmed_at` 逐一相等)
- [ ] `pytest -q` **510 + 新增全绿** - [x] `pytest -q` 全绿(**714 passed / 3 skipped** = 697 + 17)
- [ ] `test_trade_gateway.py` 既有断言按 §12 处置完毕(**R16 应零改动通过**) - [x] `test_trade_gateway.py` 既有断言按 §12 处置完毕(**R16 零改动通过**,实测)
**依赖**:T-3(排期置于 T-7 之后) **依赖**:T-3(排期置于 T-7 之后)
**执行记录(2026-09-10 · 已完成)**
*改码 4 处*
| 文件 | 动作 |
| --- | --- |
| `app/gateway/gateway_repository.py` | 【改】+模块级常量 `_SHARE_LOT_INSERT` / `_SHARE_LOT_DEDUCT`;+`insert_share_lots(lots)` / `deduct_share_lots(deductions)`(各单事务,走 `role="rw"`)。**不含任何分配逻辑** —— 只把已分配结果落库,避免出现第二份 FIFO 口径 |
| `app/repository/share_lot_repository.py` | 【改】`__init__` 增可选 `core_ro` 入参(传则复用该实例,忽略 `engine`)—— 让 `trade_gateway` 能复用**同一份** FIFO 选批口径(D18 / 自检第 13 问),而非另写一份贪心 |
| `app/gateway/trade_gateway.py` | 【改】+`SUBSCRIBE`/`REDEEM`/`SUBSCRIBE_LOT_PREFIX` 常量;+`_nav_as_of` / `_maintain_lots`;`submit_trade` 在 `insert_trade` 之后、引擎之前调用(整段 try 包住 → 绝不阻断交易) |
| `scripts/core/rebuild_lots.py` | 【新增】按 `core_holding` 快照重建批次(L-7)。**与网关 D8 兜底同调 `lot_bootstrap.bootstrap_lots`**(D18);默认「仅补建无批次的持仓行」(幂等),`--force` 先删后建,`--dry-run` 只报告;写库复用 `GatewayRepository.insert_share_lots`(INSERT 语句全仓只一份);走 `role="admin"`(`--force` 需 DELETE,业务账号刻意无) |
*实施级裁定 2 条(开发计划未点明,执行期定并留痕)*
1. **`redeem` 的份额折算 = `amount ÷ T 日净值`(2 位 HALF_UP)**。计划 §8 只写「FIFO 扣减」,
未说明扣减的**份额**从哪来 —— 而普通赎回请求体给的是**金额**、批次扣的是**份额**。
故必须折算,与 `subscribe` 对称;**取不到净值 → 降级跳过**(不猜、不阻断)。
批次只决定 FIFO 顺序与持有期,不参与定价。
> **⚠️ 2026-09-10 订正(联网查证 4 家管理人业务规则后)**:本条初稿写「按未知价法裁定」,
> **定性错误**。行业铁律是「**金额申购、份额赎回**」—— 投资者以**份额**申报,登记机构按
> T 日净值反算金额:`赎回金额 = 赎回份额 × T日净值 − 费用`。依据:睿远业务规则 §65
> 「以份额方式提出赎回申请」/ 华泰保兴 §69 / 东方基金 §57「按实际确认的有效赎回份额计算
> 赎回金额」/ 永赢招募书;货基同理(国投瑞银 §33)。**无一家公募支持「按金额赎回」**。
> 因此 `amount ÷ nav` 系**本项目侧简化建模**(`TradeRequest.amount` 为申赎共用入参,
> `app/api/simulate.py:60`),**不是行业做法** → 已记入 **PRD §10.1 已知差异第 ① 条**。
> **派生风险**:T 日净值当日不可得(T+1 公告),`_nav_as_of` 实取 **T−1 净值** →
> **折算份额只是估算值**,与最终确认份额有偏差(行业无此问题,因份额由投资者申报)。
> **处置**:折算逻辑**保留**(批次扣份额必须折算),但定性改为「项目简化」,
> `trade_gateway.py` 注释与本文档已同步订正。
2. **`rebuild_lots.py` 默认只补建「无批次」的持仓行,`--force` 才先删后建**。
计划只写「按 `core_holding` 重建批次(快照重建)」,未说是否覆盖已有批次。
默认保守(幂等、不误删转换转入的真实批次),需要全量对齐时才用 `--force`。
两者都调 `bootstrap_lots`,确定性 `lot_id` 天然防重复补建。
*测试与验证*
- `tests/test_share_lot.py` **+14 用例**(原 15 → 29):真实路径 5(申购建批次 / HALF_UP 舍入 /
赎回 FIFO 跨批 / 不足额只扣可用 / D8 兜底补建后再扣)+ 降级 3(无净值 / 无批次无持仓 / 持仓为 0)+
异常兜底 1 + `rebuild_lots` 3(默认补建 / `--force` / `--dry-run`)+ **D18 同源断言 1** +
**机制断言 1**(`rebuild.bootstrap_lots is bootstrap_lots`,防「两侧碰巧算出同值」的假绿)
- `tests/test_trade_gateway.py` **+3 用例**(原 29 → 32):申购建批次端到端 / 赎回扣减端到端 /
**阻断路径不写批次**
- `pytest -q` → **714 passed / 3 skipped**(基线 697 **+17**,**零回归**;§12 **R16 零改动通过**)
- **新增 `scripts/dev/verify_convert_lots.py` 真库验证 → 20/20 一致(退出码 0)**:
A 申购建批次(含 nav 4 位精度、`confirmed_at = T`)· B/C **赎回 FIFO 扣减真库生效 + 归零行保留 + 不超扣**
· D D8 兜底补建 · E `rebuild_lots` 真库跑通(幂等 + 补建 + **D18 同源**);
隔离前缀 `CUST-LOTT`/`PROD-LOTT*`,跑完**零残留**(已复核 5 张表计数全 0)。
**★ 本脚本的核心价值**:确认 `(:q + 0.0)` 这一**为绕开 sqlite 坑而加的写法在 MySQL 同样生效** ——
若 MySQL 上字符串绑定未被隐式转数值,扣减会静默失效(批次永不减少 → convert 超扣),**sqlite 侧看不出来**。
- **突变验证(2 组,防假绿)**:① 切断 `submit_trade` 里的 `_maintain_lots` 接线 →
**精准 2 条**(两条端到端用例)红;② 关掉 `redeem` 的 D8 兜底落库 → **精准 2 条**
(`test_redeem_bootstraps_lot_from_holding_then_deducts` + `test_d18_gateway_bootstrap_and_rebuild_lots_agree`)红。
两处均已恢复,`grep MUTATION-TEST` 无残留。
- **回归复跑**:T-6 `verify_convert_apply.py` **24/24** · T-7 `verify_convert_service.py` **35/35**(零回归)
*未顺手改(不在本任务范围)*
- **`core_holding` 不随普通申赎更新** —— T-10 只维护 `core_share_lot`(计划范围即如此)。
两者可能短暂失配,正是 `rebuild_lots.py` 的定位(快照重建)。是否让普通申赎也同步持仓,
属新范围,留待用户决定。
--- ---
## 9. 第 6 批 · 补偿脚本(T-12) ## 9. 第 6 批 · 补偿脚本(T-12)
+170
View File
@@ -0,0 +1,170 @@
#!/usr/bin/env python3
"""按 `core_holding` 快照重建份额批次(D18 · L-7 · 架构 §2 scripts/core)。
**定位**:`core_share_lot` 是 FIFO 计费的权威源,`core_holding` 是持仓快照,
两者靠交易统一出入口(`trade_gateway` 的 FR-C16 批次维护)保持一致。
当历史数据缺批次(批次机制上线前的存量、或演示库被人工改动)时,用本脚本
**按持仓快照重建**——这是**快照重建、不是交易回滚**(PRD L-7):不重放流水,
只让 `Σ core_share_lot.remain_qty` 与 `core_holding.qty` 重新对齐。
**⚠️ D18 单一副本(改本文件前必读)**
补建规则**只有一处** —— `service/convert/lot_bootstrap.bootstrap_lots`。
本脚本与 `trade_gateway` 的 D8 兜底补建**同调该函数**。任何情况下都不要在
本文件里另写一份"由 `as_of` 反推 `confirmed_at`"的逻辑:两处各写一遍必然
漂移,同一持仓在两侧算出不同费率档 —— **演示里看不出来、生产里就是错账**。
**权限**:走 `role="admin"`(= `mysql_user`)。`--force` 需要 DELETE 权限,
而业务账号 `xh_core_rw` 刻意不持有(D20 最小权限),故本运维脚本不用业务账号。
用法:
python scripts/core/rebuild_lots.py # 仅补建「无批次」的持仓行(安全、幂等)
python scripts/core/rebuild_lots.py --customer CUST-4001
python scripts/core/rebuild_lots.py --force # 先删后建:全量快照重建
python scripts/core/rebuild_lots.py --dry-run # 只报告、不写库
退出码:0 = 正常;1 = 库不可达。
"""
from __future__ import annotations
import argparse
import sys
from pathlib import Path
from typing import Any
from sqlalchemy import text
from sqlalchemy.engine import Engine
ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(ROOT))
from app.config.settings import settings # noqa: E402
from app.gateway.gateway_repository import GatewayRepository # noqa: E402
from app.service.convert.lot_bootstrap import bootstrap_lots # noqa: E402
from app.utils.db import get_engine # noqa: E402
_SELECT_HOLDINGS = "SELECT customer_id, product_id, qty, cost_amount, as_of FROM core_holding"
#: 一次取回全部已有批次的 `(客户, 产品)` 分布,避免逐行回查。
_GROUP_EXISTING_LOTS = """
SELECT customer_id, product_id, COUNT(*) AS n
FROM core_share_lot
GROUP BY customer_id, product_id
"""
_DELETE_LOTS = "DELETE FROM core_share_lot WHERE customer_id = :cid AND product_id = :pid"
def rebuild_lots(
engine: Engine,
*,
customer_id: str | None = None,
force: bool = False,
dry_run: bool = False,
) -> dict[str, int]:
"""按 `core_holding` 重建批次(CLI 与单测共用入口)。
- 缺省:**只补建「无批次」的持仓行**(安全、幂等;重复执行零写入)。
- `--force`:**先删后建**该 `(客户, 产品)` 的全部批次(全量快照重建)。
- `qty <= 0` 的持仓行**不补建**(D9/P2 归零行保留但不建批次,
与 `bootstrap_lots` 的语义一致)。
返回统计,键含义:`holdings` 扫描的持仓行数 · `bootstrapped` 实际补建的
持仓行数 · `skipped_existing` 因已有批次而跳过 · `skipped_zero` 因份额为 0
跳过 · `removed` 删除的批次行数(仅 `--force`)· `written` 写入的批次行数。
写入复用 `GatewayRepository.insert_share_lots` —— `core_share_lot` 的 INSERT
语句全仓只此一份,本脚本不另写一条(自检第 13 问)。
"""
stats = {
"holdings": 0,
"bootstrapped": 0,
"skipped_existing": 0,
"skipped_zero": 0,
"removed": 0,
"written": 0,
}
sql = _SELECT_HOLDINGS
params: dict[str, Any] = {}
if customer_id:
sql += " WHERE customer_id = :cid"
params["cid"] = customer_id
sql += " ORDER BY customer_id, product_id"
plans: list[tuple[str, str, list[Any]]] = []
with engine.connect() as conn:
holdings = [dict(r) for r in conn.execute(text(sql), params).mappings()]
existing = {
(r["customer_id"], r["product_id"]): int(r["n"])
for r in conn.execute(text(_GROUP_EXISTING_LOTS)).mappings()
}
for row in holdings:
stats["holdings"] += 1
cid, pid = str(row["customer_id"]), str(row["product_id"])
lots = bootstrap_lots(row)
if not lots:
stats["skipped_zero"] += 1
continue
has = existing.get((cid, pid), 0)
if has and not force:
stats["skipped_existing"] += 1
continue
if has:
stats["removed"] += has
stats["bootstrapped"] += 1
plans.append((cid, pid, lots))
if dry_run:
return stats
# `--force`:先删(独立事务)。放在写入之前,且不嵌套写事务 ——
# sqlite 测试库是 StaticPool 单连接,事务内再开事务会自锁。
if force:
with engine.begin() as conn:
for cid, pid, _ in plans:
conn.execute(text(_DELETE_LOTS), {"cid": cid, "pid": pid})
writer = GatewayRepository(engine=engine)
for _, _, lots in plans:
writer.insert_share_lots(lots)
stats["written"] += len(lots)
return stats
def main() -> None:
parser = argparse.ArgumentParser(
description="按 core_holding 快照重建份额批次(补建规则与网关同源,D18)"
)
parser.add_argument("--customer", help="只处理该客户(缺省全部)")
parser.add_argument(
"--force",
action="store_true",
help="先删后建(全量快照重建);缺省仅补建无批次的持仓行",
)
parser.add_argument("--dry-run", action="store_true", help="只报告、不写库")
args = parser.parse_args()
engine = get_engine(settings.mysql_core_database, "admin")
try:
with engine.connect() as conn:
conn.execute(text("SELECT 1"))
except Exception as exc: # noqa: BLE001 - CLI 兜底,给出可操作的提示
print(f"MySQL 不可达:{exc}", file=sys.stderr)
print("请确认本机 MySQL 服务已启动(见 FLOW §0 本机状态)。", file=sys.stderr)
raise SystemExit(1)
stats = rebuild_lots(
engine, customer_id=args.customer, force=args.force, dry_run=args.dry_run
)
prefix = "[dry-run] " if args.dry_run else ""
print(
f"{prefix}持仓行 {stats['holdings']} · 补建 {stats['bootstrapped']} · "
f"跳过(已有批次) {stats['skipped_existing']} · 跳过(份额为0) {stats['skipped_zero']} · "
f"删除批次 {stats['removed']} · 写入批次 {stats['written']}"
)
if __name__ == "__main__":
main()
+339
View File
@@ -0,0 +1,339 @@
"""T-10 真 MySQL 验证脚本:普通申赎批次维护 + `rebuild_lots.py` 实跑 + DoD 断言。
**为什么 sqlite 单测全绿还不够**
T-10 的扣减 SQL 写的是 `remain_qty = remain_qty - (:q + 0.0)`,其中
`+ 0.0` 是**为绕开 sqlite 坑**而加的(sqlite 在 UPDATE 的算术表达式里不把
TEXT 绑定参数转数值 → 不加就扣 0;见开发计划 §B.8 第 11 条)。
**该写法在 MySQL 是否同样成立,sqlite 侧证明不了** —— 真库若因字符串绑定
没被隐式转换成数值,扣减会静默失效(批次永远扣不掉 → convert 超扣)。
此外真库还能验证 sqlite 测不到的:DECIMAL(18,4) 存取精度、InnoDB 真实提交语义,
以及 `rebuild_lots.py` 的 DELETE/INSERT 在真库是否跑通。
用法:
python scripts/dev/verify_convert_lots.py # 建隔离数据 → 跑断言 → 清理
注意:
- 全部数据用 **LOTT 前缀**(客户 `CUST-LOTT` / 产品 `PROD-LOTT*`),
跑完 DELETE 干净,不碰既有种子;
- 建/清数据走 `role="admin"`(需 DELETE,R-e);**被测路径走生产同款仓储**
(`GatewayRepository` rw 写 + `CoreReadOnlyRepository` ro 读);
- 退出码 1 = 有断言不一致(供 CI / 人工判定)。
"""
from __future__ import annotations
import sys
from datetime import date, datetime
from decimal import ROUND_HALF_UP, Decimal
from pathlib import Path
from sqlalchemy import text
ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(ROOT))
from app.config.settings import settings # noqa: E402
from app.gateway.gateway_repository import GatewayRepository # noqa: E402
from app.gateway.trade_gateway import _maintain_lots # noqa: E402
from app.repository.core_ro import CoreReadOnlyRepository # noqa: E402
from app.service.convert.lot_bootstrap import bootstrap_lots # noqa: E402
from app.utils.db import dispose_engines, get_engine # noqa: E402
CUSTOMER = "CUST-LOTT"
PROD_SUB = "PROD-LOTTA" # A 组:申购建批次
PROD_BOOT = "PROD-LOTTB" # D 组:无批次有持仓 → D8 兜底补建
PROD_REDEEM = "PROD-LOTTC" # B/C 组:赎回 FIFO 扣减
NAV_DATE = date(2026, 9, 4)
TRADE_AT = datetime(2026, 9, 4, 14, 0, 0)
_passed = 0
_failed = 0
def check(name: str, actual, expected) -> None:
"""逐条断言并打印(与既有 verify_convert_*.py 同款输出)。"""
global _passed, _failed
ok = actual == expected
if ok:
_passed += 1
else:
_failed += 1
flag = "✅" if ok else "❌"
print(f" {flag} {name}: 实际 {actual!r}" + ("" if ok else f" / 期望 {expected!r}"))
def dec(value, places: str = "0.01") -> Decimal:
"""真库读回是 `Decimal`(MySQL DECIMAL)→ 统一量化后比较。"""
return Decimal(str(value)).quantize(Decimal(places), rounding=ROUND_HALF_UP)
def q1(engine, sql: str, **params):
with engine.connect() as conn:
return conn.execute(text(sql), params).scalar_one_or_none()
# ── 数据准备 / 清理 ─────────────────────────────────────────────────
def seed(engine) -> None:
with engine.begin() as conn:
conn.execute(
text(
"INSERT INTO core_customer (customer_id, display_name, open_date) "
"VALUES (:c, 'T10真库验证', :d)"
),
{"c": CUSTOMER, "d": NAV_DATE},
)
for pid, name in [
(PROD_SUB, "T10申购基金"),
(PROD_BOOT, "T10兜底基金"),
(PROD_REDEEM, "T10赎回基金"),
]:
conn.execute(
text(
"INSERT INTO core_product (product_id, product_name, min_risk_code, "
"product_type, can_subscribe, can_redeem) "
"VALUES (:p, :n, 'R2', 'bond', 1, 1)"
),
{"p": pid, "n": name},
)
# 三个产品的 T 日净值:A 用 1.2 验折算,B/C 用 1.0 让份额=金额、断言直观
for pid, nav in [(PROD_SUB, "1.2000"), (PROD_BOOT, "1.0000"), (PROD_REDEEM, "1.0000")]:
conn.execute(
text(
"INSERT INTO core_product_nav (product_id, nav, daily_chg_pct, nav_date) "
"VALUES (:p, :nav, 0, :d)"
),
{"p": pid, "nav": nav, "d": NAV_DATE},
)
# C 组:两个已存在批次(FIFO 扣减用)
for lot_id, qty, nav, confirmed in [
("LOT-LOTT-C1", "100", "1.0000", datetime(2026, 9, 1, 10, 0, 0)),
("LOT-LOTT-C2", "50", "1.1000", datetime(2026, 9, 2, 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_REDEEM, "q": qty, "nav": nav,
"cat": confirmed},
)
# D 组:只有持仓、没有批次(D8 兜底补建)
# A 组:持仓与申购批次同额(供 rebuild_lots 重建验证)
for pid, qty, cost in [(PROD_BOOT, "100", "100.00"), (PROD_SUB, "10000", "12000.00")]:
conn.execute(
text(
"INSERT INTO core_holding (customer_id, product_id, qty, cost_amount, "
"market_value, pnl_pct, as_of) VALUES (:c, :p, :q, :cost, :cost, 0, :d)"
),
{"c": CUSTOMER, "p": pid, "q": qty, "cost": cost, "d": date(2026, 8, 1)},
)
def cleanup(engine) -> None:
"""倒序删(子表 → 父表),只删 LOTT 前缀数据。"""
with engine.begin() as conn:
conn.execute(text("DELETE FROM core_share_lot WHERE customer_id = :c"), {"c": CUSTOMER})
conn.execute(text("DELETE FROM core_holding WHERE customer_id = :c"), {"c": CUSTOMER})
conn.execute(text("DELETE FROM core_product_nav WHERE product_id LIKE 'PROD-LOTT%'"))
conn.execute(text("DELETE FROM core_product WHERE product_id LIKE 'PROD-LOTT%'"))
conn.execute(text("DELETE FROM core_customer WHERE customer_id = :c"), {"c": CUSTOMER})
# ── 断言组 ──────────────────────────────────────────────────────────
def assert_subscribe(admin, writer, core) -> None:
print("【A】申购建批次(FR-C16 · qty = amount ÷ T 日净值)")
_maintain_lots(
writer=writer, core=core, trade_id="TRD-LOTT-SUB",
customer_id=CUSTOMER, product_id=PROD_SUB,
trade_type="subscribe", amount=Decimal("12000"), traded_at=TRADE_AT,
)
check(
"批次行数",
q1(admin, "SELECT COUNT(*) FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB),
1,
)
check(
"lot_id",
q1(admin, "SELECT lot_id FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB),
"LOT-SUB-TRD-LOTT-SUB",
)
check(
"qty = 12000 ÷ 1.2000",
dec(q1(admin, "SELECT qty FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB)),
dec("10000.00"),
)
check(
"remain_qty = qty",
dec(q1(admin, "SELECT remain_qty FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB)),
dec("10000.00"),
)
check(
"nav 4 位精度存取",
dec(q1(admin, "SELECT nav FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB), "0.0001"),
Decimal("1.2000"),
)
check(
"confirmed_at = T",
str(q1(admin, "SELECT confirmed_at FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB)),
"2026-09-04 14:00:00",
)
def assert_redeem_fifo(admin, writer, core) -> None:
print("【B/C】赎回 FIFO 扣减(★ 真库验证 `(:q + 0.0)` 方言语义)")
# 120 元 ÷ 1.0 = 120 份 → C1 扣满 100 归零,C2 再扣 20
_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,
)
check(
"C1(最老批)remain_qty 归零 —— 扣减真库生效",
dec(q1(admin, "SELECT remain_qty FROM core_share_lot WHERE lot_id = 'LOT-LOTT-C1'")),
dec("0.00"),
)
check(
"C2 remain_qty = 50 − 20",
dec(q1(admin, "SELECT remain_qty FROM core_share_lot WHERE lot_id = 'LOT-LOTT-C2'")),
dec("30.00"),
)
check(
"归零批次**保留行不删**(D9/P2)",
q1(admin, "SELECT COUNT(*) FROM core_share_lot WHERE lot_id = 'LOT-LOTT-C1'"),
1,
)
check(
"不超扣:Σ remain = 150 − 120",
dec(q1(admin, "SELECT SUM(remain_qty) FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_REDEEM)),
dec("30.00"),
)
def assert_bootstrap(admin, writer, core) -> None:
print("【D】D8 兜底补建:无批次有持仓 → 按 as_of 补建后再扣")
_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,
)
check(
"补建批次行数",
q1(admin, "SELECT COUNT(*) FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_BOOT),
1,
)
check(
"lot_id = 确定性兜底 id",
q1(admin, "SELECT lot_id FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_BOOT),
"LOT-BOOT-CUST-LOTT-PROD-LOTTB",
)
check(
"qty = 持仓快照 100",
dec(q1(admin, "SELECT qty FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_BOOT)),
dec("100.00"),
)
check(
"补建 100 后扣 30 → remain 70",
dec(q1(admin, "SELECT remain_qty FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_BOOT)),
dec("70.00"),
)
def assert_rebuild_lots(admin) -> None:
print("【E】rebuild_lots 真库跑通(D18 同源 + 幂等)")
import importlib.util
spec = importlib.util.spec_from_file_location(
"rebuild_lots_under_test", ROOT / "scripts" / "core" / "rebuild_lots.py"
)
rebuild = importlib.util.module_from_spec(spec)
spec.loader.exec_module(rebuild)
# 幂等:此刻三个产品都已有批次 → 零写入
stats = rebuild.rebuild_lots(admin, customer_id=CUSTOMER)
check("默认模式:已有批次全跳过", stats["written"], 0)
check("默认模式:补建 0 行", stats["bootstrapped"], 0)
# 人为删掉 A 组批次 → 模拟「批次缺失」,验证真库补建
with admin.begin() as conn:
conn.execute(
text("DELETE FROM core_share_lot WHERE customer_id = :c AND product_id = :p"),
{"c": CUSTOMER, "p": PROD_SUB},
)
stats = rebuild.rebuild_lots(admin, customer_id=CUSTOMER)
check("补建 1 行", stats["written"], 1)
rebuilt_lot_id = q1(
admin, "SELECT lot_id FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB,
)
check("补建 lot_id", rebuilt_lot_id, "LOT-BOOT-CUST-LOTT-PROD-LOTTA")
check(
"补建份额 = 持仓快照",
dec(q1(admin, "SELECT remain_qty FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB)),
dec("10000.00"),
)
# ★ D18 同源:真库补出的 confirmed_at 必须与 bootstrap_lots 直算一致
holding = {
"customer_id": CUSTOMER,
"product_id": PROD_SUB,
"qty": q1(admin, "SELECT qty FROM core_holding WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB),
"cost_amount": q1(admin, "SELECT cost_amount FROM core_holding WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB),
"as_of": q1(admin, "SELECT as_of FROM core_holding WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB),
}
expected_confirmed = bootstrap_lots(holding)[0].confirmed_at
check(
"D18 同源:真库补建 confirmed_at == bootstrap_lots 直算",
str(q1(admin, "SELECT confirmed_at FROM core_share_lot WHERE customer_id = :c AND product_id = :p",
c=CUSTOMER, p=PROD_SUB)),
str(expected_confirmed),
)
def main() -> int:
admin = get_engine(settings.mysql_core_database, "admin")
try:
with admin.connect() as conn:
conn.execute(text("SELECT 1"))
except Exception as exc: # noqa: BLE001 - CLI 兜底
print(f"MySQL 不可达:{exc}", file=sys.stderr)
print("请确认本机 MySQL 服务已启动(见 FLOW §0 本机状态)。", file=sys.stderr)
return 1
cleanup(admin)
try:
seed(admin)
# 被测路径走生产同款仓储(rw 写 / ro 读)
writer = GatewayRepository()
core = CoreReadOnlyRepository()
assert_subscribe(admin, writer, core)
assert_redeem_fifo(admin, writer, core)
assert_bootstrap(admin, writer, core)
assert_rebuild_lots(admin)
finally:
cleanup(admin)
dispose_engines()
print(f"\n{'=' * 60}")
print(f"真库验证:{_passed} 项一致 / {_failed} 项不一致")
return 1 if _failed else 0
if __name__ == "__main__":
sys.exit(main())
+386 -1
View File
@@ -1,4 +1,4 @@
"""T-3 `core_ro` 五个新方法 + `share_lot_repository` 单测(开发计划 §5.1 DoD)。 """T-3 读侧单测 + T-10 普通申赎批次维护(开发计划 §5.1 / §8 DoD)。
sqlite 内存库(conftest.sqlite_engine 单一事实源建表 + C×R 矩阵种子), sqlite 内存库(conftest.sqlite_engine 单一事实源建表 + C×R 矩阵种子),
数据由本文件自建,不依赖真 MySQL。 数据由本文件自建,不依赖真 MySQL。
@@ -10,18 +10,32 @@ sqlite 内存库(conftest.sqlite_engine 单一事实源建表 + C×R 矩阵种
4. sum_remain_qty(仅统计 remain_qty > 0,归零批不计入) 4. sum_remain_qty(仅统计 remain_qty > 0,归零批不计入)
5. get_holding(单行 / 不存在 None) 5. get_holding(单行 / 不存在 None)
6. ShareLotRepository.select_for_convert(FIFO 贪心选批、不足额返回全部可用) 6. ShareLotRepository.select_for_convert(FIFO 贪心选批、不足额返回全部可用)
7. **T-10 普通申赎批次维护(FR-C16)**:申购建批次 / 赎回 FIFO 扣减 /
无批次有持仓 → D8 兜底补建后再扣 / 降级两条(无净值、无批次无持仓)/
`rebuild_lots.py` 入口 / **D18 同源断言**(两侧 confirmed_at 逐一相等)
R-c(2) 覆盖率补偿:第 7 组的真实路径**全部自建完整种子**
(`core_holding` + `core_share_lot` + `core_product_nav`),
**不得让降级路径充当测试覆盖**(降级是数据不全时的兜底,不是被测对象)。
""" """
from __future__ import annotations from __future__ import annotations
import importlib.util
import logging
from datetime import date, datetime from datetime import date, datetime
from decimal import Decimal from decimal import Decimal
from pathlib import Path
import pytest import pytest
from sqlalchemy import text from sqlalchemy import text
from _ddl import create_sqlite_engine
from app.gateway import trade_gateway as tg
from app.gateway.gateway_repository import GatewayRepository
from app.repository.core_ro import CoreReadOnlyRepository, _as_date from app.repository.core_ro import CoreReadOnlyRepository, _as_date
from app.repository.share_lot_repository import ShareLotRepository from app.repository.share_lot_repository import ShareLotRepository
from app.service.convert.lot_bootstrap import bootstrap_lot_id, bootstrap_lots
def _as_dt(value: object) -> datetime: def _as_dt(value: object) -> datetime:
@@ -315,3 +329,374 @@ def test_select_for_convert_shortfall_returns_all_available(sqlite_engine):
def test_select_for_convert_zero_qty_returns_empty(sqlite_engine): def test_select_for_convert_zero_qty_returns_empty(sqlite_engine):
repo = ShareLotRepository(engine=sqlite_engine) repo = ShareLotRepository(engine=sqlite_engine)
assert repo.select_for_convert(CUSTOMER, PRODUCT_A, Decimal("0")) == [] assert repo.select_for_convert(CUSTOMER, PRODUCT_A, Decimal("0")) == []
# ═══════════════════════════════════════════════════════════════════════
# 7. T-10 普通申赎批次维护(FR-C16 · 含 D8 兜底补建)
# ═══════════════════════════════════════════════════════════════════════
#: 交易日(T 日):批次 confirmed_at 与净值取数基准
NOW_T = datetime(2026, 9, 4, 14, 0, 0)
NAV_DATE = date(2026, 9, 4)
def _dec(value: object) -> Decimal:
"""sqlite 读回 DECIMAL 列是 float —— 统一转 Decimal 再比,避免二进制误差误判。"""
return Decimal(str(value))
def _rows(engine, sql: str, params: dict | None = None) -> list[dict]:
with engine.connect() as conn:
return [dict(r) for r in conn.execute(text(sql), params or {}).mappings()]
def _seed_nav(engine, pid: str, nav: str, nav_date: date = NAV_DATE) -> None:
with engine.begin() as conn:
conn.execute(
text(
"INSERT INTO core_product_nav (product_id, nav, daily_chg_pct, nav_date) "
"VALUES (:pid, :nav, 0, :nd)"
),
{"pid": pid, "nav": float(nav), "nd": nav_date},
)
def _seed_holding(
engine, cid: str, pid: str, qty: str, cost: str, as_of: date
) -> None:
with engine.begin() as conn:
conn.execute(
text(
"INSERT INTO core_holding (customer_id, product_id, qty, cost_amount, "
"market_value, pnl_pct, as_of) "
"VALUES (:cid, :pid, :qty, :cost, :cost, 0, :as_of)"
),
{"cid": cid, "pid": pid, "qty": float(qty), "cost": float(cost), "as_of": as_of},
)
def _maintain(
engine,
*,
trade_id: str,
trade_type: str,
amount: str,
product_id: str = PRODUCT_A,
customer_id: str = CUSTOMER,
traded_at: datetime = NOW_T,
) -> None:
"""直调批次维护入口(单元级:可控、可断言语义,端到端接线见 test_trade_gateway.py)。"""
tg._maintain_lots(
writer=GatewayRepository(engine=engine),
core=CoreReadOnlyRepository(engine=engine),
trade_id=trade_id,
customer_id=customer_id,
product_id=product_id,
trade_type=trade_type,
amount=Decimal(amount),
traded_at=traded_at,
)
# ── 7.1 真实路径:申购建批次 ───────────────────────────────────────────
def test_subscribe_creates_lot_with_amount_over_nav(sqlite_engine):
"""申购:`qty = amount ÷ T 日净值`(2 位 HALF_UP),`confirmed_at = T`。"""
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A)
_seed_nav(sqlite_engine, PRODUCT_A, "1.2500")
_maintain(
sqlite_engine, trade_id="TRD-T10-SUB", trade_type="subscribe", amount="10000"
)
rows = _rows(sqlite_engine, "SELECT * FROM core_share_lot")
assert len(rows) == 1
row = rows[0]
assert row["lot_id"] == "LOT-SUB-TRD-T10-SUB"
assert row["customer_id"] == CUSTOMER and row["product_id"] == PRODUCT_A
assert _dec(row["qty"]) == Decimal("8000.00") # 10000 / 1.25
assert _dec(row["remain_qty"]) == Decimal("8000.00")
assert _dec(row["nav"]) == Decimal("1.25")
assert _as_dt(row["confirmed_at"]) == NOW_T
assert row["source_trade_id"] == "TRD-T10-SUB"
def test_subscribe_rounds_qty_half_up(sqlite_engine):
"""份额 2 位 HALF_UP:1000 ÷ 1.2 = 833.333… → 833.33(不是银行家舍入)。"""
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A)
_seed_nav(sqlite_engine, PRODUCT_A, "1.2000")
_maintain(
sqlite_engine, trade_id="TRD-T10-RND", trade_type="subscribe", amount="1000"
)
row = _rows(sqlite_engine, "SELECT * FROM core_share_lot")[0]
assert _dec(row["remain_qty"]) == Decimal("833.33")
# ── 7.2 真实路径:赎回 FIFO 扣减 ───────────────────────────────────────
def test_redeem_deducts_fifo_across_lots(sqlite_engine):
"""赎回:按 `amount ÷ T 日净值` 折算份额后 FIFO 扣减,最老批次先扣。"""
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A)
_seed_nav(sqlite_engine, PRODUCT_A, "1.0000")
_seed_lot(sqlite_engine, "L-OLD", CUSTOMER, PRODUCT_A, "100", "1.0000",
datetime(2026, 9, 1, 10, 0, 0))
_seed_lot(sqlite_engine, "L-NEW", CUSTOMER, PRODUCT_A, "50", "1.1000",
datetime(2026, 9, 2, 10, 0, 0))
# 120 元 ÷ 1.0 = 120 份 → L-OLD 扣满 100 归零,L-NEW 再扣 20
_maintain(
sqlite_engine, trade_id="TRD-T10-RED", trade_type="redeem", amount="120"
)
rows = {
r["lot_id"]: _dec(r["remain_qty"])
for r in _rows(sqlite_engine, "SELECT * FROM core_share_lot")
}
assert rows == {"L-OLD": Decimal("0.00"), "L-NEW": Decimal("30.00")}
# 归零批次**保留行不删**(D9/P2)
assert len(rows) == 2
# 不超扣:Σ remain = 初始 150 − 实际扣减 120
assert sum(rows.values()) == Decimal("30.00")
def test_redeem_shortfall_deducts_available_only(sqlite_engine):
"""普通赎回不阻断:请求份额 > 可用份额时,按可用额度全部扣减、不报错。"""
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A)
_seed_nav(sqlite_engine, PRODUCT_A, "1.0000")
_seed_lot(sqlite_engine, "L1", CUSTOMER, PRODUCT_A, "30", "1.0000",
datetime(2026, 9, 1, 10, 0, 0))
_maintain(
sqlite_engine, trade_id="TRD-T10-SHORT", trade_type="redeem", amount="5000"
)
row = _rows(sqlite_engine, "SELECT * FROM core_share_lot")[0]
assert _dec(row["remain_qty"]) == Decimal("0.00")
# ── 7.3 真实路径:无批次有持仓 → D8 兜底补建后再扣 ─────────────────────
def test_redeem_bootstraps_lot_from_holding_then_deducts(sqlite_engine):
"""D8 主场景:无批次但有持仓 → 按 `core_holding.as_of` 补建初始批次再扣。"""
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A)
_seed_nav(sqlite_engine, PRODUCT_A, "1.0000")
_seed_holding(sqlite_engine, CUSTOMER, PRODUCT_A, "100", "100.00", date(2026, 8, 1))
_maintain(
sqlite_engine, trade_id="TRD-T10-D8", trade_type="redeem", amount="30"
)
rows = _rows(sqlite_engine, "SELECT * FROM core_share_lot")
assert len(rows) == 1
row = rows[0]
assert row["lot_id"] == bootstrap_lot_id(CUSTOMER, PRODUCT_A) # 确定性 id,防重复补建
assert _dec(row["qty"]) == Decimal("100") # 原份额 = 持仓快照
assert _dec(row["remain_qty"]) == Decimal("70.00") # 补建 100 后扣 30
assert row["source_trade_id"] is None # 兜底补建无来源流水(与 08 种子一致)
def test_redeem_zero_holding_skips_bootstrap(sqlite_engine, caplog):
"""持仓份额为 0(D9 归零行保留)→ 不补建、不抛异常。"""
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A)
_seed_nav(sqlite_engine, PRODUCT_A, "1.0000")
_seed_holding(sqlite_engine, CUSTOMER, PRODUCT_A, "0", "0.00", date(2026, 8, 1))
with caplog.at_level(logging.WARNING, logger="app.gateway.trade_gateway"):
_maintain(
sqlite_engine, trade_id="TRD-T10-ZERO", trade_type="redeem", amount="30"
)
assert _rows(sqlite_engine, "SELECT * FROM core_share_lot") == []
assert "无可补建批次" in caplog.text
# ── 7.4 降级路径(R-c(1) · 保住既有 510 用例的关键)────────────────────
def test_subscribe_without_nav_skips_with_warning(sqlite_engine, caplog):
"""申购取不到 T 日净值 → warning + 不建批次 + 不抛异常。"""
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A) # 刻意不灌 core_product_nav
with caplog.at_level(logging.WARNING, logger="app.gateway.trade_gateway"):
_maintain(
sqlite_engine, trade_id="TRD-T10-NONAV", trade_type="subscribe",
amount="10000",
)
assert _rows(sqlite_engine, "SELECT * FROM core_share_lot") == []
assert "申购取不到 T 日净值" in caplog.text
def test_redeem_without_lot_and_holding_skips_with_warning(sqlite_engine, caplog):
"""赎回既无批次也无持仓 → warning + 跳过扣减 + 不抛异常(R16 的兜底)。"""
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A)
_seed_nav(sqlite_engine, PRODUCT_A, "1.0000") # 有净值,但仍无批次无持仓
with caplog.at_level(logging.WARNING, logger="app.gateway.trade_gateway"):
_maintain(
sqlite_engine, trade_id="TRD-T10-EMPTY", trade_type="redeem", amount="1000"
)
assert _rows(sqlite_engine, "SELECT * FROM core_share_lot") == []
assert "既无批次也无持仓" in caplog.text
def test_maintain_lots_never_raises_on_broken_writer(sqlite_engine, caplog):
"""兜底:维护过程抛任何异常都被吞掉(交易主流程不受影响)。"""
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A)
_seed_nav(sqlite_engine, PRODUCT_A, "1.0000")
class BoomGateway(GatewayRepository):
def insert_share_lots(self, lots): # noqa: D102
raise RuntimeError("boom")
with caplog.at_level(logging.WARNING, logger="app.gateway.trade_gateway"):
tg._maintain_lots(
writer=BoomGateway(engine=sqlite_engine),
core=CoreReadOnlyRepository(engine=sqlite_engine),
trade_id="TRD-T10-BOOM",
customer_id=CUSTOMER,
product_id=PRODUCT_A,
trade_type="subscribe",
amount=Decimal("10000"),
traded_at=NOW_T,
)
assert "批次维护失败" in caplog.text
# ═══════════════════════════════════════════════════════════════════════
# 8. `scripts/core/rebuild_lots.py`(D18 同源 + 幂等)
# ═══════════════════════════════════════════════════════════════════════
_REBUILD_PATH = Path(__file__).resolve().parents[1] / "scripts" / "core" / "rebuild_lots.py"
def _load_rebuild_lots():
"""按路径加载脚本模块(`scripts/` 非包,无法直接 import)。"""
spec = importlib.util.spec_from_file_location("rebuild_lots_under_test", _REBUILD_PATH)
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
return module
def test_rebuild_lots_bootstraps_missing_only(sqlite_engine):
"""缺省模式:只补建「无批次」的持仓行;已有批次的持仓不动。"""
rebuild = _load_rebuild_lots()
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A)
_seed_product(sqlite_engine, PRODUCT_B)
_seed_holding(sqlite_engine, CUSTOMER, PRODUCT_A, "100", "100.00", date(2026, 8, 1))
_seed_holding(sqlite_engine, CUSTOMER, PRODUCT_B, "20", "22.00", date(2026, 8, 2))
# PRODUCT_B 已有批次 → 应被跳过
_seed_lot(sqlite_engine, "L-EXIST", CUSTOMER, PRODUCT_B, "20", "1.1000",
datetime(2026, 8, 2, 10, 0, 0))
stats = rebuild.rebuild_lots(sqlite_engine, customer_id=CUSTOMER)
assert stats["holdings"] == 2
assert stats["bootstrapped"] == 1
assert stats["skipped_existing"] == 1
assert stats["written"] == 1
lot_ids = {r["lot_id"] for r in _rows(sqlite_engine, "SELECT * FROM core_share_lot")}
assert lot_ids == {"L-EXIST", bootstrap_lot_id(CUSTOMER, PRODUCT_A)}
# 幂等:再跑一次零写入
again = rebuild.rebuild_lots(sqlite_engine, customer_id=CUSTOMER)
assert again["written"] == 0 and again["bootstrapped"] == 0
def test_rebuild_lots_force_replaces_existing(sqlite_engine):
"""`--force`:先删后建,批次按持仓快照重建(L-7)。"""
rebuild = _load_rebuild_lots()
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A)
_seed_holding(sqlite_engine, CUSTOMER, PRODUCT_A, "100", "100.00", date(2026, 8, 1))
_seed_lot(sqlite_engine, "L-STALE", CUSTOMER, PRODUCT_A, "77", "9.9999",
datetime(2026, 1, 1, 10, 0, 0))
stats = rebuild.rebuild_lots(sqlite_engine, customer_id=CUSTOMER, force=True)
assert stats["removed"] == 1 and stats["written"] == 1
rows = _rows(sqlite_engine, "SELECT * FROM core_share_lot")
assert [r["lot_id"] for r in rows] == [bootstrap_lot_id(CUSTOMER, PRODUCT_A)]
assert _dec(rows[0]["remain_qty"]) == Decimal("100")
def test_rebuild_lots_dry_run_writes_nothing(sqlite_engine):
"""`--dry-run`:只报告,不写库。"""
rebuild = _load_rebuild_lots()
_seed_customer(sqlite_engine, CUSTOMER)
_seed_product(sqlite_engine, PRODUCT_A)
_seed_holding(sqlite_engine, CUSTOMER, PRODUCT_A, "100", "100.00", date(2026, 8, 1))
stats = rebuild.rebuild_lots(sqlite_engine, customer_id=CUSTOMER, dry_run=True)
assert stats["bootstrapped"] == 1 and stats["written"] == 0
assert _rows(sqlite_engine, "SELECT * FROM core_share_lot") == []
def test_d18_gateway_bootstrap_and_rebuild_lots_agree():
"""**D18 同源断言**:同一 `core_holding` 行,网关兜底补建与 `rebuild_lots.py`
算出的批次必须**完全一致**(`confirmed_at` 逐一相等)。
这是 D18 设立的唯一目的 —— 两处各写一份「由 `as_of` 反推」必然漂移,
同一持仓会落到不同费率档(演示看不出、生产是错账)。故本用例同时验证
「两侧都在调 `bootstrap_lots`」这一事实:任一侧改成自己实现即变红。
"""
rebuild = _load_rebuild_lots()
gateway_engine = create_sqlite_engine()
script_engine = create_sqlite_engine()
try:
for engine in (gateway_engine, script_engine):
_seed_customer(engine, CUSTOMER)
_seed_product(engine, PRODUCT_A)
_seed_nav(engine, PRODUCT_A, "1.0000")
_seed_holding(engine, CUSTOMER, PRODUCT_A, "100", "123.45", date(2026, 8, 1))
# 侧 A:网关 D8 兜底补建(经 redeem 触发;补建 100 后会扣 30)
_maintain(
gateway_engine, trade_id="TRD-T10-D18", trade_type="redeem", amount="30"
)
# 侧 B:rebuild_lots 快照重建(不扣减)
rebuild.rebuild_lots(script_engine, customer_id=CUSTOMER)
got = _rows(gateway_engine, "SELECT * FROM core_share_lot")[0]
want = _rows(script_engine, "SELECT * FROM core_share_lot")[0]
assert got["lot_id"] == want["lot_id"]
assert got["confirmed_at"] == want["confirmed_at"] # ← D18 的核心断言
assert _dec(got["nav"]) == _dec(want["nav"])
assert _dec(got["qty"]) == _dec(want["qty"])
# 差异只应来自「侧 A 扣了 30 份」这件事本身
assert _dec(got["remain_qty"]) == _dec(want["remain_qty"]) - Decimal("30")
finally:
gateway_engine.dispose()
script_engine.dispose()
def test_rebuild_lots_uses_same_pure_function_as_gateway():
"""机制验证(防「两侧碰巧算出同值」的假绿):脚本模块内的补建**必须**来自
`bootstrap_lots`,而不是自带的第二份实现。"""
rebuild = _load_rebuild_lots()
assert rebuild.bootstrap_lots is bootstrap_lots
row = {
"customer_id": CUSTOMER,
"product_id": PRODUCT_A,
"qty": "100",
"cost_amount": "123.45",
"as_of": date(2026, 8, 1),
}
assert [lot.lot_id for lot in rebuild.bootstrap_lots(row)] == [
lot.lot_id for lot in bootstrap_lots(row)
]
+89 -1
View File
@@ -4,7 +4,7 @@
_normalize_trades 兜底)。API 层经 monkeypatch 注入 sqlite 仓储。 _normalize_trades 兜底)。API 层经 monkeypatch 注入 sqlite 仓储。
""" """
from datetime import datetime, timedelta from datetime import date, datetime, timedelta
from decimal import Decimal from decimal import Decimal
import pytest import pytest
@@ -189,6 +189,94 @@ def test_redeem_accepted_without_alert(env):
assert _counts(engine, "risk_alert") == 0 assert _counts(engine, "risk_alert") == 0
# ---------- T-10:普通申赎批次维护(FR-C16 · R7 补断言) ----------
def _seed_nav(engine, pid: str, nav: str, nav_date: date) -> None:
with engine.begin() as conn:
conn.execute(
text(
"INSERT INTO core_product_nav (product_id, nav, daily_chg_pct, nav_date) "
"VALUES (:pid, :nav, 0, :nd)"
),
{"pid": pid, "nav": float(nav), "nd": nav_date},
)
def _seed_lot(engine, lot_id: str, cid: str, pid: str, remain: str, nav: str,
confirmed_at: datetime) -> None:
with engine.begin() as conn:
conn.execute(
text(
"INSERT INTO core_share_lot (lot_id, customer_id, product_id, qty, "
"remain_qty, nav, confirmed_at) "
"VALUES (:lot, :cid, :pid, :q, :q, :nav, :cat)"
),
{
"lot": lot_id, "cid": cid, "pid": pid, "q": float(remain),
"nav": float(nav), "cat": confirmed_at,
},
)
def _norm_dt(value: object) -> datetime:
if isinstance(value, datetime):
return value
return datetime.fromisoformat(str(value)[:19])
def test_subscribe_maintains_share_lot(env):
"""T-10(R7):申购放行后按 T 日净值建批次 —— 批次表由交易统一出入口维护。"""
core, repo, writer, _, engine = env
_seed_nav(engine, "PROD-510300", "1.2000", date(2026, 9, 6))
resp = submit_trade(
_req(customer="CUST-3001", product="PROD-510300", amount="12000"),
core_ro=core, risk_repo=repo, gateway_repo=writer,
now=datetime(2026, 9, 6, 14, 0, 0),
)
assert resp["blocked"] is False
with engine.connect() as conn:
row = conn.execute(
text("SELECT * FROM core_share_lot WHERE source_trade_id = :t"),
{"t": resp["trade_id"]},
).mappings().one()
assert row["lot_id"] == f"LOT-SUB-{resp['trade_id']}"
assert Decimal(str(row["qty"])) == Decimal("10000.00") # 12000 / 1.2
assert _norm_dt(row["confirmed_at"]) == datetime(2026, 9, 6, 14, 0, 0)
def test_redeem_deducts_share_lot_fifo(env):
"""T-10(R7):赎回放行后 FIFO 扣减既有批次(`amount ÷ T 日净值` 折算份额)。"""
core, repo, writer, _, engine = env
_seed_nav(engine, "PROD-510300", "1.0000", date(2026, 9, 6))
_seed_lot(engine, "LOT-R16", "CUST-3001", "PROD-510300", "100", "1.0000",
datetime(2026, 8, 1, 10, 0, 0))
resp = submit_trade(
_req(customer="CUST-3001", product="PROD-510300", ttype="redeem", amount="50"),
core_ro=core, risk_repo=repo, gateway_repo=writer,
now=datetime(2026, 9, 6, 14, 0, 0),
)
assert resp["blocked"] is False
with engine.connect() as conn:
remain = conn.execute(
text("SELECT remain_qty FROM core_share_lot WHERE lot_id = 'LOT-R16'")
).scalar_one()
assert Decimal(str(remain)) == Decimal("50.00") # 100 − 50/1.0
def test_blocked_trade_does_not_maintain_lots(env):
"""阻断路径不写批次:批次维护在放行分支内,适当性不匹配时两颗表都不动。"""
core, repo, writer, _, engine = env
_seed_nav(engine, "PROD-161725", "1.0000", date(2026, 9, 6))
resp = submit_trade(
_req(customer="CUST-1001", product="PROD-161725", amount="10000"), # C1 买 R4 → 阻断
core_ro=core, risk_repo=repo, gateway_repo=writer,
now=datetime(2026, 9, 6, 14, 0, 0),
)
assert resp["blocked"] is True
assert _counts(engine, "core_share_lot") == 0
def test_engine_failure_is_audited_and_degraded(env): def test_engine_failure_is_audited_and_degraded(env):
"""评审 P1-1:引擎异常 → 审计 risk_engine_error + 响应 engine_error=true(交易已成立)。""" """评审 P1-1:引擎异常 → 审计 risk_engine_error + 响应 engine_error=true(交易已成立)。"""
core, repo, writer, pub, engine = env core, repo, writer, pub, engine = env