2026-08-06 12:01:38 +08:00
|
|
|
|
"""PDD 任务页的事件绑定。
|
|
|
|
|
|
|
|
|
|
|
|
把 `pdd_ui.py` 里控件的信号,接到查询、任务引擎等实际动作上。
|
|
|
|
|
|
|
|
|
|
|
|
改动本文件前必读 `client/AGENTS.md`。以下几条最容易踩:
|
|
|
|
|
|
|
|
|
|
|
|
- **禁止在 Qt 主线程做任何阻塞的事**:Admin 请求、uiautomator2/ADB 调用、
|
|
|
|
|
|
`time.sleep()`、轮询、大 XML 解析、批量写文件。做了界面就会卡死转圈。
|
|
|
|
|
|
- 长任务统一用 `QObject` + `moveToThread` 的 Worker 写法,
|
|
|
|
|
|
模板照抄 `docs/client/02-architecture.md` §5.1。
|
|
|
|
|
|
**不要**用 `QThread` 子类、`QRunnable` 或 Python 的 `threading`。
|
|
|
|
|
|
- 后台结果只能通过**信号**回到主线程,Worker 里一行界面代码都不许有。
|
|
|
|
|
|
- 窗口关闭时要断开信号并置标志位,否则迟到的后台结果会访问
|
|
|
|
|
|
已经销毁的控件、直接崩溃。做法见同文档 §5.2。
|
|
|
|
|
|
- 数据库读写走 Repository,**不要在这里拼业务 SQL**。
|
2026-08-07 17:38:29 +08:00
|
|
|
|
- “获取任务”会真的去操作手机采集商品;“搜索”只读本地数据库。
|
2026-08-06 12:01:38 +08:00
|
|
|
|
两者必须分开,不得共用入口。
|
|
|
|
|
|
- 普通成功不弹窗,更新界面即可;可恢复错误用 `InfoBar`
|
|
|
|
|
|
(模板见 `docs/client/05-ui-specification.md` §9.1);
|
|
|
|
|
|
只有必须让用户当场做决定时才用模态对话框。
|
|
|
|
|
|
"""
|
2026-08-06 16:31:34 +08:00
|
|
|
|
|
2026-08-10 09:35:46 +08:00
|
|
|
|
from typing import Callable, Dict, Optional
|
2026-08-06 16:31:34 +08:00
|
|
|
|
|
2026-08-09 10:54:18 +08:00
|
|
|
|
from PyQt5.QtCore import (
|
|
|
|
|
|
QCoreApplication,
|
|
|
|
|
|
QObject,
|
2026-08-10 18:00:38 +08:00
|
|
|
|
Qt,
|
2026-08-09 10:54:18 +08:00
|
|
|
|
QThread,
|
|
|
|
|
|
QTimer,
|
|
|
|
|
|
pyqtSignal,
|
|
|
|
|
|
pyqtSlot,
|
|
|
|
|
|
)
|
2026-08-10 09:08:54 +08:00
|
|
|
|
from qfluentwidgets import InfoBar, InfoBarPosition, MessageBox, PushButton
|
2026-08-06 16:31:34 +08:00
|
|
|
|
|
2026-08-10 09:35:46 +08:00
|
|
|
|
from .android_device_service import (
|
|
|
|
|
|
AndroidDeviceSearchError,
|
|
|
|
|
|
AndroidDeviceService,
|
|
|
|
|
|
)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
from .admin_gateway import (
|
|
|
|
|
|
AdminGatewayError,
|
|
|
|
|
|
ClientInfo,
|
2026-08-07 17:38:29 +08:00
|
|
|
|
AdminGateway,
|
2026-08-07 16:42:31 +08:00
|
|
|
|
)
|
2026-08-07 17:38:29 +08:00
|
|
|
|
from .collect_task_service import CollectServiceFactory, CollectTaskService
|
2026-08-07 16:42:31 +08:00
|
|
|
|
from .current_client_service import CurrentClientService
|
|
|
|
|
|
from .http_admin_gateway import DEFAULT_ADMIN_BASE_URL, HttpAdminGateway
|
2026-08-06 16:31:34 +08:00
|
|
|
|
from .pdd_ui import PDDTaskPage, TaskRow
|
2026-08-10 18:00:38 +08:00
|
|
|
|
from .pdd_device_service import (
|
|
|
|
|
|
PersistentPddDeviceService,
|
|
|
|
|
|
bind_thread_device_service,
|
|
|
|
|
|
)
|
2026-08-10 00:05:40 +08:00
|
|
|
|
from .purchase_task_service import PurchaseAdapterFactory
|
2026-08-10 16:35:38 +08:00
|
|
|
|
from .purchase_task_service import LivePurchaseAdapterFactory
|
2026-08-10 00:27:52 +08:00
|
|
|
|
from .purchase_reconcile_service import PurchaseReconcileFactory
|
2026-08-07 16:42:31 +08:00
|
|
|
|
from .selected_android_device_service import SelectedAndroidDeviceService
|
|
|
|
|
|
from .settings_repository import SettingsRepository
|
|
|
|
|
|
from .task_models import (
|
2026-08-10 11:18:06 +08:00
|
|
|
|
OutboxEventType,
|
2026-08-10 10:56:48 +08:00
|
|
|
|
OutboxStatus,
|
2026-08-07 16:42:31 +08:00
|
|
|
|
TaskFilters,
|
|
|
|
|
|
TaskStatus,
|
|
|
|
|
|
TaskSummary,
|
|
|
|
|
|
TaskType,
|
|
|
|
|
|
)
|
2026-08-10 17:23:24 +08:00
|
|
|
|
from .task_repository import (
|
|
|
|
|
|
CollectRerunError,
|
|
|
|
|
|
TaskRemovalError,
|
|
|
|
|
|
TaskRepository,
|
|
|
|
|
|
)
|
2026-08-10 00:05:40 +08:00
|
|
|
|
from .task_dispatcher import (
|
|
|
|
|
|
TaskDispatcher,
|
|
|
|
|
|
admin_task_to_new_claimed_task,
|
|
|
|
|
|
)
|
2026-08-08 09:46:58 +08:00
|
|
|
|
from .task_detail_view import TaskDetailWindow
|
2026-08-06 16:31:34 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
TASK_TYPE_BY_TEXT = {
|
|
|
|
|
|
"采集": TaskType.COLLECT,
|
|
|
|
|
|
"采购": TaskType.PURCHASE,
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
TASK_STATUS_BY_TEXT = {
|
|
|
|
|
|
"待执行": TaskStatus.CLAIMED,
|
|
|
|
|
|
"执行中": TaskStatus.RUNNING,
|
|
|
|
|
|
"结果待提交": TaskStatus.RESULT_PENDING,
|
|
|
|
|
|
"等待重试": TaskStatus.RETRY_WAIT,
|
|
|
|
|
|
"需要人工处理": TaskStatus.MANUAL_REVIEW,
|
|
|
|
|
|
"已完成": TaskStatus.SUCCEEDED,
|
|
|
|
|
|
"失败": TaskStatus.FAILED,
|
|
|
|
|
|
"已取消": TaskStatus.CANCELLED,
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
TASK_TYPE_TEXT = {
|
|
|
|
|
|
TaskType.COLLECT: "采集",
|
|
|
|
|
|
TaskType.PURCHASE: "采购",
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
TASK_STATUS_TEXT = {
|
|
|
|
|
|
TaskStatus.CLAIMED: "待执行",
|
|
|
|
|
|
TaskStatus.RUNNING: "执行中",
|
|
|
|
|
|
TaskStatus.RESULT_PENDING: "结果待提交",
|
|
|
|
|
|
TaskStatus.RETRY_WAIT: "等待重试",
|
|
|
|
|
|
TaskStatus.MANUAL_REVIEW: "需要人工处理",
|
|
|
|
|
|
TaskStatus.SUCCEEDED: "已完成",
|
|
|
|
|
|
TaskStatus.FAILED: "失败",
|
|
|
|
|
|
TaskStatus.CANCELLED: "已取消",
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-10 09:59:34 +08:00
|
|
|
|
ERROR_FEEDBACK_DURATION_MS = 5_000
|
|
|
|
|
|
ERROR_FEEDBACK_MIN_WIDTH = 520
|
|
|
|
|
|
ERROR_FEEDBACK_MIN_HEIGHT = 112
|
|
|
|
|
|
ERROR_FEEDBACK_BUTTON_MIN_WIDTH = 120
|
|
|
|
|
|
ERROR_FEEDBACK_BUTTON_MIN_HEIGHT = 40
|
2026-08-10 19:18:26 +08:00
|
|
|
|
PURCHASE_RECONCILE_DELAY_MS = 5_000
|
2026-08-10 09:59:34 +08:00
|
|
|
|
|
2026-08-06 16:31:34 +08:00
|
|
|
|
|
2026-08-07 16:42:31 +08:00
|
|
|
|
class ClaimTaskWorker(QObject):
|
2026-08-10 18:00:38 +08:00
|
|
|
|
"""在固定后台线程中串行补交、领取或执行至多一条任务。"""
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
|
|
|
|
|
noTask = pyqtSignal()
|
|
|
|
|
|
taskSaved = pyqtSignal(str)
|
|
|
|
|
|
duplicateTask = pyqtSignal(str)
|
|
|
|
|
|
localSaveFailed = pyqtSignal(str, str)
|
2026-08-09 10:54:18 +08:00
|
|
|
|
retryableFailed = pyqtSignal(str)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
failed = pyqtSignal(str)
|
2026-08-10 09:35:46 +08:00
|
|
|
|
deviceUnavailable = pyqtSignal(str)
|
2026-08-07 17:38:29 +08:00
|
|
|
|
outcome = pyqtSignal(str, str, str)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
completed = pyqtSignal()
|
2026-08-10 18:00:38 +08:00
|
|
|
|
shutdownCompleted = pyqtSignal()
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
|
|
|
|
|
def __init__(
|
|
|
|
|
|
self,
|
2026-08-07 17:38:29 +08:00
|
|
|
|
gateway: AdminGateway,
|
2026-08-07 16:42:31 +08:00
|
|
|
|
task_repository: TaskRepository,
|
|
|
|
|
|
client_service: CurrentClientService,
|
|
|
|
|
|
android_device_service: SelectedAndroidDeviceService,
|
2026-08-07 17:38:29 +08:00
|
|
|
|
collect_service_factory: Optional[CollectServiceFactory] = None,
|
2026-08-10 00:05:40 +08:00
|
|
|
|
purchase_adapter_factory: Optional[PurchaseAdapterFactory] = None,
|
2026-08-10 16:35:38 +08:00
|
|
|
|
live_purchase_adapter_factory: Optional[
|
|
|
|
|
|
LivePurchaseAdapterFactory
|
|
|
|
|
|
] = None,
|
2026-08-10 00:27:52 +08:00
|
|
|
|
purchase_reconcile_factory: Optional[PurchaseReconcileFactory] = None,
|
2026-08-09 21:44:29 +08:00
|
|
|
|
selected_task_id: str = "",
|
2026-08-10 09:35:46 +08:00
|
|
|
|
device_connection_checker: Optional[Callable[[str], None]] = None,
|
2026-08-07 16:42:31 +08:00
|
|
|
|
) -> None:
|
|
|
|
|
|
super().__init__()
|
|
|
|
|
|
self._gateway = gateway
|
|
|
|
|
|
self._task_repository = task_repository
|
|
|
|
|
|
self._client_service = client_service
|
|
|
|
|
|
self._android_device_service = android_device_service
|
|
|
|
|
|
self._cancelled = False
|
2026-08-07 17:38:29 +08:00
|
|
|
|
self._collect_service_factory = collect_service_factory
|
2026-08-10 00:05:40 +08:00
|
|
|
|
self._purchase_adapter_factory = purchase_adapter_factory
|
2026-08-10 16:35:38 +08:00
|
|
|
|
self._live_purchase_adapter_factory = live_purchase_adapter_factory
|
2026-08-10 00:27:52 +08:00
|
|
|
|
self._purchase_reconcile_factory = purchase_reconcile_factory
|
2026-08-09 21:44:29 +08:00
|
|
|
|
self._selected_task_id = selected_task_id
|
2026-08-10 09:35:46 +08:00
|
|
|
|
self._device_connection_checker = (
|
|
|
|
|
|
device_connection_checker
|
|
|
|
|
|
or AndroidDeviceService().require_connected
|
|
|
|
|
|
)
|
2026-08-10 18:00:38 +08:00
|
|
|
|
self._device_service = PersistentPddDeviceService()
|
|
|
|
|
|
self._idle_timer: Optional[QTimer] = None
|
|
|
|
|
|
self._accepting = True
|
|
|
|
|
|
|
|
|
|
|
|
def prepare_run(self) -> None:
|
|
|
|
|
|
"""主线程提交新命令前清除上一次的协作停止标志。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._cancelled = False
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
|
|
|
|
|
def cancel(self) -> None:
|
|
|
|
|
|
"""阻止尚未开始的领取;已领取的任务仍必须保存到本地。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._cancelled = True
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def run(self) -> None:
|
2026-08-10 18:00:38 +08:00
|
|
|
|
"""兼容旧调用;新代码通过 ``run_selected`` 提交命令。"""
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
2026-08-10 18:00:38 +08:00
|
|
|
|
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,
|
2026-08-10 09:35:46 +08:00
|
|
|
|
)
|
2026-08-10 18:00:38 +08:00
|
|
|
|
if selected_task_id:
|
|
|
|
|
|
self._device_connection_checker(android_serial or "")
|
|
|
|
|
|
self._task_repository.prepare_collect_rerun(
|
|
|
|
|
|
selected_task_id
|
2026-08-10 16:35:38 +08:00
|
|
|
|
)
|
2026-08-10 18:00:38 +08:00
|
|
|
|
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:
|
2026-08-10 18:37:35 +08:00
|
|
|
|
purchase_mode = (
|
|
|
|
|
|
"live"
|
|
|
|
|
|
if android_serial
|
|
|
|
|
|
and self._live_purchase_adapter_factory is not None
|
|
|
|
|
|
else "dry_run"
|
|
|
|
|
|
)
|
2026-08-10 18:00:38 +08:00
|
|
|
|
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()
|
2026-08-07 17:38:29 +08:00
|
|
|
|
if not self._cancelled or result.kind == "cancelled":
|
|
|
|
|
|
self.outcome.emit(result.kind, result.message, result.task_id)
|
2026-08-10 09:35:46 +08:00
|
|
|
|
except AndroidDeviceSearchError as exc:
|
|
|
|
|
|
if not self._cancelled:
|
|
|
|
|
|
self.deviceUnavailable.emit(str(exc))
|
2026-08-07 16:42:31 +08:00
|
|
|
|
except AdminGatewayError as exc:
|
|
|
|
|
|
if not self._cancelled:
|
|
|
|
|
|
request_hint = (
|
|
|
|
|
|
f",请求编号:{exc.request_id}" if exc.request_id else ""
|
|
|
|
|
|
)
|
2026-08-09 10:54:18 +08:00
|
|
|
|
message = f"{exc}{request_hint}"
|
|
|
|
|
|
if exc.retryable:
|
|
|
|
|
|
self.retryableFailed.emit(message)
|
|
|
|
|
|
else:
|
|
|
|
|
|
self.failed.emit(message)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
except Exception as exc:
|
|
|
|
|
|
if not self._cancelled:
|
2026-08-10 00:05:40 +08:00
|
|
|
|
self.failed.emit(f"执行任务失败:{exc}")
|
2026-08-07 16:42:31 +08:00
|
|
|
|
finally:
|
2026-08-10 18:00:38 +08:00
|
|
|
|
if self._accepting and self._idle_timer is not None:
|
|
|
|
|
|
self._idle_timer.start(
|
|
|
|
|
|
max(1, int(self._device_service.ttl_seconds * 1000))
|
|
|
|
|
|
)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
self.completed.emit()
|
|
|
|
|
|
|
2026-08-10 18:00:38 +08:00
|
|
|
|
@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)
|
|
|
|
|
|
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
2026-08-10 10:56:48 +08:00
|
|
|
|
class ResultResubmitWorker(QObject):
|
2026-08-10 11:18:06 +08:00
|
|
|
|
"""在后台逐条重发既有 Outbox,不执行任何手机操作。"""
|
2026-08-10 10:56:48 +08:00
|
|
|
|
|
|
|
|
|
|
progress = pyqtSignal(int, int, str)
|
|
|
|
|
|
finished = pyqtSignal(object, object, object)
|
|
|
|
|
|
completed = pyqtSignal()
|
|
|
|
|
|
|
|
|
|
|
|
def __init__(
|
|
|
|
|
|
self,
|
|
|
|
|
|
gateway: AdminGateway,
|
|
|
|
|
|
repository: TaskRepository,
|
|
|
|
|
|
task_ids: tuple[str, ...],
|
|
|
|
|
|
) -> None:
|
|
|
|
|
|
super().__init__()
|
|
|
|
|
|
self._gateway = gateway
|
|
|
|
|
|
self._repository = repository
|
|
|
|
|
|
self._task_ids = task_ids
|
|
|
|
|
|
self._cancelled = False
|
|
|
|
|
|
|
|
|
|
|
|
def cancel(self) -> None:
|
|
|
|
|
|
"""当前网络请求结束后停止处理后续任务。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._cancelled = True
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def run(self) -> None:
|
|
|
|
|
|
succeeded: list[str] = []
|
|
|
|
|
|
failed: list[str] = []
|
|
|
|
|
|
skipped: list[str] = []
|
|
|
|
|
|
total = len(self._task_ids)
|
|
|
|
|
|
try:
|
|
|
|
|
|
for current, task_id in enumerate(self._task_ids, 1):
|
|
|
|
|
|
if self._cancelled:
|
|
|
|
|
|
break
|
|
|
|
|
|
self.progress.emit(current, total, task_id)
|
|
|
|
|
|
try:
|
2026-08-10 11:18:06 +08:00
|
|
|
|
event = self._repository.outbox_for_resubmit(task_id)
|
2026-08-10 10:56:48 +08:00
|
|
|
|
except Exception:
|
|
|
|
|
|
skipped.append(task_id)
|
|
|
|
|
|
continue
|
|
|
|
|
|
if event is None or event.status is OutboxStatus.SENDING:
|
|
|
|
|
|
skipped.append(task_id)
|
|
|
|
|
|
continue
|
|
|
|
|
|
|
|
|
|
|
|
self._repository.mark_outbox_sending(event.id)
|
|
|
|
|
|
try:
|
2026-08-10 11:18:06 +08:00
|
|
|
|
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,
|
|
|
|
|
|
)
|
2026-08-10 10:56:48 +08:00
|
|
|
|
if not receipt.accepted:
|
|
|
|
|
|
raise AdminGatewayError(
|
|
|
|
|
|
"ADMIN_RESULT_NOT_ACCEPTED",
|
|
|
|
|
|
"Admin 未确认接收任务结果",
|
|
|
|
|
|
False,
|
|
|
|
|
|
)
|
|
|
|
|
|
except AdminGatewayError as exc:
|
|
|
|
|
|
if exc.retryable:
|
|
|
|
|
|
self._repository.mark_outbox_retry(event.id, str(exc))
|
|
|
|
|
|
else:
|
|
|
|
|
|
self._repository.mark_outbox_failed(event.id, str(exc))
|
|
|
|
|
|
failed.append(task_id)
|
|
|
|
|
|
continue
|
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
|
self._repository.mark_outbox_retry(event.id, str(exc))
|
|
|
|
|
|
failed.append(task_id)
|
|
|
|
|
|
continue
|
|
|
|
|
|
|
|
|
|
|
|
self._repository.mark_outbox_sent(event.id)
|
|
|
|
|
|
succeeded.append(task_id)
|
|
|
|
|
|
self.finished.emit(succeeded, failed, skipped)
|
|
|
|
|
|
finally:
|
|
|
|
|
|
self.completed.emit()
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-08-10 17:23:24 +08:00
|
|
|
|
class TaskRemoveWorker(QObject):
|
|
|
|
|
|
"""在后台原子校验并软移除勾选任务,不访问任何界面控件。"""
|
|
|
|
|
|
|
|
|
|
|
|
succeeded = pyqtSignal(object)
|
|
|
|
|
|
failed = pyqtSignal(str)
|
|
|
|
|
|
completed = pyqtSignal()
|
|
|
|
|
|
|
|
|
|
|
|
def __init__(
|
|
|
|
|
|
self,
|
|
|
|
|
|
repository: TaskRepository,
|
|
|
|
|
|
task_ids: tuple[str, ...],
|
|
|
|
|
|
) -> None:
|
|
|
|
|
|
super().__init__()
|
|
|
|
|
|
self._repository = repository
|
|
|
|
|
|
self._task_ids = task_ids
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def run(self) -> None:
|
|
|
|
|
|
try:
|
|
|
|
|
|
self._repository.remove_tasks_from_list(self._task_ids)
|
|
|
|
|
|
except TaskRemovalError as exc:
|
|
|
|
|
|
self.failed.emit(str(exc))
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
self.failed.emit("删除本地任务失败,请检查数据库后重试。")
|
|
|
|
|
|
else:
|
|
|
|
|
|
self.succeeded.emit(self._task_ids)
|
|
|
|
|
|
finally:
|
|
|
|
|
|
self.completed.emit()
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-08-06 16:31:34 +08:00
|
|
|
|
class PDDTaskPageEvent(QObject):
|
|
|
|
|
|
"""把 PDD 页面只读操作连接到本地任务 Repository。"""
|
|
|
|
|
|
|
2026-08-10 18:00:38 +08:00
|
|
|
|
_claimRunRequested = pyqtSignal(str)
|
|
|
|
|
|
_claimReleaseRequested = pyqtSignal()
|
|
|
|
|
|
_claimShutdownRequested = pyqtSignal()
|
|
|
|
|
|
|
2026-08-06 16:31:34 +08:00
|
|
|
|
def __init__(
|
|
|
|
|
|
self,
|
|
|
|
|
|
page: PDDTaskPage,
|
|
|
|
|
|
repository: Optional[TaskRepository] = None,
|
|
|
|
|
|
parent=None,
|
2026-08-07 17:38:29 +08:00
|
|
|
|
claim_gateway: Optional[AdminGateway] = None,
|
2026-08-07 16:42:31 +08:00
|
|
|
|
settings_repository: Optional[SettingsRepository] = None,
|
2026-08-07 17:38:29 +08:00
|
|
|
|
collect_service_factory: Optional[CollectServiceFactory] = None,
|
2026-08-10 00:05:40 +08:00
|
|
|
|
purchase_adapter_factory: Optional[PurchaseAdapterFactory] = None,
|
2026-08-10 16:35:38 +08:00
|
|
|
|
live_purchase_adapter_factory: Optional[
|
|
|
|
|
|
LivePurchaseAdapterFactory
|
|
|
|
|
|
] = None,
|
2026-08-10 00:27:52 +08:00
|
|
|
|
purchase_reconcile_factory: Optional[PurchaseReconcileFactory] = None,
|
2026-08-10 09:35:46 +08:00
|
|
|
|
device_connection_checker: Optional[Callable[[str], None]] = None,
|
2026-08-09 10:54:18 +08:00
|
|
|
|
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),
|
2026-08-06 16:31:34 +08:00
|
|
|
|
):
|
|
|
|
|
|
super().__init__(parent or page)
|
|
|
|
|
|
self._page = page
|
|
|
|
|
|
self._repository = repository or TaskRepository()
|
|
|
|
|
|
self._filters = TaskFilters()
|
2026-08-07 16:42:31 +08:00
|
|
|
|
self._closing = False
|
|
|
|
|
|
self._claim_busy = False
|
|
|
|
|
|
self._claim_thread: Optional[QThread] = None
|
|
|
|
|
|
self._claim_worker: Optional[ClaimTaskWorker] = None
|
2026-08-10 18:00:38 +08:00
|
|
|
|
self._claim_operation = ""
|
2026-08-10 09:08:54 +08:00
|
|
|
|
self._rerun_cancel_requested = False
|
|
|
|
|
|
self._rerun_feedback: Optional[InfoBar] = None
|
2026-08-10 09:59:34 +08:00
|
|
|
|
self._claim_feedback: Optional[InfoBar] = None
|
2026-08-10 09:35:46 +08:00
|
|
|
|
self._device_feedback: Optional[InfoBar] = None
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._auto_fetch_running = False
|
2026-08-10 10:56:48 +08:00
|
|
|
|
self._resubmit_busy = False
|
|
|
|
|
|
self._resubmit_thread: Optional[QThread] = None
|
|
|
|
|
|
self._resubmit_worker: Optional[ResultResubmitWorker] = None
|
2026-08-10 17:23:24 +08:00
|
|
|
|
self._remove_busy = False
|
|
|
|
|
|
self._remove_thread: Optional[QThread] = None
|
|
|
|
|
|
self._remove_worker: Optional[TaskRemoveWorker] = None
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._stop_requested = False
|
|
|
|
|
|
self._cycle_next_delay_ms: Optional[int] = None
|
|
|
|
|
|
self._stop_status = "自动获取:已停止 · 当前没有执行中的任务"
|
|
|
|
|
|
self._retry_count = 0
|
2026-08-08 09:46:58 +08:00
|
|
|
|
self._detail_windows: Dict[str, TaskDetailWindow] = {}
|
2026-08-07 17:38:29 +08:00
|
|
|
|
self._collect_service_factory = collect_service_factory
|
2026-08-10 00:05:40 +08:00
|
|
|
|
self._purchase_adapter_factory = purchase_adapter_factory
|
2026-08-10 16:35:38 +08:00
|
|
|
|
self._live_purchase_adapter_factory = live_purchase_adapter_factory
|
2026-08-10 00:27:52 +08:00
|
|
|
|
self._purchase_reconcile_factory = purchase_reconcile_factory
|
2026-08-10 09:35:46 +08:00
|
|
|
|
self._device_connection_checker = (
|
|
|
|
|
|
device_connection_checker
|
|
|
|
|
|
or AndroidDeviceService().require_connected
|
|
|
|
|
|
)
|
2026-08-09 10:54:18 +08:00
|
|
|
|
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):
|
|
|
|
|
|
raise ValueError("自动获取重试等待时间必须全部大于 0")
|
|
|
|
|
|
self._next_task_delay_ms = next_task_delay_ms
|
|
|
|
|
|
self._no_task_delay_ms = no_task_delay_ms
|
|
|
|
|
|
self._retry_delays_ms = retry_delays_ms
|
|
|
|
|
|
self._next_cycle_timer = QTimer(self)
|
|
|
|
|
|
self._next_cycle_timer.setSingleShot(True)
|
|
|
|
|
|
self._next_cycle_timer.timeout.connect(self._start_claim_cycle)
|
2026-08-07 17:38:29 +08:00
|
|
|
|
|
|
|
|
|
|
try:
|
|
|
|
|
|
self._repository.recover_interrupted_work()
|
|
|
|
|
|
except AttributeError:
|
|
|
|
|
|
# 测试用的只读 Repository 可以不实现恢复接口。
|
|
|
|
|
|
pass
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
|
|
|
|
|
settings = settings_repository or SettingsRepository()
|
|
|
|
|
|
self._client_service = CurrentClientService(settings)
|
|
|
|
|
|
self._selected_android_device_service = SelectedAndroidDeviceService(
|
|
|
|
|
|
settings
|
|
|
|
|
|
)
|
|
|
|
|
|
self._claim_gateway = claim_gateway
|
|
|
|
|
|
self._claim_gateway_error = ""
|
|
|
|
|
|
if self._claim_gateway is None:
|
|
|
|
|
|
base_url = settings.get("admin.base_url", DEFAULT_ADMIN_BASE_URL)
|
|
|
|
|
|
timeout_value = settings.get("admin.request_timeout_seconds", 3.0)
|
|
|
|
|
|
try:
|
|
|
|
|
|
timeout_seconds = float(timeout_value)
|
|
|
|
|
|
self._claim_gateway = HttpAdminGateway(
|
|
|
|
|
|
base_url if isinstance(base_url, str) else "",
|
|
|
|
|
|
timeout_seconds=timeout_seconds,
|
2026-08-07 17:38:29 +08:00
|
|
|
|
client_id=self._client_service.load().client_id,
|
2026-08-07 16:42:31 +08:00
|
|
|
|
)
|
|
|
|
|
|
except (TypeError, ValueError) as exc:
|
|
|
|
|
|
self._claim_gateway_error = str(exc)
|
2026-08-06 16:31:34 +08:00
|
|
|
|
|
|
|
|
|
|
page.searchRequested.connect(self.search_tasks)
|
|
|
|
|
|
page.refreshRequested.connect(self.refresh_tasks)
|
2026-08-09 21:44:29 +08:00
|
|
|
|
page.rerunRequested.connect(self.request_rerun)
|
2026-08-10 09:08:54 +08:00
|
|
|
|
page.rerunCancelRequested.connect(self.request_cancel_rerun)
|
2026-08-10 10:56:48 +08:00
|
|
|
|
page.resubmitRequested.connect(self.request_resubmit)
|
2026-08-10 17:23:24 +08:00
|
|
|
|
page.removeRequested.connect(self.request_remove)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
page.autoFetchRequested.connect(self._request_claim_task)
|
2026-08-08 09:46:58 +08:00
|
|
|
|
page.detailRequested.connect(self.show_task_detail)
|
2026-08-06 16:31:34 +08:00
|
|
|
|
page.taskModel.loadMoreRequested.connect(self._load_page)
|
2026-08-09 10:54:18 +08:00
|
|
|
|
page.set_auto_fetch_state("stopped")
|
2026-08-07 16:42:31 +08:00
|
|
|
|
page.destroyed.connect(self.shutdown)
|
|
|
|
|
|
application = QCoreApplication.instance()
|
|
|
|
|
|
if application is not None:
|
|
|
|
|
|
application.aboutToQuit.connect(self.shutdown)
|
2026-08-06 16:31:34 +08:00
|
|
|
|
|
|
|
|
|
|
def load_initial_tasks(self) -> None:
|
|
|
|
|
|
"""应用启动后读取第一页本地任务。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._reload()
|
|
|
|
|
|
|
|
|
|
|
|
def search_tasks(self, values: Dict[str, str]) -> None:
|
|
|
|
|
|
"""把界面中文筛选值转换成领域筛选,并重新查询。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._filters = TaskFilters(
|
|
|
|
|
|
task_type=TASK_TYPE_BY_TEXT.get(values.get("task_type", "全部")),
|
|
|
|
|
|
status=TASK_STATUS_BY_TEXT.get(values.get("status", "全部")),
|
|
|
|
|
|
keyword=values.get("keyword", "").strip(),
|
|
|
|
|
|
)
|
|
|
|
|
|
self._reload()
|
|
|
|
|
|
|
|
|
|
|
|
def refresh_tasks(self) -> None:
|
|
|
|
|
|
"""使用当前筛选条件刷新列表。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._reload()
|
|
|
|
|
|
|
2026-08-09 21:44:29 +08:00
|
|
|
|
@pyqtSlot(str)
|
|
|
|
|
|
def request_rerun(self, task_id: str) -> None:
|
|
|
|
|
|
"""确认后在后台重新采集当前选中的一条终态任务。"""
|
|
|
|
|
|
|
|
|
|
|
|
if self._closing or not task_id:
|
|
|
|
|
|
return
|
2026-08-10 17:23:24 +08:00
|
|
|
|
if (
|
|
|
|
|
|
self._auto_fetch_running
|
|
|
|
|
|
or self._claim_busy
|
|
|
|
|
|
or self._resubmit_busy
|
|
|
|
|
|
or self._remove_busy
|
|
|
|
|
|
):
|
2026-08-09 21:44:29 +08:00
|
|
|
|
self._show_rerun_warning(
|
|
|
|
|
|
"暂时不能重新执行",
|
|
|
|
|
|
"自动获取或其他采集正在运行,请停止并等待当前任务结束。",
|
|
|
|
|
|
)
|
|
|
|
|
|
return
|
|
|
|
|
|
if self._claim_gateway_error or self._claim_gateway is None:
|
|
|
|
|
|
self._show_rerun_warning(
|
|
|
|
|
|
"重新执行不可用",
|
|
|
|
|
|
self._claim_gateway_error or "Admin 提交服务未初始化",
|
|
|
|
|
|
)
|
|
|
|
|
|
return
|
|
|
|
|
|
try:
|
|
|
|
|
|
client_settings = self._client_service.load()
|
|
|
|
|
|
android_serial = self._selected_android_device_service.load()
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
self._show_rerun_warning(
|
|
|
|
|
|
"重新执行不可用", "无法读取本地设备设置,请到设置页检查。"
|
|
|
|
|
|
)
|
|
|
|
|
|
return
|
|
|
|
|
|
if not client_settings.client_id:
|
|
|
|
|
|
self._show_rerun_warning(
|
|
|
|
|
|
"重新执行不可用", "请先在设置页保存当前设备号和设备名。"
|
|
|
|
|
|
)
|
|
|
|
|
|
return
|
|
|
|
|
|
if not android_serial:
|
2026-08-10 09:35:46 +08:00
|
|
|
|
self._show_device_unavailable(
|
|
|
|
|
|
"请先在设置页选择并保存 Android 设备"
|
2026-08-09 21:44:29 +08:00
|
|
|
|
)
|
|
|
|
|
|
return
|
|
|
|
|
|
try:
|
|
|
|
|
|
detail = self._repository.validate_collect_rerun(task_id)
|
|
|
|
|
|
except (CollectRerunError, ValueError) as exc:
|
|
|
|
|
|
self._show_rerun_warning("不能重新执行", str(exc))
|
|
|
|
|
|
return
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
self._show_rerun_warning(
|
|
|
|
|
|
"不能重新执行", "无法读取任务状态,请检查数据库后重试。"
|
|
|
|
|
|
)
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
|
|
title = detail.title or "尚未获取标题"
|
|
|
|
|
|
dialog = MessageBox(
|
|
|
|
|
|
"确认重新采集",
|
|
|
|
|
|
f"任务:{detail.remote_task_id}\n商品:{title}\n\n"
|
|
|
|
|
|
"新结果会覆盖 Client 和 Admin 的当前采集数据,旧结果仍保留在执行记录中。",
|
|
|
|
|
|
self._page.window(),
|
|
|
|
|
|
)
|
|
|
|
|
|
dialog.yesButton.setText("重新采集")
|
2026-08-10 09:08:54 +08:00
|
|
|
|
dialog.cancelButton.setText("暂不重新采集")
|
2026-08-09 21:44:29 +08:00
|
|
|
|
dialog.cancelButton.setFocus()
|
|
|
|
|
|
if not dialog.exec():
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
|
|
self._start_rerun_worker(task_id)
|
|
|
|
|
|
|
2026-08-10 10:56:48 +08:00
|
|
|
|
@pyqtSlot(object)
|
|
|
|
|
|
def request_resubmit(self, task_ids) -> None:
|
2026-08-10 11:18:06 +08:00
|
|
|
|
"""确认后在后台重发勾选任务尚未发送或最新的 Outbox。"""
|
2026-08-10 10:56:48 +08:00
|
|
|
|
|
|
|
|
|
|
stable_ids = tuple(
|
|
|
|
|
|
dict.fromkeys(str(value) for value in task_ids if value)
|
|
|
|
|
|
)
|
|
|
|
|
|
if self._closing or not stable_ids:
|
|
|
|
|
|
return
|
2026-08-10 17:23:24 +08:00
|
|
|
|
if (
|
|
|
|
|
|
self._auto_fetch_running
|
|
|
|
|
|
or self._claim_busy
|
|
|
|
|
|
or self._resubmit_busy
|
|
|
|
|
|
or self._remove_busy
|
|
|
|
|
|
):
|
2026-08-10 10:56:48 +08:00
|
|
|
|
self._show_rerun_warning(
|
|
|
|
|
|
"暂时不能重新上报",
|
|
|
|
|
|
"自动获取、重新执行或另一批上报正在运行,请等待当前操作结束。",
|
|
|
|
|
|
)
|
|
|
|
|
|
return
|
|
|
|
|
|
if self._claim_gateway_error or self._claim_gateway is None:
|
|
|
|
|
|
self._show_rerun_warning(
|
|
|
|
|
|
"重新上报不可用",
|
|
|
|
|
|
self._claim_gateway_error or "Admin 提交服务未初始化",
|
|
|
|
|
|
)
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
|
|
dialog = MessageBox(
|
2026-08-10 11:18:06 +08:00
|
|
|
|
f"重新上报 {len(stable_ids)} 条任务数据?",
|
|
|
|
|
|
"会优先提交尚未发送的结果或失败信息;没有待发送数据时,"
|
|
|
|
|
|
"才重新提交本地已经保存的最新结果。"
|
2026-08-10 10:56:48 +08:00
|
|
|
|
"不会重新采集、采购或操作 Android 手机,也不会修改本地结果。",
|
|
|
|
|
|
self._page.window(),
|
|
|
|
|
|
)
|
|
|
|
|
|
dialog.yesButton.setText(f"上报 {len(stable_ids)} 条")
|
|
|
|
|
|
dialog.cancelButton.setText("暂不上报")
|
|
|
|
|
|
dialog.cancelButton.setFocus()
|
|
|
|
|
|
if not dialog.exec():
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
|
|
self._start_resubmit_worker(stable_ids)
|
|
|
|
|
|
|
2026-08-10 17:23:24 +08:00
|
|
|
|
@pyqtSlot(object)
|
|
|
|
|
|
def request_remove(self, task_ids) -> None:
|
|
|
|
|
|
"""确认后在后台从普通列表软移除勾选的本地任务。"""
|
|
|
|
|
|
|
|
|
|
|
|
stable_ids = tuple(
|
|
|
|
|
|
dict.fromkeys(
|
|
|
|
|
|
str(value).strip()
|
|
|
|
|
|
for value in task_ids
|
|
|
|
|
|
if value is not None and str(value).strip()
|
|
|
|
|
|
)
|
|
|
|
|
|
)
|
|
|
|
|
|
if self._closing or not stable_ids:
|
|
|
|
|
|
return
|
|
|
|
|
|
if (
|
|
|
|
|
|
self._auto_fetch_running
|
|
|
|
|
|
or self._claim_busy
|
|
|
|
|
|
or self._resubmit_busy
|
|
|
|
|
|
or self._remove_busy
|
|
|
|
|
|
):
|
|
|
|
|
|
self._show_rerun_warning(
|
|
|
|
|
|
"暂时不能删除",
|
|
|
|
|
|
"自动获取、重新执行、重新上报或另一批删除正在运行,"
|
|
|
|
|
|
"请等待当前操作结束。",
|
|
|
|
|
|
)
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
|
|
count = len(stable_ids)
|
|
|
|
|
|
dialog = MessageBox(
|
|
|
|
|
|
f"从列表删除 {count} 条任务?",
|
|
|
|
|
|
"只会从普通任务列表中隐藏这些本地记录,不会删除 Admin 任务;"
|
|
|
|
|
|
"执行记录和上报历史仍会永久保留。\n\n"
|
|
|
|
|
|
"只有已结束、没有待上报数据且没有进入不可逆阶段的任务才允许删除。",
|
|
|
|
|
|
self._page.window(),
|
|
|
|
|
|
)
|
|
|
|
|
|
dialog.yesButton.setText(f"删除 {count} 条")
|
|
|
|
|
|
dialog.cancelButton.setText("暂不删除")
|
|
|
|
|
|
dialog.cancelButton.setFocus()
|
|
|
|
|
|
if not dialog.exec():
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
|
|
self._start_remove_worker(stable_ids)
|
|
|
|
|
|
|
|
|
|
|
|
def _start_remove_worker(self, task_ids: tuple[str, ...]) -> None:
|
|
|
|
|
|
"""启动本地任务软删除工作线程。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._remove_busy = True
|
|
|
|
|
|
self._page.set_remove_running(True)
|
|
|
|
|
|
self._page.set_engine_status(
|
|
|
|
|
|
f"正在检查并删除 {len(task_ids)} 条本地任务…"
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
thread = QThread(self)
|
|
|
|
|
|
worker = TaskRemoveWorker(self._repository, task_ids)
|
|
|
|
|
|
worker.moveToThread(thread)
|
|
|
|
|
|
thread.started.connect(worker.run)
|
|
|
|
|
|
worker.succeeded.connect(self._on_remove_succeeded)
|
|
|
|
|
|
worker.failed.connect(self._on_remove_failed)
|
|
|
|
|
|
worker.completed.connect(thread.quit)
|
|
|
|
|
|
worker.completed.connect(worker.deleteLater)
|
|
|
|
|
|
thread.finished.connect(thread.deleteLater)
|
|
|
|
|
|
thread.finished.connect(self._on_remove_thread_finished)
|
|
|
|
|
|
self._remove_thread = thread
|
|
|
|
|
|
self._remove_worker = worker
|
|
|
|
|
|
thread.start()
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot(object)
|
|
|
|
|
|
def _on_remove_succeeded(self, task_ids) -> None:
|
|
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
|
|
|
|
|
stable_ids = tuple(task_ids)
|
|
|
|
|
|
self._page.taskModel.uncheck_tasks(stable_ids)
|
|
|
|
|
|
self._page.set_engine_status(
|
|
|
|
|
|
f"已从列表删除 {len(stable_ids)} 条本地任务;"
|
|
|
|
|
|
"执行记录和上报历史仍会永久保留。"
|
|
|
|
|
|
)
|
|
|
|
|
|
self._reload()
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot(str)
|
|
|
|
|
|
def _on_remove_failed(self, message: str) -> None:
|
|
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
|
|
|
|
|
self._page.set_engine_status(f"删除未完成:{message}")
|
|
|
|
|
|
self._show_claim_error("不能删除所选任务", message)
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def _on_remove_thread_finished(self) -> None:
|
|
|
|
|
|
self._remove_worker = None
|
|
|
|
|
|
self._remove_thread = None
|
|
|
|
|
|
self._remove_busy = False
|
|
|
|
|
|
if not self._closing:
|
|
|
|
|
|
self._page.set_remove_running(False)
|
|
|
|
|
|
|
2026-08-10 10:56:48 +08:00
|
|
|
|
def _start_resubmit_worker(self, task_ids: tuple[str, ...]) -> None:
|
2026-08-10 11:18:06 +08:00
|
|
|
|
"""启动只处理指定任务 Outbox 的工作线程。"""
|
2026-08-10 10:56:48 +08:00
|
|
|
|
|
|
|
|
|
|
assert self._claim_gateway is not None
|
|
|
|
|
|
self._resubmit_busy = True
|
|
|
|
|
|
self._page.set_resubmit_running(True)
|
|
|
|
|
|
self._page.set_engine_status(
|
2026-08-10 11:18:06 +08:00
|
|
|
|
f"正在准备重新上报 {len(task_ids)} 条任务数据…"
|
2026-08-10 10:56:48 +08:00
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
thread = QThread(self)
|
|
|
|
|
|
worker = ResultResubmitWorker(
|
|
|
|
|
|
self._claim_gateway,
|
|
|
|
|
|
self._repository,
|
|
|
|
|
|
task_ids,
|
|
|
|
|
|
)
|
|
|
|
|
|
worker.moveToThread(thread)
|
|
|
|
|
|
thread.started.connect(worker.run)
|
|
|
|
|
|
worker.progress.connect(self._on_resubmit_progress)
|
|
|
|
|
|
worker.finished.connect(self._on_resubmit_finished)
|
|
|
|
|
|
worker.completed.connect(thread.quit)
|
|
|
|
|
|
worker.completed.connect(worker.deleteLater)
|
|
|
|
|
|
thread.finished.connect(thread.deleteLater)
|
|
|
|
|
|
thread.finished.connect(self._on_resubmit_thread_finished)
|
|
|
|
|
|
self._resubmit_thread = thread
|
|
|
|
|
|
self._resubmit_worker = worker
|
|
|
|
|
|
thread.start()
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot(int, int, str)
|
|
|
|
|
|
def _on_resubmit_progress(self, current: int, total: int, task_id: str) -> None:
|
|
|
|
|
|
if not self._closing:
|
|
|
|
|
|
self._page.set_engine_status(
|
|
|
|
|
|
f"正在重新上报 {current}/{total}:{task_id}"
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot(object, object, object)
|
|
|
|
|
|
def _on_resubmit_finished(self, succeeded, failed, skipped) -> None:
|
|
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
|
|
|
|
|
succeeded_ids = tuple(succeeded)
|
|
|
|
|
|
self._page.taskModel.uncheck_tasks(succeeded_ids)
|
|
|
|
|
|
message = (
|
|
|
|
|
|
f"重新上报完成:成功 {len(succeeded_ids)} 条,"
|
|
|
|
|
|
f"失败 {len(failed)} 条,跳过 {len(skipped)} 条"
|
|
|
|
|
|
)
|
|
|
|
|
|
self._page.set_engine_status(message)
|
|
|
|
|
|
if failed or skipped:
|
|
|
|
|
|
self._show_claim_error("重新上报未全部完成", message)
|
|
|
|
|
|
self._reload()
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def _on_resubmit_thread_finished(self) -> None:
|
|
|
|
|
|
self._resubmit_worker = None
|
|
|
|
|
|
self._resubmit_thread = None
|
|
|
|
|
|
self._resubmit_busy = False
|
|
|
|
|
|
if not self._closing:
|
|
|
|
|
|
self._page.set_resubmit_running(False)
|
|
|
|
|
|
|
2026-08-09 21:44:29 +08:00
|
|
|
|
def _start_rerun_worker(self, task_id: str) -> None:
|
2026-08-10 18:00:38 +08:00
|
|
|
|
"""把指定任务提交到持久设备工作线程。"""
|
2026-08-09 21:44:29 +08:00
|
|
|
|
|
|
|
|
|
|
assert self._claim_gateway is not None
|
2026-08-10 09:08:54 +08:00
|
|
|
|
self._close_rerun_feedback()
|
2026-08-10 09:59:34 +08:00
|
|
|
|
self._close_claim_feedback()
|
2026-08-10 09:35:46 +08:00
|
|
|
|
self._close_device_feedback()
|
2026-08-10 09:08:54 +08:00
|
|
|
|
self._rerun_cancel_requested = False
|
2026-08-09 21:44:29 +08:00
|
|
|
|
self._claim_busy = True
|
|
|
|
|
|
self._page.set_rerun_running(True)
|
2026-08-10 09:35:46 +08:00
|
|
|
|
self._page.set_engine_status(
|
|
|
|
|
|
f"正在检查 Android 设备并准备重新采集任务 {task_id}…"
|
|
|
|
|
|
)
|
2026-08-09 21:44:29 +08:00
|
|
|
|
self._reload()
|
|
|
|
|
|
|
2026-08-10 18:00:38 +08:00
|
|
|
|
worker = self._ensure_claim_executor()
|
|
|
|
|
|
self._claim_operation = "rerun"
|
|
|
|
|
|
worker.prepare_run()
|
|
|
|
|
|
self._claimRunRequested.emit(task_id)
|
2026-08-09 21:44:29 +08:00
|
|
|
|
|
2026-08-10 09:08:54 +08:00
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def request_cancel_rerun(self) -> None:
|
|
|
|
|
|
"""请求当前重新采集在下一个安全点停止。"""
|
|
|
|
|
|
|
|
|
|
|
|
worker = self._claim_worker
|
|
|
|
|
|
if (
|
|
|
|
|
|
self._closing
|
|
|
|
|
|
or self._rerun_cancel_requested
|
|
|
|
|
|
or not self._claim_busy
|
|
|
|
|
|
or worker is None
|
|
|
|
|
|
):
|
|
|
|
|
|
return
|
|
|
|
|
|
self._rerun_cancel_requested = True
|
|
|
|
|
|
self._page.set_rerun_cancelling()
|
|
|
|
|
|
self._page.set_engine_status(
|
|
|
|
|
|
"正在停止重新采集…正在等待手机当前操作结束"
|
|
|
|
|
|
)
|
|
|
|
|
|
# uiautomator2/ADB 的单次调用无法安全强制中断。
|
|
|
|
|
|
# Worker 会在调用返回后通过 cancelled 回调在安全点停止。
|
|
|
|
|
|
worker.cancel()
|
|
|
|
|
|
|
2026-08-09 21:44:29 +08:00
|
|
|
|
@pyqtSlot(str, str, str)
|
|
|
|
|
|
def _on_rerun_outcome(self, kind: str, message: str, _task_id: str) -> None:
|
|
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
|
|
|
|
|
self._reload()
|
2026-08-10 09:08:54 +08:00
|
|
|
|
if self._rerun_cancel_requested:
|
|
|
|
|
|
self._page.set_engine_status(
|
|
|
|
|
|
"停止请求已处理,请查看任务最新状态"
|
|
|
|
|
|
)
|
|
|
|
|
|
return
|
2026-08-09 21:44:29 +08:00
|
|
|
|
self._page.set_engine_status(message)
|
|
|
|
|
|
if kind != "succeeded":
|
2026-08-10 09:08:54 +08:00
|
|
|
|
self._show_rerun_error("重新采集需要处理", message)
|
2026-08-09 21:44:29 +08:00
|
|
|
|
|
|
|
|
|
|
@pyqtSlot(str)
|
|
|
|
|
|
def _on_rerun_failed(self, message: str) -> None:
|
|
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
2026-08-10 09:08:54 +08:00
|
|
|
|
if self._rerun_cancel_requested:
|
|
|
|
|
|
return
|
2026-08-09 21:44:29 +08:00
|
|
|
|
self._page.set_engine_status(message)
|
2026-08-10 09:08:54 +08:00
|
|
|
|
self._show_rerun_error("重新采集失败", message)
|
2026-08-09 21:44:29 +08:00
|
|
|
|
|
2026-08-10 09:35:46 +08:00
|
|
|
|
@pyqtSlot(str)
|
|
|
|
|
|
def _on_rerun_device_unavailable(self, message: str) -> None:
|
|
|
|
|
|
if self._closing or self._rerun_cancel_requested:
|
|
|
|
|
|
return
|
|
|
|
|
|
content = message or "Android 设备未连接,任务未开始"
|
|
|
|
|
|
self._page.set_engine_status(
|
|
|
|
|
|
f"Android 设备不可用,重新采集未开始:{content}"
|
|
|
|
|
|
)
|
|
|
|
|
|
self._show_device_unavailable(content)
|
|
|
|
|
|
|
2026-08-09 21:44:29 +08:00
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def _on_rerun_thread_finished(self) -> None:
|
2026-08-10 09:08:54 +08:00
|
|
|
|
cancel_requested = self._rerun_cancel_requested
|
2026-08-09 21:44:29 +08:00
|
|
|
|
self._claim_busy = False
|
2026-08-10 09:08:54 +08:00
|
|
|
|
self._rerun_cancel_requested = False
|
2026-08-09 21:44:29 +08:00
|
|
|
|
if not self._closing:
|
|
|
|
|
|
self._page.set_rerun_running(False)
|
2026-08-10 09:08:54 +08:00
|
|
|
|
if cancel_requested:
|
|
|
|
|
|
self._reload()
|
|
|
|
|
|
self._page.set_engine_status(
|
|
|
|
|
|
"停止请求已处理,请查看任务最新状态"
|
|
|
|
|
|
)
|
2026-08-09 21:44:29 +08:00
|
|
|
|
|
|
|
|
|
|
def _show_rerun_warning(self, title: str, content: str) -> None:
|
2026-08-10 09:08:54 +08:00
|
|
|
|
bar = InfoBar.warning(
|
2026-08-09 21:44:29 +08:00
|
|
|
|
title=title,
|
|
|
|
|
|
content=content,
|
|
|
|
|
|
isClosable=True,
|
|
|
|
|
|
duration=5000,
|
2026-08-10 10:16:11 +08:00
|
|
|
|
position=InfoBarPosition.TOP,
|
2026-08-09 21:44:29 +08:00
|
|
|
|
parent=self._page,
|
|
|
|
|
|
)
|
2026-08-10 09:08:54 +08:00
|
|
|
|
self._replace_rerun_feedback(bar)
|
|
|
|
|
|
|
|
|
|
|
|
def _show_rerun_error(self, title: str, content: str) -> None:
|
2026-08-10 09:59:34 +08:00
|
|
|
|
"""显示一条较大、可明确关闭且会自动消失的重新采集错误。"""
|
2026-08-10 09:08:54 +08:00
|
|
|
|
|
|
|
|
|
|
bar = InfoBar.error(
|
|
|
|
|
|
title=title,
|
|
|
|
|
|
content=content,
|
|
|
|
|
|
isClosable=True,
|
2026-08-10 09:59:34 +08:00
|
|
|
|
duration=ERROR_FEEDBACK_DURATION_MS,
|
2026-08-10 10:05:45 +08:00
|
|
|
|
position=InfoBarPosition.TOP,
|
2026-08-10 09:08:54 +08:00
|
|
|
|
parent=self._page,
|
|
|
|
|
|
)
|
|
|
|
|
|
self._replace_rerun_feedback(bar)
|
|
|
|
|
|
|
|
|
|
|
|
def _replace_rerun_feedback(self, bar: InfoBar) -> None:
|
|
|
|
|
|
"""用新提示替换旧提示,避免连续失败后堆叠。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._close_rerun_feedback()
|
|
|
|
|
|
self._rerun_feedback = bar
|
2026-08-10 09:59:34 +08:00
|
|
|
|
self._set_large_error_feedback_size(bar)
|
2026-08-10 09:08:54 +08:00
|
|
|
|
close_button = PushButton("关闭提示", bar)
|
2026-08-10 09:59:34 +08:00
|
|
|
|
self._set_large_feedback_button_size(close_button)
|
2026-08-10 09:08:54 +08:00
|
|
|
|
close_button.setAccessibleName("关闭重新采集提示")
|
|
|
|
|
|
close_button.clicked.connect(
|
|
|
|
|
|
lambda _checked=False, current=bar: self._close_rerun_feedback(current)
|
|
|
|
|
|
)
|
|
|
|
|
|
bar.addWidget(close_button)
|
|
|
|
|
|
bar.destroyed.connect(
|
|
|
|
|
|
lambda _object=None, current=bar: self._forget_rerun_feedback(current)
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
def _close_rerun_feedback(self, expected: Optional[InfoBar] = None) -> None:
|
|
|
|
|
|
"""关闭当前重新采集提示;旧提示不得关闭新提示。"""
|
|
|
|
|
|
|
|
|
|
|
|
bar = self._rerun_feedback
|
|
|
|
|
|
if bar is None or (expected is not None and bar is not expected):
|
|
|
|
|
|
return
|
|
|
|
|
|
self._rerun_feedback = None
|
|
|
|
|
|
try:
|
|
|
|
|
|
bar.close()
|
|
|
|
|
|
except RuntimeError:
|
|
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
|
|
|
def _forget_rerun_feedback(self, bar: InfoBar) -> None:
|
|
|
|
|
|
if self._rerun_feedback is bar:
|
|
|
|
|
|
self._rerun_feedback = None
|
2026-08-09 21:44:29 +08:00
|
|
|
|
|
2026-08-07 16:42:31 +08:00
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def _request_claim_task(self) -> None:
|
2026-08-09 10:54:18 +08:00
|
|
|
|
"""切换持续自动获取;每一轮仍只启动一个后台 Worker。"""
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
2026-08-09 10:54:18 +08:00
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
2026-08-10 17:23:24 +08:00
|
|
|
|
if self._resubmit_busy or self._remove_busy:
|
2026-08-10 10:56:48 +08:00
|
|
|
|
self._show_rerun_warning(
|
|
|
|
|
|
"暂时不能获取任务",
|
2026-08-10 17:23:24 +08:00
|
|
|
|
"本地任务正在重新上报或删除,请等待当前操作结束。",
|
2026-08-10 10:56:48 +08:00
|
|
|
|
)
|
|
|
|
|
|
return
|
2026-08-09 10:54:18 +08:00
|
|
|
|
if self._auto_fetch_running:
|
|
|
|
|
|
self._request_stop_auto_fetch()
|
2026-08-07 16:42:31 +08:00
|
|
|
|
return
|
|
|
|
|
|
if self._claim_gateway_error or self._claim_gateway is None:
|
|
|
|
|
|
message = self._claim_gateway_error or "Admin 领取服务未初始化"
|
|
|
|
|
|
self._page.set_engine_status(f"领取失败:{message}")
|
|
|
|
|
|
self._show_claim_error("领取任务失败", message)
|
|
|
|
|
|
return
|
|
|
|
|
|
|
2026-08-10 09:35:46 +08:00
|
|
|
|
self._close_device_feedback()
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._auto_fetch_running = True
|
|
|
|
|
|
self._stop_requested = False
|
|
|
|
|
|
self._retry_count = 0
|
|
|
|
|
|
self._page.set_auto_fetch_state("starting")
|
|
|
|
|
|
self._page.set_engine_status("自动获取:正在启动…")
|
|
|
|
|
|
self._start_claim_cycle()
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def _start_claim_cycle(self) -> None:
|
2026-08-10 18:00:38 +08:00
|
|
|
|
"""开始一轮任务处理;设备命令由固定工作线程串行执行。"""
|
2026-08-09 10:54:18 +08:00
|
|
|
|
|
|
|
|
|
|
if (
|
|
|
|
|
|
self._closing
|
|
|
|
|
|
or not self._auto_fetch_running
|
|
|
|
|
|
or self._stop_requested
|
|
|
|
|
|
or self._claim_busy
|
|
|
|
|
|
):
|
|
|
|
|
|
return
|
2026-08-07 16:42:31 +08:00
|
|
|
|
self._claim_busy = True
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._cycle_next_delay_ms = None
|
|
|
|
|
|
self._page.set_auto_fetch_state("running")
|
2026-08-10 00:05:40 +08:00
|
|
|
|
self._page.set_engine_status("自动获取:运行中 · 正在处理一条任务…")
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
2026-08-10 18:00:38 +08:00
|
|
|
|
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
|
2026-08-07 16:42:31 +08:00
|
|
|
|
thread = QThread(self)
|
|
|
|
|
|
worker = ClaimTaskWorker(
|
2026-08-10 00:05:40 +08:00
|
|
|
|
gateway=self._claim_gateway,
|
|
|
|
|
|
task_repository=self._repository,
|
|
|
|
|
|
client_service=self._client_service,
|
|
|
|
|
|
android_device_service=self._selected_android_device_service,
|
|
|
|
|
|
collect_service_factory=self._collect_service_factory,
|
|
|
|
|
|
purchase_adapter_factory=self._purchase_adapter_factory,
|
2026-08-10 16:35:38 +08:00
|
|
|
|
live_purchase_adapter_factory=self._live_purchase_adapter_factory,
|
2026-08-10 00:27:52 +08:00
|
|
|
|
purchase_reconcile_factory=self._purchase_reconcile_factory,
|
2026-08-10 09:35:46 +08:00
|
|
|
|
device_connection_checker=self._device_connection_checker,
|
2026-08-07 16:42:31 +08:00
|
|
|
|
)
|
|
|
|
|
|
worker.moveToThread(thread)
|
2026-08-10 18:00:38 +08:00
|
|
|
|
self._claimRunRequested.connect(worker.run_selected)
|
|
|
|
|
|
self._claimReleaseRequested.connect(worker.release_device)
|
|
|
|
|
|
self._claimShutdownRequested.connect(worker.shutdown_worker)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
worker.taskSaved.connect(self._on_claimed_task_saved)
|
2026-08-10 18:00:38 +08:00
|
|
|
|
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)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
thread.finished.connect(thread.deleteLater)
|
|
|
|
|
|
self._claim_thread = thread
|
|
|
|
|
|
self._claim_worker = worker
|
|
|
|
|
|
thread.start()
|
2026-08-10 18:00:38 +08:00
|
|
|
|
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()
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
2026-08-09 10:54:18 +08:00
|
|
|
|
def _request_stop_auto_fetch(self) -> None:
|
|
|
|
|
|
"""取消下一轮,并请求当前 Worker 在安全点停止。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._auto_fetch_running = False
|
|
|
|
|
|
self._stop_requested = True
|
|
|
|
|
|
self._next_cycle_timer.stop()
|
|
|
|
|
|
self._stop_status = "自动获取:已停止 · 当前没有执行中的任务"
|
|
|
|
|
|
worker = self._claim_worker
|
2026-08-10 18:00:38 +08:00
|
|
|
|
if worker is not None and self._claim_busy:
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._page.set_auto_fetch_state("stopping")
|
|
|
|
|
|
self._page.set_engine_status(
|
|
|
|
|
|
"自动获取:正在停止 · 等待当前任务安全结束"
|
|
|
|
|
|
)
|
|
|
|
|
|
try:
|
|
|
|
|
|
worker.cancel()
|
|
|
|
|
|
except RuntimeError:
|
|
|
|
|
|
pass
|
|
|
|
|
|
return
|
|
|
|
|
|
self._finish_auto_fetch_stop()
|
|
|
|
|
|
|
|
|
|
|
|
def _finish_auto_fetch_stop(self) -> None:
|
|
|
|
|
|
"""把按钮和状态区恢复为已停止。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._page.set_auto_fetch_state("stopped")
|
|
|
|
|
|
self._page.set_engine_status(self._stop_status)
|
|
|
|
|
|
self._stop_requested = False
|
2026-08-10 18:00:38 +08:00
|
|
|
|
if self._claim_worker is not None:
|
|
|
|
|
|
self._claimReleaseRequested.emit()
|
2026-08-09 10:54:18 +08:00
|
|
|
|
|
|
|
|
|
|
def _continue_after(self, delay_ms: int, status: str) -> None:
|
|
|
|
|
|
"""记录本轮结束后的下一次调度,不在 Worker 结束前启动。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._cycle_next_delay_ms = delay_ms
|
|
|
|
|
|
self._page.set_engine_status(status)
|
|
|
|
|
|
|
|
|
|
|
|
def _stop_after_current(self, status: str) -> None:
|
|
|
|
|
|
"""遇到不可自动处理的问题时,完成线程收尾后停止。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._auto_fetch_running = False
|
|
|
|
|
|
self._stop_requested = True
|
|
|
|
|
|
self._stop_status = status
|
|
|
|
|
|
self._next_cycle_timer.stop()
|
|
|
|
|
|
self._page.set_auto_fetch_state("stopping")
|
|
|
|
|
|
|
2026-08-07 16:42:31 +08:00
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def _on_no_claimed_task(self) -> None:
|
|
|
|
|
|
if not self._closing:
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._retry_count = 0
|
|
|
|
|
|
seconds = self._no_task_delay_ms / 1000
|
|
|
|
|
|
self._continue_after(
|
|
|
|
|
|
self._no_task_delay_ms,
|
|
|
|
|
|
f"自动获取:运行中 · 暂无任务,{seconds:g} 秒后再次领取",
|
|
|
|
|
|
)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
|
|
|
|
|
@pyqtSlot(str)
|
|
|
|
|
|
def _on_claimed_task_saved(self, task_id: str) -> None:
|
|
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
|
|
|
|
|
self._page.set_engine_status(
|
|
|
|
|
|
f"已领取任务 {task_id},已保存到本地任务列表"
|
|
|
|
|
|
)
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._cycle_next_delay_ms = self._next_task_delay_ms
|
2026-08-07 16:42:31 +08:00
|
|
|
|
self._reload()
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot(str)
|
|
|
|
|
|
def _on_duplicate_claimed_task(self, task_id: str) -> None:
|
|
|
|
|
|
if not self._closing:
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._continue_after(
|
|
|
|
|
|
self._next_task_delay_ms,
|
|
|
|
|
|
f"任务 {task_id} 本地已有,未重复保存",
|
2026-08-07 16:42:31 +08:00
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot(str, str)
|
|
|
|
|
|
def _on_claimed_task_save_failed(self, task_id: str, message: str) -> None:
|
|
|
|
|
|
if not self._closing:
|
|
|
|
|
|
content = (
|
|
|
|
|
|
f"任务 {task_id} 已在服务端领取,但本地保存失败:{message}。"
|
|
|
|
|
|
"请记下这个任务号联系维护者。"
|
|
|
|
|
|
)
|
|
|
|
|
|
self._page.set_engine_status(content)
|
|
|
|
|
|
self._show_claim_error("本地任务保存失败", content)
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._stop_after_current(content)
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot(str)
|
|
|
|
|
|
def _on_claim_retryable_failed(self, message: str) -> None:
|
|
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
|
|
|
|
|
index = min(self._retry_count, len(self._retry_delays_ms) - 1)
|
|
|
|
|
|
delay_ms = self._retry_delays_ms[index]
|
|
|
|
|
|
self._retry_count += 1
|
|
|
|
|
|
seconds = delay_ms / 1000
|
|
|
|
|
|
status = (
|
|
|
|
|
|
f"自动获取:退避等待 · {message},{seconds:g} 秒后重试"
|
|
|
|
|
|
)
|
|
|
|
|
|
self._continue_after(delay_ms, status)
|
|
|
|
|
|
if self._retry_count == 1:
|
|
|
|
|
|
self._show_claim_error("Admin 暂时不可用", status)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
|
|
|
|
|
@pyqtSlot(str)
|
|
|
|
|
|
def _on_claim_failed(self, message: str) -> None:
|
|
|
|
|
|
if not self._closing:
|
|
|
|
|
|
content = message or "领取任务失败,请稍后重试"
|
|
|
|
|
|
self._page.set_engine_status(content)
|
|
|
|
|
|
self._show_claim_error("领取任务失败", content)
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._stop_after_current(content)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
2026-08-10 09:35:46 +08:00
|
|
|
|
@pyqtSlot(str)
|
|
|
|
|
|
def _on_claim_device_unavailable(self, message: str) -> None:
|
|
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
|
|
|
|
|
content = message or "Android 设备未连接,任务未开始"
|
|
|
|
|
|
status = f"自动获取:已停止 · Android 设备不可用:{content}"
|
|
|
|
|
|
self._page.set_engine_status(status)
|
|
|
|
|
|
self._show_device_unavailable(content)
|
|
|
|
|
|
self._stop_after_current(status)
|
|
|
|
|
|
|
2026-08-07 17:38:29 +08:00
|
|
|
|
@pyqtSlot(str, str, str)
|
2026-08-09 22:47:44 +08:00
|
|
|
|
def _on_collect_outcome(self, kind: str, message: str, task_id: str) -> None:
|
2026-08-07 17:38:29 +08:00
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
|
|
|
|
|
self._reload()
|
2026-08-09 10:54:18 +08:00
|
|
|
|
if kind == "no_task":
|
|
|
|
|
|
self._on_no_claimed_task()
|
2026-08-10 18:20:34 +08:00
|
|
|
|
elif kind in {"succeeded", "business_failed"}:
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._retry_count = 0
|
|
|
|
|
|
self._continue_after(
|
|
|
|
|
|
self._next_task_delay_ms,
|
|
|
|
|
|
f"自动获取:运行中 · {message}",
|
|
|
|
|
|
)
|
|
|
|
|
|
elif kind == "result_pending":
|
|
|
|
|
|
self._on_claim_retryable_failed(message)
|
2026-08-10 19:18:26 +08:00
|
|
|
|
elif kind == "reconcile_pending":
|
|
|
|
|
|
self._retry_count = 0
|
|
|
|
|
|
seconds = PURCHASE_RECONCILE_DELAY_MS / 1000
|
|
|
|
|
|
self._continue_after(
|
|
|
|
|
|
PURCHASE_RECONCILE_DELAY_MS,
|
|
|
|
|
|
(
|
|
|
|
|
|
f"自动获取:运行中 · {message};"
|
|
|
|
|
|
f"{seconds:g} 秒后自动核对订单"
|
|
|
|
|
|
),
|
|
|
|
|
|
)
|
2026-08-09 22:47:44 +08:00
|
|
|
|
elif kind == "failed" and self._is_retry_wait_task(task_id):
|
|
|
|
|
|
content = (
|
|
|
|
|
|
f"自动获取:已停止 · 任务 {task_id} 重试已暂停;"
|
|
|
|
|
|
"请选择该任务点击“重新执行”,或重新启动获取任务"
|
|
|
|
|
|
)
|
|
|
|
|
|
self._show_retry_paused(content)
|
|
|
|
|
|
self._stop_after_current(content)
|
2026-08-09 10:54:18 +08:00
|
|
|
|
elif kind in {"manual_review", "failed"}:
|
2026-08-10 00:05:40 +08:00
|
|
|
|
self._show_claim_error("任务需要处理", message)
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._stop_after_current(message)
|
|
|
|
|
|
elif kind == "cancelled":
|
|
|
|
|
|
self._stop_after_current(message or "自动获取已停止")
|
|
|
|
|
|
else:
|
|
|
|
|
|
content = message or f"任务返回未知结果:{kind}"
|
|
|
|
|
|
self._show_claim_error("自动获取已停止", content)
|
|
|
|
|
|
self._stop_after_current(content)
|
2026-08-07 17:38:29 +08:00
|
|
|
|
|
2026-08-09 22:47:44 +08:00
|
|
|
|
def _is_retry_wait_task(self, task_id: str) -> bool:
|
|
|
|
|
|
"""判断失败结果对应的任务是否仍可在本地重试。"""
|
|
|
|
|
|
|
|
|
|
|
|
if not task_id:
|
|
|
|
|
|
return False
|
|
|
|
|
|
try:
|
|
|
|
|
|
task = self._repository.get_task(task_id)
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
return False
|
2026-08-10 00:05:40 +08:00
|
|
|
|
return (
|
|
|
|
|
|
task is not None
|
|
|
|
|
|
and task.task_type is TaskType.COLLECT
|
|
|
|
|
|
and task.status is TaskStatus.RETRY_WAIT
|
|
|
|
|
|
)
|
2026-08-09 22:47:44 +08:00
|
|
|
|
|
|
|
|
|
|
def _show_retry_paused(self, content: str) -> None:
|
|
|
|
|
|
"""用持久警告说明任务不会自行倒计时重试。"""
|
|
|
|
|
|
|
|
|
|
|
|
InfoBar.warning(
|
|
|
|
|
|
title="重试已暂停",
|
|
|
|
|
|
content=content,
|
|
|
|
|
|
isClosable=True,
|
|
|
|
|
|
duration=-1,
|
|
|
|
|
|
position=InfoBarPosition.TOP_RIGHT,
|
|
|
|
|
|
parent=self._page,
|
|
|
|
|
|
)
|
|
|
|
|
|
|
2026-08-07 16:42:31 +08:00
|
|
|
|
def _show_claim_error(self, title: str, content: str) -> None:
|
2026-08-10 09:59:34 +08:00
|
|
|
|
"""显示较大的临时错误,同时在底部保留完整状态。"""
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
2026-08-10 09:59:34 +08:00
|
|
|
|
self._close_claim_feedback()
|
|
|
|
|
|
bar = InfoBar.error(
|
2026-08-07 16:42:31 +08:00
|
|
|
|
title=title,
|
|
|
|
|
|
content=content,
|
|
|
|
|
|
isClosable=True,
|
2026-08-10 09:59:34 +08:00
|
|
|
|
duration=ERROR_FEEDBACK_DURATION_MS,
|
2026-08-10 10:05:45 +08:00
|
|
|
|
position=InfoBarPosition.TOP,
|
2026-08-07 16:42:31 +08:00
|
|
|
|
parent=self._page,
|
|
|
|
|
|
)
|
2026-08-10 09:59:34 +08:00
|
|
|
|
self._claim_feedback = bar
|
|
|
|
|
|
self._set_large_error_feedback_size(bar)
|
|
|
|
|
|
close_button = PushButton("关闭提示", bar)
|
|
|
|
|
|
self._set_large_feedback_button_size(close_button)
|
|
|
|
|
|
close_button.setAccessibleName("关闭采集采购错误提示")
|
|
|
|
|
|
close_button.clicked.connect(
|
|
|
|
|
|
lambda _checked=False, current=bar: self._close_claim_feedback(current)
|
|
|
|
|
|
)
|
|
|
|
|
|
bar.addWidget(close_button)
|
|
|
|
|
|
bar.destroyed.connect(
|
|
|
|
|
|
lambda _object=None, current=bar: self._forget_claim_feedback(current)
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
def _close_claim_feedback(self, expected: Optional[InfoBar] = None) -> None:
|
|
|
|
|
|
bar = self._claim_feedback
|
|
|
|
|
|
if bar is None or (expected is not None and bar is not expected):
|
|
|
|
|
|
return
|
|
|
|
|
|
self._claim_feedback = None
|
|
|
|
|
|
try:
|
|
|
|
|
|
bar.close()
|
|
|
|
|
|
except RuntimeError:
|
|
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
|
|
|
def _forget_claim_feedback(self, bar: InfoBar) -> None:
|
|
|
|
|
|
if self._claim_feedback is bar:
|
|
|
|
|
|
self._claim_feedback = None
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
2026-08-10 09:35:46 +08:00
|
|
|
|
def _show_device_unavailable(self, content: str) -> None:
|
|
|
|
|
|
"""显示一条可进入设置且不会堆叠的设备错误。"""
|
|
|
|
|
|
|
|
|
|
|
|
self._close_device_feedback()
|
|
|
|
|
|
bar = InfoBar.error(
|
|
|
|
|
|
title="Android 设备不可用",
|
|
|
|
|
|
content=content,
|
|
|
|
|
|
isClosable=True,
|
2026-08-10 09:59:34 +08:00
|
|
|
|
duration=ERROR_FEEDBACK_DURATION_MS,
|
2026-08-10 10:05:45 +08:00
|
|
|
|
position=InfoBarPosition.TOP,
|
2026-08-10 09:35:46 +08:00
|
|
|
|
parent=self._page,
|
|
|
|
|
|
)
|
|
|
|
|
|
self._device_feedback = bar
|
2026-08-10 09:59:34 +08:00
|
|
|
|
self._set_large_error_feedback_size(bar)
|
2026-08-10 09:35:46 +08:00
|
|
|
|
settings_button = PushButton("打开设置", bar)
|
2026-08-10 09:59:34 +08:00
|
|
|
|
self._set_large_feedback_button_size(settings_button)
|
2026-08-10 09:35:46 +08:00
|
|
|
|
settings_button.setAccessibleName("打开 Android 设备设置")
|
|
|
|
|
|
settings_button.clicked.connect(
|
|
|
|
|
|
lambda _checked=False, current=bar: self._open_device_settings(current)
|
|
|
|
|
|
)
|
|
|
|
|
|
close_button = PushButton("关闭提示", bar)
|
2026-08-10 09:59:34 +08:00
|
|
|
|
self._set_large_feedback_button_size(close_button)
|
2026-08-10 09:35:46 +08:00
|
|
|
|
close_button.setAccessibleName("关闭 Android 设备提示")
|
|
|
|
|
|
close_button.clicked.connect(
|
|
|
|
|
|
lambda _checked=False, current=bar: self._close_device_feedback(current)
|
|
|
|
|
|
)
|
|
|
|
|
|
bar.addWidget(settings_button)
|
|
|
|
|
|
bar.addWidget(close_button)
|
|
|
|
|
|
bar.destroyed.connect(
|
|
|
|
|
|
lambda _object=None, current=bar: self._forget_device_feedback(current)
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
def _open_device_settings(self, bar: InfoBar) -> None:
|
|
|
|
|
|
self._page.openSettingsRequested.emit()
|
|
|
|
|
|
self._close_device_feedback(bar)
|
|
|
|
|
|
|
|
|
|
|
|
def _close_device_feedback(self, expected: Optional[InfoBar] = None) -> None:
|
|
|
|
|
|
bar = self._device_feedback
|
|
|
|
|
|
if bar is None or (expected is not None and bar is not expected):
|
|
|
|
|
|
return
|
|
|
|
|
|
self._device_feedback = None
|
|
|
|
|
|
try:
|
|
|
|
|
|
bar.close()
|
|
|
|
|
|
except RuntimeError:
|
|
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
|
|
|
def _forget_device_feedback(self, bar: InfoBar) -> None:
|
|
|
|
|
|
if self._device_feedback is bar:
|
|
|
|
|
|
self._device_feedback = None
|
|
|
|
|
|
|
2026-08-10 09:59:34 +08:00
|
|
|
|
@staticmethod
|
|
|
|
|
|
def _set_large_error_feedback_size(bar: InfoBar) -> None:
|
|
|
|
|
|
"""设置适合长中文错误和高缩放环境的最小尺寸。"""
|
|
|
|
|
|
|
|
|
|
|
|
bar.setMinimumSize(
|
|
|
|
|
|
ERROR_FEEDBACK_MIN_WIDTH,
|
|
|
|
|
|
ERROR_FEEDBACK_MIN_HEIGHT,
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
@staticmethod
|
|
|
|
|
|
def _set_large_feedback_button_size(button: PushButton) -> None:
|
|
|
|
|
|
"""扩大操作按钮,避免用户只能点击右上角的小关闭图标。"""
|
|
|
|
|
|
|
|
|
|
|
|
button.setMinimumSize(
|
|
|
|
|
|
ERROR_FEEDBACK_BUTTON_MIN_WIDTH,
|
|
|
|
|
|
ERROR_FEEDBACK_BUTTON_MIN_HEIGHT,
|
|
|
|
|
|
)
|
|
|
|
|
|
|
2026-08-07 16:42:31 +08:00
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def _on_claim_thread_finished(self) -> None:
|
|
|
|
|
|
self._claim_busy = False
|
2026-08-09 10:54:18 +08:00
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
|
|
|
|
|
if not self._auto_fetch_running or self._stop_requested:
|
|
|
|
|
|
self._finish_auto_fetch_stop()
|
|
|
|
|
|
return
|
|
|
|
|
|
delay_ms = self._cycle_next_delay_ms
|
|
|
|
|
|
if delay_ms is None:
|
|
|
|
|
|
self._stop_after_current("自动获取:已停止 · 本轮没有返回有效结果")
|
|
|
|
|
|
self._finish_auto_fetch_stop()
|
|
|
|
|
|
return
|
|
|
|
|
|
self._page.set_auto_fetch_state("running")
|
|
|
|
|
|
self._next_cycle_timer.start(delay_ms)
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
2026-08-06 16:31:34 +08:00
|
|
|
|
def _reload(self) -> None:
|
|
|
|
|
|
self._page.begin_task_reload()
|
|
|
|
|
|
self._page.taskModel.fetchMore()
|
|
|
|
|
|
|
|
|
|
|
|
def _load_page(self, offset: int, limit: int) -> None:
|
|
|
|
|
|
try:
|
|
|
|
|
|
summaries = self._repository.list_tasks(
|
|
|
|
|
|
filters=self._filters,
|
|
|
|
|
|
limit=limit,
|
|
|
|
|
|
offset=offset,
|
|
|
|
|
|
)
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
self._page.set_load_error("无法读取本地任务,请检查数据库后重试。")
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
|
|
rows = [summary_to_row(summary) for summary in summaries]
|
|
|
|
|
|
self._page.append_task_page(rows, has_more=len(rows) == limit)
|
|
|
|
|
|
|
2026-08-08 09:46:58 +08:00
|
|
|
|
@pyqtSlot(str)
|
|
|
|
|
|
def show_task_detail(self, task_id: str) -> None:
|
|
|
|
|
|
"""读取本地完整任务,并打开一个非模态详情窗口。"""
|
|
|
|
|
|
|
|
|
|
|
|
if self._closing or not task_id:
|
|
|
|
|
|
return
|
|
|
|
|
|
current = self._detail_windows.get(task_id)
|
|
|
|
|
|
if current is not None:
|
|
|
|
|
|
current.showNormal()
|
|
|
|
|
|
current.raise_()
|
|
|
|
|
|
current.activateWindow()
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
|
|
try:
|
|
|
|
|
|
detail = self._repository.get_task(task_id)
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
InfoBar.error(
|
|
|
|
|
|
title="无法打开任务详情",
|
|
|
|
|
|
content="无法读取本地任务,请检查数据库后重试。",
|
|
|
|
|
|
isClosable=True,
|
|
|
|
|
|
duration=-1,
|
|
|
|
|
|
position=InfoBarPosition.TOP_RIGHT,
|
|
|
|
|
|
parent=self._page,
|
|
|
|
|
|
)
|
|
|
|
|
|
return
|
|
|
|
|
|
if detail is None:
|
|
|
|
|
|
InfoBar.warning(
|
|
|
|
|
|
title="任务不存在",
|
|
|
|
|
|
content=f"本地找不到任务 {task_id},请刷新任务列表。",
|
|
|
|
|
|
isClosable=True,
|
|
|
|
|
|
duration=5000,
|
|
|
|
|
|
position=InfoBarPosition.TOP_RIGHT,
|
|
|
|
|
|
parent=self._page,
|
|
|
|
|
|
)
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
|
|
window = TaskDetailWindow(detail, self._page.window())
|
|
|
|
|
|
window.closed.connect(self._on_detail_closed)
|
|
|
|
|
|
self._detail_windows[task_id] = window
|
|
|
|
|
|
window.show()
|
|
|
|
|
|
window.raise_()
|
|
|
|
|
|
window.activateWindow()
|
|
|
|
|
|
|
|
|
|
|
|
@pyqtSlot(str)
|
|
|
|
|
|
def _on_detail_closed(self, task_id: str) -> None:
|
|
|
|
|
|
self._detail_windows.pop(task_id, None)
|
|
|
|
|
|
if not self._closing:
|
|
|
|
|
|
self._page.taskTable.setFocus()
|
|
|
|
|
|
|
2026-08-10 18:00:38 +08:00
|
|
|
|
@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()
|
|
|
|
|
|
|
2026-08-07 16:42:31 +08:00
|
|
|
|
@pyqtSlot()
|
|
|
|
|
|
def shutdown(self) -> None:
|
|
|
|
|
|
"""窗口关闭时停止新的领取,并等待已领取任务完成本地保存。"""
|
|
|
|
|
|
|
|
|
|
|
|
if self._closing:
|
|
|
|
|
|
return
|
|
|
|
|
|
self._closing = True
|
2026-08-09 10:54:18 +08:00
|
|
|
|
self._auto_fetch_running = False
|
|
|
|
|
|
self._stop_requested = True
|
|
|
|
|
|
self._next_cycle_timer.stop()
|
2026-08-10 09:08:54 +08:00
|
|
|
|
self._close_rerun_feedback()
|
2026-08-10 09:59:34 +08:00
|
|
|
|
self._close_claim_feedback()
|
2026-08-10 09:35:46 +08:00
|
|
|
|
self._close_device_feedback()
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
2026-08-08 09:46:58 +08:00
|
|
|
|
for window in list(self._detail_windows.values()):
|
|
|
|
|
|
window.close()
|
|
|
|
|
|
self._detail_windows.clear()
|
|
|
|
|
|
|
2026-08-07 16:42:31 +08:00
|
|
|
|
worker = self._claim_worker
|
|
|
|
|
|
thread = self._claim_thread
|
|
|
|
|
|
if worker is not None:
|
|
|
|
|
|
try:
|
|
|
|
|
|
worker.cancel()
|
|
|
|
|
|
except RuntimeError:
|
2026-08-10 18:00:38 +08:00
|
|
|
|
pass
|
|
|
|
|
|
try:
|
|
|
|
|
|
self._claimShutdownRequested.emit()
|
|
|
|
|
|
except RuntimeError:
|
|
|
|
|
|
pass
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
|
|
|
|
|
if thread is not None and thread.isRunning():
|
2026-08-07 17:38:29 +08:00
|
|
|
|
# uiautomator2/ADB 的单次调用可能需要数秒才返回。先通过
|
|
|
|
|
|
# cancelled 标志让采集在下一个安全点退出,再等待工作线程收尾,
|
|
|
|
|
|
# 避免窗口销毁时出现 "QThread destroyed while running"。
|
|
|
|
|
|
thread.wait(60_000)
|
2026-08-10 18:00:38 +08:00
|
|
|
|
self._claim_worker = None
|
|
|
|
|
|
self._claim_thread = None
|
2026-08-07 16:42:31 +08:00
|
|
|
|
|
2026-08-10 10:56:48 +08:00
|
|
|
|
resubmit_worker = self._resubmit_worker
|
|
|
|
|
|
resubmit_thread = self._resubmit_thread
|
|
|
|
|
|
if resubmit_worker is not None:
|
|
|
|
|
|
try:
|
|
|
|
|
|
resubmit_worker.cancel()
|
|
|
|
|
|
resubmit_worker.progress.disconnect(self._on_resubmit_progress)
|
|
|
|
|
|
resubmit_worker.finished.disconnect(self._on_resubmit_finished)
|
|
|
|
|
|
except (TypeError, RuntimeError):
|
|
|
|
|
|
pass
|
|
|
|
|
|
if resubmit_thread is not None and resubmit_thread.isRunning():
|
|
|
|
|
|
resubmit_thread.quit()
|
|
|
|
|
|
resubmit_thread.wait(60_000)
|
|
|
|
|
|
|
2026-08-10 17:23:24 +08:00
|
|
|
|
remove_worker = self._remove_worker
|
|
|
|
|
|
remove_thread = self._remove_thread
|
|
|
|
|
|
if remove_worker is not None:
|
|
|
|
|
|
try:
|
|
|
|
|
|
remove_worker.succeeded.disconnect(self._on_remove_succeeded)
|
|
|
|
|
|
remove_worker.failed.disconnect(self._on_remove_failed)
|
|
|
|
|
|
except (TypeError, RuntimeError):
|
|
|
|
|
|
pass
|
|
|
|
|
|
if remove_thread is not None and remove_thread.isRunning():
|
|
|
|
|
|
# 删除只包含一个很短的 SQLite 事务,不中途取消,避免部分更新。
|
|
|
|
|
|
remove_thread.quit()
|
|
|
|
|
|
remove_thread.wait(60_000)
|
|
|
|
|
|
|
2026-08-06 16:31:34 +08:00
|
|
|
|
|
|
|
|
|
|
def summary_to_row(summary: TaskSummary) -> TaskRow:
|
|
|
|
|
|
"""把领域摘要转换成只供表格显示的轻量行。"""
|
|
|
|
|
|
|
|
|
|
|
|
return TaskRow(
|
|
|
|
|
|
remote_task_id=summary.remote_task_id,
|
|
|
|
|
|
task_type=TASK_TYPE_TEXT[summary.task_type],
|
|
|
|
|
|
goods_id=summary.goods_id or "",
|
|
|
|
|
|
title=summary.title or "",
|
|
|
|
|
|
color=summary.target_color or "",
|
|
|
|
|
|
size=summary.target_size or "",
|
|
|
|
|
|
price_cents=summary.price_cent,
|
|
|
|
|
|
quantity=summary.quantity,
|
|
|
|
|
|
status=TASK_STATUS_TEXT[summary.status],
|
|
|
|
|
|
updated_at=summary.updated_at,
|
|
|
|
|
|
)
|