diff --git a/client/src/admin_gateway.py b/client/src/admin_gateway.py index eb88073..f856639 100644 --- a/client/src/admin_gateway.py +++ b/client/src/admin_gateway.py @@ -56,8 +56,8 @@ class ClaimCapabilities: 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 self.purchase_mode != "dry_run": + raise ValueError("当前版本只允许 purchase_mode=dry_run") if not self.schema_versions or any( version <= 0 for version in self.schema_versions ): diff --git a/client/src/pdd_purchase_reconcile_adapter.py b/client/src/pdd_purchase_reconcile_adapter.py new file mode 100644 index 0000000..be24f00 --- /dev/null +++ b/client/src/pdd_purchase_reconcile_adapter.py @@ -0,0 +1,48 @@ +"""采购结果只读核对适配层。 + +该接口只允许读取订单候选,故意不提供商品选择、提交订单或付款方法。 +""" + +from abc import ABC, abstractmethod +from dataclasses import dataclass, field +from typing import Mapping, Optional + + +@dataclass(frozen=True) +class PurchaseReconcileQuery: + """从本地任务与执行记录构造的核对条件。""" + + goods_id: str + options: Mapping[str, str] + quantity: int + irreversible_action_at: str + + +@dataclass(frozen=True) +class PurchaseReconcileObservation: + """只读核对看到的结果。""" + + match_status: str + order_no: Optional[str] = None + ordered_at: Optional[str] = None + diagnostics: Mapping[str, object] = field(default_factory=dict) + + def __post_init__(self) -> None: + if self.match_status not in { + "matched", "not_found", "ambiguous", "unknown" + }: + raise ValueError("采购核对结果无效") + + +class PddPurchaseReconcileAdapter(ABC): + """已进入不可逆阶段后使用的只读订单核对会话。""" + + @abstractmethod + def read_order_match( + self, query: PurchaseReconcileQuery + ) -> PurchaseReconcileObservation: + """读取并核对订单候选;不得点击下单或付款。""" + + @abstractmethod + def close(self) -> None: + """释放读取会话;不得在此方法中产生点击。""" diff --git a/client/src/pdd_ui_event.py b/client/src/pdd_ui_event.py index fab9846..fbe9023 100644 --- a/client/src/pdd_ui_event.py +++ b/client/src/pdd_ui_event.py @@ -42,6 +42,7 @@ from .current_client_service import CurrentClientService from .http_admin_gateway import DEFAULT_ADMIN_BASE_URL, HttpAdminGateway from .pdd_ui import PDDTaskPage, TaskRow from .purchase_task_service import PurchaseAdapterFactory +from .purchase_reconcile_service import PurchaseReconcileFactory from .selected_android_device_service import SelectedAndroidDeviceService from .settings_repository import SettingsRepository from .task_models import ( @@ -111,6 +112,7 @@ class ClaimTaskWorker(QObject): android_device_service: SelectedAndroidDeviceService, collect_service_factory: Optional[CollectServiceFactory] = None, purchase_adapter_factory: Optional[PurchaseAdapterFactory] = None, + purchase_reconcile_factory: Optional[PurchaseReconcileFactory] = None, selected_task_id: str = "", ) -> None: super().__init__() @@ -121,6 +123,7 @@ class ClaimTaskWorker(QObject): self._cancelled = False self._collect_service_factory = collect_service_factory self._purchase_adapter_factory = purchase_adapter_factory + self._purchase_reconcile_factory = purchase_reconcile_factory self._selected_task_id = selected_task_id def cancel(self) -> None: @@ -161,6 +164,7 @@ class ClaimTaskWorker(QObject): android_serial or "", collect_service_factory=self._collect_service_factory, purchase_adapter_factory=self._purchase_adapter_factory, + purchase_reconcile_factory=self._purchase_reconcile_factory, cancelled=lambda: self._cancelled, ) result = dispatcher.execute_one() @@ -195,6 +199,7 @@ class PDDTaskPageEvent(QObject): settings_repository: Optional[SettingsRepository] = None, collect_service_factory: Optional[CollectServiceFactory] = None, purchase_adapter_factory: Optional[PurchaseAdapterFactory] = None, + purchase_reconcile_factory: Optional[PurchaseReconcileFactory] = None, next_task_delay_ms: int = 500, no_task_delay_ms: int = 5_000, retry_delays_ms: tuple[int, ...] = (5_000, 10_000, 20_000, 30_000), @@ -215,6 +220,7 @@ class PDDTaskPageEvent(QObject): self._detail_windows: Dict[str, TaskDetailWindow] = {} self._collect_service_factory = collect_service_factory self._purchase_adapter_factory = purchase_adapter_factory + self._purchase_reconcile_factory = purchase_reconcile_factory if next_task_delay_ms < 0 or no_task_delay_ms <= 0: raise ValueError("自动获取等待时间配置无效") if not retry_delays_ms or any(value <= 0 for value in retry_delays_ms): @@ -468,6 +474,7 @@ class PDDTaskPageEvent(QObject): android_device_service=self._selected_android_device_service, collect_service_factory=self._collect_service_factory, purchase_adapter_factory=self._purchase_adapter_factory, + purchase_reconcile_factory=self._purchase_reconcile_factory, ) worker.moveToThread(thread) thread.started.connect(worker.run) diff --git a/client/src/purchase_reconcile_service.py b/client/src/purchase_reconcile_service.py new file mode 100644 index 0000000..5189526 --- /dev/null +++ b/client/src/purchase_reconcile_service.py @@ -0,0 +1,138 @@ +"""不可逆阶段中断后的只读采购结果核对。""" + +from dataclasses import dataclass +from typing import Callable, Mapping + +from .pdd_purchase_reconcile_adapter import ( + PddPurchaseReconcileAdapter, + PurchaseReconcileObservation, + PurchaseReconcileQuery, +) +from .task_models import TaskDetail +from .task_repository import TaskRepository + + +PurchaseReconcileFactory = Callable[ + [str, Callable[[], bool]], PddPurchaseReconcileAdapter +] + + +@dataclass(frozen=True) +class PurchaseReconcileOutcome: + """只读核对的简短结果。""" + + kind: str + message: str + task_id: str = "" + + +class PurchaseReconcileService: + """只读核对一条已进入不可逆阶段的采购运行。""" + + def __init__( + self, + repository: TaskRepository, + device_address: str, + adapter_factory: PurchaseReconcileFactory, + *, + cancelled: Callable[[], bool] = lambda: False, + ) -> None: + self._repository = repository + self._device_address = str(device_address or "").strip() + self._factory = adapter_factory + self._cancelled = cancelled + + def execute_selected( + self, remote_task_id: str + ) -> PurchaseReconcileOutcome: + """执行一次只读核对,任何结果都交给人工最终确认。""" + + if not self._device_address: + raise ValueError("请先在设置页选择并保存 Android 设备") + if self._cancelled(): + return PurchaseReconcileOutcome( + "cancelled", "采购结果核对已取消", remote_task_id + ) + task = self._repository.get_task(remote_task_id) + run = self._repository.latest_task_run(remote_task_id) + if task is None or run is None: + raise ValueError(f"任务 {remote_task_id} 或执行记录不存在") + if run.irreversible_action_at is None: + raise ValueError("未进入不可逆阶段,不应启动订单核对") + + query = self._query(task, run.irreversible_action_at) + adapter = None + close_error = "" + try: + adapter = self._factory(self._device_address, self._cancelled) + observation = adapter.read_order_match(query) + if not isinstance(observation, PurchaseReconcileObservation): + raise TypeError("采购核对 Adapter 返回值无效") + except Exception as exc: + observation = PurchaseReconcileObservation( + "unknown", diagnostics={"error": str(exc)} + ) + finally: + if adapter is not None: + try: + adapter.close() + except Exception as exc: + close_error = str(exc) + + diagnostics = dict(observation.diagnostics) + diagnostics.update( + { + "order_no": observation.order_no, + "ordered_at": observation.ordered_at, + "mode": "reconcile_only", + } + ) + if close_error: + diagnostics["close_error"] = close_error + self._repository.save_purchase_reconciliation( + remote_task_id, + run.attempt_id, + observation.match_status, + diagnostics, + ) + if observation.match_status == "matched": + message = ( + f"任务 {remote_task_id} 仅核对到唯一候选订单;" + "请人工确认,程序不会重新下单" + ) + else: + message = ( + f"任务 {remote_task_id} 核对结果不确定;" + "需人工处理,程序不会重新下单" + ) + return PurchaseReconcileOutcome( + "manual_review", message, remote_task_id + ) + + @staticmethod + def _query( + task: TaskDetail, irreversible_action_at: str + ) -> PurchaseReconcileQuery: + payload_root = task.admin_payload + payload = payload_root.get("payload") + if not isinstance(payload, Mapping): + raise ValueError("采购任务缺少 payload") + options = payload.get("options") + if not isinstance(options, Mapping) or not options: + raise ValueError("采购任务缺少 options") + goods_id = str(payload.get("goods_id") or "").strip() + quantity = payload.get("quantity") + if not goods_id: + raise ValueError("采购任务缺少 goods_id") + if ( + isinstance(quantity, bool) + or not isinstance(quantity, int) + or quantity <= 0 + ): + raise ValueError("采购任务缺少有效 quantity") + return PurchaseReconcileQuery( + goods_id=goods_id, + options={str(k): str(v) for k, v in options.items()}, + quantity=quantity, + irreversible_action_at=irreversible_action_at, + ) diff --git a/client/src/purchase_task_service.py b/client/src/purchase_task_service.py index d69598f..fba6c5a 100644 --- a/client/src/purchase_task_service.py +++ b/client/src/purchase_task_service.py @@ -164,7 +164,7 @@ class PurchaseTaskService: exc.code, exc.message, exc.retryable, - exc.step or step, + step, exc.diagnostics, ) except Exception as exc: @@ -367,9 +367,9 @@ class PurchaseTaskService: ) -> 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, diff --git a/client/src/task_detail_view.py b/client/src/task_detail_view.py index 34e80ed..58f918d 100644 --- a/client/src/task_detail_view.py +++ b/client/src/task_detail_view.py @@ -59,6 +59,20 @@ CURRENT_STEP_TEXT = { "completed": "已完成", "failed": "执行失败", "interrupted": "上次执行中断", + "purchase_prepare": "正在准备采购演练", + "purchase_open_goods": "正在打开采购商品页", + "purchase_verify_goods": "正在核对商品", + "purchase_select_options": "正在选择采购规格", + "purchase_verify_options": "正在核对采购规格", + "purchase_set_quantity": "正在设置采购数量", + "purchase_verify_quantity_price": "正在核对数量和价格", + "purchase_enter_confirmation": "正在进入提交前确认页", + "purchase_verify_confirmation": "正在核对提交前确认页", + "purchase_dry_run_stopped": "采购演练已在提交前停止", + "purchase_recovery_ready": "上次演练中断,已安全等待恢复", + "reconcile_purchase": "只允许核对订单", + "reconcile_completed": "只读核对完成,等待人工确认", + "reconcile_manual_review": "核对结果不确定,需人工处理", } diff --git a/client/src/task_dispatcher.py b/client/src/task_dispatcher.py index 238396f..f2bacbf 100644 --- a/client/src/task_dispatcher.py +++ b/client/src/task_dispatcher.py @@ -18,6 +18,10 @@ from .purchase_task_service import ( PurchaseAdapterFactory, PurchaseTaskService, ) +from .purchase_reconcile_service import ( + PurchaseReconcileFactory, + PurchaseReconcileService, +) from .task_models import ( NewClaimedTask, OutboxEventRecord, @@ -120,6 +124,7 @@ class TaskDispatcher: *, collect_service_factory: Optional[CollectServiceFactory] = None, purchase_adapter_factory: Optional[PurchaseAdapterFactory] = None, + purchase_reconcile_factory: Optional[PurchaseReconcileFactory] = None, cancelled: Callable[[], bool] = lambda: False, ) -> None: self._gateway = gateway @@ -128,6 +133,7 @@ class TaskDispatcher: self._device_address = str(device_address or "").strip() self._collect_factory = collect_service_factory self._purchase_factory = purchase_adapter_factory + self._reconcile_factory = purchase_reconcile_factory self._cancelled = cancelled @property @@ -165,6 +171,23 @@ class TaskDispatcher: if not self._device_address: raise ValueError("请先在设置页选择并保存 Android 设备") + reconcile_task = self._repository.next_purchase_reconcile_task() + if reconcile_task is not None: + 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 + ) + if not self.purchase_ready: pending_purchase = self._repository.next_purchase_task() if pending_purchase is not None: diff --git a/client/src/task_repository.py b/client/src/task_repository.py index 01c3c49..6dd781d 100644 --- a/client/src/task_repository.py +++ b/client/src/task_repository.py @@ -23,6 +23,7 @@ from .task_models import ( TaskStatus, TaskSummary, TaskType, + TaskRunRecord, ) @@ -156,7 +157,7 @@ class TaskRepository: return self._to_detail(row) if row is not None else None def recover_interrupted_work(self) -> None: - """恢复上次异常退出留下的可重试状态。""" + """恢复上次异常退出留下的任务,不猜测内存状态。""" now = utc_now_iso() connection = open_database(self._db_path) @@ -193,6 +194,64 @@ class TaskRepository: " AND task_id IN (SELECT id FROM pdd_tasks WHERE task_type = 'collect')", (now, now), ) + purchase_runs = connection.execute( + "SELECT t.id AS task_id, r.irreversible_action_at" + " FROM pdd_tasks t JOIN task_runs r ON r.task_id = t.id" + " WHERE t.task_type = 'purchase' AND t.status = 'running'" + " AND r.run_status = 'running'" + " ORDER BY r.attempt_no DESC" + ).fetchall() + recovered_task_ids = set() + for run in purchase_runs: + task_id = int(run["task_id"]) + if task_id in recovered_task_ids: + continue + recovered_task_ids.add(task_id) + irreversible = connection.execute( + "SELECT 1 FROM task_runs WHERE task_id = ?" + " AND run_status = 'running'" + " AND irreversible_action_at IS NOT NULL LIMIT 1", + (task_id,), + ).fetchone() + if irreversible is not None: + message = ( + "上次采购在不可逆阶段中断," + "只允许核对订单" + ) + connection.execute( + "UPDATE task_runs SET run_status = 'manual_review'," + " current_step = 'reconcile_purchase'," + " error_code = 'PURCHASE_OUTCOME_UNKNOWN'," + " error_message = ?, finished_at = ?, updated_at = ?" + " WHERE task_id = ? AND run_status = 'running'", + (message, now, now, task_id), + ) + connection.execute( + "UPDATE pdd_tasks SET status = 'manual_review'," + " current_step = 'reconcile_purchase'," + " last_error_code = 'PURCHASE_OUTCOME_UNKNOWN'," + " last_error_message = ?, finished_at = ?, updated_at = ?" + " WHERE id = ?", + (message, now, now, task_id), + ) + else: + message = "采购演练上次执行中断,已关闭旧执行记录" + connection.execute( + "UPDATE task_runs SET run_status = 'failed'," + " error_code = 'CLIENT_INTERRUPTED'," + " error_message = ?, finished_at = ?, updated_at = ?" + " WHERE task_id = ? AND run_status = 'running'", + (message, now, now, task_id), + ) + connection.execute( + "UPDATE pdd_tasks SET status = 'claimed'," + " current_step = 'purchase_recovery_ready'," + " retry_count = retry_count + 1," + " last_error_code = 'CLIENT_INTERRUPTED'," + " last_error_message = ?, finished_at = NULL," + " updated_at = ? WHERE id = ?", + (message, now, task_id), + ) finally: connection.close() @@ -224,6 +283,39 @@ class TaskRepository: connection.close() return self._to_detail(row) if row is not None else None + def next_purchase_reconcile_task(self) -> Optional[TaskDetail]: + """返回最早一条只允许读取核对的采购任务。""" + + connection = open_database(self._db_path) + try: + row = connection.execute( + "SELECT * FROM pdd_tasks WHERE task_type = 'purchase'" + " AND status = 'manual_review'" + " AND current_step = 'reconcile_purchase'" + " 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 latest_task_run( + self, remote_task_id: str + ) -> Optional[TaskRunRecord]: + """返回任务最新的执行记录。""" + + connection = open_database(self._db_path) + try: + row = connection.execute( + "SELECT r.* FROM task_runs r" + " JOIN pdd_tasks t ON t.id = r.task_id" + " WHERE t.remote_task_id = ?" + " ORDER BY r.attempt_no DESC LIMIT 1", + (remote_task_id,), + ).fetchone() + finally: + connection.close() + return self._to_task_run(row) if row is not None else None + def next_runnable_task( self, *, include_purchase: bool ) -> Optional[TaskDetail]: @@ -479,6 +571,92 @@ class TaskRepository: finally: connection.close() + def save_purchase_reconciliation( + self, + remote_task_id: str, + attempt_id: str, + match_status: str, + diagnostics: Dict[str, object], + ) -> None: + """保存一次只读订单核对结果,任务仍留给人工确认。""" + + allowed = {"matched", "not_found", "ambiguous", "unknown"} + if match_status not in allowed: + raise ValueError("采购核对结果无效") + now = utc_now_iso() + step = ( + "reconcile_completed" + if match_status == "matched" + else "reconcile_manual_review" + ) + messages = { + "matched": "只读核对发现唯一候选订单,请人工确认", + "not_found": "只读核对未找到订单,不得重新下单", + "ambiguous": "只读核对发现多个候选订单,请人工确认", + "unknown": "无法确定采购结果,不得重新下单", + } + message = messages[match_status] + saved_diagnostics = dict(diagnostics) + saved_diagnostics["match_status"] = match_status + connection = open_database(self._db_path) + try: + with connection: + task = connection.execute( + "SELECT id, task_type, status, current_step 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.MANUAL_REVIEW.value + or task["current_step"] != "reconcile_purchase" + ): + raise ValueError("采购任务当前不在待核对状态") + run = connection.execute( + "SELECT irreversible_action_at, run_status, current_step" + " FROM task_runs" + " WHERE task_id = ? AND attempt_id = ?", + (task["id"], attempt_id), + ).fetchone() + if run is None or run["irreversible_action_at"] is None: + raise ValueError("只有已进入不可逆阶段的运行才能核对") + if ( + run["run_status"] != RunStatus.MANUAL_REVIEW.value + or run["current_step"] != "reconcile_purchase" + ): + raise ValueError("采购执行记录已经核对或状态已变更") + run_cursor = connection.execute( + "UPDATE task_runs SET run_status = 'manual_review'," + " current_step = ?, error_code = 'PURCHASE_RECONCILED'," + " error_message = ?, diagnostics_json = ?, updated_at = ?" + " WHERE task_id = ? AND attempt_id = ?", + ( + step, + message, + json.dumps(saved_diagnostics, ensure_ascii=False), + now, + task["id"], + attempt_id, + ), + ) + if run_cursor.rowcount != 1: + raise ValueError("采购执行记录核对保存失败") + task_cursor = connection.execute( + "UPDATE pdd_tasks SET status = 'manual_review'," + " current_step = ?, last_error_code = 'PURCHASE_RECONCILED'," + " last_error_message = ?, updated_at = ? WHERE id = ?" + " AND status = 'manual_review'" + " AND current_step = 'reconcile_purchase'", + (step, message, now, task["id"]), + ) + if task_cursor.rowcount != 1: + raise ValueError("采购任务状态已变更,核对结果未保存") + finally: + connection.close() + def save_purchase_result( self, remote_task_id: str, @@ -628,7 +806,8 @@ class TaskRepository: }.get(status, RunStatus.FAILED) cursor = connection.execute( "UPDATE task_runs SET run_status = ?, current_step = ?," - " error_code = ?, error_message = ?, finished_at = ?," + " error_code = ?, error_message = ?, diagnostics_json = ?," + " finished_at = ?," " updated_at = ? WHERE task_id = ? AND attempt_id = ?" " AND run_status = 'running'", ( @@ -636,6 +815,7 @@ class TaskRepository: checked_step, error_code, error_message, + json.dumps(diagnostics or {}, ensure_ascii=False), now, now, task["id"], @@ -935,6 +1115,39 @@ class TaskRepository: sent_at=row["sent_at"], ) + @staticmethod + def _to_task_run(row: sqlite3.Row) -> TaskRunRecord: + diagnostics = ( + TaskRepository._load_json_object(row["diagnostics_json"]) + if row["diagnostics_json"] is not None + else None + ) + result_data = ( + TaskRepository._load_json_object(row["result_data"]) + if row["result_data"] is not None + else None + ) + return TaskRunRecord( + id=row["id"], + task_id=row["task_id"], + attempt_id=row["attempt_id"], + attempt_no=row["attempt_no"], + device_address=row["device_address"], + run_status=RunStatus(row["run_status"]), + current_step=row["current_step"], + started_at=row["started_at"], + finished_at=row["finished_at"], + irreversible_action_at=row["irreversible_action_at"], + order_submitted_at=row["order_submitted_at"], + error_code=row["error_code"], + error_message=row["error_message"], + diagnostics_json=diagnostics, + result_data=result_data, + artifact_directory=row["artifact_directory"], + created_at=row["created_at"], + updated_at=row["updated_at"], + ) + @staticmethod def _validate_page(limit: int, offset: int) -> None: if not 1 <= limit <= MAX_PAGE_SIZE: diff --git a/client/test/test_pdd_ui_event.py b/client/test/test_pdd_ui_event.py index 208b5b3..5513c69 100644 --- a/client/test/test_pdd_ui_event.py +++ b/client/test/test_pdd_ui_event.py @@ -52,6 +52,9 @@ class BrokenSaveRepository(BrokenRepository): def next_purchase_task(self): return None + def next_purchase_reconcile_task(self): + return None + class RecordingClaimGateway: """记录领取参数并返回预设结果。""" diff --git a/client/test/test_purchase_recovery.py b/client/test/test_purchase_recovery.py new file mode 100644 index 0000000..14cec46 --- /dev/null +++ b/client/test/test_purchase_recovery.py @@ -0,0 +1,289 @@ +"""采购崩溃恢复和安全门禁测试;不连接真实手机。""" + +import tempfile +import unittest +from pathlib import Path + +from src.admin_gateway import AdminTask, ClaimCapabilities, ClientInfo +from src.db import open_database +from src.mock_admin_gateway import MockAdminGateway +from src.pdd_purchase_adapter import ( + PddPurchaseAdapter, + PddPurchaseError, + PurchasePageState, +) +from src.pdd_purchase_reconcile_adapter import ( + PddPurchaseReconcileAdapter, + PurchaseReconcileObservation, +) +from src.purchase_task_service import PurchaseTaskService +from src.task_dispatcher import TaskDispatcher, admin_task_to_new_claimed_task +from src.task_models import RunStatus, TaskStatus, TaskType +from src.task_repository import TaskRepository + + +OPTIONS = {"color": "黑色", "size": "L"} + + +def purchase_admin_task(task_id: str = "PUR-RECOVER") -> AdminTask: + return AdminTask( + task_id, + TaskType.PURCHASE, + 1, + 0, + { + "goods_url": "https://example.test/goods/737116531267", + "goods_id": "737116531267", + "options": dict(OPTIONS), + "quantity": 2, + "max_price_cent": 5000, + }, + ) + + +class FaultAdapter(PddPurchaseAdapter): + """在指定动作抛错,用来模拟设备断开或进程崩溃。""" + + def __init__(self, fail_action: str = "") -> None: + self.fail_action = fail_action + self.options = {} + self.quantity = 0 + self.page_kind = "goods" + + def _fail(self, action: str) -> None: + if self.fail_action == action: + raise PddPurchaseError( + "DEVICE_DISCONNECTED", + "Android 设备连接中断", + step=action, + retryable=True, + diagnostics={"action": action}, + ) + + def open_goods(self, _goods_url: str) -> None: + self._fail("open_goods") + + def read_state(self) -> PurchasePageState: + self._fail("read_state") + return PurchasePageState( + self.page_kind, + "737116531267", + dict(self.options), + self.quantity, + 4200, + 1, + ) + + def select_options(self, options) -> None: + self._fail("select_options") + self.options = dict(options) + + def set_quantity(self, quantity: int) -> None: + self._fail("set_quantity") + self.quantity = quantity + + def enter_confirmation(self) -> None: + self._fail("enter_confirmation") + self.page_kind = "order_confirmation" + + def stop_before_submit(self) -> None: + self._fail("stop_before_submit") + + def close(self) -> None: + pass + + +class ReadOnlyReconcileAdapter(PddPurchaseReconcileAdapter): + def __init__(self, calls) -> None: + self.calls = calls + + def read_order_match(self, query): + self.calls.append(("reconcile", query.goods_id)) + return PurchaseReconcileObservation( + "matched", "ORDER-001", "2026-08-10T08:00:00Z" + ) + + def close(self) -> None: + self.calls.append(("reconcile_close",)) + + +class PurchaseRecoveryTest(unittest.TestCase): + def setUp(self) -> None: + self.temporary = tempfile.TemporaryDirectory() + self.db_path = Path(self.temporary.name) / "client.db" + self.repository = TaskRepository(self.db_path) + self.gateway = MockAdminGateway() + self.client = ClientInfo("CLIENT-001") + + def tearDown(self) -> None: + self.temporary.cleanup() + + def _add(self, task_id: str = "PUR-RECOVER") -> None: + task = purchase_admin_task(task_id) + self.gateway.enqueue_task(task, self.client.client_id) + claimed = self.gateway.claim_next( + self.client, + ClaimCapabilities(supported_types=(TaskType.PURCHASE,)), + ) + assert claimed is not None + self.repository.add_claimed_task( + admin_task_to_new_claimed_task(claimed) + ) + + def _service(self, adapter, *, cancelled=lambda: False): + return PurchaseTaskService( + self.gateway, + self.repository, + self.client, + "USB-001", + lambda _address, _cancelled: adapter, + cancelled=cancelled, + ) + + def _interrupt_after_irreversible(self, task_id: str) -> None: + self._add(task_id) + started = self.repository.start_purchase_run(task_id, "USB-001") + connection = open_database(self.db_path) + try: + with connection: + connection.execute( + "UPDATE task_runs SET irreversible_action_at = ?" + " WHERE attempt_id = ?", + ("2026-08-10T08:00:00Z", started.attempt_id), + ) + finally: + connection.close() + self.repository.recover_interrupted_work() + + def test_critical_action_failure_keeps_last_persisted_step(self): + cases = { + "open_goods": "purchase_open_goods", + "select_options": "purchase_select_options", + "set_quantity": "purchase_set_quantity", + "enter_confirmation": "purchase_enter_confirmation", + "stop_before_submit": "purchase_dry_run_stopped", + } + for index, (action, expected_step) in enumerate(cases.items()): + with self.subTest(action=action): + task_id = f"PUR-FAULT-{index}" + self._add(task_id) + + self._service(FaultAdapter(action)).execute_selected(task_id) + + run = self.repository.latest_task_run(task_id) + assert run is not None + self.assertEqual(run.current_step, expected_step) + self.assertEqual(run.error_code, "DEVICE_DISCONNECTED") + self.assertEqual(run.diagnostics_json["action"], action) + detail = self.repository.get_task(task_id) + assert detail is not None + self.assertEqual(detail.status, TaskStatus.MANUAL_REVIEW) + + def test_stop_request_is_saved_before_any_device_action(self): + self._add("PUR-STOP") + checks = iter((False, True)) + adapter = FaultAdapter() + + outcome = self._service( + adapter, cancelled=lambda: next(checks, True) + ).execute_selected("PUR-STOP") + + self.assertEqual(outcome.kind, "failed") + detail = self.repository.get_task("PUR-STOP") + run = self.repository.latest_task_run("PUR-STOP") + assert detail is not None and run is not None + self.assertEqual(detail.status, TaskStatus.CANCELLED) + self.assertEqual(run.current_step, "purchase_open_goods") + self.assertEqual(run.run_status, RunStatus.CANCELLED) + + def test_restart_before_irreversible_closes_old_run_then_retries(self): + self._add() + first = self.repository.start_purchase_run("PUR-RECOVER", "USB-001") + self.repository.update_purchase_step( + "PUR-RECOVER", first.attempt_id, "purchase_select_options" + ) + + self.repository.recover_interrupted_work() + + old_run = self.repository.latest_task_run("PUR-RECOVER") + detail = self.repository.get_task("PUR-RECOVER") + assert old_run is not None and detail is not None + self.assertEqual(old_run.run_status, RunStatus.FAILED) + self.assertEqual(old_run.current_step, "purchase_select_options") + self.assertEqual(detail.status, TaskStatus.CLAIMED) + second = self.repository.start_purchase_run( + "PUR-RECOVER", "USB-001" + ) + self.assertEqual(second.attempt_no, 2) + with self.assertRaisesRegex(ValueError, "不能开始演练"): + self.repository.start_purchase_run("PUR-RECOVER", "USB-001") + + def test_irreversible_restart_only_reconciles_and_never_purchases(self): + self._interrupt_after_irreversible("PUR-RECOVER") + calls = [] + + def purchase_factory(_address, _cancelled): + calls.append(("purchase",)) + return FaultAdapter() + + dispatcher = TaskDispatcher( + self.gateway, + self.repository, + self.client, + "USB-001", + purchase_adapter_factory=purchase_factory, + purchase_reconcile_factory=( + lambda _address, _cancelled: ReadOnlyReconcileAdapter(calls) + ), + ) + + first = dispatcher.execute_one() + second = dispatcher.execute_one() + + self.assertEqual(first.kind, "manual_review") + self.assertEqual(second.kind, "no_task") + self.assertNotIn(("purchase",), calls) + self.assertEqual(calls.count(("reconcile", "737116531267")), 1) + detail = self.repository.get_task("PUR-RECOVER") + run = self.repository.latest_task_run("PUR-RECOVER") + assert detail is not None and run is not None + self.assertEqual(detail.current_step, "reconcile_completed") + self.assertEqual(run.diagnostics_json["mode"], "reconcile_only") + + def test_reconcile_device_failure_is_recorded_as_unknown(self): + task_id = "PUR-RECONCILE-OFFLINE" + self._interrupt_after_irreversible(task_id) + purchase_calls = [] + + def unavailable_reconcile(_address, _cancelled): + raise ConnectionError("核对设备已断开") + + outcome = TaskDispatcher( + self.gateway, + self.repository, + self.client, + "USB-001", + purchase_adapter_factory=( + lambda _address, _cancelled: purchase_calls.append("purchase") + ), + purchase_reconcile_factory=unavailable_reconcile, + ).execute_one() + + self.assertEqual(outcome.kind, "manual_review") + self.assertEqual(purchase_calls, []) + detail = self.repository.get_task(task_id) + run = self.repository.latest_task_run(task_id) + assert detail is not None and run is not None + self.assertEqual(detail.current_step, "reconcile_manual_review") + self.assertIn("核对设备已断开", run.diagnostics_json["error"]) + + def test_live_mode_and_order_submission_methods_are_unavailable(self): + with self.assertRaisesRegex(ValueError, "dry_run"): + ClaimCapabilities(purchase_mode="live") + for name in ("submit_order", "pay", "payment"): + self.assertFalse(hasattr(PddPurchaseAdapter, name)) + self.assertFalse(hasattr(PddPurchaseReconcileAdapter, name)) + + +if __name__ == "__main__": + unittest.main() diff --git a/client/test/test_task_detail_view.py b/client/test/test_task_detail_view.py index 6c85fb6..bcb5345 100644 --- a/client/test/test_task_detail_view.py +++ b/client/test/test_task_detail_view.py @@ -119,6 +119,20 @@ class TaskDetailViewTest(unittest.TestCase): self.assertEqual(data.current_step, "未知步骤") + def test_purchase_recovery_steps_are_clear_chinese(self): + cases = { + "purchase_dry_run_stopped": "采购演练已在提交前停止", + "reconcile_purchase": "只允许核对订单", + "reconcile_manual_review": "核对结果不确定,需人工处理", + } + for step, expected in cases.items(): + with self.subTest(step=step): + detail = replace(make_detail(), current_step=step) + self.assertEqual( + build_task_detail_view_data(detail).current_step, + expected, + ) + def test_window_builds_color_section_before_size_section(self): window = TaskDetailWindow(make_detail(collected_data())) diff --git a/docs/client/02-architecture.md b/docs/client/02-architecture.md index da233b4..14aaf80 100644 --- a/docs/client/02-architecture.md +++ b/docs/client/02-architecture.md @@ -340,9 +340,12 @@ Admin 侧必须无条件接受,见 [04](04-admin-api-contract.md) §6.1。 ### 采购崩溃恢复 -- 在最终提交订单前写入不可逆阶段标记。 -- 如果进程在该阶段退出,重启后进入订单核对流程。 -- 核对不到订单时进入人工处理,禁止自动重新下单。 +- 每个采购关键动作前,先在同一 SQLite 事务中更新任务和执行记录的步骤。 +- 演练运行中断且没有不可逆标记时,原执行记录先结束为失败,任务再回到 + `claimed`;恢复执行必须创建新的 `attempt_id`,不会存在两条并发运行。 +- `irreversible_action_at` 有值时,启动恢复立即转为 `manual_review / reconcile_purchase`。 + `PurchaseReconcileService` 只能调用独立的只读 Adapter,不会调用采购 Adapter。 +- 核对到唯一候选也仍需人工最终确认;未找到、多候选或结果不确定均保持人工处理。 - Admin 提交失败只重试 Outbox,不再次操作拼多多。 ## 9. PDD 适配边界 @@ -359,6 +362,8 @@ reconcile_purchase(task, run) -> PurchaseResult | ManualReview 打开商品、读取状态、精确选择动态规格、设置数量、进入提交前确认页和停止, **不提供提交订单或付款方法**。这样即使应用层调用错误,也没有可误触的真实下单入口。 采购规格使用完整 `options` 对象精确比较,不假定只有颜色和尺码两个维度。 +只读核对另用 `PddPurchaseReconcileAdapter`,只暴露 `read_order_match` +和 `close`,不暴露选规格、设数量、下单或付款方法。 现有 `wait_goods_page`、规格面板坐标、颜色尺码选择和下单按钮定位函数可以迁移到该适配层。实验脚本中的硬编码商品、设备、文件路径和 `print` 不得进入正式服务。 diff --git a/docs/client/03-data-model.md b/docs/client/03-data-model.md index ad2b193..0aa692a 100644 --- a/docs/client/03-data-model.md +++ b/docs/client/03-data-model.md @@ -451,14 +451,18 @@ CREATE TABLE app_settings ( | 重启时发现 | 怎么处理 | |---|---| -| `status = 'running'`,且 `task_runs.irreversible_action_at` **为空** | 说明还没下单,改成 `retry_wait`,正常重试 | -| `status = 'running'`,且 `irreversible_action_at` **不为空** | 说明可能已经下单了。**保持 `running`**,把 `current_step` 改成 `reconcile_order` 去核对订单。核对不到就转 `manual_review`。**任何情况下都不许重新下单** | +| 采集 `status = 'running'` | 改成 `retry_wait`,原执行记录结束为失败 | +| 采购 `status = 'running'`,且 `task_runs.irreversible_action_at` **为空** | 原执行记录先结束为失败,任务改回 `claimed / purchase_recovery_ready`;再次演练必须新建执行记录 | +| 采购 `status = 'running'`,且 `irreversible_action_at` **不为空** | 原执行记录和任务都改成 `manual_review / reconcile_purchase`;只调用只读核对,**任何情况下都不许重新下单** | | `status = 'claimed'` | 说明领到了还没开始动手,保持不变,等协调器重新调度执行。**不需要再向 Admin 领一次** | | `status = 'result_pending'` | 不动状态,交给 Outbox 继续重试提交 | | `outbox_events.status = 'sending'` | 改回 `pending`,让它能被重新发送 | 一句话记住:**`irreversible_action_at` 有值 = 只准查,不准买。** +可重试的设备错误在采购演练中也先进入 `manual_review`,不会因为技术上 +“可重试”就隐式重跑采购。这与采集任务的 `retry_wait` 规则不同。 + ## 8. `pdd_data` JSON ### 8.1 通用商品数据 diff --git a/docs/client/05-ui-specification.md b/docs/client/05-ui-specification.md index 764d2ca..db205c3 100644 --- a/docs/client/05-ui-specification.md +++ b/docs/client/05-ui-specification.md @@ -202,6 +202,8 @@ class TaskTableModel(QAbstractTableModel): 当前 Client 使用可调整大小的非模态详情窗口。用户可以点击表格“详情”、双击任务行,或选中任务后按 Enter 打开。相同任务只保留一个详情窗口;再次打开时切回已有窗口。按 Esc 或“关闭”按钮返回任务列表。 详情面向采购人员,任务类型、状态、当前步骤和错误说明优先使用中文。数据库中的 `current_step` 等内部值只用于业务判断和日志,不直接显示;遇到尚未适配的新步骤时显示“未知步骤”。 +采购演练和恢复至少要明确显示“演练已在提交前停止”、“结果待提交”、 +“只允许核对订单”和“需要人工处理”;不能只用颜色表达安全状态。 采集规格按下面顺序展示: diff --git a/docs/client/06-quality-security.md b/docs/client/06-quality-security.md index 6c46d41..a7dc077 100644 --- a/docs/client/06-quality-security.md +++ b/docs/client/06-quality-security.md @@ -82,6 +82,20 @@ PDD 解析测试优先使用脱敏的 XML 固件,不要求每次连接真实 8. 操作人员明确确认真实下单范围和安全设置。 真实下单不得通过调试参数、默认配置或界面误操作意外开启。 +当前版本的 `ClaimCapabilities` 仅接受 `purchase_mode=dry_run`,采购 +Adapter 不提供提交订单或付款方法。开发者不能把自动化测试通过当成真实 +下单门禁放行;仍需项目负责人和实际操作人员共同确认。 + +### 3.1 当前自动化安全检查 + +- 关键手机动作前已持久化 `current_step`。 +- 无不可逆标记的中断会先关闭旧运行,新运行使用新 `attempt_id`。 +- 有不可逆标记的中断只调用只读核对,采购 Adapter 调用次数为 0。 +- 已落库 Outbox 只补交,不重跑手机流程。 +- 停止、设备断开、验证码、登录失效和结果不确定都保留稳定错误和诊断。 +- `purchase_mode=live` 会在数据对象创建时被拒绝。 + +上述只证明演练和恢复代码的默认安全性,**不代表真实下单已放行**。 ## 4. 自动化防护