- POST /api/orders 新增可选 callback_url(仅 http/https,其余 422) - 仅终结类事件各通知一次:terminal report 推入终态(succeeded/failed/ needs_human)、租约过期被 sweep 置 stale;中间态与终结后的监控上报不通知 - 幂等重发不更新既有任务的回调地址;投递 best-effort 单次尝试,失败只记日志 - CallbackNotifier 发后不管(create_task + 在途任务强引用),关停等在途发完; 生产路径按回调地址逐次构造客户端,满足统一出站代理策略(test_proxy.py) - tasks 表加 callback_url 列,GatewayDB.start() 内置迁移兼容既有库 - 新增 RAKUTEN_CALLBACK_TIMEOUT_SECONDS(默认 10);文档补 §4.8; openapi.json 重新导出(gitignore 未跟踪);新增 17 条测试,全量 492 通过
243 lines
9.6 KiB
Python
243 lines
9.6 KiB
Python
"""下单任务网关路由
|
|
|
|
承担规格 §4 的全部接口:
|
|
|
|
- POST /api/orders 上游提交下单意图(幂等)
|
|
- GET /api/orders/lease 本地长轮询领取
|
|
- POST /api/orders/{id}/renew 续租
|
|
- POST /api/orders/{id}/report 本地回报状态
|
|
- POST /api/orders/{id}/reclaim 把 stale 任务重新租给 worker(绝不自动重投)
|
|
- GET /api/orders/{id} 任务详情 + 完整状态历史
|
|
- GET /api/orders 任务列表(运维与上游对账)
|
|
|
|
handler 风格与 app/scraping/api/routes/scrape.py 保持一致:不在 handler 里写
|
|
try/except,业务异常通过 AppError 自动套进 register_exception_handlers。
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
from fastapi import APIRouter, Depends, Query
|
|
|
|
from app.gateway.container import GatewayContainer
|
|
from app.gateway.models import (
|
|
LeaseData,
|
|
RenewData,
|
|
ReclaimData,
|
|
ReportData,
|
|
SubmitOrderData,
|
|
SubmitOrderRequest,
|
|
TaskDetail,
|
|
TaskListData,
|
|
RenewRequest,
|
|
ReclaimRequest,
|
|
ReportRequest,
|
|
)
|
|
from app.shared.api import ApiResponse, get_container, require_bearer_token
|
|
|
|
router = APIRouter(prefix="/api/orders", tags=["orders"])
|
|
|
|
|
|
@router.post(
|
|
"",
|
|
response_model=ApiResponse[SubmitOrderData],
|
|
dependencies=[Depends(require_bearer_token)],
|
|
)
|
|
async def submit_order(
|
|
payload: SubmitOrderRequest,
|
|
container: GatewayContainer = Depends(get_container),
|
|
) -> ApiResponse[SubmitOrderData]:
|
|
"""上游提交下单意图,返回 task_id
|
|
|
|
网关只做任务队列与状态镜像:intent 原文存库,由本地 worker 出站长轮询领走
|
|
(GET /api/orders/lease),真正的加购、下单、付款都在本地机完成。
|
|
|
|
幂等:同一个 task_id 重复提交不新建任务,返回既有任务且 created=false——
|
|
上游重发不会变成两单。下单不可逆,这是防重复下单的第一道闸。
|
|
|
|
可带 callback_url:任务到达终态(succeeded / failed / needs_human)或被置
|
|
stale 时,网关向该地址 POST 一条 JSON 通知(best-effort 单次投递,§4.8)。
|
|
幂等重发不会更新既有任务的回调地址。
|
|
|
|
提交后任务为 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,
|
|
callback_url=payload.callback_url,
|
|
)
|
|
return ApiResponse[SubmitOrderData](
|
|
success=True, msg="success", data=data, code=0
|
|
)
|
|
|
|
|
|
@router.get(
|
|
"/lease",
|
|
response_model=ApiResponse[LeaseData | None],
|
|
dependencies=[Depends(require_bearer_token)],
|
|
)
|
|
async def lease_order(
|
|
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 专用)
|
|
|
|
- 每次调用刷新 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,
|
|
wait=wait,
|
|
site=site,
|
|
max_wait=container.settings.lease_max_wait_seconds,
|
|
)
|
|
return ApiResponse[LeaseData | None](
|
|
success=True, msg="success", data=data, code=0
|
|
)
|
|
|
|
|
|
@router.post(
|
|
"/{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 专用)
|
|
|
|
执行时间可能超过租约 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)
|
|
|
|
|
|
@router.post(
|
|
"/{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 专用)
|
|
|
|
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,
|
|
worker_id=payload.worker_id,
|
|
state=payload.state,
|
|
payable_yen=payload.payable_yen,
|
|
pay_deadline=payload.pay_deadline,
|
|
site_order_id=payload.site_order_id,
|
|
evidence_ref=payload.evidence_ref,
|
|
detail=payload.detail,
|
|
terminal=payload.terminal,
|
|
terminal_status=payload.terminal_status.value if payload.terminal_status else None,
|
|
)
|
|
return ApiResponse[ReportData](success=True, msg="success", data=data, code=0)
|
|
|
|
|
|
@router.post(
|
|
"/{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 并告警,**不回 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)
|
|
|
|
|
|
@router.get(
|
|
"/{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)
|
|
|
|
|
|
@router.get(
|
|
"",
|
|
response_model=ApiResponse[TaskListData],
|
|
dependencies=[Depends(require_bearer_token)],
|
|
)
|
|
async def list_orders(
|
|
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
|
|
)
|
|
return ApiResponse[TaskListData](success=True, msg="success", data=data, code=0)
|