上游要的不只是网关记的任务状态镜像,还有「已登录账号在站点上的真实订单」,
但账号只在 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>
213 lines
7.7 KiB
Python
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"]
|