"""串行领取、落库并按任务类型分派一条任务。""" from __future__ import annotations from dataclasses import dataclass from typing import Callable, Mapping, Optional from .admin_gateway import ( AdminGateway, AdminGatewayError, AdminTask, AndroidDeviceInfo, ClaimCapabilities, ClientInfo, ) from .collect_task_service import CollectServiceFactory, CollectTaskService from .android_device_service import ( AndroidDeviceSearchError, AndroidDeviceService, ) from .purchase_task_service import ( LivePurchaseAdapterFactory, PurchaseAdapterFactory, PurchaseTaskService, ) from .purchase_reconcile_service import ( PurchaseReconcileFactory, PurchaseReconcileService, ) from .performance_timing import TaskPerformanceTrace from .task_models import ( NewClaimedTask, OutboxEventRecord, OutboxEventType, TaskType, ) from .task_repository import DuplicateTaskError, TaskRepository @dataclass(frozen=True) class TaskDispatchOutcome: """一轮串行任务处理的结果。""" kind: str message: str task_id: str = "" def admin_task_to_new_claimed_task(task: AdminTask) -> NewClaimedTask: """严格校验 Admin 任务并映射为本地已领取任务。""" payload = dict(task.payload) goods_url = payload.get("goods_url") if not isinstance(goods_url, str) or not goods_url.strip(): raise ValueError("payload.goods_url 不能为空") goods_id = payload.get("goods_id") if goods_id is not None and not isinstance(goods_id, str): raise ValueError("payload.goods_id 必须是文本") target_color = None target_size = None price_cent = None quantity = None if task.task_type is TaskType.PURCHASE: if not isinstance(goods_id, str) or not goods_id.strip(): raise ValueError("采购任务 payload.goods_id 不能为空") options = payload.get("options") if not isinstance(options, Mapping) or not options: raise ValueError("采购任务 payload.options 必须是非空对象") normalized_options: dict[str, str] = {} for key, value in options.items(): if not isinstance(key, str) or not key.strip(): raise ValueError("采购任务 options 的名称不能为空") if not isinstance(value, str) or not value.strip(): raise ValueError("采购任务 options 的值必须是非空文本") normalized_options[key.strip()] = value.strip() payload["options"] = normalized_options quantity = payload.get("quantity") price_cent = payload.get("max_price_cent") if ( isinstance(quantity, bool) or not isinstance(quantity, int) or quantity <= 0 ): raise ValueError("采购任务 payload.quantity 必须是大于 0 的整数") if ( isinstance(price_cent, bool) or not isinstance(price_cent, int) or price_cent <= 0 ): raise ValueError( "采购任务 payload.max_price_cent 必须是大于 0 的人民币订单总价上限(整数分)" ) target_color = normalized_options.get("color") target_size = normalized_options.get("size") original_task = { "id": task.task_id, "type": task.task_type.value, "execution_mode": task.execution_mode, "version": task.version, "priority": task.priority, "payload": payload, "created_at": task.created_at, "updated_at": task.updated_at, } return NewClaimedTask( remote_task_id=task.task_id, task_type=task.task_type, goods_url=goods_url.strip(), execution_mode=task.execution_mode, goods_id=goods_id.strip() if isinstance(goods_id, str) else None, target_color=target_color, target_size=target_size, price_cent=price_cent, quantity=quantity, priority=task.priority, version=task.version, admin_payload=original_task, ) class TaskDispatcher: """一次只补交、执行或领取并执行一条任务。""" def __init__( self, gateway: AdminGateway, repository: TaskRepository, client: ClientInfo, device_address: str, *, collect_service_factory: Optional[CollectServiceFactory] = None, purchase_adapter_factory: Optional[PurchaseAdapterFactory] = None, live_purchase_adapter_factory: Optional[ LivePurchaseAdapterFactory ] = None, purchase_mode: str = "dry_run", purchase_reconcile_factory: Optional[PurchaseReconcileFactory] = None, cancelled: Callable[[], bool] = lambda: False, device_connection_checker: Optional[Callable[[str], None]] = None, task_saved: Callable[[str], None] = lambda _task_id: None, performance_trace_factory: Callable[[], TaskPerformanceTrace] = ( TaskPerformanceTrace ), ) -> None: self._gateway = gateway self._repository = repository self._client = client self._device_address = str(device_address or "").strip() self._collect_factory = collect_service_factory self._purchase_factory = purchase_adapter_factory self._live_purchase_factory = live_purchase_adapter_factory if purchase_mode not in {"dry_run", "live"}: raise ValueError("purchase_mode 只允许 dry_run 或 live") self._purchase_mode = purchase_mode self._reconcile_factory = purchase_reconcile_factory self._cancelled = cancelled self._task_saved = task_saved self._performance_trace_factory = performance_trace_factory self._device_connection_checker = ( device_connection_checker or AndroidDeviceService().require_connected ) @property def purchase_ready(self) -> bool: """只有显式提供采购演练适配器时才声明采购能力。""" return self._purchase_factory is not None and bool( self._device_address ) def claim_capabilities(self) -> ClaimCapabilities: """集中构造不会意外开放 live 的领取能力。""" supported = [TaskType.COLLECT] if self.purchase_ready: supported.append(TaskType.PURCHASE) device = ( AndroidDeviceInfo(self._device_address) if self._device_address else None ) return ClaimCapabilities( device=device, supported_types=tuple(supported), purchase_mode=( "live" if self.purchase_ready and self._purchase_mode == "live" and self._live_purchase_factory is not None else "dry_run" ), schema_versions=(1,), ) def execute_one(self) -> TaskDispatchOutcome: """不可逆采购优先核单,其余按 Outbox、本地任务、Admin 顺序处理。""" trace = self._performance_trace_factory() with trace.activate(): return self._execute_one_traced(trace) def _execute_one_traced( self, trace: TaskPerformanceTrace ) -> TaskDispatchOutcome: """在同一任务计时上下文中串行处理一轮任务。""" unresolved_reader = getattr( self._repository, "unresolved_irreversible_purchase", None ) unresolved = ( unresolved_reader() if unresolved_reader is not None else None ) if unresolved is not None: raise RuntimeError( f"任务 {unresolved.remote_task_id} 已进入不可逆阶段但尚未转入核单;" "已停止所有采购,绝不重新下单" ) reconcile_task = self._repository.next_purchase_reconcile_task() if reconcile_task is not None: if not self._device_address: raise AndroidDeviceSearchError( "存在只允许核对的采购任务;请先连接并保存原 Android 设备" ) trace.bind_task(reconcile_task.remote_task_id) with trace.stage("adb_device_check"): self._device_connection_checker(self._device_address) if self._reconcile_factory is None: raise RuntimeError( f"任务 {reconcile_task.remote_task_id} 只允许核对订单," "但只读核对执行器未就绪;不会重新下单" ) outcome = PurchaseReconcileService( self._repository, self._device_address, self._reconcile_factory, cancelled=self._cancelled, ).execute_selected(reconcile_task.remote_task_id) return TaskDispatchOutcome( outcome.kind, outcome.message, outcome.task_id ) pending = self._repository.next_pending_outbox() if pending is not None: return self._submit_pending(pending) if not self._device_address: raise AndroidDeviceSearchError( "请先在设置页选择并保存 Android 设备" ) with trace.stage("adb_device_check"): self._device_connection_checker(self._device_address) if not self.purchase_ready: pending_purchase = self._repository.next_purchase_task() if pending_purchase is not None: raise RuntimeError( f"本地采购任务 {pending_purchase.remote_task_id} 等待执行," "但采购演练执行器未就绪;已停止领取新任务" ) task = self._repository.next_runnable_task( include_purchase=self.purchase_ready ) if task is None: with trace.stage("admin_claim_request"): remote = self._gateway.claim_next( self._client, self.claim_capabilities() ) if remote is None: return TaskDispatchOutcome("no_task", "暂无可领取的任务") trace.bind_task(remote.task_id) saved_now = False try: with trace.stage("sqlite_local_save"): self._repository.add_claimed_task( admin_task_to_new_claimed_task(remote) ) saved_now = True except DuplicateTaskError: pass except Exception as exc: raise RuntimeError( f"任务 {remote.task_id} 已领取,但本地保存失败:{exc}" ) from exc task = self._repository.get_task(remote.task_id) if task is None: raise RuntimeError( f"任务 {remote.task_id} 已领取,但未能保存到本地" ) if saved_now: self._task_saved(remote.task_id) trace.bind_task(task.remote_task_id) if task.task_type is TaskType.COLLECT: service = CollectTaskService( self._gateway, self._repository, self._client, self._device_address, cancelled=self._cancelled, collect_service_factory=self._collect_factory, ) outcome = service.execute_selected(task.remote_task_id) else: factory = self._purchase_factory if task.execution_mode == "live": if self.claim_capabilities().purchase_mode != "live": raise RuntimeError( f"真实采购任务 {task.remote_task_id} 的设备或真实采购执行器未就绪;" "任务保持待执行,不会降级为演练" ) factory = self._live_purchase_factory if factory is None: raise RuntimeError( f"采购任务 {task.remote_task_id} 的 {task.execution_mode} " "执行器未就绪,已停止且不会改变任务模式" ) purchase_service = PurchaseTaskService( self._gateway, self._repository, self._client, self._device_address, factory, cancelled=self._cancelled, ) outcome = purchase_service.execute_selected(task.remote_task_id) return TaskDispatchOutcome( outcome.kind, outcome.message, outcome.task_id ) def _submit_pending( self, event: OutboxEventRecord ) -> TaskDispatchOutcome: """只补交已落库事件,不重新访问 PDD 页面。""" task_id = self._repository.outbox_task_id(event.id) self._repository.mark_outbox_sending(event.id) try: if event.event_type is OutboxEventType.TASK_FAILURE: receipt = self._gateway.submit_failure( task_id, event.idempotency_key, event.payload_json ) else: receipt = self._gateway.submit_result( task_id, event.idempotency_key, event.payload_json ) if not receipt.accepted: raise AdminGatewayError( "ADMIN_RESULT_NOT_ACCEPTED", "Admin 未确认接收任务结果", False, ) except AdminGatewayError as exc: if exc.retryable: self._repository.mark_outbox_retry(event.id, str(exc)) return TaskDispatchOutcome( "result_pending", f"任务 {task_id} 结果已保存在本地,等待提交 Admin:{exc}", task_id, ) self._repository.mark_outbox_failed(event.id, str(exc)) return TaskDispatchOutcome( "manual_review", f"任务 {task_id} 结果被 Admin 拒绝:{exc}", task_id, ) self._repository.mark_outbox_sent(event.id) if event.event_type is OutboxEventType.TASK_FAILURE: return TaskDispatchOutcome( "failed", f"任务 {task_id} 失败信息已提交 Admin", task_id ) return TaskDispatchOutcome( "succeeded", f"任务 {task_id} 结果已提交 Admin", task_id )