diff --git a/app/repository/convert_repository.py b/app/repository/convert_repository.py index 16e34a4..22a9eb3 100644 --- a/app/repository/convert_repository.py +++ b/app/repository/convert_repository.py @@ -12,16 +12,20 @@ S2(评审):清理**标记不硬删**——`mark_expired` 置 `status='expi from __future__ import annotations +import logging from datetime import date, datetime, timedelta from decimal import Decimal from typing import Any from sqlalchemy import text from sqlalchemy.engine import Engine +from sqlalchemy.exc import IntegrityError from app.config.settings import settings from app.utils.db import get_engine +logger = logging.getLogger(__name__) + def _to_bind(value: Any) -> Any: """Decimal 转 float 再绑定(sqlite 不支持直接绑定 Decimal;MySQL DECIMAL 列自动收口)。 @@ -45,25 +49,67 @@ class ConvertRepository: # ---------- 阶段零:占位(uk_idem 兜底) ---------- - def insert_placeholder(self, group_id: str, client_request_id: str | None) -> None: - """阶段零占位:插一行 `status='pending'`。 + def insert_placeholder(self, group_id: str, client_request_id: str | None) -> bool: + """阶段零占位:插一行 `status='pending'`;**已存在则置回 pending**。 + + 返回 **True = 本笔持有该占位,可以继续**;**False = 该 `client_request_id` + 已被另一个 `group_id` 占住**(同键并发,本笔必须让路,由调用方回 202)。 `client_request_id` 为 None 时绑 NULL——MySQL / sqlite 的 UNIQUE 约束均允许多个 NULL, 故「无幂等键的请求」可重复占位、互不冲突(uk_idem 仅对非空键兜底)。 `estimated=0`:convert 占位是真实请求,非风控预估单(与 risk_alert 语义区分)。 + + **三步法(R-a:弃用方言 UPSERT,改「先查再 INSERT 或 UPDATE」)**。 + 为什么必须容错「行已存在」:阶段一失败(典型是 `LotConflict` 409)时占位已被 + `mark_failed` 置为 `failed`,而架构 §8.3 要求调用方带**同一** `client_request_id` + 退避重试(≤3 次、100/200/400ms);重试会走「复用原 group_id 重跑」分支再次进入 + 阶段零——此时 `uk_group` 与 `uk_idem` 都已被那一行占用,朴素的 INSERT 必撞唯一键, + 把**可重试的 409 升级成 `IdempotencyUnavailable`(503)**,且是**确定性的**(重试永不成功), + 与 errors.py 里「瞬时状态、恢复后重试即可成功」的注释相反。(T-13 压测前置修复) + + 置回 `pending` 是正确语义:同一 `group_id` 的这一次尝试正在进行中。 + 到达此处时既有行的状态只可能是 `pending`(上一轮中途崩溃)/ `failed`(阶段一或阶段二失败); + `completed` 已在 `convert_fund` 幂等前置分支返回,不会走到这里。 + + **为什么还要 catch IntegrityError**:`convert:idem:{cid}` 锁只包住幂等判定 + (出块即释放),两笔同键请求可能**都判定为"无占位"**、各自生成了不同的 `group_id` + (T-13 真库实测:8 路并发下偶发);后插入的那笔撞 `uk_idem` → 此前会直穿 503。 + 这里把它收敛为 False(让路),交由调用方回 202(架构 §9「同键并发 → 202」)。 """ - sql = text( - """ - INSERT INTO risk_convert_detail - (convert_group_id, client_request_id, status, estimated) - VALUES (:gid, :cid_req, :status, 0) - """ - ) - with self._engine.begin() as conn: - conn.execute( - sql, - {"gid": group_id, "cid_req": client_request_id, "status": _STATUS_PENDING}, + try: + with self._engine.begin() as conn: + existing = conn.execute( + text("SELECT 1 FROM risk_convert_detail WHERE convert_group_id = :gid"), + {"gid": group_id}, + ).first() + if existing is not None: + conn.execute( + text( + "UPDATE risk_convert_detail" + " SET status = :s WHERE convert_group_id = :gid" + ), + {"s": _STATUS_PENDING, "gid": group_id}, + ) + return True + conn.execute( + text( + """ + INSERT INTO risk_convert_detail + (convert_group_id, client_request_id, status, estimated) + VALUES (:gid, :cid_req, :status, 0) + """ + ), + {"gid": group_id, "cid_req": client_request_id, "status": _STATUS_PENDING}, + ) + except IntegrityError: + # uk_idem 竞态:同键的另一笔刚刚占位成功 → 本笔不持有,让路 + logger.info( + "幂等键已被并发请求占用(本次让路):cid=%s gid=%s", + client_request_id, + group_id, ) + return False + return True # ---------- 阶段二:回写 completed + 详情 ---------- diff --git a/app/service/convert/convert_service.py b/app/service/convert/convert_service.py index 97cca33..05b196c 100644 --- a/app/service/convert/convert_service.py +++ b/app/service/convert/convert_service.py @@ -427,8 +427,22 @@ def convert_fund( ) if rebuilt is not None: return rebuilt - # 阶段一未成 → 复用同一 group_id 重跑,杜绝第二组流水 - group_id = hit_gid + # 阶段一未成:必须区分「确定没跑成」与「在飞/未知」,否则同键并发会双跑 + # —— `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 @@ -496,11 +510,17 @@ def convert_fund( # ⑤ 阶段零:占位(仅带幂等键时;在全部 4xx 之后,故 4xx 不留占位) if placeholder: try: - crepo.insert_placeholder(group_id, cid_req) + 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) diff --git a/docs/PRD/PRD-基金转换交易.md b/docs/PRD/PRD-基金转换交易.md index f10fd89..0b4c5e0 100644 --- a/docs/PRD/PRD-基金转换交易.md +++ b/docs/PRD/PRD-基金转换交易.md @@ -1,6 +1,6 @@ # PRD · 基金转换(convert)交易 -> 版本:**v0.9.2(展示位数补全 · 定稿)** · 日期:2026-09-10 +> 版本:**v0.9.3(§9 第 18 条实测补录 · T-13)** · 日期:2026-09-10 > 分支:`risk-control-agent` > 状态:第 2 步产出。**v0.7 已冻结**(架构 v0.2 §0 查证发现 4 项合规硬伤,外审判定「❌ 不建议进入编码」) > → 回退修订出 **v0.8**(按处置表全量修订,33 条闭环)→ 架构 v1.0 独立评审**通过**后,回填 2 处契约出 **v0.9**。 @@ -904,11 +904,26 @@ SELECT 1 FROM core_trade WHERE convert_group_id = :gid LIMIT 1; 不得返回 `BELOW_MIN_QTY`(第三轮第 5 条回归点) 17. **阶段二失败可补偿**:`rebuild_alerts.py` 能仅凭 Core 侧数据 (含 `core_convert_lot_detail.nav`/`nav_date`)补出完整详情(第三轮第 2 条回归点) -18. **性能(P2,第五轮第 11 条 · v0.7 自查问题 D 修订)**: +18. **性能(P2,第五轮第 11 条 · v0.7 自查问题 D 修订 · v0.9.3 实测补录)**: 一次转换端到端响应 **< 2s**(本地模拟库环境); 阶段一单库事务 **预估 < 100ms**(**预估值,非验收硬指标**; 第 6 步集成测试**实测补录真实耗时**;若实测超阈值,需优化索引 / 锁策略后再定阈值, **不得反向修改实测数据迁就指标**) + + > **实测补录(T-13 · 2026-09-10 · n=59 剔除首个冷启动样本)** + > 测量条件:本机 MySQL **8.0.46**(127.0.0.1,InnoDB / REPEATABLE-READ)、模拟库、 + > **单线程顺序**、份额池充足的稳态、阶段 1.5 用空引擎(隔离引擎开销); + > 用例:`tests/test_convert_concurrency.py::test_performance_probe_stage1_and_end_to_end` + > (`CONVERT_STRESS=1` 开启)。 + > + > | 指标 | P50 | P95 | max | 结论 | + > | --- | --- | --- | --- | --- | + > | **阶段一单库事务**(`apply_convert`,含事务开提交) | **9.8 ms** | **12.8 ms** | **16.7 ms** | 远低于「预估 < 100ms」,**无需优化索引/锁策略** | + > | **端到端**(`convert_fund` 八步,含两次库往返 + 阶段二) | **44.2 ms** | **51.5 ms** | **59.8 ms** | 远低于 **2s** 硬指标(约 1/33) | + > + > 口径说明:端到端为 **service 层**耗时,未含 ASGI/HTTP 栈;真库 HTTP 端到端由 + > `tests/test_convert_integration.py` 覆盖(同环境同数量级)。三次复跑 P50 波动 + > ≤ 0.4ms,数据稳定。**未修改任何实测值以迁就指标。** 19. **补差费非零场景有覆盖**(v0.7 自查问题 A 衍生 · **Q8 已定稿**): 存在两端 `subscribe_fee_rate` 不同的产品对(`PROD-110022` 0.30% / `PROD-003095` 0.80%,且**同属一个管理人 + 同一 TA**),且有用例验证 `diff_fee > 0` 时 §2.1 公式成立(`in_amount = convert_amount − diff_fee`); diff --git a/docs/memory/2026-09-10.md b/docs/memory/2026-09-10.md index 406360c..dcbbc4c 100644 --- a/docs/memory/2026-09-10.md +++ b/docs/memory/2026-09-10.md @@ -610,4 +610,40 @@ B.1 状态 / B.6 拓扑与基线 / **新增 B.6.3**)· `docs/memory/{TODO,MEMO **文档回写**:开发计划 §9(DoD 全勾 + **新增 §9.1 执行记录**)+ §7.3 收口段 · `交接文档.md` **v2.3** (表头 / §0 三线状态 / B.1 / B.6 拓扑与基线 / **新增 B.6.5**)· `docs/memory/{TODO,MEMORY}` · 本条。 -**下一步 = T-13(全量回归 + 50 并发压测 + 性能补录 · 本线最后一个任务)**。 +## ✅ T-13 全量回归 + 50 并发压测 + 性能补录 已完成(2026-09-10 · 本线最后一个任务 · 已结项) + +**执行期发现并修复 2 处并发缺陷**(压测暴露,非需求变更;sqlite 单测均不可见) + +1. **阶段一失败后同键重试 → 永久 503**:阶段一 `LotConflict`(409) → 占位 `mark_failed` → 同键重试 + 复用 group_id 重入阶段零 → `insert_placeholder` 朴素 INSERT 撞 `uk_group`/`uk_idem` → + `IdempotencyUnavailable`(503)。§8.3 明说 409 可重试,此路径**重试永不成功**。 + 修法:`ConvertRepository.insert_placeholder` 改**三步法**(R-a),已存在则置回 `pending`。 +2. **同键并发的两个竞态子窗口**(**资金安全**): + - 子窗口 A(T2 在占位后进入):各自 `group_id` 不同 → 旧代码把 `pending` 当"未成"重跑阶段一 + → **真·双扣**(突变实测扣 120→**240**、两组流水)。 + - 子窗口 B(T2 在占位前进入):撞 `uk_idem` → `IntegrityError` 直穿 503。 + - 修法:① 幂等判定 `pending` → **202**(不当"未成"重跑);② `insert_placeholder` 返回"是否持有占位", + 撞键即**让路** → 调用方 **202**。**两处缺一不可**。 + - 复现用**确定性交错**(`threading.Event` / `Barrier` 拦住 T1 阶段一)——把"碰运气"变成"必现"。 + +**50 并发压测结论(评审 Q6)**:50 × 2000 争 50000 → **不超卖硬不变量成立**(实测成交 20~25,理论上限 25)· +**紧池退避 40/50 = 80%(不够用)** · **松池 50/50 = 100%** · `convert_group_id` 无重复、`remain_qty` 守恒。 +> **⚠️ 新增待评估项**:跨客户并发触发 InnoDB **1213 死锁**(唯一索引上插入意向锁互斥,8 路实测 3~5 次); +> `apply_convert` **不重试死锁**,上层按 §8.3 当可重试异常收敛后 8/8 成功。**服务端是否补重试待用户评估**。 +> **退避不够用的根因不在间隔**:一笔需求被拆成 `400+600` 碎片后,碎片所在批次随时被清零 → 重试仍抢不到; +> 要收敛到 100% 需**增加重试次数**或**冲突后换批次重规划**,单纯拉长间隔无效。 + +**性能实测(PRD §9 第 18 条,v0.9.3)**:端到端 **P50 44.2 / P95 51.5 / max 59.8 ms**(阈值 2s)· +阶段一 **P50 9.8 / P95 12.8 / max 16.7 ms**(预估 <100ms)→ **远未触阈值、未改任何实测值、无需重定阈值**。 + +**验证**:`pytest -q` **736 passed / 10 skipped**(731 **+5**,两遍稳定)· `CONVERT_STRESS=1` 并发 **9 passed** · +按 SOP 重灌双库后复跑零回归 · 7 个真库脚本复跑零回归 · **突变 4 组全部精准命中**(超卖 50/50+余 −50000 / +IntegrityError / 真·双扣 240 份 / UNIQUE constraint failed),已恢复无残留。 + +**门禁语义**:`CONVERT_STRESS=1` 才跑 7 条重载用例(另 2 条常跑);常规 CI 逐条 skip —— 与架构 §10「不进 CI」一致。 + +**文档回写**:开发计划 **§10**(10.1~10.5)· PRD **v0.9.3** · `交接文档.md` **v2.4**(新增 **B.6.6** + B.8 补 3 条坑)· +`docs/memory/{TODO,MEMORY,FRAMEWORK,ITERATION}` · 本条。 + +**✅ 基金转换线全部结项(T-0 ~ T-13)**,基线 **510 → 736 passed / 10 skipped**,**无待办**。 +**T-13 尚未提交**(8 个文件在工作区);提交与否由用户决定,**不擅自 push**。 diff --git a/docs/memory/FRAMEWORK.md b/docs/memory/FRAMEWORK.md index 57d1862..7c6f651 100644 --- a/docs/memory/FRAMEWORK.md +++ b/docs/memory/FRAMEWORK.md @@ -39,7 +39,7 @@ | 代理人助手 Agent | L2 画像、RAG、草稿 | L1 只读、Milvus | 空壳 service(chat 骨架已通) | | 数据分析 Agent | NL→SQL→解读 | Core RO、画像只读 | 空壳 service(chat 骨架已通) | | 风控监测 Agent | 预警、L3、R-02 适当性 | 交易事件、AML 名单 | **已实现 B1~B7 + C1~C3(对话线四 Tool + risk 分支 StateGraph,A-6 验收)**;追加 FR-8/9/10(C4~C6:集中度/时效升级/行为链)方案已定稿待编码 | -| **基金转换(convert)交易** | 转出赎回 + 转入申购、逐批次计费、补差费、批次维护 | Core 写(仅 `app/gateway/`)、`core_share_lot` 批次表、规则引擎 | 🟡 **第 5 步进行中**(**T-0 / T-0b / T-1 / T-2 / T-2b 已完成(609 passed),下一步 = T-3**):PRD **v0.9.1** + 架构 **v1.0** + 开发计划 **v1.0**(T-3~T-13 待实现);入口 项目根 `交接文档.md` §B | +| **基金转换(convert)交易** | 转出赎回 + 转入申购、逐批次计费、补差费、批次维护 | Core 写(仅 `app/gateway/`)、`core_share_lot` 批次表、规则引擎 | ✅ **已实现(T-0 ~ T-13 全部完成 · 2026-09-10,736 passed / 10 skipped)**:PRD **v0.9.3** + 架构 **v1.0** + 开发计划 **v1.0**;HTTP 层 convert 已走通(`/api/simulate` 三型分池 → `trade_gateway._submit_convert` → 八步编排 → 阶段 1.5 引擎 → 补偿脚本);入口 项目根 `交接文档.md` §B | | Core 只读层 | L0 事实查询 | `jinrong_core` | **已实现 + 已接对话 Tool(T-04)**:core_ro 经 app/tool/core_tools.py 三只读 Tool(L0/持仓/流水)进 chat;风控扩展查询照旧 | | 共用底座 | 会话、审计、输入防护 | MySQL 11 表 + Redis | **已接入(2026-09-07)**:会话(T-06 session_repository + memory_service 窗口)、审计中间件(T-02 http_access + input_guard_log 双写)、agent_tool_call Tool 留痕(T-04)、输入防护(T-03 input_guard:注入词表纯函数检测 + oversize + Redis 固定窗口限流,chat 链路 限流→注入/超长→归属) | | 对话编排 | LangGraph StateGraph + DeepSeek | langgraph/langchain-openai | **已实现(T-07 骨架 + T-04 Tool 节点 + C1 风控四 Tool + C5 追加 query_overdue_alerts 共五只 + C2 risk 分支)**:tool(分组关键词意图→Tool,customer/advisor/risk)→llm→guard(免责声明);analyst 分支与 LLM intent 待后续 | diff --git a/docs/memory/ITERATION.md b/docs/memory/ITERATION.md index f2940a3..b846b04 100644 --- a/docs/memory/ITERATION.md +++ b/docs/memory/ITERATION.md @@ -22,3 +22,4 @@ | 2026-09-10 | **基金转换线第 5 步开工 · 第 0 批门禁 T-0 + T-0b 完成**:基线 510 → **516 passed / 3 skipped**。**T-0**:`tests/_ddl.py` 重写 `core_holding`(`qty`/`cost_amount`/`as_of`/`pnl_pct` + PK + `uk_cust_product`)+ 新增 `core_product_nav` + `_assert_ddl_aligned()` 建库自校验;2 处测试 INSERT 改 `qty` **并补 3 个 NOT NULL 列**(**原计划漏项**);`test_db.py` +3 用例(含**门禁反向验证的自动化用例**)。**T-0b(D20 账号分离)**:新增 `scripts/core/00-grant.sql`(3 账号逐表授权,不进 reset.ps1);`settings.py` +6 项;`db.py` `get_engine(db, role)` + `_resolve_credentials`(缓存键 `(db, role)`,未配置回退 `mysql_user`);`core_ro`→`ro` / `gateway_repository`→`rw` / `risk_repository`·`session_repository`→`rw` 显式;`conftest.py` 4 处→`admin`(R-e)。**两处口径修正**:① `audit_log` 实授 **`SELECT, INSERT`**(字面「只授 INSERT」会连读一起剥夺,致 `has_engine_error_audit`/`list_audit_events` 失权)② conftest 必须 admin。**遗留环境操作**:`00-grant.sql` 需管理员执行 + 写 `.env`(3 条权限断言未配置时自动 skip) | AIcoding 第 5 步(T-0/T-0b 为阻断前置);用户「你进行下一步」 | tests/_ddl.py · tests/test_db.py · tests/conftest.py · tests/test_chat_tools.py · tests/test_concentration_c4.py · app/utils/db.py · app/config/settings.py · 4 个 Repository · **scripts/core/00-grant.sql(新)** · 开发计划-基金转换(§3.3) · 架构设计-基金转换(§11.1 实施记录) · 交接文档-基金转换(v1.2) · MEMORY / TODO / 当日日志 | | 2026-09-10 | **交接文档合并为单一入口 + 全仓指向统一**:三份交接文档合并为项目根 `交接文档.md`(**v3.0**,545 行)—— **§0 公共层** + **§A 风控模块主线** + **§B 基金转换线**(含工作区最新版 v1.4:T-0~T-2b 完成 / 609 passed)+ **§C 架构改进线**(含工作区最新版 v1.2:结项态 / 7/7 冒烟)。**全仓 10 处指向统一改到 `交接文档.md` §A/§B/§C**(`AGENTS.md` · `docs/memory/{MEMORY,TODO,FRAMEWORK}` · `PRD-架构改进与稳定性加固` · `开发计划-架构改进` · `TODO-架构改进` · `开发计划-基金转换交易`);两份 `docs/交接文档-*.md` **保留为历史留档并加醒目「已废弃·勿读」横幅**(未删除)。⚠️ `交接文档.md` 在 `.gitignore:47` 内、**不入库**(用户 2026-09-10 明确「不进提交」,交接文档只留本地) | 用户「统一一下,不是早都说过交接文档是给你 AI 接手用的吗」—— 交接文档的定位是 **AI 接手入口**,因此全仓指向必须唯一且不可指向过期文件,否则新会话会被误导到旧版(`docs/` 两份留档停在 v1.0/v1.1,不含 T-2/T-2b 等进度) | `交接文档.md`(**v3.0**,未入库) · `AGENTS.md` · `docs/memory/{MEMORY,TODO,FRAMEWORK,ITERATION}` · `PRD-架构改进与稳定性加固` · `开发计划-架构改进` · `TODO-架构改进` · `开发计划-基金转换交易` · `docs/交接文档-基金转换.md` / `docs/交接文档-架构改进.md`(加废弃横幅,保留留档) | | 2026-09-10 | **【事故·已恢复】`git rm` 会静默抹除整个 `docs/` 目录**:对 `docs/` 下两个中文名文件执行 `git rm`(带不带 `-f` 均如此)→ **`docs/` 全部 54 个文件从工作区消失**,且**返回退出码 0 不报错**(**复现 2 次**)。已排除中文路径(同路径 `ls` 正常)与自定义 hooks(`core.hooksPath` 空)→ **触发条件即 `git rm` 本身**,根因未查明。**零内容丢失**:`git checkout HEAD -- docs/` 完整还原 54 文件。改判:删文件用**纯 `rm -f <精确路径>`** + `git add -A` | 工具链缺陷(Git Bash on Windows + 中文名路径);教训:批量删除后必须立即复核文件数,异常立即从 HEAD 还原 | `docs/`(还原) · 规避纪律写入 `交接文档.md` §0.4 第 10 条 + `.workbuddy/memory/MEMORY.md` | +| 2026-09-10 | **基金转换线第 5 步 + 第 6 步全部完成(T-3 ~ T-13 结项)**:批 2 仓储/代理/锁(`core_ro` 五方法 + `share_lot_repository` + `convert_repository` + `locks.try_lock`)→ **T-6** 阶段一单事务(`convert_core_repository.apply_convert`)→ **T-7** 八步编排(关键路径)→ **T-8** 规则引擎双视图分流(`amount_view` + `process_convert_event`,阶段 1.5 由「静默跳过」变「真跑」)→ **T-9** API 三型分池 + 网关分派(**HTTP 层 convert 走通** + 展示位数口径修复)→ **T-10** 普通申赎批次维护(FR-C16 + `rebuild_lots.py`)→ **T-11** 工具与 SQL 汇总去重(FR-C15 / R-d,`_amount_view` 提升为公开 `amount_view`)→ **T-12** 补偿脚本(`compensate_convert` 单点入口 + `cleanup_pending_convert.py`)→ **T-13** 全量回归 + 50 并发压测 + 性能补录。**基线 510 → 736 passed / 10 skipped**;7 个 convert 真库脚本复跑零回归;**T-13 压测暴露并修复 2 处并发缺陷**(阶段一失败后同键重试永久 503 · 同键竞态子窗口 A 真·双扣 / B 双扣+503);性能端到端 max 59.8ms / 阶段一 max 16.7ms(远未触阈值);跨客户并发 1213 死锁登记待评估 | AIcoding 第 5 步逐任务开发 + 第 6 步集成测试;用户「开」指令推进 T-13 收尾;方言语义一律上真库(用户 2026-09-10 拍板) | `app/service/convert/*` · `app/gateway/convert_core_repository.py` · `app/repository/{core_ro,convert_repository,share_lot_repository}.py` · `app/service/risk/{rules,engine,alert_service}.py` · `app/api/simulate.py` · `app/gateway/trade_gateway.py` · `tests/test_convert_{calc,service,core,integration,engine,concurrency}.py` · 7 个 `scripts/dev/verify_convert_*.py` · `scripts/agent/cleanup_pending_convert.py`(新) · `scripts/core/rebuild_lots.py`(新) · PRD-基金转换(v0.9.3) · 开发计划-基金转换(§10) · `交接文档.md`(§B v2.4,未入库) · MEMORY / TODO / FRAMEWORK / ITERATION / 当日日志 | diff --git a/docs/memory/MEMORY.md b/docs/memory/MEMORY.md index aab8045..f97245f 100644 --- a/docs/memory/MEMORY.md +++ b/docs/memory/MEMORY.md @@ -11,7 +11,7 @@ **当前进度:** 需求与表设计已定 · **风控模块 B1~B9b 全部完成(M2 tag risk-m2),M4 复核已闭环(2026-09-07),风控阶段 B 正式完结** · **Wave 0 已完成(2026-09-07,经独立 AI 评审闭环)**:T-01 JWT 鉴权(auth_service Auth SDK + deps 工厂替换 + X-Agent-Type 准入矩阵)/ T-02 审计中间件(http_access + 独立 request_id + 4xx/500 统一错误体 + input_guard_log 双写)/ T-06 chat 最小闭环(POST /api/chat + 会话落库 + Redis 窗口)/ T-07 LangGraph StateGraph 骨架 + DeepSeek(无 key 降级)· **T-04 Core RO Tool 节点已完成(2026-09-07)**:app/tool/core_tools.py 三只读 Tool + tool_service(意图/归属校验/run_tool)+ 图 tool 节点 + agent_tool_call 落库 + `utils/authz.py` 公共鉴权留痕,**321 测试绿(首评+复审双闭环)**。风控阶段 C 已完成(2026-09-07,345 绿,tag `risk-m3`,A-6 对话线验收通过,独立 AI 评审 PASS P0=0)。T-03 输入防护已完成(2026-09-07,378 绿,独立 AI 评审 PASS with findings P0=0):app/service/input_guard.py 注入词表 45 条纯函数检测 + oversize 4000 + actor 级 Redis 固定窗口限流 30 次/分(fail-open);chat 链路顺序 = 鉴权→准入→空白→限流 429→注入/超长 400→归属→会话,被拒 fail-fast 不建会话,blocked 落 input_guard_log(ENUM 四值已用满)。**阶段一「对齐 main 基准」AL-01~AL-08 已完成(2026-09-07,逐项独立 commit bb244f4~b5fd52e):适当性判定换核为 main 的 core_ro.check_suitability(C×R 矩阵表数据驱动,match_result 五值/JR-AST-012/FM-01/FM-03/JR-AST-PRO 契约,SUIT-001~008 退役);risk_suitability_log 重建 21 列;Core 表加 is_hnw/风评七新列(expires_at);种子 33 客户/14 产品;全量 406 passed 0 failed 0 skipped + uvicorn 冒烟三端点通过;risk-m1 已补打(指向 3c07de6)**。**下一步:阶段一验收门(**合并 main 前的最终交付检查点,非分支开发阻塞**;用户浏览器目视确认 UI)→ AL-09 合并 main → AL-10 PRD v1.2 → AL-11 docx 登记 → 阶段二 C4~C6(**已完成并打 `risk-m4` tag**)/ 前端 React 多 Agent 入口(HashRouter `web/` init)。****开发在分支 `risk-control-agent`(与 origin/main 已分叉:领先 84 提交 / 落后 0,origin/main 为分支祖先;**分支已推送远程,origin/risk-control-agent 同步于 `1d00e53`,本地跟踪已建立**;**架构改进与稳定性加固 T-101~T-202 已于 2026-09-09 完成并推送 origin/risk-control-agent(2d0e2fa..f733fc2 快进,pytest 510 绿,较 503 基线 +7 例:T-201.3 双层锁 6 例 + T-202 trace 顺序守卫 1 例),含 Redis 双层分布式锁 T-201 与审计中间件 trace 顺序守卫 T-202,均为文档/告警/锁原语修正、接口契约与表结构零变更;远程已提示走 PR 合并 main;**合并(risk-control-agent → main)由合并执行人负责,不归用户管**——AI 只负责把分支推到远程 + 更新文档,不得主动发起 PR / 执行合并,移交执行人按《合并注意事项-风控模块并入main.md》操作。** 旧文档中的 `feature/risk` 为过时口径)。** -**⚡ 并行新线 · 基金转换(convert)交易(2026-09-10 设计+开发计划闭环,**第 5 步进行中:T-0 ~ T-12 已完成**):** `PRD-风控监测Agent` FR-1 一期显式拒收 convert,本线将其放开。**AIcoding 第 1~4 步已完成**:PRD **v0.9.2 定稿**(33 条外审闭环 + 2 处架构回填 + **费率分类修正** + **展示位数补全**)→ 架构 **v1.0.1 定稿**(独立评审通过,13 条建议 **0 悬空**,接受 10 / 修正性接受 3 / 驳回 0)→ 门控 **M-7 已满足** → **开发计划 v1.0 已产出(`docs/项目框架设计/开发计划-基金转换交易.md`,**1,047 行**)并经**四轮**独立子代理审核收敛(**12 条意见 → 接受 11 / 驳回 1(附实测证据)/ 0 悬空**)。**第 5 步进行中:T-0 + T-0b + T-1 + T-2 + T-2b + T-3 + T-4 + T-5 + T-6 + T-7 + T-8 + T-9 已于 2026-09-10 完成(T-9 里程碑时 697 passed / 3 skipped;T-1 断言 8/8 PASS;T-2 纯函数包 7 文件 / 93 用例;T-2b 实算脚本 15/15 一致;T-6 真库 24/24;T-7 真库 35/35;T-8 真库 31/31;T-9 真 MySQL 集成 8 条 + 11 条错误码映射 + 4 组突变验证 + **展示位数口径修复**)**,**T-10 ✅(普通申赎批次维护 · 回归风险最大 · +17 用例 / 真库 20/20)、T-11 ✅(工具与 SQL 汇总去重 · +4 用例 / 真库 14/14)与 T-12 ✅(补偿脚本 · +12 用例 / 真库 34/34 · 全套 7 真库脚本复跑零回归)均已完成**,**下一步 = T-13(全量回归 + 50 并发压测 + 性能补录 · 本线最后一个任务)**。**开工前置两个阻断项(✅ 2026-09-10 均已完成)**:**T-0**(sqlite/MySQL 结构对齐:`core_holding` 列名+PK **+ 补 `core_product_nav`**,建库自校验,**同步改 2 处测试 INSERT 并补 3 个 NOT NULL 列**)与 **T-0b**(**DB 账号分离 D20**:`xh_core_ro` SELECT 全库 / `xh_core_rw` **4 表写**无 DELETE·DDL / `xh_agent_rw` **`audit_log` 只授 SELECT+INSERT**(不可改删);**`conftest.py` 四处 engine 已显式 `role="admin"`**)。**遗留环境操作**:`scripts/core/00-grant.sql` 需管理员执行一次 + 写 `.env`,未执行时 3 条权限断言自动 skip,不阻塞 T-1。设计资产五件套:`docs/PRD/PRD-基金转换交易.md` · `docs/项目框架设计/架构设计-基金转换交易.md` · **`docs/项目框架设计/开发计划-基金转换交易.md`(新)** · `评审待办-风控主架构与基金转换.md` · `基金转换-审查意见处置表.md`。**开工前必读开发计划 §1.4(15 条代码事实)/ §1.5(7 条实现级裁定 R-a~R-g)/ §12(16 条回归面)** —— 尤其是 **R-a(弃用方言 UPSERT)· R-b(流水写 redeem/subscribe 不写 convert)· R-e(conftest 用 admin)** 三条,不读必踩。**交接入口:项目根 `交接文档.md` §B(三线合并版唯一入口)**。合规基准 = 证监会公告〔2025〕22 号。 +**⚡ 并行新线 · 基金转换(convert)交易(2026-09-10 全流程闭环,**T-0 ~ T-13 已全部完成,无待办**):** `PRD-风控监测Agent` FR-1 一期显式拒收 convert,本线将其放开。**AIcoding 第 1~4 步已完成**:PRD **v0.9.3 定稿**(33 条外审闭环 + 2 处架构回填 + **费率分类修正** + **展示位数补全** + **§9 第 18 条实测补录**)→ 架构 **v1.0.1 定稿**(独立评审通过,13 条建议 **0 悬空**,接受 10 / 修正性接受 3 / 驳回 0)→ 门控 **M-7 已满足** → **开发计划 v1.0 已产出(`docs/项目框架设计/开发计划-基金转换交易.md`,**1,047 行**)并经**四轮**独立子代理审核收敛(**12 条意见 → 接受 11 / 驳回 1(附实测证据)/ 0 悬空**)。**第 5 步 + 第 6 步均已完成:T-0 ~ T-13 已于 2026-09-10 全部完成(T-9 里程碑时 697 passed / 3 skipped;T-1 断言 8/8 PASS;T-2 纯函数包 7 文件 / 93 用例;T-2b 实算脚本 15/15 一致;T-6 真库 24/24;T-7 真库 35/35;T-8 真库 31/31;T-9 真 MySQL 集成 8 条 + 11 条错误码映射 + 4 组突变验证 + **展示位数口径修复**)**,**T-10 ✅(普通申赎批次维护 · 回归风险最大 · +17 用例 / 真库 20/20)、T-11 ✅(工具与 SQL 汇总去重 · +4 用例 / 真库 14/14)与 T-12 ✅(补偿脚本 · +12 用例 / 真库 34/34 · 全套 7 真库脚本复跑零回归)均已完成**,**T-13 ✅(全量回归 + 50 并发压测 + 性能补录)亦已完成** —— 压测暴露并修复 **2 处并发缺陷**(阶段一失败后同键重试**永久 503** · 同键竞态子窗口 A **真·双扣** / B 双扣+503,修法 = `insert_placeholder` 三步法 + 幂等判定 `pending`→202 + 撞键让路→202)、**50 并发不超卖**(实测成交 20~25,理论上限 25)、**紧池退避 40/50 = 80%(不够用,根因近空批次反复争抢)**、**松池 50/50 = 100%**、**跨客户并发触发 InnoDB 1213 死锁**(登记待评估)、**性能端到端 P50 44.2 / P95 51.5 / max 59.8 ms、阶段一 P50 9.8 / max 16.7 ms**(远未触阈值、未改任何实测值),实测值已回填 PRD §9 第 18 条(**v0.9.3**);**基线 736 passed / 10 skipped**,`CONVERT_STRESS=1` 并发 9/9。**本线全部任务闭环,无待办**。**开工前置两个阻断项(✅ 2026-09-10 均已完成)**:**T-0**(sqlite/MySQL 结构对齐:`core_holding` 列名+PK **+ 补 `core_product_nav`**,建库自校验,**同步改 2 处测试 INSERT 并补 3 个 NOT NULL 列**)与 **T-0b**(**DB 账号分离 D20**:`xh_core_ro` SELECT 全库 / `xh_core_rw` **4 表写**无 DELETE·DDL / `xh_agent_rw` **`audit_log` 只授 SELECT+INSERT**(不可改删);**`conftest.py` 四处 engine 已显式 `role="admin"`**)。**遗留环境操作**:`scripts/core/00-grant.sql` 需管理员执行一次 + 写 `.env`,未执行时 3 条权限断言自动 skip,不阻塞 T-1。设计资产五件套:`docs/PRD/PRD-基金转换交易.md` · `docs/项目框架设计/架构设计-基金转换交易.md` · **`docs/项目框架设计/开发计划-基金转换交易.md`(新)** · `评审待办-风控主架构与基金转换.md` · `基金转换-审查意见处置表.md`。**开工前必读开发计划 §1.4(15 条代码事实)/ §1.5(7 条实现级裁定 R-a~R-g)/ §12(16 条回归面)** —— 尤其是 **R-a(弃用方言 UPSERT)· R-b(流水写 redeem/subscribe 不写 convert)· R-e(conftest 用 admin)** 三条,不读必踩。**交接入口:项目根 `交接文档.md` §B(三线合并版唯一入口)**。合规基准 = 证监会公告〔2025〕22 号。 **仓库地图:** @@ -44,9 +44,9 @@ | `docs/项目框架设计/表设计/` | 已定 | Agent 共用 11 表 + agent 专用 SQL | | `docs/项目框架设计/Core模拟底座/` | 已定 | 无真实 Core 时的 L0 方案 | | `web/` | **不存在** | 前端 React 待 init | -| `docs/PRD/PRD-基金转换交易.md` | **已定稿(v0.9.2)** | 基金转换线需求权威(FR-C1~C16 / §4 表结构 / §5 接口 / §9 验收);⚠️ **代码未实现** | -| `docs/项目框架设计/架构设计-基金转换交易.md` | **已定稿(v1.0)** | 基金转换实现依据(D1~D20 / §5 事务 / §11.1 DB 账号 / §15 任务 T-0~T-13);⚠️ **代码未实现** | -| `docs/项目框架设计/开发计划-基金转换交易.md` | **已定稿(v1.0 · 2026-09-10)** | **基金转换实现计划**(T-0~T-13 逐任务改法 + DoD + §1.4 **十五条代码事实核对表** + §1.5 **七条实现级裁定 R-a~R-g** + §12 **十六条回归面** + 四轮审核记录 + **§3.3 第 0 批 + §4.1 第 1 批执行记录**);**T-0/T-0b/T-1 已实现**,T-2 起待实现 · **开工必读** | +| `docs/PRD/PRD-基金转换交易.md` | **已定稿(v0.9.3)** | 基金转换线需求权威(FR-C1~C16 / §4 表结构 / §5 接口 / §9 验收 · 第 18 条 T-13 实测补录);✅ **代码已实现** | +| `docs/项目框架设计/架构设计-基金转换交易.md` | **已定稿(v1.0)** | 基金转换实现依据(D1~D20 / §5 事务 / §11.1 DB 账号 / §15 任务 T-0~T-13);✅ **代码已实现** | +| `docs/项目框架设计/开发计划-基金转换交易.md` | **已定稿(v1.0 · 2026-09-10)** | **基金转换实现计划**(T-0~T-13 逐任务改法 + DoD + §1.4 **十五条代码事实核对表** + §1.5 **七条实现级裁定 R-a~R-g** + §12 **十六条回归面** + 四轮审核记录 + **§3.3 第 0 批 + §4.1 第 1 批执行记录**);**T-0 ~ T-13 已全部实现** · **开工必读** | | 项目根 `交接文档.md` | **v3.0(2026-09-10 三线合并)** | **全仓唯一交接入口**(§0 公共层 + §A 风控主线 + §B 基金转换线 + §C 架构改进线)——**给下一会话 AI,读完即可开工**;⚠️ 在 `.gitignore:47` 内、**不入库**,是本地文件。`docs/交接文档-基金转换.md` / `docs/交接文档-架构改进.md` 为**历史留档,内容已过期,勿读** | **本地 bootstrap(首次):** 完整步骤与前置说明见 `FLOW.md` §0(权威),速览: @@ -70,7 +70,7 @@ **下一步开发(见 TODO):** **模块侧交付完毕(2026-09-07:全量 pytest 482 绿 + 接口实调验收通过——suitability/check 阻断+放行、simulate/trade 阻断、三条鉴权边界 401/403/403 契约零偏差;`risk-m1~m4` tag 齐)。合并 main 已移交合并执行人,操作手册《docs/项目框架设计/合并注意事项-风控模块并入main.md》(含基底锁定/20 冲突裁决/14 静默文件/三硬伤/合并后必测,实测数据编制)。模块侧开放项:chat 链路 risk_suitability_log.actor_id 落 SYSTEM 待评估 / 前端 React 多 Agent 入口(`web/` 未 init,归属待拍板)。**演示走查按 `docs/项目框架设计/演示SOP-风控模块.md`(debug 头通道仍有效;JWT 通道签发用 `scripts/dev/issue_dev_token.py`;演示库已按 AL-08 expires_at 新口径重灌)。知识库入库:`python scripts/kb/build_kb.py`(先启 Ollama;**Milvus 数据路径必须纯英文**——faiss 不支持中文路径,本机 .env 已配 C:/Users/YUAN/.jinrong/milvus/)。 -**另(2026-09-10 待办)**:① **基金转换线**第 5 步进行中(**T-0~T-12 已完成、731 绿,下一步 = T-13(全量回归 + 50 并发压测 + 性能补录 · 本线最后一个任务)**;设计 + 开发计划均已闭环,入口 项目根 `交接文档.md` §B);② ~~架构改进线收尾~~ —— **2026-09-10 已闭环结项**:§7.2 七项手工冒烟补跑 **7/7 PASS**、冒烟残留按 SOP §2 重灌双库清除、全量 **510 passed** 复绿;**「`037ce7e` 未 push」的旧表述已作废**(实测 `git ls-remote`:远程 `risk-control-agent` = `fffb78a` = 本地 HEAD,`037ce7e` 在其祖先链上,早已推送;本地 `git branch -vv` 显示 `origin/risk-control-agent: gone` 是本仓 `refs/remotes/**` 写入不落盘所致(见交接清单第 7 问,**不是**「`git fetch` 即恢复」的普通 stale ref),非远程分支被删)。入口 **项目根 `交接文档.md` §C**。 +**另(2026-09-10 待办)**:① **基金转换线** ✅ **全流程闭环(T-0 ~ T-13 全部完成、736 绿,无待办)**;设计与开发计划均已闭环,入口 项目根 `交接文档.md` §B);② ~~架构改进线收尾~~ —— **2026-09-10 已闭环结项**:§7.2 七项手工冒烟补跑 **7/7 PASS**、冒烟残留按 SOP §2 重灌双库清除、全量 **510 passed** 复绿;**「`037ce7e` 未 push」的旧表述已作废**(实测 `git ls-remote`:远程 `risk-control-agent` = `fffb78a` = 本地 HEAD,`037ce7e` 在其祖先链上,早已推送;本地 `git branch -vv` 显示 `origin/risk-control-agent: gone` 是本仓 `refs/remotes/**` 写入不落盘所致(见交接清单第 7 问,**不是**「`git fetch` 即恢复」的普通 stale ref),非远程分支被删)。入口 **项目根 `交接文档.md` §C**。 **禁止(改代码前必记):** Core 正式 C1~C5 不可被画像覆盖 · 审计表只 INSERT · 代理人草稿不外发 · 仅 R-02 可阻断交易 · 四 Agent 不互调 LLM。 @@ -181,7 +181,7 @@ RBAC 联调账号:scripts/dev/rbac-seed-reference.md | `docs/项目框架设计/` | 表结构、JWT 手册、Core 模拟、技术版本 | | `docs/业务记忆管理/` | Redis 短期 vs MySQL/Milvus/Neo4j 权威记忆 | | 项目根 `交接文档.md` §B | **基金转换线开工入口** —— 接手本线先读这一节(不必重读代码) | -| `docs/PRD/PRD-基金转换交易.md` · `docs/项目框架设计/架构设计-基金转换交易.md` | 基金转换需求与实现依据(**v0.9.2 / v1.0.1,代码未实现**) | +| `docs/PRD/PRD-基金转换交易.md` · `docs/项目框架设计/架构设计-基金转换交易.md` | 基金转换需求与实现依据(**v0.9.3 / v1.0.1,代码已实现**) | 缺 `docs/memory/*` 文件:按 project-memory-kit 同名补回,**禁止空模板盖进度**。 @@ -199,8 +199,8 @@ RBAC 联调账号:scripts/dev/rbac-seed-reference.md 2. 改动属于 api / service / tool / repository 哪一层? 3. 是否需 customer_id 归属与 JWT RBAC? 4. Core 是模拟库只读还是 agent 库读写? -5. 如何验证?(`python -m pytest` 全量(当前 **731 passed / 3 skipped**,基线 510)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) -6. **当前有哪两条并行线?**(① 风控/架构改进线:**已结项**(510 基线绿、§7.2 七项冒烟 7/7 PASS、`037ce7e` 已核实早已推送);② **基金转换线**:设计闭环,**第 5 步进行中 —— T-0~T-12 已完成(731 passed),下一步 = T-13(最后一个任务)**)——动代码前先确认自己属于哪条线,别混淆前置条件。 +5. 如何验证?(`python -m pytest` 全量(当前 **736 passed / 10 skipped**,基线 510)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) +6. **当前有哪两条并行线?**(① 风控/架构改进线:**已结项**(510 基线绿、§7.2 七项冒烟 7/7 PASS、`037ce7e` 已核实早已推送);② **基金转换线**:设计与开发计划闭环,**全流程闭环 —— T-0~T-13 已全部完成(736 passed),无待办**)——动代码前先确认自己属于哪条线,别混淆前置条件。 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/**` **写入不落盘** —— `git update-ref refs/remotes/origin/X ` 返回 0,但松散引用消失,**且整个 `refs/remotes/origin/` 目录被删** diff --git a/docs/memory/TODO.md b/docs/memory/TODO.md index a03fc2d..994d27f 100644 --- a/docs/memory/TODO.md +++ b/docs/memory/TODO.md @@ -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》,随分支上传),后续模块侧待办见下方。 -**⚡ 并行新线 · 基金转换(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-12 均已于 2026-09-10 完成**(**731 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 新增 4 用例 + 真库 14/14 + 3 组突变验证 · `amount_view` 提升为公开**、**T-12 新增 12 用例 + 真库 34/34 + 3 组突变验证 · 补偿逻辑收敛到 `compensate_convert` 单点**),**下一步 = T-13(全量回归 + 50 并发压测 + 性能补录 · 本线最后一个任务)**。两个阻断前置(**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.3** + 架构 **v1.0.1** + 独立评审 13 条 **0 悬空**(接受 10 / 修正性接受 3 / 驳回 0),门控 **M-7 已满足**;**第 4 步开发计划 v1.0 已产出并经独立审核**(4 条意见全接受、**驳回 0**,含新增 2 条回归面 R15/R16 + R-c 双条修订);**第 5 步 + 第 6 步均已完成**:**T-0 ~ T-13 已于 2026-09-10 全部完成**(**736 passed / 10 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 新增 4 用例 + 真库 14/14 + 3 组突变验证 · `amount_view` 提升为公开**、**T-12 新增 12 用例 + 真库 34/34 + 3 组突变验证 · 补偿逻辑收敛到 `compensate_convert` 单点**),**T-13 ✅ 已完成(全量回归 + 50 并发压测 + 性能补录):压测暴露并修复 2 处并发缺陷(阶段一失败后同键重试永久 503 · 同键竞态子窗口 A 真·双扣 / B 双扣+503)、50 并发不超卖(成交 20~25)、紧池退避 40/50=80%(不够用,根因近空批次反复争抢)、松池 100%、跨客户 1213 死锁已登记、端到端 P50 44.2ms/max 59.8ms · 阶段一 P50 9.8ms/max 16.7ms**。**本线全部任务闭环,无待办**。两个阻断前置(**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 条回归面)**。 ### 基金转换线待办(推荐顺序) @@ -18,11 +18,11 @@ - [x] **【T-2 · 第 2 批】`service/convert/` 纯函数包** —— **完成(2026-09-10)**:新建 `app/service/convert/` **7 文件**(`__init__` / `types`(`Lot`/`FeeRule`/`LotAllocation`/`PlanResult` frozen dataclass + 3 个归一工具)/ `calc`(`plan_lots`/`lot_amount`/`lot_fee`/`convert_amount`/`in_qty`/`rounding_diff`/`diff_fee`/`hold_days`/`ensure_batch_limit`)/ `fee`(`pick_fee_rate` 左闭右开)/ `nav`(`ensure_nav_ready`→503 / `is_stale`)/ `lot_bootstrap`(D18 单点,`crc32` 确定性偏移)/ `errors`(`ConvertError` + 11 子类))。**新增 `tests/test_convert_calc.py` 93 用例**(12 类:精度 HALF_UP 反向自证 / 分档边界 6-7-29-30-179-180-364-365 / FIFO 含同 `confirmed_at` tiebreak / 跨批计费 / 双口径 252.40 vs 253.91 / 强制全转与强制赎回 / 恰好等于阈值不触发 / **零剩余不触发(新裁定 R-h)** / PRD §5.3 全链自证 / T+1 起算 / 净值 503 与 stale 分家 / D18 确定性 / §8.3 错误码 / **纯函数零 IO 依赖断言**)。**验证**:`pytest -q` → **609 passed / 3 skipped(+93,零回归)**;`calc_convert_demo.py` → **15/15 与 PRD §5.3 一致**(退出码 0) - [x] **【T-3 ~ T-9 已完成】** 仓储与锁(并行组 A:T-3 / T-4 / T-5 ✅)→ T-6 ✅(阶段一事务 · 真库 24/24)→ T-7 ✅(八步编排 · 真库 35/35)→ **T-8 ✅(规则引擎改造 · 真库 31/31;⭐ 阶段 1.5 从「跳过」变「真跑」)** → **T-9 ✅(API 模型 + 网关分派 · **HTTP 层 convert 已走通**;新增真 MySQL 集成 7 条 + 11 条错误码映射 + 3 组突变验证)** - [x] **【T-10 ✅ · T-11 ✅ · T-12 ✅】均已结项(2026-09-10)**:**T-10** 普通申赎批次维护(改 `trade_gateway` 主流程 + D8 兜底补建 + `rebuild_lots.py`;**+17 用例** · 真库 20/20 · 突变 2 组 · R16 零改动通过)→ **T-11** 工具与 SQL 汇总去重(`_amount_view` **提升为公开 `amount_view`** + `query_recent_trades` 汇总去重 FR-C15 + `list_holdings` 过滤 `qty > 0` + `sum_trades_on_date` 加 convert 去重 **R-d**;**+4 用例** · 真库 14/14 · 突变 3 组)→ **T-12** 补偿脚本(公开 **`compensate_convert`** 单点入口 + `rebuild_alerts --convert-group` 薄封装 + 新增 `scripts/agent/cleanup_pending_convert.py`;**+12 用例** · 真库 34/34 · 突变 3 组 · 全套 7 真库脚本复跑零回归;顺带收口 `concentration_profile` 的 `qty > 0`) -- [ ] **【T-13】**(最后一个):**全量回归 → 50 并发压测 → 性能实测补录 → PRD §5.3 / §9 第 18 条数字回填**。**T-13 内部顺序**:先 50 并发压测 → 再性能实测补录 → 最后 PRD 数字回填(详见架构 §15 + 开发计划 §2~§10)。**DoD**:① 性能实测值**回填 PRD §9 第 18 条**(替换「预估 < 100ms」);② 50 并发结论(是否超卖 / `100/200/400ms` 退避间隔够不够)写入 `交接文档.md` §B;③ 实测超阈值 → 优化索引/锁策略后**重定阈值**,**不得反向修改实测数据迁就指标** +- [x] **【T-13 ✅】**(最后一个 · **已完成 2026-09-10**):**全量回归 → 50 并发压测 → 性能实测补录 → PRD §9 第 18 条数字回填**。**DoD 三条**:① 性能实测值**已回填** PRD §9 第 18 条(v0.9.3,替换「预估 < 100ms」);② 50 并发结论(不超卖 / `100/200/400ms` 退避**不够用**)**已写入** `交接文档.md` §B.6.6;③ 实测**未超阈值** → 无需重定阈值(**未修改实测数据**)。**产出**:`tests/test_convert_concurrency.py`(9 条 · `CONVERT_STRESS` 门禁)+ `test_convert_service` +3 + 修复 `insert_placeholder` 三步法 / 幂等判定 `pending`→202;**736 passed / 10 skipped**、`CONVERT_STRESS=1` 并发 **9/9**、重灌双库复跑零回归、4 组突变验证全命中。 - [ ] **【待用户裁定】幂等重放响应的数值位数偏差**(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)。 -> **下一步 = T-13**(全量回归 + 50 并发压测 + 性能补录 **· 本线最后一个任务**)。基线 **731 passed / 3 skipped**(T-10 后 714 → T-11 后 718 → **T-12 后 731**,+13)。 +> **本线已全部闭环(T-0 ~ T-13)**。基线 **736 passed / 10 skipped**(T-10 后 714 → T-11 后 718 → T-12 后 731 → **T-13 后 736**,+5);`CONVERT_STRESS=1` 另跑 9 条并发用例全绿。 > **T-11 执行期两条裁定**:① `_amount_view` **提升为公开 `amount_view`**(跨层复用,全仓唯一金额聚合口径); > ② SQL 去重条件从计划的 `IS NULL` 扩为 **`IS NULL OR convert_group_id = ''`** —— 真库实证「只写 `IS NULL` 会漏掉空串那笔」(400000 vs 450000,差 50000)。 > **T-12 执行期三条裁定**:① 补偿逻辑收敛到 **`convert_service.compensate_convert` 单点**,`rebuild_alerts --convert-group` 只是薄封装; diff --git a/docs/项目框架设计/开发计划-基金转换交易.md b/docs/项目框架设计/开发计划-基金转换交易.md index 0018a5d..3488f29 100644 --- a/docs/项目框架设计/开发计划-基金转换交易.md +++ b/docs/项目框架设计/开发计划-基金转换交易.md @@ -125,16 +125,19 @@ T-7 幂等窗口 · T-13 的 50 并发压测与性能补录 · PRD §5.3 实算 | | **T-7** ✅ | `convert_service` 编排(八步 + 执行权 + 幂等 + 三阶段 + 阶段 1.5)—— **17 用例 + 真库 35/35** | T-2~T-6 | **高(关键路径)** | | **第 4 批 · 引擎与网关** | **T-8** ✅ | `_amount_view` + `engine.process_convert_event` + `alert_service.events` —— **2026-09-10 完成(15 用例 + 真库 31/31)** | 无(可与 T-2 并行) | 中 | | | **T-9** ✅ | `api/simulate.py` 模型与错误码 + `trade_gateway` convert 分派 + **展示位数口径修复** —— **2026-09-10 完成(新增集成 8 条 / 11 条错误码映射)** | T-7 | 中 | -| | T-11 | `core_tools` 汇总去重 + 持仓 `qty <= 0` 过滤 + `sum_trades_on_date` 去重 | T-8 | 中 | -| **第 5 批 · 高风险专项** | **T-10** | 普通申赎批次维护(FR-C16,含 D8 兜底补建)+ `rebuild_lots.py` | T-3(排在 T-7 后) | **最高(打穿 510)** | -| **第 6 批 · 补偿** | T-12 | `rebuild_alerts --convert-group` + `cleanup_pending_convert.py` | T-4/T-7 | 低 | -| **第 7 批 · 收口** | T-13 | 全量回归 + 集成测试 + 50 并发压测 + 性能实测补录 | 全部 | 中 | +| | **T-11** ✅ | `core_tools` 汇总去重 + 持仓 `qty <= 0` 过滤 + `sum_trades_on_date` 去重 —— **2026-09-10 完成(+4 用例 · 真库 14/14)** | T-8 | 中 | +| **第 5 批 · 高风险专项** | **T-10** ✅ | 普通申赎批次维护(FR-C16,含 D8 兜底补建)+ `rebuild_lots.py` —— **2026-09-10 完成(+17 用例 · 真库 20/20)** | T-3(排在 T-7 后) | **最高(打穿 510)** | +| **第 6 批 · 补偿** | **T-12** ✅ | `rebuild_alerts --convert-group` + `cleanup_pending_convert.py` —— **2026-09-10 完成(+12 用例 · 真库 34/34)** | T-4/T-7 | 低 | +| **第 7 批 · 收口** | **T-13** ✅ | 全量回归 + 集成测试 + 50 并发压测 + 性能实测补录 —— **2026-09-10 完成(+5 用例 · 真库并发 9/9 · 突变 4 组 · 实测补录 PRD v0.9.3)** | 全部 | 中 | -**关键路径**:`T-0 → T-1 → T-2 → T-6 → T-7 → T-13`(**T-7 已通;T-9 已完成**) +**关键路径**:`T-0 → T-1 → T-2 → T-6 → T-7 → T-13`(**全程已通,T-13 收口完成**) **并行组 A**:T-3 / T-4 / T-5(✅ 全部完成) **并行组 B**:T-8 全程可与 T-2 之后任意任务并行(✅ 已完成) **硬门禁**:`T-0` 与 `T-0b` **双双绿**才允许启动 T-1 及之后(T-0 用例 = `test_db.py::test_core_holding_columns`) -**测试基线**:**697**(2026-09-10 T-9 后;批 0~3 路线 510 → 516 → 609 → 634 → 639 → 656 → 672 → **697**)→ 剩余任务(T-10~T-13)预计再加 **10~40** → **707~737**(估算) +**测试基线**:**736 passed / 10 skipped**(2026-09-10 T-13 后;批 0~3 路线 510 → 516 → 609 → 634 → 639 → 656 → 672 → 697 → 714 → 718 → 719 → 731 → **736**) + +> **✅ T-0 ~ T-13 全部结项(2026-09-10)。** 本线第 5 步(todo_list 逐步开发)走完; +> 剩余仅为「第 6 步:最终集成测试」的独立复跑(SOP 重灌后全链路已验证,见 §10.1)。 --- @@ -1441,25 +1444,72 @@ sqlite 无 gap lock,故该分支由 `tests/test_convert_core.py` 用注入点 ## 10. 第 7 批 · 回归与实测(T-13) ### 10.1 全量回归 -- `pytest -q` 全绿,计数 ≥ 585(CI 内预计约 597,压测 10 条不进门禁) -- 集成测试(真 MySQL,`CNV-TEST-`/`TRD-TEST-` 前缀)绿 -- 按 SOP 重灌双库后复跑一次(避免残留数据污染) +- ✅ `pytest -q` → **736 passed / 10 skipped**(基线 731 **+5**;10 skip = 7 条压测门禁 + 3 条 D20 账号未配) +- ✅ 真库集成测试(`test_convert_integration.py`,`CNV-TEST-`/`TRD-TEST-` 前缀)绿 +- ✅ **按 SOP 重灌双库后复跑一次**:`DROP jinrong_core` → `scripts/core/00~09*.sql` 全 OK → + `scripts/demo/prepare_risk_demo.sql` → `scripts/sync/sync_advisor_rel.py`(33 行 upsert)→ 复跑仍 **736/10** + (Neo4j 未启动,`sync_neo4j.py` 按 reset.ps1 的 `-SkipNeo4j` 语义跳过) +- ✅ 全部 convert 真库脚本复跑零回归:seed 全 PASS / apply 24 / service 35 / engine 31 / lots 20 / tools 14 / compensate 34 ### 10.2 50 并发压测(评审 Q6 · **不进 CI 门禁**) -同一 `(customer_id, product_id)` 上 50 并发争抢同一批份额,断言三件事: -1. `LotConflict`(409) 命中数与剩余可转份额一致 → **不许超卖** -2. 按 `100/200/400ms` 退避重试 ≤3 次后的**最终成功率**(验证 §8.3 建议间隔是否够) -3. `core_trade` 中 `convert_group_id` **无重复**;`core_share_lot.remain_qty` 之和 = 初始值 − 实际成交份额 +新建 `tests/test_convert_concurrency.py`(9 条)。门禁语义用环境变量实现: +常规 `pytest -q` 下 7 条重载用例逐条 skip;`CONVERT_STRESS=1` 显式开启全套。 +真库缺失时整模块 skip(模块级 `ensure_risk_demo_ready()`)。 + +同 `(customer_id, product_id)` 上 50 并发争抢同一批份额,三条断言: +1. **不许超卖** ✅ 硬不变量成立:50 × 2000 争 50000 → 扣减恒 ≤ 50000、余 = 池 − 扣、 + 每组恰好 2 条流水;**实测成交 20~25(理论上限 25)**——浮动原因是 + `plan_lots` 会把一笔跨批次拆成多笔扣减、`_deduct_lots` 要求本次全部成功否则整体回滚 +2. **退避重试最终成功率** ✅ 已量化(**结论见下**) +3. `convert_group_id` 无重复、`remain_qty` 守恒 ✅ + +**退避间隔是否够用的结论(评审 Q6 的答案)** + +| 场景 | 成功率 | 结论 | +| --- | --- | --- | +| **紧池**(需求 50000 = 供给 50000,最紧张) | **40/50 = 80%** | **3 次 × 100/200/400ms 不够**。根因不是"间隔太短",而是**近空批次的反复争抢**:一笔需求被拆成 `400+600` 这类碎片后,碎片所在批次随时被别人清零 → 重试仍可能抢不到。要收敛到 100% 需**增加重试次数**或**冲突后换批次重规划**,单纯拉长间隔无效 | +| **松池**(需求 25000 < 供给 50000,对照) | **50/50 = 100%** | 非耗尽场景下哨兵 + 退避链路完全可靠 | + +> 24 小时口径外的补充发现:**跨客户并发会触发 InnoDB 1213 死锁** +> (不同客户在 `core_holding`/`core_share_lot` 唯一索引上的插入意向锁互斥), +> 8 路并发实测累计 3~5 次;`apply_convert` **不重试死锁**,按 §8.3 阶梯把它当可重试异常 +> 收敛后 8/8 成功。**"服务端不重试死锁"已登记为待评估项**(见交接文档 §B)。 ### 10.3 性能实测补录(PRD §9 第 18 条) -- 端到端响应 **< 2s**(本地模拟库) -- 阶段一单库事务实测耗时**补录真实值**(PRD 原文的「< 100ms」是**预估值、非验收硬指标**) -- **若实测超阈值 → 优化索引/锁策略后重定阈值;不得反向修改实测数据迁就指标** +- ✅ 端到端 **P50 44.2ms / P95 51.5ms / max 59.8ms** → 远低于 **2s** 硬指标 +- ✅ 阶段一单库事务实测 **P50 9.8ms / P95 12.8ms / max 16.7ms** → 远低于「预估 < 100ms」 +- ✅ **未触发**"超阈值 → 优化索引/锁策略"分支;**未修改任何实测值**;数据已回填 PRD §9 第 18 条(v0.9.3) **DoD** -- [ ] 三项完成且**落点明确**:① 性能实测数据**回填 `docs/PRD/PRD-基金转换交易.md` §9 第 18 条**(用实测值替换「预估 < 100ms」); - ② 50 并发结论(是否超卖 / 退避间隔够不够)写入 项目根 `交接文档.md` §B;③ 若实测超阈值 → 优化索引/锁策略后**重定阈值**(**不得反向修改实测数据迁就指标**) -- [ ] 50 并发结果写入交接文档(含退避间隔是否够用的结论) +- [x] 三项完成且**落点明确**:① 性能实测数据**已回填** `docs/PRD/PRD-基金转换交易.md` §9 第 18 条(v0.9.3,实测值替换预估值); + ② 50 并发结论(不超卖 / 退避 80% 不足)写入 项目根 `交接文档.md` §B;③ 实测**未超阈值** → 无需重定阈值(**未修改实测数据**) +- [x] 50 并发结果写入交接文档(含退避间隔**不够用**及根因、1213 死锁登记) + +### 10.4 执行记录(2026-09-10 · 已落地) + +**执行期发现并修复的 2 处缺陷**(均为压测暴露,非需求变更;详见 §10.5) + +| # | 缺陷 | 影响面 | 修法 | +| --- | --- | --- | --- | +| 1 | 阶段一失败(`LotConflict` 409)后带**同一** `client_request_id` 重试 → 幂等占位撞 `uk_group`/`uk_idem` → `IdempotencyUnavailable`(503) | **确定性永久失败**:§8.3 明说 409 可重试,此路径**重试永不成功** —— T-13 退避成功率实测的直接前置阻断 | `ConvertRepository.insert_placeholder` 改**三步法**(R-a):已存在则置回 `pending` | +| 2 | 同键并发的**两个竞态子窗口**:
· 子窗口 A(T2 在占位后进入)→ 各自 `group_id` 不同 → 旧代码重跑阶段一 → **真·双扣**(突变实测 120→**240**、两组流水)
· 子窗口 B(T2 在占位前进入)→ 撞 `uk_idem` → `IntegrityError` **直穿 503** | **资金安全**:A 会重复扣客户份额 | ① 幂等判定区分「确定没跑成」(`failed`/`expired` → 复用 group_id 重跑)与「在飞/未知」(`pending` → **202**);② `insert_placeholder` 返回"是否持有占位",撞键即让路 → 调用方 **202**(架构 §9「同键并发 → 202」) | + +**验证** +- `pytest -q` → **736 passed / 10 skipped**(+5,零回归;跑两遍稳定) +- `CONVERT_STRESS=1 pytest tests/test_convert_concurrency.py` → **9 passed** +- **突变验证 4 组,全部精准命中**(见 §10.5) +- **重灌双库后复跑**:736/10 + 并发 9/9,零回归 + +### 10.5 突变验证明细(防假绿 · 逐组留痕) + +| 组 | 突变 | 命中用例 | 实测失效形态 | +| --- | --- | --- | --- | +| ① | 拆掉 `_deduct_lots` 的并发哨兵 `remain_qty >= :q` | `test_50_concurrent_requests_never_oversell`(1 红) | 成交 **50/50**、余 **−50000**、扣减 **100000** → 超卖被抓 | +| ② | 幂等判定把 `pending` 也当"未成"重跑 | `test_same_client_request_id_interleaving_must_not_leak_5xx`(1 红) | T1 直穿 `IntegrityError` | +| ③ | 忽略 `insert_placeholder` 的让路信号 | `test_..._placeholder_race_must_not_leak_5xx`(并发文件,1 红)**+** `test_placeholder_uk_idem_race_returns_processing_without_writing`(sqlite 版,1 红) | **真·双扣:扣 240 份、两组流水** | +| ④ | `insert_placeholder` 退回朴素 INSERT | `test_retry_after_phase_one_failure_*`(2 红) | `IdempotencyUnavailable`:`UNIQUE constraint failed` | + +> 已恢复全部突变,`grep -rn "MUTATION\|1 = 1" app/ tests/ scripts/` **无残留**。 --- diff --git a/tests/test_convert_concurrency.py b/tests/test_convert_concurrency.py new file mode 100644 index 0000000..3fb6bdc --- /dev/null +++ b/tests/test_convert_concurrency.py @@ -0,0 +1,900 @@ +"""T-13 真 MySQL 并发测试(架构 §9 测试清单 · §10「50 并发压测口径(评审 Q6)」)。 + +**为什么必须是真 MySQL**:本文件验的是 InnoDB **行锁 + 条件 UPDATE rowcount** 这套 +数据库级并发原语。sqlite 内存库(`conftest.sqlite_engine`)共用单连接,多线程跑不出 +真争抢——在那里"绿"是**假绿**(与本项目「方言语义一律上真库」的既有铁律一致)。 + +**为什么不算 CI 常规门禁**(架构 §10 原话:该用例不进 CI,归入 T-13 一并跑): +重载用例要 50 线程 + 真库,耗时长。本文件用环境变量 **`CONVERT_STRESS=1`** 实现该语义: +未设置时重载用例逐条 skip(给出可操作的原因),常规 `pytest -q` 不受拖累。 +T-13 实测:`CONVERT_STRESS=1 pytest tests/test_convert_concurrency.py -q -s`。 + +**隔离策略**(与 `test_convert_integration.py` 同款,三条同时成立,缺一即污染种子): +1. id 前缀:`CNV-CONC-` / `TRD-CONC-`(monkeypatch `convert_service._new_id`); +2. 数据自建:客户 `CUST-CONCTEST` / 产品 `PROD-CONCTEST?` 全自建,**绝不碰种子**—— + `risk_demo_env` 的 teardown 不还原 `core_share_lot`/`core_holding`,借种子客户跑转换 + 会跨用例污染 `test_integration_risk.py`; +3. teardown 全清:函数级 fixture 按客户 + 前缀删两库全部自建行(幂等)。 + +**连接池**:并发用例用**专用 Engine**(`pool_size=60`)。默认池 `5+10` 会把 50 个线程 +压成 15 路串行,削弱争抢强度、让"不超卖"变成弱证据;专用池只在压测用,常规用例仍走 +`get_engine` 生产口径。 +""" + +from __future__ import annotations + +import os +import threading +import time +from concurrent.futures import ThreadPoolExecutor +from datetime import date, datetime, time as dtime, timedelta +from decimal import Decimal +from statistics import mean +from uuid import uuid4 + +import pytest +from sqlalchemy import create_engine, text +from sqlalchemy.exc import OperationalError + +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.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.errors import LotConflict # noqa: E402 +from app.service.risk.rules import RiskThresholds # noqa: E402 +from app.utils.db import _resolve_credentials # noqa: E402 + +ensure_risk_demo_ready() + +# ── 隔离种子常量(全部带 CONCTEST 前缀)────────────────────────────── +CUSTOMER = "CUST-CONCTEST" +PROD_OUT = "PROD-CONCTESTO" # 转出:债基,申购费率 0.0030 +PROD_IN = "PROD-CONCTESTI" # 转入:股基,申购费率 0.0080 +COMPANY = "华夏模拟基金" +TA = "TA-CN-001" +OUT_RATE = Decimal("0.0030") +IN_RATE = Decimal("0.0080") +OUT_NAV = Decimal("1.0300") +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"), +] +_LOT_DAYS = 100 # 全部批次同一费率档(0.0050)→ 断言只看份额守恒,不被费用档干扰 + +_GROUP_LIKE = "CNV-CONC-%" + +#: 架构 §8.3 建议的退避间隔(`LOT_CONFLICT` → 调用方重试 ≤3 次) +BACKOFF_MS = (100, 200, 400) +#: 可重试异常:409 是明确设计为可重试的;死锁/锁等待超时同属瞬时状态 +_RETRYABLE = (LotConflict, OperationalError) + + +def _stress_enabled() -> bool: + return os.environ.get("CONVERT_STRESS") == "1" + + +requires_stress = pytest.mark.skipif( + not _stress_enabled(), + reason="重载并发用例不进 CI 常规门禁(架构 §10);设 CONVERT_STRESS=1 显式开启", +) + + +def _noop_engine(out_trade: dict, in_trade: dict) -> dict: + """阶段 1.5 的空引擎:并发用例只验**份额争抢**,不掺引擎/出单噪声。""" + return {"triggered_rules": [], "alert_ids": [], "aml_hit": False} + + +def _new_id(prefix: str, now: datetime) -> str: + """`_new_id` 替换:id 带 `-CONC-` 段(与 `CNV-CONC-%` 清理前缀一致)。""" + return f"{prefix}-CONC-{uuid4().hex[:10].upper()}" + + +# ── 专用 Engine(放大连接池)───────────────────────────────────────── +def _big_pool_engine(database: str, role: str, size: int = 60): + """建一个连接池放大的 Engine(与 `get_engine` 同 URL/凭据,仅池参数不同)。""" + user, password = _resolve_credentials(database, role) + auth = f"{user}:{password}" if password else user + url = ( + f"mysql+pymysql://{auth}@{settings.mysql_host}:{settings.mysql_port}" + f"/{database}?charset=utf8mb4" + ) + return create_engine(url, pool_pre_ping=True, pool_size=size, max_overflow=size) + + +# ── 种子与清理 ────────────────────────────────────────────────────── +def _seed(core, lots: list[str]) -> Decimal: + """自建隔离种子;`lots` = 各批次份额。返回份额池总量。""" + today = date.today() + base = datetime.combine(today, dtime(10, 0)) + total = sum(Decimal(q) for q in lots) + with core.begin() as conn: + conn.execute( + text( + "INSERT INTO core_customer (customer_id, display_name, open_date)" + " VALUES (:c, 'T13并发测试', :d)" + ), + {"c": CUSTOMER, "d": today}, + ) + conn.execute( + text( + "INSERT INTO core_customer_risk (customer_id, risk_code, evaluated_at, expires_at)" + " VALUES (:c, 'C5', :t, :exp)" + ), + { + "c": CUSTOMER, + "t": base - timedelta(days=30), + "exp": base + timedelta(days=300), + }, + ) + for pid, name, ptype, rate in [ + (PROD_OUT, "T13转出基金", "bond", OUT_RATE), + (PROD_IN, "T13转入基金", "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," + " min_redeem_qty, min_hold_qty, min_hold_action, fund_company, ta_code)" + " VALUES (:p, :n, 'R2', :t, 1, 1, :r, 0, 0, 'force_transfer', :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, :mm, :r)" + ), + {"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 i, qty in enumerate(lots, start=1): + 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, :n, :cat)" + ), + { + "l": f"LOT-CONCTEST-{i}", + "c": CUSTOMER, + "p": PROD_OUT, + "q": qty, + "n": str(OUT_NAV), + "cat": base - timedelta(days=_LOT_DAYS), + }, + ) + 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, :mv, 0, :d)" + ), + { + "c": CUSTOMER, + "p": PROD_OUT, + "q": str(total), + "cost": str(total * OUT_NAV), + "mv": str(total * OUT_NAV), + "d": today, + }, + ) + return total + + +_CUST_LIKE = "CUST-CONCTEST%" +_PROD_LIKE = "PROD-CONCTEST%" +#: 删除顺序按真库 FK 依赖(`information_schema.KEY_COLUMN_USAGE` 实测): +#: `core_holding`/`core_share_lot`/`core_trade` 同时引用 `core_customer` 与 +#: `core_product`,故必须先删它们,再删 `core_product`,最后删 `core_customer`。 +_CORE_CLEANUP_BY_CUSTOMER = [ + "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", + "DELETE FROM core_cash_flow WHERE customer_id LIKE :cl", + "DELETE FROM core_customer_advisor WHERE customer_id LIKE :cl", + "DELETE FROM core_customer_risk WHERE customer_id LIKE :cl", +] +_CORE_CLEANUP_BY_PRODUCT = [ + f"DELETE FROM core_fee_rule WHERE product_id LIKE '{_PROD_LIKE}'", + f"DELETE FROM core_product_nav WHERE product_id LIKE '{_PROD_LIKE}'", + f"DELETE FROM core_product WHERE product_id LIKE '{_PROD_LIKE}'", +] + + +def _cleanup(core, agent) -> None: + """按客户 + 前缀清两库全部自建行(幂等;seed 失败亦可清)。""" + params = {"cl": _CUST_LIKE} + with core.begin() as conn: + conn.execute( + text( + "DELETE FROM core_convert_lot_detail" + f" WHERE convert_group_id LIKE '{_GROUP_LIKE}'" + ) + ) + for sql in _CORE_CLEANUP_BY_CUSTOMER: + conn.execute(text(sql), params) + for sql in _CORE_CLEANUP_BY_PRODUCT: + conn.execute(text(sql)) + conn.execute(text("DELETE FROM core_customer WHERE customer_id LIKE :cl"), params) + with agent.begin() as conn: + conn.execute( + text(f"DELETE FROM risk_convert_detail WHERE convert_group_id LIKE '{_GROUP_LIKE}'") + ) + for table in ("risk_alert", "risk_suitability_log", "audit_log"): + conn.execute(text(f"DELETE FROM {table} WHERE customer_id LIKE :cl"), params) + + +@pytest.fixture() +def conc_env(risk_demo_env, monkeypatch): + """真库并发环境:admin 引擎(teardown 需 DELETE)+ 专用大池引擎 + id 前缀注入。 + + 种子由用例自己 `_seed`(不同用例份额池不同),fixture 只负责清场与收尾。 + """ + core_admin = risk_demo_env["core"] + agent_admin = risk_demo_env["agent"] + core_ro = _big_pool_engine(settings.mysql_core_database, "ro") + core_rw = _big_pool_engine(settings.mysql_core_database, "rw") + agent_rw = _big_pool_engine(settings.mysql_database, "rw") + monkeypatch.setattr(cs, "_new_id", _new_id) + _cleanup(core_admin, agent_admin) + try: + yield { + "core": core_admin, + "agent": agent_admin, + "core_ro": core_ro, + "core_rw": core_rw, + "agent_rw": agent_rw, + } + finally: + _cleanup(core_admin, agent_admin) + for eng in (core_ro, core_rw, agent_rw): + eng.dispose() + + +# ── 调用与度量 ────────────────────────────────────────────────────── +_NOW = datetime.combine(date.today(), dtime(10, 0)) + + +def _convert(env, qty: str, cid: str | None, hook=_noop_engine): + """一次转换(生产账号、生产口径;仅 id 工厂与引擎 hook 被替换)。""" + return convert_fund( + { + "customer_id": CUSTOMER, + "from_product_id": PROD_OUT, + "to_product_id": PROD_IN, + "qty": qty, + "client_request_id": cid, + }, + core_ro=CoreReadOnlyRepository(engine=env["core_ro"]), + risk_repo=RiskRepository(engine=env["agent_rw"]), + convert_repo=ConvertRepository(engine=env["agent_rw"]), + core_writer=ConvertCoreRepository(engine=env["core_rw"]), + thresholds=RiskThresholds.from_settings(), + now=_NOW, + engine_hook=hook, + ) + + +class Outcome: + """一个并发请求的终态(是否成功、重试了几次、最终异常)。""" + + __slots__ = ("ok", "retries", "error", "processing", "group_id") + + def __init__(self, ok, retries, error=None, processing=False, 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` 全程不变。""" + attempt = 0 + while True: + try: + resp = _convert(env, qty, cid) + 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__) + if resp.get("status") == PROCESSING: + return Outcome(False, attempt, "PROCESSING", processing=True) + return Outcome(True, attempt, group_id=resp["convert_group_id"]) + + +def _concurrently(fn, n: int, *, stagger_ms: float = 0.0) -> list: + """n 个线程同时起跑;`stagger_ms` 为逐线程错开步长(0 = 全同时)。 + + 起跑线用 `threading.Barrier` 对齐,避免"线程创建开销"把并发抹平。 + """ + barrier = threading.Barrier(n) + + def _worker(i: int): + barrier.wait() + if stagger_ms: + time.sleep(stagger_ms * i / 1000.0) + return fn(i) + + with ThreadPoolExecutor(max_workers=n) as pool: + return list(pool.map(_worker, range(n))) + + +# ── 度量(全部直接读库,不看响应自证)──────────────────────────────── +def _remain_total(core) -> Decimal: + with core.connect() as conn: + return Decimal( + str( + conn.execute( + text( + "SELECT COALESCE(SUM(remain_qty), 0) FROM core_share_lot" + " WHERE customer_id = :c AND product_id = :p" + ), + {"c": CUSTOMER, "p": PROD_OUT}, + ).scalar_one() + ) + ) + + +def _deducted_total(core) -> Decimal: + """实际成交的转出份额 = 全部 `redeem` 流水的 qty 之和(转出端=客户被扣的部分)。""" + with core.connect() as conn: + return Decimal( + str( + conn.execute( + text( + "SELECT COALESCE(SUM(qty), 0) FROM core_trade" + " WHERE customer_id = :c AND trade_type = 'redeem'" + ), + {"c": CUSTOMER}, + ).scalar_one() + ) + ) + + +def _trades(core) -> list[dict]: + with core.connect() as conn: + return [ + dict(r) + for r in conn.execute( + text("SELECT * FROM core_trade WHERE customer_id = :c ORDER BY trade_id"), + {"c": CUSTOMER}, + ).mappings() + ] + + +def _group_counts(core) -> dict[str, int]: + """convert_group_id → 该组流水条数(正常恒为 2)。""" + counts: dict[str, int] = {} + for row in _trades(core): + gid = row["convert_group_id"] + counts[gid] = counts.get(gid, 0) + 1 + return counts + + +def _assert_no_oversell(core, pool: Decimal, ok_count: int) -> 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)→ **不许超卖**。 + + 典型失效形态:若 `_deduct_lots` 的条件 UPDATE 丢掉 `remain_qty >= :q` + (并发哨兵),50 笔会**全部**成交 → 扣减 100000 > 池 50000 → 本用例红。 + + ⚠️ **为什么成交笔数不是恒等于 25(理论上限)**——实测在 **20~25** 间浮动: + `plan_lots` 会把一笔需求跨批次拆成多笔扣减(如 `400 + 1600`),而 + `_deduct_lots` 要求**本次全部扣减都成功**,否则整个阶段一事务回滚 → + 只要其中任一批次被别人清零,这一笔就整体 `LotConflict` 并进入重试; + 重试窗口重叠时彼此让路,落点因此不定。**不超卖是硬不变量,成交笔数是观测量。** + """ + 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"]) + n = 8 + cid = f"CONC-SAME-{uuid4().hex[:8]}" + + def _one(_i: int) -> Outcome: + try: + resp = _convert(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"]) + + 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) + print( + f"\n[同键竞态·确定性交错] T1={got.get('t1_err') or 'ok'} " + f"T2={got.get('t2_err') or 'ok'} · 扣减 {deducted} · 组→流水数 {counts}" + ) + # ① 资金安全:只扣一次(当前由 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}" + + +# ── 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 的并发面)──────────────── +def test_concurrent_compensation_writes_detail_once(conc_env, monkeypatch): + """阶段二失败后并发补偿(`convert:rerun:` 锁)→ 详情只补一次、只有一组流水。""" + from app.service.convert.convert_service import compensate_convert + + core = conc_env["core"] + agent = conc_env["agent"] + _seed(core, ["50000"]) + + original = ConvertRepository.complete_convert + + def boom(self, group_id, **kwargs): # noqa: ANN001 + raise RuntimeError("阶段二写失败(制造待补偿态)") + + monkeypatch.setattr(ConvertRepository, "complete_convert", boom) + gid = _convert(conc_env, "120", f"CONC-CMP-{uuid4().hex[:8]}")["convert_group_id"] + monkeypatch.setattr(ConvertRepository, "complete_convert", original) + + 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"] + + states = _concurrently(_comp, 6, stagger_ms=1.0) + print(f"\n[并发补偿 6 路] 状态分布 {states}") + assert states.count("rebuilt") == 1, f"应恰好 1 路真补写,实际 {states}" + with agent.connect() as conn: + rows = [ + dict(r) + for r in conn.execute( + text("SELECT status FROM risk_convert_detail WHERE convert_group_id = :g"), + {"g": gid}, + ).mappings() + ] + assert rows == [{"status": "completed"}], f"详情应只有一行 completed:{rows}" + assert len(_group_counts(core)) == 1 and len(_trades(core)) == 2 + + +# ── 5. 多客户并发:互不阻塞,但会出 1213 死锁(须靠退避重试收敛)───────── +@requires_stress +def test_concurrent_different_customers_need_deadlock_retry(conc_env): + """8 个客户并发转换同一转出产品 → 全部成功,但**必须靠重试死锁**。 + + **实测发现(2026-09-10)**:跨客户并发在 `core_holding`/`core_share_lot` 的 + 唯一索引上会触发 InnoDB **1213 死锁**(RR 隔离级下 INSERT/UPDATE 的 + 插入意向锁与间隙锁互斥);`apply_convert` **不重试死锁**,死锁以未映射的 + `OperationalError` 直穿(HTTP 500)。本用例按 §8.3 的退避阶梯把死锁当 + 可重试异常收敛,验证「阶梯对死锁同样有效」。 + + 注意与「50 并发不超卖」的区别:那批是**同一** (客户, 产品),走的是行锁阻塞; + 这批是**不同**客户,走的是间隙锁 → 死锁而非阻塞。 + """ + core = conc_env["core"] + _seed(core, ["5000"] * 2) + + others = [f"CUST-CONCTEST-{i}" for i in range(1, 8)] + with core.begin() as conn: + for c in others: + conn.execute( + text( + "INSERT INTO core_customer (customer_id, display_name, open_date)" + " VALUES (:c, 'T13并发客户', CURDATE())" + ), + {"c": c}, + ) + conn.execute( + text( + "INSERT INTO core_customer_risk (customer_id, risk_code, evaluated_at," + " expires_at) VALUES (:c, 'C5', NOW(), DATE_ADD(NOW(), INTERVAL 300 DAY))" + ), + {"c": c}, + ) + conn.execute( + text( + "INSERT INTO core_share_lot (lot_id, customer_id, product_id, qty," + " remain_qty, nav, confirmed_at)" + " VALUES (:l, :c, :p, 5000, 5000, :n, DATE_SUB(NOW(), INTERVAL 100 DAY))" + ), + {"l": f"LOT-CONCTEST-OTH-{c}", "c": c, "p": PROD_OUT, "n": str(OUT_NAV)}, + ) + + customers = [CUSTOMER, *others] + + 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"]) + + 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"失败 {[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]}" + with core.connect() as conn: + total = Decimal( + str( + conn.execute( + text( + "SELECT COALESCE(SUM(remain_qty), 0) FROM core_share_lot" + " WHERE product_id = :p AND lot_id LIKE 'LOT-CONCTEST%'" + ), + {"p": PROD_OUT}, + ).scalar_one() + ) + ) + # 8 个客户各转出 1000;主客户有 2 批 × 5000 = 10000,其余 7 个各 5000 + # → 初始合计 45000,扣 8000 → 37000 + expected = Decimal("45000") - Decimal("1000") * len(customers) + assert total == expected, f"剩余合计应为 {expected},实际 {total}" + assert len({o.group_id for o in ok}) == 8, "8 个客户的 group_id 必须两两不同" + + +# ── 6. 性能实测(PRD §9 第 18 条)──────────────────────────────────── +class _TimingWriter: + """`ConvertCoreRepository` 代理:只测**阶段一单库事务**耗时,不改任何行为。""" + + def __init__(self, inner): + self._inner = inner + self.samples: list[float] = [] + + def apply_convert(self, req): + t0 = time.perf_counter() + try: + return self._inner.apply_convert(req) + finally: + self.samples.append((time.perf_counter() - t0) * 1000.0) + + +def _pct(values: list[float], p: float) -> float: + ordered = sorted(values) + idx = min(len(ordered) - 1, max(0, int(round((p / 100.0) * len(ordered) + 0.5)) - 1)) + return ordered[idx] + + +@requires_stress +def test_performance_probe_stage1_and_end_to_end(conc_env): + """实测:阶段一单库事务 + 端到端一次转换(PRD §9 第 18 条补录来源)。 + + 测量条件:本机 MySQL 8.0.46(127.0.0.1)、模拟库、**单线程顺序**、 + 份额池充足的稳态(每笔 1 份,避免 `plan_lots` 跨批波动)。 + 端到端含阶段 1.5 空引擎(真实引擎的额外开销在 `test_convert_integration` 另行体现)。 + """ + core = conc_env["core"] + n = 60 + _seed(core, [str(n)]) + + writer = _TimingWriter(ConvertCoreRepository(engine=conc_env["core_rw"])) + totals: 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, + 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=writer, + thresholds=RiskThresholds.from_settings(), + now=_NOW, + engine_hook=_noop_engine, + ) + totals.append((time.perf_counter() - t0) * 1000.0) + + # 首个样本含连接池冷启动,剔除(否则把"建连"算进"业务耗时") + stage1 = writer.samples[1:] + e2e = totals[1:] + 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"max={max(e2e):.1f} mean={mean(e2e):.1f}" + ) + # PRD §9 第 18 条的验收硬指标是端到端 < 2s(阶段一 <100ms 是预估、非硬指标) + 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 cc30ffd..5ddaa17 100644 --- a/tests/test_convert_service.py +++ b/tests/test_convert_service.py @@ -47,6 +47,7 @@ from app.service.convert.errors import ( BelowMinQty, CrossEntityNotSupported, InsufficientShares, + LotConflict, NavNotReady, ProductNotRedeemable, SameProduct, @@ -388,6 +389,112 @@ def test_phase_two_failure_keeps_trade_and_marks_failed(sqlite_engine, monkeypat 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)