Files
q792602257andClaude Opus 5 8381896eeb feat(observability): 失败响应与网关信封进链路,手工埋点尊重 suppressed
失败此前在 trace 里近乎不可见:异常处理器把异常吃掉换成信封响应,自动
instrumentation 只看到一个 HTTP 状态码,而 AppError 默认 400、信封里
success=false,跟正常返回分不出来。

- api.py:四个异常处理器(对外失败的唯一出口)各记一次 span;兜底处理器额外
  record_error——对外只回一句无信息量的错误文案,异常类型与栈只在本地日志里
- telemetry.py:新增 record_envelope / record_parse_failure /
  span_unless_suppressed;record_error 补 error.message / retryable /
  status_code;snapshot 支持 extra 带上「这份 HTML 是哪来的」
- worker/client.py:_request 自建 span,活到解信封之后。httpx 那个 CLIENT span
  在 request() 返回时就结束,此时信封还没解——success=false code=6002(租约
  无效)在它看来是完成的 200 请求
- span_unless_suppressed:suppress_instrumentation 只被 instrumentation 库尊重,
  手工 span 不看它,lease 空转长轮询会从这个口子把孤立 trace 放回来
- 两个 scraping client 的解析失败分支收拢到 record_parse_failure;shop_items
  显式标 stage=delegate 且不落快照(它自己不抓页面,按 html 推断只会得出
  「fetch 失败」的错误结论)

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-28 16:07:00 +08:00

254 lines
9.5 KiB
Python

"""网关 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"]