Files
cmautobuy/client/test/test_http_admin_gateway.py
T

444 lines
15 KiB
Python

"""Admin 登记 HTTP 契约测试,不访问真实网络。"""
import io
import json
import socket
import unittest
from http.client import RemoteDisconnected
from urllib.error import HTTPError, URLError
from src.admin_gateway import (
AdminGatewayError,
AndroidDeviceInfo,
ClaimCapabilities,
ClientInfo,
)
from src.http_admin_gateway import DEFAULT_ADMIN_BASE_URL, HttpAdminGateway
from src.task_models import TaskType
class FakeResponse:
def __init__(self, status: int, payload):
self.status = status
self._body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
def getcode(self):
return self.status
def read(self):
return self._body
def __enter__(self):
return self
def __exit__(self, *_args):
return False
class RecordingOpener:
def __init__(self, response):
self.response = response
self.request = None
self.timeout = None
def __call__(self, request, timeout):
self.request = request
self.timeout = timeout
if isinstance(self.response, Exception):
raise self.response
return self.response
class HttpAdminGatewayTest(unittest.TestCase):
def _capabilities(self, with_device=True):
device = (
AndroidDeviceInfo("192.168.0.173:5555") if with_device else None
)
return ClaimCapabilities(device=device)
def test_register_sends_confirmed_contract_and_parses_response(self):
opener = RecordingOpener(
FakeResponse(
200,
{
"registered": True,
"client_id": "CLIENT-001",
"registered_at": "2026-08-06T10:00:00Z",
},
)
)
gateway = HttpAdminGateway(
"http://127.0.0.1:8080/", "secret-token", 2.5, opener
)
receipt = gateway.register_client(
ClientInfo("CLIENT-001", "办公室电脑"), self._capabilities()
)
self.assertTrue(receipt.registered)
self.assertEqual(opener.request.get_method(), "PUT")
self.assertEqual(
opener.request.full_url,
"http://127.0.0.1:8080/api/v1/client/registration",
)
headers = {
key.lower(): value for key, value in opener.request.header_items()
}
self.assertEqual(headers["x-client-id"], "CLIENT-001")
self.assertTrue(headers["x-request-id"])
self.assertEqual(headers["authorization"], "Bearer secret-token")
body = json.loads(opener.request.data.decode("utf-8"))
self.assertEqual(body["client"]["name"], "办公室电脑")
self.assertEqual(body["supported_types"], ["collect"])
self.assertEqual(body["device"]["platform"], "android")
self.assertEqual(body["capabilities"]["purchase_mode"], "dry_run")
self.assertEqual(opener.timeout, 2.5)
def test_optional_device_is_omitted(self):
opener = RecordingOpener(
FakeResponse(
200,
{
"registered": True,
"client_id": "CLIENT-001",
"registered_at": "2026-08-06T10:00:00Z",
},
)
)
gateway = HttpAdminGateway(opener=opener, client_id="CLIENT-001")
gateway.register_client(
ClientInfo("CLIENT-001"), self._capabilities(False)
)
body = json.loads(opener.request.data.decode("utf-8"))
self.assertNotIn("device", body)
self.assertNotIn("authorization", {
key.lower(): value for key, value in opener.request.header_items()
})
def test_submit_result_sends_idempotency_key_and_parses_receipt(self):
opener = RecordingOpener(
FakeResponse(
200,
{
"accepted": True,
"result_id": "RESULT-001",
"accepted_at": "2026-08-07T08:00:01Z",
},
)
)
gateway = HttpAdminGateway(opener=opener, client_id="CLIENT-001")
payload = {
"task_version": 1,
"attempt_id": "ATTEMPT-001",
"result_type": "collect",
"completed_at": "2026-08-07T08:00:00Z",
"pdd_data": {"goods_id": "123"},
}
receipt = gateway.submit_result(
"COL-001", "COL-001:ATTEMPT-001:result-v1", payload
)
self.assertTrue(receipt.accepted)
self.assertEqual(
opener.request.full_url,
f"{DEFAULT_ADMIN_BASE_URL}/api/v1/client/tasks/COL-001/result",
)
headers = {
key.lower(): value for key, value in opener.request.header_items()
}
self.assertEqual(
headers["idempotency-key"], "COL-001:ATTEMPT-001:result-v1"
)
self.assertEqual(headers["x-client-id"], "CLIENT-001")
self.assertEqual(json.loads(opener.request.data), payload)
def test_submit_failure_uses_failure_endpoint(self):
opener = RecordingOpener(
FakeResponse(
201,
{
"accepted": True,
"result_id": "FAILURE-001",
"accepted_at": "2026-08-07T08:00:01Z",
},
)
)
gateway = HttpAdminGateway(opener=opener, client_id="CLIENT-001")
gateway.submit_failure(
"COL-001",
"COL-001:ATTEMPT-001:failure-v1",
{
"task_version": 1,
"attempt_id": "ATTEMPT-001",
"status": "manual_review",
"error": {"code": "PDD_PAGE_CAPTCHA"},
"reported_at": "2026-08-07T08:00:00Z",
},
)
self.assertTrue(opener.request.full_url.endswith("/COL-001/failure"))
def test_submit_failure_accepts_legacy_admin_receipt_without_result_id(self):
gateway = HttpAdminGateway(
opener=RecordingOpener(
FakeResponse(
200,
{
"accepted": True,
"task_status": "assigned",
"accepted_at": "2026-08-07T08:00:01Z",
},
)
),
client_id="CLIENT-001",
)
receipt = gateway.submit_failure(
"COL-001",
"COL-001:ATTEMPT-001:failure-v1",
{"status": "retry_wait"},
)
self.assertTrue(receipt.accepted)
self.assertTrue(receipt.result_id.startswith("legacy-failure-"))
def test_malformed_2xx_submission_is_retryable_and_not_called_rejected(self):
gateway = HttpAdminGateway(
opener=RecordingOpener(FakeResponse(200, {"accepted": True})),
client_id="CLIENT-001",
)
with self.assertRaises(AdminGatewayError) as raised:
gateway.submit_result("COL-001", "key-1", {})
self.assertEqual(raised.exception.code, "ADMIN_INVALID_RESPONSE")
self.assertTrue(raised.exception.retryable)
self.assertIn("可能已接收", str(raised.exception))
def test_claim_maps_payload_and_serializes_confirmed_capabilities(self):
opener = RecordingOpener(
FakeResponse(
200,
{
"task": {
"id": "cj1",
"type": "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",
}
},
)
)
gateway = HttpAdminGateway(opener=opener)
task = gateway.claim_next(
ClientInfo("CLIENT-001", "办公室电脑"),
self._capabilities(),
)
self.assertIsNotNone(task)
self.assertEqual(task.task_id, "cj1")
self.assertEqual(task.task_type.value, "collect")
self.assertEqual(task.payload["goods_id"], "737116531267")
self.assertEqual(opener.request.get_method(), "POST")
self.assertEqual(
opener.request.full_url,
f"{DEFAULT_ADMIN_BASE_URL}/api/v1/client/tasks/claim",
)
headers = {
key.lower(): value for key, value in opener.request.header_items()
}
self.assertEqual(headers["x-client-id"], "CLIENT-001")
self.assertTrue(headers["x-request-id"])
self.assertIn("application/json", headers["content-type"])
body = json.loads(opener.request.data.decode("utf-8"))
self.assertEqual(body["supported_types"], ["collect"])
self.assertEqual(body["capabilities"]["purchase_mode"], "dry_run")
def test_claim_purchase_requires_all_safety_fields(self):
valid_task = {
"id": "PUR-001",
"type": "purchase",
"version": 1,
"priority": 0,
"execution_mode": "live",
"payload": {
"goods_url": "https://example.test/goods/PUR-001",
"goods_id": "737116531267",
"options": {"color": "黑色", "size": "L"},
"quantity": 2,
"max_price_cent": 4200,
},
"created_at": "2026-08-09T08:00:00Z",
"updated_at": "2026-08-09T08:00:00Z",
}
gateway = HttpAdminGateway(
opener=RecordingOpener(FakeResponse(200, {"task": valid_task}))
)
task = gateway.claim_next(
ClientInfo("CLIENT-001"),
ClaimCapabilities(
supported_types=(TaskType.COLLECT, TaskType.PURCHASE),
purchase_mode="live",
),
)
self.assertEqual(task.task_type, TaskType.PURCHASE)
self.assertEqual(task.execution_mode, "live")
self.assertEqual(task.payload["max_price_cent"], 4200)
for missing in ("goods_id", "options", "quantity", "max_price_cent"):
with self.subTest(missing=missing):
invalid_task = dict(valid_task)
invalid_payload = dict(valid_task["payload"])
invalid_payload.pop(missing)
invalid_task["payload"] = invalid_payload
invalid_gateway = HttpAdminGateway(
opener=RecordingOpener(
FakeResponse(200, {"task": invalid_task})
)
)
with self.assertRaises(AdminGatewayError) as raised:
invalid_gateway.claim_next(
ClientInfo("CLIENT-001"), self._capabilities()
)
self.assertEqual(
raised.exception.code, "ADMIN_INVALID_RESPONSE"
)
invalid_mode_task = dict(valid_task)
invalid_mode_task["execution_mode"] = "unsafe"
invalid_gateway = HttpAdminGateway(
opener=RecordingOpener(
FakeResponse(200, {"task": invalid_mode_task})
)
)
with self.assertRaises(AdminGatewayError) as raised:
invalid_gateway.claim_next(
ClientInfo("CLIENT-001"),
ClaimCapabilities(
supported_types=(TaskType.COLLECT, TaskType.PURCHASE),
purchase_mode="live",
),
)
self.assertEqual(raised.exception.code, "ADMIN_INVALID_RESPONSE")
def test_claim_204_returns_none(self):
gateway = HttpAdminGateway(opener=RecordingOpener(FakeResponse(204, {})))
task = gateway.claim_next(
ClientInfo("CLIENT-001"),
self._capabilities(),
)
self.assertIsNone(task)
def test_claim_invalid_task_fields_include_request_id(self):
gateway = HttpAdminGateway(
opener=RecordingOpener(FakeResponse(200, {"task": {"id": None}}))
)
with self.assertRaises(AdminGatewayError) as context:
gateway.claim_next(
ClientInfo("CLIENT-001"),
self._capabilities(),
)
self.assertEqual(context.exception.code, "ADMIN_INVALID_RESPONSE")
self.assertTrue(context.exception.request_id)
def test_claim_http_error_keeps_server_request_id(self):
error_body = json.dumps(
{
"error": {
"code": "TASK_CLAIM_FAILED",
"message": "领取任务失败",
"retryable": True,
"request_id": "claim-request-id",
}
}
).encode("utf-8")
error = HTTPError(
"http://admin/api/v1/client/tasks/claim",
500,
"Internal Server Error",
{},
io.BytesIO(error_body),
)
gateway = HttpAdminGateway(opener=RecordingOpener(error))
with self.assertRaises(AdminGatewayError) as context:
gateway.claim_next(
ClientInfo("CLIENT-001"),
self._capabilities(),
)
self.assertEqual(context.exception.code, "TASK_CLAIM_FAILED")
self.assertEqual(context.exception.request_id, "claim-request-id")
self.assertTrue(context.exception.retryable)
def test_admin_error_preserves_code_retry_and_request_id(self):
error_body = json.dumps(
{
"error": {
"code": "INVALID_CLIENT_PROFILE",
"message": "资料无效",
"retryable": False,
"request_id": "server-request-id",
"details": {"field": "supported_types"},
}
}
).encode("utf-8")
error = HTTPError(
"http://admin/api/v1/client/registration",
422,
"Unprocessable Entity",
{},
io.BytesIO(error_body),
)
gateway = HttpAdminGateway(opener=RecordingOpener(error))
with self.assertRaises(AdminGatewayError) as context:
gateway.register_client(
ClientInfo("CLIENT-001"), self._capabilities(False)
)
self.assertEqual(context.exception.code, "INVALID_CLIENT_PROFILE")
self.assertFalse(context.exception.retryable)
self.assertEqual(context.exception.request_id, "server-request-id")
def test_timeout_and_connection_failure_are_retryable(self):
for exception, code in (
(URLError(socket.timeout()), "ADMIN_TIMEOUT"),
(URLError("connection refused"), "ADMIN_UNAVAILABLE"),
(RemoteDisconnected("closed"), "ADMIN_UNAVAILABLE"),
):
with self.subTest(code=code):
gateway = HttpAdminGateway(
opener=RecordingOpener(exception)
)
with self.assertRaises(AdminGatewayError) as context:
gateway.register_client(
ClientInfo("CLIENT-001"), self._capabilities(False)
)
self.assertEqual(context.exception.code, code)
self.assertTrue(context.exception.retryable)
if __name__ == "__main__":
unittest.main()