diff --git a/client/src/collect_task_service.py b/client/src/collect_task_service.py index 51dc890..3aecd35 100644 --- a/client/src/collect_task_service.py +++ b/client/src/collect_task_service.py @@ -111,11 +111,26 @@ class CollectTaskService: if task is None: raise RuntimeError(f"任务 {remote.task_id} 未能保存到本地") + return self._execute_task(task.remote_task_id) + + def execute_selected(self, remote_task_id: str) -> CollectTaskOutcome: + """只执行指定的本地采集任务,不领取任务或补交其他结果。""" + + if not self._device_address.strip(): + raise ValueError("请先在设置页选择并保存 Android 设备") + task = self._repository.get_task(remote_task_id) + if task is None: + raise ValueError(f"任务 {remote_task_id} 不存在") + return self._execute_task(task.remote_task_id) + + def _execute_task(self, remote_task_id: str) -> CollectTaskOutcome: + """执行一条已经保存在本地且处于待执行状态的采集任务。""" + if self._cancelled(): - return CollectTaskOutcome("cancelled", "本次采集已取消", task.remote_task_id) + return CollectTaskOutcome("cancelled", "本次采集已取消", remote_task_id) started = self._repository.start_collect_run( - task.remote_task_id, self._device_address + remote_task_id, self._device_address ) collector = self._factory( self._device_address, self._client.client_id, self._cancelled @@ -123,13 +138,13 @@ class CollectTaskService: try: result = collector.collect(started.task) event = self._repository.save_collect_result( - task.remote_task_id, started.attempt_id, result.to_pdd_data() + remote_task_id, started.attempt_id, result.to_pdd_data() ) except PddCollectError as exc: status, retryable = self._classify_error(exc.code) report_code = self._report_error_code(exc.code) event = self._repository.save_collect_failure( - task.remote_task_id, + remote_task_id, started.attempt_id, status, report_code, diff --git a/client/src/db_schema.py b/client/src/db_schema.py index e2d1265..2efe68b 100644 --- a/client/src/db_schema.py +++ b/client/src/db_schema.py @@ -4,7 +4,7 @@ ``MIGRATIONS`` 末尾增加版本,不能修改已经发布的迁移。 """ -SCHEMA_VERSION = 1 +SCHEMA_VERSION = 2 MIGRATION_1 = ( @@ -127,6 +127,14 @@ MIGRATION_1 = ( ) +MIGRATION_2 = ( + """ + ALTER TABLE task_runs ADD COLUMN result_data TEXT + """, +) + + MIGRATIONS = { 1: MIGRATION_1, + 2: MIGRATION_2, } diff --git a/client/src/pdd_ui.py b/client/src/pdd_ui.py index b56323e..2b17c4e 100644 --- a/client/src/pdd_ui.py +++ b/client/src/pdd_ui.py @@ -21,8 +21,8 @@ - 图标用 `FluentIcon`,**不得用表情符号**。 - 页面 `objectName` 固定为 `pddTaskPage`,不要改。 -注意:本页唯一会产生外部后果(真的去操作手机、可能下单)的命令是 -“开始自动获取”。其余操作全部只读本地数据库。 +注意:“开始自动获取”和“重新执行”会真的操作手机。重新执行当前只允许 +采集任务,采购任务必须在事件层拦截,不能从本入口下单。 """ from dataclasses import dataclass @@ -276,6 +276,7 @@ class PDDTaskPage(QWidget): autoFetchRequested = pyqtSignal() searchRequested = pyqtSignal(dict) refreshRequested = pyqtSignal() + rerunRequested = pyqtSignal(str) detailRequested = pyqtSignal(str) def __init__(self, parent=None): @@ -323,6 +324,10 @@ class PDDTaskPage(QWidget): self.refreshButton = PushButton(FIF.SYNC, "刷新", self) self.refreshButton.setAccessibleName("刷新本地任务表格") + self.rerunButton = PushButton(FIF.UPDATE, "重新执行", self) + self.rerunButton.setAccessibleName("重新执行当前选中的采集任务") + self.rerunButton.setEnabled(False) + self.commandCard = CardWidget(self) commandLayout = QGridLayout(self.commandCard) commandLayout.setContentsMargins(20, 14, 20, 14) @@ -338,6 +343,7 @@ class PDDTaskPage(QWidget): commandLayout.setColumnStretch(6, 4) commandLayout.setColumnStretch(8, 1) commandLayout.addWidget(self.refreshButton, 0, 9) + commandLayout.addWidget(self.rerunButton, 0, 10) def _build_content_area(self) -> None: self.taskTable = TableView(self) @@ -410,9 +416,13 @@ class PDDTaskPage(QWidget): self.searchButton.clicked.connect(self._apply_filters) self.keywordInput.returnPressed.connect(self._apply_filters) self.refreshButton.clicked.connect(self.refreshRequested.emit) + self.rerunButton.clicked.connect(self._request_rerun) self.clearFiltersButton.clicked.connect(self.clear_filters) self.taskTable.clicked.connect(self._on_table_clicked) self.taskTable.activated.connect(self._on_table_activated) + self.taskTable.selectionModel().selectionChanged.connect( + self._update_rerun_button + ) self.findShortcut = QShortcut(QKeySequence.Find, self) self.findShortcut.activated.connect(self.keywordInput.setFocus) @@ -497,6 +507,25 @@ class PDDTaskPage(QWidget): task = self.taskModel.row_at(index.row()) if index.isValid() else None return task.remote_task_id if task else "" + def set_rerun_running(self, running: bool) -> None: + """重新采集期间锁定两个会操作手机的入口。""" + + self.rerunButton.setText("正在重新采集" if running else "重新执行") + self.autoFetchButton.setEnabled(not running) + self._update_rerun_button() + + def _request_rerun(self) -> None: + task_id = self.current_task_id() + if task_id: + self.rerunRequested.emit(task_id) + + def _update_rerun_button(self) -> None: + running = self.rerunButton.text() == "正在重新采集" + has_selection = self.taskTable.selectionModel().hasSelection() + self.rerunButton.setEnabled( + has_selection and bool(self.current_task_id()) and not running + ) + def _apply_filters(self) -> None: filters = { "task_type": self.taskTypeCombo.currentText(), diff --git a/client/src/pdd_ui_event.py b/client/src/pdd_ui_event.py index 75e4345..22a427d 100644 --- a/client/src/pdd_ui_event.py +++ b/client/src/pdd_ui_event.py @@ -30,7 +30,7 @@ from PyQt5.QtCore import ( pyqtSignal, pyqtSlot, ) -from qfluentwidgets import InfoBar, InfoBarPosition +from qfluentwidgets import InfoBar, InfoBarPosition, MessageBox from .admin_gateway import ( AdminGatewayError, @@ -51,7 +51,7 @@ from .task_models import ( TaskSummary, TaskType, ) -from .task_repository import TaskRepository +from .task_repository import CollectRerunError, TaskRepository from .task_detail_view import TaskDetailWindow @@ -141,6 +141,7 @@ class ClaimTaskWorker(QObject): client_service: CurrentClientService, android_device_service: SelectedAndroidDeviceService, collect_service_factory: Optional[CollectServiceFactory] = None, + selected_task_id: str = "", ) -> None: super().__init__() self._gateway = gateway @@ -149,6 +150,7 @@ class ClaimTaskWorker(QObject): 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: """阻止尚未开始的领取;已领取的任务仍必须保存到本地。""" @@ -166,7 +168,7 @@ class ClaimTaskWorker(QObject): return android_serial = self._android_device_service.load() - result = CollectTaskService( + service = CollectTaskService( self._gateway, self._task_repository, ClientInfo( @@ -176,7 +178,12 @@ class ClaimTaskWorker(QObject): android_serial or "", cancelled=lambda: self._cancelled, collect_service_factory=self._collect_service_factory, - ).execute_one() + ) + 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: @@ -265,6 +272,7 @@ class PDDTaskPageEvent(QObject): 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) @@ -294,6 +302,144 @@ class PDDTaskPageEvent(QObject): 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。""" diff --git a/client/src/task_models.py b/client/src/task_models.py index ab55768..f768ed8 100644 --- a/client/src/task_models.py +++ b/client/src/task_models.py @@ -161,6 +161,7 @@ class TaskRunRecord: error_code: Optional[str] error_message: Optional[str] diagnostics_json: Optional[Dict[str, Any]] + result_data: Optional[Dict[str, Any]] artifact_directory: Optional[str] created_at: str updated_at: str diff --git a/client/src/task_repository.py b/client/src/task_repository.py index 56ae7fd..3448ce8 100644 --- a/client/src/task_repository.py +++ b/client/src/task_repository.py @@ -34,6 +34,10 @@ class DuplicateTaskError(ValueError): """相同远程任务编号已经存在,不能覆盖。""" +class CollectRerunError(ValueError): + """当前任务不满足重新采集条件。""" + + def utc_now_iso() -> str: """返回精确到秒的 UTC ISO 8601 时间。""" @@ -206,6 +210,81 @@ class TaskRepository: connection.close() return self._to_detail(row) if row is not None else None + def validate_collect_rerun(self, remote_task_id: str) -> TaskDetail: + """校验任务能否重新采集,成功时返回任务详情。""" + + connection = open_database(self._db_path) + try: + row = connection.execute( + "SELECT * FROM pdd_tasks WHERE remote_task_id = ?", + (remote_task_id,), + ).fetchone() + if row is None: + raise CollectRerunError(f"任务 {remote_task_id} 不存在") + self._check_collect_rerun(connection, row) + finally: + connection.close() + return self._to_detail(row) + + def prepare_collect_rerun(self, remote_task_id: str) -> TaskDetail: + """事务内把已结束的采集任务恢复为待执行,保留旧结果。""" + + now = utc_now_iso() + connection = open_database(self._db_path) + try: + with connection: + row = connection.execute( + "SELECT * FROM pdd_tasks WHERE remote_task_id = ?", + (remote_task_id,), + ).fetchone() + if row is None: + raise CollectRerunError(f"任务 {remote_task_id} 不存在") + self._check_collect_rerun(connection, row) + connection.execute( + "UPDATE pdd_tasks SET status = 'claimed'," + " current_step = 'rerun_requested', finished_at = NULL," + " last_error_code = NULL, last_error_message = NULL," + " updated_at = ? WHERE id = ?", + (now, row["id"]), + ) + finally: + connection.close() + task = self.get_task(remote_task_id) + assert task is not None + return task + + @staticmethod + def _check_collect_rerun( + connection: sqlite3.Connection, row: sqlite3.Row + ) -> None: + """检查重新采集的类型、终态和 Outbox 约束。""" + + if row["task_type"] != TaskType.COLLECT.value: + raise CollectRerunError("采购任务不能重新执行,以免重复下单") + allowed_statuses = { + TaskStatus.SUCCEEDED.value, + TaskStatus.FAILED.value, + TaskStatus.CANCELLED.value, + TaskStatus.MANUAL_REVIEW.value, + } + if row["status"] not in allowed_statuses: + status_name = { + TaskStatus.CLAIMED.value: "待执行", + TaskStatus.RUNNING.value: "执行中", + TaskStatus.RESULT_PENDING.value: "结果待提交", + TaskStatus.RETRY_WAIT.value: "等待重试", + }.get(row["status"], row["status"]) + raise CollectRerunError(f"任务当前为“{status_name}”,不能重新采集") + unsent_count = int( + connection.execute( + "SELECT COUNT(*) FROM outbox_events" + " WHERE task_id = ? AND status != 'sent'", + (row["id"],), + ).fetchone()[0] + ) + if unsent_count: + raise CollectRerunError("任务仍有未发送的结果,请先完成提交") + def start_collect_run( self, remote_task_id: str, device_address: str ) -> StartedTaskRun: @@ -284,22 +363,24 @@ class TaskRepository: "pdd_data": pdd_data, } idempotency_key = f"{remote_task_id}:{attempt_id}:result-v1" + result_json = json.dumps(pdd_data, ensure_ascii=False) connection.execute( "UPDATE pdd_tasks SET status = 'result_pending'," " current_step = 'submit_result', pdd_data = ?, goods_id = ?," " title = ?, price_cent = ?, finished_at = ?, updated_at = ?" " WHERE id = ?", ( - json.dumps(pdd_data, ensure_ascii=False), + result_json, pdd_data.get("goods_id"), pdd_data.get("title"), self._summary_price(pdd_data), now, now, task["id"], ), ) connection.execute( "UPDATE task_runs SET run_status = 'succeeded'," - " current_step = 'submit_result', finished_at = ?, updated_at = ?" + " current_step = 'submit_result', result_data = ?," + " finished_at = ?, updated_at = ?" " WHERE attempt_id = ?", - (now, now, attempt_id), + (result_json, now, now, attempt_id), ) cursor = connection.execute( "INSERT INTO outbox_events (task_id, event_type, idempotency_key," diff --git a/client/test/test_collect_task_service.py b/client/test/test_collect_task_service.py index 8355da0..dd1755e 100644 --- a/client/test/test_collect_task_service.py +++ b/client/test/test_collect_task_service.py @@ -4,7 +4,12 @@ import tempfile import unittest from pathlib import Path -from src.admin_gateway import AdminTask, ClaimCapabilities, ClientInfo +from src.admin_gateway import ( + AdminTask, + ClaimCapabilities, + ClientInfo, + SubmissionReceipt, +) from src.collect_task_service import CollectTaskService from src.mock_admin_gateway import MockAdminGateway from src.pdd_collect_service import PddCollectError @@ -159,6 +164,36 @@ class CollectTaskServiceTest(unittest.TestCase): ) self.assertEqual(self.gateway.submission_count, 1) + def test_execute_selected_only_runs_requested_local_task(self): + for task_id in ("COL-SELECTED", "COL-OTHER"): + self.repository.add_claimed_task( + NewClaimedTask( + remote_task_id=task_id, + task_type=TaskType.COLLECT, + goods_url=f"https://example.test/{task_id}", + ) + ) + calls = [] + + class AcceptGateway: + def submit_result(self, *_args): + return SubmissionReceipt(True, "RESULT-001", "2026-08-07T08:00:00Z") + + service = CollectTaskService( + AcceptGateway(), + self.repository, + self.client, + "USB-001", + collect_service_factory=lambda *_args: FakeCollector(calls), + ) + outcome = service.execute_selected("COL-SELECTED") + + self.assertEqual(outcome.kind, "succeeded") + self.assertEqual(calls, ["COL-SELECTED"]) + self.assertEqual( + self.repository.get_task("COL-OTHER").status, TaskStatus.CLAIMED + ) + if __name__ == "__main__": unittest.main() diff --git a/client/test/test_db.py b/client/test/test_db.py index 09b9dfb..ccfddb4 100644 --- a/client/test/test_db.py +++ b/client/test/test_db.py @@ -1,4 +1,4 @@ -"""SQLite 初始化和 v1 数据库结构测试。""" +"""SQLite 初始化和数据库迁移测试。""" import sqlite3 import tempfile @@ -6,6 +6,7 @@ import unittest from pathlib import Path from src.db import DatabaseVersionError, initialize_database, open_database +from src.db_schema import MIGRATION_1 EXPECTED_TABLES = { @@ -36,7 +37,7 @@ class DatabaseInitializationTests(unittest.TestCase): def tearDown(self) -> None: self._temporary_directory.cleanup() - def test_initialize_creates_v1_tables_and_indexes(self) -> None: + def test_initialize_creates_latest_tables_and_indexes(self) -> None: result_path = initialize_database(self.db_path) self.assertEqual(result_path, self.db_path) @@ -62,7 +63,50 @@ class DatabaseInitializationTests(unittest.TestCase): self.assertTrue(EXPECTED_TABLES.issubset(tables)) self.assertTrue(EXPECTED_INDEXES.issubset(indexes)) - self.assertEqual(version, 1) + self.assertEqual(version, 2) + + def test_v1_database_is_upgraded_without_losing_task_runs(self) -> None: + connection = open_database(self.db_path) + try: + with connection: + for statement in MIGRATION_1: + connection.execute(statement) + connection.execute("PRAGMA user_version = 1") + connection.execute( + "INSERT INTO pdd_tasks" + " (remote_task_id, task_type, goods_url, status, received_at," + " created_at, updated_at) VALUES" + " ('TASK-OLD', 'collect', 'https://example.test', 'claimed'," + " '2026-08-06T00:00:00Z', '2026-08-06T00:00:00Z'," + " '2026-08-06T00:00:00Z')" + ) + connection.execute( + "INSERT INTO task_runs" + " (task_id, attempt_id, attempt_no, device_address, run_status," + " started_at, created_at, updated_at) VALUES" + " (1, 'ATTEMPT-OLD', 1, 'USB-001', 'succeeded'," + " '2026-08-06T00:00:00Z', '2026-08-06T00:00:00Z'," + " '2026-08-06T00:00:00Z')" + ) + finally: + connection.close() + + initialize_database(self.db_path) + + connection = open_database(self.db_path) + try: + columns = { + row[1] for row in connection.execute("PRAGMA table_info(task_runs)") + } + attempt_id = connection.execute( + "SELECT attempt_id FROM task_runs" + ).fetchone()[0] + version = connection.execute("PRAGMA user_version").fetchone()[0] + finally: + connection.close() + self.assertIn("result_data", columns) + self.assertEqual(attempt_id, "ATTEMPT-OLD") + self.assertEqual(version, 2) def test_initialize_can_run_twice_without_losing_data(self) -> None: initialize_database(self.db_path) diff --git a/client/test/test_pdd_ui_event.py b/client/test/test_pdd_ui_event.py index dbf0eee..5c076de 100644 --- a/client/test/test_pdd_ui_event.py +++ b/client/test/test_pdd_ui_event.py @@ -6,6 +6,7 @@ import threading import time import unittest from pathlib import Path +from unittest.mock import patch os.environ.setdefault("QT_QPA_PLATFORM", "offscreen") @@ -242,6 +243,56 @@ class PDDTaskPageEventTest(unittest.TestCase): self.assertEqual(row.price_cents, 3990) page.deleteLater() + def test_rerun_button_follows_current_row_selection(self): + self._add_task(1) + page = PDDTaskPage() + events = PDDTaskPageEvent(page, self.repository) + events.load_initial_tasks() + + page.taskTable.clearSelection() + self.app.processEvents() + self.assertFalse(page.rerunButton.isEnabled()) + page.taskTable.selectRow(0) + self.app.processEvents() + self.assertTrue(page.rerunButton.isEnabled()) + page.taskTable.clearSelection() + self.app.processEvents() + self.assertFalse(page.rerunButton.isEnabled()) + events.shutdown() + page.deleteLater() + + def test_confirmed_rerun_executes_selected_terminal_collect_task(self): + self._add_task(1) + first = self.repository.start_collect_run("PDD-001", "USB-001") + event = self.repository.save_collect_result( + "PDD-001", first.attempt_id, FakeCollectResult().to_pdd_data() + ) + self.repository.mark_outbox_sent(event.id) + page = PDDTaskPage() + gateway = RecordingClaimGateway() + events = PDDTaskPageEvent( + page, + self.repository, + claim_gateway=gateway, + settings_repository=self._saved_settings(), + collect_service_factory=fake_collect_factory, + ) + + with patch("src.pdd_ui_event.MessageBox") as message_box: + message_box.return_value.exec.return_value = True + page.rerunRequested.emit("PDD-001") + self.assertTrue( + wait_until(self.app, lambda: not events._claim_busy), + "重新采集线程没有按时结束", + ) + + detail = self.repository.get_task("PDD-001") + self.assertEqual(detail.status, TaskStatus.SUCCEEDED) + self.assertIn("采集完成", page.statusLabel.text()) + self.assertEqual(gateway.calls, []) + events.shutdown() + page.deleteLater() + def test_query_failure_shows_readable_error(self): page = PDDTaskPage() events = PDDTaskPageEvent(page, BrokenRepository()) diff --git a/client/test/test_task_repository.py b/client/test/test_task_repository.py index 7c9809c..448fa74 100644 --- a/client/test/test_task_repository.py +++ b/client/test/test_task_repository.py @@ -14,7 +14,7 @@ from src.task_models import ( TaskStatus, TaskType, ) -from src.task_repository import DuplicateTaskError, TaskRepository +from src.task_repository import CollectRerunError, DuplicateTaskError, TaskRepository class TaskRepositoryTests(unittest.TestCase): @@ -209,6 +209,64 @@ class TaskRepositoryTests(unittest.TestCase): TaskStatus.SUCCEEDED, ) + def test_prepare_rerun_preserves_old_result_and_creates_new_attempt(self): + self.repository.add_claimed_task(self._task("TASK-RERUN")) + first = self.repository.start_collect_run("TASK-RERUN", "USB-001") + first_data = {"goods_id": "10001", "title": "旧标题", "skus": []} + first_event = self.repository.save_collect_result( + "TASK-RERUN", first.attempt_id, first_data + ) + self.repository.mark_outbox_sent(first_event.id) + + prepared = self.repository.prepare_collect_rerun("TASK-RERUN") + second = self.repository.start_collect_run("TASK-RERUN", "USB-001") + second_data = {"goods_id": "10001", "title": "新标题", "skus": []} + self.repository.save_collect_result( + "TASK-RERUN", second.attempt_id, second_data + ) + + connection = open_database(self.db_path) + try: + rows = connection.execute( + "SELECT attempt_id, result_data FROM task_runs" + " WHERE task_id = ? ORDER BY attempt_no", + (prepared.id,), + ).fetchall() + finally: + connection.close() + self.assertEqual(second.attempt_no, 2) + self.assertIn("旧标题", rows[0]["result_data"]) + self.assertIn("新标题", rows[1]["result_data"]) + self.assertEqual(self.repository.get_task("TASK-RERUN").title, "新标题") + + def test_rerun_rejects_purchase_active_and_unsent_tasks(self): + self.repository.add_claimed_task( + self._task("PURCHASE-RERUN", TaskType.PURCHASE) + ) + with self.assertRaisesRegex(CollectRerunError, "采购任务"): + self.repository.validate_collect_rerun("PURCHASE-RERUN") + + self.repository.add_claimed_task(self._task("ACTIVE-RERUN")) + with self.assertRaisesRegex(CollectRerunError, "待执行"): + self.repository.validate_collect_rerun("ACTIVE-RERUN") + + self.repository.add_claimed_task(self._task("UNSENT-RERUN")) + started = self.repository.start_collect_run("UNSENT-RERUN", "USB-001") + self.repository.save_collect_result( + "UNSENT-RERUN", started.attempt_id, {"title": "结果", "skus": []} + ) + connection = open_database(self.db_path) + try: + with connection: + connection.execute( + "UPDATE pdd_tasks SET status = 'failed'" + " WHERE remote_task_id = 'UNSENT-RERUN'" + ) + finally: + connection.close() + with self.assertRaisesRegex(CollectRerunError, "未发送"): + self.repository.validate_collect_rerun("UNSENT-RERUN") + def test_recovery_restores_sending_and_interrupted_running(self): self.repository.add_claimed_task(self._task("TASK-RECOVER")) started = self.repository.start_collect_run("TASK-RECOVER", "USB-001") diff --git a/docs/client/03-data-model.md b/docs/client/03-data-model.md index c94ac4a..94911b7 100644 --- a/docs/client/03-data-model.md +++ b/docs/client/03-data-model.md @@ -210,7 +210,8 @@ LIMIT :limit OFFSET :offset; ## 4. `task_runs` -每次自动执行生成一条记录,任务重试时创建新的 `attempt_no`,不得覆盖历史。 +每次自动执行或人工确认的重新采集都生成一条记录。任务再次执行时创建新的 +`attempt_no`,不得覆盖历史。 ```sql CREATE TABLE task_runs ( @@ -232,6 +233,7 @@ CREATE TABLE task_runs ( error_code TEXT, error_message TEXT, diagnostics_json TEXT, + result_data TEXT, artifact_directory TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, @@ -244,6 +246,8 @@ CREATE INDEX idx_task_runs_task ``` `irreversible_action_at` 一旦写入,恢复逻辑不得再次下单,只能核对订单或转人工处理。 +`result_data` 保存这次成功采集的完整结果;`pdd_tasks.pdd_data` 只保存最新结果。 +这样重新采集可以更新当前数据,同时仍能按 `attempt_no` 追查旧结果。 ## 5. `outbox_events` @@ -424,8 +428,7 @@ CREATE TABLE app_settings ( | `retry_wait` 重试等待 | `running` | 退避时间到,**本地直接重跑**,新建 `task_runs` 记录(`attempt_no` 加 1) | 任务协调器 | | `retry_wait` | `failed` | 超过最大重试次数,且从未进入不可逆阶段 | 任务协调器 | | `retry_wait` | `manual_review` | 超过最大重试次数,但**曾经进入过不可逆阶段** | 任务协调器 | -| `manual_review` 需要人工 | 不自动变 | 只能由人在 Admin 侧处理;需要再执行时由 Admin 下发新任务 | 人 | -| `succeeded` / `failed` / `cancelled` | **终态,不再变化** | — | — | +| `manual_review` / `succeeded` / `failed` / `cancelled` | `claimed` | 用户在 Client 明确确认重新采集;仅限采集任务,且没有未发送 Outbox | 人 | **注意 `retry_wait` → `running` 是纯本地操作。** 任务已经在本地库里了,重试直接重跑就行, **不需要再向 Admin 要一次**。没有租约,也就没有"重新获取执行权"这回事。 @@ -435,9 +438,10 @@ CREATE TABLE app_settings ( - `[必须]` 任务只能通过 `claim` 成功进入本地库,不许本地自己造一条 `claimed` 记录。 - `[必须]` 只有 `running` 状态才允许操作设备、产生业务执行步骤。 - `[必须]` 只有 Admin 返回 `accepted: true` 才能进入 `succeeded`,本地不许自己判定成功。 -- `[必须]` `manual_review` 不自动重新执行,任何自动流程都不许把它改回 `running`。 -- `[必须]` 已经是 `succeeded`、`failed`、`cancelled` 的任务,不得被迟到的后台回调改回运行中状态。 +- `[必须]` `manual_review` 不自动重新执行;只有用户在 Client 明确确认重新采集,才允许先回到 `claimed`。 +- `[必须]` 已经是 `succeeded`、`failed`、`cancelled` 的任务,不得被迟到的后台回调改回运行中状态;人工确认的重新采集除外。 - `[必须]` 每次重试都要新建一条 `task_runs` 记录(`attempt_no` 加 1),不许覆盖上一次的记录。 +- `[必须]` 采购任务不能从“重新执行”入口启动;重新采集不处理任何采购动作。 ### 7.3 崩溃重启后怎么恢复 diff --git a/docs/client/05-ui-specification.md b/docs/client/05-ui-specification.md index 33b275e..86c7bd0 100644 --- a/docs/client/05-ui-specification.md +++ b/docs/client/05-ui-specification.md @@ -9,7 +9,7 @@ ## 1. 设计目标 - 让操作人员在一个页面完成任务监控、搜索和异常定位。 -- 界面上只有一个会产生外部后果的命令:“获取任务”。其余操作全部只读本地数据库。 +- “获取任务”和“重新执行”会产生外部后果;后者当前只允许重新采集,不允许采购。 - 长任务状态始终可找到,不使用连续模态弹窗打断工作。 - 任务表格在数据增长后仍保持响应速度、稳定选择和可访问性。 - 界面只展示任务状态,不在 Qt 主线程执行 Admin 或手机自动化。 @@ -35,7 +35,7 @@ ```text ┌─────────────────────────────────────────────────────────────┐ │ PDD 任务 │ -│ [开始自动获取] [类型▼] [状态▼] [关键词............] [搜索] │ +│ [开始自动获取] [类型▼] [状态▼] [关键词...] [搜索] [刷新] [重新执行] │ ├─────────────────────────────────────────────────────────────┤ │ 类型 │ 商品标题 │ 颜色 │ 尺码 │ 价格 │ 数量 │ 状态 │ 更新时间 │详情│ │ │ @@ -77,6 +77,16 @@ 筛选条件之间采用 AND。点击搜索或在关键词输入框按 Enter 执行本地数据库查询,不请求领取任务。活动筛选必须可见,筛选无结果时保留条件并提供清除入口。 +### 4.3 刷新与重新执行 + +- “刷新”只重新读取本地任务列表,不请求 Admin,也不操作手机。 +- “重新执行”位于刷新右侧;没有选中任务时禁用。 +- 点击后先校验任务,再显示明确的“重新采集”确认弹窗。弹窗显示任务编号、商品标题,并说明新结果会覆盖 Client 和 Admin 的当前采集数据。 +- 只允许重新采集已经结束的采集任务。采购、执行中、结果待提交、等待重试、仍有未发送 Outbox 或自动获取忙碌时必须阻止,并用中文说明原因。 +- 确认后只执行选中的稳定任务编号,不领取新任务,不先处理其他任务或 Outbox。 +- 重新采集在工作线程运行。执行期间禁用“获取任务”和“重新执行”,完成后刷新列表。 +- 每次重新采集创建新的执行记录和幂等键;当前结果更新,旧结果保存在历史执行记录中。 + ## 5. 任务表格 表格显示的是**本机已领取的全部任务**,包括正在做的和早已做完的。已完成任务永久保留,不会被清理,所以数据只增不减——增量加载和索引是必须的,不是优化。 @@ -333,7 +343,7 @@ self.show_recoverable_error( ## 10. 键盘与无障碍 -- Tab 顺序:自动获取 → 类型 → 状态 → 关键词 → 搜索 → 表格 → 状态区可操作项。 +- Tab 顺序:自动获取 → 类型 → 状态 → 关键词 → 搜索 → 刷新 → 重新执行 → 表格 → 状态区可操作项。 - `Ctrl+F` 聚焦关键词,Enter 打开当前行详情。不绑定 `F5`——界面上没有需要刷新的远端数据。 - 仅图标按钮必须设置准确的无障碍名称和工具提示。 - 表单具有可见标签,占位符不能替代标签。