Files
cmautobuy/client/src/collect_task_service.py
T

284 lines
11 KiB
Python

"""一条采集任务的应用层流程:领取、采集、落库、提交。"""
from __future__ import annotations
from dataclasses import dataclass
from typing import Callable, Optional
from .admin_gateway import (
AdminGateway,
AdminGatewayError,
AndroidDeviceInfo,
ClaimCapabilities,
ClientInfo,
)
from .pdd_collect_service import PddCollectError, PddCollectService
from .pdd_device_service import (
PddDeviceService,
current_thread_device_service,
)
from .db import data_dir
from .task_models import OutboxEventRecord, OutboxEventType, TaskStatus, TaskType
from .task_models import NewClaimedTask
from .task_repository import DuplicateTaskError, TaskRepository
@dataclass(frozen=True)
class CollectTaskOutcome:
"""工作线程返回给界面的简短结果。"""
kind: str
message: str
task_id: str = ""
CollectServiceFactory = Callable[
[str, str, Callable[[], bool]], PddCollectService
]
class CollectTaskService:
"""一次调用只处理一条本地工作或一个待提交事件。"""
def __init__(
self,
gateway: AdminGateway,
repository: TaskRepository,
client: ClientInfo,
device_address: str,
*,
cancelled: Callable[[], bool] = lambda: False,
started: Callable[[str], None] = lambda _task_id: None,
collect_service_factory: Optional[CollectServiceFactory] = None,
) -> None:
self._gateway = gateway
self._repository = repository
self._client = client
self._device_address = device_address
self._cancelled = cancelled
self._started = started
self._factory = collect_service_factory or self._default_factory
def execute_one(self) -> CollectTaskOutcome:
"""先补交已有结果,再执行最早的本地任务,最后才领取新任务。"""
pending = self._repository.next_pending_outbox()
if pending is not None:
return self._submit(pending)
if not self._device_address.strip():
raise ValueError("请先在设置页选择并保存 Android 设备")
task = self._repository.next_collect_task()
if task is None:
remote = self._gateway.claim_next(
self._client,
ClaimCapabilities(
device=AndroidDeviceInfo(self._device_address),
supported_types=(TaskType.COLLECT,),
purchase_mode="dry_run",
schema_versions=(1,),
),
)
if remote is None:
return CollectTaskOutcome("no_task", "暂无可领取的采集任务")
try:
self._repository.add_claimed_task(
NewClaimedTask(
remote_task_id=remote.task_id,
task_type=remote.task_type,
goods_url=str(remote.payload.get("goods_url") or ""),
goods_id=(
str(remote.payload["goods_id"])
if remote.payload.get("goods_id") is not None
else None
),
priority=remote.priority,
version=remote.version,
admin_payload={
"id": remote.task_id,
"type": remote.task_type.value,
"version": remote.version,
"priority": remote.priority,
"payload": dict(remote.payload),
"created_at": remote.created_at,
"updated_at": remote.updated_at,
},
)
)
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} 未能保存到本地")
return self._execute_task(task.remote_task_id)
def execute_selected(self, remote_task_id: str) -> CollectTaskOutcome:
"""只执行指定的本地采集任务,不领取任务或补交其他结果。"""
if not self._device_address.strip():
raise ValueError("请先在设置页选择并保存 Android 设备")
task = self._repository.get_task(remote_task_id)
if task is None:
raise ValueError(f"任务 {remote_task_id} 不存在")
return self._execute_task(task.remote_task_id)
def _execute_task(self, remote_task_id: str) -> CollectTaskOutcome:
"""执行一条已经保存在本地且处于待执行状态的采集任务。"""
started = self._repository.start_collect_run(
remote_task_id, self._device_address
)
self._started(remote_task_id)
if self._cancelled():
event = self._repository.save_collect_failure(
remote_task_id,
started.attempt_id,
TaskStatus.CANCELLED,
"PDD_CANCELLED",
"采集任务已安全取消",
False,
{},
)
return self._submit_cancelled(event, remote_task_id)
collector = self._factory(
self._device_address, self._client.client_id, self._cancelled
)
business_failure_message = ""
try:
result = collector.collect(started.task)
event = self._repository.save_collect_result(
remote_task_id, started.attempt_id, result.to_pdd_data()
)
except PddCollectError as exc:
status, retryable = self._classify_error(exc.code)
report_code = self._report_error_code(exc.code)
event = self._repository.save_collect_failure(
remote_task_id,
started.attempt_id,
status,
report_code,
exc.message,
retryable,
exc.diagnostics,
)
if exc.code == "PDD_CANCELLED":
return self._submit_cancelled(event, remote_task_id)
if exc.code == "PDD_GOODS_UNAVAILABLE":
business_failure_message = exc.message
outcome = self._submit(event)
if business_failure_message and outcome.kind == "failed":
return CollectTaskOutcome(
"business_failed", business_failure_message, remote_task_id
)
return outcome
def _submit_cancelled(
self, event: OutboxEventRecord, remote_task_id: str
) -> CollectTaskOutcome:
"""提交取消结果;Admin 接收后对界面返回明确的已取消。"""
outcome = self._submit(event)
if outcome.kind == "failed":
return CollectTaskOutcome(
"cancelled", f"任务 {remote_task_id} 已安全取消", remote_task_id
)
return outcome
def _submit(self, event: OutboxEventRecord) -> CollectTaskOutcome:
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:
message = str(exc)
if exc.retryable:
self._repository.mark_outbox_retry(event.id, message)
return CollectTaskOutcome(
"result_pending",
f"任务 {task_id} 数据已保存在本地,等待重新提交 Admin:{message}",
task_id,
)
self._repository.mark_outbox_failed(event.id, message)
return CollectTaskOutcome(
"manual_review",
f"任务 {task_id} 提交被 Admin 拒绝:{message}",
task_id,
)
self._repository.mark_outbox_sent(event.id)
if event.event_type is OutboxEventType.TASK_FAILURE:
return CollectTaskOutcome(
"failed",
f"任务 {task_id} 采集未完成,失败信息已提交 Admin",
task_id,
)
return CollectTaskOutcome(
"succeeded", f"任务 {task_id} 采集完成并已提交 Admin", task_id
)
@staticmethod
def _classify_error(code: str) -> tuple[TaskStatus, bool]:
if code == "PDD_CANCELLED":
return TaskStatus.CANCELLED, False
if code == "PDD_GOODS_UNAVAILABLE":
return TaskStatus.FAILED, False
if code in {
"PDD_PAGE_LOGIN_REQUIRED",
"PDD_PAGE_CAPTCHA",
"PDD_DATA_SPEC_INCOMPLETE",
"PDD_DATA_TITLE_MISSING",
"PDD_DATA_PRICE_MISSING",
"PDD_DATA_SKU_NAME_TRUNCATED",
}:
return TaskStatus.MANUAL_REVIEW, False
if code.startswith("PDD_DATA_GOODS_"):
return TaskStatus.FAILED, False
return TaskStatus.RETRY_WAIT, True
@staticmethod
def _report_error_code(code: str) -> str:
"""把内部解析代码映射为 Admin 可枚举的稳定代码。"""
return {
"PDD_DATA_SKU_NAME_TRUNCATED": "SKU_NAME_TRUNCATED",
"PDD_PAGE_SPEC_ENTRY_MISSING": "SKU_PANEL_NOT_FOUND",
"PDD_PAGE_SPEC_PANEL_TIMEOUT": "SKU_PANEL_NOT_FOUND",
"PDD_DATA_TITLE_MISSING": "TITLE_TOO_SHORT",
"PDD_PAGE_OVERALL_TIMEOUT": "COLLECT_TIMEOUT",
"DEVICE_OFFLINE": "DEVICE_DISCONNECTED",
"DEVICE_CONNECT_FAILED": "DEVICE_DISCONNECTED",
"DEVICE_DISCONNECTED": "DEVICE_DISCONNECTED",
}.get(code, code)
@staticmethod
def _default_factory(
device_address: str,
client_id: str,
cancelled: Callable[[], bool],
) -> PddCollectService:
return PddCollectService(
current_thread_device_service() or PddDeviceService(),
device_address,
client_id,
cancelled=cancelled,
overall_timeout=600.0,
artifact_directory=data_dir() / "artifacts",
)