diff --git a/app/gateway/api/routes/orders.py b/app/gateway/api/routes/orders.py index 7f8add2..a3bf60d 100644 --- a/app/gateway/api/routes/orders.py +++ b/app/gateway/api/routes/orders.py @@ -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 ) diff --git a/app/gateway/api/routes/queries.py b/app/gateway/api/routes/queries.py index b1ea47f..88011cc 100644 --- a/app/gateway/api/routes/queries.py +++ b/app/gateway/api/routes/queries.py @@ -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 ) diff --git a/app/gateway/models.py b/app/gateway/models.py index b26a902..af2f804 100644 --- a/app/gateway/models.py +++ b/app/gateway/models.py @@ -23,17 +23,40 @@ class SubmitOrderRequest(BaseModel): 不新建任务,返回既有任务且 created=false。 """ - task_id: str | None = None - site: str - intent: dict[str, Any] + task_id: str | None = Field( + default=None, + description="幂等键。不传则服务端生成;同一 task_id 重复提交不新建任务," + "返回既有任务(created=false)——上游重发不会变成两单", + examples=["po-20260727-0001"], + ) + site: str = Field( + description="站点标识。交易服务只覆盖乐天市场,固定 rakuten", + examples=["rakuten"], + ) + intent: dict[str, Any] = Field( + description=( + "下单意图原文。网关不解释内容,原样存库并透传给本地 worker,结构由 trading 侧定义:" + "item_url(必填,商品页 URL);quantity(可选,默认 1);" + "variant_id(多规格商品必填,取自 /api/item_detail 的 variants,不传时 worker 自动选第一个非售罄规格);" + "choice(可选,商品选项,如 \"颜色:赤\",可传字符串或字符串列表);" + "max_total_yen(可选,本次金额上限:确认页实际应付超过即中止并报 needs_human," + "缺省用服务端 RAKUTEN_ORDER_MAX_TOTAL_YEN)" + ), + examples=[{ + "item_url": "https://item.rakuten.co.jp/shop/code/", + "quantity": 1, + "variant_id": "1001", + "max_total_yen": 30000, + }], + ) class SubmitOrderData(BaseModel): """提交响应""" - task_id: str - status: TaskStatus - created: bool # True=本次新建,False=命中既有任务(幂等) + task_id: str = Field(description="任务 ID(幂等键),后续查询与 worker 回报都用它") + status: TaskStatus = Field(description="任务状态,新建恒为 queued") + created: bool = Field(description="True=本次新建,False=命中既有任务(幂等重发)") # ---- GET /api/orders/lease ---- @@ -46,12 +69,20 @@ class LeaseData(BaseModel): 领取(stale → reclaim),worker 必须先核对站点订单列表,见 runner 主循环。 """ - task_id: str - site: str - intent: dict[str, Any] - lease_expires_at: str # ISO8601 UTC - lease_count: int - known_state: OrderState | None # 之前上报过的最新订单状态;首次领取为 null + task_id: str = Field(description="任务 ID(幂等键)") + site: str = Field(description="站点标识(rakuten)") + intent: dict[str, Any] = Field(description="下单意图原文,与提交时一致") + lease_expires_at: str = Field( + description="租约过期时间(ISO8601 UTC)。TTL 由 RAKUTEN_LEASE_TTL_SECONDS 控制" + "(默认 300 秒);执行中每 60 秒调 renew 续租" + ) + lease_count: int = Field( + description="该任务被领取过的次数。>1 表示恢复领取(stale → reclaim):worker 执行前" + "必须先核对站点订单列表确认这单下没下,核对不出结论报 needs_human,绝不直接重新下单" + ) + known_state: OrderState | None = Field( + description="之前上报过的最新订单状态;首次领取为 null。恢复领取时是核对的重要线索" + ) # ---- POST /api/orders/{id}/renew ---- @@ -60,15 +91,15 @@ class LeaseData(BaseModel): class RenewRequest(BaseModel): """续租请求""" - worker_id: str + worker_id: str = Field(description="worker 标识,必须是当前租约持有者,否则 6002") class RenewData(BaseModel): """续租响应""" - task_id: str - lease_expires_at: str - lease_count: int + task_id: str = Field(description="任务 ID") + lease_expires_at: str = Field(description="续租后的新过期时间(ISO8601 UTC)") + lease_count: int = Field(description="领取次数,续租不改变") # ---- POST /api/orders/{id}/report ---- @@ -81,23 +112,39 @@ class ReportRequest(BaseModel): 不产生第二条记录。terminal=true 时释放租约并把任务推到终态。 """ - worker_id: str - state: OrderState - payable_yen: int | None = None - pay_deadline: str | None = None # ISO8601 - site_order_id: str | None = None - evidence_ref: str | None = None # 本地相对路径,不含扩展名 - detail: str = "" - terminal: bool = False - terminal_status: TaskStatus | None = None # terminal=true 时指定终态,缺省按 state 推断 + worker_id: str = Field(description="worker 标识,必须是当前租约持有者,否则 6002") + state: OrderState = Field( + description="本次上报的订单状态:created / in_cart / ordered / awaiting_payment / " + "paid / shipped / delivered / cancelled" + ) + payable_yen: int | None = Field(default=None, description="实际应付金额(日元),进入确认页后上报") + pay_deadline: str | None = Field(default=None, description="付款期限(ISO8601),コンビニ払い 常见三天") + site_order_id: str | None = Field(default=None, description="站点侧订单号(注文番号),下单成功后上报") + evidence_ref: str | None = Field( + default=None, + description="本地证据相对路径,如 po-20260727-0001/03-order-confirm(不含扩展名)。" + "证据原文(HTML/截图)留在本地,不回传", + ) + detail: str = Field(default="", description="补充说明,如付款方式与单号、失败原因") + terminal: bool = Field( + default=False, + description="True 表示任务到此结束:释放租约并推到终态。缺省按 state 推断终态:" + "paid→succeeded,cancelled→failed,其余(如 awaiting_payment 搁置)→needs_human", + ) + terminal_status: TaskStatus | None = Field( + default=None, + description="terminal=true 时显式指定终态(succeeded / failed / needs_human),缺省按 state 推断", + ) class ReportData(BaseModel): """回报响应""" - task_id: str - status: TaskStatus - recorded: bool # False=同 (task_id, state) 已存在,本次为幂等覆盖;True=新写入 + task_id: str = Field(description="任务 ID") + status: TaskStatus = Field(description="回报后的任务状态") + recorded: bool = Field( + description="True=新写入一条状态记录;False=同 (task_id, state) 已存在,本次为幂等覆盖" + ) # ---- POST /api/orders/{id}/reclaim ---- @@ -106,58 +153,66 @@ class ReportData(BaseModel): class ReclaimRequest(BaseModel): """把 stale 任务重新租给 worker""" - worker_id: str + worker_id: str = Field( + description="worker 标识。reclaim 是显式恢复动作——正常 lease 永远拿不到 stale 任务" + ) class ReclaimData(BaseModel): """reclaim 响应,结构与 LeaseData 一致,但 lease_count 必然 > 1""" - task_id: str - site: str - intent: dict[str, Any] - lease_expires_at: str - lease_count: int - known_state: OrderState | None + task_id: str = Field(description="任务 ID(幂等键)") + site: str = Field(description="站点标识(rakuten)") + intent: dict[str, Any] = Field(description="下单意图原文,与提交时一致") + lease_expires_at: str = Field(description="新租约的过期时间(ISO8601 UTC)") + lease_count: int = Field( + description="必然 > 1:worker 执行前必须先核对站点订单列表,确认这单到底下没下" + ) + known_state: OrderState | None = Field( + description="租约丢失前上报过的最新订单状态,恢复核对时的重要线索" + ) # ---- GET /api/orders/{id} 与 GET /api/orders ---- class ReportEntry(BaseModel): - """task_reports 单行""" + """task_reports 单行:一次状态上报(append-only,最新一条即当前状态)""" - state: OrderState - payable_yen: int | None = None - pay_deadline: str | None = None - site_order_id: str | None = None - evidence_ref: str | None = None - detail: str = "" - reported_at: str + state: OrderState = Field(description="上报的订单状态") + payable_yen: int | None = Field(default=None, description="实际应付金额(日元)") + pay_deadline: str | None = Field(default=None, description="付款期限(ISO8601)") + site_order_id: str | None = Field(default=None, description="站点侧订单号(注文番号)") + evidence_ref: str | None = Field(default=None, description="本地证据相对路径(不含扩展名)") + detail: str = Field(default="", description="补充说明") + reported_at: str = Field(description="上报时间(ISO8601 UTC)") class TaskDetail(BaseModel): """单个任务详情""" - task_id: str - site: str - intent: dict[str, Any] - status: TaskStatus - lease_owner: str | None = None - lease_expires_at: str | None = None - lease_count: int = 0 - created_at: str - updated_at: str - latest_state: OrderState | None = None - reports: list[ReportEntry] = Field(default_factory=list) + task_id: str = Field(description="任务 ID(幂等键)") + site: str = Field(description="站点标识(rakuten)") + intent: dict[str, Any] = Field(description="下单意图原文") + status: TaskStatus = Field( + description="任务状态:queued / leased / running / succeeded / failed / needs_human / stale" + ) + lease_owner: str | None = Field(default=None, description="当前租约持有者(worker_id);无租约为 null") + lease_expires_at: str | None = Field(default=None, description="租约过期时间(ISO8601 UTC);无租约为 null") + lease_count: int = Field(default=0, description="被领取过的次数,>1 即发生过恢复") + created_at: str = Field(description="创建时间(ISO8601 UTC)") + updated_at: str = Field(description="最近更新时间(ISO8601 UTC)") + latest_state: OrderState | None = Field(default=None, description="最新一次上报的订单状态;从未上报为 null") + reports: list[ReportEntry] = Field(default_factory=list, description="完整状态历史(append-only)") class TaskListData(BaseModel): """任务列表""" - items: list[TaskDetail] - total: int - limit: int - offset: int + items: list[TaskDetail] = Field(description="任务详情列表,按创建时间倒序") + total: int = Field(description="符合筛选条件的总条数") + limit: int = Field(description="本次分页大小") + offset: int = Field(description="本次分页偏移") # ---- 账号只读查询通道(见 docs/order-gateway.md §11)---- @@ -172,18 +227,34 @@ class SubmitQueryRequest(BaseModel): - order_detail:`{"order_number": "306087-20260813-0863947697"}` """ - query_id: str | None = None - site: str = "rakuten" - kind: AccountQueryKind - params: dict[str, Any] = Field(default_factory=dict) + query_id: str | None = Field( + default=None, + description="幂等键。不传则服务端生成;同一 query_id 重复提交返回既有单(created=false)", + examples=["aq-20260816-0001"], + ) + site: str = Field(default="rakuten", description="站点标识,固定 rakuten") + kind: AccountQueryKind = Field( + description="查询类型:order_list=按时间窗口翻页扫描订单列表;" + "order_detail=读单笔注文番号的详情(配送阶段等)" + ) + params: dict[str, Any] = Field( + default_factory=dict, + description=( + "查询参数,网关原样透传给 worker。order_list:since(可选,ISO8601 窗口下界," + "缺省不设下界)、max_pages(可选,1..20,缺省用服务端 " + "RAKUTEN_ACCOUNT_QUERY_DEFAULT_MAX_PAGES=3);" + "order_detail:order_number(必填,站点注文番号)" + ), + examples=[{"order_number": "306087-20260813-0863947697"}], + ) class SubmitQueryData(BaseModel): """提交响应。纯异步:这里只拿到单号,结果去 GET /api/account/queries/{id} 取""" - query_id: str - status: QueryStatus - created: bool # True=本次新建,False=命中既有单(幂等) + query_id: str = Field(description="查询单 ID(幂等键),取结果时用它") + status: QueryStatus = Field(description="查询单状态,新建恒为 queued") + created: bool = Field(description="True=本次新建,False=命中既有单(幂等重发)") class QueryLeaseData(BaseModel): @@ -194,12 +265,15 @@ class QueryLeaseData(BaseModel): 不需要为此改变行为(只读,重跑安全)。 """ - query_id: str - site: str - kind: AccountQueryKind - params: dict[str, Any] - lease_expires_at: str # ISO8601 UTC - attempt: int + query_id: str = Field(description="查询单 ID") + site: str = Field(description="站点标识(rakuten)") + kind: AccountQueryKind = Field(description="查询类型(order_list / order_detail)") + params: dict[str, Any] = Field(description="查询参数原文,与提交时一致") + lease_expires_at: str = Field( + description="租约过期时间(ISO8601 UTC)。TTL 由 RAKUTEN_QUERY_LEASE_TTL_SECONDS 控制" + "(默认 180 秒),过期自动重投" + ) + attempt: int = Field(description="第几次被领取。>1 表示上一轮超时被重投(只读,重跑安全)") class QueryResultRequest(BaseModel): @@ -209,25 +283,31 @@ class QueryResultRequest(BaseModel): (沿用 shared.errors 的错误码,便于上游按同一张表分支)。 """ - worker_id: str - success: bool - result: dict[str, Any] | None = None - error_code: int | None = None - error_message: str = "" + worker_id: str = Field(description="worker 标识,必须是当前租约持有者,否则 6006") + success: bool = Field(description="站点读取是否成功") + result: dict[str, Any] | None = Field( + default=None, + description="success=true 时必填:规范化字段 + 站点原始 JSON,结构按 kind 见 " + "GET /api/account/queries/{id} 的说明。体积上限 RAKUTEN_QUERY_RESULT_MAX_BYTES(默认 1MB)", + ) + error_code: int | None = Field( + default=None, description="success=false 时的错误码,沿用统一错误码表(如 5001 掉登录)" + ) + error_message: str = Field(default="", description="success=false 时必填,失败原因") class QueryResultData(BaseModel): """回报响应""" - query_id: str - status: QueryStatus + query_id: str = Field(description="查询单 ID") + status: QueryStatus = Field(description="回报后的查询单状态(succeeded / failed)") class QueryError(BaseModel): """查询失败的原因""" - code: int | None = None - message: str = "" + code: int | None = Field(default=None, description="错误码(沿用统一错误码表),可能为 null") + message: str = Field(default="", description="失败原因") class QueryDetail(BaseModel): @@ -236,28 +316,34 @@ class QueryDetail(BaseModel): result 是 worker 回的原文(站点原始 JSON + 规范化字段),网关不解释内容。 """ - query_id: str - site: str - kind: AccountQueryKind - params: dict[str, Any] - status: QueryStatus - lease_owner: str | None = None - lease_expires_at: str | None = None - attempts: int = 0 - created_at: str - updated_at: str - completed_at: str | None = None - result: dict[str, Any] | None = None - error: QueryError | None = None + query_id: str = Field(description="查询单 ID(幂等键)") + site: str = Field(description="站点标识(rakuten)") + kind: AccountQueryKind = Field(description="查询类型(order_list / order_detail)") + params: dict[str, Any] = Field(description="查询参数原文") + status: QueryStatus = Field( + description="查询单状态:queued / leased / succeeded / failed / expired" + ) + lease_owner: str | None = Field(default=None, description="当前租约持有者(worker_id);无租约为 null") + lease_expires_at: str | None = Field(default=None, description="租约过期时间(ISO8601 UTC);无租约为 null") + attempts: int = Field(default=0, description="已被领取的次数,用尽(默认 3 次)置 failed") + created_at: str = Field(description="创建时间(ISO8601 UTC)") + updated_at: str = Field(description="最近更新时间(ISO8601 UTC)") + completed_at: str | None = Field(default=None, description="到达终态的时间(ISO8601 UTC);未终结为 null") + result: dict[str, Any] | None = Field( + default=None, + description="status=succeeded 时的结果:规范化字段 + 站点原始 JSON,结构按 kind " + "见接口说明;未成功为 null", + ) + error: QueryError | None = Field(default=None, description="status=failed/expired 时的失败原因;否则为 null") class QueryListData(BaseModel): """查询单列表(运维排查用)""" - items: list[QueryDetail] - total: int - limit: int - offset: int + items: list[QueryDetail] = Field(description="查询单详情列表,按创建时间倒序") + total: int = Field(description="符合筛选条件的总条数") + limit: int = Field(description="本次分页大小") + offset: int = Field(description="本次分页偏移") # ---- 定时下派通道:账号订单编目(docs/order-gateway.md §12)----