Files
rakuten-api/app/trading/worker/client.py
T
q792602257andClaude Opus 5 c03158488b feat(gateway): 账号只读查询通道——从已登录账号取真实订单
上游要的不只是网关记的任务状态镜像,还有「已登录账号在站点上的真实订单」,
但账号只在 NAT 后本地机上,只能经网关队列走。新增独立查询通道(§11):

- gateway 单开 account_queries 表 + QueryStatus 状态机,接口
  POST /api/account/queries(幂等)/ lease / {id}/result / {id}
- 不复用下单任务队列:查询是只读,租约过期可安全重投(与下单「绝不自动
  重投」相反),且不该被全局并发度 1 堵死、task_reports 是订单镜像不能污染
- 本地交易服务起第二条常驻循环 query_runner,领到即调 SiteInteractor 真读:
  order_list 复用已实测的 list_recent_orders(规范化字段 + 站点
  orderListData 原文),order_detail 复用 fetch_order_detail(配送阶段 +
  页面 __INITIAL_STATE__ 原样透传,结构未经真实样本,不抽字段)
- 账号级串行仍由 SiteInteractor 的锁保证;每次执行套超时按失败回报
- 错误码 6005/6006(查询通道,可重试只读区别于 6001-6004);/health 暴露
  queued_query_count;结果体积上限先丢原始 JSON

openapi.json 重导,docs/order-gateway.md §11、README、.env.example 补全
配置与实测边界。全量测试 404→454 通过。

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

213 lines
7.7 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 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.trading.worker.models import LeaseTask, QueryTask
logger = logging.getLogger(__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"""
response = await self._client.request(method, path, **kwargs)
try:
body = response.json()
except ValueError as exc:
raise AppError(
message=f"网关响应不是合法 JSON:HTTP {response.status_code}",
code="GATEWAY_BAD_BODY",
err_code=3001,
retryable=True,
) from exc
if not body.get("success"):
raise AppError(
message=body.get("msg", "网关返回失败"),
code="GATEWAY_ERROR",
err_code=int(body.get("code", 1500)),
retryable=False,
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"]