test_surrender_recovery.py 9.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252
  1. from datetime import UTC, date, datetime
  2. import pytest
  3. from fastapi.testclient import TestClient
  4. from zbt.commands.seed import build_seed_manifest
  5. from zbt.core.config import Settings
  6. from zbt.core.errors import AppError
  7. from zbt.core.passwords import PasswordService
  8. from zbt.domains.catalog.repository import InMemoryCatalogRepository
  9. from zbt.domains.enrollment.models import OutboxEvent, Policy
  10. from zbt.domains.enrollment.repository import InMemoryEnrollmentRepository
  11. from zbt.domains.identity.models import AdminUser
  12. from zbt.domains.identity.repository import InMemoryIdentityRepository
  13. from zbt.domains.surrender.migration import ManualRefundMigrationService
  14. from zbt.domains.surrender.models import RefundAttempt, SurrenderRequest
  15. from zbt.domains.surrender.recovery import RefundRecoveryService
  16. from zbt.domains.surrender.repository import InMemorySurrenderRepository
  17. from zbt.main import create_app
  18. NOW = datetime(2026, 7, 30, 10, 0, tzinfo=UTC)
  19. REQUEST_ID = "01KZ0000000000000000000101"
  20. POLICY_ID = "01KZ0000000000000000000102"
  21. EVENT_ID = "01KZ0000000000000000000103"
  22. def admin() -> AdminUser:
  23. return AdminUser(
  24. id="01KZ0000000000000000000104",
  25. username="operator01",
  26. password_hash="hash",
  27. display_name="运营人员",
  28. status="ACTIVE",
  29. roles=("OPERATOR",),
  30. permissions=("surrender:review",),
  31. data_scope="MASKED_ALL",
  32. )
  33. def build_recovery(
  34. *,
  35. attempt_status: str = "SUBMITTED",
  36. ) -> tuple[
  37. RefundRecoveryService,
  38. InMemoryEnrollmentRepository,
  39. InMemorySurrenderRepository,
  40. ]:
  41. enrollments = InMemoryEnrollmentRepository()
  42. surrenders = InMemorySurrenderRepository()
  43. enrollments.save_policy(
  44. Policy(
  45. id=POLICY_ID,
  46. policy_no="POL-20260730-TEST01",
  47. order_id="01KZ0000000000000000000105",
  48. user_id="01KZ0000000000000000000106",
  49. product_id="01KZ0000000000000000000107",
  50. product_version_id="01KZ0000000000000000000108",
  51. plan_id="01KZ0000000000000000000109",
  52. premium_cents=36_500,
  53. currency="CNY",
  54. coverage_start=date(2026, 7, 1),
  55. coverage_end=date(2027, 6, 30),
  56. status="SURRENDERING",
  57. issued_at=NOW,
  58. )
  59. )
  60. surrenders.save_request(
  61. SurrenderRequest(
  62. id=REQUEST_ID,
  63. request_no="ZBT-TB-20260730-1001",
  64. user_id="01KZ0000000000000000000106",
  65. policy_id=POLICY_ID,
  66. idempotency_key="recovery-test",
  67. reason="保障计划调整",
  68. status="REFUND_PENDING",
  69. refund_amount_cents=30_000,
  70. currency="CNY",
  71. rule_version="surrender-v1",
  72. calculation={},
  73. workflow_thread_id="recovery-thread",
  74. created_at=NOW,
  75. )
  76. )
  77. surrenders.save_refund_attempt(
  78. RefundAttempt(
  79. id="01KZ0000000000000000000110",
  80. surrender_request_id=REQUEST_ID,
  81. attempt_no=1,
  82. refund_no="ZBT-RF-20260730-1001",
  83. provider="INTERNAL_REFUND",
  84. idempotency_key="recovery-refund",
  85. amount_cents=30_000,
  86. currency="CNY",
  87. status=attempt_status,
  88. provider_refund_no=None,
  89. error_code=None,
  90. error_message=None,
  91. created_at=NOW,
  92. updated_at=NOW,
  93. )
  94. )
  95. enrollments.save_outbox_event(
  96. OutboxEvent(
  97. id=EVENT_ID,
  98. aggregate_type="SURRENDER_REQUEST",
  99. aggregate_id=REQUEST_ID,
  100. event_type="SURRENDER_REFUND_REQUESTED",
  101. payload={"request_id": REQUEST_ID},
  102. status="DEAD",
  103. created_at=NOW,
  104. attempt_count=5,
  105. lease_owner="worker-old",
  106. leased_until=NOW,
  107. last_error="数据库临时不可用",
  108. )
  109. )
  110. return (
  111. RefundRecoveryService(enrollments, surrenders, lambda: NOW),
  112. enrollments,
  113. surrenders,
  114. )
  115. def test_dead_refund_task_can_be_requeued_with_a_fresh_retry_budget() -> None:
  116. service, enrollments, surrenders = build_recovery()
  117. result = service.requeue(admin(), EVENT_ID, note="数据库已恢复,重新执行")
  118. recovered_event = enrollments.get_outbox_event(EVENT_ID)
  119. request = surrenders.get_request(REQUEST_ID)
  120. policy = enrollments.get_policy(POLICY_ID)
  121. assert result["event_status"] == "PENDING"
  122. assert recovered_event is not None
  123. assert recovered_event.attempt_count == 0
  124. assert recovered_event.last_error is None
  125. assert recovered_event.payload["previous_attempt_count"] == 5
  126. assert request is not None and request.status == "REFUND_PENDING"
  127. assert policy is not None and policy.status == "SURRENDERING"
  128. assert surrenders.list_timeline_events(REQUEST_ID)[-1].event_type == (
  129. "REFUND_TASK_REQUEUED"
  130. )
  131. def test_dead_refund_task_can_be_cancelled_and_restore_policy() -> None:
  132. service, enrollments, surrenders = build_recovery()
  133. result = service.cancel(admin(), EVENT_ID, note="客户已改为继续保障")
  134. event = enrollments.get_outbox_event(EVENT_ID)
  135. request = surrenders.get_request(REQUEST_ID)
  136. policy = enrollments.get_policy(POLICY_ID)
  137. attempt = surrenders.get_latest_refund_attempt(REQUEST_ID)
  138. assert result["event_status"] == "CANCELLED"
  139. assert event is not None and event.status == "CANCELLED"
  140. assert request is not None and request.status == "COMPENSATED"
  141. assert policy is not None and policy.status == "ACTIVE"
  142. assert attempt is not None and attempt.status == "CANCELLED"
  143. assert surrenders.list_timeline_events(REQUEST_ID)[-1].event_type == (
  144. "REFUND_TASK_CANCELLED"
  145. )
  146. @pytest.mark.parametrize("operation", ["requeue", "cancel"])
  147. def test_succeeded_refund_cannot_be_recovered(operation: str) -> None:
  148. service, _, _ = build_recovery(attempt_status="SUCCEEDED")
  149. with pytest.raises(AppError) as captured:
  150. getattr(service, operation)(admin(), EVENT_ID, note="错误操作")
  151. assert captured.value.code == "REFUND_TASK_ALREADY_SETTLED"
  152. def test_admin_can_requeue_dead_refund_task_via_operations_api() -> None:
  153. _, enrollments, surrenders = build_recovery()
  154. identities = InMemoryIdentityRepository()
  155. for user in build_seed_manifest(
  156. PasswordService().hash("zaq1XSW@")
  157. ).admin_users:
  158. identities.save_admin_user(user)
  159. app = create_app(
  160. settings=Settings(
  161. app_env="test",
  162. jwt_access_secret="a" * 32,
  163. jwt_refresh_secret="b" * 32,
  164. field_encryption_key="c" * 32,
  165. ),
  166. identity_repository=identities,
  167. catalog_repository=InMemoryCatalogRepository(),
  168. enrollment_repository=enrollments,
  169. surrender_repository=surrenders,
  170. clock=lambda: NOW,
  171. )
  172. with TestClient(app) as client:
  173. login = client.post(
  174. "/api/v1/admin/auth/login",
  175. json={"username": "operator01", "password": "zaq1XSW@"},
  176. )
  177. token = login.json()["data"]["tokens"]["access_token"]
  178. response = client.post(
  179. f"/api/v1/admin/refund-operations/{EVENT_ID}/requeue",
  180. headers={"Authorization": f"Bearer {token}"},
  181. json={"note": "基础设施已经恢复"},
  182. )
  183. assert response.status_code == 200
  184. assert response.json()["data"]["event_status"] == "PENDING"
  185. assert response.json()["data"]["policy_status"] == "SURRENDERING"
  186. def test_manual_intervention_request_is_migrated_to_dead_task_once() -> None:
  187. enrollments = InMemoryEnrollmentRepository()
  188. surrenders = InMemorySurrenderRepository()
  189. surrenders.save_request(
  190. SurrenderRequest(
  191. id=REQUEST_ID,
  192. request_no="ZBT-TB-20260730-1002",
  193. user_id="01KZ0000000000000000000106",
  194. policy_id=POLICY_ID,
  195. idempotency_key="legacy-manual-intervention",
  196. reason="历史退款连续失败",
  197. status="MANUAL_INTERVENTION",
  198. refund_amount_cents=30_000,
  199. currency="CNY",
  200. rule_version="surrender-v1",
  201. calculation={},
  202. workflow_thread_id="legacy-thread",
  203. created_at=NOW,
  204. )
  205. )
  206. service = ManualRefundMigrationService(enrollments, surrenders, lambda: NOW)
  207. first = service.migrate()
  208. second = service.migrate()
  209. events = [
  210. event
  211. for event in enrollments.list_outbox_events(
  212. event_type="SURRENDER_REFUND_REQUESTED"
  213. )
  214. if event.aggregate_id == REQUEST_ID
  215. ]
  216. timeline = surrenders.list_timeline_events(REQUEST_ID)
  217. request = surrenders.get_request(REQUEST_ID)
  218. assert first == {"scanned": 1, "migrated": 1}
  219. assert second == {"scanned": 1, "migrated": 0}
  220. assert len(events) == 1
  221. assert events[0].status == "DEAD"
  222. assert events[0].attempt_count == 5
  223. assert events[0].payload["migration_key"] == "manual-intervention-v1"
  224. assert request is not None and request.status == "MANUAL_INTERVENTION"
  225. assert [event.event_type for event in timeline] == ["REFUND_TASK_MIGRATED"]