| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252 |
- from datetime import UTC, date, datetime
- import pytest
- from fastapi.testclient import TestClient
- from zbt.commands.seed import build_seed_manifest
- from zbt.core.config import Settings
- from zbt.core.errors import AppError
- from zbt.core.passwords import PasswordService
- from zbt.domains.catalog.repository import InMemoryCatalogRepository
- from zbt.domains.enrollment.models import OutboxEvent, Policy
- from zbt.domains.enrollment.repository import InMemoryEnrollmentRepository
- from zbt.domains.identity.models import AdminUser
- from zbt.domains.identity.repository import InMemoryIdentityRepository
- from zbt.domains.surrender.migration import ManualRefundMigrationService
- from zbt.domains.surrender.models import RefundAttempt, SurrenderRequest
- from zbt.domains.surrender.recovery import RefundRecoveryService
- from zbt.domains.surrender.repository import InMemorySurrenderRepository
- from zbt.main import create_app
- NOW = datetime(2026, 7, 30, 10, 0, tzinfo=UTC)
- REQUEST_ID = "01KZ0000000000000000000101"
- POLICY_ID = "01KZ0000000000000000000102"
- EVENT_ID = "01KZ0000000000000000000103"
- def admin() -> AdminUser:
- return AdminUser(
- id="01KZ0000000000000000000104",
- username="operator01",
- password_hash="hash",
- display_name="运营人员",
- status="ACTIVE",
- roles=("OPERATOR",),
- permissions=("surrender:review",),
- data_scope="MASKED_ALL",
- )
- def build_recovery(
- *,
- attempt_status: str = "SUBMITTED",
- ) -> tuple[
- RefundRecoveryService,
- InMemoryEnrollmentRepository,
- InMemorySurrenderRepository,
- ]:
- enrollments = InMemoryEnrollmentRepository()
- surrenders = InMemorySurrenderRepository()
- enrollments.save_policy(
- Policy(
- id=POLICY_ID,
- policy_no="POL-20260730-TEST01",
- order_id="01KZ0000000000000000000105",
- user_id="01KZ0000000000000000000106",
- product_id="01KZ0000000000000000000107",
- product_version_id="01KZ0000000000000000000108",
- plan_id="01KZ0000000000000000000109",
- premium_cents=36_500,
- currency="CNY",
- coverage_start=date(2026, 7, 1),
- coverage_end=date(2027, 6, 30),
- status="SURRENDERING",
- issued_at=NOW,
- )
- )
- surrenders.save_request(
- SurrenderRequest(
- id=REQUEST_ID,
- request_no="ZBT-TB-20260730-1001",
- user_id="01KZ0000000000000000000106",
- policy_id=POLICY_ID,
- idempotency_key="recovery-test",
- reason="保障计划调整",
- status="REFUND_PENDING",
- refund_amount_cents=30_000,
- currency="CNY",
- rule_version="surrender-v1",
- calculation={},
- workflow_thread_id="recovery-thread",
- created_at=NOW,
- )
- )
- surrenders.save_refund_attempt(
- RefundAttempt(
- id="01KZ0000000000000000000110",
- surrender_request_id=REQUEST_ID,
- attempt_no=1,
- refund_no="ZBT-RF-20260730-1001",
- provider="INTERNAL_REFUND",
- idempotency_key="recovery-refund",
- amount_cents=30_000,
- currency="CNY",
- status=attempt_status,
- provider_refund_no=None,
- error_code=None,
- error_message=None,
- created_at=NOW,
- updated_at=NOW,
- )
- )
- enrollments.save_outbox_event(
- OutboxEvent(
- id=EVENT_ID,
- aggregate_type="SURRENDER_REQUEST",
- aggregate_id=REQUEST_ID,
- event_type="SURRENDER_REFUND_REQUESTED",
- payload={"request_id": REQUEST_ID},
- status="DEAD",
- created_at=NOW,
- attempt_count=5,
- lease_owner="worker-old",
- leased_until=NOW,
- last_error="数据库临时不可用",
- )
- )
- return (
- RefundRecoveryService(enrollments, surrenders, lambda: NOW),
- enrollments,
- surrenders,
- )
- def test_dead_refund_task_can_be_requeued_with_a_fresh_retry_budget() -> None:
- service, enrollments, surrenders = build_recovery()
- result = service.requeue(admin(), EVENT_ID, note="数据库已恢复,重新执行")
- recovered_event = enrollments.get_outbox_event(EVENT_ID)
- request = surrenders.get_request(REQUEST_ID)
- policy = enrollments.get_policy(POLICY_ID)
- assert result["event_status"] == "PENDING"
- assert recovered_event is not None
- assert recovered_event.attempt_count == 0
- assert recovered_event.last_error is None
- assert recovered_event.payload["previous_attempt_count"] == 5
- assert request is not None and request.status == "REFUND_PENDING"
- assert policy is not None and policy.status == "SURRENDERING"
- assert surrenders.list_timeline_events(REQUEST_ID)[-1].event_type == (
- "REFUND_TASK_REQUEUED"
- )
- def test_dead_refund_task_can_be_cancelled_and_restore_policy() -> None:
- service, enrollments, surrenders = build_recovery()
- result = service.cancel(admin(), EVENT_ID, note="客户已改为继续保障")
- event = enrollments.get_outbox_event(EVENT_ID)
- request = surrenders.get_request(REQUEST_ID)
- policy = enrollments.get_policy(POLICY_ID)
- attempt = surrenders.get_latest_refund_attempt(REQUEST_ID)
- assert result["event_status"] == "CANCELLED"
- assert event is not None and event.status == "CANCELLED"
- assert request is not None and request.status == "COMPENSATED"
- assert policy is not None and policy.status == "ACTIVE"
- assert attempt is not None and attempt.status == "CANCELLED"
- assert surrenders.list_timeline_events(REQUEST_ID)[-1].event_type == (
- "REFUND_TASK_CANCELLED"
- )
- @pytest.mark.parametrize("operation", ["requeue", "cancel"])
- def test_succeeded_refund_cannot_be_recovered(operation: str) -> None:
- service, _, _ = build_recovery(attempt_status="SUCCEEDED")
- with pytest.raises(AppError) as captured:
- getattr(service, operation)(admin(), EVENT_ID, note="错误操作")
- assert captured.value.code == "REFUND_TASK_ALREADY_SETTLED"
- def test_admin_can_requeue_dead_refund_task_via_operations_api() -> None:
- _, enrollments, surrenders = build_recovery()
- identities = InMemoryIdentityRepository()
- for user in build_seed_manifest(
- PasswordService().hash("zaq1XSW@")
- ).admin_users:
- identities.save_admin_user(user)
- app = create_app(
- settings=Settings(
- app_env="test",
- jwt_access_secret="a" * 32,
- jwt_refresh_secret="b" * 32,
- field_encryption_key="c" * 32,
- ),
- identity_repository=identities,
- catalog_repository=InMemoryCatalogRepository(),
- enrollment_repository=enrollments,
- surrender_repository=surrenders,
- clock=lambda: NOW,
- )
- with TestClient(app) as client:
- login = client.post(
- "/api/v1/admin/auth/login",
- json={"username": "operator01", "password": "zaq1XSW@"},
- )
- token = login.json()["data"]["tokens"]["access_token"]
- response = client.post(
- f"/api/v1/admin/refund-operations/{EVENT_ID}/requeue",
- headers={"Authorization": f"Bearer {token}"},
- json={"note": "基础设施已经恢复"},
- )
- assert response.status_code == 200
- assert response.json()["data"]["event_status"] == "PENDING"
- assert response.json()["data"]["policy_status"] == "SURRENDERING"
- def test_manual_intervention_request_is_migrated_to_dead_task_once() -> None:
- enrollments = InMemoryEnrollmentRepository()
- surrenders = InMemorySurrenderRepository()
- surrenders.save_request(
- SurrenderRequest(
- id=REQUEST_ID,
- request_no="ZBT-TB-20260730-1002",
- user_id="01KZ0000000000000000000106",
- policy_id=POLICY_ID,
- idempotency_key="legacy-manual-intervention",
- reason="历史退款连续失败",
- status="MANUAL_INTERVENTION",
- refund_amount_cents=30_000,
- currency="CNY",
- rule_version="surrender-v1",
- calculation={},
- workflow_thread_id="legacy-thread",
- created_at=NOW,
- )
- )
- service = ManualRefundMigrationService(enrollments, surrenders, lambda: NOW)
- first = service.migrate()
- second = service.migrate()
- events = [
- event
- for event in enrollments.list_outbox_events(
- event_type="SURRENDER_REFUND_REQUESTED"
- )
- if event.aggregate_id == REQUEST_ID
- ]
- timeline = surrenders.list_timeline_events(REQUEST_ID)
- request = surrenders.get_request(REQUEST_ID)
- assert first == {"scanned": 1, "migrated": 1}
- assert second == {"scanned": 1, "migrated": 0}
- assert len(events) == 1
- assert events[0].status == "DEAD"
- assert events[0].attempt_count == 5
- assert events[0].payload["migration_key"] == "manual-intervention-v1"
- assert request is not None and request.status == "MANUAL_INTERVENTION"
- assert [event.event_type for event in timeline] == ["REFUND_TASK_MIGRATED"]
|