Files
group_xinghuo_jinrong/tests/test_integration_risk.py
T

386 lines
14 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""风控集成测试(B8 · PRD §8 验收 A-1~A-5/A-7/A-9 + trace 一致性 + 审计 JSON)。
真链路:TestClient(main app)(路由 + trace 中间件 + lifespan)→ 真本机 MySQL
(jinrong_core/jinrong_agent,演示数据就位校验见 conftest.ensure_risk_demo_ready)。
仅注入两点:trade_gateway._new_trade_id 统一 TRD-TEST- 前缀(teardown 按前缀
清理,B5 评审 P3-2);redis_gateway fake(断言推送,不依赖本机 Redis)。
用例间按演示时间线顺序耦合(A-3 产生的预警单供 A-7 处置),模块内保序。
"""
from __future__ import annotations
import json
from uuid import uuid4
import pytest
from fastapi.testclient import TestClient
from sqlalchemy import text
from conftest import ensure_risk_demo_ready
from app.config.settings import settings # noqa: E402
from app.gateway import trade_gateway # noqa: E402
from app.main import app # noqa: E402
from app.service.risk import redis_gateway # noqa: E402
ensure_risk_demo_ready()
OFFICER = {"X-Debug-Role": "risk_officer", "X-Debug-Actor": "STAFF-90001"}
COMPLIANCE = {"X-Debug-Role": "compliance", "X-Debug-Actor": "STAFF-40001"}
CUSTOMER_1001 = {"X-Debug-Role": "customer", "X-Debug-Actor": "CUST-1001"}
CUSTOMER_1002 = {"X-Debug-Role": "customer", "X-Debug-Actor": "CUST-1002"}
# simulate 交易白名单:risk_demo 演示账号或客户本人(FR-1 鉴权)
DEMO = {"X-Debug-Role": "risk_demo", "X-Debug-Actor": "STAFF-DEMO"}
# STAFF-10087 真实存在但名下无 CUST-3001(归属表 28 行分布于 5 个 advisor)
ADVISOR_10087 = {"X-Debug-Role": "advisor", "X-Debug-Actor": "STAFF-10087"}
# 演示时间线状态(A-1/A-3 产物供 A-7 审计与处置用例续用)
_state: dict[str, str] = {}
class FakePub:
def __init__(self):
self.messages = []
self.deletes = []
def publish(self, channel, payload):
self.messages.append((channel, payload))
def delete(self, *keys):
self.deletes.append(keys)
def _test_trade_id(now):
"""集成交易统一 TRD-TEST- 前缀(conftest teardown 定位清理)。"""
return f"TRD-TEST-{uuid4().hex[:8].upper()}"
@pytest.fixture()
def iclient(risk_demo_env, monkeypatch):
monkeypatch.setattr(trade_gateway, "_new_trade_id", _test_trade_id)
fake = FakePub()
with TestClient(app) as c: # 真 lifespan:dev 放行 + Redis 网关注册
monkeypatch.setattr(redis_gateway, "_gateway", fake) # 覆盖 lifespan 注册的真网关
yield c, fake
def _trade(customer_id, product_id, trade_type, amount):
return {
"customer_id": customer_id,
"product_id": product_id,
"trade_type": trade_type,
"amount": amount,
}
def _one(env, sql, **params):
with env["agent"].connect() as conn:
return conn.execute(text(sql), params).mappings().first()
def _core_one(env, sql, **params):
with env["core"].connect() as conn:
return conn.execute(text(sql), params).mappings().first()
def _audit_by_trade(env, trade_id, event_type=None):
sql = (
"SELECT * FROM audit_log WHERE event_type = :et"
" AND JSON_UNQUOTE(JSON_EXTRACT(input_summary, '$.trade_id')) = :tid"
)
return _one(env, sql, et=event_type or "trade_request", tid=trade_id)
# ---------- A-1:C1 申购 R4 → SUIT-001 阻断(含 trace 一致性) ----------
def test_a1_suit001_blocked_with_trace_consistency(iclient, risk_demo_env):
c, fake = iclient
env = risk_demo_env
trace = f"trc-integ-{uuid4().hex[:12]}"
r = c.post(
"/api/simulate/trade",
json=_trade("CUST-1001", "PROD-161725", "subscribe", 10000),
headers={**CUSTOMER_1001, "X-Trace-Id": trace},
)
assert r.status_code == 200
body = r.json()
assert body["blocked"] is True
assert "SUIT-001" in body["block_reason"] # C1 仅可购 R1
assert r.headers["X-Trace-Id"] == trace # 中间件透传(PRD §7.5 全链路环①)
trade_id = body["trade_id"]
assert trade_id.startswith("TRD-TEST-")
# 阻断不落交易(FR-1)
assert _core_one(env, "SELECT * FROM core_trade WHERE trade_id = :t", t=trade_id) is None
# 校验日志 blocked,trace 一致(环②)
slog = _one(
env,
"SELECT * FROM risk_suitability_log WHERE request_ref = :t",
t=trade_id,
)
assert slog is not None and slog["is_blocked"] == 1 and slog["is_matched"] == 0
assert slog["trace_id"] == trace
# suitability 预警单生成(R-02 阻断单),trace 与请求一致
alert = _one(env, "SELECT * FROM risk_alert WHERE trade_id = :t", t=trade_id)
assert alert["alert_type"] == "suitability" and alert["status"] == "pending_review"
assert alert["trace_id"] == trace
# 审计可查(platform 阻断留痕),trace 一致
audit = _audit_by_trade(env, trade_id, "trade_request")
assert audit["decision"] == "suitability_blocked" and audit["trace_id"] == trace
assert audit["agent_type"] == "platform"
# Redis 推送 trace 一致
pushed = [p for ch, p in fake.messages if ch == "risk:pub:alert"]
assert any(p["alert_id"] == alert["alert_id"] and p["trace_id"] == trace for p in pushed)
_state["a1_trade_id"] = trade_id
def test_a1_audit_input_summary_json(iclient, risk_demo_env):
"""B5 复审 L1:platform 审计 input_summary 可 JSON 解析(含阻断 reasons)。"""
c, _ = iclient
env = risk_demo_env
audit = _audit_by_trade(env, _state["a1_trade_id"], "trade_request")
summary = json.loads(audit["input_summary"])
assert summary["trade_id"] == _state["a1_trade_id"]
assert summary["product_id"] == "PROD-161725"
assert summary["trade_type"] == "subscribe"
assert summary["amount"] == "10000"
assert "SUIT-001" in summary["block_reason"]
assert isinstance(summary["reasons"], list) and summary["reasons"]
# ---------- A-2:70 岁 C5 → SUIT-006 封顶 + SUIT-003,无 SUIT-008 ----------
def test_a2_age70_cap_suit006_no_suit008(iclient, risk_demo_env):
c, _ = iclient
env = risk_demo_env
r = c.post(
"/api/simulate/trade",
json=_trade("CUST-4001", "PROD-161725", "subscribe", 20000),
headers={**DEMO, "X-Trace-Id": f"trc-integ-{uuid4().hex[:12]}"},
)
assert r.status_code == 200
body = r.json()
assert body["blocked"] is True
reasons = body["reasons"]
assert any("SUIT-006" in x for x in reasons), reasons # ≥70 按 C3 封顶
assert any("SUIT-003" in x for x in reasons), reasons # C3 < R4 不匹配
assert not any("SUIT-008" in x for x in reasons), reasons # 测评已刷新,不得干扰
trade_id = body["trade_id"]
assert _core_one(env, "SELECT * FROM core_trade WHERE trade_id = :t", t=trade_id) is None
_state["a2_trade_id"] = trade_id
# ---------- A-3:C3 单笔 50 万 R3 → 放行 + RISK-001/002 预警 + 推送 ----------
def test_a3_large_amount_alert_and_publish(iclient, risk_demo_env):
c, fake = iclient
env = risk_demo_env
r = c.post(
"/api/simulate/trade",
json=_trade("CUST-3001", "PROD-510300", "subscribe", 500000),
headers=DEMO,
)
assert r.status_code == 200
body = r.json()
assert body["blocked"] is False
assert body["triggered_rules"] == ["RISK-001", "RISK-002"]
assert body["aml_hit"] is False
trade_id = body["trade_id"]
alert_id = body["alert_ids"][0]
assert _core_one(env, "SELECT * FROM core_trade WHERE trade_id = :t", t=trade_id) is not None
alert = _one(env, "SELECT * FROM risk_alert WHERE alert_id = :a", a=alert_id)
assert alert["status"] == "pending_review" and alert["risk_score"] == 70
assert json.loads(alert["triggered_rules"]) == ["RISK-001", "RISK-002"]
pushed = [p for ch, p in fake.messages if ch == "risk:pub:alert"]
assert any(p["alert_id"] == alert_id for p in pushed)
# 放行审计全量引擎输出(B5 复审 L1)
audit = _audit_by_trade(env, trade_id, "trade_request")
assert audit["decision"] == "trade_accepted"
summary = json.loads(audit["input_summary"])
assert summary["triggered_rules"] == ["RISK-001", "RISK-002"]
assert summary["alert_ids"] == [alert_id]
_state["a3_alert_id"] = alert_id
_state["a3_trade_id"] = trade_id
# ---------- A-4:同产品当日第 3 笔 → freq 并入当日预警单 ----------
def test_a4_freq_trade_merged_into_same_alert(iclient, risk_demo_env):
c, _ = iclient
env = risk_demo_env
alert_id = None
for _ in range(3):
r = c.post(
"/api/simulate/trade",
json=_trade("CUST-9527", "PROD-510300", "subscribe", 1000),
headers=DEMO,
)
assert r.status_code == 200
body = r.json()
assert body["blocked"] is False
if body["triggered_rules"]:
assert body["triggered_rules"] == ["RISK-003"]
alert_id = body["alert_ids"][0]
assert alert_id, "第 3 笔应触发 RISK-003"
# 第 4 笔:并入既有单,不另开新单(PRD FR-4)
r = c.post(
"/api/simulate/trade",
json=_trade("CUST-9527", "PROD-510300", "subscribe", 1000),
headers=DEMO,
)
body = r.json()
assert body["blocked"] is False and body["triggered_rules"] == ["RISK-003"]
assert body["alert_ids"] == [alert_id]
alerts = _one(
env,
"SELECT COUNT(*) AS n FROM risk_alert WHERE customer_id = 'CUST-9527'"
" AND alert_type = 'freq_trade'",
)
assert alerts["n"] == 1
row = _one(env, "SELECT payload FROM risk_alert WHERE alert_id = :a", a=alert_id)
assert len(json.loads(row["payload"])["events"]) == 2
# ---------- A-5:AML 命中 → 独立单 + L3 high + compliance 可见 + 不冻户 ----------
def test_a5_aml_hit_independent_alert(iclient, risk_demo_env):
c, fake = iclient
env = risk_demo_env
r = c.post(
"/api/simulate/trade",
json=_trade("CUST-1002", "PROD-005828", "subscribe", 100), # C2+R2 匹配,不阻断
headers=DEMO,
)
assert r.status_code == 200
body = r.json()
assert body["blocked"] is False and body["aml_hit"] is True
trade_id = body["trade_id"]
alert = _one(
env,
"SELECT * FROM risk_alert WHERE customer_id = 'CUST-1002' AND alert_type = 'aml'",
)
assert alert is not None
assert alert["risk_score"] == 95 and alert["status"] == "pending_review"
assert json.loads(alert["triggered_rules"]) == ["AML-001"]
# L3 置 high(FR-7)
l3 = _one(env, "SELECT * FROM customer_profile_l3 WHERE customer_id = 'CUST-1002'")
assert l3["monitor_tier"] == "high"
# 紧急推送含 compliance
pushed = [p for ch, p in fake.messages if ch == "risk:pub:alert"]
aml_push = [p for p in pushed if p["alert_id"] == alert["alert_id"]]
assert aml_push and "compliance" in aml_push[-1]["notify_role"]
# compliance 账号台账可见该单(A-7 强制 aml 过滤的另一面)
r = c.get("/api/risk/alerts", headers=COMPLIANCE)
assert r.status_code == 200
items = r.json()["items"]
assert items and all(i["alert_type"] == "aml" for i in items)
assert any(i["alert_id"] == alert["alert_id"] for i in items)
# 账户未被冻结(R-03 禁止自动冻户)
row = _core_one(env, "SELECT is_active FROM core_customer WHERE customer_id = 'CUST-1002'")
assert row["is_active"] == 1
_state["a5_alert_id"] = alert["alert_id"]
# ---------- A-7:人工处置状态机 + compliance 403 + GET 仅 aml ----------
def test_a7_handle_state_machine_compliance_forbidden(iclient, risk_demo_env):
c, _ = iclient
env = risk_demo_env
alert_id = _state["a3_alert_id"]
r = c.post(
f"/api/risk/alerts/{alert_id}/handle",
json={"handler_result": "confirmed_suspicious", "handler_comment": "确认可疑"},
headers=OFFICER,
)
assert r.status_code == 200
body = r.json()
assert body["status"] == "confirmed_suspicious" and body["handler_id"] == "STAFF-90001"
# 状态机:已处置单禁止跳改
r2 = c.post(
f"/api/risk/alerts/{alert_id}/handle",
json={"handler_result": "confirmed_normal"},
headers=OFFICER,
)
assert r2.status_code == 409
# compliance 无处置权
r3 = c.post(
f"/api/risk/alerts/{alert_id}/handle",
json={"handler_result": "confirmed_normal"},
headers=COMPLIANCE,
)
assert r3.status_code == 403
# 处置审计留痕(同事务,agent_type='risk')
audit = _one(
env,
"SELECT * FROM audit_log WHERE event_type = 'alert_handle' AND actor_id = 'STAFF-90001'"
" AND decision = 'alert_handled' AND handler_result = 'confirmed_suspicious'",
)
assert audit is not None
# ---------- A-9:越权 403 + 审计 ----------
def test_a9_cross_customer_and_unassigned_advisor_403(iclient, risk_demo_env):
c, _ = iclient
env = risk_demo_env
# customer 查他人
r = c.post(
"/api/risk/suitability/check",
json={"customer_id": "CUST-3001", "product_id": "PROD-510300"},
headers=CUSTOMER_1002,
)
assert r.status_code == 403 and r.json()["error_code"] == "AUTH_403_NOT_OWNER"
# advisor 查非名下(STAFF-10086 名下无 CUST-3001)
r = c.post(
"/api/risk/suitability/check",
json={"customer_id": "CUST-3001", "product_id": "PROD-510300"},
headers=ADVISOR_10087, # 名下无 CUST-3001
)
assert r.status_code == 403 and r.json()["error_code"] == "AUTH_403_NOT_ASSIGNED"
denials = _one(
env,
"SELECT COUNT(*) AS n FROM audit_log WHERE event_type = 'authz' AND decision = 'forbidden'"
" AND actor_id IN ('CUST-1002', 'STAFF-10087')",
)
assert denials["n"] == 2
# ---------- 参数校验与无 trace 头兜底 ----------
def test_convert_400_and_no_new_trade_audit(iclient, risk_demo_env):
c, _ = iclient
env = risk_demo_env
before = _one(env, "SELECT COUNT(*) AS n FROM audit_log WHERE event_type = 'trade_request'")
r = c.post(
"/api/simulate/trade",
json=_trade("CUST-3001", "PROD-510300", "convert", 1000),
headers=DEMO,
)
assert r.status_code == 400 and r.json()["error_code"] == "BAD_REQUEST"
after = _one(env, "SELECT COUNT(*) AS n FROM audit_log WHERE event_type = 'trade_request'")
assert after["n"] == before["n"] # convert 属参数校验失败,不落审计
def test_missing_trace_header_generates_one(iclient):
c, _ = iclient
r = c.get("/api/risk/alerts", headers=OFFICER)
assert r.status_code == 200
assert r.headers["X-Trace-Id"].startswith("trc-")
def test_settings_thresholds_loaded():
"""冒烟:阈值配置与冻结规则一致(.env 未覆盖时)。"""
assert settings.risk_large_amount == 500000
assert settings.risk_assessment_valid_days == 365