1637 lines
58 KiB
Python
1637 lines
58 KiB
Python
"""PDD 任务列表接入 SQLite 的离屏测试。"""
|
|
|
|
import os
|
|
import tempfile
|
|
import threading
|
|
import time
|
|
import unittest
|
|
from pathlib import Path
|
|
from unittest.mock import patch
|
|
|
|
os.environ.setdefault("QT_QPA_PLATFORM", "offscreen")
|
|
|
|
from PyQt5.QtCore import Qt
|
|
from PyQt5.QtWidgets import QApplication
|
|
from qfluentwidgets import InfoBarPosition
|
|
|
|
from src.android_device_service import AndroidDeviceSearchError
|
|
from src.db import open_database
|
|
from src.pdd_ui import PDDTaskPage
|
|
from src.admin_gateway import AdminGatewayError, AdminTask, SubmissionReceipt
|
|
from src.pdd_ui_event import (
|
|
ERROR_FEEDBACK_BUTTON_MIN_HEIGHT,
|
|
ERROR_FEEDBACK_BUTTON_MIN_WIDTH,
|
|
ERROR_FEEDBACK_DURATION_MS,
|
|
ERROR_FEEDBACK_MIN_HEIGHT,
|
|
ERROR_FEEDBACK_MIN_WIDTH,
|
|
PURCHASE_RECONCILE_DELAY_MS,
|
|
PDDTaskPageEvent,
|
|
admin_task_to_new_claimed_task,
|
|
summary_to_row,
|
|
)
|
|
from src.mock_admin_gateway import MockAdminGateway
|
|
from src.pdd_purchase_adapter import PddPurchaseAdapter, PurchasePageState
|
|
from src.pdd_collect_service import PddCollectError
|
|
from src.pdd_u2_purchase_adapter import create_u2_purchase_adapter
|
|
from src.settings_repository import SettingsRepository
|
|
from src.task_models import (
|
|
NewClaimedTask,
|
|
OutboxEventType,
|
|
TaskStatus,
|
|
TaskSummary,
|
|
TaskType,
|
|
)
|
|
from src.task_repository import TaskRepository
|
|
from src.ui_main import MainWindow
|
|
|
|
|
|
class BrokenRepository:
|
|
"""模拟无法读取数据库的 Repository。"""
|
|
|
|
def list_tasks(self, filters=None, limit=50, offset=0):
|
|
raise RuntimeError("database is unavailable")
|
|
|
|
|
|
class BrokenSaveRepository(BrokenRepository):
|
|
"""模拟任务已经在 Admin 领取,但本地写入失败。"""
|
|
|
|
def add_claimed_task(self, _task):
|
|
raise RuntimeError("disk is full")
|
|
|
|
def next_pending_outbox(self):
|
|
return None
|
|
|
|
def next_collect_task(self):
|
|
return None
|
|
|
|
def next_runnable_task(self, *, include_purchase):
|
|
return None
|
|
|
|
def next_purchase_task(self):
|
|
return None
|
|
|
|
def next_purchase_reconcile_task(self):
|
|
return None
|
|
|
|
|
|
class RecordingClaimGateway:
|
|
"""记录领取参数并返回预设结果。"""
|
|
|
|
def __init__(self, response=None):
|
|
self.response = response
|
|
self.calls = []
|
|
self.thread_ids = []
|
|
|
|
def claim_next(self, client, capabilities):
|
|
self.calls.append((client, capabilities))
|
|
self.thread_ids.append(threading.get_ident())
|
|
return self.response
|
|
|
|
def submit_result(self, task_id, idempotency_key, result):
|
|
return SubmissionReceipt(True, "RESULT-001", "2026-08-07T08:00:00Z")
|
|
|
|
def submit_failure(self, task_id, idempotency_key, failure):
|
|
return SubmissionReceipt(True, "FAILURE-001", "2026-08-07T08:00:00Z")
|
|
|
|
|
|
class SlowClaimGateway(RecordingClaimGateway):
|
|
"""让关闭测试能稳定发生在 HTTP 返回之前。"""
|
|
|
|
def __init__(self, response):
|
|
super().__init__(response)
|
|
self.started = threading.Event()
|
|
|
|
def claim_next(self, client, capabilities):
|
|
self.started.set()
|
|
time.sleep(0.1)
|
|
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 RecordingResubmitGateway(RecordingClaimGateway):
|
|
"""记录结果重新上报,并可按任务编号模拟失败。"""
|
|
|
|
def __init__(self, failures=None, delay=0.0):
|
|
super().__init__()
|
|
self.failures = dict(failures or {})
|
|
self.delay = delay
|
|
self.submit_calls = []
|
|
self.failure_calls = []
|
|
|
|
def submit_result(self, task_id, idempotency_key, result):
|
|
self.submit_calls.append((task_id, idempotency_key, result))
|
|
if self.delay:
|
|
time.sleep(self.delay)
|
|
error = self.failures.get(task_id)
|
|
if error is not None:
|
|
raise error
|
|
return SubmissionReceipt(True, f"RESULT-{task_id}", "2026-08-10T08:00:00Z")
|
|
|
|
def submit_failure(self, task_id, idempotency_key, failure):
|
|
self.failure_calls.append((task_id, idempotency_key, failure))
|
|
if self.delay:
|
|
time.sleep(self.delay)
|
|
error = self.failures.get(task_id)
|
|
if error is not None:
|
|
raise error
|
|
return SubmissionReceipt(True, f"FAILURE-{task_id}", "2026-08-10T08:00:00Z")
|
|
|
|
|
|
class FakeCollectResult:
|
|
def to_pdd_data(self):
|
|
return {
|
|
"schema_version": 1,
|
|
"goods_id": "737116531267",
|
|
"title": "测试商品",
|
|
"shop_name": "测试店铺",
|
|
"price_granularity": "color",
|
|
"dimensions": [],
|
|
"skus": [{"price_cent": 990}],
|
|
}
|
|
|
|
|
|
class FakeCollector:
|
|
def collect(self, _task):
|
|
return FakeCollectResult()
|
|
|
|
|
|
class BlockingCollector:
|
|
"""让测试能在新任务落库后、采集结束前检查界面。"""
|
|
|
|
def __init__(self, started, release):
|
|
self.started = started
|
|
self.release = release
|
|
|
|
def collect(self, _task):
|
|
self.started.set()
|
|
if not self.release.wait(timeout=3):
|
|
raise RuntimeError("测试没有及时结束阻塞采集")
|
|
return FakeCollectResult()
|
|
|
|
|
|
def fake_collect_factory(*_args):
|
|
return FakeCollector()
|
|
|
|
|
|
def collect_admin_task(task_id="COL-001"):
|
|
return AdminTask(
|
|
task_id=task_id,
|
|
task_type=TaskType.COLLECT,
|
|
version=1,
|
|
priority=0,
|
|
payload={
|
|
"goods_id": "737116531267",
|
|
"goods_url": (
|
|
"https://mobile.yangkeduo.com/goods.html"
|
|
"?goods_id=737116531267"
|
|
),
|
|
},
|
|
created_at="2026-08-07T03:19:49Z",
|
|
updated_at="2026-08-07T03:19:49Z",
|
|
)
|
|
|
|
|
|
def purchase_admin_task(task_id="PUR-001"):
|
|
return AdminTask(
|
|
task_id=task_id,
|
|
task_type=TaskType.PURCHASE,
|
|
version=1,
|
|
priority=10,
|
|
payload={
|
|
"goods_id": "737116531267",
|
|
"goods_url": "https://example.test/737116531267",
|
|
"options": {"color": "黑色", "size": "L"},
|
|
"quantity": 2,
|
|
"max_price_cent": 1200,
|
|
},
|
|
created_at="2026-08-09T08:00:00Z",
|
|
updated_at="2026-08-09T08:00:00Z",
|
|
)
|
|
|
|
|
|
class FakePurchaseAdapter(PddPurchaseAdapter):
|
|
def __init__(self):
|
|
self.options = {}
|
|
self.quantity = 0
|
|
self.page_kind = "goods"
|
|
|
|
def open_goods(self, _goods_url):
|
|
pass
|
|
|
|
def read_state(self):
|
|
return PurchasePageState(
|
|
self.page_kind,
|
|
"737116531267",
|
|
dict(self.options),
|
|
self.quantity,
|
|
990,
|
|
1,
|
|
)
|
|
|
|
def select_options(self, options):
|
|
self.options = dict(options)
|
|
|
|
def set_quantity(self, quantity):
|
|
self.quantity = quantity
|
|
|
|
def enter_confirmation(self):
|
|
self.page_kind = "order_confirmation"
|
|
|
|
def stop_before_submit(self):
|
|
pass
|
|
|
|
def close(self):
|
|
pass
|
|
|
|
|
|
def fake_purchase_factory(_address, _cancelled):
|
|
return FakePurchaseAdapter()
|
|
|
|
|
|
def wait_until(application, predicate, timeout=3.0):
|
|
deadline = time.monotonic() + timeout
|
|
while time.monotonic() < deadline:
|
|
application.processEvents()
|
|
if predicate():
|
|
application.processEvents()
|
|
return True
|
|
time.sleep(0.01)
|
|
application.processEvents()
|
|
return predicate()
|
|
|
|
|
|
class PDDTaskPageEventTest(unittest.TestCase):
|
|
@classmethod
|
|
def setUpClass(cls):
|
|
cls.app = QApplication.instance() or QApplication([])
|
|
|
|
def setUp(self):
|
|
self.temp_directory = tempfile.TemporaryDirectory()
|
|
self.db_path = Path(self.temp_directory.name) / "client.db"
|
|
self.repository = TaskRepository(self.db_path)
|
|
self.device_checker_patch = patch(
|
|
"src.pdd_ui_event.AndroidDeviceService.require_connected",
|
|
autospec=True,
|
|
)
|
|
self.device_checker = self.device_checker_patch.start()
|
|
|
|
def tearDown(self):
|
|
self.device_checker_patch.stop()
|
|
self.temp_directory.cleanup()
|
|
|
|
def _add_task(
|
|
self,
|
|
number: int,
|
|
task_type: TaskType = TaskType.COLLECT,
|
|
title: str = "测试商品",
|
|
) -> None:
|
|
self.repository.add_claimed_task(
|
|
NewClaimedTask(
|
|
remote_task_id=f"PDD-{number:03d}",
|
|
task_type=task_type,
|
|
goods_url=f"pdd://goods/{number}",
|
|
goods_id=str(number),
|
|
title=title,
|
|
target_color="黑色",
|
|
target_size="M",
|
|
price_cent=3990,
|
|
quantity=1 if task_type is TaskType.PURCHASE else None,
|
|
),
|
|
received_at="2026-08-06T08:00:00Z",
|
|
)
|
|
|
|
def _saved_settings(self):
|
|
settings = SettingsRepository(self.db_path)
|
|
settings.set_many(
|
|
{
|
|
"admin.client_id": "CLIENT-001",
|
|
"admin.client_name": "办公室电脑",
|
|
"android.selected_serial": "USB-001",
|
|
}
|
|
)
|
|
return settings
|
|
|
|
def test_initial_load_and_fetch_next_page(self):
|
|
for number in range(51):
|
|
self._add_task(number)
|
|
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(page, self.repository)
|
|
events.load_initial_tasks()
|
|
|
|
self.assertEqual(page.taskModel.data_row_count(), 50)
|
|
self.assertEqual(page.taskModel.row_at(0).remote_task_id, "PDD-050")
|
|
self.assertTrue(page.taskModel.canFetchMore())
|
|
|
|
page.taskModel.fetchMore()
|
|
|
|
self.assertEqual(page.taskModel.data_row_count(), 51)
|
|
self.assertEqual(page.taskModel.row_at(50).remote_task_id, "PDD-000")
|
|
self.assertFalse(page.taskModel.canFetchMore())
|
|
page.deleteLater()
|
|
|
|
def test_search_uses_type_status_and_keyword_together(self):
|
|
self._add_task(1, TaskType.COLLECT, "目标短袖")
|
|
self._add_task(2, TaskType.PURCHASE, "目标短袖")
|
|
self._add_task(3, TaskType.COLLECT, "其他商品")
|
|
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(page, self.repository)
|
|
page.searchRequested.emit(
|
|
{"task_type": "采集", "status": "待执行", "keyword": "目标"}
|
|
)
|
|
|
|
self.assertEqual(page.taskModel.data_row_count(), 1)
|
|
row = page.taskModel.row_at(0)
|
|
self.assertEqual(row.remote_task_id, "PDD-001")
|
|
self.assertEqual(row.task_type, "采集")
|
|
self.assertEqual(row.status, "待执行")
|
|
self.assertEqual(row.price_cents, 3990)
|
|
page.deleteLater()
|
|
|
|
def test_rerun_button_requires_exactly_one_checked_task(self):
|
|
self._add_task(1)
|
|
self._add_task(2)
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(page, self.repository)
|
|
events.load_initial_tasks()
|
|
|
|
self.assertFalse(page.rerunButton.isEnabled())
|
|
page.taskModel.setData(
|
|
page.taskModel.index(0, 0), Qt.Checked, Qt.CheckStateRole
|
|
)
|
|
self.app.processEvents()
|
|
self.assertTrue(page.rerunButton.isEnabled())
|
|
page.taskModel.setData(
|
|
page.taskModel.index(1, 0), Qt.Checked, Qt.CheckStateRole
|
|
)
|
|
self.app.processEvents()
|
|
self.assertFalse(page.rerunButton.isEnabled())
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_batch_resubmit_reuses_outbox_and_keeps_skipped_checked(self):
|
|
for number in range(1, 4):
|
|
self._add_task(number)
|
|
expected = {}
|
|
for task_id in ("PDD-001", "PDD-002"):
|
|
started = self.repository.start_collect_run(task_id, "USB-001")
|
|
event = self.repository.save_collect_result(
|
|
task_id, started.attempt_id, FakeCollectResult().to_pdd_data()
|
|
)
|
|
self.repository.mark_outbox_sent(event.id)
|
|
expected[task_id] = (event.idempotency_key, event.payload_json)
|
|
|
|
page = PDDTaskPage()
|
|
gateway = RecordingResubmitGateway()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
)
|
|
events.load_initial_tasks()
|
|
for row in range(3):
|
|
page.taskModel.setData(
|
|
page.taskModel.index(row, 0), Qt.Checked, Qt.CheckStateRole
|
|
)
|
|
|
|
with patch("src.pdd_ui_event.MessageBox") as message_box:
|
|
message_box.return_value.exec.return_value = True
|
|
page.resubmitButton.click()
|
|
self.assertTrue(
|
|
wait_until(self.app, lambda: not events._resubmit_busy),
|
|
"重新上报线程没有按时结束",
|
|
)
|
|
|
|
calls = {
|
|
task_id: (key, payload)
|
|
for task_id, key, payload in gateway.submit_calls
|
|
}
|
|
self.assertEqual(calls, expected)
|
|
self.assertEqual(page.taskModel.checked_task_ids(), ("PDD-003",))
|
|
self.assertIn("成功 2 条,失败 0 条,跳过 1 条", page.statusLabel.text())
|
|
self.device_checker.assert_not_called()
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_batch_resubmit_records_retryable_and_permanent_failures(self):
|
|
for number in range(1, 3):
|
|
self._add_task(number)
|
|
event_ids = {}
|
|
for task_id in ("PDD-001", "PDD-002"):
|
|
started = self.repository.start_collect_run(task_id, "USB-001")
|
|
event = self.repository.save_collect_result(
|
|
task_id, started.attempt_id, FakeCollectResult().to_pdd_data()
|
|
)
|
|
self.repository.mark_outbox_sent(event.id)
|
|
event_ids[task_id] = event.id
|
|
|
|
gateway = RecordingResubmitGateway(
|
|
{
|
|
"PDD-001": AdminGatewayError("TEMP", "暂时不可用", True),
|
|
"PDD-002": AdminGatewayError("REJECTED", "结果被拒绝", False),
|
|
}
|
|
)
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(page, self.repository, claim_gateway=gateway)
|
|
events.load_initial_tasks()
|
|
for row in range(2):
|
|
page.taskModel.setData(
|
|
page.taskModel.index(row, 0), Qt.Checked, Qt.CheckStateRole
|
|
)
|
|
|
|
with patch("src.pdd_ui_event.MessageBox") as message_box:
|
|
message_box.return_value.exec.return_value = True
|
|
page.resubmitButton.click()
|
|
self.assertTrue(wait_until(self.app, lambda: not events._resubmit_busy))
|
|
|
|
self.assertEqual(
|
|
self.repository.get_outbox_event(event_ids["PDD-001"]).status.value,
|
|
"pending",
|
|
)
|
|
self.assertEqual(
|
|
self.repository.get_outbox_event(event_ids["PDD-002"]).status.value,
|
|
"failed",
|
|
)
|
|
self.assertEqual(
|
|
set(page.taskModel.checked_task_ids()), {"PDD-001", "PDD-002"}
|
|
)
|
|
self.assertIn("失败 2 条", page.statusLabel.text())
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_resubmit_prioritizes_pending_failure_and_repairs_task_status(self):
|
|
self._add_task(1)
|
|
first = self.repository.start_collect_run("PDD-001", "USB-001")
|
|
result = self.repository.save_collect_result(
|
|
"PDD-001", first.attempt_id, FakeCollectResult().to_pdd_data()
|
|
)
|
|
self.repository.mark_outbox_sent(result.id)
|
|
self.repository.prepare_collect_rerun("PDD-001")
|
|
second = self.repository.start_collect_run("PDD-001", "USB-001")
|
|
failure = self.repository.save_collect_failure(
|
|
"PDD-001",
|
|
second.attempt_id,
|
|
TaskStatus.FAILED,
|
|
"SKU_PANEL_NOT_FOUND",
|
|
"规格面板加载超时",
|
|
False,
|
|
)
|
|
connection = open_database(self.db_path)
|
|
try:
|
|
with connection:
|
|
connection.execute(
|
|
"UPDATE pdd_tasks SET status = 'succeeded',"
|
|
" current_step = 'completed' WHERE remote_task_id = 'PDD-001'"
|
|
)
|
|
finally:
|
|
connection.close()
|
|
|
|
page = PDDTaskPage()
|
|
gateway = RecordingResubmitGateway()
|
|
events = PDDTaskPageEvent(page, self.repository, claim_gateway=gateway)
|
|
events.load_initial_tasks()
|
|
page.taskModel.setData(
|
|
page.taskModel.index(0, 0), Qt.Checked, Qt.CheckStateRole
|
|
)
|
|
|
|
with patch("src.pdd_ui_event.MessageBox") as message_box:
|
|
message_box.return_value.exec.return_value = True
|
|
page.resubmitButton.click()
|
|
self.assertTrue(wait_until(self.app, lambda: not events._resubmit_busy))
|
|
|
|
self.assertEqual(gateway.submit_calls, [])
|
|
self.assertEqual(len(gateway.failure_calls), 1)
|
|
task_id, idempotency_key, payload = gateway.failure_calls[0]
|
|
self.assertEqual(task_id, "PDD-001")
|
|
self.assertEqual(idempotency_key, failure.idempotency_key)
|
|
self.assertEqual(payload, failure.payload_json)
|
|
self.assertEqual(payload["status"], TaskStatus.FAILED.value)
|
|
self.assertEqual(
|
|
self.repository.get_outbox_event(failure.id).event_type,
|
|
OutboxEventType.TASK_FAILURE,
|
|
)
|
|
repaired = self.repository.get_task("PDD-001")
|
|
self.assertEqual(repaired.status, TaskStatus.FAILED)
|
|
self.assertEqual(repaired.last_error_code, "SKU_PANEL_NOT_FOUND")
|
|
self.device_checker.assert_not_called()
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_resubmit_confirmation_cancel_does_not_submit(self):
|
|
self._add_task(1)
|
|
page = PDDTaskPage()
|
|
gateway = RecordingResubmitGateway()
|
|
events = PDDTaskPageEvent(page, self.repository, claim_gateway=gateway)
|
|
events.load_initial_tasks()
|
|
page.taskModel.setData(
|
|
page.taskModel.index(0, 0), Qt.Checked, Qt.CheckStateRole
|
|
)
|
|
|
|
with patch("src.pdd_ui_event.MessageBox") as message_box:
|
|
message_box.return_value.exec.return_value = False
|
|
page.resubmitButton.click()
|
|
|
|
self.assertEqual(gateway.submit_calls, [])
|
|
self.assertFalse(events._resubmit_busy)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_resubmit_rejects_duplicate_start_and_shutdown_waits_worker(self):
|
|
self._add_task(1)
|
|
started = self.repository.start_collect_run("PDD-001", "USB-001")
|
|
event = self.repository.save_collect_result(
|
|
"PDD-001", started.attempt_id, FakeCollectResult().to_pdd_data()
|
|
)
|
|
self.repository.mark_outbox_sent(event.id)
|
|
page = PDDTaskPage()
|
|
gateway = RecordingResubmitGateway(delay=0.1)
|
|
events = PDDTaskPageEvent(page, self.repository, claim_gateway=gateway)
|
|
events.load_initial_tasks()
|
|
page.taskModel.setData(
|
|
page.taskModel.index(0, 0), Qt.Checked, Qt.CheckStateRole
|
|
)
|
|
|
|
with patch("src.pdd_ui_event.MessageBox") as message_box:
|
|
message_box.return_value.exec.return_value = True
|
|
page.resubmitButton.click()
|
|
events.request_resubmit(("PDD-001",))
|
|
self.assertTrue(
|
|
wait_until(self.app, lambda: bool(gateway.submit_calls))
|
|
)
|
|
events.shutdown()
|
|
|
|
self.assertEqual(len(gateway.submit_calls), 1)
|
|
self.assertFalse(
|
|
events._resubmit_thread and events._resubmit_thread.isRunning()
|
|
)
|
|
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_remove_confirmation_cancel_keeps_local_task(self):
|
|
self._add_task(1)
|
|
connection = open_database(self.db_path)
|
|
try:
|
|
with connection:
|
|
connection.execute(
|
|
"UPDATE pdd_tasks SET status = 'succeeded'"
|
|
" WHERE remote_task_id = 'PDD-001'"
|
|
)
|
|
finally:
|
|
connection.close()
|
|
page = PDDTaskPage()
|
|
gateway = RecordingClaimGateway()
|
|
events = PDDTaskPageEvent(page, self.repository, claim_gateway=gateway)
|
|
events.load_initial_tasks()
|
|
page.taskModel.setData(
|
|
page.taskModel.index(0, 0), Qt.Checked, Qt.CheckStateRole
|
|
)
|
|
|
|
with patch("src.pdd_ui_event.MessageBox") as message_box:
|
|
dialog = message_box.return_value
|
|
dialog.exec.return_value = False
|
|
page.removeButton.click()
|
|
|
|
message_box.assert_called_once()
|
|
self.assertIs(message_box.call_args.args[2], page.window())
|
|
dialog.cancelButton.setFocus.assert_called_once()
|
|
self.assertEqual(self.repository.count_tasks(), 1)
|
|
self.assertFalse(events._remove_busy)
|
|
self.assertEqual(gateway.calls, [])
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_remove_terminal_task_refreshes_list_and_keeps_audit_data(self):
|
|
self._add_task(1)
|
|
connection = open_database(self.db_path)
|
|
try:
|
|
with connection:
|
|
connection.execute(
|
|
"UPDATE pdd_tasks SET status = 'failed'"
|
|
" WHERE remote_task_id = 'PDD-001'"
|
|
)
|
|
finally:
|
|
connection.close()
|
|
page = PDDTaskPage()
|
|
gateway = RecordingClaimGateway()
|
|
events = PDDTaskPageEvent(page, self.repository, claim_gateway=gateway)
|
|
events.load_initial_tasks()
|
|
page.taskModel.setData(
|
|
page.taskModel.index(0, 0), Qt.Checked, Qt.CheckStateRole
|
|
)
|
|
|
|
with patch("src.pdd_ui_event.MessageBox") as message_box:
|
|
message_box.return_value.exec.return_value = True
|
|
page.removeButton.click()
|
|
self.assertTrue(wait_until(self.app, lambda: not events._remove_busy))
|
|
|
|
self.assertEqual(page.taskModel.data_row_count(), 0)
|
|
self.assertEqual(page.taskModel.checked_task_ids(), ())
|
|
self.assertIsNotNone(self.repository.get_task("PDD-001"))
|
|
self.assertIn("已从列表删除 1 条", page.statusLabel.text())
|
|
self.assertEqual(gateway.calls, [])
|
|
self.device_checker.assert_not_called()
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_remove_mixed_batch_reports_reason_and_changes_nothing(self):
|
|
self._add_task(1)
|
|
self._add_task(2)
|
|
connection = open_database(self.db_path)
|
|
try:
|
|
with connection:
|
|
connection.execute(
|
|
"UPDATE pdd_tasks SET status = 'cancelled'"
|
|
" WHERE remote_task_id = 'PDD-001'"
|
|
)
|
|
finally:
|
|
connection.close()
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page, self.repository, claim_gateway=RecordingClaimGateway()
|
|
)
|
|
events.load_initial_tasks()
|
|
for row in range(2):
|
|
page.taskModel.setData(
|
|
page.taskModel.index(row, 0), Qt.Checked, Qt.CheckStateRole
|
|
)
|
|
|
|
with patch("src.pdd_ui_event.MessageBox") as message_box:
|
|
message_box.return_value.exec.return_value = True
|
|
page.removeButton.click()
|
|
self.assertTrue(wait_until(self.app, lambda: not events._remove_busy))
|
|
|
|
self.assertEqual(self.repository.count_tasks(), 2)
|
|
self.assertEqual(len(page.taskModel.checked_task_ids()), 2)
|
|
self.assertIn("尚未结束", page.statusLabel.text())
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_rerun_confirmation_has_clear_safe_cancel_action(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()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=RecordingClaimGateway(),
|
|
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 = False
|
|
page.rerunRequested.emit("PDD-001")
|
|
|
|
message_box.return_value.cancelButton.setText.assert_called_once_with(
|
|
"暂不重新采集"
|
|
)
|
|
message_box.return_value.cancelButton.setFocus.assert_called_once_with()
|
|
self.assertFalse(events._claim_busy)
|
|
self.assertEqual(
|
|
self.repository.get_task("PDD-001").status, TaskStatus.SUCCEEDED
|
|
)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_running_rerun_can_request_cooperative_cancel_once(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)
|
|
started = threading.Event()
|
|
|
|
def cancellable_factory(_address, _client_id, cancelled):
|
|
class CancellableCollector:
|
|
def collect(self, _task):
|
|
started.set()
|
|
while not cancelled():
|
|
time.sleep(0.01)
|
|
raise PddCollectError("PDD_CANCELLED", "采集任务已安全取消")
|
|
|
|
return CancellableCollector()
|
|
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=RecordingClaimGateway(),
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=cancellable_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(started.wait(1.0), "重新采集没有启动")
|
|
self.assertTrue(
|
|
wait_until(
|
|
self.app,
|
|
lambda: page.rerunButton.text() == "停止重新采集",
|
|
)
|
|
)
|
|
self.assertTrue(page.rerunButton.isEnabled())
|
|
page.rerunButton.click()
|
|
self.app.processEvents()
|
|
self.assertEqual(page.rerunButton.text(), "正在停止…")
|
|
self.assertFalse(page.rerunButton.isEnabled())
|
|
self.assertIn("等待手机当前操作结束", page.statusLabel.text())
|
|
self.assertTrue(
|
|
wait_until(self.app, lambda: not events._claim_busy),
|
|
"取消后工作线程没有按时结束",
|
|
)
|
|
|
|
self.assertEqual(page.rerunButton.text(), "重新执行")
|
|
self.assertIn("停止请求已处理", page.statusLabel.text())
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_rerun_feedback_has_close_button_and_replaces_previous(self):
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(page, self.repository)
|
|
|
|
events._show_rerun_warning("暂时不能重新执行", "测试警告")
|
|
first = events._rerun_feedback
|
|
self.assertIsNotNone(first)
|
|
self.assertEqual(first.position, InfoBarPosition.TOP)
|
|
events._show_rerun_error("重新采集失败", "测试错误")
|
|
second = events._rerun_feedback
|
|
self.assertIsNot(first, second)
|
|
self.assertEqual(second.position, InfoBarPosition.TOP)
|
|
self.assertEqual(second.duration, ERROR_FEEDBACK_DURATION_MS)
|
|
|
|
close_buttons = [
|
|
button
|
|
for button in second.findChildren(type(page.rerunButton))
|
|
if button.text() == "关闭提示"
|
|
]
|
|
self.assertEqual(len(close_buttons), 1)
|
|
self.assertGreaterEqual(second.minimumWidth(), ERROR_FEEDBACK_MIN_WIDTH)
|
|
self.assertGreaterEqual(second.minimumHeight(), ERROR_FEEDBACK_MIN_HEIGHT)
|
|
self.assertGreaterEqual(
|
|
close_buttons[0].minimumWidth(),
|
|
ERROR_FEEDBACK_BUTTON_MIN_WIDTH,
|
|
)
|
|
self.assertGreaterEqual(
|
|
close_buttons[0].minimumHeight(),
|
|
ERROR_FEEDBACK_BUTTON_MIN_HEIGHT,
|
|
)
|
|
close_buttons[0].click()
|
|
self.app.processEvents()
|
|
self.assertIsNone(events._rerun_feedback)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_claim_error_is_large_closable_and_replaces_previous(self):
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(page, self.repository)
|
|
|
|
events._show_claim_error("采集失败", "第一条错误")
|
|
first = events._claim_feedback
|
|
events._show_claim_error("采购失败", "第二条错误")
|
|
second = events._claim_feedback
|
|
|
|
self.assertIsNot(first, second)
|
|
self.assertEqual(second.position, InfoBarPosition.TOP)
|
|
self.assertEqual(second.duration, ERROR_FEEDBACK_DURATION_MS)
|
|
self.assertGreaterEqual(second.minimumWidth(), ERROR_FEEDBACK_MIN_WIDTH)
|
|
self.assertGreaterEqual(second.minimumHeight(), ERROR_FEEDBACK_MIN_HEIGHT)
|
|
close_buttons = [
|
|
button
|
|
for button in second.findChildren(type(page.rerunButton))
|
|
if button.text() == "关闭提示"
|
|
]
|
|
self.assertEqual(len(close_buttons), 1)
|
|
self.assertGreaterEqual(
|
|
close_buttons[0].minimumWidth(),
|
|
ERROR_FEEDBACK_BUTTON_MIN_WIDTH,
|
|
)
|
|
self.assertGreaterEqual(
|
|
close_buttons[0].minimumHeight(),
|
|
ERROR_FEEDBACK_BUTTON_MIN_HEIGHT,
|
|
)
|
|
close_buttons[0].click()
|
|
self.app.processEvents()
|
|
self.assertIsNone(events._claim_feedback)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_late_rerun_success_does_not_override_cancel_status(self):
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(page, self.repository)
|
|
events._rerun_cancel_requested = True
|
|
|
|
events._on_rerun_outcome("succeeded", "采集成功", "PDD-001")
|
|
|
|
self.assertNotIn("采集成功", page.statusLabel.text())
|
|
self.assertIn("停止请求已处理", page.statusLabel.text())
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_disconnected_device_stops_auto_fetch_before_claim(self):
|
|
page = PDDTaskPage()
|
|
gateway = RecordingClaimGateway(collect_admin_task())
|
|
checker_threads = []
|
|
|
|
def disconnected(_serial):
|
|
checker_threads.append(threading.get_ident())
|
|
raise AndroidDeviceSearchError("USB Android 设备 USB-001 未连接")
|
|
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=fake_collect_factory,
|
|
device_connection_checker=disconnected,
|
|
)
|
|
page.autoFetchRequested.emit()
|
|
|
|
self.assertTrue(wait_until(self.app, lambda: not events._claim_busy))
|
|
self.assertEqual(gateway.calls, [])
|
|
self.assertEqual(self.repository.count_tasks(), 0)
|
|
self.assertEqual(len(checker_threads), 1)
|
|
self.assertNotEqual(checker_threads[0], threading.get_ident())
|
|
self.assertIn("Android 设备不可用", page.statusLabel.text())
|
|
self.assertIsNotNone(events._device_feedback)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_disconnected_rerun_keeps_existing_task_unchanged(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)
|
|
before = self.repository.get_task("PDD-001")
|
|
before_run = self.repository.latest_task_run("PDD-001")
|
|
page = PDDTaskPage()
|
|
|
|
def disconnected(_serial):
|
|
raise AndroidDeviceSearchError(
|
|
"USB Android 设备 USB-001 未连接"
|
|
)
|
|
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=RecordingClaimGateway(),
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=fake_collect_factory,
|
|
device_connection_checker=disconnected,
|
|
)
|
|
|
|
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))
|
|
|
|
after = self.repository.get_task("PDD-001")
|
|
after_run = self.repository.latest_task_run("PDD-001")
|
|
self.assertEqual(after.status, TaskStatus.SUCCEEDED)
|
|
self.assertEqual(after.pdd_data, before.pdd_data)
|
|
self.assertEqual(after_run.attempt_id, before_run.attempt_id)
|
|
self.assertIn("重新采集未开始", page.statusLabel.text())
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_device_feedback_opens_settings_closes_and_does_not_stack(self):
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(page, self.repository)
|
|
opened = []
|
|
page.openSettingsRequested.connect(lambda: opened.append(True))
|
|
|
|
events._show_device_unavailable("第一条")
|
|
first = events._device_feedback
|
|
events._show_device_unavailable("第二条")
|
|
second = events._device_feedback
|
|
self.assertIsNot(first, second)
|
|
self.assertEqual(second.position, InfoBarPosition.TOP)
|
|
self.assertEqual(second.duration, ERROR_FEEDBACK_DURATION_MS)
|
|
self.assertGreaterEqual(second.minimumWidth(), ERROR_FEEDBACK_MIN_WIDTH)
|
|
self.assertGreaterEqual(second.minimumHeight(), ERROR_FEEDBACK_MIN_HEIGHT)
|
|
buttons = {
|
|
button.text(): button
|
|
for button in second.findChildren(type(page.rerunButton))
|
|
}
|
|
self.assertIn("打开设置", buttons)
|
|
self.assertIn("关闭提示", buttons)
|
|
for button in buttons.values():
|
|
self.assertGreaterEqual(
|
|
button.minimumWidth(),
|
|
ERROR_FEEDBACK_BUTTON_MIN_WIDTH,
|
|
)
|
|
self.assertGreaterEqual(
|
|
button.minimumHeight(),
|
|
ERROR_FEEDBACK_BUTTON_MIN_HEIGHT,
|
|
)
|
|
buttons["打开设置"].click()
|
|
self.app.processEvents()
|
|
self.assertEqual(opened, [True])
|
|
self.assertIsNone(events._device_feedback)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_query_failure_shows_readable_error(self):
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(page, BrokenRepository())
|
|
|
|
events.load_initial_tasks()
|
|
|
|
self.assertFalse(page.emptyStateCard.isHidden())
|
|
self.assertEqual(page.emptyTitleLabel.text(), "任务加载失败")
|
|
self.assertIn("检查数据库", page.emptyMessageLabel.text())
|
|
self.assertFalse(page.taskModel.canFetchMore())
|
|
page.deleteLater()
|
|
|
|
def test_detail_signal_opens_one_reusable_non_modal_window(self):
|
|
self._add_task(1)
|
|
started = self.repository.start_collect_run("PDD-001", "USB-001")
|
|
self.repository.save_collect_result(
|
|
"PDD-001",
|
|
started.attempt_id,
|
|
{
|
|
"goods_id": "1",
|
|
"title": "测试商品",
|
|
"dimensions": [
|
|
{
|
|
"key": "color",
|
|
"name": "颜色",
|
|
"values": [{"text": "黑色", "available": True}],
|
|
},
|
|
{
|
|
"key": "size",
|
|
"name": "尺码",
|
|
"values": [{"text": "M", "available": True}],
|
|
},
|
|
],
|
|
"skus": [
|
|
{
|
|
"options": {"color": "黑色", "size": "M"},
|
|
"price_cent": 3990,
|
|
}
|
|
],
|
|
},
|
|
)
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(page, self.repository)
|
|
|
|
page.detailRequested.emit("PDD-001")
|
|
self.app.processEvents()
|
|
first = events._detail_windows["PDD-001"]
|
|
self.assertTrue(first.isVisible())
|
|
self.assertFalse(first.isModal())
|
|
|
|
page.detailRequested.emit("PDD-001")
|
|
self.app.processEvents()
|
|
self.assertIs(events._detail_windows["PDD-001"], first)
|
|
|
|
first.close()
|
|
self.app.processEvents()
|
|
self.assertNotIn("PDD-001", events._detail_windows)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_missing_detail_does_not_open_empty_window(self):
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(page, self.repository)
|
|
|
|
page.detailRequested.emit("PDD-NOT-FOUND")
|
|
self.app.processEvents()
|
|
|
|
self.assertEqual(events._detail_windows, {})
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_summary_to_row_does_not_copy_detail_json(self):
|
|
summary = TaskSummary(
|
|
id=1,
|
|
remote_task_id="PDD-001",
|
|
task_type=TaskType.PURCHASE,
|
|
goods_id=None,
|
|
title=None,
|
|
target_color=None,
|
|
target_size=None,
|
|
price_cent=None,
|
|
quantity=2,
|
|
status=TaskStatus.MANUAL_REVIEW,
|
|
updated_at="2026-08-06T08:00:00Z",
|
|
)
|
|
|
|
row = summary_to_row(summary)
|
|
|
|
self.assertEqual(row.task_type, "采购")
|
|
self.assertEqual(row.status, "需要人工处理")
|
|
self.assertEqual(row.goods_id, "")
|
|
self.assertFalse(hasattr(row, "pdd_data"))
|
|
|
|
def test_main_window_keeps_event_object_alive(self):
|
|
window = MainWindow(
|
|
task_repository=self.repository,
|
|
settings_repository=SettingsRepository(self.db_path),
|
|
admin_gateway=MockAdminGateway(),
|
|
)
|
|
|
|
self.assertIsInstance(window.pddTaskPageEvent, PDDTaskPageEvent)
|
|
self.assertIs(
|
|
window.pddTaskPageEvent._purchase_adapter_factory,
|
|
create_u2_purchase_adapter,
|
|
)
|
|
self.assertEqual(window.pddTaskPage.taskModel.data_row_count(), 0)
|
|
window.pddTaskPage.openSettingsRequested.emit()
|
|
self.app.processEvents()
|
|
self.assertIs(window.stackedWidget.currentWidget(), window.settingsPage)
|
|
window.close()
|
|
window.deleteLater()
|
|
|
|
def test_admin_task_mapping_uses_explicit_real_field_names(self):
|
|
local_task = admin_task_to_new_claimed_task(
|
|
collect_admin_task("COL-8020a8729f111c15")
|
|
)
|
|
|
|
self.assertEqual(local_task.remote_task_id, "COL-8020a8729f111c15")
|
|
self.assertEqual(local_task.task_type, TaskType.COLLECT)
|
|
self.assertEqual(local_task.goods_id, "737116531267")
|
|
self.assertEqual(
|
|
local_task.admin_payload["payload"]["goods_id"],
|
|
"737116531267",
|
|
)
|
|
|
|
def test_start_claims_task_and_keeps_auto_fetch_running(self):
|
|
gateway = RecordingClaimGateway(collect_admin_task())
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=fake_collect_factory,
|
|
)
|
|
|
|
page.autoFetchRequested.emit()
|
|
|
|
self.assertTrue(
|
|
wait_until(self.app, lambda: not events._claim_busy)
|
|
)
|
|
self.assertEqual(len(gateway.calls), 1)
|
|
self.assertNotEqual(gateway.thread_ids[0], threading.get_ident())
|
|
_, capabilities = gateway.calls[0]
|
|
self.assertEqual(
|
|
[task_type.value for task_type in capabilities.supported_types],
|
|
["collect"],
|
|
)
|
|
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.assertIn("COL-001", page.statusLabel.text())
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_auto_fetch_reuses_one_fixed_worker_thread_across_cycles(self):
|
|
gateway = RecordingClaimGateway()
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=fake_collect_factory,
|
|
no_task_delay_ms=10,
|
|
)
|
|
|
|
page.autoFetchRequested.emit()
|
|
self.assertTrue(
|
|
wait_until(self.app, lambda: len(gateway.thread_ids) >= 2),
|
|
"自动获取没有完成两轮串行调度",
|
|
)
|
|
executor_thread = events._claim_thread
|
|
page.autoFetchRequested.emit()
|
|
self.assertTrue(wait_until(self.app, lambda: not events._claim_busy))
|
|
|
|
self.assertIs(events._claim_thread, executor_thread)
|
|
self.assertTrue(
|
|
executor_thread is not None and executor_thread.isRunning()
|
|
)
|
|
self.assertEqual(len(set(gateway.thread_ids)), 1)
|
|
events.shutdown()
|
|
self.assertFalse(executor_thread.isRunning())
|
|
page.deleteLater()
|
|
|
|
def test_new_claim_refreshes_table_before_collection_finishes(self):
|
|
gateway = RecordingClaimGateway(collect_admin_task("COL-EARLY"))
|
|
collect_started = threading.Event()
|
|
release_collect = threading.Event()
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=(
|
|
lambda *_args: BlockingCollector(
|
|
collect_started, release_collect
|
|
)
|
|
),
|
|
)
|
|
|
|
page.autoFetchRequested.emit()
|
|
try:
|
|
self.assertTrue(wait_until(self.app, collect_started.is_set))
|
|
self.assertTrue(events._claim_busy)
|
|
self.assertTrue(
|
|
wait_until(
|
|
self.app,
|
|
lambda: page.taskModel.data_row_count() == 1,
|
|
)
|
|
)
|
|
self.assertEqual(
|
|
page.taskModel.row_at(0).remote_task_id, "COL-EARLY"
|
|
)
|
|
finally:
|
|
release_collect.set()
|
|
|
|
self.assertTrue(wait_until(self.app, lambda: not events._claim_busy))
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_ready_runtime_claims_and_dispatches_purchase_in_worker(self):
|
|
gateway = RecordingClaimGateway(purchase_admin_task())
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=fake_collect_factory,
|
|
purchase_adapter_factory=fake_purchase_factory,
|
|
)
|
|
|
|
page.autoFetchRequested.emit()
|
|
|
|
self.assertTrue(wait_until(self.app, lambda: not events._claim_busy))
|
|
_, capabilities = gateway.calls[0]
|
|
self.assertEqual(
|
|
capabilities.supported_types,
|
|
(TaskType.COLLECT, TaskType.PURCHASE),
|
|
)
|
|
self.assertEqual(capabilities.purchase_mode, "dry_run")
|
|
detail = self.repository.get_task("PUR-001")
|
|
assert detail is not None
|
|
self.assertEqual(detail.task_type, TaskType.PURCHASE)
|
|
self.assertEqual(detail.status, TaskStatus.SUCCEEDED)
|
|
self.assertIn("演练完成", page.statusLabel.text())
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_ready_runtime_automatically_declares_live_capability(self):
|
|
gateway = RecordingClaimGateway(None)
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
purchase_adapter_factory=fake_purchase_factory,
|
|
live_purchase_adapter_factory=lambda *_args: None,
|
|
)
|
|
|
|
page.autoFetchRequested.emit()
|
|
|
|
self.assertTrue(wait_until(self.app, lambda: not events._claim_busy))
|
|
_, capabilities = gateway.calls[0]
|
|
self.assertEqual(capabilities.purchase_mode, "live")
|
|
self.assertIn(TaskType.PURCHASE, capabilities.supported_types)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_no_task_schedules_next_claim_and_stop_cancels_timer(self):
|
|
gateway = RecordingClaimGateway(None)
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=fake_collect_factory,
|
|
)
|
|
|
|
page.autoFetchRequested.emit()
|
|
|
|
self.assertTrue(wait_until(self.app, lambda: not events._claim_busy))
|
|
self.assertEqual(len(gateway.calls), 1)
|
|
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_second_click_while_claiming_requests_safe_stop(self):
|
|
gateway = SlowClaimGateway(None)
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=fake_collect_factory,
|
|
)
|
|
|
|
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()
|
|
|
|
def test_retry_wait_failure_shows_paused_recovery_actions(self):
|
|
self._add_task(1)
|
|
started = self.repository.start_collect_run("PDD-001", "USB-001")
|
|
event = self.repository.save_collect_failure(
|
|
"PDD-001",
|
|
started.attempt_id,
|
|
TaskStatus.RETRY_WAIT,
|
|
"DEVICE_OFFLINE",
|
|
"设备离线",
|
|
True,
|
|
)
|
|
self.repository.mark_outbox_sent(event.id)
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=RecordingClaimGateway(None),
|
|
settings_repository=self._saved_settings(),
|
|
)
|
|
|
|
events._on_collect_outcome(
|
|
"failed",
|
|
"任务 PDD-001 采集未完成",
|
|
"PDD-001",
|
|
)
|
|
|
|
self.assertIn("重试已暂停", events._stop_status)
|
|
self.assertIn("重新执行", events._stop_status)
|
|
self.assertIn("重新启动获取任务", events._stop_status)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_failed_task_does_not_show_retry_paused_status(self):
|
|
self._add_task(1)
|
|
started = self.repository.start_collect_run("PDD-001", "USB-001")
|
|
event = self.repository.save_collect_failure(
|
|
"PDD-001",
|
|
started.attempt_id,
|
|
TaskStatus.FAILED,
|
|
"GOODS_OFF_SHELF",
|
|
"商品已下架",
|
|
False,
|
|
)
|
|
self.repository.mark_outbox_sent(event.id)
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=RecordingClaimGateway(None),
|
|
settings_repository=self._saved_settings(),
|
|
)
|
|
|
|
events._on_collect_outcome("failed", "商品已下架", "PDD-001")
|
|
|
|
self.assertEqual(events._stop_status, "商品已下架")
|
|
self.assertNotIn("重试已暂停", events._stop_status)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_unavailable_goods_continues_auto_fetch_after_current_task(self):
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=RecordingClaimGateway(None),
|
|
settings_repository=self._saved_settings(),
|
|
next_task_delay_ms=20,
|
|
)
|
|
events._auto_fetch_running = True
|
|
|
|
events._on_collect_outcome(
|
|
"business_failed",
|
|
"商品链接已失效,PDD 无法打开商品详情页并返回了首页",
|
|
"PDD-001",
|
|
)
|
|
|
|
self.assertEqual(20, events._cycle_next_delay_ms)
|
|
self.assertTrue(events._auto_fetch_running)
|
|
self.assertFalse(events._stop_requested)
|
|
self.assertIn("商品链接已失效", page.statusLabel.text())
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_live_submit_schedules_read_only_reconcile_after_five_seconds(self):
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=RecordingClaimGateway(None),
|
|
settings_repository=self._saved_settings(),
|
|
)
|
|
events._auto_fetch_running = True
|
|
|
|
events._on_collect_outcome(
|
|
"reconcile_pending",
|
|
"任务 PUR-001 订单提交点击已执行,结果待只读核对",
|
|
"PUR-001",
|
|
)
|
|
|
|
self.assertEqual(
|
|
events._cycle_next_delay_ms, PURCHASE_RECONCILE_DELAY_MS
|
|
)
|
|
self.assertTrue(events._auto_fetch_running)
|
|
self.assertFalse(events._stop_requested)
|
|
self.assertIn("5 秒后自动核对订单", page.statusLabel.text())
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_stop_during_reconcile_wait_cancels_next_cycle(self):
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=RecordingClaimGateway(None),
|
|
settings_repository=self._saved_settings(),
|
|
)
|
|
events._auto_fetch_running = True
|
|
events._on_collect_outcome(
|
|
"reconcile_pending", "等待只读核对", "PUR-001"
|
|
)
|
|
events._on_claim_thread_finished()
|
|
self.assertTrue(events._next_cycle_timer.isActive())
|
|
self.assertEqual(
|
|
events._next_cycle_timer.interval(), PURCHASE_RECONCILE_DELAY_MS
|
|
)
|
|
|
|
events._request_stop_auto_fetch()
|
|
|
|
self.assertFalse(events._next_cycle_timer.isActive())
|
|
self.assertFalse(events._auto_fetch_running)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_duplicate_task_is_normal_status(self):
|
|
task = collect_admin_task()
|
|
self.repository.add_claimed_task(admin_task_to_new_claimed_task(task))
|
|
gateway = RecordingClaimGateway(task)
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=fake_collect_factory,
|
|
)
|
|
|
|
page.autoFetchRequested.emit()
|
|
|
|
self.assertTrue(wait_until(self.app, lambda: not events._claim_busy))
|
|
self.assertIn("采集完成", page.statusLabel.text())
|
|
self.assertEqual(gateway.calls, [])
|
|
self.assertEqual(self.repository.count_tasks(), 1)
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_local_save_failure_status_contains_claimed_task_id(self):
|
|
gateway = RecordingClaimGateway(collect_admin_task("COL-LOST"))
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
BrokenSaveRepository(),
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=fake_collect_factory,
|
|
)
|
|
|
|
page.autoFetchRequested.emit()
|
|
|
|
self.assertTrue(wait_until(self.app, lambda: not events._claim_busy))
|
|
self.assertIn("COL-LOST", page.statusLabel.text())
|
|
self.assertIn("本地保存失败", page.statusLabel.text())
|
|
self.assertIn("disk is full", page.statusLabel.text())
|
|
events.shutdown()
|
|
page.deleteLater()
|
|
|
|
def test_missing_client_settings_do_not_call_admin(self):
|
|
gateway = RecordingClaimGateway(collect_admin_task())
|
|
page = PDDTaskPage()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=SettingsRepository(self.db_path),
|
|
)
|
|
|
|
page.autoFetchRequested.emit()
|
|
|
|
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()
|
|
events = PDDTaskPageEvent(
|
|
page,
|
|
self.repository,
|
|
claim_gateway=gateway,
|
|
settings_repository=self._saved_settings(),
|
|
collect_service_factory=fake_collect_factory,
|
|
)
|
|
page.autoFetchRequested.emit()
|
|
self.assertTrue(wait_until(self.app, gateway.started.is_set))
|
|
status_before_close = page.statusLabel.text()
|
|
|
|
events.shutdown()
|
|
self.app.processEvents()
|
|
|
|
self.assertIsNotNone(self.repository.get_task("COL-CLOSE"))
|
|
self.assertEqual(page.statusLabel.text(), status_before_close)
|
|
page.deleteLater()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|