diff --git a/app/service/risk/alert_service.py b/app/service/risk/alert_service.py index 188acc1..3864729 100644 --- a/app/service/risk/alert_service.py +++ b/app/service/risk/alert_service.py @@ -93,16 +93,43 @@ def _publish_alert( redis_gateway.publish("risk:pub:alert", body) +def build_trade_event(trade: dict[str, Any], hits: list[RuleHit]) -> dict[str, Any]: + """构造预警 `payload.events` 的单条事件体(trade 为 core_trade 行或等价 dict)。 + + T-8 抽出为**公开函数**:一次基金转换要落**两条**事件(转出端 + 转入端), + 由本函数统一构造 —— 避免 `engine.process_convert_event` 另写一份副本 + (自检第 13 问:同一结构只留一个副本)。 + """ + return { + "trade_id": trade["trade_id"], + "product_id": trade["product_id"], + "trade_type": trade["trade_type"], + "amount": str(trade["amount"]), + "traded_at": str(trade["traded_at"]), + "matched_rules": [h.rule_id for h in hits], + "details": {h.rule_id: h.detail for h in hits}, + } + + def record_trade_alerts( trade: dict[str, Any], hits: list[RuleHit], risk_repo: RiskRepository | None = None, customer_context: dict[str, Any] | None = None, + events: list[dict[str, Any]] | None = None, ) -> dict[str, Any] | None: """交易事件预警入口(引擎 B4 调用)。 hits 为空 → 审计 pass;否则事件类命中聚合为一张单(PRD FR-4:同客户同日仅一张事件类 pending 单,alert_type 随最高分规则动态更新)。 + + `events`(T-8 新增,可选):本事件关联的**事件体列表**(由 `build_trade_event` 构造)。 + 缺省 `None` → 退化为单条「主流水」事件,**既有调用零改动**。 + 一次基金转换传 `[转出端, 转入端]` → **一张单、`payload.events` 两条**(PRD §6.2)。 + 约定:`trade`(主流水)必须是 `events[0]`,审计 `input_summary` 取首条。 + + 聚合追加分支(当日已有同客户 pending 单)**只追加首条**(主事件):与架构 §5.4 + 「一次转换两条流水,只认转出端」同口径 —— 否则同一转换会在同一张单里重复出现两次。 """ risk_repo = risk_repo or RiskRepository() if not hits: @@ -117,15 +144,8 @@ def record_trade_alerts( # 空集不注入 payload.alert_subtype(评审 P2-3:老单结构保持干净)。 subtypes = sorted({h.alert_subtype for h in hits if h.alert_subtype}) publish_extra = {"alert_subtype": subtypes} if subtypes else None - event = { - "trade_id": trade["trade_id"], - "product_id": trade["product_id"], - "trade_type": trade["trade_type"], - "amount": str(trade["amount"]), - "traded_at": str(trade["traded_at"]), - "matched_rules": [h.rule_id for h in hits], - "details": {h.rule_id: h.detail for h in hits}, - } + event_list = events if events is not None else [build_trade_event(trade, hits)] + event = event_list[0] # 主事件(convert 场景 = 转出端):供审计 input_summary 落痕 def _agg(locked: bool) -> dict[str, Any]: if locked: @@ -165,7 +185,7 @@ def record_trade_alerts( "status": "pending_review", "payload": { "product_id": trade["product_id"], - "events": [event], + "events": event_list, "customer_context": customer_context or {}, }, } diff --git a/app/service/risk/engine.py b/app/service/risk/engine.py index b973ac1..43184fa 100644 --- a/app/service/risk/engine.py +++ b/app/service/risk/engine.py @@ -10,22 +10,32 @@ RISK-004 窗口以 trade["traded_at"] 为事件时点(非墙钟 now):rebui 客户事件钩子 on_customer_created/on_customer_updated 为 FR-5 预留(本期 no-op, 模拟环境无开户流程)。 + +T-8 追加:`process_convert_event`(基金转换阶段 1.5 入口,两条流水**只跑一次**) +与 `process_trade_event` 共用内部 `_run()`;`process_trade_event` 签名与行为不变。 """ from __future__ import annotations +import logging from datetime import datetime, timedelta from decimal import Decimal -from typing import Any +from typing import Any, Callable from app.repository.core_ro import CoreReadOnlyRepository from app.repository.risk_repository import RiskRepository -from app.service.risk.alert_service import record_aml_alert, record_trade_alerts +from app.service.risk.alert_service import ( + build_trade_event, + record_aml_alert, + record_trade_alerts, +) from app.service.risk.aml_service import match_customer from app.service.risk.profile_l3 import upsert_profile_l3 from app.service.risk.rules import RiskThresholds, rule_concentration, run_rules from app.utils.trace import current_trace, ensure_trace, new_trace +logger = logging.getLogger(__name__) + def _as_datetime(value: Any) -> datetime: if isinstance(value, datetime): @@ -111,20 +121,38 @@ def _audit_concentration( ) -def process_trade_event( - trade: dict[str, Any], - core_ro: CoreReadOnlyRepository | None = None, - risk_repo: RiskRepository | None = None, - thresholds: RiskThresholds | None = None, -) -> dict[str, Any]: - """处理一笔已落库交易(PRD FR-1 ②③b 之后)。 +def _notify_error_hook( + on_error_hook: Callable[[dict[str, Any], Exception], None] | None, + out_trade: dict[str, Any], + exc: Exception, +) -> None: + """D19 预留钩子:引擎异常时通知一次(一期传 `None` = 仅落 `engine_error` 审计)。 - 返回 {"triggered_rules": [...], "alert_ids": [...], "aml_hit": bool}, - 网关据此拼装响应(FR-1 ⑤:blocked=false + trade_id + 触发规则列表)。 + **hook 自身失败必须被吞掉**:与「阶段 1.5 不阻断交易」同原则 —— + 补偿通道的故障不得污染主流程的异常语义(原始异常仍按原样上抛)。 + """ + if on_error_hook is None: + return + try: + on_error_hook(out_trade, exc) + except Exception: # noqa: BLE001 + logger.exception("on_error_hook 调用失败(已忽略,不影响交易主流程)") + + +def _run( + trade: dict[str, Any], + *, + core: CoreReadOnlyRepository, + repo: RiskRepository, + th: RiskThresholds, + related_trades: list[dict[str, Any]] | None = None, +) -> dict[str, Any]: + """事件处理共用实现(`process_trade_event` 单流水 / `process_convert_event` 两条流水)。 + + `trade` = **主流水**(convert 场景为转出端 `redeem`):返回体、AML、L3 均以它为准。 + `related_trades` 仅决定预警 `payload.events` 的条数(convert 传 `[转出端, 转入端]` + → 一张单承载两条事件);缺省 `None` 退化为单条,行为与改造前**完全一致**。 """ - core = core_ro or CoreReadOnlyRepository() - repo = risk_repo or RiskRepository() - th = thresholds or RiskThresholds.from_settings() ensure_trace() # 脚本/重放入口兜底归因(中间件场景保留现有 trace,评审 P3-8) event_at = _as_datetime(trade["traded_at"]) @@ -147,11 +175,18 @@ def process_trade_event( hits = [*hits, conc] if hits: + # T-8:convert 两条流水 → 两个事件体(同一份 hits),一张单承载(PRD §6.2) + events = ( + None + if related_trades is None + else [build_trade_event(t, hits) for t in related_trades] + ) alert = record_trade_alerts( trade, hits, risk_repo=repo, customer_context=_build_customer_context(core, trade["customer_id"], day_start), + events=events, ) result["triggered_rules"] = sorted({h.rule_id for h in hits}) if alert: @@ -192,6 +227,60 @@ def process_trade_event( return result +def process_trade_event( + trade: dict[str, Any], + core_ro: CoreReadOnlyRepository | None = None, + risk_repo: RiskRepository | None = None, + thresholds: RiskThresholds | None = None, +) -> dict[str, Any]: + """处理一笔已落库交易(PRD FR-1 ②③b 之后)。 + + 返回 {"triggered_rules": [...], "alert_ids": [...], "aml_hit": bool}, + 网关据此拼装响应(FR-1 ⑤:blocked=false + trade_id + 触发规则列表)。 + **签名与行为自 T-8 起完全不变**(内部改为共用 `_run`,既有 13 处调用零改动)。 + """ + return _run( + trade, + core=core_ro or CoreReadOnlyRepository(), + repo=risk_repo or RiskRepository(), + th=thresholds or RiskThresholds.from_settings(), + ) + + +def process_convert_event( + out_trade: dict[str, Any], + in_trade: dict[str, Any], + *, + core_ro: CoreReadOnlyRepository | None = None, + risk_repo: RiskRepository | None = None, + thresholds: RiskThresholds | None = None, + on_error_hook: Callable[[dict[str, Any], Exception], None] | None = None, +) -> dict[str, Any]: + """处理一次基金转换(**只跑一次**,PRD §6.2 / FR-C11,架构 §6.2)—— 阶段 1.5 入口。 + + - 主流水 = **转出端** `out_trade`:返回体、L3、AML 与幂等锚点口径均取它; + - 两条流水(转出 `redeem` + 转入 `subscribe`)各成一个事件体 → + **一张单、`payload.events` 两条**(验收 7); + - 金额聚合去重(RISK-002/005 只看转出端)在 `run_rules` 内部完成, + 本函数不做二次过滤(PRD §6.3「去重只在金额聚合处,绝不删行」)。 + + `on_error_hook`(D19 预留,一期传 `None`):异常时调用一次,签名 `(out_trade, exc)`; + 二期可注入 Redis 事件 / 补偿队列。**hook 自身抛异常被吞掉,异常仍按原样上抛** —— + 由调用方(`convert_service` 阶段 1.5)落 `engine_error` 审计且**不阻断交易**(D17)。 + """ + try: + return _run( + out_trade, + core=core_ro or CoreReadOnlyRepository(), + repo=risk_repo or RiskRepository(), + th=thresholds or RiskThresholds.from_settings(), + related_trades=[out_trade, in_trade], + ) + except Exception as exc: # noqa: BLE001 + _notify_error_hook(on_error_hook, out_trade, exc) + raise + + def on_customer_created(customer_id: str) -> None: """AML 开户触发预留(本期 no-op;模拟环境无开户流程,PRD FR-5)。""" diff --git a/app/service/risk/rules.py b/app/service/risk/rules.py index 8c5ea8f..ba1dfb6 100644 --- a/app/service/risk/rules.py +++ b/app/service/risk/rules.py @@ -3,6 +3,9 @@ 输入约定:trades 为**当日 confirmed 的 subscribe/redeem 明细**(引擎组装;函数内再做防御过滤), amount 为 Decimal,traded_at 为 datetime;输出命中规则列表(RuleHit)。 阈值由 RiskThresholds 打包(默认来自 settings,可 .env 覆盖)。 + +T-8 起:一次基金转换落**两条**同 `convert_group_id` 的流水,金额聚合类规则 +(RISK-002/005)走 `_amount_view` 只计一次,逐笔类规则(RISK-001/003/004)仍看全量。 """ from __future__ import annotations @@ -88,6 +91,38 @@ def _amount(trade: dict[str, Any]) -> Decimal: return v if isinstance(v, Decimal) else Decimal(str(v)) +def _amount_view(trades: list[dict[str, Any]]) -> list[dict[str, Any]]: + """金额口径视图(架构 §6.1 / PRD §6.3 · T-8):一次转换只计**一次**金额。 + + 一次基金转换在 `core_trade` 落**两条**流水(转出 `redeem` + 转入 `subscribe`, + 同 `convert_group_id`),存储层如实记两份金额是对的,但**金额聚合类规则必须只算一次** —— + 否则 RISK-002 当日累计翻倍、RISK-005 把转入端重复当成「铺垫」。 + + - 同 `convert_group_id` 组内**只保留 `trade_type='redeem'` 那条**(转出端 = 金额基准); + - **无 `convert_group_id` 的交易原样通过** → 非 convert 场景恒等(既有断言零影响); + - 组内**无 `redeem`** 时保留首条(防御:R-b 规定 convert 必有转出端,理论不可达); + - **绝不删行**(只在金额视图内替换,`run_rules` 的其余规则仍看全量)。 + + 返回新列表,顺序与入参一致(命中替换不改变位置),故对无 gid 输入满足恒等。 + """ + view: list[dict[str, Any]] = [] + index_of_group: dict[str, int] = {} + for t in trades: + gid = t.get("convert_group_id") + if not gid: + view.append(t) + continue + idx = index_of_group.get(gid) + if idx is None: + index_of_group[gid] = len(view) + view.append(t) + continue + # 同组已占位:转出端优先(转入端不得覆盖已占位的转出端) + if t.get("trade_type") == "redeem" and view[idx].get("trade_type") != "redeem": + view[idx] = t + return view + + def rule_large_amount(trades: list[dict[str, Any]], th: RiskThresholds) -> RuleHit | None: """RISK-001 单笔大额:amount ≥ large_amount。""" for t in trades: @@ -203,17 +238,29 @@ def run_rules( thresholds: RiskThresholds | None = None, now: datetime | None = None, ) -> list[RuleHit]: - """引擎入口的规则编排:对单笔交易事件后的当日流水跑全部规则,返回全部命中。""" + """引擎入口的规则编排:对单笔交易事件后的当日流水跑全部规则,返回全部命中。 + + **两个视图(T-8 · 架构 §6.1)**: + - `eligible`(全量,两条都进)供 RISK-001 单笔大额 / RISK-003 频繁 / RISK-004 试探 + —— 逐笔判定,转入端也是真实权益变动,必须看得见; + - `amount_view`(同组只留转出端)供 RISK-002 当日累计 / RISK-005 先小后大 + —— 金额聚合口径,一次转换只计一次,防翻倍。 + RISK-006 读持仓快照,不经本函数(见 `rule_concentration`)。 + + RISK-003 的「分产品各归各」是刻意保留全量的原因:转出/转入是**不同产品**, + 归组后互不影响,一次转换不会被自己算成 2 笔频繁。 + """ th = thresholds or RiskThresholds.from_settings() now = now or datetime.now() eligible = [t for t in trades if _eligible(t)] if not eligible: return [] + amount_view = _amount_view(eligible) hits: list[RuleHit | None] = [ rule_large_amount(eligible, th), - rule_daily_total(eligible, th), + rule_daily_total(amount_view, th), rule_freq_trade(eligible, th), rule_probe_pattern(eligible, th, now), - rule_small_then_large(eligible, th), + rule_small_then_large(amount_view, th), ] return [h for h in hits if h is not None] diff --git a/docs/memory/2026-09-10.md b/docs/memory/2026-09-10.md index 0df3a61..047def5 100644 --- a/docs/memory/2026-09-10.md +++ b/docs/memory/2026-09-10.md @@ -225,3 +225,52 @@ SOP A-4:CUST-9527 / PROD-510300 / subscribe / 1000 **连发 4 笔**: 3. **交接文档 §B 顶部加 skill 指针**:本线设计方法论已固化,避免重复维护。 **效果**:方法论跨项目可复用(新项目直接加载 skill);本项目记忆回到「状态 + 红线 + 口径」本职。 + +--- + +## 深夜 · 基金转换线 T-8(规则引擎改造)完成 + +**范围**:开发计划 §7.1 —— `rules._amount_view` + `engine.process_convert_event` + `alert_service.events`(并行组 B)。 + +### 改码三处(均在代码注释留痕) + +- `app/service/risk/rules.py`:+`_amount_view(trades)`(同 `convert_group_id` 组内只留 `redeem`;**无 gid 恒等通过**;组内无 redeem 保首条);`run_rules` **双视图分流** —— `eligible`(全量)供 RISK-001/003/004,`_amount_view(eligible)` 供 RISK-002/005。 +- `app/service/risk/alert_service.py`:+**公开** `build_trade_event(trade, hits)`(结构定义唯一副本);`record_trade_alerts(..., events=None)` 缺省退化为单条,既有调用零改动;新建单落 `payload.events` 全部,**聚合追加只追首条**(与架构 §5.4「只认转出端」同口径)。 +- `app/service/risk/engine.py`:抽 `_run()` 共用实现;`process_trade_event` 变薄封装(**签名与行为不变**);新增 `process_convert_event(out_trade, in_trade, *, core_ro, risk_repo, thresholds, on_error_hook=None)` + `_notify_error_hook`(hook 自身异常吞掉、原始异常照常上抛)。 + +### ⭐ 实质影响:阶段 1.5 从「静默跳过」变「真跑」 + +T-7 落地时 `process_convert_event` 不存在 → `_run_engine` 走 `ImportError` 分支**跳过**; +**T-8 落地后同一笔转换会真实出单 + 写 `customer_profile_l3` + 落审计**。 +→ `verify_convert_service.py` 的执行语义随之变化,复跑仍 **35/35 绿**(无连带破坏)。 + +### 验证 + +- 新增 `tests/test_convert_engine.py` **15 用例**;`test_convert_service.py` **+1 条接线回归**(大额转换真出单,防 `_run_engine` 退回跳过)。 +- `pytest -q` → **672 passed / 3 skipped**(基线 656 **+16**,零回归)。 +- **突变验证(防假绿)**:临时把 `amount_view = eligible` → **4 条变红**,其中 `detail` 直接暴露 `当日申赎累计 599000 元(共 2 笔)` 的翻倍错误;恢复后复绿。 +- **新增 `scripts/dev/verify_convert_engine.py` 真库 31/31**(A 引擎真跑出单 / B **RISK-002 不翻倍夹逼**:单条 512000 < 800000 < 两条之和 1019975.33 / C 不删行 + DECIMAL 零漂移 / D 幂等重试不重复出单 / E 无命中只落 pass 审计 / F 六表零残留)。 + +### 真库脚本两条踩坑(与 T-7 同款 + 新一条,供后续 `verify_*.py` 参考) + +1. **`id_factory` 必须注入**:默认生成 `CNV-<日期>-`,与清理口径 `LIKE 'CNV-T8M%'` 不匹配 → 重跑遗留行撞 `uk_group`。注入 `_t8m_id` 并把清理条件补 `client_request_id LIKE 'T8M-%'` 兜底。 +2. **`core_trade.trade_type` 在 MySQL 是 ENUM,`ORDER BY` 按定义序**(实测 `[subscribe, redeem]`)→ 断言改用 `{trade_type: row}` 字典定位,不依赖排序。 + +### 发现的既有行为(不在本任务范围,**未顺手改**) + +> **更正(同日)**:初次记录误写成「RISK-001 实为当日累计口径」,**错误**。RISK-001 与 RISK-002 +> 是两条独立规则——前者**单笔**(`rules.py:128-133`)、后者**当日累计**(`:137-145`)。T-8 只切后者。 + +真实根因在**引擎入参范围**:`process_trade_event` 拉**当日全量**流水重跑规则 +(`engine.py` 的 `list_trades_range(customer, day_start, day_end)`)→ **触发这笔**与**命中那笔** +可以不是同一笔:真库 E 组实测,同客户同日先有一笔大额转换时,后续一笔 1000 元小额交易也带出 RISK-001。 +既有聚合逻辑(同客户同日一张 pending 单)兜住了重复出单,属**既有设计**,T-8 不改变也不修。 + +### 文档回写 + +开发计划(§0 速览标 ✅ + 测试基线改 672 · **§7.1 DoD 全勾 + 新增执行记录 + 3 条实施级收敛 + 2 条脚本坑**)· +根 `交接文档.md` §B → **v1.9**(§0 导航 · §B 头部 · §B.1 状态与代码改动 · **新增 §B.6.1 T-8 小节** · §B.6 任务树与基线)· +`docs/memory/{MEMORY,TODO}` · 本条。 + +**下一步 = T-9(`api/simulate.py` 模型与错误码 + `trade_gateway` convert 分派 · 关键路径)/ +T-11(`core_tools` 与 `sum_trades_on_date` 汇总去重 · 依赖 T-8 已解锁)**。 diff --git a/docs/memory/MEMORY.md b/docs/memory/MEMORY.md index 14bda68..db34e42 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.1 定稿**(33 条外审闭环 + 2 处架构回填 + **费率分类修正**)→ 架构 **v1.0 定稿**(独立评审通过,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 已于 2026-09-10 完成(609 passed / 3 skipped;T-1 断言 8/8 PASS;T-2 纯函数包 7 文件 / 93 用例;T-2b 实算脚本 15/15 一致)**,**下一步 = T-3(并行组 A:T-3 / T-4 / T-5 可同时开工)**。**开工前置两个阻断项(✅ 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-0b / T-1 已完成**):** `PRD-风控监测Agent` FR-1 一期显式拒收 convert,本线将其放开。**AIcoding 第 1~4 步已完成**:PRD **v0.9.1 定稿**(33 条外审闭环 + 2 处架构回填 + **费率分类修正**)→ 架构 **v1.0 定稿**(独立评审通过,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 已于 2026-09-10 完成(672 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(`api/simulate.py` + `trade_gateway` 分派 · 关键路径)/ T-11(汇总去重 · 依赖 T-8 已解锁)**。**开工前置两个阻断项(✅ 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-0b / T-1 / T-2 / T-2b 已完成,下一步 = T-3**;设计 + 开发计划均已闭环,入口 项目根 `交接文档.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` 只是远程跟踪引用失效,`git fetch` 即恢复,非远程分支被删)。入口 **项目根 `交接文档.md` §C**。 +**另(2026-09-10 待办)**:① **基金转换线**第 5 步进行中(**T-0~T-8 已完成、672 绿,下一步 = T-9(API+网关分派)/ T-11(汇总去重)**;设计 + 开发计划均已闭环,入口 项目根 `交接文档.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` 只是远程跟踪引用失效,`git fetch` 即恢复,非远程分支被删)。入口 **项目根 `交接文档.md` §C**。 **禁止(改代码前必记):** Core 正式 C1~C5 不可被画像覆盖 · 审计表只 INSERT · 代理人草稿不外发 · 仅 R-02 可阻断交易 · 四 Agent 不互调 LLM。 @@ -200,6 +200,6 @@ RBAC 联调账号:scripts/dev/rbac-seed-reference.md 3. 是否需 customer_id 归属与 JWT RBAC? 4. Core 是模拟库只读还是 agent 库读写? 5. 如何验证?(`python -m pytest` 全量(当前 **516 passed / 3 skipped**,基线 510)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) -6. **当前有哪两条并行线?**(① 风控/架构改进线:**已结项**(510 基线绿、§7.2 七项冒烟 7/7 PASS、`037ce7e` 已核实早已推送);② **基金转换线**:设计闭环,**第 5 步进行中 —— T-0 / T-0b / T-1 / T-2 / T-2b 已完成(609 passed),下一步 = T-3**)——动代码前先确认自己属于哪条线,别混淆前置条件。 +6. **当前有哪两条并行线?**(① 风控/架构改进线:**已结项**(510 基线绿、§7.2 七项冒烟 7/7 PASS、`037ce7e` 已核实早已推送);② **基金转换线**:设计闭环,**第 5 步进行中 —— T-0~T-8 已完成(672 passed),下一步 = T-9 / T-11**)——动代码前先确认自己属于哪条线,别混淆前置条件。 大任务:FRAMEWORK/FLOW 与实现状态不符时先更新 memory 再编码(用户确认跳过除外)。 diff --git a/docs/memory/TODO.md b/docs/memory/TODO.md index f56d16e..1a1e20f 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.1** + 架构 **v1.0** + 独立评审 13 条 **0 悬空**(接受 10 / 修正性接受 3 / 驳回 0),门控 **M-7 已满足**;**第 4 步开发计划 v1.0 已产出并经独立审核**(4 条意见全接受、**驳回 0**,含新增 2 条回归面 R15/R16 + R-c 双条修订);**第 5 步进行中**:**T-0 + T-0b + T-1 + T-2 + T-2b 均于 2026-09-10 完成**(**609 passed / 3 skipped**;T-1 断言 **8/8 PASS**、T-2 纯函数 **93 用例**、T-2b 实算 **15/15 一致**),**下一步 = T-3(并行组 A:T-3 / T-4 / T-5 可同时开工)**。两个阻断前置(**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.1** + 架构 **v1.0** + 独立评审 13 条 **0 悬空**(接受 10 / 修正性接受 3 / 驳回 0),门控 **M-7 已满足**;**第 4 步开发计划 v1.0 已产出并经独立审核**(4 条意见全接受、**驳回 0**,含新增 2 条回归面 R15/R16 + R-c 双条修订);**第 5 步进行中**:**T-0 + T-0b + T-1 + T-2 + T-2b + T-3 + T-4 + T-5 + T-6 + T-7 + T-8 均已于 2026-09-10 完成**(**672 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(`api/simulate.py` + `trade_gateway` 分派 · 关键路径)/ T-11(`core_tools` 与 `sum_trades_on_date` 汇总去重 · 依赖 T-8 已解锁)**。两个阻断前置(**T-0** sqlite/MySQL 列名统一 + 建库自校验 · **T-0b** DB 账号分离 D20:`xh_core_ro`/`xh_core_rw`/`xh_agent_rw`)**均已落地**。**入口:项目根 `交接文档.md` §B(三线合并版唯一入口,读这一节即可开工)**;**开工前必读开发计划 §1.4(15 条代码事实)+ §1.5(8 条实现级裁定 R-a~R-h)+ §12(16 条回归面)**。 ### 基金转换线待办(推荐顺序) @@ -16,7 +16,8 @@ - [x] **【T-0b · 阻断前置】DB 账号分离(D20)** —— **完成(2026-09-10)**:`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 - [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) -- [ ] **【T-3 起】** 仓储与锁(并行组 A:T-3 / T-4 / T-5 可同时开工)→ T-6 / T-7 → T-8 / T-9 / T-11 → **T-10 高风险单列** → T-12 → **T-13「50 并发压测 + 性能补录」**(最后跑,产出 PRD §9 第 18 条实测值)。**T-13 内部顺序**:先 50 并发压测 → 再性能实测补录 → 最后 PRD §5.3 数字回填(详见架构 §15 + 开发计划 §2~§10) +- [x] **【T-3 ~ T-8 已完成】** 仓储与锁(并行组 A:T-3 / T-4 / T-5 ✅)→ T-6 ✅(阶段一事务 · 真库 24/24)→ T-7 ✅(八步编排 · 真库 35/35)→ **T-8 ✅(规则引擎改造 · 真库 31/31;⭐ 阶段 1.5 从「跳过」变「真跑」)** +- [ ] **【T-9 起】剩余任务**:T-9(`api/simulate.py` 模型与错误码 + `trade_gateway` convert 分派)→ T-11(`core_tools` / `sum_trades_on_date` 汇总去重)→ **T-10 高风险单列** → T-12 → **T-13「50 并发压测 + 性能补录」**(最后跑,产出 PRD §9 第 18 条实测值)。**T-13 内部顺序**:先 50 并发压测 → 再性能实测补录 → 最后 PRD §5.3 数字回填(详见架构 §15 + 开发计划 §2~§10) > **第 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-3**(`core_ro` 五个新方法:`get_nav_as_of` / `get_redeem_fee_rules` / `list_share_lots` / `sum_remain_qty` / `get_holding` + 新增 `app/repository/share_lot_repository.py`;DoD 见开发计划 §5.1)。**并行组 A 的 T-3 / T-4 / T-5 可同时开工**。 diff --git a/docs/项目框架设计/开发计划-基金转换交易.md b/docs/项目框架设计/开发计划-基金转换交易.md index 473bfe1..69babf8 100644 --- a/docs/项目框架设计/开发计划-基金转换交易.md +++ b/docs/项目框架设计/开发计划-基金转换交易.md @@ -115,26 +115,26 @@ T-7 幂等窗口 · T-13 的 50 并发压测与性能补录 · PRD §5.3 实算 | --- | --- | --- | --- | --- | | **第 0 批 · 门禁** | **T-0** ✅ | sqlite/MySQL 结构对齐(`core_holding` 列名 + PK、补 `core_product_nav`)+ 建库自校验 + 门禁用例 —— **2026-09-10 完成** | 无 | **低(但阻断)** | | | **T-0b** ✅ | DB 账号分离(D20):`00-grant.sql` + `settings` 3 组账号 + `get_engine(db, role)` + 3 个权限断言 —— **2026-09-10 完成** | 无(可并行 T-0) | **低(但阻断)** | -| **第 1 批 · 数据与纯函数** | T-1 | MySQL DDL + 3 个种子 + `reset.ps1` + sqlite DDL 同步 | **T-0 + T-0b 双绿** | 低 | -| | T-2 | `service/convert/` 纯函数包(`types`/`calc`/`fee`/`nav`/`lot_bootstrap`/`errors`) | 无 | 低 | -| | T-2b | `calc_convert_demo.py` 实算 + 回填 PRD §5.3 与验收断言 | T-2 | 低 | -| **第 2 批 · 仓储与锁** | T-3 | `core_ro` 五个新方法 + `share_lot_repository` | T-1 | 低 | -| | T-4 | `convert_repository`(占位/回写/查询/清理) | T-1 | 低 | -| | T-5 | `locks.try_lock` + 单测 | 无 | 低 | -| **第 3 批 · 事务与编排** | T-6 | `convert_core_repository.apply_convert`(阶段一单事务) | T-1/T-3 | **高(方言 + 并发)** | -| | T-7 | `convert_service` 编排(八步 + 执行权 + 幂等 + 三阶段 + 阶段 1.5) | T-2~T-6 | **高(关键路径)** | -| **第 4 批 · 引擎与网关** | T-8 | `_amount_view` + `engine.process_convert_event` + `alert_service.events` | 无(可与 T-2 并行) | 中 | +| **第 1 批 · 数据与纯函数** | **T-1** ✅ | MySQL DDL + 3 个种子 + `reset.ps1` + sqlite DDL 同步 —— **2026-09-10 完成(断言 8/8)** | **T-0 + T-0b 双绿** | 低 | +| | **T-2** ✅ | `service/convert/` 纯函数包(`types`/`calc`/`fee`/`nav`/`lot_bootstrap`/`errors`)—— **93 用例全绿** | 无 | 低 | +| | **T-2b** ✅ | `calc_convert_demo.py` 实算 + 回填 PRD §5.3 与验收断言 —— **15/15 一致** | T-2 | 低 | +| **第 2 批 · 仓储与锁** | **T-3** ✅ | `core_ro` 五个新方法 + `share_lot_repository` | T-1 | 低 | +| | **T-4** ✅ | `convert_repository`(占位/回写/查询/清理) | T-1 | 低 | +| | **T-5** ✅ | `locks.try_lock` + 单测 | 无 | 低 | +| **第 3 批 · 事务与编排** | **T-6** ✅ | `convert_core_repository.apply_convert`(阶段一单事务)—— **真库 24/24** | T-1/T-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 分派 | 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-0 → T-1 → T-2 → T-6 → T-7 → T-13` -**并行组 A**:T-3 / T-4 / T-5(T-1 完成后同时开工) -**并行组 B**:T-8 全程可与 T-2 之后任意任务并行 +**关键路径**:`T-0 → T-1 → T-2 → T-6 → T-7 → T-13`(**T-7 已通,T-9 已解锁**) +**并行组 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`) -**测试基线**:**516**(批 0 后;原 510)→ 预计 **591 ~ 616**(CI 内约 **603**;含不进 CI 的压测用例约 **613**,见 §11) +**测试基线**:**672**(2026-09-10 T-8 后;批 0~3 路线 510 → 516 → 609 → 634 → 639 → 656 → **672**)→ 剩余任务(T-9~T-13)预计再加 **25~55** → **700~730**(估算) --- @@ -1014,15 +1014,63 @@ sqlite 无 gap lock,故该分支由 `tests/test_convert_core.py` 用注入点 - `process_convert_event`:取当日全量流水(已含两条)→ `run_rules` → `record_trade_alerts(primary=out_trade, hits, events=[event_of(out), event_of(in)])` → **一张单、`payload.events` 两条** - **`on_error_hook`(D19)**:一期传 `None`;**hook 调用必须包 `try/except`**,hook 自身失败**不得**反噬主流程(与「阶段 1.5 不阻断交易」同原则)——**必须有单测** -**DoD** -- [ ] `_amount_view` 对无 gid 交易恒等(`assert _amount_view(x) == x` 型用例) -- [ ] 一组 convert 两条流水 → RISK-002 只计一次;RISK-001/RISK-003 仍看到两条(验收 5/6) -- [ ] `process_convert_event` 只出**一条**预警单,`payload.events` 长度 2(验收 7) -- [ ] hook 抛异常时主流程正常返回(断言不抛) -- [ ] `pytest -q` 全绿(引擎既有 13 处 `run_rules` 调用零改动) +**DoD(全部达成,见下方执行记录)** +- [x] `_amount_view` 对无 gid 交易恒等(`assert _amount_view(x) == x` 型用例) +- [x] 一组 convert 两条流水 → RISK-002 只计一次;RISK-001/RISK-003 仍看到两条(验收 5/6) +- [x] `process_convert_event` 只出**一条**预警单,`payload.events` 长度 2(验收 7) +- [x] hook 抛异常时主流程正常返回(断言不抛) +- [x] `pytest -q` 全绿(引擎既有 13 处 `run_rules` 调用零改动) **依赖**:无(可与 T-2 之后任意阶段并行) +**执行记录(2026-09-10)** + +| 项 | 内容 | +| --- | --- | +| 改动文件 | `app/service/risk/rules.py`(+`_amount_view` + `run_rules` 双视图分流)· `app/service/risk/alert_service.py`(+公开 `build_trade_event` + `record_trade_alerts(events=...)`)· `app/service/risk/engine.py`(抽 `_run` + `process_convert_event` + `_notify_error_hook`) | +| 新增文件 | `tests/test_convert_engine.py`(**15 用例**)· `scripts/dev/verify_convert_engine.py`(**真 MySQL 验证 31 项**) | +| 测试改动 | `tests/test_convert_service.py` +1 条「接线回归」用例(大额转换真出单,防 `_run_engine` 退回静默跳过) | +| 结果 | pytest **672 passed / 3 skipped**(基线 656 **+16**,零回归);真库验证 **31/31 一致**(退出码 0) | + +**关键实现点(三处实施级收敛,均已在代码注释留痕)** + +1. **`events` 参数的语义收敛为「已构造的事件体列表」,但构造逻辑只留一份**。 + 架构 §6.2 写 `events=[event_of(out), event_of(in)]`,若 engine 自建 `event_of`,就会与 + `record_trade_alerts` 内的 `event = {...}` 形成**两份副本**(违反自检第 13 问)。 + 收敛为:`alert_service.build_trade_event(trade, hits)` **公开导出**,engine 调它构造两条, + `record_trade_alerts` 内部单流水分支也调它 —— 结构定义**唯一副本**,且与架构措辞一致。 +2. **聚合追加分支只追加首条(主事件)**。当日已有同客户 pending 单时,`append_alert_event` + 只追加转出端事件 —— 与架构 §5.4「一次转换两条流水,**只认转出端**」同口径 + (与 T-12 补偿脚本的幂等锚点同一理由)。**新建单**才落两条(验收 7)。 +3. **`related_trades` 缺省 `None` 使 `process_trade_event` 行为逐字节等价**:`_run` 抽出后 + `process_trade_event` 变薄封装(签名零改动),既有 13 处 `run_rules` 调用与 510 断言零影响。 + +**真库验证的额外产出(sqlite 绿证明不了的部分)** + +| 组 | 验证内容 | 结果 | +| --- | --- | --- | +| A | **阶段 1.5 从「跳过」变「生效」**:真库实跑出一张单、`payload.events` 两条(顺序 [redeem, subscribe])、`engine_error=False` | ✅ | +| B | **RISK-002 不翻倍**:阈值夹逼(单条 512000 < 800000 < 两条之和 1019975.33)→ 不含 RISK-002 | ✅ | +| C | 去重不删行:`core_trade` 2 条同 `convert_group_id`;两端 DECIMAL(18,2) 精度零漂移 | ✅ | +| D | 幂等重试(同 `client_request_id`)→ 同 group_id、流水仍 2 条、**单仍 1 张** | ✅ | +| E | 无命中场景(换日隔离)→ 不建单、仅落 1 条 pass 审计 | ✅ | +| F | 清理后 6 张表零残留(core_trade / core_customer / risk_alert / customer_profile_l3 / audit_log / risk_convert_detail) | ✅ | + +**真库验证两处脚本侧坑(已修正,与 T-7 同款)** + +1. **`convert_fund` 默认 `id_factory` 生成 `CNV-<日期>-`**,与清理口径 `LIKE 'CNV-T8M%'` + 不匹配 → 重跑时遗留行撞 `uk_group`。已注入 `id_factory=_t8m_id`(`CNV-T8M-`), + 并把清理条件补上 `client_request_id LIKE 'T8M-%'` 兜底。 +2. **`core_trade.trade_type` 在 MySQL 是 ENUM,`ORDER BY` 按定义序而非字母序** + (实测 `[subscribe, redeem]`)—— 断言改用 `{trade_type: row}` 字典定位,不依赖排序。 + +**发现但不在本任务范围(记下来,未顺手改)**:口径先钉死 —— **RISK-001 单笔 / RISK-002 当日累计, +是两条独立规则**(`rules.py:128-133` 逐笔比对 vs `:137-145` 求和比对),T-8 只把后者切到金额视图。 +真库 E 组现象(后续 1000 元小额交易也带出 RISK-001)的根因是**引擎入参范围**: +`process_trade_event` 拉**当日全量**流水重跑规则,故**触发这笔**与**命中那笔**可能不是同一笔。 +既有聚合逻辑(同客户同日一张 pending 单)兜住了重复出单,属**既有设计**,T-8 不改变、 +也不修(见 §B.9 纪律「不做清单外改动」)。 + --- ### 7.2 T-9 · API 模型与网关分派 diff --git a/scripts/dev/verify_convert_engine.py b/scripts/dev/verify_convert_engine.py new file mode 100644 index 0000000..2fbb27f --- /dev/null +++ b/scripts/dev/verify_convert_engine.py @@ -0,0 +1,412 @@ +"""T-8 真 MySQL 验证脚本:`process_convert_event` 引擎改造在真库/真账号下实跑 + DoD 断言。 + +**为什么 sqlite 单测全绿还不够(T-8 视角)** + +1. **「阶段 1.5 从跳过变生效」是主流程行为变化**:T-7 落地时 `process_convert_event` 尚不存在, + `convert_service._run_engine` 走 `ImportError` 分支**静默跳过**;T-8 落地后同一笔转换会 + **真实出单 + 写 L3 + 落审计**。真库跑一遍才能确认没有连带破坏(审计条数、L3 写入、单落库)。 +2. **`convert_group_id` 的 NULL 值域**:MySQL 里普通交易该列为 NULL、convert 两条为同一串值; + `_amount_view` 的 `if not gid` 必须在真库值域下成立(NULL / 空串都不得误聚合)。 +3. **DECIMAL 精度**:`core_trade.amount` 在 MySQL 是 DECIMAL,进 Python 后与生产纯函数 + 逐项比对(sqlite 用 REAL 无此保证,需 `float()` 绑定)。 +4. **JSON 列反解**:`risk_alert.payload` 真库落库后反解出的 `events` 长度必须为 2。 + +**断言分组** +A 引擎真跑出单:1 张单 / `payload.events` 两条 / 顺序 [转出, 转入] +B **RISK-002 不翻倍**(阈值夹逼:单条 < 阈值 < 两条之和) +C 去重不删行:`core_trade` 仍 2 条且同 `convert_group_id` +D 幂等重试(同 `client_request_id`)**不产生第二张单** +E 无命中场景(小额)不建单、仅 pass 审计 +F 清理后残留为零(自检) + +用法: + python scripts/dev/verify_convert_engine.py + +约定(与 T-6/T-7 脚本一致): +- 用 **T8M 前缀**隔离数据(客户/产品/批次/group_id),跑完**两个库**(core + agent)全清; +- 建/清数据走 `role="admin"`(需 DELETE);业务本身走 `convert_fund` 默认账号; +- 固定 `trace_id = "T8M-VERIFY-TRACE"`,便于精准清理 `audit_log`。 +""" + +from __future__ import annotations + +import json +import sys +import uuid +from datetime import date, datetime, timedelta +from decimal import ROUND_HALF_UP, Decimal +from pathlib import Path + +from sqlalchemy import text + +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT)) + +from app.config.settings import settings # noqa: E402 +from app.repository.core_ro import CoreReadOnlyRepository # noqa: E402 +from app.service.convert.calc import ( # noqa: E402 + convert_amount, + diff_fee, + hold_days, + in_qty, + lot_amount, + lot_fee, + plan_lots, +) +from app.service.convert.convert_service import convert_fund # noqa: E402 +from app.service.convert.fee import pick_fee_rate # noqa: E402 +from app.service.convert.types import FeeRule, Lot # noqa: E402 +from app.service.risk.rules import RiskThresholds # noqa: E402 +from app.utils.db import dispose_engines, get_engine # noqa: E402 +from app.utils.trace import new_trace # noqa: E402 + +CUSTOMER = "CUST-T8M" +PROD_OUT = "PROD-T8MO" +PROD_IN = "PROD-T8MI" +COMPANY = "华夏模拟基金" +TA = "TA-CN-001" +TRADE_AT = datetime(2026, 9, 4, 10, 0, 0) +TRADE_DATE = date(2026, 9, 4) +NOW = TRADE_AT +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 = "T8M-IDEM-0001" +TRACE = "T8M-VERIFY-TRACE" + +# 夹逼阈值:单条转出 512000 < daily_total 800000 < 两条之和(约 1020000) +TH = RiskThresholds( + large_amount=Decimal("100000"), + daily_total=Decimal("800000"), + 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 + + +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 _t8m_id(prefix: str, now: datetime) -> str: + """固定前缀的 id 工厂:`convert_fund` 默认生成 `CNV-<日期>-`,与清理口径 + `LIKE 'CNV-T8M%'` 不匹配 → 会在重跑时留下撞唯一键的残留行(T-7 脚本踩过同款坑)。""" + return f"{prefix}-T8M-{uuid.uuid4().hex[:8].upper()}" + + +def q1(engine, sql: str, **params): + with engine.connect() as conn: + return conn.execute(text(sql), params).scalar() + + +def rows(engine, sql: str, **params): + with engine.connect() as conn: + return [dict(r) for r in conn.execute(text(sql), params).mappings()] + + +# ── 数据准备 / 清理 ───────────────────────────────────────────────── +def cleanup(core_engine, agent_engine) -> None: + """按 T8M 前缀清理两个库(含幂等重跑前置清理)。""" + with core_engine.begin() as conn: + for sql, params in [ + ("DELETE FROM core_convert_lot_detail WHERE convert_group_id LIKE 'CNV-T8M%'", {}), + ("DELETE FROM core_trade WHERE customer_id = :c", {"c": CUSTOMER}), + ("DELETE FROM core_share_lot WHERE customer_id = :c", {"c": CUSTOMER}), + ("DELETE FROM core_holding WHERE customer_id = :c", {"c": CUSTOMER}), + ("DELETE FROM core_customer_risk WHERE customer_id = :c", {"c": CUSTOMER}), + ("DELETE FROM core_fee_rule WHERE product_id LIKE 'PROD-T8M%'", {}), + ("DELETE FROM core_product_nav WHERE product_id LIKE 'PROD-T8M%'", {}), + ("DELETE FROM core_product WHERE product_id LIKE 'PROD-T8M%'", {}), + ("DELETE FROM core_customer WHERE customer_id = :c", {"c": CUSTOMER}), + ]: + 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-T8M%' OR client_request_id LIKE 'T8M-%'" + ) + ) + conn.execute(text("DELETE FROM risk_alert WHERE customer_id = :c"), {"c": CUSTOMER}) + conn.execute( + text("DELETE FROM customer_profile_l3 WHERE customer_id = :c"), {"c": CUSTOMER} + ) + conn.execute( + text("DELETE FROM audit_log WHERE trace_id LIKE 'T8M-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, 'T8真库验证', :d)" + ), + {"c": CUSTOMER, "d": TRADE_DATE}, + ) + conn.execute( + text( + "INSERT INTO core_customer_risk (customer_id, risk_code, evaluated_at, expires_at) " + "VALUES (:c, 'C3', :t, :exp)" + ), + {"c": CUSTOMER, "t": TRADE_AT - timedelta(days=30), + "exp": TRADE_AT + timedelta(days=300)}, + ) + for pid, name, ptype, rate in [ + (PROD_OUT, "T8转出基金", "bond", OUT_RATE), + (PROD_IN, "T8转入基金", "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}, + ) + # 两批:400000 份(持有 ~400 天,费率 0)+ 100000 份(持有 3 天,费率 1.5%) + for lot_id, qty, nav, confirmed in [ + ("LOT-T8M-A1", "400000", "1.0300", datetime(2025, 8, 1, 10, 0, 0)), + ("LOT-T8M-A2", "100000", "1.0000", datetime(2026, 9, 1, 10, 0, 0)), + ]: + 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 expected_quote(core_ro, requested: str) -> dict: + """期望折算(全部走生产纯函数,与 service 内部同口径)。""" + lots = [Lot.from_row(r) for r in core_ro.list_share_lots(CUSTOMER, PROD_OUT)] + fee_rules = [FeeRule.from_row(r) for r in core_ro.get_redeem_fee_rules(PROD_OUT)] + plan = plan_lots(lots, Decimal(requested)) + out_amount = Decimal("0") + redeem_fee = Decimal("0") + for alloc in plan.allocations: + amount = lot_amount(alloc.qty, alloc.nav) + rate = pick_fee_rate( + fee_rules, hold_days(TRADE_DATE, alloc.confirmed_at), product_id=PROD_OUT + ) + out_amount += amount + redeem_fee += lot_fee(amount, rate) + conv = convert_amount(out_amount, redeem_fee) + gap = diff_fee(conv, OUT_RATE, IN_RATE, settings.convert_diff_fee_mode) + return { + "out_amount": out_amount, + "redeem_fee": redeem_fee, + "convert_amount": conv, + "diff_fee": gap, + "in_amount": conv - gap, + } + + +def main() -> int: + core_engine = get_engine("jinrong_core", role="admin") + agent_engine = get_engine("jinrong_agent", role="admin") + print("=" * 60) + print("T-8 真库验证:规则引擎改造(_amount_view 去重 + process_convert_event 出单)") + print("=" * 60) + + cleanup(core_engine, agent_engine) + seed(core_engine) + core_ro = CoreReadOnlyRepository() + req_qty = "500000" # 全部转出 + exp = expected_quote(core_ro, req_qty) + + # ── A/B/C/D:带幂等键的大额转换 ───────────────────────────────── + print("\n【A】引擎真跑(阶段 1.5 不再跳过):一张单 + payload.events 两条") + new_trace(TRACE) + resp = convert_fund( + { + "customer_id": CUSTOMER, + "from_product_id": PROD_OUT, + "to_product_id": PROD_IN, + "qty": req_qty, + "client_request_id": CID_REQ, + }, + now=NOW, + thresholds=TH, + id_factory=_t8m_id, + ) + check("转换未被阻断", resp["blocked"], False) + check("引擎异常标记为 False(说明真的跑了且没炸)", resp["engine_error"], False) + check("折算 out_amount 与生产纯函数一致", Decimal(resp["out_amount"]), exp["out_amount"]) + check("折算 in_amount 与生产纯函数一致", Decimal(resp["in_amount"]), exp["in_amount"]) + check("RISK-001 命中(转出端 512000 ≥ 100000)", "RISK-001" in resp["triggered_rules"], True) + check("只有一张预警单", len(resp["alert_ids"]), 1) + + alerts = rows( + agent_engine, "SELECT * FROM risk_alert WHERE customer_id = :c", c=CUSTOMER + ) + check("risk_alert 落库 1 行", len(alerts), 1) + payload = alerts[0]["payload"] + payload = json.loads(payload) if isinstance(payload, str) else payload + check("payload.events 长度 2", len(payload["events"]), 2) + check( + "events 顺序 [转出, 转入]", + [e["trade_type"] for e in payload["events"]], + ["redeem", "subscribe"], + ) + check("events[0].trade_id 为转出端", payload["events"][0]["trade_id"], resp["out_trade_id"]) + check("events[1].trade_id 为转入端", payload["events"][1]["trade_id"], resp["in_trade_id"]) + + print("\n【B】RISK-002 不翻倍(单条 512000 < 800000 < 两条之和)") + check("triggered_rules 不含 RISK-002", "RISK-002" in resp["triggered_rules"], False) + # 同日两条流水之和确实超过阈值 → 证明「不命中」来自去重而非阈值太松 + total_two = exp["out_amount"] + exp["in_amount"] + check("两条之和确实 ≥ 阈值(阈值夹逼成立)", total_two >= TH.daily_total, True) + print(f" └ 转出 {exp['out_amount']} + 转入 {exp['in_amount']} = {total_two}") + + print("\n【C】去重不删行:core_trade 仍 2 条且同组") + gid = resp["convert_group_id"] + trades = rows( + core_engine, + "SELECT * FROM core_trade WHERE convert_group_id = :g", + g=gid, + ) + by_type = {t["trade_type"]: t for t in trades} + check("core_trade 2 条", len(trades), 2) + check("类型覆盖 [redeem, subscribe]", sorted(by_type), ["redeem", "subscribe"]) + check("两条同 convert_group_id", {t["convert_group_id"] for t in trades}, {gid}) + check( + "转出端金额精度零漂移(DECIMAL 18,2)", + Decimal(str(by_type["redeem"]["amount"])), + exp["out_amount"], + ) + check( + "转入端金额精度零漂移(DECIMAL 18,2)", + Decimal(str(by_type["subscribe"]["amount"])), + exp["in_amount"], + ) + + print("\n【D】幂等重试(同 client_request_id)不产生第二张单") + resp2 = convert_fund( + { + "customer_id": CUSTOMER, + "from_product_id": PROD_OUT, + "to_product_id": PROD_IN, + "qty": req_qty, + "client_request_id": CID_REQ, + }, + now=NOW + timedelta(minutes=1), + thresholds=TH, + id_factory=_t8m_id, + ) + check("重试返回同一 group_id", resp2["convert_group_id"], gid) + check("重试 core_trade 仍 2 条", q1( + core_engine, "SELECT COUNT(*) FROM core_trade WHERE customer_id = :c", c=CUSTOMER), 2) + check("重试后预警单仍 1 张", q1( + agent_engine, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", c=CUSTOMER), 1) + + # ── E:无命中场景 ──────────────────────────────────────────────── + print("\n【E】无命中场景(换日 → 当日无其他流水)不建单、仅 pass 审计") + before_alerts = q1( + agent_engine, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", c=CUSTOMER) + before_pass = q1( + agent_engine, + "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'pass'", t=TRACE) + # 该客户份额已全转出 → 4xx 不落库;故先补一批小额份额再造一笔 + with core_engine.begin() as conn: + conn.execute( + text( + "INSERT INTO core_share_lot (lot_id, customer_id, product_id, qty, remain_qty, " + "nav, confirmed_at) VALUES ('LOT-T8M-B1', :c, :p, 1000, 1000, 1.0000, :cat)" + ), + {"c": CUSTOMER, "p": PROD_OUT, "cat": datetime(2025, 8, 1, 10, 0, 0)}, + ) + # 换到次日:`list_trades_range` 按 [日初, 次日) 取数,前一日的大额流水不再进本批规则输入 + resp3 = convert_fund( + { + "customer_id": CUSTOMER, + "from_product_id": PROD_OUT, + "to_product_id": PROD_IN, + "qty": "1000", + "client_request_id": "T8M-IDEM-0002", + }, + now=NOW + timedelta(days=1), + thresholds=TH, + id_factory=_t8m_id, + ) + check("小额转换成功(非 blocked)", resp3["blocked"], False) + check("未命中任何规则", resp3["triggered_rules"], []) + check("未新建预警单", q1( + agent_engine, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", c=CUSTOMER), + before_alerts) + check("落 pass 审计 1 条", q1( + agent_engine, + "SELECT COUNT(*) FROM audit_log WHERE trace_id = :t AND decision = 'pass'", t=TRACE) + - before_pass, 1) + + # ── F:清理自检 ───────────────────────────────────────────────── + print("\n【F】清理 T8M 前缀数据并自检无残留") + cleanup(core_engine, agent_engine) + check("core_trade 零残留", q1( + core_engine, "SELECT COUNT(*) FROM core_trade WHERE customer_id = :c", c=CUSTOMER), 0) + check("core 客户零残留", q1( + core_engine, "SELECT COUNT(*) FROM core_customer WHERE customer_id = :c", c=CUSTOMER), 0) + check("risk_alert 零残留", q1( + agent_engine, "SELECT COUNT(*) FROM risk_alert WHERE customer_id = :c", c=CUSTOMER), 0) + check("customer_profile_l3 零残留", q1( + agent_engine, + "SELECT COUNT(*) FROM customer_profile_l3 WHERE customer_id = :c", c=CUSTOMER), 0) + check("audit_log 零残留", q1( + agent_engine, "SELECT COUNT(*) FROM audit_log WHERE trace_id LIKE 'T8M-VERIFY%'"), 0) + check("risk_convert_detail 零残留", q1( + agent_engine, + "SELECT COUNT(*) FROM risk_convert_detail " + "WHERE convert_group_id LIKE 'CNV-T8M%' OR client_request_id LIKE 'T8M-%'"), 0) + + dispose_engines() + print("\n" + "=" * 60) + print(f"真库验证(T-8 规则引擎改造):{_passed} 项一致 / {_failed} 项不一致") + print("=" * 60) + return 1 if _failed else 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_convert_engine.py b/tests/test_convert_engine.py new file mode 100644 index 0000000..5b28664 --- /dev/null +++ b/tests/test_convert_engine.py @@ -0,0 +1,310 @@ +"""T-8 规则引擎改造单测(开发计划 §7.1 DoD / 验收 5·6·7)。 + +覆盖三层: +1. `rules._amount_view` 纯函数 —— 同组只留转出端 / 无 gid 恒等 / 组内无 redeem 防御 / 顺序保持 +2. `run_rules` 视图分流 —— RISK-002 与 RISK-005 走金额视图(**不翻倍**), + RISK-001 / RISK-003 仍看**全量**(证明去重未删行,验收 6) +3. `engine.process_convert_event` —— **一张单 + `payload.events` 两条**(验收 7)、 + 两条流水只跑一次、D19 `on_error_hook` 容错(hook 自身抛异常不得反噬主流程) + +sqlite 内存库(`conftest.sqlite_engine`),数据自建。阈值一律**显式传入**, +不读 settings(避免与 conftest 的 autouse 隔离相互干扰,也让去重差异可精确归因)。 +""" + +from __future__ import annotations + +from datetime import datetime, timedelta +from decimal import Decimal + +import pytest +from sqlalchemy import text + +from app.repository.core_ro import CoreReadOnlyRepository +from app.repository.risk_repository import RiskRepository +from app.service.risk.engine import process_convert_event, process_trade_event +from app.service.risk.rules import RiskThresholds, _amount_view, run_rules + +CUST = "CUST-T8" +PROD_A = "PROD-T8A" # 转出方 +PROD_B = "PROD-T8B" # 转入方 +GID = "CNV-T8-0001" +NOW = datetime(2026, 9, 4, 10, 0, 0) + + +def _th(**over) -> RiskThresholds: + """显式阈值(不让用例依赖 settings;concentration 推到不可达避免 RISK-006 干扰)。""" + base = dict( + large_amount=Decimal("200000"), + daily_total=Decimal("500000"), + 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, + ) + base.update(over) + return RiskThresholds(**base) + + +def _trade( + trade_id: str, + trade_type: str, + amount: str, + at: datetime = NOW, + *, + gid: str | None = None, + product: str = PROD_A, +) -> dict: + """构造一条规则引擎口径的流水(等价 core_trade 行)。""" + return { + "trade_id": trade_id, + "customer_id": CUST, + "product_id": product, + "trade_type": trade_type, + "amount": Decimal(amount), + "qty": Decimal("1"), + "trade_status": "confirmed", + "traded_at": at, + "convert_group_id": gid, + } + + +def _convert_legs(out_amount: str = "300000", in_amount: str = "299000") -> tuple[dict, dict]: + """一次转换的两条流水(同组,转出在前)。""" + out = _trade("TRD-T8-OUT", "redeem", out_amount, gid=GID, product=PROD_A) + inn = _trade( + "TRD-T8-IN", "subscribe", in_amount, NOW + timedelta(seconds=1), gid=GID, product=PROD_B + ) + return out, inn + + +def _ids(hits) -> set[str]: + return {h.rule_id for h in hits} + + +# ── 1. `_amount_view` 纯函数 ───────────────────────────────────────── +def test_amount_view_identity_without_group(): + """无 convert_group_id 的交易**原样通过**(非 convert 场景恒等 → 既有断言零影响)。""" + raw = [ + _trade("T-1", "redeem", "1000"), + _trade("T-2", "subscribe", "900", NOW + timedelta(seconds=1), product=PROD_B), + ] + view = _amount_view(raw) + assert view == raw, "无 gid 输入必须恒等(内容与顺序都不变)" + assert view is not raw, "返回独立列表,不得共享可变状态" + + +def test_amount_view_keeps_redeem_when_subscribe_comes_first(): + out, inn = _convert_legs() + view = _amount_view([inn, out]) + assert [t["trade_id"] for t in view] == ["TRD-T8-OUT"], "同组只留转出端,且位置不变" + + +def test_amount_view_keeps_redeem_when_redeem_comes_first(): + out, inn = _convert_legs() + view = _amount_view([out, inn]) + assert [t["trade_id"] for t in view] == ["TRD-T8-OUT"], "转入端不得覆盖已占位的转出端" + + +def test_amount_view_group_without_redeem_keeps_first_row(): + """防御分支:组内无 redeem 时保留首条(R-b 规定 convert 必有转出端,理论不可达)。""" + a = _trade("T-IN-1", "subscribe", "100", gid="CNV-T8-X", product=PROD_B) + b = _trade("T-IN-2", "subscribe", "200", NOW + timedelta(seconds=1), gid="CNV-T8-X", + product=PROD_B) + assert _amount_view([a, b]) == [a] + + +def test_amount_view_mixes_groups_and_plain_trades_in_order(): + """多组 + 无 gid 混合:顺序保持,各组各留一条。""" + n1 = _trade("T-N1", "redeem", "1000") + n2 = _trade("T-N2", "subscribe", "2000", NOW + timedelta(seconds=3), product=PROD_B) + out1, in1 = _convert_legs() + out2 = _trade("TRD-T8-OUT2", "redeem", "500", NOW + timedelta(seconds=4), gid="CNV-T8-0002") + in2 = _trade("TRD-T8-IN2", "subscribe", "480", NOW + timedelta(seconds=5), + gid="CNV-T8-0002", product=PROD_B) + view = _amount_view([n1, in1, out1, n2, in2, out2]) + assert [t["trade_id"] for t in view] == ["T-N1", "TRD-T8-OUT", "T-N2", "TRD-T8-OUT2"] + + +# ── 2. `run_rules` 视图分流(验收 5/6)─────────────────────────────── +def test_daily_total_counts_convert_only_once(): + """RISK-002 走金额视图:一次转换只计转出一端,不翻倍(验收 5)。""" + out, inn = _convert_legs() + # 阈值 500000:全量口径 300000+299000=599000 会命中;去重后 300000 不命中 + assert "RISK-002" not in _ids(run_rules([out, inn], _th(), now=NOW)) + + # 反证:同一对流水若**没有 gid**(视为两笔独立交易)→ 命中,证明差异来自去重而非阈值 + loose = [_trade("TRD-T8-OUT", "redeem", "300000"), _trade("TRD-T8-IN", "subscribe", "299000")] + assert "RISK-002" in _ids(run_rules(loose, _th(), now=NOW)) + + +def test_daily_total_detail_reflects_deduped_amount(): + """合计金额与笔数都只算转出端(阈值下调到 290000 让 RISK-002 命中以便读 detail)。""" + out, inn = _convert_legs() + hits = run_rules([out, inn], _th(daily_total=Decimal("290000")), now=NOW) + hit = next(h for h in hits if h.rule_id == "RISK-002") + assert "300000" in hit.detail + assert "599000" not in hit.detail, "detail 不得出现两条流水之和" + assert "共 1 笔" in hit.detail + + +def test_freq_trade_still_sees_both_legs(): + """RISK-003 看全量(去重未删行,验收 6):转入端也参与分产品计数。""" + out, inn = _convert_legs("1000", "900") + extra = _trade("TRD-T8-B2", "subscribe", "800", NOW + timedelta(seconds=2), product=PROD_B) + hits = run_rules([out, inn, extra], _th(freq_count=2), now=NOW) + assert "RISK-003" in _ids(hits), "产品 B 应有 2 笔(convert 转入端 + 普通申购)" + assert any(PROD_B in h.detail and "2 笔" in h.detail for h in hits if h.rule_id == "RISK-003") + + +def test_large_amount_still_sees_both_legs(): + """RISK-001 看全量:转入端大额也能命中(逐笔判定,不走去重视图)。""" + out, inn = _convert_legs("100000", "250000") # 仅转入端 ≥ 200000 + assert "RISK-001" in _ids(run_rules([out, inn], _th(), now=NOW)) + + +def test_small_then_large_ignores_converted_leg_as_buildup(): + """RISK-005 走金额视图:转入端不得被当成「铺垫」小额(否则虚增 small_count 误报)。""" + small1 = _trade("T-S1", "redeem", "5000", product=PROD_A) + small2 = _trade("T-S2", "redeem", "6000", NOW + timedelta(seconds=1), product=PROD_A) + # 转入端本身是小额(5000)且时间早于转出端 → 若全量参与,会被算作第 3 笔铺垫 + inn = _trade("TRD-T8-IN", "subscribe", "5000", NOW + timedelta(seconds=2), + gid=GID, product=PROD_B) + out = _trade("TRD-T8-OUT", "redeem", "300000", NOW + timedelta(seconds=3), + gid=GID, product=PROD_A) + assert "RISK-005" not in _ids(run_rules([small1, small2, inn, out], _th(), now=NOW)) + + # 反证:去掉 gid(转入端成为独立小额)→ 铺垫凑够 3 笔 → 命中 + inn_loose = _trade("TRD-T8-IN", "subscribe", "5000", NOW + timedelta(seconds=2), + product=PROD_B) + assert "RISK-005" in _ids(run_rules([small1, small2, inn_loose, out], _th(), now=NOW)) + + +# ── 3. `process_convert_event` 集成(验收 7)───────────────────────── +def _exec(engine, sql: str, **params) -> None: + with engine.begin() as conn: + conn.execute(text(sql), params) + + +def _rows(engine, sql: str, **params) -> list[dict]: + with engine.connect() as conn: + return [dict(r) for r in conn.execute(text(sql), params).mappings()] + + +def _seed_trades(engine, *trades: dict) -> None: + for i, t in enumerate(trades): + _exec( + engine, + "INSERT INTO core_trade (trade_id, customer_id, product_id, trade_type, amount, " + "qty, convert_group_id, trade_status, traded_at) " + "VALUES (:tid, :c, :p, :tt, :amt, 1, :gid, :st, :at)", + tid=t["trade_id"], c=t["customer_id"], p=t["product_id"], tt=t["trade_type"], + amt=float(t["amount"]), gid=t["convert_group_id"], st=t["trade_status"], + at=t["traded_at"], + ) + + +def _services(engine) -> dict: + return dict( + core_ro=CoreReadOnlyRepository(engine=engine), + risk_repo=RiskRepository(engine=engine), + ) + + +def test_process_convert_event_creates_one_alert_with_two_events(sqlite_engine): + """一次转换 → 一张单、`payload.events` 两条、RISK-002 不翻倍(验收 5/7)。""" + out, inn = _convert_legs() + _seed_trades(sqlite_engine, out, inn) + th = _th(daily_total=Decimal("400000")) # 全量 599000 会命中;去重后 300000 不命中 + + result = process_convert_event(out, inn, thresholds=th, **_services(sqlite_engine)) + + assert result["triggered_rules"] == ["RISK-001"], "只应命中单笔大额,累计不翻倍" + assert result["aml_hit"] is False + assert len(result["alert_ids"]) == 1, "一次转换只出一张单" + + alerts = _rows(sqlite_engine, "SELECT * FROM risk_alert") + assert len(alerts) == 1 + import json + + payload = json.loads(alerts[0]["payload"]) + assert len(payload["events"]) == 2, "一张单承载两条事件" + assert [e["trade_id"] for e in payload["events"]] == ["TRD-T8-OUT", "TRD-T8-IN"] + assert payload["events"][0]["trade_type"] == "redeem" + assert payload["events"][1]["trade_type"] == "subscribe" + # 主流水口径:单主体与 payload.product_id 取转出端 + assert alerts[0]["trade_id"] == "TRD-T8-OUT" + assert payload["product_id"] == PROD_A + + # 去重不删行:两条流水仍在 + assert len(_rows(sqlite_engine, "SELECT * FROM core_trade")) == 2 + + +def test_process_convert_event_no_hit_writes_pass_audit_only(sqlite_engine): + """无命中 → 不建单,仅落 pass 审计(events 参数不参与该分支)。""" + out, inn = _convert_legs("1000", "900") + _seed_trades(sqlite_engine, out, inn) + + result = process_convert_event(out, inn, thresholds=_th(), **_services(sqlite_engine)) + + assert result["alert_ids"] == [] + assert _rows(sqlite_engine, "SELECT * FROM risk_alert") == [] + rows = _rows(sqlite_engine, "SELECT * FROM audit_log WHERE decision = 'pass'") + assert len(rows) == 1 + + +def test_process_trade_event_still_writes_single_event(sqlite_engine): + """回归保护:`process_trade_event` 签名与行为不变(payload.events 仍为 1 条)。""" + t = _trade("TRD-T8-PLAIN", "redeem", "300000") + _seed_trades(sqlite_engine, t) + + result = process_trade_event(t, thresholds=_th(), **_services(sqlite_engine)) + + assert result["triggered_rules"] == ["RISK-001"] + import json + + alerts = _rows(sqlite_engine, "SELECT * FROM risk_alert") + payload = json.loads(alerts[0]["payload"]) + assert len(payload["events"]) == 1 + assert payload["events"][0]["trade_id"] == "TRD-T8-PLAIN" + + +# ── 4. D19 `on_error_hook` 容错 ────────────────────────────────────── +def _boom(*_args, **_kwargs): + raise ZeroDivisionError("引擎内部炸了") + + +def test_error_hook_receives_exception_and_swallows_its_own(sqlite_engine, monkeypatch): + """hook 被调用一次;**hook 自身抛异常必须被吞掉**,原始异常照常上抛。""" + out, inn = _convert_legs() + core = CoreReadOnlyRepository(engine=sqlite_engine) + monkeypatch.setattr(core, "list_trades_range", _boom) + calls: list[tuple[dict, Exception]] = [] + + def hook(out_trade: dict, exc: Exception) -> None: + calls.append((out_trade, exc)) + raise RuntimeError("hook 自己也炸了") + + with pytest.raises(ZeroDivisionError, match="引擎内部炸了"): + process_convert_event( + out, inn, core_ro=core, risk_repo=RiskRepository(engine=sqlite_engine), + thresholds=_th(), on_error_hook=hook, + ) + assert len(calls) == 1, "hook 必须被调用一次" + assert calls[0][0]["trade_id"] == "TRD-T8-OUT", "hook 收到主流水(转出端)" + assert isinstance(calls[0][1], ZeroDivisionError) + + +def test_error_hook_none_keeps_original_exception(sqlite_engine, monkeypatch): + """一期默认 `on_error_hook=None`:异常原样上抛(由阶段 1.5 落 engine_error)。""" + out, inn = _convert_legs() + core = CoreReadOnlyRepository(engine=sqlite_engine) + monkeypatch.setattr(core, "list_trades_range", _boom) + + with pytest.raises(ZeroDivisionError): + process_convert_event( + out, inn, core_ro=core, risk_repo=RiskRepository(engine=sqlite_engine), + thresholds=_th(), + ) diff --git a/tests/test_convert_service.py b/tests/test_convert_service.py index 0dddde5..e3f8f60 100644 --- a/tests/test_convert_service.py +++ b/tests/test_convert_service.py @@ -18,6 +18,7 @@ sqlite 内存库(`conftest.sqlite_engine`,含 C×R 矩阵种子),数据 from __future__ import annotations +import json from datetime import date, datetime, timedelta from decimal import Decimal @@ -49,6 +50,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.rules import RiskThresholds CUST = "CUST-T7" # C3 客户 → R4 产品 allowed_with_disclosure(放行) CUST_LOW = "CUST-T7L" # C1 客户 → R4 forbidden(用于 blocked 分支) @@ -434,3 +436,35 @@ def test_too_many_lots_rejected(sqlite_engine, monkeypatch): with pytest.raises(TooManyLots) as exc: convert_fund(_req(qty="120"), now=NOW, **_services(sqlite_engine)) assert exc.value.extra == {"batch_count": 2, "max_lots": 1} + + +# ── 10. 阶段 1.5 接线(T-8):引擎真跑并出单 ───────────────────────── +def test_engine_wired_produces_single_alert_with_two_events(sqlite_engine): + """T-8 落地后阶段 1.5 不再跳过:大额转换**真出一张单**、`payload.events` 两条(验收 7)。 + + 本用例是「接线回归」:若 `_run_engine` 又被改回静默跳过(或签名对不上被 + ImportError 吞掉),这里会因 `triggered_rules` 为空而变红。 + """ + _seed(sqlite_engine) + th = RiskThresholds( + large_amount=Decimal("100"), # 调低以让 120 份的折算额命中 RISK-001 + daily_total=Decimal("1000000"), + freq_count=3, + probe_window_minutes=5, + probe_count=3, + probe_amount=Decimal("400000"), + small_amount=Decimal("10000"), + small_count=3, + concentration_threshold=1.01, + ) + resp = convert_fund(_req(), now=NOW, thresholds=th, **_services(sqlite_engine)) + + assert resp["engine_error"] is False, "引擎真的跑了且没炸(未走 ImportError 跳过分支)" + assert "RISK-001" in resp["triggered_rules"] + assert len(resp["alert_ids"]) == 1, "一次转换只出一张单" + + alerts = _rows(sqlite_engine, "SELECT * FROM risk_alert") + assert len(alerts) == 1 + payload = json.loads(alerts[0]["payload"]) + assert len(payload["events"]) == 2, "一张单承载两条事件(转出 + 转入)" + assert [e["trade_type"] for e in payload["events"]] == ["redeem", "subscribe"]