444 lines
15 KiB
Python
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": "COL-8020a8729f111c15",
|
|
"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, "COL-8020a8729f111c15")
|
|
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()
|