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"]