diff --git a/app/api/simulate.py b/app/api/simulate.py index a42d078..85ecde4 100644 --- a/app/api/simulate.py +++ b/app/api/simulate.py @@ -3,6 +3,9 @@ 鉴权:依赖 T-01 JWT(risk_demo 或客户本人 customer_id == subject); `get_auth_context` B6 落地后在此挂 `Depends`,本路由签名不变(开发计划 B5/B6 边界)。 trace:B7 中间件贯通;B7 前由 service 层 ensure_trace 兜底。 +挂载:B7 集成 main.py(当前仅 TestClient 独立挂 router 验证)。 +统一响应外壳:utils/response.py 为占位(P1 任务),落地点挂账 B7(届时 +simulate/risk 一并包裹,本路由返回体不变)。 """ from __future__ import annotations @@ -20,7 +23,7 @@ router = APIRouter(prefix="/api/simulate", tags=["simulate"]) class TradeRequest(BaseModel): customer_id: str = Field(..., min_length=1) product_id: str = Field(..., min_length=1) - trade_type: str = Field(..., description="subscribe | redeem;convert 显式拒绝") + trade_type: str = Field(..., max_length=16, description="subscribe | redeem;convert 显式拒绝") amount: Decimal = Field(..., gt=0, description="交易金额(元),必须为正数") @@ -31,5 +34,5 @@ def submit_trade_api(req: TradeRequest) -> dict: return submit_trade(req.model_dump()) except UnsupportedTradeType as exc: raise HTTPException(status_code=400, detail=str(exc)) from exc - except LookupError as exc: + except LookupError as exc: # KeyError 亦为其子类;B6 统一 NotFoundError 时收敛为精确类型 raise HTTPException(status_code=404, detail=str(exc)) from exc diff --git a/app/gateway/trade_gateway.py b/app/gateway/trade_gateway.py index 002007a..2500af4 100644 --- a/app/gateway/trade_gateway.py +++ b/app/gateway/trade_gateway.py @@ -3,7 +3,12 @@ submit_trade 为唯一入口:参数校验 → suitability_check(FR-2,落校验日志; 不匹配→ suitability 预警单 + 阻断响应,交易不落 core_trade)→ 匹配 → INSERT core_trade → 同步调规则引擎(FR-3)→ 返回 blocked + trade_id + -触发规则。阻断/放行全量审计(agent_type='platform',FR-1 §6)。 +触发规则。阻断/放行全量审计(agent_type='platform',FR-1 §6); +convert/未知类型属参数校验失败(400),不落审计(PRD 审计口径仅阻断/放行)。 + +引擎异常兜底(架构 §5.3):交易已成立(core_trade 已提交),审计 +decision='risk_engine_error' + logger.exception,响应带 engine_error=true +供 B9a rebuild_alerts 按 trade_id 补偿重放。 鉴权归路由层(T-01/B6 的 get_auth_context:risk_demo 或客户本人); trace 由调用方中间件贯通,本层 ensure_trace 兜底(脚本/测试直调场景)。 @@ -11,9 +16,10 @@ trace 由调用方中间件贯通,本层 ensure_trace 兜底(脚本/测试 from __future__ import annotations +import logging from datetime import datetime from decimal import Decimal -from typing import Any +from typing import Any, Callable from uuid import uuid4 from app.gateway.gateway_repository import GatewayRepository @@ -25,6 +31,8 @@ from app.service.risk.alert_service import record_suitability_alert from app.service.suitability import SuitabilityResult, suitability_check from app.utils.trace import ensure_trace +logger = logging.getLogger(__name__) + SUPPORTED_TRADE_TYPES = ("subscribe", "redeem") CONVERT_MESSAGE = "转换交易暂不支持,请分别发起申购/赎回" ADVICE = "请联系持证投资顾问" @@ -81,12 +89,16 @@ def submit_trade( gateway_repo: GatewayRepository | None = None, thresholds: RiskThresholds | None = None, now: datetime | None = None, + trade_id_factory: Callable[[datetime], str] | None = None, ) -> dict[str, Any]: """处理一笔模拟交易请求(PRD FR-1 流程 ①~⑤)。 req:{customer_id, product_id, trade_type, amount};amount 转 Decimal。 + trade_id_factory:测试注入点(架构 §7 约定集成测试交易用 TRD-TEST- 前缀, + B8 teardown 按前缀清理;缺省 TRD-{date}-{uuid8})。 返回 FR-1 ⑤ 响应体:阻断 {blocked, trade_id, block_reason, reasons, advice, - notice};放行 {blocked, trade_id, triggered_rules, alert_ids, aml_hit}。 + notice};放行 {blocked, trade_id, triggered_rules, alert_ids, aml_hit} + (引擎异常时附 engine_error=true)。 """ core = core_ro or CoreReadOnlyRepository() repo = risk_repo or RiskRepository() @@ -101,7 +113,7 @@ def submit_trade( raise UnsupportedTradeType(f"不支持的交易类型: {trade_type}(仅 subscribe/redeem)") now = now or datetime.now() - trade_id = _new_trade_id(now) + trade_id = (trade_id_factory or _new_trade_id)(now) amount = Decimal(str(req["amount"])) traded_at = now @@ -128,6 +140,7 @@ def submit_trade( trade_id=trade_id, req=req, rule_id=result.rule_id, + detail={"block_reason": result.block_reason, "reasons": list(result.reasons)}, ) return { "blocked": True, @@ -141,19 +154,44 @@ def submit_trade( writer.insert_trade( trade_id, req["customer_id"], req["product_id"], trade_type, amount, traded_at ) - engine_result = process_trade_event( - { + try: + engine_result = process_trade_event( + { + "trade_id": trade_id, + "customer_id": req["customer_id"], + "product_id": req["product_id"], + "trade_type": trade_type, + "amount": amount, + "trade_status": "confirmed", + "traded_at": traded_at, + }, + core_ro=core, + risk_repo=repo, + thresholds=th, + ) + except Exception: + # 交易已成立(core_trade 已提交):留审计与日志供 B9a rebuild_alerts 补偿 + logger.exception("risk engine failed after trade accepted: %s", trade_id) + _audit( + repo, + decision="risk_engine_error", + trade_id=trade_id, + req=req, + detail={"error_stage": "process_trade_event"}, + ) + return { + "blocked": False, "trade_id": trade_id, - "customer_id": req["customer_id"], - "product_id": req["product_id"], - "trade_type": trade_type, - "amount": amount, - "trade_status": "confirmed", - "traded_at": traded_at, - }, - core_ro=core, - risk_repo=repo, - thresholds=th, + "triggered_rules": [], + "alert_ids": [], + "aml_hit": False, + "engine_error": True, + } + _audit( + repo, + decision="trade_accepted", + trade_id=trade_id, + req=req, + detail=dict(engine_result), # 全量输出(评审 P2-3) ) - _audit(repo, decision="trade_accepted", trade_id=trade_id, req=req) return {"blocked": False, "trade_id": trade_id, **engine_result} diff --git a/docs/memory/TODO.md b/docs/memory/TODO.md index 2bdc9be..287b131 100644 --- a/docs/memory/TODO.md +++ b/docs/memory/TODO.md @@ -20,7 +20,7 @@ ### 风控模块(PRD v1.0 已冻结 · `docs/PRD/PRD-风控监测Agent.md`,事件驱动线不依赖 T-07 可先行) -- [ ] T-30 风控事件线:`app/gateway/` 交易网关 + `service/risk/` 规则引擎(RISK-001~005)+ 预警单聚合 + L3 最小写入 + `risk:pub:alert` 推送 + AML(含 `risk_aml_list` 种子)+ 演示数据脚本验收 A-1~A-5/A-9 —— **进度(2026-09-06):B1 规则纯函数 / B2 预警服务 / B3 L3 写入(risk_score 一期不写)/ B4 AML+引擎编排 均已完成并经独立 AI 评审闭环;剩 B5 网关、B6 鉴权+API、B7 main 集成、B8 conftest+集成测试、B9a 脚本、B9b 演示走查** +- [ ] T-30 风控事件线:`app/gateway/` 交易网关 + `service/risk/` 规则引擎(RISK-001~005)+ 预警单聚合 + L3 最小写入 + `risk:pub:alert` 推送 + AML(含 `risk_aml_list` 种子)+ 演示数据脚本验收 A-1~A-5/A-9 —— **进度(2026-09-06):B1 规则纯函数 / B2 预警服务 / B3 L3 写入(risk_score 一期不写)/ B4 AML+引擎编排 / B5 交易网关(convert 400、阻断不落 trade、引擎异常审计降级)均已完成并经独立 AI 评审闭环;剩 B6 鉴权+API、B7 main 集成、B8 conftest+集成测试、B9a 脚本、B9b 演示走查** - [ ] T-31 `service/suitability.py` 公共校验(SUIT-001~008)+ `POST /api/risk/suitability/check` + 单测验收 A-8 —— **代码已完成(A1~A4 · 77 测试绿);MySQL 手工 SQL 对照挂账至 B9b 执行(阶段 A 评审 P2-8)** - [ ] T-32 预警台账与人工处置 API(`GET /alerts`、`POST /handle`,risk_officer/compliance 权限)+ 对话线(依赖 T-01/T-03/T-07)验收 A-6/A-7 diff --git a/docs/项目框架设计/开发计划-风控模块.md b/docs/项目框架设计/开发计划-风控模块.md index 15e0a6e..3e9235b 100644 --- a/docs/项目框架设计/开发计划-风控模块.md +++ b/docs/项目框架设计/开发计划-风控模块.md @@ -26,8 +26,8 @@ | B4 | `aml_service.py`(归一化+相似度匹配、scan_all)+ `engine.py`(process_trade_event 组装 + 预留客户事件钩子)+ **`scoring.py` 占位签名(FR-7 预留)** | 引擎完整 | 单测:AML 阈值边界;引擎集成冒烟 | B1、B2、B3。**B4 评审修复(2026-09-06):FR-4 payload 客户上下文(L0+近 30 天统计)补组装;core_ro 增 list_trades_range(升序)/list_active_customers;种子名单 3 条改名维持唯一演示命中;core_cash_flow 上下文一期未接(PRD「如有」),B8 前评估;真实名单数据接入时 payload/审计出口接 utils/desensitize(B9b 前核查项);复审观察项:30 天窗口口径(现 31 自然日)/申赎混合求和方向,R-05 接入前统一** | | B5 | `app/gateway/`(trade_gateway + gateway_repository 仅 INSERT core_trade)+ `api/simulate.py` 薄路由 | 网关 | 集成:convert 400、阻断不落 trade | A4、B4 | | B6 | **`app/api/deps.py`:`AuthContext`(actor_id/roles/customer_id,字段按 JWT 手册冻结)+ `get_auth_context()` 工厂**——dev 模式从 `X-Debug-Role`/`X-Debug-Actor` 请求头构造、`app_env != development` 启动时检测 debug 头直接拒绝;T-01 就绪后仅替换工厂内部为 JWT 解析,签名不变。另:`api/risk.py` 4 个 API(GET alerts / POST handle / POST suitability/check / POST aml/scan)+ 归属校验(含 compliance 强制 aml 过滤)。**备注:依赖层须校验 handler_result 枚举(repo 不校验);本阶段顺手统一 `NotFoundError` 异常(utils/exceptions.py 现为占位)**(评审 P2-7①②) | 鉴权依赖 + 4 个 API | Swagger 手测 + **权限矩阵(按 debug 头切换角色/身份执行 A-7/A-9 用例)** | A4、B2、**B4**(aml/scan 依赖 scan_all) | -| B7 | `main.py` 集成:路由挂载 + lifespan(双 Engine 单例注入 + Redis 单例 + trace 中间件)。**备注:顺手提取 `utils/db.py` 引擎工厂收敛 core_ro/risk_repository 双份 _default_engine**(评审 P2-6);**B3 挂账(B3 评审 P2-5):① `_run_locked` 锁原语公共化(alert_service/profile_l3 现复用私有实现)② L3 写侧 Redis 缓存 DEL 钩子(PRD §5.1 `profile:l3:{customer_id}` 更新时 DEL,`profile_l3.upsert_profile_l3` 已留痕)③ 删除 core_ro.list_trades 死代码(B4 改用 list_trades_range 后无调用方,复审 N3)** | 可运行应用 | `uvicorn` 启动 + `/health` + 全路由可达 | B5、B6 | -| B8 | **`tests/conftest.py`**:a) session fixture 启动校验演示数据就位(CUST-4001 测评 <365 天、risk_aml_list ≥8),缺失则中止并提示先跑 FLOW §0 ③④;b) fixture 幂等代跑 `prepare_risk_demo.sql`;c) teardown 按 `TRD-TEST-` 清 core_trade + 关联 risk_alert/risk_suitability_log/audit_log + 还原 L3 行。**顺手集中 sqlite 测试 DDL 为单一事实源(B4 评审 P3-12,各测试文件手写 DDL 收敛)**。集成测试:A-1~A-5、A-7(状态机/compliance 403/GET 强制 aml)、A-9 越权、**trace 一致性断言** | 测试套件 + fixture | `pytest` 全绿 | B7 | +| B7 | `main.py` 集成:路由挂载 + lifespan(双 Engine 单例注入 + Redis 单例 + trace 中间件)。**备注:顺手提取 `utils/db.py` 引擎工厂收敛 core_ro/risk_repository 双份 _default_engine**(评审 P2-6);**B3 挂账(B3 评审 P2-5):① `_run_locked` 锁原语公共化(alert_service/profile_l3 现复用私有实现)② L3 写侧 Redis 缓存 DEL 钩子(PRD §5.1 `profile:l3:{customer_id}` 更新时 DEL,`profile_l3.upsert_profile_l3` 已留痕)③ 删除 core_ro.list_trades 死代码(B4 改用 list_trades_range 后无调用方,复审 N3)④ 统一响应外壳落地(utils/response.py 现占位,simulate/risk 路由届时一并包裹,B5 评审 P2-2)** | 可运行应用 | `uvicorn` 启动 + `/health` + 全路由可达 | B5、B6 | +| B8 | **`tests/conftest.py`**:a) session fixture 启动校验演示数据就位(CUST-4001 测评 <365 天、risk_aml_list ≥8),缺失则中止并提示先跑 FLOW §0 ③④;b) fixture 幂等代跑 `prepare_risk_demo.sql`;c) teardown 按 `TRD-TEST-` 清 core_trade + 关联 risk_alert/risk_suitability_log/audit_log + 还原 L3 行。**顺手集中 sqlite 测试 DDL 为单一事实源(B4 评审 P3-12,各测试文件手写 DDL 收敛;CURRENT_TIMESTAMP 改 localtime 或 fixture 固定时间,防 UTC/本地日界错位——B5 评审 P3-4)**;集成测试交易统一走 `trade_id_factory` 注入 `TRD-TEST-` 前缀(trade_gateway 已留参数,B5 评审 P3-2)。集成测试:A-1~A-5、A-7(状态机/compliance 403/GET 强制 aml)、A-9 越权、**trace 一致性断言** | 测试套件 + fixture | `pytest` 全绿 | B7 | | B9a | 演示/运维脚本开发:`scripts/demo/subscribe_alerts.py`(订阅演示)+ `scripts/demo/rebuild_alerts.py`(按 trade_id 幂等重放补偿) | 2 个脚本 | 手工执行验证 | B2、B4(可与 B5~B8 并行) | | B9b | 演示链路走查:`reset.ps1` → `prepare_risk_demo.sql` → agent 库建表 → `seed-aml-list.sql` → Swagger 逐条过 **A-1~A-5、A-7~A-9(A-6 归 M3)** | 演示 SOP | 按 PRD §8 验收表逐条打勾 | B8、B9a | diff --git a/tests/test_trade_gateway.py b/tests/test_trade_gateway.py index 9841164..f7612a1 100644 --- a/tests/test_trade_gateway.py +++ b/tests/test_trade_gateway.py @@ -174,10 +174,49 @@ def test_accepted_trades_and_engine_fires(env): assert row["trade_status"] == "confirmed" and Decimal(str(row["amount"])) == Decimal("600000") alert = repo.get_alert(resp["alert_ids"][0]) assert alert["alert_type"] == "large_amount" and alert["risk_score"] == 70 + assert alert["status"] == "pending_review" # 评审 P3-3 加固 assert _counts(engine, "risk_suitability_log", "is_blocked=0") == 1 assert _counts(engine, "audit_log", "agent_type='platform' AND decision='trade_accepted'") == 1 assert _counts(engine, "customer_profile_l3", "monitor_tier='watch'") == 1 - assert len(pub.messages) == 1 + (channel, payload), = pub.messages + assert channel == "risk:pub:alert" and payload["risk_score"] == 70 + + +def test_redeem_accepted_without_alert(env): + """redeem 正向路径(评审 P3-3):小额赎回放行,无预警。""" + core, repo, writer, pub, engine = env + resp = submit_trade( + _req(customer="CUST-3001", product="PROD-510300", ttype="redeem", amount="1000"), + core_ro=core, risk_repo=repo, gateway_repo=writer, now=datetime(2026, 9, 6, 14, 0, 0), + ) + assert resp["blocked"] is False and resp["triggered_rules"] == [] + assert _counts(engine, "core_trade", "trade_type='redeem'") == 1 + assert _counts(engine, "risk_alert") == 0 + + +def test_engine_failure_is_audited_and_degraded(env): + """评审 P1-1:引擎异常 → 审计 risk_engine_error + 响应 engine_error=true(交易已成立)。""" + core, repo, writer, pub, engine = env + + class Boom(Exception): + pass + + monkey_patch = lambda *a, **k: (_ for _ in ()).throw(Boom()) + saved = tg.process_trade_event + tg.process_trade_event = monkey_patch + try: + resp = submit_trade( + _req(customer="CUST-3001", product="PROD-510300", amount="600000"), + core_ro=core, risk_repo=repo, gateway_repo=writer, + now=datetime(2026, 9, 6, 14, 0, 0), + ) + finally: + tg.process_trade_event = saved + assert resp["blocked"] is False and resp["engine_error"] is True + assert _counts(engine, "core_trade") == 1 # 交易已成立 + assert _counts(engine, "audit_log", "agent_type='platform' AND decision='risk_engine_error'") == 1 + assert _counts(engine, "risk_alert") == 0 # 引擎未跑,无预警 + assert pub.messages == [] def test_missing_customer_returns_lookup_error(env):