"""投保领域异步任务边界。 第一阶段使用进程内调度器保证本地一键运行;接口可以在后续替换为 RabbitMQ 发布器和独立 Worker,支付回调与保单签发规则无需改写。 """ from collections.abc import Callable from dataclasses import replace from datetime import datetime, timedelta from typing import Protocol from app.core.errors import AppError from app.core.identifiers import new_ulid from app.domains.enrollment.models import AsyncTask, Policy from app.domains.enrollment.repository import EnrollmentRepository class TaskDispatcher(Protocol): def dispatch(self, task: AsyncTask) -> Policy: ... class InlinePolicyTaskDispatcher: """本地运行适配器:真实记录任务生命周期,然后在当前进程执行。""" def __init__( self, repository: EnrollmentRepository, clock: Callable[[], datetime], ) -> None: self._repository = repository self._clock = clock def dispatch(self, task: AsyncTask) -> Policy: existing = self._repository.get_policy_by_order(task.business_key) if existing is not None: self._repository.save_task( replace(task, status="SUCCEEDED", attempt_count=task.attempt_count + 1) ) return existing running = replace( task, status="RUNNING", attempt_count=task.attempt_count + 1, ) self._repository.save_task(running) try: policy = self._issue_policy(task.business_key) except Exception: self._repository.save_task(replace(running, status="FAILED")) raise self._repository.save_task(replace(running, status="SUCCEEDED")) return policy def _issue_policy(self, order_id: str) -> Policy: order = self._repository.get_order(order_id) if order is None: raise AppError("ORDER_NOT_FOUND", "未找到待签发保单的订单", 404) if order.status not in {"PAID", "ISSUED"}: raise AppError("ORDER_NOT_PAID", "订单尚未支付,不能签发保单", 409) existing = self._repository.get_policy_by_order(order.id) if existing is not None: return existing now = self._clock() policy = Policy( id=new_ulid(), policy_no=f"POL-{now:%Y%m%d}-{new_ulid()[-8:]}", order_id=order.id, user_id=order.user_id, product_id=order.product_id, product_version_id=order.product_version_id, plan_id=order.plan_id, premium_cents=order.amount_cents, currency=order.currency, coverage_start=now.date() + timedelta(days=1), coverage_end=now.date() + timedelta(days=365), status="ACTIVE", issued_at=now, ) self._repository.save_policy(policy) self._repository.save_order(replace(order, status="ISSUED")) return policy