From 5fd71981e094884b8f4343b3071520c4cc20fa75 Mon Sep 17 00:00:00 2001 From: chengma Date: Sun, 9 Aug 2026 10:54:18 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=94=AF=E6=8C=81=E6=8C=81=E7=BB=AD?= =?UTF-8?q?=E8=87=AA=E5=8A=A8=E8=8E=B7=E5=8F=96=E4=BB=BB=E5=8A=A1=20(#45)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- client/src/pdd_ui_event.py | 181 ++++++++++++++++++++++++++--- client/test/test_pdd_ui_event.py | 147 ++++++++++++++++++++++- docs/client/02-architecture.md | 2 +- docs/client/05-ui-specification.md | 17 +-- 4 files changed, 317 insertions(+), 30 deletions(-) diff --git a/client/src/pdd_ui_event.py b/client/src/pdd_ui_event.py index 58543b4..75e4345 100644 --- a/client/src/pdd_ui_event.py +++ b/client/src/pdd_ui_event.py @@ -22,7 +22,14 @@ from typing import Dict, Optional -from PyQt5.QtCore import QCoreApplication, QObject, QThread, pyqtSignal, pyqtSlot +from PyQt5.QtCore import ( + QCoreApplication, + QObject, + QThread, + QTimer, + pyqtSignal, + pyqtSlot, +) from qfluentwidgets import InfoBar, InfoBarPosition from .admin_gateway import ( @@ -122,6 +129,7 @@ class ClaimTaskWorker(QObject): taskSaved = pyqtSignal(str) duplicateTask = pyqtSignal(str) localSaveFailed = pyqtSignal(str, str) + retryableFailed = pyqtSignal(str) failed = pyqtSignal(str) outcome = pyqtSignal(str, str, str) completed = pyqtSignal() @@ -176,7 +184,11 @@ class ClaimTaskWorker(QObject): request_hint = ( f",请求编号:{exc.request_id}" if exc.request_id else "" ) - self.failed.emit(f"{exc}{request_hint}") + 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}") @@ -195,6 +207,9 @@ class PDDTaskPageEvent(QObject): 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 @@ -204,8 +219,23 @@ class PDDTaskPageEvent(QObject): 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() @@ -238,7 +268,7 @@ class PDDTaskPageEvent(QObject): page.autoFetchRequested.connect(self._request_claim_task) page.detailRequested.connect(self.show_task_detail) page.taskModel.loadMoreRequested.connect(self._load_page) - page.autoFetchButton.setText("获取任务") + page.set_auto_fetch_state("stopped") page.destroyed.connect(self.shutdown) application = QCoreApplication.instance() if application is not None: @@ -266,9 +296,12 @@ class PDDTaskPageEvent(QObject): @pyqtSlot() def _request_claim_task(self) -> None: - """启动一次后台领取;重复点击和关闭期间直接忽略。""" + """切换持续自动获取;每一轮仍只启动一个后台 Worker。""" - if self._closing or self._claim_busy: + 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 领取服务未初始化" @@ -276,9 +309,28 @@ class PDDTaskPageEvent(QObject): 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._page.autoFetchButton.setEnabled(False) - self._page.set_engine_status("正在处理一条采集任务,请稍候…") + self._cycle_next_delay_ms = None + self._page.set_auto_fetch_state("running") + self._page.set_engine_status("自动获取:运行中 · 正在处理一条采集任务…") thread = QThread(self) worker = ClaimTaskWorker( @@ -294,6 +346,7 @@ class PDDTaskPageEvent(QObject): 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) @@ -304,10 +357,57 @@ class PDDTaskPageEvent(QObject): 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._page.set_engine_status("暂无可领取的采集任务") + 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: @@ -316,13 +416,15 @@ class PDDTaskPageEvent(QObject): 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._page.set_engine_status( - f"任务 {task_id} 本地已有,未重复保存" + self._continue_after( + self._next_task_delay_ms, + f"任务 {task_id} 本地已有,未重复保存", ) @pyqtSlot(str, str) @@ -334,6 +436,22 @@ class PDDTaskPageEvent(QObject): ) 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: @@ -341,15 +459,32 @@ class PDDTaskPageEvent(QObject): 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._page.set_engine_status(message) self._reload() - if kind in {"result_pending", "manual_review", "failed"}: + 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 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 _show_claim_error(self, title: str, content: str) -> None: """显示不会自动消失的可恢复错误,同时保留底部状态文字。""" @@ -368,9 +503,18 @@ class PDDTaskPageEvent(QObject): self._claim_worker = None self._claim_thread = None self._claim_busy = False - if not self._closing: - self._page.autoFetchButton.setText("获取任务") - self._page.autoFetchButton.setEnabled(True) + 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() @@ -446,6 +590,9 @@ class PDDTaskPageEvent(QObject): 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() @@ -465,6 +612,10 @@ class PDDTaskPageEvent(QObject): 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: diff --git a/client/test/test_pdd_ui_event.py b/client/test/test_pdd_ui_event.py index 6e804f7..dbf0eee 100644 --- a/client/test/test_pdd_ui_event.py +++ b/client/test/test_pdd_ui_event.py @@ -12,7 +12,7 @@ os.environ.setdefault("QT_QPA_PLATFORM", "offscreen") from PyQt5.QtWidgets import QApplication from src.pdd_ui import PDDTaskPage -from src.admin_gateway import AdminTask, SubmissionReceipt +from src.admin_gateway import AdminGatewayError, AdminTask, SubmissionReceipt from src.pdd_ui_event import ( PDDTaskPageEvent, admin_task_to_new_claimed_task, @@ -78,6 +78,35 @@ class SlowClaimGateway(RecordingClaimGateway): return super().claim_next(client, capabilities) +class SequenceClaimGateway(RecordingClaimGateway): + """依次返回多条任务,并记录是否出现并发领取。""" + + def __init__(self, responses): + super().__init__() + self.responses = list(responses) + self.active_calls = 0 + self.max_active_calls = 0 + + def claim_next(self, client, capabilities): + self.active_calls += 1 + self.max_active_calls = max(self.max_active_calls, self.active_calls) + try: + time.sleep(0.01) + response = self.responses.pop(0) if self.responses else None + self.response = response + return super().claim_next(client, capabilities) + finally: + self.active_calls -= 1 + + +class RetryableClaimGateway(RecordingClaimGateway): + """模拟 Admin 暂时不可用。""" + + def claim_next(self, client, capabilities): + self.calls.append((client, capabilities)) + raise AdminGatewayError("ADMIN_UNAVAILABLE", "Admin 暂时不可用", True) + + class FakeCollectResult: def to_pdd_data(self): return { @@ -331,7 +360,7 @@ class PDDTaskPageEventTest(unittest.TestCase): "737116531267", ) - def test_click_claims_one_collect_task_saves_and_refreshes_table(self): + def test_start_claims_task_and_keeps_auto_fetch_running(self): gateway = RecordingClaimGateway(collect_admin_task()) page = PDDTaskPage() events = PDDTaskPageEvent( @@ -357,12 +386,12 @@ class PDDTaskPageEventTest(unittest.TestCase): self.assertEqual(self.repository.count_tasks(), 1) self.assertEqual(page.taskModel.data_row_count(), 1) self.assertEqual(page.taskModel.row_at(0).remote_task_id, "COL-001") - self.assertEqual(page.autoFetchButton.text(), "获取任务") + self.assertEqual(page.autoFetchButton.text(), "停止自动获取") self.assertIn("COL-001", page.statusLabel.text()) events.shutdown() page.deleteLater() - def test_204_is_neutral_and_each_click_only_calls_once(self): + def test_no_task_schedules_next_claim_and_stop_cancels_timer(self): gateway = RecordingClaimGateway(None) page = PDDTaskPage() events = PDDTaskPageEvent( @@ -377,12 +406,18 @@ class PDDTaskPageEventTest(unittest.TestCase): self.assertTrue(wait_until(self.app, lambda: not events._claim_busy)) self.assertEqual(len(gateway.calls), 1) - self.assertEqual(page.statusLabel.text(), "暂无可领取的采集任务") + self.assertIn("5 秒后再次领取", page.statusLabel.text()) + self.assertTrue(events._next_cycle_timer.isActive()) self.assertEqual(self.repository.count_tasks(), 0) + + page.autoFetchRequested.emit() + + self.assertFalse(events._next_cycle_timer.isActive()) + self.assertEqual(page.autoFetchButton.text(), "开始自动获取") events.shutdown() page.deleteLater() - def test_repeated_click_while_claiming_does_not_start_second_request(self): + def test_second_click_while_claiming_requests_safe_stop(self): gateway = SlowClaimGateway(None) page = PDDTaskPage() events = PDDTaskPageEvent( @@ -394,11 +429,86 @@ class PDDTaskPageEventTest(unittest.TestCase): ) page.autoFetchRequested.emit() + self.assertTrue(wait_until(self.app, gateway.started.is_set)) page.autoFetchRequested.emit() self.assertFalse(page.autoFetchButton.isEnabled()) self.assertTrue(wait_until(self.app, lambda: not events._claim_busy)) self.assertEqual(len(gateway.calls), 1) + self.assertEqual(page.autoFetchButton.text(), "开始自动获取") + events.shutdown() + page.deleteLater() + + def test_one_start_processes_multiple_tasks_strictly_in_sequence(self): + gateway = SequenceClaimGateway( + [collect_admin_task("COL-001"), collect_admin_task("COL-002")] + ) + page = PDDTaskPage() + events = PDDTaskPageEvent( + page, + self.repository, + claim_gateway=gateway, + settings_repository=self._saved_settings(), + collect_service_factory=fake_collect_factory, + next_task_delay_ms=10, + no_task_delay_ms=1_000, + ) + + page.autoFetchRequested.emit() + + self.assertTrue( + wait_until(self.app, lambda: self.repository.count_tasks() == 2) + ) + self.assertEqual(gateway.max_active_calls, 1) + self.assertEqual(page.autoFetchButton.text(), "停止自动获取") + + page.autoFetchRequested.emit() + self.assertTrue(wait_until(self.app, lambda: not events._claim_busy)) + self.assertEqual(page.autoFetchButton.text(), "开始自动获取") + events.shutdown() + page.deleteLater() + + def test_retryable_admin_error_schedules_backoff(self): + gateway = RetryableClaimGateway() + page = PDDTaskPage() + events = PDDTaskPageEvent( + page, + self.repository, + claim_gateway=gateway, + settings_repository=self._saved_settings(), + retry_delays_ms=(20, 40, 80, 100), + ) + + page.autoFetchRequested.emit() + + self.assertTrue(wait_until(self.app, lambda: not events._claim_busy)) + self.assertEqual(events._retry_count, 1) + self.assertEqual(events._next_cycle_timer.interval(), 20) + self.assertIn("退避等待", page.statusLabel.text()) + + page.autoFetchRequested.emit() + self.assertFalse(events._next_cycle_timer.isActive()) + events.shutdown() + page.deleteLater() + + def test_retry_delay_grows_to_cap_and_success_resets_it(self): + page = PDDTaskPage() + events = PDDTaskPageEvent( + page, + self.repository, + claim_gateway=RecordingClaimGateway(None), + settings_repository=self._saved_settings(), + retry_delays_ms=(5_000, 10_000, 20_000, 30_000), + ) + + observed = [] + for _ in range(5): + events._on_claim_retryable_failed("Admin 暂时不可用") + observed.append(events._cycle_next_delay_ms) + + self.assertEqual(observed, [5_000, 10_000, 20_000, 30_000, 30_000]) + events._on_collect_outcome("succeeded", "任务完成", "COL-001") + self.assertEqual(events._retry_count, 0) events.shutdown() page.deleteLater() @@ -459,9 +569,34 @@ class PDDTaskPageEventTest(unittest.TestCase): self.assertTrue(wait_until(self.app, lambda: not events._claim_busy)) self.assertEqual(gateway.calls, []) self.assertIn("设置页", page.statusLabel.text()) + self.assertFalse(events._auto_fetch_running) + self.assertEqual(page.autoFetchButton.text(), "开始自动获取") events.shutdown() page.deleteLater() + def test_shutdown_cancels_scheduled_next_claim(self): + gateway = RecordingClaimGateway(None) + page = PDDTaskPage() + events = PDDTaskPageEvent( + page, + self.repository, + claim_gateway=gateway, + settings_repository=self._saved_settings(), + no_task_delay_ms=20, + ) + page.autoFetchRequested.emit() + self.assertTrue(wait_until(self.app, lambda: not events._claim_busy)) + self.assertTrue(events._next_cycle_timer.isActive()) + + events.shutdown() + deadline = time.monotonic() + 0.08 + while time.monotonic() < deadline: + self.app.processEvents() + time.sleep(0.005) + + self.assertEqual(len(gateway.calls), 1) + page.deleteLater() + def test_shutdown_after_claim_started_still_saves_task_without_ui_callback(self): gateway = SlowClaimGateway(collect_admin_task("COL-CLOSE")) page = PDDTaskPage() diff --git a/docs/client/02-architecture.md b/docs/client/02-architecture.md index 1caf7f5..647f76c 100644 --- a/docs/client/02-architecture.md +++ b/docs/client/02-architecture.md @@ -87,7 +87,7 @@ client/ |---|---|---| | 主窗口文件 | `src/ui/main_window.py` | `src/ui_main.py`。**注意**:`client/AGENTS.md` 已按 `src/ui_main.py` 写规则,两处不一致,见下方说明 | | 目录分层 | domain / application / infrastructure / workers / ui | 平铺在 `src/` 下:`db.py`、`db_schema.py`、`task_repository.py`、`settings_repository.py`、`task_models.py` 等,另有 `src/util/`、`src/demo1/` | -| 主按钮文案 | 「获取任务」⇄「停止获取」 | #32 暂为「获取任务」单次模式,运行期间禁用 | +| 主按钮文案 | 「开始自动获取」⇄「停止自动获取」 | 已按持续串行模式实现 | | Admin 网关 | `AdminGateway` + Mock/HTTP 两实现 | 登记、领取、结果和失败提交已实现 | | 采集任务协调器 | `CollectTaskService` | 已实现单条采集任务流程;采购仍未接入 | | Outbox 提交 | 从 `outbox_events` 取件重试 | 采集结果与失败已实现,重试不重复采集 | diff --git a/docs/client/05-ui-specification.md b/docs/client/05-ui-specification.md index 8ec502b..33b275e 100644 --- a/docs/client/05-ui-specification.md +++ b/docs/client/05-ui-specification.md @@ -35,7 +35,7 @@ ```text ┌─────────────────────────────────────────────────────────────┐ │ PDD 任务 │ -│ [获取任务] [类型▼] [状态▼] [关键词................] [搜索] │ +│ [开始自动获取] [类型▼] [状态▼] [关键词............] [搜索] │ ├─────────────────────────────────────────────────────────────┤ │ 类型 │ 商品标题 │ 颜色 │ 尺码 │ 价格 │ 数量 │ 状态 │ 更新时间 │详情│ │ │ @@ -50,20 +50,21 @@ ## 4. 顶部命令区 -### 4.1 获取任务 / 停止获取 +### 4.1 开始自动获取 / 停止自动获取 - 使用 `PrimaryPushButton`,是页面唯一主要强调操作。 -- 初始文本为“获取任务”,图标表达开始。 -- 启动成功后文本变为“停止获取”,按钮位置和宽度尽量稳定。 +- 初始文本为“开始自动获取”,图标表达开始。 +- 启动成功后文本变为“停止自动获取”,按钮位置和宽度尽量稳定。 - `[必须]` 这两个文案是**按钮标签**。状态栏里的“自动获取:已开启/已停止”说的是**功能状态**,两者不是一回事,不要混改。 - 启动过程中禁用重复点击,显示“正在启动”。 - 停止表示不再领取新任务;当前任务进入安全停止流程。 - 如果设备、Admin 或必要设置无效,按钮可禁用,但附近必须说明缺少的条件。 -**当前单任务模式(#32):**按钮文字固定为“获取任务”,每次点击先补交一条本地 -Outbox;若无待提交结果则执行最早的本地采集任务,再没有才向 Admin 领取一条。 -领取、手机采集和提交都在工作线程完成,运行期间按钮禁用,结束后刷新任务表。 -当前不循环领取,也不执行采购任务。 +每一轮先补交一条本地 Outbox;若无待提交结果则执行最早的本地采集任务,再没有 +才向 Admin 领取一条。领取、手机采集和提交都在单独的工作线程完成,一轮结束且 +线程完全退出后,由主线程的单次定时器安排下一轮。暂无任务时默认 5 秒后重试; +可恢复的 Admin 错误按 5、10、20、30 秒退避。需要人工处理、不可恢复错误、设备 +或配置错误会停止自动获取。当前仍不执行采购任务。 ### 4.2 搜索与筛选