"""网关 API 数据模型:请求体与响应体 intent 字段刻意保留成 `dict[str, Any]`——网关不解释下单意图,结构由 trading 侧 定义。网关只负责把它存下来、原样吐给 worker,避免业务规则悄悄渗进任务队列。 """ from __future__ import annotations from typing import Any from pydantic import BaseModel, Field from app.shared.task_state import OrderState, TaskStatus # ---- POST /api/orders ---- class SubmitOrderRequest(BaseModel): """上游提交下单意图 task_id 可选:上游自带的幂等键。不传则服务端生成。同一个 task_id 重复提交 不新建任务,返回既有任务且 created=false。 """ task_id: str | None = None site: str intent: dict[str, Any] class SubmitOrderData(BaseModel): """提交响应""" task_id: str status: TaskStatus created: bool # True=本次新建,False=命中既有任务(幂等) # ---- GET /api/orders/lease ---- class LeaseData(BaseModel): """lease 响应 无任务可领时整个 data 为 null(HTTP 仍 200)。lease_count > 1 表示这是恢复 领取(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 # ---- POST /api/orders/{id}/renew ---- class RenewRequest(BaseModel): """续租请求""" worker_id: str class RenewData(BaseModel): """续租响应""" task_id: str lease_expires_at: str lease_count: int # ---- POST /api/orders/{id}/report ---- class ReportRequest(BaseModel): """本地回报订单状态 同一 (task_id, state) 重复上报是幂等的——网络抖动导致 worker 重发时覆盖同一行, 不产生第二条记录。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 推断 class ReportData(BaseModel): """回报响应""" task_id: str status: TaskStatus recorded: bool # False=同 (task_id, state) 已存在,本次为幂等覆盖;True=新写入 # ---- POST /api/orders/{id}/reclaim ---- class ReclaimRequest(BaseModel): """把 stale 任务重新租给 worker""" worker_id: str 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 # ---- GET /api/orders/{id} 与 GET /api/orders ---- class ReportEntry(BaseModel): """task_reports 单行""" 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 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) class TaskListData(BaseModel): """任务列表""" items: list[TaskDetail] total: int limit: int offset: int # ---- GET /health ---- class WorkerHealthEntry(BaseModel): """单个 worker 的健康指标""" worker_id: str last_seen_at: str last_seen_seconds: int class QueuedAlertEntry(BaseModel): """长时间无人领的任务告警项""" task_id: str site: str created_at: str age_seconds: int class GatewayHealthData(BaseModel): """网关健康状态 只做规格 §4.7 的两条兜底告警:worker 失联、任务长时间无人领。付款期限监控 不在网关——本地机 7×24 在线,那套逻辑放本地。 """ status: str # "ok" 或 "degraded" queued_count: int active_tasks: list[TaskDetail] = Field(default_factory=list) workers: list[WorkerHealthEntry] = Field(default_factory=list) offline_workers: list[WorkerHealthEntry] = Field(default_factory=list) stale_queued_tasks: list[QueuedAlertEntry] = Field(default_factory=list)