Files
cmautobuy/client/test/test_pdd_ui_event.py
T

624 lines
21 KiB
Python

"""PDD 任务列表接入 SQLite 的离屏测试。"""
import os
import tempfile
import threading
import time
import unittest
from pathlib import Path
os.environ.setdefault("QT_QPA_PLATFORM", "offscreen")
from PyQt5.QtWidgets import QApplication
from src.pdd_ui import PDDTaskPage
from src.admin_gateway import AdminGatewayError, AdminTask, SubmissionReceipt
from src.pdd_ui_event import (
PDDTaskPageEvent,
admin_task_to_new_claimed_task,
summary_to_row,
)
from src.mock_admin_gateway import MockAdminGateway
from src.settings_repository import SettingsRepository
from src.task_models import NewClaimedTask, 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
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 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()
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 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)
def tearDown(self):
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_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.assertEqual(window.pddTaskPage.taskModel.data_row_count(), 0)
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_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_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()