"""下单任务网关路由 承担规格 §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)