247 lines
8.3 KiB
Python
247 lines
8.3 KiB
Python
"""采集任务应用流程测试,不连接真机和网络。"""
|
|
|
|
import tempfile
|
|
import unittest
|
|
from pathlib import Path
|
|
|
|
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
|
|
from src.task_models import NewClaimedTask, TaskStatus, TaskType
|
|
from src.task_repository import TaskRepository
|
|
|
|
|
|
class FakeResult:
|
|
def to_pdd_data(self):
|
|
return {
|
|
"schema_version": 1,
|
|
"goods_id": "737116531267",
|
|
"title": "测试商品",
|
|
"shop_name": "测试店铺",
|
|
"price_granularity": "color",
|
|
"dimensions": [
|
|
{"key": "color", "name": "颜色分类"},
|
|
{"key": "size", "name": "尺码"},
|
|
],
|
|
"skus": [
|
|
{
|
|
"options": {"color": "黑色", "size": "M"},
|
|
"price_cent": 990,
|
|
"price_observed_at": {"color": "黑色", "size": "M"},
|
|
"available": True,
|
|
}
|
|
],
|
|
}
|
|
|
|
|
|
class FakeCollector:
|
|
def __init__(self, calls, error=None):
|
|
self.calls = calls
|
|
self.error = error
|
|
|
|
def collect(self, task):
|
|
self.calls.append(task.remote_task_id)
|
|
if self.error is not None:
|
|
raise self.error
|
|
return FakeResult()
|
|
|
|
|
|
class CollectTaskServiceTest(unittest.TestCase):
|
|
def setUp(self):
|
|
self.temporary = tempfile.TemporaryDirectory()
|
|
self.repository = TaskRepository(Path(self.temporary.name) / "client.db")
|
|
self.gateway = MockAdminGateway()
|
|
self.client = ClientInfo("CLIENT-001", "测试电脑")
|
|
self.task = AdminTask(
|
|
task_id="COL-001",
|
|
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-07T08:00:00Z",
|
|
updated_at="2026-08-07T08:00:00Z",
|
|
)
|
|
self.gateway.enqueue_task(self.task, self.client.client_id)
|
|
|
|
def tearDown(self):
|
|
self.temporary.cleanup()
|
|
|
|
def _service(self, calls, error=None):
|
|
return CollectTaskService(
|
|
self.gateway,
|
|
self.repository,
|
|
self.client,
|
|
"USB-001",
|
|
collect_service_factory=lambda *_args: FakeCollector(calls, error),
|
|
)
|
|
|
|
def test_claim_collect_persist_submit_completes_one_task(self):
|
|
calls = []
|
|
|
|
outcome = self._service(calls).execute_one()
|
|
|
|
self.assertEqual(outcome.kind, "succeeded")
|
|
self.assertEqual(calls, ["COL-001"])
|
|
detail = self.repository.get_task("COL-001")
|
|
self.assertEqual(detail.status, TaskStatus.SUCCEEDED)
|
|
self.assertEqual(detail.pdd_data["price_granularity"], "color")
|
|
self.assertEqual(self.gateway.submission_count, 1)
|
|
|
|
def test_submit_timeout_retries_stored_outbox_without_recollecting(self):
|
|
remote = self.gateway.claim_next(
|
|
self.client,
|
|
ClaimCapabilities(supported_types=(TaskType.COLLECT,)),
|
|
)
|
|
self.repository.add_claimed_task(
|
|
NewClaimedTask(
|
|
remote_task_id=remote.task_id,
|
|
task_type=remote.task_type,
|
|
goods_url=remote.payload["goods_url"],
|
|
goods_id=remote.payload["goods_id"],
|
|
version=remote.version,
|
|
admin_payload={"payload": dict(remote.payload)},
|
|
)
|
|
)
|
|
calls = []
|
|
self.gateway.timeout_next_call()
|
|
|
|
first = self._service(calls).execute_one()
|
|
second = self._service(calls).execute_one()
|
|
|
|
self.assertEqual(first.kind, "result_pending")
|
|
self.assertEqual(second.kind, "succeeded")
|
|
self.assertEqual(calls, ["COL-001"])
|
|
self.assertEqual(self.gateway.submission_count, 1)
|
|
|
|
def test_pending_outbox_can_submit_without_android_device(self):
|
|
remote = self.gateway.claim_next(
|
|
self.client,
|
|
ClaimCapabilities(supported_types=(TaskType.COLLECT,)),
|
|
)
|
|
self.repository.add_claimed_task(
|
|
NewClaimedTask(
|
|
remote_task_id=remote.task_id,
|
|
task_type=remote.task_type,
|
|
goods_url=remote.payload["goods_url"],
|
|
goods_id=remote.payload["goods_id"],
|
|
version=remote.version,
|
|
)
|
|
)
|
|
started = self.repository.start_collect_run(remote.task_id, "USB-001")
|
|
self.repository.save_collect_result(
|
|
remote.task_id, started.attempt_id, FakeResult().to_pdd_data()
|
|
)
|
|
|
|
outcome = CollectTaskService(
|
|
self.gateway, self.repository, self.client, ""
|
|
).execute_one()
|
|
|
|
self.assertEqual(outcome.kind, "succeeded")
|
|
self.assertEqual(
|
|
self.repository.get_task(remote.task_id).status,
|
|
TaskStatus.SUCCEEDED,
|
|
)
|
|
|
|
def test_captcha_becomes_manual_review_and_is_reported(self):
|
|
calls = []
|
|
outcome = self._service(
|
|
calls, PddCollectError("PDD_PAGE_CAPTCHA", "需要验证")
|
|
).execute_one()
|
|
|
|
self.assertEqual(outcome.kind, "failed")
|
|
self.assertEqual(
|
|
self.repository.get_task("COL-001").status,
|
|
TaskStatus.MANUAL_REVIEW,
|
|
)
|
|
self.assertEqual(self.gateway.submission_count, 1)
|
|
|
|
def test_unavailable_goods_is_terminal_but_allows_next_task(self):
|
|
message = "商品链接已失效,PDD 无法打开商品详情页并返回了首页"
|
|
outcome = self._service(
|
|
[], PddCollectError("PDD_GOODS_UNAVAILABLE", message)
|
|
).execute_one()
|
|
|
|
self.assertEqual("business_failed", outcome.kind)
|
|
self.assertEqual(message, outcome.message)
|
|
detail = self.repository.get_task("COL-001")
|
|
self.assertEqual(TaskStatus.FAILED, detail.status)
|
|
self.assertEqual("PDD_GOODS_UNAVAILABLE", detail.last_error_code)
|
|
self.assertEqual(1, self.gateway.submission_count)
|
|
|
|
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
|
|
)
|
|
|
|
def test_selected_task_cancelled_before_collector_starts_is_persisted(self):
|
|
self.repository.add_claimed_task(
|
|
NewClaimedTask(
|
|
remote_task_id="COL-CANCEL",
|
|
task_type=TaskType.COLLECT,
|
|
goods_url="https://example.test/COL-CANCEL",
|
|
)
|
|
)
|
|
calls = []
|
|
|
|
class AcceptGateway:
|
|
def submit_failure(self, *_args):
|
|
return SubmissionReceipt(
|
|
True, "FAILURE-001", "2026-08-07T08:00:00Z"
|
|
)
|
|
|
|
service = CollectTaskService(
|
|
AcceptGateway(),
|
|
self.repository,
|
|
self.client,
|
|
"USB-001",
|
|
cancelled=lambda: True,
|
|
collect_service_factory=lambda *_args: FakeCollector(calls),
|
|
)
|
|
|
|
outcome = service.execute_selected("COL-CANCEL")
|
|
|
|
self.assertEqual(outcome.kind, "cancelled")
|
|
self.assertEqual(calls, [])
|
|
self.assertEqual(
|
|
self.repository.get_task("COL-CANCEL").status,
|
|
TaskStatus.CANCELLED,
|
|
)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|