diff --git a/client/src/collect_task_service.py b/client/src/collect_task_service.py index 13bb58a..255c444 100644 --- a/client/src/collect_task_service.py +++ b/client/src/collect_task_service.py @@ -48,6 +48,7 @@ class CollectTaskService: device_address: str, *, cancelled: Callable[[], bool] = lambda: False, + started: Callable[[str], None] = lambda _task_id: None, collect_service_factory: Optional[CollectServiceFactory] = None, ) -> None: self._gateway = gateway @@ -55,6 +56,7 @@ class CollectTaskService: self._client = client self._device_address = device_address self._cancelled = cancelled + self._started = started self._factory = collect_service_factory or self._default_factory def execute_one(self) -> CollectTaskOutcome: @@ -132,6 +134,7 @@ class CollectTaskService: started = self._repository.start_collect_run( remote_task_id, self._device_address ) + self._started(remote_task_id) if self._cancelled(): event = self._repository.save_collect_failure( remote_task_id, diff --git a/client/src/pdd_ui_event.py b/client/src/pdd_ui_event.py index 9a13ad9..f6de38f 100644 --- a/client/src/pdd_ui_event.py +++ b/client/src/pdd_ui_event.py @@ -133,6 +133,7 @@ class ClaimTaskWorker(QObject): retryableFailed = pyqtSignal(str) failed = pyqtSignal(str) deviceUnavailable = pyqtSignal(str) + taskStarted = pyqtSignal(str, str, int, int) outcome = pyqtSignal(str, str, str) completed = pyqtSignal() shutdownCompleted = pyqtSignal() @@ -309,11 +310,23 @@ class ClaimTaskWorker(QObject): 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) else: - result = self._execute_purchase_rerun(task_id, client, android_serial) + 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) continue @@ -328,7 +341,12 @@ class ClaimTaskWorker(QObject): break def _execute_purchase_rerun( - self, task_id: str, client: ClientInfo, android_serial: str + self, + task_id: str, + client: ClientInfo, + android_serial: str, + current: int = 1, + total: int = 1, ): """执行一次安全采购,并在真实下单后完成核单和 Admin 上报。""" @@ -347,6 +365,12 @@ class ClaimTaskWorker(QObject): 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": @@ -976,8 +1000,8 @@ class PDDTaskPageEvent(QObject): self._claim_busy = True self._page.set_rerun_running(True) self._page.set_engine_status( - f"正在检查 Android 设备并准备重新{'采集' if command.target_type is TaskType.COLLECT else '采购'}" - f" {len(command.task_ids)} 条任务…" + f"已加入队列 {len(command.task_ids)} 条,正在检查 Android 设备,等待开始重新" + f"{'采集' if command.target_type is TaskType.COLLECT else '采购'}…" ) self._reload() @@ -1025,6 +1049,24 @@ class PDDTaskPageEvent(QObject): 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: @@ -1203,6 +1245,7 @@ class PDDTaskPageEvent(QObject): 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 完成释放后 @@ -1704,6 +1747,10 @@ class PDDTaskPageEvent(QObject): 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 "采购中" + return TaskRow( remote_task_id=summary.remote_task_id, task_type=TASK_TYPE_TEXT[summary.task_type], @@ -1714,7 +1761,7 @@ def summary_to_row(summary: TaskSummary) -> TaskRow: size=summary.target_size or "", price_cents=summary.price_cent, quantity=summary.quantity, - status=TASK_STATUS_TEXT[summary.status], + status=status_text, latest_run_status=( summary.latest_run_status.value if summary.latest_run_status else "" ), diff --git a/client/src/purchase_task_service.py b/client/src/purchase_task_service.py index d15bdd3..5d45dc0 100644 --- a/client/src/purchase_task_service.py +++ b/client/src/purchase_task_service.py @@ -69,6 +69,7 @@ class PurchaseTaskService: adapter_factory: PurchaseAdapterFactory, *, cancelled: Callable[[], bool] = lambda: False, + started: Callable[[str], None] = lambda _task_id: None, ) -> None: self._gateway = gateway self._repository = repository @@ -76,6 +77,7 @@ class PurchaseTaskService: self._device_address = str(device_address or "").strip() self._factory = adapter_factory self._cancelled = cancelled + self._started = started def execute_one_local(self) -> PurchaseTaskOutcome: """先补交结果,再执行最早的本地待执行采购任务。""" @@ -103,6 +105,7 @@ class PurchaseTaskService: started = self._repository.start_purchase_run( remote_task_id, self._device_address ) + self._started(remote_task_id) adapter: PddPurchaseAdapter | None = None step = "purchase_prepare" business_failure_message = "" diff --git a/client/test/test_collect_task_service.py b/client/test/test_collect_task_service.py index 107908c..2886fae 100644 --- a/client/test/test_collect_task_service.py +++ b/client/test/test_collect_task_service.py @@ -207,6 +207,30 @@ class CollectTaskServiceTest(unittest.TestCase): self.repository.get_task("COL-OTHER").status, TaskStatus.CLAIMED ) + def test_started_callback_runs_after_database_enters_running(self): + self.repository.add_claimed_task( + NewClaimedTask( + remote_task_id="COL-STARTED", + task_type=TaskType.COLLECT, + goods_url="https://example.test/COL-STARTED", + ) + ) + observed = [] + service = CollectTaskService( + self.gateway, + self.repository, + self.client, + "USB-001", + started=lambda task_id: observed.append( + (task_id, self.repository.get_task(task_id).status) + ), + collect_service_factory=lambda *_args: FakeCollector([]), + ) + + service.execute_selected("COL-STARTED") + + self.assertEqual(observed, [("COL-STARTED", TaskStatus.RUNNING)]) + def test_selected_task_cancelled_before_collector_starts_is_persisted(self): self.repository.add_claimed_task( NewClaimedTask( diff --git a/client/test/test_pdd_ui_event.py b/client/test/test_pdd_ui_event.py index 74f79b6..e6ce4ec 100644 --- a/client/test/test_pdd_ui_event.py +++ b/client/test/test_pdd_ui_event.py @@ -1094,6 +1094,37 @@ class PDDTaskPageEventTest(unittest.TestCase): self.assertEqual(row.goods_id, "") self.assertFalse(hasattr(row, "pdd_data")) + def test_summary_to_row_uses_task_specific_running_text(self): + common = dict( + id=1, + goods_id=None, + title=None, + target_color=None, + target_size=None, + price_cent=None, + quantity=None, + status=TaskStatus.RUNNING, + updated_at="2026-08-11T08:00:00Z", + ) + + collect = summary_to_row( + TaskSummary( + remote_task_id="COL-RUNNING", + task_type=TaskType.COLLECT, + **common, + ) + ) + purchase = summary_to_row( + TaskSummary( + remote_task_id="PUR-RUNNING", + task_type=TaskType.PURCHASE, + **common, + ) + ) + + self.assertEqual(collect.status, "采集中") + self.assertEqual(purchase.status, "采购中") + def test_main_window_keeps_event_object_alive(self): window = MainWindow( task_repository=self.repository, diff --git a/client/test/test_purchase_task_service.py b/client/test/test_purchase_task_service.py index 5336072..f7a8605 100644 --- a/client/test/test_purchase_task_service.py +++ b/client/test/test_purchase_task_service.py @@ -217,6 +217,24 @@ class PurchaseTaskServiceTest(unittest.TestCase): self.assertFalse(purchase["order_submitted"]) self.assertFalse(purchase["payment_attempted"]) + def test_started_callback_runs_after_database_enters_running(self): + self._prepare_task(task_id="PUR-STARTED") + observed = [] + service = PurchaseTaskService( + self.gateway, + self.repository, + self.client, + "USB-001", + lambda _address, _cancelled: RecordingDryRunAdapter(), + started=lambda task_id: observed.append( + (task_id, self.repository.get_task(task_id).status) + ), + ) + + service.execute_selected("PUR-STARTED") + + self.assertEqual(observed, [("PUR-STARTED", TaskStatus.RUNNING)]) + def test_price_above_limit_stops_before_confirmation(self): self._prepare_task() adapter = RecordingDryRunAdapter(price_cent=5001) diff --git a/docs/client/05-ui-specification.md b/docs/client/05-ui-specification.md index 398e673..6669ad8 100644 --- a/docs/client/05-ui-specification.md +++ b/docs/client/05-ui-specification.md @@ -102,6 +102,9 @@ - 重新采集的前置校验警告和后台错误 `InfoBar` 都显示在软件窗口顶部水平居中,并有可见的“关闭提示”按钮;提示 5 秒后自动关闭,同一时间只保留一条,新提示替换旧提示。 - 每次重新采集创建新的执行记录和幂等键;当前结果更新,旧结果保存在历史执行记录中。 - 重新采购逐条串行执行。任一任务未形成“执行、核单、上报”成功闭环时立即停止剩余队列。历史上任何一次执行已有 `irreversible_action_at` 的任务永久禁止重新采购,只能只读核单。 +- 用户确认批量操作后立即刷新列表,并在状态区显示已加入队列的数量。排队中的任务保持原状态; + 只有工作线程已经把当前任务写成 `running` 后,才通知主线程刷新该行并显示“采集中”或 + “采购中”。每条完成、失败、跳过或取消后再次刷新,再开始下一条。 - “等待重试”当前没有倒计时。自动获取因可恢复采集错误停止时,底部状态显示 “重试已暂停”,并提示选择任务点击“重新执行”或重新启动获取任务。 - “重新上报”作用于当前已经加载并勾选的任务。执行前显示任务数量,并明确说明不会重新采集、采购或操作手机。