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>
This commit is contained in:
@@ -45,10 +45,17 @@ async def submit_order(
|
||||
payload: SubmitOrderRequest,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[SubmitOrderData]:
|
||||
"""上游提交下单意图
|
||||
"""上游提交下单意图,返回 task_id
|
||||
|
||||
重复提交同一 task_id 不新建任务,返回既有任务且 created=false——上游重发
|
||||
不会变成两单。
|
||||
网关只做任务队列与状态镜像:intent 原文存库,由本地 worker 出站长轮询领走
|
||||
(GET /api/orders/lease),真正的加购、下单、付款都在本地机完成。
|
||||
|
||||
幂等:同一个 task_id 重复提交不新建任务,返回既有任务且 created=false——
|
||||
上游重发不会变成两单。下单不可逆,这是防重复下单的第一道闸。
|
||||
|
||||
提交后任务为 queued。进度跟踪用 GET /api/orders/{task_id}(任务状态 +
|
||||
完整状态历史);任务状态词汇表:queued / leased / running / succeeded /
|
||||
failed / needs_human / stale。
|
||||
"""
|
||||
data = await container.task_queue.submit(
|
||||
task_id=payload.task_id, site=payload.site, intent=payload.intent
|
||||
@@ -64,15 +71,22 @@ async def submit_order(
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def lease_order(
|
||||
worker_id: str = Query(...),
|
||||
wait: int = Query(default=30, ge=0, le=300),
|
||||
site: str | None = Query(default=None),
|
||||
worker_id: str = Query(..., description="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[LeaseData | None]:
|
||||
"""本地长轮询领取
|
||||
"""本地长轮询领取下单任务(worker 专用)
|
||||
|
||||
无可领任务时挂起到 wait 秒后返回 data: null,HTTP 仍 200。已有 leased/running
|
||||
任务时立即返回空(全局并发度 1)。
|
||||
- 每次调用刷新 worker 心跳(workers.last_seen_at),不另建心跳接口。
|
||||
- **全局并发度 1**:已存在 leased / running 任务时立即返回空,绝不发第二个任务。
|
||||
- 有可领任务时原子置 queued → leased,写入 lease_owner 与过期时间
|
||||
(TTL 由 RAKUTEN_LEASE_TTL_SECONDS 控制,默认 300 秒),lease_count += 1。
|
||||
- 无任务时挂起至多 wait 秒再返回,HTTP 仍 200、data 为 null。
|
||||
- 响应 lease_count > 1 只可能来自 reclaim(正常 lease 拿不到 stale 任务)。
|
||||
"""
|
||||
data = await container.task_queue.lease(
|
||||
worker_id=worker_id,
|
||||
@@ -89,13 +103,21 @@ async def lease_order(
|
||||
"/{task_id}/renew",
|
||||
response_model=ApiResponse[RenewData],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
responses={
|
||||
404: {"description": "任务不存在(错误码 6001)"},
|
||||
409: {"description": "租约无效:不是持有者或任务已终结(错误码 6002)"},
|
||||
},
|
||||
)
|
||||
async def renew_order(
|
||||
task_id: str,
|
||||
payload: RenewRequest,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[RenewData]:
|
||||
"""续租。worker 在长任务中每 60 秒调一次,避免租约 TTL 误判"""
|
||||
"""续租(worker 专用)
|
||||
|
||||
执行时间可能超过租约 TTL(下单页面慢、付款要等),worker 在长任务中每 60 秒
|
||||
调一次,避免租约过期被误判成 stale。
|
||||
"""
|
||||
data = await container.task_queue.renew(task_id, payload.worker_id)
|
||||
return ApiResponse[RenewData](success=True, msg="success", data=data, code=0)
|
||||
|
||||
@@ -104,15 +126,27 @@ async def renew_order(
|
||||
"/{task_id}/report",
|
||||
response_model=ApiResponse[ReportData],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
responses={
|
||||
404: {"description": "任务不存在(错误码 6001)"},
|
||||
409: {"description": "租约无效或终态冲突(错误码 6002 / 6003)"},
|
||||
},
|
||||
)
|
||||
async def report_order(
|
||||
task_id: str,
|
||||
payload: ReportRequest,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[ReportData]:
|
||||
"""本地回报订单状态
|
||||
"""本地回报订单状态(worker 专用)
|
||||
|
||||
同一 (task_id, state) 重复上报幂等。terminal=true 时释放租约并推到终态。
|
||||
worker 每一步遵循「动作 → 落证据 → 写本地库 → 回报」,先落证据再回报,
|
||||
保证服务器上看到的状态一定有本地证据可查。
|
||||
|
||||
- worker_id 必须是当前租约持有者,否则 6002。
|
||||
- 首次 report 把任务从 leased 推到 running。
|
||||
- 同一 (task_id, state) 重复上报幂等(覆盖同一行,recorded=false),
|
||||
网络抖动重发安全。
|
||||
- terminal=true 时释放租约并推到终态:terminal_status 显式指定优先,否则按
|
||||
state 推断(paid→succeeded,cancelled→failed,其余→needs_human)。
|
||||
"""
|
||||
data = await container.task_queue.report(
|
||||
task_id,
|
||||
@@ -133,16 +167,27 @@ async def report_order(
|
||||
"/{task_id}/reclaim",
|
||||
response_model=ApiResponse[ReclaimData],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
responses={
|
||||
404: {"description": "任务不存在(错误码 6001)"},
|
||||
409: {"description": "任务不是 stale 状态(错误码 6003)"},
|
||||
},
|
||||
)
|
||||
async def reclaim_order(
|
||||
task_id: str,
|
||||
payload: ReclaimRequest,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[ReclaimData]:
|
||||
"""把 stale 任务重新租给 worker(恢复领取)
|
||||
"""把 stale 任务重新租给 worker(恢复领取,绝不自动重投)
|
||||
|
||||
返回体的 lease_count 必然 > 1,worker 必须先核对站点订单列表再决定是否
|
||||
继续执行——核对不出结论时报 needs_human,绝不重新提交。
|
||||
租约过期后任务置 stale 并告警,**不回 queued、不会被正常 lease 取到**——本地
|
||||
可能已经下单成功、只是回报那一步断网,自动重投等于再买一次。恢复只能显式
|
||||
调本接口。
|
||||
|
||||
返回体的 lease_count 必然 > 1 且带 known_state(租约丢失前的最新订单状态)。
|
||||
worker 执行前**必须先核对站点订单列表**(用 intent 里的商品 + 时间窗口比对):
|
||||
确认已下单则直接补报状态,核对不出结论报 needs_human,绝不重新提交。
|
||||
|
||||
仅 stale 状态可 reclaim;其他状态报 6003。
|
||||
"""
|
||||
data = await container.task_queue.reclaim(task_id, payload.worker_id)
|
||||
return ApiResponse[ReclaimData](success=True, msg="success", data=data, code=0)
|
||||
@@ -152,12 +197,18 @@ async def reclaim_order(
|
||||
"/{task_id}",
|
||||
response_model=ApiResponse[TaskDetail],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
responses={404: {"description": "任务不存在(错误码 6001)"}},
|
||||
)
|
||||
async def get_order(
|
||||
task_id: str,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[TaskDetail]:
|
||||
"""任务详情 + 完整状态历史"""
|
||||
"""任务详情 + 完整状态历史
|
||||
|
||||
返回任务本体(status / 租约信息 / intent 原文)、latest_state(最新一次上报的
|
||||
订单状态)与 reports(append-only 状态历史,含金额、付款期限、站点订单号、
|
||||
证据路径)。任务不存在报 6001。
|
||||
"""
|
||||
data = await container.task_queue.get_task_detail(task_id)
|
||||
return ApiResponse[TaskDetail](success=True, msg="success", data=data, code=0)
|
||||
|
||||
@@ -168,13 +219,16 @@ async def get_order(
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def list_orders(
|
||||
status: str | None = Query(default=None),
|
||||
site: str | None = Query(default=None),
|
||||
limit: int = Query(default=50, ge=1, le=500),
|
||||
offset: int = Query(default=0, ge=0),
|
||||
status: str | None = Query(
|
||||
default=None,
|
||||
description="按任务状态筛选:queued / leased / running / succeeded / failed / needs_human / stale",
|
||||
),
|
||||
site: str | None = Query(default=None, description="按站点筛选(rakuten)"),
|
||||
limit: int = Query(default=50, ge=1, le=500, description="分页大小"),
|
||||
offset: int = Query(default=0, ge=0, description="分页偏移"),
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[TaskListData]:
|
||||
"""任务列表(运维与上游对账用)"""
|
||||
"""任务列表,按创建时间倒序(运维与上游对账用)"""
|
||||
data = await container.task_queue.list_tasks(
|
||||
status=status, site=site, limit=limit, offset=offset
|
||||
)
|
||||
|
||||
@@ -10,7 +10,7 @@
|
||||
- GET /api/account/queries/{id} 取查询单状态与结果
|
||||
- GET /api/account/queries 查询单列表(运维排查)
|
||||
|
||||
**纯异步**:提交只拿 query_id,结果去 GET 取。网关不提供「挂起等结果」的接口——
|
||||
**纯异步**:提交只拿单号,结果去 GET 取。网关不提供「挂起等结果」的接口——
|
||||
下单任务可能占着账号锁跑好几分钟,同步等待会把上游的连接一起卡住。
|
||||
|
||||
注意路由声明顺序:`/lease` 必须在 `/{query_id}` 之前,否则 lease 会被当成 query_id。
|
||||
@@ -43,10 +43,19 @@ async def submit_query(
|
||||
payload: SubmitQueryRequest,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[SubmitQueryData]:
|
||||
"""提交一张账号只读查询单
|
||||
"""提交一张账号只读查询单(纯异步:这里只拿单号)
|
||||
|
||||
重复提交同一 query_id 不新建,返回既有单且 created=false。查询是只读的,
|
||||
重复提交本身无害,幂等键的意义在于上游重发时不会拿到两个单号。
|
||||
上游要的是「已登录账号在站点上的真实订单」,而不是网关的任务状态镜像——两者
|
||||
会不一致(站点侧被商家取消、走别的渠道下的单,网关镜像都看不见)。查询单进
|
||||
队列后由本地 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,
|
||||
@@ -63,15 +72,24 @@ async def submit_query(
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def lease_query(
|
||||
worker_id: str = Query(...),
|
||||
wait: int = Query(default=30, ge=0, le=300),
|
||||
site: str | None = Query(default=None),
|
||||
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 专用)
|
||||
|
||||
无单可领时挂起到 wait 秒后返回 data: null,HTTP 仍 200。与下单 lease 不同,
|
||||
这里没有「已有任务在执行就一律返回空」的闸门。
|
||||
- 与下单 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,
|
||||
@@ -86,16 +104,24 @@ async def lease_query(
|
||||
"/{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 专用)
|
||||
|
||||
租约不是自己的(多半是超时后被重投给了下一轮)时报 6006 并丢弃本次结果,
|
||||
避免迟到的旧结果覆盖新结果。
|
||||
- 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,
|
||||
@@ -112,6 +138,7 @@ async def submit_query_result(
|
||||
"/{query_id}",
|
||||
response_model=ApiResponse[QueryDetail],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
responses={404: {"description": "查询单不存在或已过保留期被清理(错误码 6005)"}},
|
||||
)
|
||||
async def get_query(
|
||||
query_id: str,
|
||||
@@ -119,8 +146,25 @@ async def get_query(
|
||||
) -> ApiResponse[QueryDetail]:
|
||||
"""取查询单状态与结果
|
||||
|
||||
status=succeeded 时 result 里是 worker 从站点读到的原文 + 规范化字段;
|
||||
failed/expired 时看 error。仍在 queued/leased 就是还没读到,稍后再来。
|
||||
仍在 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)
|
||||
@@ -132,13 +176,19 @@ async def get_query(
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def list_queries(
|
||||
status: str | None = Query(default=None),
|
||||
kind: str | None = Query(default=None),
|
||||
limit: int = Query(default=50, ge=1, le=500),
|
||||
offset: int = Query(default=0, ge=0),
|
||||
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
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user