"""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 Dict, Optional from PyQt5.QtCore import ( QCoreApplication, QObject, QThread, QTimer, pyqtSignal, pyqtSlot, ) from qfluentwidgets import InfoBar, InfoBarPosition, MessageBox from .admin_gateway import ( AdminGatewayError, AdminTask, ClientInfo, AdminGateway, ) from .collect_task_service import CollectServiceFactory, CollectTaskService from .current_client_service import CurrentClientService from .http_admin_gateway import DEFAULT_ADMIN_BASE_URL, HttpAdminGateway from .pdd_ui import PDDTaskPage, TaskRow from .selected_android_device_service import SelectedAndroidDeviceService from .settings_repository import SettingsRepository from .task_models import ( NewClaimedTask, TaskFilters, TaskStatus, TaskSummary, TaskType, ) from .task_repository import CollectRerunError, TaskRepository 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: "已取消", } def admin_task_to_new_claimed_task(task: AdminTask) -> NewClaimedTask: """把 Admin 任务显式映射为本地任务,避免字段名自动展开出错。""" if task.task_type is not TaskType.COLLECT: raise ValueError(f"任务 {task.task_id} 不是采集任务") payload = dict(task.payload) goods_url = payload.get("goods_url") goods_id = payload.get("goods_id") if not isinstance(goods_url, str) or not goods_url.strip(): raise ValueError("payload.goods_url 不能为空") if goods_id is not None and not isinstance(goods_id, str): raise ValueError("payload.goods_id 必须是文本") original_task = { "id": task.task_id, "type": task.task_type.value, "version": task.version, "priority": task.priority, "payload": payload, "created_at": task.created_at, "updated_at": task.updated_at, } return NewClaimedTask( remote_task_id=task.task_id, task_type=task.task_type, goods_url=goods_url, goods_id=goods_id.strip() if isinstance(goods_id, str) else None, priority=task.priority, version=task.version, admin_payload=original_task, ) class ClaimTaskWorker(QObject): """在后台补交或执行至多一条采集任务。""" noTask = pyqtSignal() taskSaved = pyqtSignal(str) duplicateTask = pyqtSignal(str) localSaveFailed = pyqtSignal(str, str) retryableFailed = pyqtSignal(str) failed = pyqtSignal(str) outcome = pyqtSignal(str, str, str) completed = pyqtSignal() def __init__( self, gateway: AdminGateway, task_repository: TaskRepository, client_service: CurrentClientService, android_device_service: SelectedAndroidDeviceService, collect_service_factory: Optional[CollectServiceFactory] = None, selected_task_id: str = "", ) -> 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._selected_task_id = selected_task_id def cancel(self) -> None: """阻止尚未开始的领取;已领取的任务仍必须保存到本地。""" self._cancelled = True @pyqtSlot() def run(self) -> None: try: if self._cancelled: return client_settings = self._client_service.load() if not client_settings.client_id: self.failed.emit("请先在设置页保存当前设备号和设备名") return android_serial = self._android_device_service.load() service = CollectTaskService( self._gateway, self._task_repository, ClientInfo( client_settings.client_id, client_settings.client_name, ), android_serial or "", cancelled=lambda: self._cancelled, collect_service_factory=self._collect_service_factory, ) result = ( service.execute_selected(self._selected_task_id) if self._selected_task_id else service.execute_one() ) if not self._cancelled or result.kind == "cancelled": self.outcome.emit(result.kind, result.message, result.task_id) 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: self.completed.emit() class PDDTaskPageEvent(QObject): """把 PDD 页面只读操作连接到本地任务 Repository。""" 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, 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._auto_fetch_running = False 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 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._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, 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.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) 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(str) def request_rerun(self, task_id: str) -> None: """确认后在后台重新采集当前选中的一条终态任务。""" if self._closing or not task_id: return if self._auto_fetch_running or self._claim_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_rerun_warning( "重新执行不可用", "请先在设置页选择并保存 Android 设备。" ) 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("重新采集") dialog.cancelButton.setText("取消") dialog.cancelButton.setFocus() if not dialog.exec(): return try: self._repository.prepare_collect_rerun(task_id) except (CollectRerunError, ValueError) as exc: self._show_rerun_warning("不能重新执行", str(exc)) self._reload() return except Exception: self._show_rerun_warning( "重新执行失败", "无法更新本地任务状态,请检查数据库后重试。" ) return self._start_rerun_worker(task_id) def _start_rerun_worker(self, task_id: str) -> None: """启动只处理指定任务的工作线程。""" assert self._claim_gateway is not None self._claim_busy = True self._page.set_rerun_running(True) self._page.set_engine_status(f"正在重新采集任务 {task_id}…") self._reload() thread = QThread(self) worker = ClaimTaskWorker( self._claim_gateway, self._repository, self._client_service, self._selected_android_device_service, self._collect_service_factory, selected_task_id=task_id, ) worker.moveToThread(thread) thread.started.connect(worker.run) worker.retryableFailed.connect(self._on_rerun_failed) worker.failed.connect(self._on_rerun_failed) worker.outcome.connect(self._on_rerun_outcome) worker.completed.connect(thread.quit) worker.completed.connect(worker.deleteLater) thread.finished.connect(thread.deleteLater) thread.finished.connect(self._on_rerun_thread_finished) self._claim_thread = thread self._claim_worker = worker thread.start() @pyqtSlot(str, str, str) def _on_rerun_outcome(self, kind: str, message: str, _task_id: str) -> None: if self._closing: return self._reload() self._page.set_engine_status(message) if kind != "succeeded": self._show_claim_error("重新采集需要处理", message) @pyqtSlot(str) def _on_rerun_failed(self, message: str) -> None: if self._closing: return self._page.set_engine_status(message) self._show_claim_error("重新采集失败", message) @pyqtSlot() def _on_rerun_thread_finished(self) -> None: self._claim_worker = None self._claim_thread = None self._claim_busy = False if not self._closing: self._page.set_rerun_running(False) def _show_rerun_warning(self, title: str, content: str) -> None: InfoBar.warning( title=title, content=content, isClosable=True, duration=5000, position=InfoBarPosition.TOP_RIGHT, parent=self._page, ) @pyqtSlot() def _request_claim_task(self) -> None: """切换持续自动获取;每一轮仍只启动一个后台 Worker。""" if self._closing: 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._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("自动获取:运行中 · 正在处理一条采集任务…") thread = QThread(self) worker = ClaimTaskWorker( self._claim_gateway, self._repository, self._client_service, self._selected_android_device_service, self._collect_service_factory, ) worker.moveToThread(thread) thread.started.connect(worker.run) worker.noTask.connect(self._on_no_claimed_task) worker.taskSaved.connect(self._on_claimed_task_saved) worker.duplicateTask.connect(self._on_duplicate_claimed_task) worker.localSaveFailed.connect(self._on_claimed_task_save_failed) worker.retryableFailed.connect(self._on_claim_retryable_failed) worker.failed.connect(self._on_claim_failed) worker.outcome.connect(self._on_collect_outcome) worker.completed.connect(thread.quit) worker.completed.connect(worker.deleteLater) thread.finished.connect(thread.deleteLater) thread.finished.connect(self._on_claim_thread_finished) self._claim_thread = thread self._claim_worker = worker thread.start() 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: 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 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, 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 == "succeeded": 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 == "failed" and self._is_retry_wait_task(task_id): content = ( f"自动获取:已停止 · 任务 {task_id} 重试已暂停;" "请选择该任务点击“重新执行”,或重新启动获取任务" ) self._show_retry_paused(content) self._stop_after_current(content) elif kind in {"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 _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 return task is not None and task.status is TaskStatus.RETRY_WAIT 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: """显示不会自动消失的可恢复错误,同时保留底部状态文字。""" InfoBar.error( title=title, content=content, isClosable=True, duration=-1, position=InfoBarPosition.TOP_RIGHT, parent=self._page, ) @pyqtSlot() def _on_claim_thread_finished(self) -> None: self._claim_worker = None self._claim_thread = 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() 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() 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() signal_slots = ( (worker.noTask, self._on_no_claimed_task), (worker.taskSaved, self._on_claimed_task_saved), (worker.duplicateTask, self._on_duplicate_claimed_task), ( worker.localSaveFailed, self._on_claimed_task_save_failed, ), (worker.failed, self._on_claim_failed), ( worker.retryableFailed, self._on_claim_retryable_failed, ), (worker.outcome, self._on_collect_outcome), ) except RuntimeError: signal_slots = () for signal, slot in signal_slots: try: signal.disconnect(slot) except (TypeError, RuntimeError): pass if thread is not None and thread.isRunning(): thread.quit() # uiautomator2/ADB 的单次调用可能需要数秒才返回。先通过 # cancelled 标志让采集在下一个安全点退出,再等待工作线程收尾, # 避免窗口销毁时出现 "QThread destroyed while running"。 thread.wait(60_000) 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, )