379 lines
14 KiB
Python
379 lines
14 KiB
Python
"""串行领取、落库并按任务类型分派一条任务。"""
|
|
|
|
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
|
|
)
|