"""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_resolve_purchase_spec_posts_once_and_parses_whitelisted_match(self): response = { "schema_version": 1, "resolution_id": "psr-001", "outcome": "matched", "source": "rule", "candidate_snapshot_hash": "a" * 64, "match": { "candidate_id": "c1", "raw_text": "120斤", "options": {"color": "黑色", "size": "120斤"}, }, "confidence_bps": 10000, "reason": "唯一等价", "resolved_at": "2026-08-17T08:00:01Z", "future_field": "ignored", } opener = RecordingOpener(FakeResponse(200, response)) gateway = HttpAdminGateway( opener=opener, client_id="CLIENT-001", timeout_seconds=5 ) receipt = gateway.resolve_purchase_spec( "PUR-001", "spec-resolution-v1:key", {"schema_version": 1} ) self.assertEqual(receipt.resolution_id, "psr-001") self.assertEqual(receipt.match.raw_text, "120斤") self.assertTrue(opener.request.full_url.endswith( "/api/v1/client/tasks/PUR-001/spec-resolution" )) headers = { key.lower(): value for key, value in opener.request.header_items() } self.assertEqual( headers["idempotency-key"], "spec-resolution-v1:key" ) self.assertEqual(opener.timeout, 5) def test_resolve_purchase_spec_accepts_all_nonmatched_business_outcomes(self): for outcome in ("uncertain", "rejected", "failed"): with self.subTest(outcome=outcome): gateway = HttpAdminGateway( opener=RecordingOpener( FakeResponse( 200, { "schema_version": 1, "resolution_id": f"psr-{outcome}", "outcome": outcome, "source": None, "candidate_snapshot_hash": "b" * 64, "match": None, "confidence_bps": None, "reason": "没有唯一安全候选", "resolved_at": "2026-08-17T08:00:01Z", }, ) ), client_id="CLIENT-001", ) receipt = gateway.resolve_purchase_spec( "PUR-001", "spec-resolution-v1:key", {} ) self.assertEqual(receipt.outcome, outcome) self.assertIsNone(receipt.match) def test_resolve_purchase_spec_rejects_match_on_nonmatched_outcome(self): gateway = HttpAdminGateway( opener=RecordingOpener( FakeResponse( 200, { "schema_version": 1, "resolution_id": "psr-invalid", "outcome": "uncertain", "source": "ai", "candidate_snapshot_hash": "c" * 64, "match": { "candidate_id": "c1", "raw_text": "120斤", "options": {"color": "黑色", "size": "120斤"}, }, "confidence_bps": 5000, "reason": "响应自相矛盾", "resolved_at": "2026-08-17T08:00:01Z", }, ) ), client_id="CLIENT-001", ) with self.assertRaises(AdminGatewayError) as raised: gateway.resolve_purchase_spec( "PUR-001", "spec-resolution-v1:key", {} ) self.assertEqual(raised.exception.code, "ADMIN_INVALID_RESPONSE") 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()