"""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**。 - “获取任务”会真的去操作手机采集商品;“搜索”只读本地数据库。 两者必须分开,不得共用入口。 - 普通成功不弹窗,更新界面即可;可恢复错误用 `InfoBar` (模板见 `docs/client/05-ui-specification.md` §9.1); 只有必须让用户当场做决定时才用模态对话框。 """ from typing import Callable, Dict, Optional from PyQt5.QtCore import ( QCoreApplication, QObject, Qt, QThread, QTimer, pyqtSignal, pyqtSlot, ) from qfluentwidgets import InfoBar, InfoBarPosition, MessageBox, PushButton from .android_device_service import ( AndroidDeviceSearchError, AndroidDeviceService, ) from .admin_gateway import ( AdminGatewayError, ClientInfo, AdminGateway, ) from .collect_task_service import CollectServiceFactory, CollectTaskService from .current_client_service import CurrentClientService, load_admin_base_url from .http_admin_gateway import HttpAdminGateway from .pdd_ui import PDDTaskPage, TaskRow from .pdd_device_service import ( PersistentPddDeviceService, bind_thread_device_service, ) from .purchase_task_service import ( LivePurchaseAdapterFactory, PurchaseAdapterFactory, PurchaseTaskService, ) from .purchase_reconcile_service import ( PurchaseReconcileFactory, PurchaseReconcileService, ) from .selected_android_device_service import SelectedAndroidDeviceService from .settings_repository import SettingsRepository from .task_models import ( OutboxEventType, OutboxStatus, TaskFilters, TaskRerunBatch, TaskStatus, TaskSummary, TaskType, ) from .task_repository import ( CollectRerunError, PurchaseRerunError, TaskRemovalError, TaskRepository, ) from .task_dispatcher import ( TaskDispatcher, admin_task_to_new_claimed_task, ) from .task_detail_view import TaskDetailWindow 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: "已取消", } 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 PURCHASE_RECONCILE_DELAY_MS = 5_000 def _reject_reconcile_only_purchase_adapter(*_args): """继续核单分支绝不允许创建采购 Adapter。""" raise RuntimeError("只读核单禁止创建采购执行器") class ClaimTaskWorker(QObject): """在固定后台线程中串行补交、领取或执行至多一条任务。""" noTask = pyqtSignal() taskSaved = pyqtSignal(str) duplicateTask = pyqtSignal(str) localSaveFailed = pyqtSignal(str, str) retryableFailed = pyqtSignal(str) failed = pyqtSignal(str) deviceUnavailable = pyqtSignal(str) taskStarted = pyqtSignal(str, str, int, int) outcome = pyqtSignal(str, str, str) completed = pyqtSignal() shutdownCompleted = pyqtSignal() def __init__( self, gateway: AdminGateway, task_repository: TaskRepository, client_service: CurrentClientService, android_device_service: SelectedAndroidDeviceService, collect_service_factory: Optional[CollectServiceFactory] = None, purchase_adapter_factory: Optional[PurchaseAdapterFactory] = None, live_purchase_adapter_factory: Optional[ LivePurchaseAdapterFactory ] = None, purchase_reconcile_factory: Optional[PurchaseReconcileFactory] = None, selected_task_id: str = "", device_connection_checker: Optional[Callable[[str], None]] = None, ) -> 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 self._collect_service_factory = collect_service_factory self._purchase_adapter_factory = purchase_adapter_factory self._live_purchase_adapter_factory = live_purchase_adapter_factory self._purchase_reconcile_factory = purchase_reconcile_factory self._selected_task_id = selected_task_id self._device_connection_checker = ( 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: """阻止尚未开始的领取;已领取的任务仍必须保存到本地。""" self._cancelled = True @pyqtSlot(object) def replace_gateway(self, gateway: AdminGateway) -> None: """当前命令结束后,为后续领取和提交切换 Admin Gateway。""" self._gateway = gateway @pyqtSlot() def run(self) -> None: """兼容旧调用;新代码通过 ``run_selected`` 提交命令。""" self.run_selected(self._selected_task_id) @pyqtSlot(object) def run_selected(self, command="") -> 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 command: 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, ) if isinstance(command, TaskRerunBatch): self._device_connection_checker(android_serial or "") self._execute_rerun_batch( command, client, android_serial or "" ) return selected_task_id = str(command or "") if selected_task_id: self._device_connection_checker(android_serial or "") self._task_repository.prepare_collect_rerun( selected_task_id ) 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 = ( "live" if android_serial and self._live_purchase_adapter_factory is not None else "dry_run" ) 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: if not self._cancelled: self.deviceUnavailable.emit(str(exc)) except AdminGatewayError as exc: if not self._cancelled: request_hint = ( f",请求编号:{exc.request_id}" if exc.request_id else "" ) message = f"{exc}{request_hint}" if exc.retryable: self.retryableFailed.emit(message) else: self.failed.emit(message) except Exception as exc: 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() def _execute_rerun_batch( self, command: TaskRerunBatch, client: ClientInfo, android_serial: str ) -> None: """在同一设备线程中逐条执行;采购不确定时停止剩余队列。""" for index, task_id in enumerate(command.task_ids): if self._cancelled: self.outcome.emit("cancelled", f"任务 {task_id} 尚未开始", task_id) continue try: if command.target_type is TaskType.COLLECT: self._task_repository.prepare_collect_rerun(task_id) service = CollectTaskService( self._gateway, self._task_repository, client, android_serial, cancelled=lambda: self._cancelled, started=lambda started_id, current=index + 1: self.taskStarted.emit( started_id, command.target_type.value, current, len(command.task_ids), ), collect_service_factory=self._collect_service_factory, ) result = service.execute_selected(task_id) elif task_id in command.reconcile_task_ids: result = self._execute_purchase_reconcile_retry( task_id, client, android_serial, current=index + 1, total=len(command.task_ids), ) else: result = self._execute_purchase_rerun( task_id, client, android_serial, current=index + 1, total=len(command.task_ids), ) except (CollectRerunError, PurchaseRerunError, ValueError) as exc: self.outcome.emit("skipped", str(exc), task_id) if command.target_type is TaskType.PURCHASE: for remaining_id in command.task_ids[index + 1 :]: self.outcome.emit( "skipped", f"前一采购任务未形成成功闭环,任务 {remaining_id} 未开始", remaining_id, ) break continue self.outcome.emit(result.kind, result.message, result.task_id or task_id) if command.target_type is TaskType.PURCHASE and result.kind != "succeeded": for remaining_id in command.task_ids[index + 1 :]: self.outcome.emit( "skipped", f"前一采购任务未形成成功闭环,任务 {remaining_id} 未开始", remaining_id, ) break def _execute_purchase_reconcile_retry( self, task_id: str, client: ClientInfo, android_serial: str, current: int = 1, total: int = 1, ): """恢复原不可逆运行并只读核单,绝不创建采购运行。""" if self._purchase_reconcile_factory is None: raise ValueError("采购核单执行器未就绪;绝不重新下单") self._task_repository.prepare_purchase_reconcile_retry(task_id) self.taskStarted.emit( task_id, TaskType.PURCHASE.value, current, total, ) result = PurchaseReconcileService( self._task_repository, android_serial, self._purchase_reconcile_factory, cancelled=lambda: False, ).execute_selected(task_id) if result.kind != "result_pending": return result event = self._task_repository.outbox_for_resubmit(task_id) if event is None: return type(result)( "manual_review", "核单结果未写入待上报队列", task_id ) submitter = PurchaseTaskService( self._gateway, self._task_repository, client, android_serial, _reject_reconcile_only_purchase_adapter, cancelled=lambda: False, ) return submitter.submit_saved_event(event) def _execute_purchase_rerun( self, task_id: str, client: ClientInfo, android_serial: str, current: int = 1, total: int = 1, ): """执行一次安全采购,并在真实下单后完成核单和 Admin 上报。""" task = self._task_repository.prepare_purchase_rerun(task_id) factory = ( self._live_purchase_adapter_factory if task.execution_mode == "live" else self._purchase_adapter_factory ) if factory is None: raise ValueError("采购执行器未就绪") service = PurchaseTaskService( self._gateway, self._task_repository, client, android_serial, factory, cancelled=lambda: self._cancelled, started=lambda started_id: self.taskStarted.emit( started_id, TaskType.PURCHASE.value, current, total, ), ) result = service.execute_selected(task_id) if result.kind != "reconcile_pending": return result # 已进入不可逆阶段后,即使用户停止,也必须完成当前任务的只读核单。 remaining_ms = PURCHASE_RECONCILE_DELAY_MS while remaining_ms > 0: step_ms = min(250, remaining_ms) QThread.msleep(step_ms) remaining_ms -= step_ms if self._purchase_reconcile_factory is None: return type(result)( "manual_review", "采购核单执行器未就绪;绝不重新下单", task_id ) reconcile = PurchaseReconcileService( self._task_repository, android_serial, self._purchase_reconcile_factory, cancelled=lambda: False, ).execute_selected(task_id) if reconcile.kind != "result_pending": return type(result)(reconcile.kind, reconcile.message, task_id) event = self._task_repository.outbox_for_resubmit(task_id) if event is None: return type(result)("manual_review", "核单结果未写入待上报队列", task_id) return service.submit_saved_event(event) @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,不执行任何手机操作。""" 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: event = self._repository.outbox_for_resubmit(task_id) 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: if event.event_type is OutboxEventType.TASK_FAILURE: receipt = self._gateway.submit_failure( task_id, event.idempotency_key, event.payload_json, ) else: receipt = self._gateway.submit_result( task_id, event.idempotency_key, event.payload_json, ) if not receipt.accepted: raise AdminGatewayError( "ADMIN_RESULT_NOT_ACCEPTED", "Admin 未确认接收任务结果", False, ) except AdminGatewayError as exc: 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() 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() class PDDTaskPageEvent(QObject): """把 PDD 页面只读操作连接到本地任务 Repository。""" _claimRunRequested = pyqtSignal(object) _claimReleaseRequested = pyqtSignal() _claimShutdownRequested = pyqtSignal() _claimGatewayChanged = pyqtSignal(object) def __init__( self, page: PDDTaskPage, repository: Optional[TaskRepository] = None, parent=None, claim_gateway: Optional[AdminGateway] = None, settings_repository: Optional[SettingsRepository] = None, collect_service_factory: Optional[CollectServiceFactory] = None, purchase_adapter_factory: Optional[PurchaseAdapterFactory] = None, live_purchase_adapter_factory: Optional[ LivePurchaseAdapterFactory ] = None, purchase_reconcile_factory: Optional[PurchaseReconcileFactory] = None, device_connection_checker: Optional[Callable[[str], None]] = None, next_task_delay_ms: int = 500, no_task_delay_ms: int = 5_000, retry_delays_ms: tuple[int, ...] = (5_000, 10_000, 20_000, 30_000), ): super().__init__(parent or page) self._page = page self._repository = repository or TaskRepository() self._filters = TaskFilters() self._closing = False 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_target_type = TaskType.COLLECT self._rerun_total = 0 self._rerun_succeeded = 0 self._rerun_failed = 0 self._rerun_skipped = 0 self._rerun_feedback: Optional[InfoBar] = None self._claim_feedback: Optional[InfoBar] = None self._device_feedback: Optional[InfoBar] = None self._auto_fetch_running = False self._resubmit_busy = False self._resubmit_thread: Optional[QThread] = None self._resubmit_worker: Optional[ResultResubmitWorker] = None self._remove_busy = False self._remove_thread: Optional[QThread] = None self._remove_worker: Optional[TaskRemoveWorker] = None self._stop_requested = False self._cycle_next_delay_ms: Optional[int] = None self._stop_status = "自动获取:已停止 · 当前没有执行中的任务" self._retry_count = 0 self._detail_windows: Dict[str, TaskDetailWindow] = {} self._collect_service_factory = collect_service_factory self._purchase_adapter_factory = purchase_adapter_factory self._live_purchase_adapter_factory = live_purchase_adapter_factory self._purchase_reconcile_factory = purchase_reconcile_factory self._device_connection_checker = ( device_connection_checker or AndroidDeviceService().require_connected ) 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) try: self._repository.recover_interrupted_work() except AttributeError: # 测试用的只读 Repository 可以不实现恢复接口。 pass settings = settings_repository or SettingsRepository() self._settings_repository = settings self._client_service = CurrentClientService(settings) self._selected_android_device_service = SelectedAndroidDeviceService( settings ) self._claim_gateway = claim_gateway self._owns_claim_gateway = claim_gateway is None self._claim_gateway_error = "" if self._claim_gateway is None: base_url = load_admin_base_url(settings) 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, client_id=self._client_service.load().client_id, ) except (TypeError, ValueError) as exc: self._claim_gateway_error = str(exc) page.searchRequested.connect(self.search_tasks) page.refreshRequested.connect(self.refresh_tasks) page.rerunRequested.connect(self.request_rerun) page.purchaseRerunRequested.connect(self.request_purchase_rerun) page.rerunCancelRequested.connect(self.request_cancel_rerun) page.resubmitRequested.connect(self.request_resubmit) page.removeRequested.connect(self.request_remove) page.autoFetchRequested.connect(self._request_claim_task) page.detailRequested.connect(self.show_task_detail) page.taskModel.loadMoreRequested.connect(self._load_page) page.set_auto_fetch_state("stopped") page.destroyed.connect(self.shutdown) application = QCoreApplication.instance() if application is not None: application.aboutToQuit.connect(self.shutdown) @pyqtSlot(str) def update_admin_base_url(self, base_url: str) -> None: """让下一次领取、执行和上报统一使用新的 Admin 地址。""" if self._closing or not self._owns_claim_gateway: return timeout_value = self._settings_repository.get( "admin.request_timeout_seconds", 3.0 ) try: timeout_seconds = float(timeout_value) gateway = HttpAdminGateway( base_url, timeout_seconds=timeout_seconds, client_id=self._client_service.load().client_id, ) except (TypeError, ValueError) as exc: self._claim_gateway_error = str(exc) return self._claim_gateway = gateway self._claim_gateway_error = "" if self._claim_worker is not None: self._claimGatewayChanged.emit(gateway) 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() @pyqtSlot(object) def request_rerun(self, task_ids) -> None: """预检并确认批量重新采集。""" self._request_rerun_batch(task_ids, TaskType.COLLECT) @pyqtSlot(object) def request_purchase_rerun(self, task_ids) -> None: """预检并确认带不可逆门禁的批量重新采购。""" self._request_rerun_batch(task_ids, TaskType.PURCHASE) def _request_rerun_batch(self, task_ids, target_type: TaskType) -> None: values = (task_ids,) if isinstance(task_ids, str) else task_ids stable_ids = tuple(dict.fromkeys(str(value) for value in values if value)) 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 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: self._show_device_unavailable( "请先在设置页选择并保存 Android 设备" ) return try: plan = self._repository.plan_rerun_batch(stable_ids, target_type) except (CollectRerunError, PurchaseRerunError, ValueError) as exc: self._show_rerun_warning("不能重新执行", str(exc)) return except Exception: self._show_rerun_warning( "不能重新执行", "无法读取任务状态,请检查数据库后重试。" ) return action = "采集" if target_type is TaskType.COLLECT else "采购" processable_task_ids = ( plan.reconcile_task_ids + plan.eligible_task_ids ) if not processable_task_ids: self._show_rerun_warning( f"没有可重新{action}的任务", f"已选 {plan.selected_count} 条,过滤 {len(plan.filtered_task_ids)} 条其他类型," f"安全条件阻止 {len(plan.blocked)} 条。", ) return reconcile_count = len(plan.reconcile_task_ids) purchase_count = len(plan.eligible_task_ids) if target_type is TaskType.PURCHASE: title = f"确认处理 {len(processable_task_ids)} 条采购任务?" detail = ( f"已选 {plan.selected_count} 条;可重新采购 {purchase_count} 条;" f"只继续核单 {reconcile_count} 条;" f"过滤其他类型 {len(plan.filtered_task_ids)} 条;" f"安全条件阻止 {len(plan.blocked)} 条。\n\n" "已提交订单的任务只读取待付款订单,绝不重新下单;" "只读核单优先执行,结果不确定时立即停止剩余任务。" ) confirm_text = ( "继续核单" if reconcile_count and not purchase_count else "继续处理" ) else: title = f"确认重新{action} {len(plan.eligible_task_ids)} 条任务?" detail = ( f"已选 {plan.selected_count} 条;可执行 {len(plan.eligible_task_ids)} 条;" f"过滤其他类型 {len(plan.filtered_task_ids)} 条;" f"安全条件阻止 {len(plan.blocked)} 条。\n\n" "新结果会覆盖 Client 和 Admin 当前采集数据,旧执行记录仍保留。" ) confirm_text = f"重新{action}" dialog = MessageBox( title, detail, self._page.window(), ) dialog.yesButton.setText(confirm_text) dialog.cancelButton.setText("取消") dialog.cancelButton.setFocus() if not dialog.exec(): return self._start_rerun_worker( TaskRerunBatch( target_type, processable_task_ids, plan.reconcile_task_ids, ) ) @pyqtSlot(object) def request_resubmit(self, task_ids) -> None: """确认后在后台重发勾选任务尚未发送或最新的 Outbox。""" stable_ids = tuple( dict.fromkeys(str(value) for value in task_ids if value) ) 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 if self._claim_gateway_error or self._claim_gateway is None: self._show_rerun_warning( "重新上报不可用", self._claim_gateway_error or "Admin 提交服务未初始化", ) return dialog = MessageBox( f"重新上报 {len(stable_ids)} 条任务数据?", "会优先提交尚未发送的结果或失败信息;没有待发送数据时," "才重新提交本地已经保存的最新结果。" "不会重新采集、采购或操作 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) @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) def _start_resubmit_worker(self, task_ids: tuple[str, ...]) -> None: """启动只处理指定任务 Outbox 的工作线程。""" assert self._claim_gateway is not None self._resubmit_busy = True self._page.set_resubmit_running(True) self._page.set_engine_status( f"正在准备重新上报 {len(task_ids)} 条任务数据…" ) 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) def _start_rerun_worker(self, command: TaskRerunBatch) -> None: """把指定任务提交到持久设备工作线程。""" assert self._claim_gateway is not None self._close_rerun_feedback() self._close_claim_feedback() self._close_device_feedback() self._rerun_cancel_requested = False self._rerun_target_type = command.target_type self._rerun_total = len(command.task_ids) self._rerun_succeeded = self._rerun_failed = self._rerun_skipped = 0 self._claim_busy = True self._page.set_rerun_running(True) self._page.set_engine_status( self._rerun_queue_status(command) ) self._reload() worker = self._ensure_claim_executor() self._claim_operation = "rerun" worker.prepare_run() self._claimRunRequested.emit(command) @staticmethod def _rerun_queue_status(command: TaskRerunBatch) -> str: total = len(command.task_ids) if command.target_type is TaskType.COLLECT: action = "重新采集" elif len(command.reconcile_task_ids) == total: action = "继续核单" elif command.reconcile_task_ids: action = "处理采购任务" else: action = "重新采购" return ( f"已加入队列 {total} 条,正在检查 Android 设备," f"等待开始{action}…" ) @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() @pyqtSlot(str, str, str) def _on_rerun_outcome(self, kind: str, message: str, _task_id: str) -> None: if self._closing: return self._reload() if self._rerun_cancel_requested: self._page.set_engine_status( "停止请求已处理,请查看任务最新状态" ) return if kind == "succeeded": self._rerun_succeeded += 1 elif kind in {"skipped", "cancelled"}: self._rerun_skipped += 1 else: self._rerun_failed += 1 self._page.set_engine_status(message) @pyqtSlot(str, str, int, int) def _on_rerun_task_started( self, task_id: str, task_type: str, current: int, total: int, ) -> None: """数据库进入 running 后,立即让主线程刷新当前任务。""" if self._closing or self._claim_operation != "rerun": return action = "采集" if task_type == TaskType.COLLECT.value else "采购" self._reload() self._page.set_engine_status( f"批量重新{action}:正在{action} {task_id}({current}/{total})" ) @pyqtSlot(str) def _on_rerun_failed(self, message: str) -> None: if self._closing: return if self._rerun_cancel_requested: return self._page.set_engine_status(message) self._show_rerun_error("重新采集失败", message) @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) @pyqtSlot() def _on_rerun_thread_finished(self) -> None: cancel_requested = self._rerun_cancel_requested self._claim_busy = False self._rerun_cancel_requested = False if not self._closing: self._page.set_rerun_running(False) if cancel_requested: self._reload() self._page.set_engine_status( "停止请求已处理,请查看任务最新状态" ) elif self._rerun_succeeded + self._rerun_failed + self._rerun_skipped: action = "采集" if self._rerun_target_type is TaskType.COLLECT else "采购" summary = ( f"批量重新{action}完成:成功 {self._rerun_succeeded} 条," f"失败 {self._rerun_failed} 条,跳过 {self._rerun_skipped} 条" ) self._page.set_engine_status(summary) if self._rerun_failed or self._rerun_skipped: self._show_rerun_warning(f"重新{action}未全部完成", summary) def _show_rerun_warning(self, title: str, content: str) -> None: bar = InfoBar.warning( title=title, content=content, isClosable=True, duration=5000, position=InfoBarPosition.TOP, parent=self._page, ) self._replace_rerun_feedback(bar) def _show_rerun_error(self, title: str, content: str) -> None: """显示一条较大、可明确关闭且会自动消失的重新采集错误。""" bar = InfoBar.error( title=title, content=content, isClosable=True, duration=ERROR_FEEDBACK_DURATION_MS, position=InfoBarPosition.TOP, parent=self._page, ) self._replace_rerun_feedback(bar) def _replace_rerun_feedback(self, bar: InfoBar) -> None: """用新提示替换旧提示,避免连续失败后堆叠。""" self._close_rerun_feedback() self._rerun_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_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 @pyqtSlot() def _request_claim_task(self) -> None: """切换持续自动获取;每一轮仍只启动一个后台 Worker。""" if self._closing: return if self._resubmit_busy or self._remove_busy: self._show_rerun_warning( "暂时不能获取任务", "本地任务正在重新上报或删除,请等待当前操作结束。", ) return if self._auto_fetch_running: self._request_stop_auto_fetch() 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 self._close_device_feedback() 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: """开始一轮任务处理;设备命令由固定工作线程串行执行。""" if ( self._closing or not self._auto_fetch_running or self._stop_requested or self._claim_busy ): return self._claim_busy = True self._cycle_next_delay_ms = None 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, 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, live_purchase_adapter_factory=self._live_purchase_adapter_factory, purchase_reconcile_factory=self._purchase_reconcile_factory, device_connection_checker=self._device_connection_checker, ) worker.moveToThread(thread) self._claimRunRequested.connect(worker.run_selected) self._claimReleaseRequested.connect(worker.release_device) self._claimShutdownRequested.connect(worker.shutdown_worker) self._claimGatewayChanged.connect(worker.replace_gateway) worker.taskSaved.connect(self._on_claimed_task_saved) worker.retryableFailed.connect(self._route_claim_retryable_failed) worker.failed.connect(self._route_claim_failed) worker.deviceUnavailable.connect(self._route_device_unavailable) worker.taskStarted.connect(self._on_rerun_task_started) 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) 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 在安全点停止。""" self._auto_fetch_running = False self._stop_requested = True self._next_cycle_timer.stop() self._stop_status = "自动获取:已停止 · 当前没有执行中的任务" worker = self._claim_worker if worker is not None and self._claim_busy: 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 if self._claim_worker is not None: self._claimReleaseRequested.emit() 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") @pyqtSlot() def _on_no_claimed_task(self) -> None: if not self._closing: self._retry_count = 0 seconds = self._no_task_delay_ms / 1000 self._continue_after( self._no_task_delay_ms, f"自动获取:运行中 · 暂无任务,{seconds:g} 秒后再次领取", ) @pyqtSlot(str) def _on_claimed_task_saved(self, task_id: str) -> None: if self._closing: return self._page.set_engine_status( f"已领取任务 {task_id},已保存到本地任务列表" ) self._cycle_next_delay_ms = self._next_task_delay_ms self._reload() @pyqtSlot(str) def _on_duplicate_claimed_task(self, task_id: str) -> None: if not self._closing: self._continue_after( self._next_task_delay_ms, f"任务 {task_id} 本地已有,未重复保存", ) @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) 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) @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) self._stop_after_current(content) @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) @pyqtSlot(str, str, str) def _on_collect_outcome(self, kind: str, message: str, task_id: str) -> None: if self._closing: return self._reload() if kind == "no_task": self._on_no_claimed_task() elif kind in {"succeeded", "business_failed", "task_failed"}: self._retry_count = 0 self._continue_after( self._next_task_delay_ms, f"自动获取:运行中 · {message}", ) elif kind == "result_pending": self._on_claim_retryable_failed(message) 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} 秒后自动核对订单" ), ) elif kind in {"global_failed", "manual_review", "failed"}: self._show_claim_error("任务需要处理", message) 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) def _show_retry_paused(self, content: str) -> None: """用持久警告说明任务不会自行倒计时重试。""" InfoBar.warning( title="重试已暂停", content=content, isClosable=True, duration=-1, position=InfoBarPosition.TOP_RIGHT, parent=self._page, ) def _show_claim_error(self, title: str, content: str) -> None: """显示较大的临时错误,同时在底部保留完整状态。""" self._close_claim_feedback() bar = InfoBar.error( title=title, content=content, isClosable=True, duration=ERROR_FEEDBACK_DURATION_MS, position=InfoBarPosition.TOP, parent=self._page, ) 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 def _show_device_unavailable(self, content: str) -> None: """显示一条可进入设置且不会堆叠的设备错误。""" self._close_device_feedback() bar = InfoBar.error( title="Android 设备不可用", content=content, isClosable=True, duration=ERROR_FEEDBACK_DURATION_MS, position=InfoBarPosition.TOP, parent=self._page, ) self._device_feedback = bar self._set_large_error_feedback_size(bar) settings_button = PushButton("打开设置", bar) self._set_large_feedback_button_size(settings_button) settings_button.setAccessibleName("打开 Android 设备设置") settings_button.clicked.connect( lambda _checked=False, current=bar: self._open_device_settings(current) ) close_button = PushButton("关闭提示", bar) self._set_large_feedback_button_size(close_button) 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 @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, ) @pyqtSlot() def _on_claim_thread_finished(self) -> None: self._claim_busy = False 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) 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) @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() @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: """窗口关闭时停止新的领取,并等待已领取任务完成本地保存。""" if self._closing: return self._closing = True self._auto_fetch_running = False self._stop_requested = True self._next_cycle_timer.stop() self._close_rerun_feedback() self._close_claim_feedback() self._close_device_feedback() for window in list(self._detail_windows.values()): window.close() self._detail_windows.clear() worker = self._claim_worker thread = self._claim_thread if worker is not None: try: worker.cancel() except RuntimeError: pass try: self._claimShutdownRequested.emit() except RuntimeError: pass if thread is not None and thread.isRunning(): # 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 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) 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) def summary_to_row(summary: TaskSummary) -> TaskRow: """把领域摘要转换成只供表格显示的轻量行。""" status_text = TASK_STATUS_TEXT[summary.status] if summary.status is TaskStatus.RUNNING: status_text = "采集中" if summary.task_type is TaskType.COLLECT else "采购中" elif summary.status in {TaskStatus.RETRY_WAIT, TaskStatus.FAILED}: status_text = ( "采集失败" if summary.task_type is TaskType.COLLECT else "采购失败" ) 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 "", shop_name=summary.shop_name or "", color=summary.target_color or "", size=summary.target_size or "", price_cents=summary.price_cent, quantity=summary.quantity, status=status_text, latest_run_status=( summary.latest_run_status.value if summary.latest_run_status else "" ), duration_seconds=summary.duration_seconds, order_no=summary.order_no or "", updated_at=summary.updated_at, )