tasks.py 2.9 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283
  1. """投保领域异步任务边界。
  2. 第一阶段使用进程内调度器保证本地一键运行;接口可以在后续替换为 RabbitMQ
  3. 发布器和独立 Worker,支付回调与保单签发规则无需改写。
  4. """
  5. from collections.abc import Callable
  6. from dataclasses import replace
  7. from datetime import datetime, timedelta
  8. from typing import Protocol
  9. from app.core.errors import AppError
  10. from app.core.identifiers import new_ulid
  11. from app.domains.enrollment.models import AsyncTask, Policy
  12. from app.domains.enrollment.repository import EnrollmentRepository
  13. class TaskDispatcher(Protocol):
  14. def dispatch(self, task: AsyncTask) -> Policy: ...
  15. class InlinePolicyTaskDispatcher:
  16. """本地运行适配器:真实记录任务生命周期,然后在当前进程执行。"""
  17. def __init__(
  18. self,
  19. repository: EnrollmentRepository,
  20. clock: Callable[[], datetime],
  21. ) -> None:
  22. self._repository = repository
  23. self._clock = clock
  24. def dispatch(self, task: AsyncTask) -> Policy:
  25. existing = self._repository.get_policy_by_order(task.business_key)
  26. if existing is not None:
  27. self._repository.save_task(
  28. replace(task, status="SUCCEEDED", attempt_count=task.attempt_count + 1)
  29. )
  30. return existing
  31. running = replace(
  32. task,
  33. status="RUNNING",
  34. attempt_count=task.attempt_count + 1,
  35. )
  36. self._repository.save_task(running)
  37. try:
  38. policy = self._issue_policy(task.business_key)
  39. except Exception:
  40. self._repository.save_task(replace(running, status="FAILED"))
  41. raise
  42. self._repository.save_task(replace(running, status="SUCCEEDED"))
  43. return policy
  44. def _issue_policy(self, order_id: str) -> Policy:
  45. order = self._repository.get_order(order_id)
  46. if order is None:
  47. raise AppError("ORDER_NOT_FOUND", "未找到待签发保单的订单", 404)
  48. if order.status not in {"PAID", "ISSUED"}:
  49. raise AppError("ORDER_NOT_PAID", "订单尚未支付,不能签发保单", 409)
  50. existing = self._repository.get_policy_by_order(order.id)
  51. if existing is not None:
  52. return existing
  53. now = self._clock()
  54. policy = Policy(
  55. id=new_ulid(),
  56. policy_no=f"POL-{now:%Y%m%d}-{new_ulid()[-8:]}",
  57. order_id=order.id,
  58. user_id=order.user_id,
  59. product_id=order.product_id,
  60. product_version_id=order.product_version_id,
  61. plan_id=order.plan_id,
  62. premium_cents=order.amount_cents,
  63. currency=order.currency,
  64. coverage_start=now.date() + timedelta(days=1),
  65. coverage_end=now.date() + timedelta(days=365),
  66. status="ACTIVE",
  67. issued_at=now,
  68. )
  69. self._repository.save_policy(policy)
  70. self._repository.save_order(replace(order, status="ISSUED"))
  71. return policy