docs/05 §5.1 把"业务写接口"列为必须携带 Idempotency-Key 的接口,风控 6 个 POST (手工扫描、确认接收、进入调查、关闭误报、完成结案、升级处理)此前一个都没带, 重复提交会二次驱动状态机。 复用平台的 api_request_receipt 与 ApiTransactionService,但新增 execute_in:原来的 execute 自己开 SessionFactory() 和 session.begin(),而 RiskActionService._finish 会在内部 commit,套进去就成了"内层提交外层事务"。execute_in 改为在调用方传入的 session 上读写幂等记录,幂等记录因此与业务写入同处一个事务(§5.2)。 scope 用实际路径(含 alert_no)而不是路由模板:§5.1 的幂等范围是 user_id + method + normalized_path + idempotency_key,把路径参数折成模板会让同一个键 在不同预警之间互相回放 —— 那是把两次不同资源的操作当成一次。 测试:单测加幂等透传替身与"缺键即 422"用例;新增 tests/integration/test_risk_idempotency_mysql.py,覆盖同键回放不重复执行、同键不同 正文 409、缺键/非 ASCII 拒绝,以及端到端"重复 POST 只调用一次处置逻辑"。
179 lines
7.1 KiB
Python
179 lines
7.1 KiB
Python
"""风控写接口的幂等语义(`docs/05` §5.1、§5.2;`docs/25` P3 #22)。
|
||
|
||
两级覆盖:
|
||
|
||
1. `ApiTransactionService.execute_in` 本身 —— 同键同正文回放原响应且**不重复执行**、
|
||
同键不同正文返回 `409 IDEMPOTENCY_CONFLICT`、缺键直接拒绝;
|
||
2. Controller 接线 —— 重复 POST 同一个风控处置端点时,Action Service 只被调用一次。
|
||
|
||
刻意不去驱动真实状态机:处置动作会改预警状态,测试不该污染演示数据。这里用计数替身
|
||
验证"第二次请求没有落到业务逻辑上",这才是幂等要保证的事情。
|
||
"""
|
||
|
||
from collections.abc import AsyncIterator
|
||
from typing import Any
|
||
from uuid import uuid4
|
||
|
||
import httpx
|
||
import pytest
|
||
from sqlalchemy import delete
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
from app.api.controllers import risk as risk_controller
|
||
from app.api.dependencies.auth import build_request_context
|
||
from app.api.dependencies.database import get_session
|
||
from app.core.contracts import RequestContext
|
||
from app.core.errors import IdempotencyConflictError, ValidationAgentError
|
||
from app.infrastructure.db import SessionFactory
|
||
from app.main import app
|
||
from app.repository.platform_repository import PlatformRepository
|
||
from app.service.api_transaction_service import ApiTransactionService
|
||
|
||
SCOPE = "POST /api/v1/risk/alerts/ALERT-IDEM-TEST/acknowledgements"
|
||
EXCLUSION_SCOPE = "POST /api/v1/risk/alerts/ALERT-IDEM-TEST/exclusions"
|
||
ACK_PATH = "/api/v1/risk/alerts/ALERT-IDEM-TEST/acknowledgements"
|
||
EXCLUSION_PATH = "/api/v1/risk/alerts/ALERT-IDEM-TEST/exclusions"
|
||
|
||
|
||
async def override_context() -> RequestContext:
|
||
return RequestContext(
|
||
user_id="990000002",
|
||
trace_id=str(uuid4()),
|
||
roles=("risk_operator",),
|
||
permissions=("risk:alert:read", "risk:alert:write"),
|
||
data_scope="all",
|
||
)
|
||
|
||
|
||
async def override_session() -> AsyncIterator[AsyncSession]:
|
||
async with SessionFactory() as session:
|
||
yield session
|
||
|
||
|
||
async def purge(*keys: str) -> None:
|
||
async with SessionFactory() as session, session.begin():
|
||
table = await PlatformRepository(session).table("api_request_receipt")
|
||
for key in keys:
|
||
await session.execute(delete(table).where(table.c.idempotency_key == key))
|
||
|
||
|
||
@pytest.mark.integration
|
||
@pytest.mark.asyncio
|
||
async def test_execute_in_replays_without_running_action_twice() -> None:
|
||
key = f"risk-idem-{uuid4()}"
|
||
context = await override_context()
|
||
calls: list[int] = []
|
||
|
||
async def action(_session: AsyncSession) -> dict[str, Any]:
|
||
calls.append(1)
|
||
return {"alert_no": "ALERT-IDEM-TEST", "status": "待处理"}
|
||
|
||
try:
|
||
async with SessionFactory() as session:
|
||
service = ApiTransactionService()
|
||
first = await service.execute_in(session, context, SCOPE, key, {}, action)
|
||
second = await service.execute_in(session, context, SCOPE, key, {}, action)
|
||
|
||
assert first == second == {"alert_no": "ALERT-IDEM-TEST", "status": "待处理"}
|
||
assert len(calls) == 1, "重复请求必须回放 response_json,而不是再次执行 action"
|
||
finally:
|
||
await purge(key)
|
||
|
||
|
||
@pytest.mark.integration
|
||
@pytest.mark.asyncio
|
||
async def test_execute_in_conflicts_on_same_key_with_different_body() -> None:
|
||
key = f"risk-idem-{uuid4()}"
|
||
context = await override_context()
|
||
|
||
async def action(_session: AsyncSession) -> dict[str, Any]:
|
||
return {"ok": True}
|
||
|
||
try:
|
||
async with SessionFactory() as session:
|
||
service = ApiTransactionService()
|
||
await service.execute_in(session, context, SCOPE, key, {"reason": "第一次"}, action)
|
||
with pytest.raises(IdempotencyConflictError):
|
||
await service.execute_in(
|
||
session, context, SCOPE, key, {"reason": "第二次"}, action
|
||
)
|
||
finally:
|
||
await purge(key)
|
||
|
||
|
||
@pytest.mark.integration
|
||
@pytest.mark.asyncio
|
||
async def test_execute_in_rejects_missing_or_short_key() -> None:
|
||
context = await override_context()
|
||
|
||
async def action(_session: AsyncSession) -> dict[str, Any]:
|
||
raise AssertionError("缺键时不应执行 action")
|
||
|
||
async with SessionFactory() as session:
|
||
service = ApiTransactionService()
|
||
for bad in (None, "too-short", "带中文字符的-key-1234567890"):
|
||
with pytest.raises(ValidationAgentError):
|
||
await service.execute_in(session, context, SCOPE, bad, {}, action)
|
||
|
||
|
||
@pytest.mark.integration
|
||
@pytest.mark.asyncio
|
||
async def test_repeated_risk_post_calls_action_service_once(monkeypatch) -> None:
|
||
"""端到端:同一 `Idempotency-Key` 重复 POST 只触发一次处置逻辑。
|
||
|
||
用 `httpx.ASGITransport` 而不是 `TestClient`:后者自建事件循环,测试结束后的
|
||
`SessionFactory` 清理会落在另一个循环上,连接池析构时报
|
||
`AttributeError: 'NoneType' object has no attribute 'send'`。
|
||
"""
|
||
calls: list[str] = []
|
||
|
||
class CountingActionService:
|
||
def __init__(self, _session: Any) -> None:
|
||
pass
|
||
|
||
async def acknowledge(self, alert_no: str, _context: RequestContext) -> dict[str, Any]:
|
||
calls.append(alert_no)
|
||
return {"alert_no": alert_no, "status": "待处理", "ack_status": "已确认"}
|
||
|
||
async def exclude(
|
||
self, alert_no: str, reason: str, _context: RequestContext
|
||
) -> dict[str, Any]:
|
||
calls.append(f"exclude:{reason}")
|
||
return {"alert_no": alert_no, "status": "已排除", "handle_result": reason}
|
||
|
||
monkeypatch.setattr(risk_controller, "RiskActionService", CountingActionService)
|
||
ack_key = f"risk-ack-{uuid4()}"
|
||
exclude_key = f"risk-exclude-{uuid4()}"
|
||
app.dependency_overrides[build_request_context] = override_context
|
||
app.dependency_overrides[get_session] = override_session
|
||
try:
|
||
transport = httpx.ASGITransport(app=app)
|
||
async with httpx.AsyncClient(
|
||
transport=transport, base_url="http://testserver"
|
||
) as client:
|
||
first = await client.post(ACK_PATH, headers={"Idempotency-Key": ack_key})
|
||
replay = await client.post(ACK_PATH, headers={"Idempotency-Key": ack_key})
|
||
excluded = await client.post(
|
||
EXCLUSION_PATH,
|
||
json={"reason": "客户本人确认"},
|
||
headers={"Idempotency-Key": exclude_key},
|
||
)
|
||
conflict = await client.post(
|
||
EXCLUSION_PATH,
|
||
json={"reason": "换了理由"},
|
||
headers={"Idempotency-Key": exclude_key},
|
||
)
|
||
|
||
assert first.status_code == 200
|
||
assert replay.status_code == 200
|
||
# 只比 `data`:`meta.trace_id` 标识的是**本次**请求,重放也必须换一个新的。
|
||
assert replay.json()["data"] == first.json()["data"]
|
||
assert excluded.status_code == 200
|
||
assert conflict.status_code == 409
|
||
assert conflict.json()["error"]["code"] == "IDEMPOTENCY_CONFLICT"
|
||
assert calls == ["ALERT-IDEM-TEST", "exclude:客户本人确认"]
|
||
finally:
|
||
app.dependency_overrides.pop(build_request_context, None)
|
||
app.dependency_overrides.pop(get_session, None)
|
||
await purge(ack_key, exclude_key)
|