Files
q792602257andClaude Opus 5 d54af141eb docs(gateway): 补齐任务/查询接口的 docstring 与请求响应字段描述
orders.py / queries.py 各接口补响应状态码(6001/6002/6003/6005/6006)、查询
单两种 kind(order_list / order_detail)结果结构、窗口未覆盖/体积上限等边界
说明;models.py 给 Pydantic 请求/响应模型补 Field 描述并与 Query 参数文档对齐。
纯文档与元数据改动,不含逻辑变更;gateway 相关测试全绿。

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

196 lines
9.1 KiB
Python

"""账号只读查询通道路由(docs/order-gateway.md §11)
上游要的是「已登录账号在站点上的真实订单」,而不是网关自己记录的任务状态镜像;
但本地机在 NAT 后没有公网入口,网关推不进去。所以这条通道与下单任务同构:
上游把查询意图放进队列,本地 worker 出站长轮询领走、去站点上真读一次、回结果。
- POST /api/account/queries 上游提交查询单(幂等)
- GET /api/account/queries/lease 本地长轮询领取
- POST /api/account/queries/{id}/result 本地回结果
- GET /api/account/queries/{id} 取查询单状态与结果
- GET /api/account/queries 查询单列表(运维排查)
**纯异步**:提交只拿单号,结果去 GET 取。网关不提供「挂起等结果」的接口——
下单任务可能占着账号锁跑好几分钟,同步等待会把上游的连接一起卡住。
注意路由声明顺序:`/lease` 必须在 `/{query_id}` 之前,否则 lease 会被当成 query_id。
"""
from __future__ import annotations
from fastapi import APIRouter, Depends, Query
from app.gateway.container import GatewayContainer
from app.gateway.models import (
QueryDetail,
QueryLeaseData,
QueryListData,
QueryResultData,
QueryResultRequest,
SubmitQueryData,
SubmitQueryRequest,
)
from app.shared.api import ApiResponse, get_container, require_bearer_token
router = APIRouter(prefix="/api/account/queries", tags=["account-queries"])
@router.post(
"",
response_model=ApiResponse[SubmitQueryData],
dependencies=[Depends(require_bearer_token)],
)
async def submit_query(
payload: SubmitQueryRequest,
container: GatewayContainer = Depends(get_container),
) -> ApiResponse[SubmitQueryData]:
"""提交一张账号只读查询单(纯异步:这里只拿单号)
上游要的是「已登录账号在站点上的真实订单」,而不是网关的任务状态镜像——两者
会不一致(站点侧被商家取消、走别的渠道下的单,网关镜像都看不见)。查询单进
队列后由本地 worker 出站领走、去站点上真读一次;结果去
GET /api/account/queries/{query_id} 取。
幂等:同一 query_id 重复提交返回既有单且 created=false。查询只读,重发无
副作用,幂等键的意义是上游重发时不会拿到两个单号。
kind=order_list:按时间窗口翻页扫描订单列表(params 支持 since / max_pages)。
kind=order_detail:读单笔注文番号详情(params 必填 order_number)。
kind 非法在提交时就 422,不会等 worker 领走才失败。
"""
data = await container.query_queue.submit(
query_id=payload.query_id,
site=payload.site,
kind=payload.kind.value,
params=payload.params,
)
return ApiResponse[SubmitQueryData](success=True, msg="success", data=data, code=0)
@router.get(
"/lease",
response_model=ApiResponse[QueryLeaseData | None],
dependencies=[Depends(require_bearer_token)],
)
async def lease_query(
worker_id: str = Query(..., description="查询 worker 标识(不刷下单 worker 心跳)"),
wait: int = Query(
default=30, ge=0, le=300,
description="无单可领时的挂起秒数。服务端上限 RAKUTEN_LEASE_MAX_WAIT_SECONDS(默认 60)",
),
site: str | None = Query(default=None, description="限定站点(rakuten);不传不限"),
container: GatewayContainer = Depends(get_container),
) -> ApiResponse[QueryLeaseData | None]:
"""本地长轮询领取查询单(查询 worker 专用)
- 与下单 lease 不同:**没有**「已有任务在执行就一律返回空」的闸门,多张查询单
可以同时在飞;账号级串行由本地 worker 的账号锁保证,下单跑着时查询排队等锁。
- 不刷下单 worker 心跳——心跳的语义是「下单 worker 还活着」,查询循环活着不
代表它活着。
- 无单可领时挂起至多 wait 秒再返回,HTTP 仍 200、data 为 null。
- 租约过期自动重投(查询只读,重投安全,与下单「绝不自动重投」相反):
attempt 递增,用尽(RAKUTEN_QUERY_MAX_ATTEMPTS,默认 3)置 failed;
超过查询单整体 TTL(RAKUTEN_QUERY_TTL_SECONDS,默认 900 秒)置 expired。
"""
data = await container.query_queue.lease(
worker_id=worker_id,
wait=wait,
site=site,
max_wait=container.settings.lease_max_wait_seconds,
)
return ApiResponse[QueryLeaseData | None](success=True, msg="success", data=data, code=0)
@router.post(
"/{query_id}/result",
response_model=ApiResponse[QueryResultData],
dependencies=[Depends(require_bearer_token)],
responses={
404: {"description": "查询单不存在(错误码 6005)"},
409: {"description": "查询单租约无效:不是持有者、已被重投或已终结(错误码 6006)"},
},
)
async def submit_query_result(
query_id: str,
payload: QueryResultRequest,
container: GatewayContainer = Depends(get_container),
) -> ApiResponse[QueryResultData]:
"""本地回报查询结果(查询 worker 专用)
- success=true 时 result 必填;false 时 error_message 必填,error_code 沿用
统一错误码表(如 5001 掉登录),原样带给上游按同一张表分支。
- 租约不是自己的(多半是执行超时、单子已被重投给下一轮)时报 6006——
**丢弃本次结果即可,不要重发**,否则迟到的旧结果会覆盖新一轮的结果。
- result 体积上限 RAKUTEN_QUERY_RESULT_MAX_BYTES(默认 1MB):超限由 worker 侧
先丢 raw / raw_pages(置 raw_omitted)再判失败,让上游缩小窗口重来。
"""
data = await container.query_queue.submit_result(
query_id,
worker_id=payload.worker_id,
success=payload.success,
result=payload.result,
error_code=payload.error_code,
error_message=payload.error_message,
)
return ApiResponse[QueryResultData](success=True, msg="success", data=data, code=0)
@router.get(
"/{query_id}",
response_model=ApiResponse[QueryDetail],
dependencies=[Depends(require_bearer_token)],
responses={404: {"description": "查询单不存在或已过保留期被清理(错误码 6005)"}},
)
async def get_query(
query_id: str,
container: GatewayContainer = Depends(get_container),
) -> ApiResponse[QueryDetail]:
"""取查询单状态与结果
仍在 queued / leased 就是还没读到,稍后再来;failed / expired 看 error;
status=succeeded 时 result 为「规范化字段 + 站点原始 JSON」,结构按 kind:
**kind=order_list**:`{kind, since, max_pages, window_fully_covered, orders[], raw_pages[]}`
— orders[] 含 order_number / order_date / shop_id / shop_name / items[](稳定契约字段);
raw_pages 是站点 __INITIAL_STATE__.orderListData 原文(按翻页顺序,可能因体积上限
被省略并置 raw_omitted=true)。
**window_fully_covered=false 时,「列表里没有某笔订单」不等于「账号里没有」**——
翻页在覆盖完窗口前就停了(命中 max_pages、页面改版、掉登录),请缩小窗口重查或
交人工,不要当成确定性的否定结论。
**kind=order_detail**:`{kind, order_number, found, stage_label, order_state,
delivery_status, raw, raw_available}`
— found=false 不是错误:站点自己说下单后要约 10 分钟才反映;stage_label 是站点
原文(如「ご注文確認中」),order_state 是映射后的订单状态(映射不到为 null),
delivery_status 是站点结构化配送状态码(原样带给上游对账);金额、收货地址、
付款方式等在 raw.orderData 里原样透传,本服务不抽取。
终态查询单保留 RAKUTEN_QUERY_RETENTION_SECONDS(默认 7 天),过期被清理后查也报 6005。
"""
data = await container.query_queue.get_detail(query_id)
return ApiResponse[QueryDetail](success=True, msg="success", data=data, code=0)
@router.get(
"",
response_model=ApiResponse[QueryListData],
dependencies=[Depends(require_bearer_token)],
)
async def list_queries(
status: str | None = Query(
default=None,
description="按状态筛选:queued / leased / succeeded / failed / expired",
),
kind: str | None = Query(default=None, description="按类型筛选:order_list / order_detail"),
limit: int = Query(default=50, ge=1, le=500, description="分页大小"),
offset: int = Query(default=0, ge=0, description="分页偏移"),
container: GatewayContainer = Depends(get_container),
) -> ApiResponse[QueryListData]:
"""查询单列表,按创建时间倒序(运维排查用)
注意定时下派通道(collector)自主派发的 discover-* 查询单也会出现在这里。
"""
data = await container.query_queue.list_queries(
status=status, kind=kind, limit=limit, offset=offset
)
return ApiResponse[QueryListData](success=True, msg="success", data=data, code=0)