Files
cmautobuy/client/src/admin_gateway.py
T

207 lines
6.1 KiB
Python
Raw Normal View History

"""Client 访问 Admin 的稳定边界和简单数据对象。
AdminGateway 提供登记、领取、结果提交和一次性采购规格解析命令。
业务层不应直接依赖 HTTP 请求或 Mock 的内部实现。
"""
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import Any, Mapping, Optional, Tuple
from .task_models import TaskType
@dataclass(frozen=True)
class ClientInfo:
2026-08-06 17:57:49 +08:00
"""Client 身份和可选名称,不保存访问令牌。"""
client_id: str
2026-08-06 17:57:49 +08:00
name: str = ""
def __post_init__(self) -> None:
if not self.client_id.strip():
raise ValueError("client_id 不能为空")
2026-08-06 17:57:49 +08:00
if len(self.name.strip()) > 50:
raise ValueError("client.name 最多 50 个字")
@dataclass(frozen=True)
class AndroidDeviceInfo:
"""领取任务时上报的 Android 设备信息。"""
address: str
platform: str = "android"
pdd_package: str = "com.xunmeng.pinduoduo"
def __post_init__(self) -> None:
if not self.address.strip():
raise ValueError("设备地址不能为空")
if self.platform != "android":
raise ValueError("当前只支持 android 平台")
if not self.pdd_package.strip():
raise ValueError("PDD 包名不能为空")
@dataclass(frozen=True)
class ClaimCapabilities:
"""Client 领取任务时声明的设备与执行能力。"""
2026-08-06 17:57:49 +08:00
device: Optional[AndroidDeviceInfo] = None
supported_types: Tuple[TaskType, ...] = (TaskType.COLLECT,)
purchase_mode: str = "dry_run"
schema_versions: Tuple[int, ...] = (1,)
def __post_init__(self) -> None:
if not self.supported_types:
raise ValueError("supported_types 不能为空")
if any(not isinstance(value, TaskType) for value in self.supported_types):
raise ValueError("supported_types 必须使用 TaskType")
if self.purchase_mode not in {"dry_run", "live"}:
raise ValueError("purchase_mode 只允许 dry_run 或 live")
if not self.schema_versions or any(
version <= 0 for version in self.schema_versions
):
raise ValueError("schema_versions 必须是正整数")
@dataclass(frozen=True)
class AdminTask:
"""Admin 派发给 Client 的一个任务。"""
task_id: str
task_type: TaskType
version: int
priority: int
payload: Mapping[str, Any] = field(default_factory=dict)
created_at: str = ""
updated_at: str = ""
execution_mode: str = "dry_run"
def __post_init__(self) -> None:
if not self.task_id.strip():
raise ValueError("task_id 不能为空")
if not isinstance(self.task_type, TaskType):
raise ValueError("task_type 必须使用 TaskType")
if self.version <= 0:
raise ValueError("version 必须大于 0")
if self.execution_mode not in {"dry_run", "live"}:
raise ValueError("execution_mode 只允许 dry_run 或 live")
if (
self.task_type is not TaskType.PURCHASE
and self.execution_mode != "dry_run"
):
raise ValueError("只有采购任务允许 execution_mode=live")
if not isinstance(self.payload, Mapping):
raise ValueError("payload 必须是对象")
@dataclass(frozen=True)
class SubmissionReceipt:
"""Admin 已接收并保存一次提交的确认。"""
accepted: bool
result_id: str
accepted_at: str
2026-08-06 17:57:49 +08:00
@dataclass(frozen=True)
class RegistrationReceipt:
"""Admin 已登记当前 Client 的确认。"""
registered: bool
client_id: str
registered_at: str
@dataclass(frozen=True)
class SpecResolutionMatch:
"""Admin 从本次候选白名单中选中的原始规格。"""
candidate_id: str
raw_text: str
options: Mapping[str, str]
@dataclass(frozen=True)
class SpecResolutionReceipt:
"""一次采购规格解析的完整业务响应。"""
schema_version: int
resolution_id: str
outcome: str
source: Optional[str]
candidate_snapshot_hash: str
match: Optional[SpecResolutionMatch]
confidence_bps: Optional[int]
reason: str
resolved_at: str
class AdminGatewayError(RuntimeError):
"""带稳定错误代码和可重试标志的 Admin 边界错误。"""
def __init__(
self,
code: str,
message: str,
retryable: bool,
request_id: str = "",
details: Optional[Mapping[str, Any]] = None,
):
super().__init__(message)
self.code = code
self.retryable = retryable
self.request_id = request_id
self.details = dict(details or {})
2026-08-06 17:57:49 +08:00
class ClientRegistrationGateway(ABC):
"""设置页只依赖登记能力,不依赖任务领取和提交。"""
@abstractmethod
def register_client(
self, client: ClientInfo, capabilities: ClaimCapabilities
) -> RegistrationReceipt:
"""幂等登记或更新当前 Client。"""
class TaskClaimGateway(ABC):
"""任务页只依赖领取能力,不依赖登记和结果提交。"""
@abstractmethod
def claim_next(
self, client: ClientInfo, capabilities: ClaimCapabilities
) -> Optional[AdminTask]:
"""领取至多一个任务;暂无任务时返回 ``None``。"""
class AdminGateway(ClientRegistrationGateway, TaskClaimGateway):
"""完整 Admin 边界;不得增加状态查询、心跳或租约方法。"""
@abstractmethod
def submit_result(
self,
task_id: str,
idempotency_key: str,
result: Mapping[str, Any],
) -> SubmissionReceipt:
"""幂等提交采集或采购成功结果。"""
@abstractmethod
def submit_failure(
self,
task_id: str,
idempotency_key: str,
failure: Mapping[str, Any],
) -> SubmissionReceipt:
"""幂等提交失败、取消或人工处理结果。"""
@abstractmethod
def resolve_purchase_spec(
self,
task_id: str,
idempotency_key: str,
observation: Mapping[str, Any],
) -> SpecResolutionReceipt:
"""一次性解析当前真机尺码候选;不得用于查询任务状态。"""