"""网关 HTTP 客户端:出站长轮询领任务、回报状态、续租、恢复领取 worker 与网关之间是出站单向通信。所有方法都包装成「成功 → data,失败 → 抛 AppError」的形式,runner 拿到 AppError 直接打日志或转报 needs_human。 请求与响应严格使用 shared.api 的 ApiResponse 信封;网关侧错误码(6xxx)原样上抛, 不在这里翻译。鉴权头从 settings.bearer_token 取,与抓取/交易服务共用同一份。 **不 import app.gateway**:worker 只看 HTTP 响应 JSON,本地用 LeaseTask 表达领到的 任务,与网关侧的 Pydantic 模型解耦。 """ from __future__ import annotations import logging from typing import Any import httpx from opentelemetry import trace from opentelemetry.trace import SpanKind from app.shared.config import Settings from app.shared.errors import AppError from app.shared.proxy import httpx_client_options from app.shared.task_state import OrderState, TaskStatus from app.shared.telemetry import ( record_envelope, record_error, set_attributes, span_unless_suppressed, ) from app.trading.worker.models import LeaseTask, QueryTask logger = logging.getLogger(__name__) tracer = trace.get_tracer(__name__) class GatewayClient: """网关 HTTP 客户端 一份 AsyncClient 实例贯穿 worker 整个生命周期,连接池由 httpx 管理。 """ def __init__( self, base_url: str, bearer_token: str, *, settings: Settings, timeout: float = 60.0 ): # 末尾去斜杠,避免 base + "/api/..." 拼出双斜杠 self._base_url = base_url.rstrip("/") self._client = httpx.AsyncClient( base_url=self._base_url, headers={"Authorization": f"Bearer {bearer_token}"}, timeout=timeout, **httpx_client_options(settings, target_url=self._base_url), ) async def aclose(self) -> None: await self._client.aclose() # ---- 基础封装 ---- async def _request(self, method: str, path: str, **kwargs: Any) -> dict[str, Any]: """发起请求并解信封。失败(success=False)抛 AppError 整段套一个自己的 span,而不是依赖 httpx 自动 instrumentation 那个 CLIENT span:**网关的失败在信封里,不在 HTTP 状态码上**。httpx 那个 span 在 `request()` 返回时就结束了,此时信封还没解——一次 `success=false, code=6002` (租约无效)的调用在它看来是完成的 200 请求,链路里跟成功毫无区别。 本 span 活到解信封之后,所以能把「这次调用的结论」记下来。 """ with span_unless_suppressed( tracer, f"gateway.{path.strip('/').replace('/', '.')}", kind=SpanKind.CLIENT, ) as span: set_attributes(span, {"gateway.method": method, "gateway.path": path}) response = await self._client.request(method, path, **kwargs) try: body = response.json() except ValueError as exc: err = AppError( message=f"网关响应不是合法 JSON:HTTP {response.status_code}", code="GATEWAY_BAD_BODY", err_code=3001, retryable=True, ) record_error(span, err) span.set_attribute("gateway.status_code", response.status_code) raise err from exc if not body.get("success"): err_code = int(body.get("code", 1500)) msg = body.get("msg", "网关返回失败") # 记成信封结果而不是只抛异常:上报被网关拒时,链路里能直接按 # api.code 筛出是哪一类拒绝(6002 租约无效 / 6003 状态不允许…), # 不必回头翻 worker 日志。 record_envelope( span, success=False, err_code=err_code, msg=msg, status_code=response.status_code, ) raise AppError( message=msg, code="GATEWAY_ERROR", err_code=err_code, retryable=False, status_code=response.status_code, ) record_envelope( span, success=True, status_code=response.status_code ) return body # ---- 接口 ---- async def lease( self, worker_id: str, *, wait: int = 30, site: str | None = None ) -> LeaseTask | None: """长轮询领取。无可领任务时返回 None""" params: dict[str, Any] = {"worker_id": worker_id, "wait": wait} if site: params["site"] = site body = await self._request("GET", "/api/orders/lease", params=params) data = body.get("data") if not data: return None return LeaseTask( task_id=data["task_id"], site=data["site"], intent=data.get("intent") or {}, lease_expires_at=data.get("lease_expires_at", ""), lease_count=data.get("lease_count", 1), known_state=data.get("known_state"), ) async def renew(self, task_id: str, worker_id: str) -> dict[str, Any]: """续租,返回网关响应里的 data 字段""" body = await self._request( "POST", f"/api/orders/{task_id}/renew", json={"worker_id": worker_id}, ) return body["data"] async def report( self, task_id: str, worker_id: str, *, state: OrderState, payable_yen: int | None = None, pay_deadline: str | None = None, site_order_id: str | None = None, evidence_ref: str | None = None, detail: str = "", terminal: bool = False, terminal_status: TaskStatus | None = None, ) -> dict[str, Any]: """回报状态,返回网关响应里的 data 字段""" payload: dict[str, Any] = { "worker_id": worker_id, "state": state.value, "payable_yen": payable_yen, "pay_deadline": pay_deadline, "site_order_id": site_order_id, "evidence_ref": evidence_ref, "detail": detail, "terminal": terminal, } if terminal_status is not None: payload["terminal_status"] = terminal_status.value body = await self._request("POST", f"/api/orders/{task_id}/report", json=payload) return body["data"] async def get_task(self, task_id: str) -> dict[str, Any]: """任务详情,返回网关响应里的 data 字段(含 created_at) 供 verify.verify_on_site 恢复核对用:LeaseTask(lease/reclaim 的响应) 不带 created_at,规格 §5「按创建时间 ~ stale 之间」的窗口核对必须单独 查一次这个接口才能拿到。 """ body = await self._request("GET", f"/api/orders/{task_id}") return body["data"] async def reclaim(self, task_id: str, worker_id: str) -> LeaseTask: """恢复领取。返回的 LeaseTask 必然 lease_count > 1""" body = await self._request( "POST", f"/api/orders/{task_id}/reclaim", json={"worker_id": worker_id}, ) data = body["data"] return LeaseTask( task_id=data["task_id"], site=data["site"], intent=data.get("intent") or {}, lease_expires_at=data.get("lease_expires_at", ""), lease_count=data.get("lease_count", 1), known_state=data.get("known_state"), ) # ---- 账号只读查询通道(docs/order-gateway.md §11)---- async def lease_query(self, worker_id: str, *, wait: int = 30) -> QueryTask | None: """长轮询领取一张只读查询单。无单可领时返回 None 与 lease() 走的是两条互不干扰的通道:下单任务在跑的时候,查询照样能领到 (账号级串行由 SiteInteractor 那把锁保证,不靠网关排队)。 """ body = await self._request( "GET", "/api/account/queries/lease", params={"worker_id": worker_id, "wait": wait}, ) data = body.get("data") if not data: return None return QueryTask( query_id=data["query_id"], site=data["site"], kind=data["kind"], params=data.get("params") or {}, lease_expires_at=data.get("lease_expires_at", ""), attempt=data.get("attempt", 1), ) async def report_query_result( self, query_id: str, worker_id: str, *, success: bool, result: dict[str, Any] | None = None, error_code: int | None = None, error_message: str = "", ) -> dict[str, Any]: """回报查询结果,返回网关响应里的 data 字段 网关可能回 6006(租约已被重投给下一轮)——那说明本次结果已经作废, 调用方按「丢弃」处理即可,不要重发。 """ payload: dict[str, Any] = { "worker_id": worker_id, "success": success, "result": result, "error_code": error_code, "error_message": error_message, } body = await self._request( "POST", f"/api/account/queries/{query_id}/result", json=payload ) return body["data"]