fix: 引擎异常审计降级/审计全量输入输出/TRD-TEST 注入点(B5 评审 P1-1、P2-1~3、P3 批量)

This commit is contained in:
2026-09-06 18:17:44 +08:00
parent 2a95b65ec8
commit 46a1f24870
5 changed files with 103 additions and 23 deletions
+5 -2
View File
@@ -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
+55 -17
View File
@@ -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}
+1 -1
View File
@@ -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
@@ -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 |
+40 -1
View File
@@ -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):