refactor: 复用持久设备工作线程 (#110)

This commit is contained in:
chengma
2026-08-10 18:00:38 +08:00
parent 553659208a
commit 70cedc34d4
11 changed files with 563 additions and 149 deletions
+212 -125
View File
@@ -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