Files
rakuten-api/app/gateway/api/routes/orders.py
T
q792602257andClaude Opus 4.6 07107094a7 实现下单任务网关与本地 worker
按 docs/order-gateway.md 落地:第三个部署单元 app.gateway(:31109)承担任务队列
+ 状态镜像;本地 worker 在 app.trading.worker 内,按 RAKUTEN_ORDER_GATEWAY_URL
决定是否启动。规格 §5 最关键约束已守:租约过期绝不自动重投,恢复只能 reclaim,
worker 收到 lease_count>1 时先核对站点订单。

站点交互(加购/下单/付款/订单列表反查)按规格 §10 留接口缝,site_interact.py
全部 NotImplementedError,verify.py 恒返回 unknown——等真实账号实测后再填,
不写猜测的提交逻辑。

310 个测试全绿,覆盖规格 §9 验收清单 12 条;架构测试守住三方互不 import。

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-07-27 16:25:20 +08:00

182 lines
5.9 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 不新建任务,返回既有任务且 created=false——上游重发
不会变成两单。
"""
data = await container.task_queue.submit(
task_id=payload.task_id, site=payload.site, intent=payload.intent
)
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(...),
wait: int = Query(default=30, ge=0, le=300),
site: str | None = Query(default=None),
container: GatewayContainer = Depends(get_container),
) -> ApiResponse[LeaseData | None]:
"""本地长轮询领取
无可领任务时挂起到 wait 秒后返回 data: null,HTTP 仍 200。已有 leased/running
任务时立即返回空(全局并发度 1)。
"""
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)],
)
async def renew_order(
task_id: str,
payload: RenewRequest,
container: GatewayContainer = Depends(get_container),
) -> ApiResponse[RenewData]:
"""续租。worker 在长任务中每 60 秒调一次,避免租约 TTL 误判"""
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)],
)
async def report_order(
task_id: str,
payload: ReportRequest,
container: GatewayContainer = Depends(get_container),
) -> ApiResponse[ReportData]:
"""本地回报订单状态
同一 (task_id, state) 重复上报幂等。terminal=true 时释放租约并推到终态。
"""
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)],
)
async def reclaim_order(
task_id: str,
payload: ReclaimRequest,
container: GatewayContainer = Depends(get_container),
) -> ApiResponse[ReclaimData]:
"""把 stale 任务重新租给 worker(恢复领取)
返回体的 lease_count 必然 > 1,worker 必须先核对站点订单列表再决定是否
继续执行——核对不出结论时报 needs_human,绝不重新提交。
"""
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)],
)
async def get_order(
task_id: str,
container: GatewayContainer = Depends(get_container),
) -> ApiResponse[TaskDetail]:
"""任务详情 + 完整状态历史"""
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),
site: str | None = Query(default=None),
limit: int = Query(default=50, ge=1, le=500),
offset: int = Query(default=0, ge=0),
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)