From 70cedc34d4c2454569a5df427d38dd8c5af87ae5 Mon Sep 17 00:00:00 2001 From: chengma Date: Mon, 10 Aug 2026 18:00:38 +0800 Subject: [PATCH] =?UTF-8?q?refactor:=20=E5=A4=8D=E7=94=A8=E6=8C=81?= =?UTF-8?q?=E4=B9=85=E8=AE=BE=E5=A4=87=E5=B7=A5=E4=BD=9C=E7=BA=BF=E7=A8=8B?= =?UTF-8?q?=20(#110)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- client/src/collect_task_service.py | 7 +- client/src/pdd_device_service.py | 189 +++++++++- client/src/pdd_u2_purchase_adapter.py | 8 +- .../src/pdd_u2_purchase_reconcile_adapter.py | 17 +- client/src/pdd_ui_event.py | 337 +++++++++++------- client/src/settings_ui_event.py | 3 + client/src/ui_main.py | 3 + client/test/test_pdd_device_service.py | 89 ++++- client/test/test_pdd_ui_event.py | 30 ++ client/test/test_settings_ui_event.py | 5 + docs/client/02-architecture.md | 24 +- 11 files changed, 563 insertions(+), 149 deletions(-) diff --git a/client/src/collect_task_service.py b/client/src/collect_task_service.py index b63bcf6..39ac4c0 100644 --- a/client/src/collect_task_service.py +++ b/client/src/collect_task_service.py @@ -13,7 +13,10 @@ from .admin_gateway import ( ClientInfo, ) from .pdd_collect_service import PddCollectError, PddCollectService -from .pdd_device_service import PddDeviceService +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 @@ -258,7 +261,7 @@ class CollectTaskService: cancelled: Callable[[], bool], ) -> PddCollectService: return PddCollectService( - PddDeviceService(), + current_thread_device_service() or PddDeviceService(), device_address, client_id, cancelled=cancelled, diff --git a/client/src/pdd_device_service.py b/client/src/pdd_device_service.py index 66c71ae..60ba2c9 100644 --- a/client/src/pdd_device_service.py +++ b/client/src/pdd_device_service.py @@ -1,14 +1,16 @@ """uiautomator2 设备连接边界。 -每次自动化任务通过 ``connect`` 获取一个会话。会话只能在创建它的线程中 -使用,退出 ``with`` 后立即释放,避免多个任务同时控制同一台手机。 +普通服务每个任务建立一次连接;持久服务只复用 Device 连接,不复用页面数据。 +两种会话都只能在创建它们的线程中使用。 """ from __future__ import annotations import threading -from contextlib import AbstractContextManager -from typing import Any, Callable, Optional +import time +from contextlib import AbstractContextManager, contextmanager +from contextvars import ContextVar +from typing import Any, Callable, Iterator, Optional from .performance_timing import current_performance_trace @@ -151,3 +153,182 @@ class PddDeviceService: with self._state_lock: if self._active_serial == serial: self._active_serial = None + + +class PersistentPddDeviceService(PddDeviceService): + """在固定工作线程内复用同一台设备的 Device 连接。 + + ``connect`` 每次仍调用 ``app_current`` 做健康检查并建立全新的任务会话。 + 缓存连接失效时最多重连一次;页面 XML、选择器和业务数据均不缓存。 + """ + + def __init__( + self, + connector: Optional[Callable[[str], Any]] = None, + *, + ttl_seconds: float = 90.0, + monotonic: Callable[[], float] = time.monotonic, + ) -> None: + super().__init__(connector) + if ttl_seconds <= 0: + raise ValueError("设备会话 TTL 必须大于 0") + self._ttl_seconds = float(ttl_seconds) + self._monotonic = monotonic + self._owner_thread_id: Optional[int] = None + self._cached_serial: Optional[str] = None + self._cached_device: Any = None + self._last_used_at = 0.0 + + @property + def ttl_seconds(self) -> float: + return self._ttl_seconds + + @property + def has_cached_device(self) -> bool: + """只供生命周期管理和测试判断,不返回 Device 对象。""" + + return self._cached_device is not None + + def connect(self, serial: str) -> PddDeviceSession: + """在所有者线程取得独占任务会话,并按需复用底层连接。""" + + self._assert_owner_thread() + checked_serial = self._validate_serial(serial) + with self._state_lock: + if self._active_serial is not None: + raise PddDeviceError( + "DEVICE_IN_USE", + f"设备 {self._active_serial} 正在执行其他自动化任务", + ) + self._active_serial = checked_serial + + try: + device, current = self._healthy_device(checked_serial) + except PddDeviceError: + self._release(checked_serial) + raise + except Exception as exc: + self._release(checked_serial) + details = str(exc).lower() + code = ( + "DEVICE_OFFLINE" + if any( + word in details + for word in ("offline", "not found", "disconnected") + ) + else "DEVICE_CONNECT_FAILED" + ) + raise PddDeviceError( + code, + f"无法连接 Android 设备 {checked_serial}:{exc}", + ) from exc + + return PddDeviceSession( + device, + checked_serial, + current, + lambda: self._release_persistent(checked_serial), + ) + + def release_cached(self) -> None: + """在所有者线程丢弃缓存连接;可重复调用。""" + + self._assert_owner_thread() + with self._state_lock: + if self._active_serial is not None: + raise PddDeviceError( + "DEVICE_IN_USE", + "设备任务尚未到达安全结束点,暂不能释放连接", + ) + self._discard_cached() + + def _healthy_device(self, serial: str) -> tuple[Any, dict[str, Any]]: + now = self._monotonic() + expired = ( + self._cached_device is not None + and now - self._last_used_at >= self._ttl_seconds + ) + if self._cached_serial != serial or expired: + self._discard_cached() + + if self._cached_device is not None: + try: + current = self._read_app_state(self._cached_device) + return self._cached_device, current + except Exception: + # 旧连接失效时只丢弃并重连一次,不能在这里无限重试。 + self._discard_cached() + + trace = current_performance_trace() + if trace is None: + device = self._connector(serial) + else: + with trace.stage("uiautomator2_connect"): + device = self._connector(serial) + current = self._read_app_state(device) + self._cached_serial = serial + self._cached_device = device + self._last_used_at = now + return device, current + + @staticmethod + def _read_app_state_without_trace(device: Any) -> dict[str, Any]: + current = device.app_current() + if not isinstance(current, dict): + raise RuntimeError("uiautomator2 未返回有效的设备状态") + return current + + def _read_app_state(self, device: Any) -> dict[str, Any]: + trace = current_performance_trace() + if trace is None: + return self._read_app_state_without_trace(device) + with trace.stage("first_app_current"): + return self._read_app_state_without_trace(device) + + def _release_persistent(self, serial: str) -> None: + self._assert_owner_thread() + with self._state_lock: + if self._active_serial == serial: + self._active_serial = None + self._last_used_at = self._monotonic() + + def _discard_cached(self) -> None: + # uiautomator2 Device 没有稳定的 close API;清除唯一引用即可让底层 + # HTTP 客户端按库自身生命周期回收,不能猜测调用私有方法。 + self._cached_serial = None + self._cached_device = None + self._last_used_at = 0.0 + + def _assert_owner_thread(self) -> None: + current = threading.get_ident() + if self._owner_thread_id is None: + self._owner_thread_id = current + elif self._owner_thread_id != current: + raise PddDeviceError( + "DEVICE_THREAD_VIOLATION", + "持久 uiautomator2 Device 只能在固定工作线程中使用和释放", + ) + + +_CURRENT_DEVICE_SERVICE: ContextVar[Optional[PersistentPddDeviceService]] = ( + ContextVar("pdd_device_service", default=None) +) + + +@contextmanager +def bind_thread_device_service( + service: PersistentPddDeviceService, +) -> Iterator[None]: + """让当前工作线程创建的采集、采购和核单适配器使用同一连接服务。""" + + token = _CURRENT_DEVICE_SERVICE.set(service) + try: + yield + finally: + _CURRENT_DEVICE_SERVICE.reset(token) + + +def current_thread_device_service() -> Optional[PersistentPddDeviceService]: + """返回当前工作线程绑定的服务;未绑定时返回 ``None``。""" + + return _CURRENT_DEVICE_SERVICE.get() diff --git a/client/src/pdd_u2_purchase_adapter.py b/client/src/pdd_u2_purchase_adapter.py index ebc8b5f..1dcdd03 100644 --- a/client/src/pdd_u2_purchase_adapter.py +++ b/client/src/pdd_u2_purchase_adapter.py @@ -18,6 +18,7 @@ from .pdd_device_service import ( PDD_PACKAGE_NAME, PddDeviceError, PddDeviceService, + current_thread_device_service, ) from .performance_timing import current_performance_trace from .pdd_purchase_adapter import ( @@ -561,9 +562,6 @@ class U2PddPurchaseAdapter(PddPurchaseAdapter): ) from error -_PURCHASE_DEVICE_SERVICE = PddDeviceService() - - def create_u2_purchase_adapter( device_address: str, cancelled: Callable[[], bool] ) -> PddPurchaseAdapter: @@ -571,7 +569,7 @@ def create_u2_purchase_adapter( return U2PddPurchaseAdapter( device_address, - device_service=_PURCHASE_DEVICE_SERVICE, + device_service=current_thread_device_service() or PddDeviceService(), cancelled=cancelled, ) @@ -625,6 +623,6 @@ def create_u2_live_purchase_adapter( return U2PddLivePurchaseAdapter( device_address, - device_service=_PURCHASE_DEVICE_SERVICE, + device_service=current_thread_device_service() or PddDeviceService(), cancelled=cancelled, ) diff --git a/client/src/pdd_u2_purchase_reconcile_adapter.py b/client/src/pdd_u2_purchase_reconcile_adapter.py index 8aa33e3..b26f9fc 100644 --- a/client/src/pdd_u2_purchase_reconcile_adapter.py +++ b/client/src/pdd_u2_purchase_reconcile_adapter.py @@ -13,7 +13,11 @@ from datetime import datetime, timedelta, timezone from decimal import Decimal, InvalidOperation, ROUND_HALF_UP from typing import Any, Callable, Optional -from .pdd_device_service import PDD_PACKAGE_NAME, PddDeviceService +from .pdd_device_service import ( + PDD_PACKAGE_NAME, + PddDeviceService, + current_thread_device_service, +) from .pdd_purchase_reconcile_adapter import ( PddPurchaseReconcileAdapter, PurchaseOrderCandidate, @@ -215,7 +219,11 @@ class U2PddPurchaseReconcileAdapter(PddPurchaseReconcileAdapter): max_pages: int = 5, ) -> None: self._device_address = str(device_address or "").strip() - self._device_service = device_service or _RECONCILE_DEVICE_SERVICE + self._device_service = ( + device_service + or current_thread_device_service() + or PddDeviceService() + ) self._cancelled = cancelled self._settle_seconds = max(0.0, float(settle_seconds)) self._max_pages = max(1, int(max_pages)) @@ -373,9 +381,6 @@ class U2PddPurchaseReconcileAdapter(PddPurchaseReconcileAdapter): time.sleep(self._settle_seconds) -_RECONCILE_DEVICE_SERVICE = PddDeviceService() - - def create_u2_purchase_reconcile_adapter( device_address: str, cancelled: Callable[[], bool] ) -> PddPurchaseReconcileAdapter: @@ -383,6 +388,6 @@ def create_u2_purchase_reconcile_adapter( return U2PddPurchaseReconcileAdapter( device_address, - device_service=_RECONCILE_DEVICE_SERVICE, + device_service=current_thread_device_service() or PddDeviceService(), cancelled=cancelled, ) diff --git a/client/src/pdd_ui_event.py b/client/src/pdd_ui_event.py index a4c9e58..e5bb484 100644 --- a/client/src/pdd_ui_event.py +++ b/client/src/pdd_ui_event.py @@ -25,6 +25,7 @@ from typing import Callable, Dict, Optional from PyQt5.QtCore import ( QCoreApplication, QObject, + Qt, QThread, QTimer, pyqtSignal, @@ -45,6 +46,10 @@ from .collect_task_service import CollectServiceFactory, CollectTaskService from .current_client_service import CurrentClientService from .http_admin_gateway import DEFAULT_ADMIN_BASE_URL, HttpAdminGateway from .pdd_ui import PDDTaskPage, TaskRow +from .pdd_device_service import ( + PersistentPddDeviceService, + bind_thread_device_service, +) from .purchase_task_service import PurchaseAdapterFactory from .purchase_task_service import LivePurchaseAdapterFactory from .live_purchase_authorization import LivePurchaseAuthorizationService @@ -111,7 +116,7 @@ ERROR_FEEDBACK_BUTTON_MIN_HEIGHT = 40 class ClaimTaskWorker(QObject): - """在后台补交、领取或执行至多一条任务。""" + """在固定后台线程中串行补交、领取或执行至多一条任务。""" noTask = pyqtSignal() taskSaved = pyqtSignal(str) @@ -122,6 +127,7 @@ class ClaimTaskWorker(QObject): deviceUnavailable = pyqtSignal(str) outcome = pyqtSignal(str, str, str) completed = pyqtSignal() + shutdownCompleted = pyqtSignal() def __init__( self, @@ -157,6 +163,14 @@ class ClaimTaskWorker(QObject): device_connection_checker or AndroidDeviceService().require_connected ) + self._device_service = PersistentPddDeviceService() + self._idle_timer: Optional[QTimer] = None + self._accepting = True + + def prepare_run(self) -> None: + """主线程提交新命令前清除上一次的协作停止标志。""" + + self._cancelled = False def cancel(self) -> None: """阻止尚未开始的领取;已领取的任务仍必须保存到本地。""" @@ -165,62 +179,82 @@ class ClaimTaskWorker(QObject): @pyqtSlot() def run(self) -> None: - try: - if self._cancelled and not self._selected_task_id: - return - client_settings = self._client_service.load() - if not client_settings.client_id: - self.failed.emit("请先在设置页保存当前设备号和设备名") - return - android_serial = self._android_device_service.load() + """兼容旧调用;新代码通过 ``run_selected`` 提交命令。""" - client = ClientInfo( - client_settings.client_id, - client_settings.client_name, - ) - if self._selected_task_id: - self._device_connection_checker(android_serial or "") - self._task_repository.prepare_collect_rerun( - self._selected_task_id + self.run_selected(self._selected_task_id) + + @pyqtSlot(str) + def run_selected(self, selected_task_id: str = "") -> None: + """执行一条命令;该槽始终由固定设备 QThread 串行调用。""" + + if not self._accepting: + self.completed.emit() + return + self._ensure_idle_timer() + assert self._idle_timer is not None + self._idle_timer.stop() + try: + if self._cancelled and not selected_task_id: + return + with bind_thread_device_service(self._device_service): + client_settings = self._client_service.load() + if not client_settings.client_id: + self.failed.emit("请先在设置页保存当前设备号和设备名") + return + android_serial = self._android_device_service.load() + + client = ClientInfo( + client_settings.client_id, + client_settings.client_name, ) - service = CollectTaskService( - self._gateway, - self._task_repository, - client, - android_serial or "", - cancelled=lambda: self._cancelled, - collect_service_factory=self._collect_service_factory, - ) - result = service.execute_selected(self._selected_task_id) - else: - purchase_mode = "dry_run" - if self._live_authorization_service is not None: - purchase_mode = ( - self._live_authorization_service.purchase_mode_for( - client.client_id, - android_serial or "", - live_adapter_ready=( - self._live_purchase_adapter_factory is not None - ), - ) + if selected_task_id: + self._device_connection_checker(android_serial or "") + self._task_repository.prepare_collect_rerun( + selected_task_id ) - dispatcher = TaskDispatcher( - self._gateway, - self._task_repository, - client, - android_serial or "", - collect_service_factory=self._collect_service_factory, - purchase_adapter_factory=self._purchase_adapter_factory, - live_purchase_adapter_factory=( - self._live_purchase_adapter_factory - ), - purchase_mode=purchase_mode, - purchase_reconcile_factory=self._purchase_reconcile_factory, - cancelled=lambda: self._cancelled, - device_connection_checker=self._device_connection_checker, - task_saved=self.taskSaved.emit, - ) - result = dispatcher.execute_one() + service = CollectTaskService( + self._gateway, + self._task_repository, + client, + android_serial or "", + cancelled=lambda: self._cancelled, + collect_service_factory=self._collect_service_factory, + ) + result = service.execute_selected(selected_task_id) + else: + purchase_mode = "dry_run" + if self._live_authorization_service is not None: + purchase_mode = ( + self._live_authorization_service.purchase_mode_for( + client.client_id, + android_serial or "", + live_adapter_ready=( + self._live_purchase_adapter_factory + is not None + ), + ) + ) + dispatcher = TaskDispatcher( + self._gateway, + self._task_repository, + client, + android_serial or "", + collect_service_factory=self._collect_service_factory, + purchase_adapter_factory=self._purchase_adapter_factory, + live_purchase_adapter_factory=( + self._live_purchase_adapter_factory + ), + purchase_mode=purchase_mode, + purchase_reconcile_factory=( + self._purchase_reconcile_factory + ), + cancelled=lambda: self._cancelled, + device_connection_checker=( + self._device_connection_checker + ), + task_saved=self.taskSaved.emit, + ) + result = dispatcher.execute_one() if not self._cancelled or result.kind == "cancelled": self.outcome.emit(result.kind, result.message, result.task_id) except AndroidDeviceSearchError as exc: @@ -240,8 +274,38 @@ class ClaimTaskWorker(QObject): if not self._cancelled: self.failed.emit(f"执行任务失败:{exc}") finally: + if self._accepting and self._idle_timer is not None: + self._idle_timer.start( + max(1, int(self._device_service.ttl_seconds * 1000)) + ) self.completed.emit() + @pyqtSlot() + def release_device(self) -> None: + """停止、换设备或配置变化后,在设备线程安全释放缓存。""" + + self._ensure_idle_timer() + assert self._idle_timer is not None + self._idle_timer.stop() + self._device_service.release_cached() + + @pyqtSlot() + def shutdown_worker(self) -> None: + """停止接收命令并释放设备;调用方随后再 quit/wait。""" + + self._accepting = False + self._cancelled = True + try: + self.release_device() + finally: + self.shutdownCompleted.emit() + + def _ensure_idle_timer(self) -> None: + if self._idle_timer is None: + self._idle_timer = QTimer(self) + self._idle_timer.setSingleShot(True) + self._idle_timer.timeout.connect(self.release_device) + class ResultResubmitWorker(QObject): """在后台逐条重发既有 Outbox,不执行任何手机操作。""" @@ -359,6 +423,10 @@ class TaskRemoveWorker(QObject): class PDDTaskPageEvent(QObject): """把 PDD 页面只读操作连接到本地任务 Repository。""" + _claimRunRequested = pyqtSignal(str) + _claimReleaseRequested = pyqtSignal() + _claimShutdownRequested = pyqtSignal() + def __init__( self, page: PDDTaskPage, @@ -385,6 +453,7 @@ class PDDTaskPageEvent(QObject): self._claim_busy = False self._claim_thread: Optional[QThread] = None self._claim_worker: Optional[ClaimTaskWorker] = None + self._claim_operation = "" self._rerun_cancel_requested = False self._rerun_feedback: Optional[InfoBar] = None self._claim_feedback: Optional[InfoBar] = None @@ -744,7 +813,7 @@ class PDDTaskPageEvent(QObject): self._page.set_resubmit_running(False) def _start_rerun_worker(self, task_id: str) -> None: - """启动只处理指定任务的工作线程。""" + """把指定任务提交到持久设备工作线程。""" assert self._claim_gateway is not None self._close_rerun_feedback() @@ -758,29 +827,10 @@ class PDDTaskPageEvent(QObject): ) self._reload() - thread = QThread(self) - worker = ClaimTaskWorker( - self._claim_gateway, - self._repository, - self._client_service, - self._selected_android_device_service, - self._collect_service_factory, - selected_task_id=task_id, - device_connection_checker=self._device_connection_checker, - ) - worker.moveToThread(thread) - thread.started.connect(worker.run) - worker.retryableFailed.connect(self._on_rerun_failed) - worker.failed.connect(self._on_rerun_failed) - worker.deviceUnavailable.connect(self._on_rerun_device_unavailable) - worker.outcome.connect(self._on_rerun_outcome) - worker.completed.connect(thread.quit) - worker.completed.connect(worker.deleteLater) - thread.finished.connect(thread.deleteLater) - thread.finished.connect(self._on_rerun_thread_finished) - self._claim_thread = thread - self._claim_worker = worker - thread.start() + worker = self._ensure_claim_executor() + self._claim_operation = "rerun" + worker.prepare_run() + self._claimRunRequested.emit(task_id) @pyqtSlot() def request_cancel_rerun(self) -> None: @@ -839,8 +889,6 @@ class PDDTaskPageEvent(QObject): @pyqtSlot() def _on_rerun_thread_finished(self) -> None: cancel_requested = self._rerun_cancel_requested - self._claim_worker = None - self._claim_thread = None self._claim_busy = False self._rerun_cancel_requested = False if not self._closing: @@ -939,7 +987,7 @@ class PDDTaskPageEvent(QObject): @pyqtSlot() def _start_claim_cycle(self) -> None: - """开始一轮任务处理;只允许上一轮线程完全结束后进入。""" + """开始一轮任务处理;设备命令由固定工作线程串行执行。""" if ( self._closing @@ -953,6 +1001,20 @@ class PDDTaskPageEvent(QObject): self._page.set_auto_fetch_state("running") self._page.set_engine_status("自动获取:运行中 · 正在处理一条任务…") + worker = self._ensure_claim_executor() + self._claim_operation = "auto" + worker.prepare_run() + self._claimRunRequested.emit("") + + def _ensure_claim_executor(self) -> ClaimTaskWorker: + """懒创建唯一的设备 Worker 和固定 QThread。""" + + current = self._claim_worker + thread = self._claim_thread + if current is not None and thread is not None and thread.isRunning(): + return current + + assert self._claim_gateway is not None thread = QThread(self) worker = ClaimTaskWorker( gateway=self._claim_gateway, @@ -967,22 +1029,63 @@ class PDDTaskPageEvent(QObject): device_connection_checker=self._device_connection_checker, ) worker.moveToThread(thread) - thread.started.connect(worker.run) - worker.noTask.connect(self._on_no_claimed_task) + self._claimRunRequested.connect(worker.run_selected) + self._claimReleaseRequested.connect(worker.release_device) + self._claimShutdownRequested.connect(worker.shutdown_worker) worker.taskSaved.connect(self._on_claimed_task_saved) - worker.duplicateTask.connect(self._on_duplicate_claimed_task) - worker.localSaveFailed.connect(self._on_claimed_task_save_failed) - worker.retryableFailed.connect(self._on_claim_retryable_failed) - worker.failed.connect(self._on_claim_failed) - worker.deviceUnavailable.connect(self._on_claim_device_unavailable) - worker.outcome.connect(self._on_collect_outcome) - worker.completed.connect(thread.quit) - worker.completed.connect(worker.deleteLater) + worker.retryableFailed.connect(self._route_claim_retryable_failed) + worker.failed.connect(self._route_claim_failed) + worker.deviceUnavailable.connect(self._route_device_unavailable) + worker.outcome.connect(self._route_claim_outcome) + worker.completed.connect(self._on_claim_cycle_completed) + # QThread 对象属于主线程;显式直连可在 shutdown_worker 完成释放后 + # 立即停止工作线程事件循环,避免主线程 wait() 时互相等待。 + worker.shutdownCompleted.connect(thread.quit, Qt.DirectConnection) + thread.finished.connect(worker.deleteLater) thread.finished.connect(thread.deleteLater) - thread.finished.connect(self._on_claim_thread_finished) self._claim_thread = thread self._claim_worker = worker thread.start() + return worker + + @pyqtSlot(str) + def _route_claim_retryable_failed(self, message: str) -> None: + if self._claim_operation == "rerun": + self._on_rerun_failed(message) + else: + self._on_claim_retryable_failed(message) + + @pyqtSlot(str) + def _route_claim_failed(self, message: str) -> None: + if self._claim_operation == "rerun": + self._on_rerun_failed(message) + else: + self._on_claim_failed(message) + + @pyqtSlot(str) + def _route_device_unavailable(self, message: str) -> None: + if self._claim_operation == "rerun": + self._on_rerun_device_unavailable(message) + else: + self._on_claim_device_unavailable(message) + + @pyqtSlot(str, str, str) + def _route_claim_outcome( + self, kind: str, message: str, task_id: str + ) -> None: + if self._claim_operation == "rerun": + self._on_rerun_outcome(kind, message, task_id) + else: + self._on_collect_outcome(kind, message, task_id) + + @pyqtSlot() + def _on_claim_cycle_completed(self) -> None: + operation = self._claim_operation + self._claim_operation = "" + if operation == "rerun": + self._on_rerun_thread_finished() + else: + self._on_claim_thread_finished() def _request_stop_auto_fetch(self) -> None: """取消下一轮,并请求当前 Worker 在安全点停止。""" @@ -992,7 +1095,7 @@ class PDDTaskPageEvent(QObject): self._next_cycle_timer.stop() self._stop_status = "自动获取:已停止 · 当前没有执行中的任务" worker = self._claim_worker - if worker is not None: + if worker is not None and self._claim_busy: self._page.set_auto_fetch_state("stopping") self._page.set_engine_status( "自动获取:正在停止 · 等待当前任务安全结束" @@ -1010,6 +1113,8 @@ class PDDTaskPageEvent(QObject): self._page.set_auto_fetch_state("stopped") self._page.set_engine_status(self._stop_status) self._stop_requested = False + if self._claim_worker is not None: + self._claimReleaseRequested.emit() def _continue_after(self, delay_ms: int, status: str) -> None: """记录本轮结束后的下一次调度,不在 Worker 结束前启动。""" @@ -1266,8 +1371,6 @@ class PDDTaskPageEvent(QObject): @pyqtSlot() def _on_claim_thread_finished(self) -> None: - self._claim_worker = None - self._claim_thread = None self._claim_busy = False if self._closing: return @@ -1349,6 +1452,13 @@ class PDDTaskPageEvent(QObject): if not self._closing: self._page.taskTable.setFocus() + @pyqtSlot(str) + def android_device_configuration_changed(self, _serial: str = "") -> None: + """设置页保存、删除或切换设备后释放旧的持久连接。""" + + if not self._closing and self._claim_worker is not None: + self._claimReleaseRequested.emit() + @pyqtSlot() def shutdown(self) -> None: """窗口关闭时停止新的领取,并等待已领取任务完成本地保存。""" @@ -1372,43 +1482,20 @@ class PDDTaskPageEvent(QObject): if worker is not None: try: worker.cancel() - signal_slots = ( - (worker.noTask, self._on_no_claimed_task), - (worker.taskSaved, self._on_claimed_task_saved), - (worker.duplicateTask, self._on_duplicate_claimed_task), - ( - worker.localSaveFailed, - self._on_claimed_task_save_failed, - ), - (worker.failed, self._on_claim_failed), - ( - worker.deviceUnavailable, - self._on_claim_device_unavailable, - ), - ( - worker.deviceUnavailable, - self._on_rerun_device_unavailable, - ), - ( - worker.retryableFailed, - self._on_claim_retryable_failed, - ), - (worker.outcome, self._on_collect_outcome), - ) except RuntimeError: - signal_slots = () - for signal, slot in signal_slots: - try: - signal.disconnect(slot) - except (TypeError, RuntimeError): - pass + pass + try: + self._claimShutdownRequested.emit() + except RuntimeError: + pass if thread is not None and thread.isRunning(): - thread.quit() # uiautomator2/ADB 的单次调用可能需要数秒才返回。先通过 # cancelled 标志让采集在下一个安全点退出,再等待工作线程收尾, # 避免窗口销毁时出现 "QThread destroyed while running"。 thread.wait(60_000) + self._claim_worker = None + self._claim_thread = None resubmit_worker = self._resubmit_worker resubmit_thread = self._resubmit_thread diff --git a/client/src/settings_ui_event.py b/client/src/settings_ui_event.py index f37ba23..dc86c68 100644 --- a/client/src/settings_ui_event.py +++ b/client/src/settings_ui_event.py @@ -388,6 +388,7 @@ class SettingsPageEventBinder(QObject): saveRequested = pyqtSignal(str, str) deleteRequested = pyqtSignal(str) liveAuthorizationRequested = pyqtSignal(str, str) + androidDeviceConfigurationChanged = pyqtSignal(str) def __init__( self, @@ -1355,12 +1356,14 @@ class SettingsPageEventBinder(QObject): self._page.deviceStatusLabel.setText( f"当前使用设备:{serial}(已保存)" ) + self.androidDeviceConfigurationChanged.emit(serial) else: self._saved_android_serial = "" self._page.deviceTableModel.set_checked_serial("") self._page.deviceStatusLabel.setText( "当前使用设备:未选择(本地配置已删除)" ) + self.androidDeviceConfigurationChanged.emit("") @pyqtSlot(str, str) def _on_android_device_setting_failed( diff --git a/client/src/ui_main.py b/client/src/ui_main.py index 19f0422..a8058fe 100644 --- a/client/src/ui_main.py +++ b/client/src/ui_main.py @@ -71,6 +71,9 @@ class MainWindow(FluentWindow): self.pddTaskPage.openSettingsRequested.connect( lambda: self.switchTo(self.settingsPage) ) + self.settingsPage.eventBinder.androidDeviceConfigurationChanged.connect( + self.pddTaskPageEvent.android_device_configuration_changed + ) self.addSubInterface(self.pddTaskPage, FIF.HOME, "pdd") self.addSubInterface( self.settingsPage, diff --git a/client/test/test_pdd_device_service.py b/client/test/test_pdd_device_service.py index fe29340..f3f431b 100644 --- a/client/test/test_pdd_device_service.py +++ b/client/test/test_pdd_device_service.py @@ -3,7 +3,13 @@ import threading import unittest -from src.pdd_device_service import PddDeviceError, PddDeviceService +from src.pdd_device_service import ( + PddDeviceError, + PddDeviceService, + PersistentPddDeviceService, + bind_thread_device_service, + current_thread_device_service, +) from src.performance_timing import TaskPerformanceTrace @@ -89,5 +95,86 @@ class PddDeviceServiceTest(unittest.TestCase): self.assertEqual(errors, ["DEVICE_THREAD_VIOLATION"]) +class PersistentPddDeviceServiceTest(unittest.TestCase): + def test_same_serial_within_ttl_connects_once_but_checks_each_task(self): + connected = [] + device = FakeDevice() + service = PersistentPddDeviceService( + lambda serial: connected.append(serial) or device, + ttl_seconds=90, + ) + + with service.connect("USB-001"): + pass + with service.connect("USB-001"): + pass + + self.assertEqual(connected, ["USB-001"]) + self.assertEqual(device.current_calls, 2) + + def test_expired_or_changed_serial_discards_cached_connection(self): + now = [10.0] + connected = [] + service = PersistentPddDeviceService( + lambda serial: connected.append(serial) or FakeDevice(), + ttl_seconds=5, + monotonic=lambda: now[0], + ) + + with service.connect("USB-001"): + pass + now[0] = 16.0 + with service.connect("USB-001"): + pass + with service.connect("USB-002"): + pass + + self.assertEqual(connected, ["USB-001", "USB-001", "USB-002"]) + + def test_failed_cached_health_check_reconnects_only_once(self): + class UnhealthyDevice(FakeDevice): + def app_current(self): + self.current_calls += 1 + if self.current_calls >= 2: + raise OSError("offline") + return {"package": "com.xunmeng.pinduoduo"} + + devices = [UnhealthyDevice(), FakeDevice()] + connected = [] + service = PersistentPddDeviceService( + lambda serial: connected.append(serial) or devices.pop(0) + ) + + with service.connect("USB-001"): + pass + with service.connect("USB-001") as current: + self.assertIsInstance(current, FakeDevice) + + self.assertEqual(connected, ["USB-001", "USB-001"]) + + def test_release_and_context_binding_stay_on_owner_thread(self): + service = PersistentPddDeviceService(lambda _serial: FakeDevice()) + with service.connect("USB-001"): + pass + with bind_thread_device_service(service): + self.assertIs(current_thread_device_service(), service) + self.assertIsNone(current_thread_device_service()) + + errors = [] + + def release_from_other_thread(): + try: + service.release_cached() + except PddDeviceError as exc: + errors.append(exc.code) + + worker = threading.Thread(target=release_from_other_thread) + worker.start() + worker.join() + self.assertEqual(errors, ["DEVICE_THREAD_VIOLATION"]) + service.release_cached() + self.assertFalse(service.has_cached_device) + + if __name__ == "__main__": unittest.main() diff --git a/client/test/test_pdd_ui_event.py b/client/test/test_pdd_ui_event.py index 0b36ac5..f0102a7 100644 --- a/client/test/test_pdd_ui_event.py +++ b/client/test/test_pdd_ui_event.py @@ -1153,6 +1153,36 @@ class PDDTaskPageEventTest(unittest.TestCase): events.shutdown() page.deleteLater() + def test_auto_fetch_reuses_one_fixed_worker_thread_across_cycles(self): + gateway = RecordingClaimGateway() + page = PDDTaskPage() + events = PDDTaskPageEvent( + page, + self.repository, + claim_gateway=gateway, + settings_repository=self._saved_settings(), + collect_service_factory=fake_collect_factory, + no_task_delay_ms=10, + ) + + page.autoFetchRequested.emit() + self.assertTrue( + wait_until(self.app, lambda: len(gateway.thread_ids) >= 2), + "自动获取没有完成两轮串行调度", + ) + executor_thread = events._claim_thread + page.autoFetchRequested.emit() + self.assertTrue(wait_until(self.app, lambda: not events._claim_busy)) + + self.assertIs(events._claim_thread, executor_thread) + self.assertTrue( + executor_thread is not None and executor_thread.isRunning() + ) + self.assertEqual(len(set(gateway.thread_ids)), 1) + events.shutdown() + self.assertFalse(executor_thread.isRunning()) + page.deleteLater() + def test_new_claim_refreshes_table_before_collection_finishes(self): gateway = RecordingClaimGateway(collect_admin_task("COL-EARLY")) collect_started = threading.Event() diff --git a/client/test/test_settings_ui_event.py b/client/test/test_settings_ui_event.py index 7a4dd83..ed050fa 100644 --- a/client/test/test_settings_ui_event.py +++ b/client/test/test_settings_ui_event.py @@ -513,6 +513,10 @@ class SettingsPageEventTest(unittest.TestCase): ) page.deviceTableModel.set_checked_serial("USB-001") self._wait_until(lambda: page.eventBinder._pdd_check_thread is None) + changed_serials = [] + page.eventBinder.androidDeviceConfigurationChanged.connect( + changed_serials.append + ) page.saveButton.click() self.assertFalse(page.saveButton.isEnabled()) @@ -530,6 +534,7 @@ class SettingsPageEventTest(unittest.TestCase): "当前使用设备:USB-001(已保存)", ) self.assertTrue(page.deleteButton.isEnabled()) + self.assertEqual(changed_serials, ["USB-001"]) page.eventBinder.shutdown() page.deleteLater() diff --git a/docs/client/02-architecture.md b/docs/client/02-architecture.md index b187864..1d5fa9e 100644 --- a/docs/client/02-architecture.md +++ b/docs/client/02-architecture.md @@ -156,11 +156,10 @@ Qt 主线程 ├── 窗口、页面、表格模型和用户事件 └── 接收后台信号并更新界面 -任务工作线程 -├── Admin 请求 -├── uiautomator2 调用 +持久设备工作线程(单设备、固定 QThread) +├── uiautomator2 Device 的创建、健康检查、调用和释放 ├── 页面等待与 XML 解析 -└── 单个任务的串行执行 +└── 采集、采购、订单核对和手动重新执行的串行命令 结果提交工作线程或同一任务线程的独立队列 └── Outbox 重试,不重复执行 PDD 操作 @@ -172,6 +171,12 @@ Qt 主线程 - `[必须]` QWidget 只能在 Qt 主线程创建和访问。 - `[必须]` 一个 Android 设备由一个工作线程独占,不跨线程共享 uiautomator2 Device 对象。 +- `[必须]` 同一设备在 90 秒空闲期内只复用 Device 连接;每个任务必须重新读取 + 当前应用、页面、控件树、商品、价格、规格和采购安全状态。 +- `[必须]` 停止自动获取、切换或删除已保存设备、设备断连、配置变化、窗口关闭 + 和空闲超时都会在设备线程释放连接。关闭只使用协作停止和 `quit()/wait()`。 +- `[必须]` `task_runs.irreversible_action_at` 有值后只允许把核单命令放入队列, + 不得重新执行采购下单步骤。 - `[必须]` 后台信号只传递不可变数据、稳定编号或轻量视图模型。 - `[必须]` 点“停止获取”后不再领取新任务,当前任务在定义的安全点退出。 - `[建议]` 关闭窗口时应选择停止、等待或后台继续;MVP 默认安全停止并持久化状态。 @@ -259,6 +264,12 @@ def start_task(self, remote_task_id: str) -> None: | 用 `thread.terminate()` 停任务 | 数据库写一半、订单状态不明 | 用 `cancel()` 置标志位,在安全点退出 | | 在 `run()` 里改界面 | 随机崩溃,且很难复现 | 只 emit 信号 | +PDD 设备执行器是这个模板的长生命周期版本:QThread 在第一次设备任务时创建, +后续通过队列信号提交一条命令,命令完成后线程继续等待;不能在 Worker 内写无限 +轮询。设备 Worker 仍不访问界面。停止接单后等待当前命令到安全点结束,再在线程内 +释放 Device,最后执行 `quit()/wait()`。Outbox 批量重新上报不操作手机,继续使用 +独立的网络 Worker,不进入设备命令队列。 + ### 5.2 后台结果回主线程 / 迟到结果 窗口关掉了,后台任务还在跑,跑完再发信号——这时槽函数去访问已经销毁的控件,程序就崩了。 @@ -403,8 +414,9 @@ reconcile_purchase(task, run) -> PurchaseResult | ManualReview 或电话。最终匹配由领域服务完成,不能让页面解析层单独决定采购成功。 正式 Client 由 `pdd_u2_purchase_adapter.py` 分别实现演练和 live 接口,并在 -`ui_main.py` 注入工厂。Adapter 通过 `PddDeviceService` 独占连接,每次判断 -都重新读取当前包名和控件树。当前 Admin 下发的 `color` 和 `size` +`ui_main.py` 注入工厂。Adapter 通过设备工作线程绑定的 +`PersistentPddDeviceService` 独占一次任务会话;同一设备可在 90 秒内复用底层 +Device 连接,但每次判断都重新读取当前包名和控件树。当前 Admin 下发的 `color` 和 `size` 使用精确文字匹配;其他动态维度直接停止,不做相似匹配。原生 PDD 控件树不暴露商品编号,因此商品编号来自 Adapter 本次已校验并打开的 PDD URL,包名和页面类型仍以最新控件树确认。