From 311d55f0da8caa88728add1f7e7955d9bd984143 Mon Sep 17 00:00:00 2001 From: chengma Date: Sun, 9 Aug 2026 23:39:37 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=AE=9E=E7=8E=B0=E9=87=87=E8=B4=AD?= =?UTF-8?q?=E4=BB=BB=E5=8A=A1=E5=AE=89=E5=85=A8=E6=BC=94=E7=BB=83=20(#71)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- client/src/pdd_purchase_adapter.py | 75 ++++ client/src/purchase_task_service.py | 466 ++++++++++++++++++++++ client/src/task_models.py | 2 +- client/src/task_repository.py | 319 ++++++++++++++- client/test/test_purchase_task_service.py | 278 +++++++++++++ docs/client/02-architecture.md | 7 +- docs/client/03-data-model.md | 22 +- 7 files changed, 1162 insertions(+), 7 deletions(-) create mode 100644 client/src/pdd_purchase_adapter.py create mode 100644 client/src/purchase_task_service.py create mode 100644 client/test/test_purchase_task_service.py diff --git a/client/src/pdd_purchase_adapter.py b/client/src/pdd_purchase_adapter.py new file mode 100644 index 0000000..2fd4e44 --- /dev/null +++ b/client/src/pdd_purchase_adapter.py @@ -0,0 +1,75 @@ +"""PDD 采购演练适配层的稳定边界。 + +应用服务只依赖本文件中的小接口,不直接认识 uiautomator2。接口故意不提供 +“提交订单”或“付款”方法,因此演练代码没有可误调用的真实下单入口。 +""" + +from __future__ import annotations + +from abc import ABC, abstractmethod +from dataclasses import dataclass, field +from typing import Mapping + + +@dataclass(frozen=True) +class PurchasePageState: + """一次重新读取页面后得到的采购相关事实。""" + + page_kind: str + goods_id: str + selected_options: Mapping[str, str] = field(default_factory=dict) + quantity: int = 0 + price_cent: int = 0 + candidate_count: int = 1 + + +class PddPurchaseError(RuntimeError): + """适配层失败,包含稳定代码、步骤和是否可安全重试。""" + + def __init__( + self, + code: str, + message: str, + *, + step: str, + retryable: bool = False, + diagnostics: Mapping[str, object] | None = None, + ) -> None: + super().__init__(message) + self.code = str(code or "PURCHASE_ADAPTER_ERROR") + self.message = str(message or "PDD 采购演练失败") + self.step = str(step or "purchase_prepare") + self.retryable = bool(retryable) + self.diagnostics = dict(diagnostics or {}) + + +class PddPurchaseAdapter(ABC): + """一台设备的一次采购演练会话;只能由创建它的工作线程使用。""" + + @abstractmethod + def open_goods(self, goods_url: str) -> None: + """通过商品链接打开 PDD 商品页。""" + + @abstractmethod + def read_state(self) -> PurchasePageState: + """重新读取当前页面,不返回缓存状态。""" + + @abstractmethod + def select_options(self, options: Mapping[str, str]) -> None: + """按完整动态规格对象精确选择,不做相似匹配。""" + + @abstractmethod + def set_quantity(self, quantity: int) -> None: + """设置采购数量。""" + + @abstractmethod + def enter_confirmation(self) -> None: + """进入最终提交前确认页,但不得提交订单。""" + + @abstractmethod + def stop_before_submit(self) -> None: + """在提交订单按钮之前停止并保持页面可供人工核对。""" + + @abstractmethod + def close(self) -> None: + """释放设备会话;不得在此方法中产生页面点击。""" diff --git a/client/src/purchase_task_service.py b/client/src/purchase_task_service.py new file mode 100644 index 0000000..d69598f --- /dev/null +++ b/client/src/purchase_task_service.py @@ -0,0 +1,466 @@ +"""一条本地采购任务的安全演练流程。 + +本模块固定为 ``dry_run``。它会核对商品、动态规格、数量和价格,并停在最终 +提交订单之前;代码中没有提交订单或付款入口。 +""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import datetime, timezone +from typing import Callable, Mapping + +from .admin_gateway import AdminGateway, AdminGatewayError, ClientInfo +from .pdd_purchase_adapter import ( + PddPurchaseAdapter, + PddPurchaseError, + PurchasePageState, +) +from .task_models import OutboxEventRecord, OutboxEventType, TaskStatus +from .task_repository import TaskRepository + + +def utc_now_iso() -> str: + """返回精确到秒的 UTC ISO 8601 时间。""" + + return datetime.now(timezone.utc).isoformat(timespec="seconds").replace( + "+00:00", "Z" + ) + + +@dataclass(frozen=True) +class PurchaseTarget: + """从已领取任务中校验出的采购目标。""" + + goods_id: str + goods_url: str + options: Mapping[str, str] + quantity: int + max_price_cent: int + + +@dataclass(frozen=True) +class PurchaseTaskOutcome: + """采购工作线程返回给协调层的简短结果。""" + + kind: str + message: str + task_id: str = "" + + +PurchaseAdapterFactory = Callable[ + [str, Callable[[], bool]], PddPurchaseAdapter +] + + +class PurchaseTaskService: + """只执行本地已保存的采购任务,不领取新任务。""" + + def __init__( + self, + gateway: AdminGateway, + repository: TaskRepository, + client: ClientInfo, + device_address: str, + adapter_factory: PurchaseAdapterFactory, + *, + cancelled: Callable[[], bool] = lambda: False, + ) -> None: + self._gateway = gateway + self._repository = repository + self._client = client + self._device_address = str(device_address or "").strip() + self._factory = adapter_factory + self._cancelled = cancelled + + def execute_one_local(self) -> PurchaseTaskOutcome: + """先补交结果,再执行最早的本地待执行采购任务。""" + + pending = self._repository.next_pending_outbox() + if pending is not None: + return self._submit(pending) + if not self._device_address: + raise ValueError("请先在设置页选择并保存 Android 设备") + task = self._repository.next_purchase_task() + if task is None: + return PurchaseTaskOutcome("no_task", "暂无本地待执行的采购任务") + return self.execute_selected(task.remote_task_id) + + def execute_selected(self, remote_task_id: str) -> PurchaseTaskOutcome: + """对指定的本地采购任务执行一次演练。""" + + if not self._device_address: + raise ValueError("请先在设置页选择并保存 Android 设备") + if self._cancelled(): + return PurchaseTaskOutcome( + "cancelled", "采购演练已在开始前取消", remote_task_id + ) + + started = self._repository.start_purchase_run( + remote_task_id, self._device_address + ) + adapter: PddPurchaseAdapter | None = None + step = "purchase_prepare" + try: + target = self._target_from_task(started.task) + adapter = self._factory(self._device_address, self._cancelled) + + step = "purchase_open_goods" + self._enter_step(remote_task_id, started.attempt_id, step) + adapter.open_goods(target.goods_url) + + step = "purchase_verify_goods" + self._enter_step(remote_task_id, started.attempt_id, step) + state = adapter.read_state() + self._validate_identity(state, target, "goods") + + step = "purchase_select_options" + self._enter_step(remote_task_id, started.attempt_id, step) + self._validate_identity(adapter.read_state(), target, "goods") + adapter.select_options(target.options) + + step = "purchase_verify_options" + self._enter_step(remote_task_id, started.attempt_id, step) + state = adapter.read_state() + self._validate_selected_options(state, target) + + step = "purchase_set_quantity" + self._enter_step(remote_task_id, started.attempt_id, step) + self._validate_selected_options(adapter.read_state(), target) + adapter.set_quantity(target.quantity) + + step = "purchase_verify_quantity_price" + self._enter_step(remote_task_id, started.attempt_id, step) + state = adapter.read_state() + self._validate_checkout_values(state, target) + + step = "purchase_enter_confirmation" + self._enter_step(remote_task_id, started.attempt_id, step) + self._validate_checkout_values(adapter.read_state(), target) + adapter.enter_confirmation() + + step = "purchase_verify_confirmation" + self._enter_step(remote_task_id, started.attempt_id, step) + state = adapter.read_state() + if state.page_kind != "order_confirmation": + raise PddPurchaseError( + "PURCHASE_CONFIRMATION_NOT_REACHED", + "没有到达最终提交前确认页,已停止演练", + step=step, + ) + self._validate_checkout_values(state, target) + + step = "purchase_dry_run_stopped" + self._enter_step(remote_task_id, started.attempt_id, step) + adapter.stop_before_submit() + result = self._result_data(target, state) + event = self._repository.save_purchase_result( + remote_task_id, started.attempt_id, result + ) + except PddPurchaseError as exc: + event = self._save_failure( + remote_task_id, + started.attempt_id, + exc.code, + exc.message, + exc.retryable, + exc.step or step, + exc.diagnostics, + ) + except Exception as exc: + event = self._save_failure( + remote_task_id, + started.attempt_id, + "PURCHASE_UNEXPECTED", + f"采购演练在“{step}”发生未知错误:{exc}", + False, + step, + {}, + ) + finally: + if adapter is not None: + adapter.close() + return self._submit(event) + + def _enter_step( + self, remote_task_id: str, attempt_id: str, step: str + ) -> None: + """先持久化步骤,再检查停止请求。""" + + self._repository.update_purchase_step( + remote_task_id, attempt_id, step + ) + if self._cancelled(): + raise PddPurchaseError( + "PURCHASE_CANCELLED", + "用户已请求停止采购演练", + step=step, + ) + + @staticmethod + def _target_from_task(task: object) -> PurchaseTarget: + payload_root = getattr(task, "admin_payload", None) + if not isinstance(payload_root, Mapping): + raise PddPurchaseError( + "PURCHASE_TASK_INVALID", + "采购任务原始数据不是对象", + step="purchase_prepare", + ) + payload = payload_root.get("payload") + if not isinstance(payload, Mapping): + raise PddPurchaseError( + "PURCHASE_TASK_INVALID", + "采购任务缺少 payload 对象", + step="purchase_prepare", + ) + goods_url = str(payload.get("goods_url") or "").strip() + goods_id = str(payload.get("goods_id") or "").strip() + options = payload.get("options") + quantity = payload.get("quantity") + max_price_cent = payload.get("max_price_cent") + if not goods_url or not goods_id: + raise PddPurchaseError( + "PURCHASE_TASK_INVALID", + "采购任务缺少商品链接或商品编号", + step="purchase_prepare", + ) + if not isinstance(options, Mapping) or not options: + raise PddPurchaseError( + "PURCHASE_TASK_INVALID", + "采购任务缺少动态规格 options", + step="purchase_prepare", + ) + normalized_options: dict[str, str] = {} + for key, value in options.items(): + checked_key = str(key or "").strip() + checked_value = str(value or "").strip() + if not checked_key or not checked_value: + raise PddPurchaseError( + "PURCHASE_TASK_INVALID", + "采购任务 options 的名称和值都不能为空", + step="purchase_prepare", + ) + normalized_options[checked_key] = checked_value + if isinstance(quantity, bool) or not isinstance(quantity, int) or quantity <= 0: + raise PddPurchaseError( + "PURCHASE_TASK_INVALID", + "采购数量必须是大于 0 的整数", + step="purchase_prepare", + ) + if ( + isinstance(max_price_cent, bool) + or not isinstance(max_price_cent, int) + or max_price_cent <= 0 + ): + raise PddPurchaseError( + "PURCHASE_TASK_INVALID", + "人民币价格上限必须是大于 0 的整数分", + step="purchase_prepare", + ) + return PurchaseTarget( + goods_id, + goods_url, + normalized_options, + quantity, + max_price_cent, + ) + + @staticmethod + def _validate_identity( + state: PurchasePageState, target: PurchaseTarget, expected_page: str + ) -> None: + PurchaseTaskService._validate_common_state(state) + if state.page_kind != expected_page: + raise PddPurchaseError( + "PDD_PAGE_UNKNOWN", + f"当前页面不是预期的 {expected_page} 页面,已停止演练", + step="purchase_verify_goods", + ) + if state.goods_id != target.goods_id: + raise PddPurchaseError( + "PURCHASE_GOODS_MISMATCH", + "当前 PDD 商品与采购任务不一致,已停止演练", + step="purchase_verify_goods", + ) + + @staticmethod + def _validate_selected_options( + state: PurchasePageState, target: PurchaseTarget + ) -> None: + PurchaseTaskService._validate_common_state(state) + if state.goods_id != target.goods_id: + raise PddPurchaseError( + "PURCHASE_GOODS_MISMATCH", + "选择规格后商品编号发生变化,已停止演练", + step="purchase_verify_options", + ) + if dict(state.selected_options) != dict(target.options): + raise PddPurchaseError( + "PURCHASE_OPTIONS_MISMATCH", + "当前选中规格与采购任务不完全一致,已停止演练", + step="purchase_verify_options", + ) + + @staticmethod + def _validate_checkout_values( + state: PurchasePageState, target: PurchaseTarget + ) -> None: + PurchaseTaskService._validate_selected_options(state, target) + if state.quantity != target.quantity: + raise PddPurchaseError( + "PURCHASE_QUANTITY_MISMATCH", + "当前数量与采购任务不一致,已停止演练", + step="purchase_verify_quantity_price", + ) + if state.price_cent <= 0: + raise PddPurchaseError( + "PURCHASE_PRICE_MISSING", + "无法读取稳定的当前人民币价格,已停止演练", + step="purchase_verify_quantity_price", + ) + if state.price_cent > target.max_price_cent: + raise PddPurchaseError( + "PURCHASE_PRICE_EXCEEDED", + ( + f"当前价格 {state.price_cent} 分超过上限 " + f"{target.max_price_cent} 分,已停止演练" + ), + step="purchase_verify_quantity_price", + ) + + @staticmethod + def _validate_common_state(state: PurchasePageState) -> None: + if state.page_kind == "captcha": + raise PddPurchaseError( + "PDD_PAGE_CAPTCHA", + "PDD 出现安全验证,请人工处理", + step="purchase_page_check", + ) + if state.page_kind == "login_required": + raise PddPurchaseError( + "PDD_PAGE_LOGIN_REQUIRED", + "PDD 登录状态失效,请人工登录后再处理", + step="purchase_page_check", + ) + if state.page_kind in {"unknown", "risk_control", "payment"}: + raise PddPurchaseError( + "PDD_PAGE_UNKNOWN", + "PDD 出现未知、风控或支付页面,已停止演练", + step="purchase_page_check", + ) + if state.candidate_count != 1: + raise PddPurchaseError( + "PURCHASE_AMBIGUOUS_TARGET", + "页面存在多个候选目标,无法安全确认,已停止演练", + step="purchase_page_check", + ) + + def _save_failure( + self, + remote_task_id: str, + attempt_id: str, + code: str, + message: str, + retryable: bool, + step: str, + diagnostics: Mapping[str, object], + ) -> OutboxEventRecord: + if code == "PURCHASE_CANCELLED": + status = TaskStatus.CANCELLED + elif retryable: + status = TaskStatus.RETRY_WAIT + else: + status = TaskStatus.MANUAL_REVIEW + return self._repository.save_purchase_failure( + remote_task_id, + attempt_id, + status, + code, + message, + retryable, + step, + dict(diagnostics), + ) + + def _submit(self, event: OutboxEventRecord) -> PurchaseTaskOutcome: + 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 PurchaseTaskOutcome( + "result_pending", + f"任务 {task_id} 演练结果已保存在本地,等待提交 Admin:{message}", + task_id, + ) + self._repository.mark_outbox_failed(event.id, message) + return PurchaseTaskOutcome( + "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 PurchaseTaskOutcome( + "failed", + f"任务 {task_id} 采购演练未完成,失败信息已提交 Admin", + task_id, + ) + return PurchaseTaskOutcome( + "succeeded", + f"任务 {task_id} 演练完成并已提交 Admin;没有提交订单", + task_id, + ) + + def _result_data( + self, target: PurchaseTarget, state: PurchasePageState + ) -> dict[str, object]: + return { + "schema_version": 1, + "goods_id": target.goods_id, + "goods_url": target.goods_url, + "purchase": { + "mode": "dry_run", + "requested": { + "options": dict(target.options), + "quantity": target.quantity, + "max_price_cent": target.max_price_cent, + }, + "confirmed": { + "options": dict(state.selected_options), + "quantity": state.quantity, + "unit_price_cent": state.price_cent, + "total_price_cent": state.price_cent * state.quantity, + }, + "confirmation_reached": True, + "order_submitted": False, + "payment_attempted": False, + "order_no": None, + "ordered_at": None, + "ordered_at_raw": None, + "match_status": "not_submitted", + }, + "captured_at": utc_now_iso(), + "source": { + "client_id": self._client.client_id, + "device_address": self._device_address, + "mode": "dry_run", + }, + } diff --git a/client/src/task_models.py b/client/src/task_models.py index f768ed8..f0ec5cb 100644 --- a/client/src/task_models.py +++ b/client/src/task_models.py @@ -187,7 +187,7 @@ class OutboxEventRecord: @dataclass(frozen=True) class StartedTaskRun: - """已进入执行状态的一次采集尝试。""" + """已进入执行状态的一次采集或采购尝试。""" task: TaskDetail attempt_id: str diff --git a/client/src/task_repository.py b/client/src/task_repository.py index e0ae4a7..6bbd0d2 100644 --- a/client/src/task_repository.py +++ b/client/src/task_repository.py @@ -210,6 +210,20 @@ class TaskRepository: connection.close() return self._to_detail(row) if row is not None else None + def next_purchase_task(self) -> Optional[TaskDetail]: + """返回最早的一条本地待执行采购任务。""" + + connection = open_database(self._db_path) + try: + row = connection.execute( + "SELECT * FROM pdd_tasks" + " WHERE task_type = 'purchase' AND status = 'claimed'" + " ORDER BY received_at ASC, id ASC LIMIT 1" + ).fetchone() + finally: + connection.close() + return self._to_detail(row) if row is not None else None + def validate_collect_rerun(self, remote_task_id: str) -> TaskDetail: """校验任务能否重新采集,成功时返回任务详情。""" @@ -341,6 +355,291 @@ class TaskRepository: finally: connection.close() + def start_purchase_run( + self, remote_task_id: str, device_address: str + ) -> StartedTaskRun: + """原子地开始一次采购演练并创建独立运行记录。""" + + now = utc_now_iso() + attempt_id = str(uuid4()) + connection = open_database(self._db_path) + try: + with connection: + row = connection.execute( + "SELECT * FROM pdd_tasks WHERE remote_task_id = ?", + (remote_task_id,), + ).fetchone() + if row is None: + raise ValueError(f"任务 {remote_task_id} 不存在") + if row["task_type"] != TaskType.PURCHASE.value: + raise ValueError("当前任务不是采购任务") + if row["status"] != TaskStatus.CLAIMED.value: + raise ValueError( + f"采购任务状态 {row['status']} 不能开始演练" + ) + attempt_no = int( + connection.execute( + "SELECT COALESCE(MAX(attempt_no), 0) + 1" + " FROM task_runs WHERE task_id = ?", + (row["id"],), + ).fetchone()[0] + ) + connection.execute( + "UPDATE pdd_tasks SET status = 'running'," + " current_step = 'purchase_prepare'," + " started_at = COALESCE(started_at, ?)," + " last_error_code = NULL, last_error_message = NULL," + " updated_at = ? WHERE id = ?", + (now, now, row["id"]), + ) + connection.execute( + "INSERT INTO task_runs (task_id, attempt_id, attempt_no," + " device_address, run_status, current_step, started_at," + " created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)", + ( + row["id"], + attempt_id, + attempt_no, + device_address, + RunStatus.RUNNING.value, + "purchase_prepare", + now, + now, + now, + ), + ) + task = self.get_task(remote_task_id) + assert task is not None + return StartedTaskRun(task, attempt_id, attempt_no) + finally: + connection.close() + + def update_purchase_step( + self, remote_task_id: str, attempt_id: str, step: str + ) -> None: + """在进入采购关键步骤前,同时持久化任务和本次运行的步骤。""" + + checked_step = str(step or "").strip() + if not checked_step: + raise ValueError("采购步骤不能为空") + now = utc_now_iso() + connection = open_database(self._db_path) + try: + with connection: + task = connection.execute( + "SELECT id, task_type, status FROM pdd_tasks" + " WHERE remote_task_id = ?", + (remote_task_id,), + ).fetchone() + if task is None: + raise ValueError(f"任务 {remote_task_id} 不存在") + if task["task_type"] != TaskType.PURCHASE.value: + raise ValueError("当前任务不是采购任务") + if task["status"] != TaskStatus.RUNNING.value: + raise ValueError("采购任务当前不在执行中") + cursor = connection.execute( + "UPDATE task_runs SET current_step = ?, updated_at = ?" + " WHERE task_id = ? AND attempt_id = ?" + " AND run_status = 'running'", + (checked_step, now, task["id"], attempt_id), + ) + if cursor.rowcount != 1: + raise ValueError("采购执行记录不存在或已经结束") + connection.execute( + "UPDATE pdd_tasks SET current_step = ?, updated_at = ?" + " WHERE id = ?", + (checked_step, now, task["id"]), + ) + finally: + connection.close() + + def save_purchase_result( + self, + remote_task_id: str, + attempt_id: str, + pdd_data: Dict[str, object], + ) -> OutboxEventRecord: + """原子保存采购演练结果,并创建采购结果 Outbox。""" + + now = utc_now_iso() + connection = open_database(self._db_path) + try: + with connection: + task = connection.execute( + "SELECT id, version, task_type FROM pdd_tasks" + " WHERE remote_task_id = ?", + (remote_task_id,), + ).fetchone() + if task is None: + raise ValueError(f"任务 {remote_task_id} 不存在") + if task["task_type"] != TaskType.PURCHASE.value: + raise ValueError("当前任务不是采购任务") + payload = { + "task_version": task["version"], + "attempt_id": attempt_id, + "result_type": "purchase", + "completed_at": now, + "pdd_data": pdd_data, + } + result_json = json.dumps(pdd_data, ensure_ascii=False) + idempotency_key = ( + f"{remote_task_id}:{attempt_id}:result-v1" + ) + connection.execute( + "UPDATE pdd_tasks SET status = 'result_pending'," + " current_step = 'submit_result', pdd_data = ?," + " price_cent = ?, finished_at = ?, updated_at = ?" + " WHERE id = ?", + ( + result_json, + self._purchase_result_price(pdd_data), + now, + now, + task["id"], + ), + ) + cursor = connection.execute( + "UPDATE task_runs SET run_status = 'succeeded'," + " current_step = 'submit_result', result_data = ?," + " finished_at = ?, updated_at = ?" + " WHERE task_id = ? AND attempt_id = ?" + " AND run_status = 'running'", + ( + result_json, + now, + now, + task["id"], + attempt_id, + ), + ) + if cursor.rowcount != 1: + raise ValueError("采购执行记录不存在或已经结束") + event_cursor = connection.execute( + "INSERT INTO outbox_events (task_id, event_type," + " idempotency_key, payload_json, status, created_at, updated_at)" + " VALUES (?, 'purchase_result', ?, ?, 'pending', ?, ?)", + ( + task["id"], + idempotency_key, + json.dumps(payload, ensure_ascii=False), + now, + now, + ), + ) + event_id = int(event_cursor.lastrowid) + event = self.get_outbox_event(event_id) + assert event is not None + return event + finally: + connection.close() + + def save_purchase_failure( + self, + remote_task_id: str, + attempt_id: str, + status: TaskStatus, + error_code: str, + error_message: str, + retryable: bool, + step: str, + diagnostics: Optional[Dict[str, object]] = None, + ) -> OutboxEventRecord: + """原子保存采购演练失败,并创建失败 Outbox。""" + + if status not in { + TaskStatus.RETRY_WAIT, + TaskStatus.MANUAL_REVIEW, + TaskStatus.FAILED, + TaskStatus.CANCELLED, + }: + raise ValueError("失败状态无效") + checked_step = str(step or "purchase_prepare").strip() + now = utc_now_iso() + connection = open_database(self._db_path) + try: + with connection: + task = connection.execute( + "SELECT id, version, task_type FROM pdd_tasks" + " WHERE remote_task_id = ?", + (remote_task_id,), + ).fetchone() + if task is None: + raise ValueError(f"任务 {remote_task_id} 不存在") + if task["task_type"] != TaskType.PURCHASE.value: + raise ValueError("当前任务不是采购任务") + payload = { + "task_version": task["version"], + "attempt_id": attempt_id, + "status": status.value, + "error": { + "code": error_code, + "message": error_message, + "retryable": retryable, + "step": checked_step, + }, + "diagnostics": diagnostics or {"artifacts": []}, + "reported_at": now, + } + connection.execute( + "UPDATE pdd_tasks SET status = ?, current_step = ?," + " retry_count = retry_count + ?, last_error_code = ?," + " last_error_message = ?, finished_at = ?, updated_at = ?" + " WHERE id = ?", + ( + status.value, + "failed", + 1 if status is TaskStatus.RETRY_WAIT else 0, + error_code, + error_message, + now, + now, + task["id"], + ), + ) + run_status = { + TaskStatus.CANCELLED: RunStatus.CANCELLED, + TaskStatus.MANUAL_REVIEW: RunStatus.MANUAL_REVIEW, + }.get(status, RunStatus.FAILED) + cursor = connection.execute( + "UPDATE task_runs SET run_status = ?, current_step = ?," + " error_code = ?, error_message = ?, finished_at = ?," + " updated_at = ? WHERE task_id = ? AND attempt_id = ?" + " AND run_status = 'running'", + ( + run_status.value, + checked_step, + error_code, + error_message, + now, + now, + task["id"], + attempt_id, + ), + ) + if cursor.rowcount != 1: + raise ValueError("采购执行记录不存在或已经结束") + idempotency_key = ( + f"{remote_task_id}:{attempt_id}:failure-v1" + ) + event_cursor = connection.execute( + "INSERT INTO outbox_events (task_id, event_type," + " idempotency_key, payload_json, status, created_at, updated_at)" + " VALUES (?, 'task_failure', ?, ?, 'pending', ?, ?)", + ( + task["id"], + idempotency_key, + json.dumps(payload, ensure_ascii=False), + now, + now, + ), + ) + event_id = int(event_cursor.lastrowid) + event = self.get_outbox_event(event_id) + assert event is not None + return event + finally: + connection.close() + def save_collect_result( self, remote_task_id: str, @@ -543,7 +842,10 @@ class TaskRepository: " WHERE id = ?", (now, now, event_id), ) - if row["event_type"] == OutboxEventType.COLLECT_RESULT.value: + if row["event_type"] in { + OutboxEventType.COLLECT_RESULT.value, + OutboxEventType.PURCHASE_RESULT.value, + }: connection.execute( "UPDATE pdd_tasks SET status = 'succeeded'," " current_step = 'completed', updated_at = ? WHERE id = ?", @@ -579,6 +881,21 @@ class TaskRepository: ] return min(prices) if prices else None + @staticmethod + def _purchase_result_price( + pdd_data: Dict[str, object] + ) -> Optional[int]: + purchase = pdd_data.get("purchase") + if not isinstance(purchase, dict): + return None + confirmed = purchase.get("confirmed") + if not isinstance(confirmed, dict): + return None + price = confirmed.get("unit_price_cent") + if isinstance(price, bool) or not isinstance(price, int): + return None + return price + @staticmethod def _to_outbox(row: sqlite3.Row) -> OutboxEventRecord: return OutboxEventRecord( diff --git a/client/test/test_purchase_task_service.py b/client/test/test_purchase_task_service.py new file mode 100644 index 0000000..4346206 --- /dev/null +++ b/client/test/test_purchase_task_service.py @@ -0,0 +1,278 @@ +"""采购任务演练应用服务测试;全部使用 Mock,不连接真实手机。""" + +import tempfile +import unittest +from pathlib import Path + +from src.admin_gateway import ( + AdminTask, + AndroidDeviceInfo, + ClaimCapabilities, + ClientInfo, +) +from src.mock_admin_gateway import MockAdminGateway +from src.pdd_purchase_adapter import PddPurchaseAdapter, PurchasePageState +from src.purchase_task_service import PurchaseTaskService +from src.task_models import NewClaimedTask, TaskStatus, TaskType +from src.task_repository import TaskRepository + + +OPTIONS = {"color": "黑色", "size": "L", "bundle": "标准版"} + + +class RecordingDryRunAdapter(PddPurchaseAdapter): + """只记录调用的演练适配器,不提供提交订单能力。""" + + def __init__( + self, + *, + goods_id: str = "737116531267", + price_cent: int = 4200, + candidate_count: int = 1, + forced_page: str = "", + wrong_options: bool = False, + ) -> None: + self.goods_id = goods_id + self.price_cent = price_cent + self.candidate_count = candidate_count + self.forced_page = forced_page + self.wrong_options = wrong_options + self.options = {} + self.quantity = 0 + self.page_kind = "goods" + self.calls = [] + + def open_goods(self, goods_url: str) -> None: + self.calls.append(("open_goods", goods_url)) + + def read_state(self) -> PurchasePageState: + self.calls.append(("read_state",)) + options = self.options + if self.wrong_options and options: + options = {**options, "size": "XL"} + return PurchasePageState( + page_kind=self.forced_page or self.page_kind, + goods_id=self.goods_id, + selected_options=dict(options), + quantity=self.quantity, + price_cent=self.price_cent, + candidate_count=self.candidate_count, + ) + + def select_options(self, options) -> None: + self.calls.append(("select_options", dict(options))) + self.options = dict(options) + + def set_quantity(self, quantity: int) -> None: + self.calls.append(("set_quantity", quantity)) + self.quantity = quantity + + def enter_confirmation(self) -> None: + self.calls.append(("enter_confirmation",)) + self.page_kind = "order_confirmation" + + def stop_before_submit(self) -> None: + self.calls.append(("stop_before_submit",)) + + def close(self) -> None: + self.calls.append(("close",)) + + +class PurchaseTaskServiceTest(unittest.TestCase): + def setUp(self) -> None: + self.temp_dir = tempfile.TemporaryDirectory() + self.repository = TaskRepository( + Path(self.temp_dir.name) / "client.db" + ) + self.gateway = MockAdminGateway() + self.client = ClientInfo("CLIENT-001", "测试客户端") + + def tearDown(self) -> None: + self.temp_dir.cleanup() + + def _prepare_task(self, *, task_id: str = "PUR-001") -> None: + task = AdminTask( + task_id=task_id, + task_type=TaskType.PURCHASE, + version=1, + priority=10, + payload={ + "goods_url": ( + "https://mobile.yangkeduo.com/goods.html?" + "goods_id=737116531267" + ), + "goods_id": "737116531267", + "options": dict(OPTIONS), + "quantity": 2, + "max_price_cent": 5000, + }, + ) + self.gateway.enqueue_task(task, self.client.client_id) + claimed = self.gateway.claim_next( + self.client, + ClaimCapabilities( + device=AndroidDeviceInfo("USB-001"), + supported_types=(TaskType.PURCHASE,), + purchase_mode="dry_run", + ), + ) + assert claimed is not None + self.repository.add_claimed_task( + NewClaimedTask( + remote_task_id=claimed.task_id, + task_type=claimed.task_type, + goods_url=str(claimed.payload["goods_url"]), + goods_id=str(claimed.payload["goods_id"]), + quantity=int(claimed.payload["quantity"]), + priority=claimed.priority, + version=claimed.version, + admin_payload={ + "id": claimed.task_id, + "type": claimed.task_type.value, + "version": claimed.version, + "priority": claimed.priority, + "payload": dict(claimed.payload), + }, + ) + ) + + def _service(self, adapter: RecordingDryRunAdapter): + return PurchaseTaskService( + self.gateway, + self.repository, + self.client, + "USB-001", + lambda _address, _cancelled: adapter, + ) + + def test_dynamic_options_dry_run_stops_before_order_submission(self): + self._prepare_task() + adapter = RecordingDryRunAdapter() + + outcome = self._service(adapter).execute_one_local() + + self.assertEqual(outcome.kind, "succeeded") + self.assertIn("没有提交订单", outcome.message) + self.assertIn(("select_options", OPTIONS), adapter.calls) + self.assertIn(("set_quantity", 2), adapter.calls) + self.assertIn(("enter_confirmation",), adapter.calls) + self.assertIn(("stop_before_submit",), adapter.calls) + detail = self.repository.get_task("PUR-001") + assert detail is not None + self.assertEqual(detail.status, TaskStatus.SUCCEEDED) + purchase = detail.pdd_data["purchase"] + self.assertEqual(purchase["mode"], "dry_run") + self.assertEqual(purchase["requested"]["options"], OPTIONS) + self.assertEqual(purchase["confirmed"]["options"], OPTIONS) + self.assertFalse(purchase["order_submitted"]) + self.assertFalse(purchase["payment_attempted"]) + + def test_price_above_limit_stops_before_confirmation(self): + self._prepare_task() + adapter = RecordingDryRunAdapter(price_cent=5001) + + outcome = self._service(adapter).execute_one_local() + + self.assertEqual(outcome.kind, "failed") + self.assertNotIn(("enter_confirmation",), adapter.calls) + detail = self.repository.get_task("PUR-001") + assert detail is not None + self.assertEqual(detail.status, TaskStatus.MANUAL_REVIEW) + self.assertEqual(detail.last_error_code, "PURCHASE_PRICE_EXCEEDED") + + def test_result_submit_timeout_does_not_run_adapter_twice(self): + self._prepare_task() + adapter = RecordingDryRunAdapter() + self.gateway.timeout_next_call() + + first = self._service(adapter).execute_one_local() + first_call_count = len(adapter.calls) + second = self._service(adapter).execute_one_local() + + self.assertEqual(first.kind, "result_pending") + self.assertEqual(second.kind, "succeeded") + self.assertEqual(len(adapter.calls), first_call_count) + self.assertEqual(self.gateway.submission_count, 1) + + def test_invalid_pages_and_ambiguous_target_fail_safely(self): + cases = ( + ("captcha", 1, "PDD_PAGE_CAPTCHA"), + ("login_required", 1, "PDD_PAGE_LOGIN_REQUIRED"), + ("unknown", 1, "PDD_PAGE_UNKNOWN"), + ("", 2, "PURCHASE_AMBIGUOUS_TARGET"), + ) + for index, (page, candidates, expected_code) in enumerate(cases): + with self.subTest(page=page, candidates=candidates): + task_id = f"PUR-BLOCK-{index}" + self._prepare_task(task_id=task_id) + adapter = RecordingDryRunAdapter( + forced_page=page, + candidate_count=candidates, + ) + + outcome = self._service(adapter).execute_selected(task_id) + + self.assertEqual(outcome.kind, "failed") + self.assertNotIn(("enter_confirmation",), adapter.calls) + detail = self.repository.get_task(task_id) + assert detail is not None + self.assertEqual(detail.last_error_code, expected_code) + + def test_wrong_goods_or_options_never_reaches_confirmation(self): + cases = ( + ({"goods_id": "OTHER"}, "PURCHASE_GOODS_MISMATCH"), + ({"wrong_options": True}, "PURCHASE_OPTIONS_MISMATCH"), + ) + for index, (kwargs, expected_code) in enumerate(cases): + with self.subTest(expected_code=expected_code): + task_id = f"PUR-MISMATCH-{index}" + self._prepare_task(task_id=task_id) + adapter = RecordingDryRunAdapter(**kwargs) + + self._service(adapter).execute_selected(task_id) + + self.assertNotIn(("enter_confirmation",), adapter.calls) + detail = self.repository.get_task(task_id) + assert detail is not None + self.assertEqual(detail.last_error_code, expected_code) + + def test_purchase_task_with_missing_safety_fields_is_reported(self): + task_id = "PUR-INVALID" + task = AdminTask( + task_id=task_id, + task_type=TaskType.PURCHASE, + version=1, + priority=0, + payload={"goods_url": "https://example.test/goods"}, + ) + self.gateway.enqueue_task(task, self.client.client_id) + self.gateway.claim_next( + self.client, + ClaimCapabilities(supported_types=(TaskType.PURCHASE,)), + ) + self.repository.add_claimed_task( + NewClaimedTask( + remote_task_id=task_id, + task_type=TaskType.PURCHASE, + goods_url="https://example.test/goods", + admin_payload={ + "id": task_id, + "type": "purchase", + "version": 1, + "payload": dict(task.payload), + }, + ) + ) + adapter = RecordingDryRunAdapter() + + outcome = self._service(adapter).execute_selected(task_id) + + self.assertEqual(outcome.kind, "failed") + self.assertNotIn(("open_goods", "https://example.test/goods"), adapter.calls) + detail = self.repository.get_task(task_id) + assert detail is not None + self.assertEqual(detail.last_error_code, "PURCHASE_TASK_INVALID") + + +if __name__ == "__main__": + unittest.main() diff --git a/docs/client/02-architecture.md b/docs/client/02-architecture.md index 647f76c..5bf86ed 100644 --- a/docs/client/02-architecture.md +++ b/docs/client/02-architecture.md @@ -89,7 +89,7 @@ client/ | 目录分层 | domain / application / infrastructure / workers / ui | 平铺在 `src/` 下:`db.py`、`db_schema.py`、`task_repository.py`、`settings_repository.py`、`task_models.py` 等,另有 `src/util/`、`src/demo1/` | | 主按钮文案 | 「开始自动获取」⇄「停止自动获取」 | 已按持续串行模式实现 | | Admin 网关 | `AdminGateway` + Mock/HTTP 两实现 | 登记、领取、结果和失败提交已实现 | -| 采集任务协调器 | `CollectTaskService` | 已实现单条采集任务流程;采购仍未接入 | +| 任务应用服务 | `CollectTaskService` / `PurchaseTaskService` | 采集已接入自动获取;采购已实现本地 `dry_run` 演练,领取与分派由后续工单接入 | | Outbox 提交 | 从 `outbox_events` 取件重试 | 采集结果与失败已实现,重试不重复采集 | | PDD 自动化 | `infrastructure/pdd/` 适配层 | `pdd_device_service.py` 与 `pdd_collect_service.py` 已接入采集主链 | @@ -350,6 +350,11 @@ purchase(task, mode) -> PurchaseResult reconcile_purchase(task, run) -> PurchaseResult | ManualReview ``` +采购演练通过 `PddPurchaseAdapter` 的窄接口逐步读取最新页面状态。该接口只提供 +打开商品、读取状态、精确选择动态规格、设置数量、进入提交前确认页和停止, +**不提供提交订单或付款方法**。这样即使应用层调用错误,也没有可误触的真实下单入口。 +采购规格使用完整 `options` 对象精确比较,不假定只有颜色和尺码两个维度。 + 现有 `wait_goods_page`、规格面板坐标、颜色尺码选择和下单按钮定位函数可以迁移到该适配层。实验脚本中的硬编码商品、设备、文件路径和 `print` 不得进入正式服务。 PDD 页面可能出现登录失效、验证码、控件树不完整、A/B 页面、库存变化和价格变化。适配层必须返回结构化错误,不得把这些情况统一返回 `False`。 diff --git a/docs/client/03-data-model.md b/docs/client/03-data-model.md index 53d28d5..ad2b193 100644 --- a/docs/client/03-data-model.md +++ b/docs/client/03-data-model.md @@ -544,18 +544,27 @@ CREATE TABLE app_settings ( "purchase": { "mode": "dry_run", "requested": { - "color": "黑色", - "size": "L", + "options": { + "color": "黑色", + "size": "L", + "bundle": "标准版" + }, "quantity": 2, "max_price_cent": 4200 }, "confirmed": { - "color": "黑色", - "size": "L", + "options": { + "color": "黑色", + "size": "L", + "bundle": "标准版" + }, "quantity": 2, "unit_price_cent": 3990, "total_price_cent": 7980 }, + "confirmation_reached": true, + "order_submitted": false, + "payment_attempted": false, "order_no": null, "ordered_at": null, "ordered_at_raw": null, @@ -564,6 +573,11 @@ CREATE TABLE app_settings ( } ``` +`requested.options` 和 `confirmed.options` 都是动态对象,键名来自 Admin 和 PDD +页面,不允许写死成颜色、尺码,也不做相似匹配。`dry_run` 必须同时满足 +`order_submitted=false`、`payment_attempted=false` 和 +`match_status=not_submitted`;只表示已经安全到达最终提交前确认页。 + 真实下单后 `mode` 为 `live`,`match_status` 只能是 `matched`、`ambiguous` 或 `not_found`。只有 `matched` 可以自动报告采购成功。 ## 9. JSON 兼容规则