diff --git a/app/repository/core_ro.py b/app/repository/core_ro.py index eb03a02..7621d59 100644 --- a/app/repository/core_ro.py +++ b/app/repository/core_ro.py @@ -317,6 +317,13 @@ class CoreReadOnlyRepository: 与 list_holdings 同源同口径(市值降序、min_risk_code 判定 R4/R5), 只是把聚合收进仓储,避免调用方取回全量明细再自行求和(N+1)。 + ⚠️ **`qty > 0` 过滤(2026-09-10 T-11 收尾)**:原实现漏了这条,与 + `list_holdings` 出现口径落差。convert 转出全部份额后归零行保留(`qty=0`、 + 市值 0),它不构成持仓——**必须与列表口径一致**,否则「持仓画像」与 + 「持仓列表」两个出口对同一客户给出不一致的持仓集合。 + 对 `ratio` 数值无影响(归零行市值为 0,分子分母同增 0),但 `rows` + 明细会多出已清仓产品,且 `holdings_truncated` 的判定位会被虚占。 + 采用「SQL 取明细 + Python 聚合」而非 SUM(CASE WHEN):明细行还要进审计 摘要,且 sqlite/MySQL 的 min_risk_code 类型差异(VARCHAR(8) vs CHAR(2)) 下 Python 端判定最稳。limit+1 多取一行探测截断,省掉二次 COUNT 查询。 @@ -334,6 +341,7 @@ class CoreReadOnlyRepository: FROM core_holding h JOIN core_product p ON p.product_id = h.product_id WHERE h.customer_id = :cid + AND h.qty > 0 ORDER BY h.market_value DESC LIMIT :lim """ diff --git a/app/repository/risk_repository.py b/app/repository/risk_repository.py index ec67756..22ed767 100644 --- a/app/repository/risk_repository.py +++ b/app/repository/risk_repository.py @@ -220,21 +220,30 @@ class RiskRepository: for r in conn.execute(sql, {"pat": f"%{trade_id}%"}).mappings() ] - def has_engine_error_audit(self, trade_id: str) -> bool: - """存在 decision='risk_engine_error' 的该笔审计(rebuild_alerts 警示用)。 + def has_engine_error_audit( + self, trade_id: str, decision: str = "risk_engine_error" + ) -> bool: + """存在该决策码的、指向该笔的审计(rebuild_alerts 警示用)。 audit_log 无 trade_id 列,trade_id 在 input_summary(JSON 列/TEXT JSON) 内,同 find_alerts_by_trade 的 LIKE 口径。 + + `decision` 可覆盖:convert 线的阶段 1.5 失败审计用 **`engine_error`** + (`convert_service` ⑦ 步),与普通交易的 `risk_engine_error` 不同码—— + T-12 补偿脚本据此区分两条线,避免写第二份 LIKE 查询。 """ sql = text( """ SELECT 1 FROM audit_log - WHERE decision = 'risk_engine_error' AND input_summary LIKE :pat + WHERE decision = :decision AND input_summary LIKE :pat LIMIT 1 """ ) with self._engine.connect() as conn: - return conn.execute(sql, {"pat": f"%{trade_id}%"}).first() is not None + return ( + conn.execute(sql, {"decision": decision, "pat": f"%{trade_id}%"}).first() + is not None + ) def list_alerts( self, diff --git a/app/service/convert/convert_service.py b/app/service/convert/convert_service.py index 55a8600..97cca33 100644 --- a/app/service/convert/convert_service.py +++ b/app/service/convert/convert_service.py @@ -876,6 +876,131 @@ def rebuild_convert_response( ) +def compensate_convert( + group_id: str, + *, + core_ro: CoreReadOnlyRepository | None = None, + risk_repo: RiskRepository | None = None, + convert_repo: ConvertRepository | None = None, + thresholds: RiskThresholds | None = None, + now: datetime | None = None, + actor_id: str | None = None, +) -> dict[str, Any]: + """按 `convert_group_id` 从 Core 侧补偿一次转换(T-12 · PRD §7.1 / 架构 §5.4)。 + + 补偿**两件事**:「详情」+「预警」,且在同一把 `convert:rerun:{gid}` 锁内串行 —— + 锁键与 `convert_fund` 的幂等重试路径**完全一致**,故「人工补跑」与「客户端带同键 + 重试」不会并发重复写审计、重复出单(架构 §5.4「为何不自动化」条)。 + + 两个幂等锚点,各自独立: + - **详情侧**:`risk_convert_detail.status` 已 `completed` → 不重写(阶段二本就是 + upsert,重复执行不会产生第二行); + - **预警侧**:以**转出端 `out_trade_id`** 为锚点查 `find_alerts_by_trade` + (评审 Q4)—— 一次转换有**两条**流水,只认转出端才不会重复出单。 + + Returns(`state` 四态,CLI 据此判退出码): + `missing` Core 侧不足两条流水(不是有效转换组,**未做任何写入**) + `locked` 未抢到 `convert:rerun:{gid}`(有并发重试或另一实例在跑) + `skipped` 两件事都无需做 —— **幂等语义**,重复补偿走这里 + `rebuilt` 至少补写了一项(`detail` / `engine` 各自给明细) + + 注意:本函数**只补不判** —— 不重新做参数校验、适当性检查或份额规划, + 因为阶段一(Core 侧扣减)已成是它的前提(Core 侧不足两条流水直接 `missing`)。 + """ + core = core_ro or CoreReadOnlyRepository() + repo = risk_repo or RiskRepository() + crepo = convert_repo or ConvertRepository() + th = thresholds or RiskThresholds.from_settings() + now = now or datetime.now() + ensure_trace() + + with try_lock( + f"convert:rerun:{group_id}", settings.convert_lock_ttl_seconds + ) as acquired: + if not acquired: + logger.warning("补偿未抢到执行权(有并发重试/实例在跑):%s", group_id) + return {"state": "locked", "group_id": group_id, "alert_ids": []} + + trades = core.list_convert_trades(group_id) + if len(trades) < 2: + logger.warning( + "补偿缺 Core 侧流水(需 2 条,实得 %d):%s", len(trades), group_id + ) + return { + "state": "missing", + "group_id": group_id, + "trade_count": len(trades), + "alert_ids": [], + } + out_trade = next(t for t in trades if t["trade_type"] == "redeem") + in_trade = next(t for t in trades if t["trade_type"] == "subscribe") + out_trade_id = str(out_trade["trade_id"]) + in_trade_id = str(in_trade["trade_id"]) + + # ── 详情侧:非 completed 才补写(复用阶段二补跑,不另写第二份)── + detail_row = crepo.get_by_group_id(group_id) + detail = "already_completed" + if detail_row is None or str(detail_row.get("status")) != "completed": + req = { + "customer_id": str(out_trade["customer_id"]), + "from_product_id": str(out_trade["product_id"]), + "to_product_id": str(in_trade["product_id"]), + # 原幂等键随详情一起留痕;补写时无需再校验(阶段一已成) + "client_request_id": (detail_row or {}).get("client_request_id"), + } + finalize = _finalize_from_core( + req, group_id, core, repo, crepo, now, actor_id + ) + detail = "missing" if finalize is None else "rebuilt" + + # ── 预警侧:幂等锚点 = 转出端 ── + existing = repo.find_alerts_by_trade(out_trade_id) + if existing: + engine = "skipped" + alert_ids = [str(r["alert_id"]) for r in existing] + triggered_rules: list[str] = [] + # 部分失败中间态:预警单在、但当初引擎异常中断过 → 提示人工核对 + # (复用 rebuild_alerts 的同一判定,仅决策码换成 convert 线的 engine_error) + warning = None + if repo.has_engine_error_audit(group_id, decision="engine_error"): + warning = ( + "该转换曾阶段 1.5 引擎异常中断(engine_error):已入单可能不完整" + "(aml 单/L3 或缺失),请人工核对预警台账与 audit_log" + ) + else: + result = _run_engine( + out_trade, + in_trade, + core_ro=core, + risk_repo=repo, + thresholds=th, + engine_hook=None, + ) + warning = None + if result is None: + engine = "unavailable" # 引擎模块缺失(T-8 未落地),不算补写成功 + alert_ids = [] + triggered_rules = [] + else: + engine = "rebuilt" + alert_ids = list(result.get("alert_ids") or []) + triggered_rules = list(result.get("triggered_rules") or []) + + out: dict[str, Any] = { + "state": "rebuilt" if (detail == "rebuilt" or engine == "rebuilt") else "skipped", + "group_id": group_id, + "out_trade_id": out_trade_id, + "in_trade_id": in_trade_id, + "detail": detail, + "engine": engine, + "alert_ids": alert_ids, + "triggered_rules": triggered_rules, + } + if warning: + out["warning"] = warning + return out + + def _as_date(value: Any) -> date: """sqlite 读回 DATE/TIMESTAMP 为字符串 → 统一成 date。""" if isinstance(value, datetime): @@ -885,4 +1010,4 @@ def _as_date(value: Any) -> date: return date.fromisoformat(str(value)[:10]) -__all__ = ["convert_fund", "rebuild_convert_response", "PROCESSING"] +__all__ = ["convert_fund", "compensate_convert", "rebuild_convert_response", "PROCESSING"] diff --git a/docs/memory/2026-09-10.md b/docs/memory/2026-09-10.md index 35073b2..6296d7e 100644 --- a/docs/memory/2026-09-10.md +++ b/docs/memory/2026-09-10.md @@ -546,3 +546,68 @@ B.1 状态 / B.6 拓扑与基线 / **新增 B.6.3**)· `docs/memory/{TODO,MEMO (§0 导航 / §B 状态 / B.1 / B.6 拓扑与基线 / **新增 B.6.4**)· `docs/memory/{TODO,MEMORY}` · 本条。 **下一步 = T-12(补偿脚本)→ T-13(全量回归 + 50 并发压测 + PRD §5.3 数字回填)**。 + +--- + +## ✅ T-12 补偿脚本 已完成(2026-09-10 · 本线第 6 批) + +**目标**:阶段二失败(或阶段 1.5 引擎失败)后,能**仅凭 Core 侧数据**把「详情 + 预警」两件事补回来 +(PRD §7.1 / 架构 §5.4)。 + +**交付(改 2 + 增 2 脚本 + 测试 2 文件)** + +| 文件 | 动作 | +| --- | --- | +| `app/service/convert/convert_service.py` | 【改】新增公开 **`compensate_convert(group_id, ...)`** —— 补偿的**服务端单点入口** | +| `app/repository/risk_repository.py` | 【改】`has_engine_error_audit(trade_id, decision=...)` 加 decision 参数(默认值不变) | +| `scripts/demo/rebuild_alerts.py` | 【改】新增 `--convert-group CNV-xxx`(与位置参数 trade_ids 互斥) | +| `scripts/agent/cleanup_pending_convert.py` | 【新增】超 SLA 的 `pending` → `expired`(**标记不硬删**,S2) | +| `scripts/dev/verify_convert_compensate.py` | 【新增】真 MySQL 验证脚本(34 项) | +| `tests/test_convert_service.py` / `tests/test_demo_scripts.py` | 【改】+4 / +8 用例 | + +**执行期裁定 3 条(留痕)** + +1. **补偿逻辑收敛到 service 单点**(自检第 13 问):`rebuild_alerts --convert-group` 只是**薄封装**, + 脚本单测用「注入假实现验证参数转交」钉住这条约束,防止日后长成第二份实现。 +2. **锁与幂等锚点都沿用既有口径,不新造**:锁键 `convert:rerun:{gid}` 与 `convert_fund` + 幂等重试路径**同一个键**(人工补跑 vs 客户端重试不并发重复出单);幂等锚点 = **转出端 + `out_trade_id`**(一次转换两条流水,只认转出端;复用 `find_alerts_by_trade`)。 +3. **`has_engine_error_audit` 加 `decision` 参数而非另写一份查询**:convert 线阶段 1.5 失败的 + 审计决策码是 **`engine_error`**,与普通交易的 **`risk_engine_error` 不是同一个码**; + 默认值保持 `risk_engine_error`,B9a 既有路径零改动。 + +**`state` 四态 + CLI 退出码**:`rebuilt`(0) / `skipped`(0,幂等) / `missing`(1,Core 侧不足两条流水, +零写入) / `locked`(3,未抢到 `convert:rerun` 锁)。**2 保留给 argparse 用法错误**,故 locked 取 3。 + +**验证** + +- `pytest -q` → **731 passed / 3 skipped**(基线 719 **+12**,零回归) +- **突变验证 3 组**:① 去掉预警侧幂等锚点(恒跑引擎)→ **精准 1 红**; + ② 去掉 `list_expired_candidates` 的 `status` 过滤 → **精准 2 红**;③ 补偿无视抢锁结果 → **精准 1 红**。 + 均已恢复,`grep "MUTATION\|1 = 1"` **无残留** +- **真库**:`verify_convert_compensate.py` **34/34**、跑完**零残留** +- **全套真库脚本复跑**(铁律 2:`convert_service` / `risk_repository` / `core_ro` 三处均被触及): + seed `全部 PASS` · apply `24/24` · service `35/35` · engine `31/31` · lots `20/20` · tools `14/14` · + compensate `34/34`,**均零回归** + +**真库专属证据(sqlite 给不了的,本任务核心增量)** + +1. `status='expired'` 在 MySQL **ENUM** 上被接受 —— sqlite 该列是 `VARCHAR(16)`,**写什么都收**, + 「能写进 expired」只能在真库证明。 +2. `created_at < cutoff` 在 **DATETIME(3)** 上的时间边界正确(超时进候选 / 未超时不进 / 复跑幂等)。 +3. ⭐ **`input_summary` 是真 JSON 列**,而 `has_engine_error_audit` 用 `LIKE '%gid%'` 判定。 + 脚本先断言 `information_schema` 里 `DATA_TYPE='json'` 再验命中,并**反向断言** + 「决策码不匹配则不命中」,证明过滤真的按 `decision` 走(sqlite 是 TEXT,此风险不可见)。 + +**新增目录 `scripts/agent/`**:按架构 §2 与开发计划 §9 指定路径落地。`scripts/` 下原有 +`core`/`cron`/`demo`/`dev`/`kb`/`sync`,其中 `cron` 已承载两个巡检脚本 —— 若后续要统一巡检脚本的家, +可评估合并;**本次不擅自偏离已评审的架构文档**。 + +**顺带收口(用户指示「顺手改完」)**:`core_ro.concentration_profile` 补 `h.qty > 0`, +与 `list_holdings` **真正同口径**(`ratio` 不变,但 `rows` 不再多出已清仓产品、不虚占截断位)。 +新用例含**跨出口一致性断言**;突变验证:去掉 `h.qty > 0` → **精准 1 红**(已恢复)。 + +**文档回写**:开发计划 §9(DoD 全勾 + **新增 §9.1 执行记录**)+ §7.3 收口段 · `交接文档.md` **v2.3** +(表头 / §0 三线状态 / B.1 / B.6 拓扑与基线 / **新增 B.6.5**)· `docs/memory/{TODO,MEMORY}` · 本条。 + +**下一步 = T-13(全量回归 + 50 并发压测 + 性能补录 · 本线最后一个任务)**。 diff --git a/docs/memory/MEMORY.md b/docs/memory/MEMORY.md index 59c80aa..aab8045 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-0b / T-1 已完成**):** `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 完成(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(补偿脚本)**。**开工前置两个阻断项(✅ 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 设计+开发计划闭环,**第 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 号。 **仓库地图:** @@ -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-11 已完成、718 绿,下一步 = T-12(补偿脚本)**;设计 + 开发计划均已闭环,入口 项目根 `交接文档.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 待办)**:① **基金转换线**第 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**。 **禁止(改代码前必记):** Core 正式 C1~C5 不可被画像覆盖 · 审计表只 INSERT · 代理人草稿不外发 · 仅 R-02 可阻断交易 · 四 Agent 不互调 LLM。 @@ -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` 全量(当前 **718 passed / 3 skipped**,基线 510)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) -6. **当前有哪两条并行线?**(① 风控/架构改进线:**已结项**(510 基线绿、§7.2 七项冒烟 7/7 PASS、`037ce7e` 已核实早已推送);② **基金转换线**:设计闭环,**第 5 步进行中 —— T-0~T-11 已完成(718 passed),下一步 = T-12**)——动代码前先确认自己属于哪条线,别混淆前置条件。 +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(最后一个任务)**)——动代码前先确认自己属于哪条线,别混淆前置条件。 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 6e2dda9..a03fc2d 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-11 均已于 2026-09-10 完成**(**718 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(补偿脚本)**。两个阻断前置(**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-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 条回归面)**。 ### 基金转换线待办(推荐顺序) @@ -17,14 +17,18 @@ - [x] **【T-1 · 第 1 批】DDL + 种子 + sqlite 同步** —— **完成(2026-09-10)**:`scripts/core/01-ddl.sql` 新建 `core_fee_rule`/`core_share_lot`/`core_convert_lot_detail` + `core_trade` 加 `convert_group_id`+索引 + `core_product` 加 8 列(`subscribe_fee_rate` 等)+ `fee_rate` 补 COMMENT;**新增 `07-seed-fee-rule.sql`**(14 产品 × 5 档,按 22 号文 §10)/ **`08-seed-share-lot.sql`**(58 行持仓 → 61 行批次,Σ remain_qty 恒等于 qty,CUST-9527 跨批次)/ **`09-seed-org.sql`**(管理人 + TA + 申购费率 + 最低持有余额);`reset.ps1` 追加 07/08/09;`02-mysql-agent专用.sql` 追加 `risk_convert_detail`(status ENUM 建表即 5 值);`tests/_ddl.py` 同步 4 表 + `REQUIRED_CONVERT_TABLES` 门禁。**验证**:新增 `scripts/dev/verify_convert_seed.py`(pymysql 等价 reset 流程 + 8 条断言)→ **8/8 PASS**;`pytest -q` → **516 passed / 3 skipped(零回归)**。**3 点需注意**:① mysql 不在 PATH → 用该脚本替代 reset.ps1;② `core_fee_rule` 读取走只读账号(T-6 遵守);③ ~~`PROD-005827` 费率分类口径差异(`mixed` vs 主动偏股)待裁定~~ → **已裁定并修正(PRD v0.9.1)**:`mixed` 归位 `0.0050`(其他混合型),主示例转入方改真主动偏股 `PROD-003095`,`09-seed-org.sql` 升 **v1.1** 按「管理人全产品线」重排,并新增断言 ⑧ 机器化卡口 - [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 ✅】均已结项(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 起】剩余任务**:**T-12(补偿脚本)→ T-13「50 并发压测 + 性能补录」**(最后跑,产出 PRD §9 第 18 条实测值)。**T-13 内部顺序**:先 50 并发压测 → 再性能实测补录 → 最后 PRD §5.3 数字回填(详见架构 §15 + 开发计划 §2~§10) +- [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;③ 实测超阈值 → 优化索引/锁策略后**重定阈值**,**不得反向修改实测数据迁就指标** - [ ] **【待用户裁定】幂等重放响应的数值位数偏差**(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-12**(补偿脚本,开发计划 §9)。基线 **718 passed / 3 skipped**(T-10 后 714 → **T-11 后 718**,+4)。 +> **下一步 = T-13**(全量回归 + 50 并发压测 + 性能补录 **· 本线最后一个任务**)。基线 **731 passed / 3 skipped**(T-10 后 714 → T-11 后 718 → **T-12 后 731**,+13)。 > **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` 只是薄封装; +> ② 锁键 `convert:rerun:{gid}` 与幂等锚点(**转出端 `out_trade_id`**)**沿用既有口径不新造**; +> ③ `has_engine_error_audit` **加 `decision` 参数**而非另写一份查询(convert 线失败码是 `engine_error`,与 `risk_engine_error` 不同)。 +> **T-12 真库三处 sqlite 给不了的证据**:`expired` 的 **ENUM 值域** · `DATETIME(3)` 时间边界 · **真 JSON 列 `input_summary LIKE`**。 > ⚠️ **基金转换的 T-0b 与下方「架构改进第 3/4 批」的 `core_ro` 只读账号是同一件事** —— 已由本线定案为 D20,**不再挂在架构改进线**(该线原「不要做」清单已更新)。 diff --git a/docs/项目框架设计/开发计划-基金转换交易.md b/docs/项目框架设计/开发计划-基金转换交易.md index c5e2558..0018a5d 100644 --- a/docs/项目框架设计/开发计划-基金转换交易.md +++ b/docs/项目框架设计/开发计划-基金转换交易.md @@ -1234,12 +1234,16 @@ sqlite 无 gap lock,故该分支由 `tests/test_convert_core.py` 用注入点 **回归复跑**(铁律 2,接线变化必跑):T-8 `31/31`(`amount_view` 改名影响面)· T-7 `35/35` · T-10 `20/20`,**均零回归** -*未顺手改(不在本任务范围,已上报)* +*收口(原「未顺手改」,用户 2026-09-10 指示一并改掉)* -- **`concentration_profile` 未过滤 `qty = 0`**:它走**独立 SQL**(非 `list_holdings`),当前不过滤归零行。 - 归零行 `market_value = 0`,对 R4+R5 占比**无实际影响**,但 `core_ro.py` 中 - 「与 list_holdings 同源同口径」的注释与实际存在**措辞落差**。 - 是否一并收口属新范围,留待用户决定。 +- **`concentration_profile` 已补 `qty > 0` 过滤**:它走**独立 SQL**(非 `list_holdings`), + 原实现不过滤归零行,与 `core_ro` 中「与 list_holdings 同源同口径」的注释存在措辞落差。 + 本次收口后两个出口(持仓画像 / 持仓列表)对同一客户给出**一致的持仓集合**。 + 数值影响:`ratio` 不变(归零行市值为 0,分子分母同增 0),但 `rows` 明细不再多出已清仓产品、 + 不再虚占 `holdings_truncated` 的判定位。 + - 新增用例 `tests/test_concentration_c4.py::test_concentration_profile_filters_zero_qty`, + 内含**跨出口一致性断言**(`list_holdings` 与 `concentration_profile["rows"]` 同集合); + - 突变验证:去掉 `h.qty > 0` → **精准 1 条**红(该新用例),已恢复。 --- @@ -1359,15 +1363,81 @@ sqlite 无 gap lock,故该分支由 `tests/test_convert_core.py` 用注入点 **幂等锚点(Q4)**:以**转出端 `out_trade_id`** 为锚点 —— 复用 `rebuild_alerts.py:54-66` 的 `find_alerts_by_trade(trade_id)`,命中即 `skipped`;**一次转换有两条流水,只认转出端**,避免重复出单。该脚本**本就幂等**(其 `:8-10` docstring),**不需要再加 `uk_idem`**。 **DoD** -- [ ] 阶段二失败 → `decision='convert_detail_write_failed'` 审计存在(验收 17 前置) -- [ ] `rebuild_alerts --convert-group` 能仅凭 Core 侧数据补出**完整详情** → `completed`(验收 17) -- [ ] 重复执行补偿 → `skipped`,不产生第二张预警单(幂等) -- [ ] `cleanup_pending_convert.py` 把超 24h `pending` 置 `expired`,**行仍在**(不硬删) +- [x] 阶段二失败 → `decision='convert_detail_write_failed'` 审计存在(验收 17 前置) +- [x] `rebuild_alerts --convert-group` 能仅凭 Core 侧数据补出**完整详情** → `completed`(验收 17) +- [x] 重复执行补偿 → `skipped`,不产生第二张预警单(幂等) +- [x] `cleanup_pending_convert.py` 把超 24h `pending` 置 `expired`,**行仍在**(不硬删) **依赖**:T-4 / T-7 --- +### 9.1 执行记录(2026-09-10 · 已落地) + +**改动文件** + +| 文件 | 动作 | +| --- | --- | +| `app/service/convert/convert_service.py` | 【改】新增公开 **`compensate_convert(group_id, ...)`** —— 补偿的**服务端单点入口** | +| `app/repository/risk_repository.py` | 【改】`has_engine_error_audit(trade_id, decision=...)` 加 decision 参数(默认值不变,既有调用零改动) | +| `scripts/demo/rebuild_alerts.py` | 【改】新增 `--convert-group CNV-xxx`(与位置参数 trade_ids 互斥) | +| `scripts/agent/cleanup_pending_convert.py` | 【新增】超 SLA 的 `pending` → `expired`(**标记不硬删**) | +| `scripts/dev/verify_convert_compensate.py` | 【新增】真 MySQL 验证脚本 | +| `tests/test_convert_service.py` | 【改】+4 条补偿用例 | +| `tests/test_demo_scripts.py` | 【改】+8 条(`--convert-group` 分派/接线 + cleanup 脚本) | + +**执行期裁定 3 条(留痕)** + +1. **补偿逻辑收敛到 service 单点**(自检第 13 问)。`rebuild_alerts.py --convert-group` 只是 + **薄封装**,不复制任何一条逻辑;脚本单测以「注入假实现验证参数转交」锁住这条约束。 +2. **锁与幂等锚点都沿用既有口径,不新造**: + - 锁键 `convert:rerun:{gid}` —— 与 `convert_fund` 幂等重试路径**同一个键**, + 故「人工补跑」与「客户端带同键重试」不会并发重复出单(架构 §5.4「为何不自动化」原话); + - 幂等锚点 = **转出端 `out_trade_id`**,复用 `find_alerts_by_trade`(评审 Q4 原话)。 +3. **`has_engine_error_audit` 加 `decision` 参数而非另写一份查询**:convert 线的阶段 1.5 失败 + 审计决策码是 **`engine_error`**(`convert_service` ⑦ 步),与普通交易的 `risk_engine_error` + **不是同一个码**;默认值保持 `risk_engine_error`,B9a 既有路径零改动。 + +**`state` 四态(CLI 退出码据此判定)** + +| state | 含义 | CLI 退出码 | +| --- | --- | --- | +| `rebuilt` | 至少补写了一项(`detail` / `engine` 各自给明细) | 0 | +| `skipped` | 详情已 `completed` 且预警已存在 —— **幂等语义** | 0 | +| `missing` | Core 侧不足两条流水(不是有效转换组,**零写入**) | 1 | +| `locked` | 未抢到 `convert:rerun:{gid}`(有并发重试/实例在跑) | 3 | + +> 退出码 **2 保留给 argparse 的参数用法错误**,故 `locked` 取 3,避免与用法错误混淆。 + +**验证** + +- `pytest -q` → **731 passed / 3 skipped**(基线 719 加 12,零回归) +- **突变验证 3 组(防假绿)**: + ① 去掉预警侧幂等锚点(恒跑引擎)→ **精准 1 条**红(`test_compensate_is_idempotent`); + ② 去掉 `list_expired_candidates` 的 `status` 过滤 → **精准 2 条**红(只清理超时 pending / 复跑幂等); + ③ 让补偿无视抢锁结果 → **精准 1 条**红(`test_compensate_returns_locked_...`)。 + 三处均已恢复,`grep "MUTATION\|1 = 1"` **无残留** +- **真库**:新增 `verify_convert_compensate.py` **34/34**(零残留) +- **回归复跑(铁律 2)**:`convert_service` / `risk_repository` / `core_ro` 三处均被本次改动触及 → + 全套 7 个真库脚本原样复跑:seed `全部 PASS` · apply `24/24` · service `35/35` · engine `31/31` · + lots `20/20` · tools `14/14` · compensate `34/34`,**均零回归** + +**真库专属证据(sqlite 给不了的)** + +- `status='expired'` 在 MySQL **ENUM** 上被接受(sqlite 是 `VARCHAR(16)`,写什么都收); +- `created_at < cutoff` 在 **DATETIME(3)** 上的时间边界正确(超时进候选、未超时不进、复跑幂等); +- `input_summary LIKE '%gid%'` 在**真 JSON 列**上生效 —— 脚本先断言 `DATA_TYPE='json'` 再验命中, + 并反向断言「决策码不匹配则不命中」,证明过滤真的按 `decision` 走。 + +**新增目录 `scripts/agent/`** + +架构 §2 与本文 §9 均指定 `scripts/agent/cleanup_pending_convert.py`,**按文档落地**。 +(`scripts/` 下原有 `core` / `cron` / `demo` / `dev` / `kb` / `sync`,其中 `cron` 已承载 +`agent_behavior_scan` / `escalation_scan` 两个巡检脚本 —— 若后续要统一巡检脚本的家, +可评估把三者合并到 `cron/`;本次不擅自偏离已评审的架构文档。) + +--- + ## 10. 第 7 批 · 回归与实测(T-13) ### 10.1 全量回归 diff --git a/scripts/agent/cleanup_pending_convert.py b/scripts/agent/cleanup_pending_convert.py new file mode 100644 index 0000000..3e95b99 --- /dev/null +++ b/scripts/agent/cleanup_pending_convert.py @@ -0,0 +1,103 @@ +"""超时 `pending` 转换占位清理(T-12 · 架构 §5.4 · PRD §7.1 I-3)。 + +背景:convert 阶段零先落一行 `risk_convert_detail(status='pending')` 占位,阶段二 +回写 `completed`。若进程在阶段零与阶段一之间挂掉,这行占位会**永久滞留**,既污染 +巡检又让人误判「有一笔转换在飞」。本脚本按 SLA(`convert_compensate_sla_hours`, +默认 24h)把超时占位置 `status='expired'`。 + +**标记不硬删(评审 S2)**:只 UPDATE 状态,**绝不 DELETE** —— 留痕供对账 +(`expired` 是建表即含的 ENUM 值,无需 ALTER)。这与 +`ConvertRepository.mark_expired` 是**同一份实现**,本脚本不复制 SQL。 + +用法与退出码: + python scripts/agent/cleanup_pending_convert.py # 按 SLA 24h 清理 + python scripts/agent/cleanup_pending_convert.py --hours 1 # 覆盖 SLA(演练/测试) + python scripts/agent/cleanup_pending_convert.py --dry-run # 只报告不写库 + + 0 = 清理完成(含「无候选」);1 = MySQL 不可达或写库失败。 +""" + +from __future__ import annotations + +import argparse +import json +import sys +from pathlib import Path + +# ① sys.path 引导项目根(rebuild_alerts / escalation_scan 先例) +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT)) + +from app.config.settings import settings # noqa: E402 +from app.repository.convert_repository import ConvertRepository # noqa: E402 +from app.utils.trace import new_trace # noqa: E402 + + +def cleanup(repo: ConvertRepository, hours: int, *, dry_run: bool = False) -> dict: + """扫出超 `hours` 的 pending 占位并置 `expired`;返回 JSON 友好摘要。 + + 幂等:`mark_expired` 只把 `pending` 行改状态,已 `expired` 的行**不再进候选** + (`list_expired_candidates` 硬编码 `status='pending'`),故重复执行第二次为空操作。 + """ + candidates = repo.list_expired_candidates(hours) + expired: list[str] = [] + if not dry_run: + for row in candidates: + group_id = str(row["convert_group_id"]) + repo.mark_expired(group_id) + expired.append(group_id) + return { + "sla_hours": hours, + "dry_run": dry_run, + "candidate_count": len(candidates), + "expired_count": len(expired), + "expired": expired, + "candidates": [ + { + "convert_group_id": str(r["convert_group_id"]), + "client_request_id": r.get("client_request_id"), + "created_at": r.get("created_at"), + } + for r in candidates + ], + } + + +def main() -> int: + parser = argparse.ArgumentParser( + description="超时 pending 转换占位 → status='expired'(标记不硬删,S2)" + ) + parser.add_argument( + "--hours", + type=int, + default=None, + help=f"覆盖 SLA 小时数(缺省取 settings.convert_compensate_sla_hours=" + f"{settings.convert_compensate_sla_hours})", + ) + parser.add_argument("--dry-run", action="store_true", help="只报告候选,不写库") + args = parser.parse_args() + + hours = args.hours if args.hours is not None else settings.convert_compensate_sla_hours + new_trace() + repo = ConvertRepository() + + # 连接自检(零写入探针) + try: + repo.list_expired_candidates(hours) + except Exception as exc: + print(f"MySQL 不可达:{exc}", file=sys.stderr) + print("请确认本机 MySQL 服务已启动(见 FLOW §0 本机状态)。", file=sys.stderr) + return 1 + + try: + summary = cleanup(repo, hours, dry_run=args.dry_run) + except Exception as exc: # noqa: BLE001 + print(f"清理失败:{exc}", file=sys.stderr) + return 1 + + print(json.dumps(summary, ensure_ascii=False, default=str, indent=2)) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/scripts/demo/rebuild_alerts.py b/scripts/demo/rebuild_alerts.py index 352fa72..7b347dc 100644 --- a/scripts/demo/rebuild_alerts.py +++ b/scripts/demo/rebuild_alerts.py @@ -9,6 +9,13 @@ 不重复出单、不重复 append 事件。**勿并行运行多个本脚本实例**(聚合锁为进程内 锁,跨进程并发重放同一 trade_id 无防护)。 +**基金转换补偿(T-12 · PRD §7.1 / 架构 §5.4)**:`--convert-group CNV-xxx` 从 +Core 侧(`core_trade` 两条流水 + `core_convert_lot_detail` 的 `nav`/`nav_date`) +补偿一次转换的**详情 + 预警**两件事,全部委托 `convert_service.compensate_convert` +(锁在 service 内,锁键 `convert:rerun:{gid}`,与客户端带同键重试共用 → +人工补跑与客户端重试不会并发重复出单)。幂等锚点为**转出端 out_trade_id**: +一次转换有两条流水,只认转出端才不重复出单(评审 Q4)。 + 可复现性与局限: - RISK-004/005 规则窗口以 core_trade.traded_at 为事件时点,不受脚本执行时刻影响; 聚合锚点(同日 pending 单查找)按执行日——跨日补放历史交易时事件并入执行日 @@ -19,6 +26,11 @@ 用法: python scripts/demo/rebuild_alerts.py TRD-20260907-AB12CD34 [更多 trade_id ...] + python scripts/demo/rebuild_alerts.py --convert-group CNV-20260907-AB12CD34 + +退出码:0 = 全部成功/幂等跳过;1 = 有 trade_id 未找到,或 convert 组 Core 侧 +流水不足(missing);3 = convert 组未抢到执行权(locked,稍后重试即可)。 +(2 保留给 argparse 的参数用法错误。) """ from __future__ import annotations @@ -31,11 +43,14 @@ from typing import Any ROOT = Path(__file__).resolve().parents[2] sys.path.insert(0, str(ROOT)) +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.convert_service import compensate_convert # noqa: E402 from app.service.risk.engine import process_trade_event # noqa: E402 + def rebuild_trade( trade_id: str, core: CoreReadOnlyRepository, @@ -74,11 +89,41 @@ def rebuild_trade( } +def rebuild_convert_group( + group_id: str, + core: CoreReadOnlyRepository, + repo: RiskRepository, + crepo: ConvertRepository | None = None, +) -> dict[str, Any]: + """按 convert_group_id 补偿一次转换(CLI 与单测共用入口)。 + + 薄封装:全部逻辑在 `convert_service.compensate_convert`(锁 / 详情补写 / 预警 + 补偿都是那**一份**实现,本脚本不复制任何一条 —— 自检第 13 问)。 + 返回 `{"state": "rebuilt"|"skipped"|"missing"|"locked", ...}`。 + """ + return compensate_convert( + group_id, core_ro=core, risk_repo=repo, convert_repo=crepo + ) + + def main() -> None: - parser = argparse.ArgumentParser(description="按 trade_id 幂等重放风控引擎(引擎异常补偿)") - parser.add_argument("trade_ids", nargs="+", help="core_trade 中的 trade_id,可一次多笔") + parser = argparse.ArgumentParser( + description="按 trade_id 幂等重放风控引擎;或按 --convert-group 补偿一次基金转换" + ) + parser.add_argument("trade_ids", nargs="*", help="core_trade 中的 trade_id,可一次多笔") + parser.add_argument( + "--convert-group", + metavar="CNV-xxx", + help="按 convert_group_id 从 Core 侧补偿一次转换(详情 + 预警),与 trade_id 互斥", + ) args = parser.parse_args() + if args.convert_group: + if args.trade_ids: + parser.error("--convert-group 与 trade_id 位置参数互斥,请择一使用") + elif not args.trade_ids: + parser.error("请给出 trade_id 位置参数,或使用 --convert-group CNV-xxx") + core = CoreReadOnlyRepository() repo = RiskRepository() @@ -91,6 +136,40 @@ def main() -> None: print("请确认本机 MySQL 服务已启动(见 FLOW §0 本机状态)。", file=sys.stderr) raise SystemExit(1) + if args.convert_group: + out = rebuild_convert_group(args.convert_group, core, repo, ConvertRepository()) + gid = args.convert_group + state = out["state"] + if state == "rebuilt": + rules = ",".join(out.get("triggered_rules") or []) or "-" + print( + f"[rebuilt] {gid}: detail={out['detail']} engine={out['engine']} " + f"rules={rules} alerts={out['alert_ids']} " + f"out={out['out_trade_id']} in={out['in_trade_id']}" + ) + elif state == "skipped": + print( + f"[skipped] {gid}: 详情已 completed、预警单已存在 " + f"{out['alert_ids']}(幂等跳过)" + ) + else: + print( + f"[{state}] {gid}: " + + ( + f"Core 侧流水不足(实得 {out.get('trade_count')} 条,需 2 条)" + if state == "missing" + else "未抢到 convert:rerun 执行权,稍后重试" + ), + file=sys.stderr, + ) + if out.get("warning"): + print(f"[warning] {gid}: {out['warning']}", file=sys.stderr) + if state == "missing": + raise SystemExit(1) + if state == "locked": + raise SystemExit(3) + return + missing = 0 for tid in args.trade_ids: out = rebuild_trade(tid, core, repo) diff --git a/scripts/dev/verify_convert_compensate.py b/scripts/dev/verify_convert_compensate.py new file mode 100644 index 0000000..45c67cb --- /dev/null +++ b/scripts/dev/verify_convert_compensate.py @@ -0,0 +1,469 @@ +"""T-12 真 MySQL 验证脚本:补偿脚本在真库/真账号下实跑 + DoD 逐条断言。 + +**为什么 sqlite 单测全绿还不够(T-12 视角)** + +1. **`status='expired'` 的真库 ENUM 值域**:`risk_convert_detail.status` 在 MySQL 是 + `ENUM('pending','completed','failed','cancelled','expired')`,写错值会**直接报错**; + 而 sqlite 测试库该列是 `VARCHAR(16)`,**任何字符串都收**。故「cleanup 能写进 expired」 + 这件事必须在真库证明一次。 +2. **`input_summary` 是真 JSON 列**:`has_engine_error_audit` 用 `input_summary LIKE '%gid%'` + 做幂等/警示判定。sqlite 里它是 TEXT、必定能 LIKE;MySQL 的 JSON 列能否直接 LIKE + 属方言行为,必须实跑确认(本脚本 E 组专测)。 +3. **DATETIME(3) 的时间边界**:`list_expired_candidates` 的 `created_at < cutoff` 在 + MySQL 毫秒精度列上是否按预期分流(超时进候选、未超时不进),sqlite 的字符串时间 + 给不了这个保证。 +4. **端到端**:阶段二失败 → 补偿 → `completed` + 出单 + 主审计,走真事务、真账号。 + +**断言分组** +A 待补偿态(阶段二 + 阶段 1.5 双失败)在真库如实落痕 +B 超时 `pending` 占位 → `expired`(含 ENUM 值域 / 时间边界 / 幂等 / 标记不硬删) +C `compensate_convert` 端到端:详情 `completed` + 1 张单 + 主审计 +D 重复补偿 → `skipped`(不产生第二张单、不重复主审计) +E 真 **JSON 列**上 `input_summary LIKE` 生效(`has_engine_error_audit`) +F 清理后残留为零(自检) + +用法: + python scripts/dev/verify_convert_compensate.py + +约定(与 T-6/T-7/T-8/T-11 脚本一致): +- 用 **T12C 前缀**隔离数据(客户/产品/批次/group_id),跑完**两个库**(core + agent)全清; +- 建/清数据走 `role="admin"`(需 DELETE);业务本身走 `convert_fund` 默认账号; +- 固定 `trace_id = "T12C-VERIFY-TRACE"`,便于精准清理 `audit_log`; +- 退出码 1 = 有断言不一致。 +""" + +from __future__ import annotations + +import json +import sys +import uuid +from datetime import date, datetime, timedelta +from decimal import Decimal +from pathlib import Path +from unittest.mock import patch + +from sqlalchemy import text + +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT)) + +from app.config.settings import settings # noqa: E402 +from app.repository.convert_repository import ConvertRepository # noqa: E402 +from app.repository.core_ro import CoreReadOnlyRepository # noqa: E402 +from app.repository.risk_repository import RiskRepository # noqa: E402 +from app.service.convert.convert_service import ( # noqa: E402 + compensate_convert, + convert_fund, +) +from app.service.risk.rules import RiskThresholds # noqa: E402 +from app.utils.db import dispose_engines, get_engine # noqa: E402 +from app.utils.trace import new_trace # noqa: E402 + +CUSTOMER = "CUST-T12C" +PROD_OUT = "PROD-T12CO" +PROD_IN = "PROD-T12CI" +COMPANY = "华夏模拟基金" +TA = "TA-CN-001" +TRADE_AT = datetime(2026, 9, 4, 10, 0, 0) +TRADE_DATE = date(2026, 9, 4) +NOW = TRADE_AT +IN_NAV = Decimal("0.9500") +OUT_RATE = Decimal("0.0030") +IN_RATE = Decimal("0.0080") +FEE_TIERS = [ + (0, 7, "0.0150"), (7, 30, "0.0100"), (30, 180, "0.0050"), + (180, 365, "0.0025"), (365, None, "0.0000"), +] +CID_REQ = "T12C-COMP-0001" +TRACE = "T12C-VERIFY-TRACE" + +#: 大额阈值调到 100000 —— 转出端约 412000,确保补偿时引擎**真出一张单**(可断言) +TH = RiskThresholds( + large_amount=Decimal("100000"), + daily_total=Decimal("1000000"), + freq_count=3, + probe_window_minutes=5, + probe_count=3, + probe_amount=Decimal("400000"), + small_amount=Decimal("10000"), + small_count=3, + concentration_threshold=1.01, # 持仓画像不参与(避免 RISK-006 干扰) +) + +_passed = 0 +_failed = 0 + +#: 建/清数据用 admin 引擎(需 DELETE,R-e);在 main 中初始化 +core_engine = None +agent_engine = None + + +def check(name: str, actual, expected) -> None: + global _passed, _failed + ok = actual == expected + if ok: + _passed += 1 + else: + _failed += 1 + flag = "✅" if ok else "❌" + print(f" {flag} {name}: 实际 {actual!r}" + ("" if ok else f" / 期望 {expected!r}")) + + +def _t12c_id(prefix: str, now: datetime) -> str: + """固定前缀 id 工厂:默认 `CNV-<日期>-` 与清理口径 `LIKE 'CNV-T12C%'` 不匹配, + 重跑会留撞唯一键的残留(T-7 脚本踩过同款坑)。""" + return f"{prefix}-T12C-{uuid.uuid4().hex[:8].upper()}" + + +def q1(engine, sql: str, **params): + with engine.connect() as conn: + return conn.execute(text(sql), params).scalar_one_or_none() + + +def rows(engine, sql: str, **params): + with engine.connect() as conn: + return [dict(r) for r in conn.execute(text(sql), params).mappings()] + + +# ── 数据准备 / 清理 ───────────────────────────────────────────────── +def cleanup(core_engine, agent_engine) -> None: + """按 T12C 前缀清理两个库(含幂等重跑前置清理)。""" + with core_engine.begin() as conn: + for sql, params in [ + ("DELETE FROM core_convert_lot_detail WHERE convert_group_id LIKE 'CNV-T12C%'", {}), + ("DELETE FROM core_trade WHERE customer_id = :c", {"c": CUSTOMER}), + ("DELETE FROM core_share_lot WHERE customer_id = :c", {"c": CUSTOMER}), + ("DELETE FROM core_holding WHERE customer_id = :c", {"c": CUSTOMER}), + ("DELETE FROM core_customer_risk WHERE customer_id = :c", {"c": CUSTOMER}), + ("DELETE FROM core_fee_rule WHERE product_id LIKE 'PROD-T12C%'", {}), + ("DELETE FROM core_product_nav WHERE product_id LIKE 'PROD-T12C%'", {}), + ("DELETE FROM core_product WHERE product_id LIKE 'PROD-T12C%'", {}), + ("DELETE FROM core_customer WHERE customer_id = :c", {"c": CUSTOMER}), + ]: + conn.execute(text(sql), params) + with agent_engine.begin() as conn: + conn.execute( + text( + "DELETE FROM risk_convert_detail " + "WHERE convert_group_id LIKE 'CNV-T12C%' OR client_request_id LIKE 'T12C-%'" + ) + ) + conn.execute(text("DELETE FROM risk_alert WHERE customer_id = :c"), {"c": CUSTOMER}) + conn.execute( + text("DELETE FROM customer_profile_l3 WHERE customer_id = :c"), {"c": CUSTOMER} + ) + conn.execute(text("DELETE FROM audit_log WHERE trace_id LIKE 'T12C-VERIFY%'")) + + +def seed(core_engine) -> None: + with core_engine.begin() as conn: + conn.execute( + text( + "INSERT INTO core_customer (customer_id, display_name, open_date) " + "VALUES (:c, 'T12补偿真库验证', :d)" + ), + {"c": CUSTOMER, "d": TRADE_DATE}, + ) + conn.execute( + text( + "INSERT INTO core_customer_risk (customer_id, risk_code, evaluated_at, expires_at) " + "VALUES (:c, 'C3', :t, :exp)" + ), + {"c": CUSTOMER, "t": TRADE_AT - timedelta(days=30), + "exp": TRADE_AT + timedelta(days=300)}, + ) + for pid, name, ptype, rate in [ + (PROD_OUT, "T12转出基金", "bond", OUT_RATE), + (PROD_IN, "T12转入基金", "stock", IN_RATE), + ]: + conn.execute( + text( + "INSERT INTO core_product (product_id, product_name, min_risk_code, " + "product_type, can_subscribe, can_redeem, subscribe_fee_rate, " + "fund_company, ta_code) " + "VALUES (:p, :n, 'R2', :t, 1, 1, :r, :co, :ta)" + ), + {"p": pid, "n": name, "t": ptype, "r": str(rate), "co": COMPANY, "ta": TA}, + ) + for mh, mh_max, rate in FEE_TIERS: + conn.execute( + text( + "INSERT INTO core_fee_rule (product_id, fee_type, min_hold_days, " + "max_hold_days, rate) VALUES (:p, 'redeem', :mh, :mh_max, :rate)" + ), + {"p": PROD_OUT, "mh": mh, "mh_max": mh_max, "rate": rate}, + ) + conn.execute( + text( + "INSERT INTO core_product_nav (product_id, nav, daily_chg_pct, nav_date) " + "VALUES (:p, :nav, 0, :d)" + ), + {"p": PROD_IN, "nav": str(IN_NAV), "d": TRADE_DATE}, + ) + for lot_id, qty, nav, confirmed in [ + ("LOT-T12C-A1", "400000", "1.0300", datetime(2025, 8, 1, 10, 0, 0)), + ("LOT-T12C-A2", "100000", "1.0000", datetime(2026, 9, 1, 10, 0, 0)), + ]: + conn.execute( + text( + "INSERT INTO core_share_lot (lot_id, customer_id, product_id, qty, " + "remain_qty, nav, confirmed_at) VALUES (:l, :c, :p, :q, :q, :nav, :cat)" + ), + {"l": lot_id, "c": CUSTOMER, "p": PROD_OUT, "q": qty, "nav": nav, + "cat": confirmed}, + ) + conn.execute( + text( + "INSERT INTO core_holding (customer_id, product_id, qty, cost_amount, " + "market_value, pnl_pct, as_of) VALUES (:c, :p, 500000, 512000.00, 512000.00, 0, :d)" + ), + {"c": CUSTOMER, "p": PROD_OUT, "d": TRADE_DATE}, + ) + + +def _req(cid_req: str | None = CID_REQ) -> dict: + return { + "customer_id": CUSTOMER, + "from_product_id": PROD_OUT, + "to_product_id": PROD_IN, + "qty": "120000", + "client_request_id": cid_req, + } + + +def _boom_engine(out_trade, in_trade): + raise RuntimeError("T12C 故意:阶段 1.5 引擎失败") + + +# ── 断言主体 ──────────────────────────────────────────────────────── +def build_pending_state(core, repo, crepo) -> str: + """构造待补偿态:阶段二写详情失败 + 阶段 1.5 引擎失败,但 Core 侧两条流水已成立。""" + print("\n[A] 待补偿态构造(阶段二 + 阶段 1.5 双失败)") + with patch.object( + ConvertRepository, "complete_convert", side_effect=RuntimeError("T12C 故意:阶段二失败") + ): + resp = convert_fund( + _req(), + core_ro=core, + risk_repo=repo, + convert_repo=crepo, + thresholds=TH, + now=NOW, + engine_hook=_boom_engine, + id_factory=_t12c_id, + ) + gid = resp["convert_group_id"] + trades = rows( + core_engine, "SELECT * FROM core_trade WHERE convert_group_id = :g", g=gid + ) + check("Core 侧两条流水已成立(交易不可回滚)", len(trades), 2) + check( + "两条流水 trade_type 齐备", + sorted(t["trade_type"] for t in trades), + ["redeem", "subscribe"], + ) + detail = rows( + agent_engine, "SELECT * FROM risk_convert_detail WHERE convert_group_id = :g", g=gid + ) + check("占位已置 failed(待补偿)", detail[0]["status"] if detail else None, "failed") + check( + "此时本客户无预警单", + q1( + agent_engine, + "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", + c=CUSTOMER, + ), + 0, + ) + check( + "验收 17 前置:convert_detail_write_failed 审计存在", + q1( + agent_engine, + "SELECT COUNT(*) FROM audit_log WHERE decision = 'convert_detail_write_failed' " + "AND trace_id LIKE 'T12C-VERIFY%'", + ), + 1, + ) + check( + "阶段 1.5 失败审计(decision='engine_error')存在", + q1( + agent_engine, + "SELECT COUNT(*) FROM audit_log WHERE decision = 'engine_error' " + "AND trace_id LIKE 'T12C-VERIFY%'", + ), + 1, + ) + return gid + + +def check_cleanup(core, repo, crepo) -> None: + """B 组:超时 pending → expired(ENUM 值域 / 时间边界 / 幂等 / 不硬删)。""" + print("\n[B] 超时 pending 占位 → expired(cleanup_pending_convert 同款仓储路径)") + orphan = f"CNV-T12C-ORPHAN-{uuid.uuid4().hex[:6].upper()}" + crepo.insert_placeholder(orphan, f"T12C-ORPHAN-{uuid.uuid4().hex[:6].upper()}") + + cand = [r["convert_group_id"] for r in crepo.list_expired_candidates(24)] + check("刚落下的 pending 不算超时(时间边界下界)", orphan in cand, False) + + # 回拨 created_at 到 30h 前 → 应进候选 + with agent_engine.begin() as conn: + conn.execute( + text( + "UPDATE risk_convert_detail SET created_at = :t WHERE convert_group_id = :g" + ), + {"t": datetime.now() - timedelta(hours=30), "g": orphan}, + ) + cand = [r["convert_group_id"] for r in crepo.list_expired_candidates(24)] + check("回拨 30h 后进候选(时间边界上界)", orphan in cand, True) + + crepo.mark_expired(orphan) + check( + "真库 ENUM 接受 status='expired'", + q1( + agent_engine, + "SELECT status FROM risk_convert_detail WHERE convert_group_id = :g", + g=orphan, + ), + "expired", + ) + cand = [r["convert_group_id"] for r in crepo.list_expired_candidates(24)] + check("复跑不再进候选(幂等)", orphan in cand, False) + check( + "S2 标记不硬删:行仍在", + q1( + agent_engine, + "SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id = :g", + g=orphan, + ), + 1, + ) + + +def check_compensate(gid: str, core, repo, crepo) -> None: + """C/D/E 组:补偿端到端 + 幂等 + 真 JSON 列 LIKE。""" + print("\n[C] compensate_convert 端到端(详情 + 预警一次补齐)") + out = compensate_convert( + gid, + core_ro=core, + risk_repo=repo, + convert_repo=crepo, + thresholds=TH, + now=NOW, + ) + check("state", out["state"], "rebuilt") + check("detail 分支", out["detail"], "rebuilt") + check("engine 分支", out["engine"], "rebuilt") + check("补偿后出一张单", len(out["alert_ids"]), 1) + check("命中 RISK-001", "RISK-001" in out["triggered_rules"], True) + + detail = rows( + agent_engine, "SELECT * FROM risk_convert_detail WHERE convert_group_id = :g", g=gid + )[0] + check("详情置 completed", detail["status"], "completed") + check("回填 out_trade_id 与 Core 一致", detail["out_trade_id"], out["out_trade_id"]) + check("回填 in_trade_id 与 Core 一致", detail["in_trade_id"], out["in_trade_id"]) + check("回填 nav 非空", detail["nav"] is not None, True) + check("回填 nav_date 非空", detail["nav_date"] is not None, True) + check( + "主审计(convert_accepted)已补", + q1( + agent_engine, + "SELECT COUNT(*) FROM audit_log WHERE decision = 'convert_accepted' " + "AND trace_id LIKE 'T12C-VERIFY%'", + ), + 1, + ) + + alerts = rows(agent_engine, "SELECT * FROM risk_alert WHERE customer_id = :c", c=CUSTOMER) + check("预警单 1 张", len(alerts), 1) + payload = json.loads(alerts[0]["payload"]) if isinstance(alerts[0]["payload"], str) \ + else alerts[0]["payload"] + check("一张单两条事件(转出 + 转入,真 JSON 列反解)", len(payload["events"]), 2) + + print("\n[D] 重复补偿 → skipped(幂等,不产生第二张单)") + again = compensate_convert( + gid, core_ro=core, risk_repo=repo, convert_repo=crepo, thresholds=TH, now=NOW + ) + check("state", again["state"], "skipped") + check("detail 分支", again["detail"], "already_completed") + check("engine 分支", again["engine"], "skipped") + check("alert_ids 与首次一致", again["alert_ids"], out["alert_ids"]) + check( + "预警单仍只 1 张", + q1(agent_engine, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", c=CUSTOMER), + 1, + ) + check( + "主审计不重复", + q1( + agent_engine, + "SELECT COUNT(*) FROM audit_log WHERE decision = 'convert_accepted' " + "AND trace_id LIKE 'T12C-VERIFY%'", + ), + 1, + ) + check("幂等跳过时附人工核对提示", "人工核对" in (again.get("warning") or ""), True) + + print("\n[E] 真 JSON 列上 input_summary LIKE 生效(has_engine_error_audit)") + col_type = q1( + agent_engine, + "SELECT DATA_TYPE FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = DATABASE() " + "AND TABLE_NAME = 'audit_log' AND COLUMN_NAME = 'input_summary'", + ) + check("前提:input_summary 真库列类型为 json", str(col_type).lower(), "json") + check( + "按 group_id 命中(sqlite 是 TEXT,此断言只在真库有意义)", + repo.has_engine_error_audit(gid, decision="engine_error"), + True, + ) + check( + "决策码不匹配则不命中(证明判定真的按 decision 过滤)", + repo.has_engine_error_audit(gid, decision="risk_engine_error"), + False, + ) + + +def main() -> int: + global core_engine, agent_engine + core_engine = get_engine(settings.mysql_core_database, "admin") + agent_engine = get_engine(settings.mysql_database, "admin") + print("T-12 真库验证:补偿脚本(详情 + 预警)+ 超时占位清理") + print(f"客户={CUSTOMER} 产品={PROD_OUT}/{PROD_IN} 交易日={TRADE_DATE}") + + cleanup(core_engine, agent_engine) + new_trace(TRACE) # 固定 trace:审计行才能按 'T12C-VERIFY%' 精准清理与断言 + try: + seed(core_engine) + core = CoreReadOnlyRepository(engine=core_engine) + # 业务路径用生产默认账号(ro 读 / rw 写),仅建清数据用 admin + core_ro = CoreReadOnlyRepository(engine=get_engine(settings.mysql_core_database, "ro")) + repo = RiskRepository(engine=get_engine(settings.mysql_database, "rw")) + crepo = ConvertRepository(engine=get_engine(settings.mysql_database, "rw")) + + gid = build_pending_state(core_ro, repo, crepo) + check_cleanup(core_ro, repo, crepo) + check_compensate(gid, core_ro, repo, crepo) + finally: + cleanup(core_engine, agent_engine) + left = q1( + core_engine, + "SELECT (SELECT COUNT(*) FROM core_trade WHERE customer_id LIKE 'CUST-T12C%') " + "+ (SELECT COUNT(*) FROM core_share_lot WHERE customer_id LIKE 'CUST-T12C%') " + "+ (SELECT COUNT(*) FROM core_holding WHERE customer_id LIKE 'CUST-T12C%') " + "+ (SELECT COUNT(*) FROM core_product WHERE product_id LIKE 'PROD-T12C%') " + "+ (SELECT COUNT(*) FROM core_customer WHERE customer_id LIKE 'CUST-T12C%')", + ) + q1( + agent_engine, + "SELECT (SELECT COUNT(*) FROM risk_convert_detail WHERE convert_group_id LIKE 'CNV-T12C%') " + "+ (SELECT COUNT(*) FROM risk_alert WHERE customer_id = 'CUST-T12C') " + "+ (SELECT COUNT(*) FROM audit_log WHERE trace_id LIKE 'T12C-VERIFY%')", + ) + print(f"\n清理后残留行数:{left}") + dispose_engines() + + print(f"\n结果:{_passed} 通过 / {_failed} 失败") + return 1 if _failed else 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_concentration_c4.py b/tests/test_concentration_c4.py index 006cc3f..ed96084 100644 --- a/tests/test_concentration_c4.py +++ b/tests/test_concentration_c4.py @@ -157,6 +157,43 @@ def test_concentration_profile_empty_customer(env): assert rule_concentration(profile, RiskThresholds()) is None +def test_concentration_profile_filters_zero_qty(env): + """归零行(convert 转出全部份额后保留的 qty=0 行)不构成持仓,须过滤。 + + T-11 收尾:原实现走独立 SQL、漏了 `qty > 0`,与 `list_holdings` 口径落差 —— + ratio 数值不受影响(归零行市值 0),但 `rows` 会多出已清仓产品、虚占截断位。 + """ + core, repo, pub, engine = env + with engine.begin() as conn: + conn.execute( + text( + "INSERT INTO core_product (product_id, product_name, min_risk_code," + " product_type) VALUES ('P3', '已清仓基金', 'R5', 'equity')" + ) + ) + conn.execute( + text( + "INSERT INTO core_holding (customer_id, product_id, market_value, qty," + " cost_amount, pnl_pct, as_of)" + " VALUES ('C1', 'P3', 0, 0, 800000, 0, '2026-09-05')" + ) + ) + profile = core.concentration_profile("C1") + assert profile["total_value"] == Decimal(1000000) # 不含归零行 + assert profile["r45_value"] == Decimal(900000) + assert profile["ratio"] == 0.9 + assert len(profile["rows"]) == 2 # 归零行 P3 未进明细 + assert profile["holdings_truncated"] is False + # 与 list_holdings 同口径(T-11 收尾的核心诉求):持仓集合与市值一致 + listed = core.list_holdings("C1") + assert len(listed) == len(profile["rows"]) + assert {(r["product_id"], r["market_value"]) for r in listed} == { + ("P2", Decimal(900000)), + ("P1", Decimal(100000)), + } + assert {r["product_id"] for r in listed} == {"P1", "P2"} + + # ---------- 引擎接入 ---------- diff --git a/tests/test_convert_service.py b/tests/test_convert_service.py index e3f8f60..cc30ffd 100644 --- a/tests/test_convert_service.py +++ b/tests/test_convert_service.py @@ -38,7 +38,11 @@ from app.service.convert.calc import ( lot_fee, plan_lots, ) -from app.service.convert.convert_service import PROCESSING, convert_fund +from app.service.convert.convert_service import ( + PROCESSING, + compensate_convert, + convert_fund, +) from app.service.convert.errors import ( BelowMinQty, CrossEntityNotSupported, @@ -50,6 +54,7 @@ from app.service.convert.errors import ( ) from app.service.convert.fee import pick_fee_rate from app.service.convert.types import FeeRule, Lot +from app.service.risk.locks import try_lock from app.service.risk.rules import RiskThresholds CUST = "CUST-T7" # C3 客户 → R4 产品 allowed_with_disclosure(放行) @@ -468,3 +473,150 @@ def test_engine_wired_produces_single_alert_with_two_events(sqlite_engine): payload = json.loads(alerts[0]["payload"]) assert len(payload["events"]) == 2, "一张单承载两条事件(转出 + 转入)" assert [e["trade_type"] for e in payload["events"]] == ["redeem", "subscribe"] + + +# ── 11. T-12 补偿:阶段二失败 → 按 group 补写详情 + 预警(FR-C17 / 架构 §5.4)── +def _low_thresholds() -> RiskThresholds: + """调低大额阈值,让 120 份的折算额足以命中 RISK-001(与 §10 接线用例同口径)。""" + return RiskThresholds( + large_amount=Decimal("100"), + daily_total=Decimal("1000000"), + freq_count=3, + probe_window_minutes=5, + probe_count=3, + probe_amount=Decimal("400000"), + small_amount=Decimal("10000"), + small_count=3, + concentration_threshold=1.01, + ) + + +def _prepare_failed_convert(sqlite_engine, monkeypatch, cid_req: str) -> str: + """构造**真实待补偿态**:Core 侧两条流水已成立,但详情与预警双双缺失。 + + - 阶段二:`complete_convert` 打桩抛异常 → 占位 `failed` + `convert_detail_write_failed` 审计; + - 阶段 1.5:注入会抛的 `engine_hook` → 无预警单 + `engine_error` 审计。 + + 这正是 PRD §7.1 描述的补偿场景(agent 侧两件事都没落)。 + 返回 convert_group_id;`complete_convert` 已被还原(补偿才能真的写进去)。 + """ + _seed(sqlite_engine) + original = ConvertRepository.complete_convert + + def boom(self, group_id, **kwargs): # noqa: ANN001 + raise RuntimeError("阶段二写失败") + + def engine_boom(out_trade, in_trade): # noqa: ANN001 + raise RuntimeError("阶段 1.5 引擎失败") + + monkeypatch.setattr(ConvertRepository, "complete_convert", boom) + resp = convert_fund( + _req(cid_req=cid_req), + now=NOW, + engine_hook=engine_boom, + thresholds=_low_thresholds(), + **_services(sqlite_engine), + ) + monkeypatch.setattr(ConvertRepository, "complete_convert", original) + return resp["convert_group_id"] + + +def _compensate(sqlite_engine, gid: str) -> dict: + svc = _services(sqlite_engine) + return compensate_convert( + gid, + core_ro=svc["core_ro"], + risk_repo=svc["risk_repo"], + convert_repo=svc["convert_repo"], + thresholds=_low_thresholds(), + now=NOW, + ) + + +def _detail_rows(engine, gid: str) -> list[dict]: + return _rows( + engine, "SELECT * FROM risk_convert_detail WHERE convert_group_id = :g", g=gid + ) + + +def test_compensate_rebuilds_detail_and_alerts(sqlite_engine, monkeypatch): + """补偿把「详情 + 预警」两件事都补齐,且详情直接置 completed(验收 17)。""" + gid = _prepare_failed_convert(sqlite_engine, monkeypatch, "CLI-T12-001") + # 前置态自检:交易已成立、详情 failed、无预警单 + assert len(_rows(sqlite_engine, "SELECT * FROM core_trade")) == 2 + assert _detail_rows(sqlite_engine, gid)[0]["status"] == "failed" + assert _rows(sqlite_engine, "SELECT * FROM risk_alert") == [] + assert len( + _rows( + sqlite_engine, + "SELECT * FROM audit_log WHERE decision = 'convert_detail_write_failed'", + ) + ) == 1 + + out = _compensate(sqlite_engine, gid) + + assert out["state"] == "rebuilt" + assert out["detail"] == "rebuilt" and out["engine"] == "rebuilt" + assert "RISK-001" in out["triggered_rules"] + assert len(out["alert_ids"]) == 1 + row = _detail_rows(sqlite_engine, gid)[0] + assert row["status"] == "completed" + # 补偿回填的指针与 Core 侧一致(幂等锚点 = 转出端) + assert row["out_trade_id"] == out["out_trade_id"] + assert row["in_trade_id"] == out["in_trade_id"] + assert row["nav"] is not None and row["nav_date"] is not None + assert len(_rows(sqlite_engine, "SELECT * FROM risk_alert")) == 1 + # 补偿也补了主审计(与首次成功路径同一份 `_write_main_audit`) + assert len( + _rows( + sqlite_engine, + "SELECT * FROM audit_log WHERE decision = 'convert_accepted'", + ) + ) == 1 + + +def test_compensate_is_idempotent(sqlite_engine, monkeypatch): + """重复补偿 → `skipped`,不产生第二张预警单、不重复写详情与主审计。""" + gid = _prepare_failed_convert(sqlite_engine, monkeypatch, "CLI-T12-002") + first = _compensate(sqlite_engine, gid) + second = _compensate(sqlite_engine, gid) + + assert first["state"] == "rebuilt" + assert second["state"] == "skipped" + assert second["detail"] == "already_completed" + assert second["engine"] == "skipped" + assert second["alert_ids"] == first["alert_ids"] + assert second["triggered_rules"] == [] + assert len(_rows(sqlite_engine, "SELECT * FROM risk_alert")) == 1 # 只有一张单 + assert len( + _rows( + sqlite_engine, + "SELECT * FROM audit_log WHERE decision = 'convert_accepted'", + ) + ) == 1 # 主审计不重复 + # 曾阶段 1.5 中断过 → 幂等跳过时附「人工核对」提示(与 rebuild_alerts 同款警示) + assert "人工核对" in second["warning"] + + +def test_compensate_missing_group_touches_nothing(sqlite_engine): + """Core 侧不足两条流水 → missing,零写入(不是有效转换组)。""" + _seed(sqlite_engine) + out = _compensate(sqlite_engine, "CNV-T12-NOPE") + assert out["state"] == "missing" + assert out["trade_count"] == 0 + assert _rows(sqlite_engine, "SELECT * FROM risk_alert") == [] + assert _rows(sqlite_engine, "SELECT * FROM audit_log") == [] + assert _rows(sqlite_engine, "SELECT * FROM risk_convert_detail") == [] + + +def test_compensate_returns_locked_when_execution_right_taken( + sqlite_engine, monkeypatch +): + """未抢到 `convert:rerun:{gid}` → locked(有并发重试/实例在跑),零写入。""" + gid = _prepare_failed_convert(sqlite_engine, monkeypatch, "CLI-T12-004") + with try_lock(f"convert:rerun:{gid}", 30) as acquired: + assert acquired is True + out = _compensate(sqlite_engine, gid) + assert out["state"] == "locked" and out["alert_ids"] == [] + assert _detail_rows(sqlite_engine, gid)[0]["status"] == "failed" # 仍待补偿 + assert _rows(sqlite_engine, "SELECT * FROM risk_alert") == [] diff --git a/tests/test_demo_scripts.py b/tests/test_demo_scripts.py index 7fbd138..b215a76 100644 --- a/tests/test_demo_scripts.py +++ b/tests/test_demo_scripts.py @@ -2,11 +2,13 @@ rebuild_alerts:补偿重放出单并推送、幂等跳过(防重复 append/aml 重复出单)、 missing 不落库;subscribe_alerts:payload → 单行可读文本。 +**T-12 追加**:`rebuild_alerts --convert-group` 的分派与接线、 +`cleanup_pending_convert` 超时占位清理(标记不硬删)。 脚本目录非包,动态入 sys.path 后按模块名导入。 """ import sys -from datetime import datetime +from datetime import datetime, timedelta from decimal import Decimal from pathlib import Path @@ -16,16 +18,21 @@ from sqlalchemy import text from _ddl import create_sqlite_engine DEMO_DIR = Path(__file__).resolve().parents[1] / "scripts" / "demo" +ROOT = Path(__file__).resolve().parents[1] sys.path.insert(0, str(DEMO_DIR)) +sys.path.insert(0, str(ROOT / "scripts" / "agent")) -from rebuild_alerts import rebuild_trade # noqa: E402 +import rebuild_alerts # noqa: E402 +from rebuild_alerts import main, rebuild_convert_group, rebuild_trade # noqa: E402 from subscribe_alerts import format_alert # noqa: E402 +from cleanup_pending_convert import cleanup # noqa: E402 -from app.repository.core_ro import CoreReadOnlyRepository -from app.repository.risk_repository import RiskRepository -from app.service.risk import alert_service -from app.service.risk.alert_service import handle_alert -from app.service.risk.engine import process_trade_event +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.risk import alert_service # noqa: E402 +from app.service.risk.alert_service import handle_alert # noqa: E402 +from app.service.risk.engine import process_trade_event # noqa: E402 class FakePublisher: @@ -211,3 +218,139 @@ def test_format_alert_renders_payload_fields(): assert frag in line assert "CUST-9527" not in line # payload 只有脱敏掩码,无原始 id assert format_alert({"raw": "not-a-json-dict"}) == "[raw] not-a-json-dict" + + +# ---------- T-12 · rebuild_alerts --convert-group 分派与接线 ---------- + + +def test_convert_group_delegates_to_service(monkeypatch): + """`--convert-group` 只是薄封装:所有逻辑必须落在 compensate_convert(不复制实现)。""" + seen: dict = {} + + def fake(group_id, **kwargs): + seen["gid"] = group_id + seen["kwargs"] = kwargs + return {"state": "skipped", "alert_ids": ["ALT-9"]} + + monkeypatch.setattr(rebuild_alerts, "compensate_convert", fake) + out = rebuild_convert_group("CNV-T12-1", "CORE", "REPO", "CREPO") + assert out["state"] == "skipped" + assert seen["gid"] == "CNV-T12-1" + assert seen["kwargs"] == {"core_ro": "CORE", "risk_repo": "REPO", "convert_repo": "CREPO"} + + +def test_convert_group_missing_is_reported_without_writes(env): + """未知 group(Core 侧无流水)→ missing,且不写任何东西。""" + core, repo, pub, engine = env + out = rebuild_convert_group("CNV-T12-NOPE", core, repo, ConvertRepository(engine=engine)) + assert out["state"] == "missing" and out["trade_count"] == 0 + assert _counts(engine, "risk_alert") == 0 + assert _counts(engine, "audit_log") == 0 + + +def test_cli_rejects_convert_group_with_trade_ids(monkeypatch): + """`--convert-group` 与 trade_id 位置参数互斥(argparse 用法错误 → 退出码 2)。""" + monkeypatch.setattr( + sys, "argv", ["rebuild_alerts.py", "TRD-X", "--convert-group", "CNV-Y"] + ) + with pytest.raises(SystemExit) as exc: + main() + assert exc.value.code == 2 + + +def test_cli_requires_one_of_the_two_forms(monkeypatch): + """两者都不给 → 用法错误(不是静默什么都不做)。""" + monkeypatch.setattr(sys, "argv", ["rebuild_alerts.py"]) + with pytest.raises(SystemExit) as exc: + main() + assert exc.value.code == 2 + + +# ---------- T-12 · cleanup_pending_convert:超时占位 → expired ---------- + + +def _seed_convert_detail(engine, rows) -> None: + """rows: (group_id, status, created_at)。""" + with engine.begin() as conn: + for gid, status, created in rows: + conn.execute( + text( + "INSERT INTO risk_convert_detail" + " (convert_group_id, status, estimated, created_at)" + " VALUES (:g, :s, 0, :c)" + ), + {"g": gid, "s": status, "c": created}, + ) + + +def _status_map(engine) -> dict: + with engine.connect() as conn: + return { + r["convert_group_id"]: r["status"] + for r in conn.execute( + text("SELECT convert_group_id, status FROM risk_convert_detail") + ).mappings() + } + + +def test_cleanup_expires_only_overdue_pending(sqlite_engine): + """只动超时的 `pending`:fresh pending / completed / failed 一律不碰。""" + overdue = datetime.now() - timedelta(hours=30) + fresh = datetime.now() - timedelta(hours=1) + _seed_convert_detail( + sqlite_engine, + [ + ("CNV-OLD", "pending", overdue), + ("CNV-NEW", "pending", fresh), + ("CNV-DONE", "completed", overdue), + ("CNV-FAIL", "failed", overdue), + ], + ) + summary = cleanup(ConvertRepository(engine=sqlite_engine), 24) + + assert summary["sla_hours"] == 24 and summary["dry_run"] is False + assert summary["candidate_count"] == 1 + assert summary["expired"] == ["CNV-OLD"] + assert _status_map(sqlite_engine) == { + "CNV-OLD": "expired", + "CNV-NEW": "pending", + "CNV-DONE": "completed", + "CNV-FAIL": "failed", + } + # S2:标记不硬删 —— 行仍在(4 行一行不少) + assert _counts(sqlite_engine, "risk_convert_detail") == 4 + + +def test_cleanup_idempotent_second_run_finds_nothing(sqlite_engine): + """复跑 → 候选为空(已 expired 的行不再进候选),行数不变。""" + _seed_convert_detail( + sqlite_engine, [("CNV-OLD", "pending", datetime.now() - timedelta(hours=30))] + ) + repo = ConvertRepository(engine=sqlite_engine) + first = cleanup(repo, 24) + second = cleanup(repo, 24) + assert first["expired"] == ["CNV-OLD"] + assert second["candidate_count"] == 0 and second["expired"] == [] + assert _status_map(sqlite_engine) == {"CNV-OLD": "expired"} + assert _counts(sqlite_engine, "risk_convert_detail") == 1 + + +def test_cleanup_dry_run_writes_nothing(sqlite_engine): + """--dry-run 只报告候选,状态不动。""" + _seed_convert_detail( + sqlite_engine, [("CNV-OLD", "pending", datetime.now() - timedelta(hours=30))] + ) + summary = cleanup(ConvertRepository(engine=sqlite_engine), 24, dry_run=True) + assert summary["dry_run"] is True + assert summary["candidate_count"] == 1 and summary["expired_count"] == 0 + assert summary["candidates"][0]["convert_group_id"] == "CNV-OLD" + assert _status_map(sqlite_engine) == {"CNV-OLD": "pending"} + + +def test_cleanup_honours_custom_sla(sqlite_engine): + """--hours 覆盖 SLA:2h 前落下的占位在 1h 口径下不算超时、在 3h 口径下算。""" + created = datetime.now() - timedelta(hours=2) + _seed_convert_detail(sqlite_engine, [("CNV-MID", "pending", created)]) + repo = ConvertRepository(engine=sqlite_engine) + assert cleanup(repo, 3, dry_run=True)["candidate_count"] == 0 + assert cleanup(repo, 1, dry_run=True)["candidate_count"] == 1