feat(gateway): 账号只读查询通道——从已登录账号取真实订单
上游要的不只是网关记的任务状态镜像,还有「已登录账号在站点上的真实订单」,
但账号只在 NAT 后本地机上,只能经网关队列走。新增独立查询通道(§11):
- gateway 单开 account_queries 表 + QueryStatus 状态机,接口
POST /api/account/queries(幂等)/ lease / {id}/result / {id}
- 不复用下单任务队列:查询是只读,租约过期可安全重投(与下单「绝不自动
重投」相反),且不该被全局并发度 1 堵死、task_reports 是订单镜像不能污染
- 本地交易服务起第二条常驻循环 query_runner,领到即调 SiteInteractor 真读:
order_list 复用已实测的 list_recent_orders(规范化字段 + 站点
orderListData 原文),order_detail 复用 fetch_order_detail(配送阶段 +
页面 __INITIAL_STATE__ 原样透传,结构未经真实样本,不抽字段)
- 账号级串行仍由 SiteInteractor 的锁保证;每次执行套超时按失败回报
- 错误码 6005/6006(查询通道,可重试只读区别于 6001-6004);/health 暴露
queued_query_count;结果体积上限先丢原始 JSON
openapi.json 重导,docs/order-gateway.md §11、README、.env.example 补全
配置与实测边界。全量测试 404→454 通过。
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -16,8 +16,11 @@ async def health(
|
||||
|
||||
只做规格 §4.7 的两条兜底告警:worker 失联、任务长时间无人领。
|
||||
付款期限监控不在网关——本地机 7×24 在线,那套逻辑放本地。
|
||||
另外带上账号只读查询的积压数量:它不参与告警判定(查询有自己的 TTL 会自动
|
||||
expired),只是给运维一个可见信号。
|
||||
"""
|
||||
data = await container.task_queue.health_snapshot()
|
||||
data.queued_query_count = await container.query_queue.queued_count()
|
||||
return ApiResponse[GatewayHealthData](
|
||||
success=True, msg="success", data=data, code=0
|
||||
)
|
||||
|
||||
@@ -0,0 +1,145 @@
|
||||
"""账号只读查询通道路由(docs/order-gateway.md §11)
|
||||
|
||||
上游要的是「已登录账号在站点上的真实订单」,而不是网关自己记录的任务状态镜像;
|
||||
但本地机在 NAT 后没有公网入口,网关推不进去。所以这条通道与下单任务同构:
|
||||
上游把查询意图放进队列,本地 worker 出站长轮询领走、去站点上真读一次、回结果。
|
||||
|
||||
- POST /api/account/queries 上游提交查询单(幂等)
|
||||
- GET /api/account/queries/lease 本地长轮询领取
|
||||
- POST /api/account/queries/{id}/result 本地回结果
|
||||
- GET /api/account/queries/{id} 取查询单状态与结果
|
||||
- GET /api/account/queries 查询单列表(运维排查)
|
||||
|
||||
**纯异步**:提交只拿 query_id,结果去 GET 取。网关不提供「挂起等结果」的接口——
|
||||
下单任务可能占着账号锁跑好几分钟,同步等待会把上游的连接一起卡住。
|
||||
|
||||
注意路由声明顺序:`/lease` 必须在 `/{query_id}` 之前,否则 lease 会被当成 query_id。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from fastapi import APIRouter, Depends, Query
|
||||
|
||||
from app.gateway.container import GatewayContainer
|
||||
from app.gateway.models import (
|
||||
QueryDetail,
|
||||
QueryLeaseData,
|
||||
QueryListData,
|
||||
QueryResultData,
|
||||
QueryResultRequest,
|
||||
SubmitQueryData,
|
||||
SubmitQueryRequest,
|
||||
)
|
||||
from app.shared.api import ApiResponse, get_container, require_bearer_token
|
||||
|
||||
router = APIRouter(prefix="/api/account/queries", tags=["account-queries"])
|
||||
|
||||
|
||||
@router.post(
|
||||
"",
|
||||
response_model=ApiResponse[SubmitQueryData],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def submit_query(
|
||||
payload: SubmitQueryRequest,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[SubmitQueryData]:
|
||||
"""提交一张账号只读查询单
|
||||
|
||||
重复提交同一 query_id 不新建,返回既有单且 created=false。查询是只读的,
|
||||
重复提交本身无害,幂等键的意义在于上游重发时不会拿到两个单号。
|
||||
"""
|
||||
data = await container.query_queue.submit(
|
||||
query_id=payload.query_id,
|
||||
site=payload.site,
|
||||
kind=payload.kind.value,
|
||||
params=payload.params,
|
||||
)
|
||||
return ApiResponse[SubmitQueryData](success=True, msg="success", data=data, code=0)
|
||||
|
||||
|
||||
@router.get(
|
||||
"/lease",
|
||||
response_model=ApiResponse[QueryLeaseData | None],
|
||||
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),
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[QueryLeaseData | None]:
|
||||
"""本地长轮询领取查询单
|
||||
|
||||
无单可领时挂起到 wait 秒后返回 data: null,HTTP 仍 200。与下单 lease 不同,
|
||||
这里没有「已有任务在执行就一律返回空」的闸门。
|
||||
"""
|
||||
data = await container.query_queue.lease(
|
||||
worker_id=worker_id,
|
||||
wait=wait,
|
||||
site=site,
|
||||
max_wait=container.settings.lease_max_wait_seconds,
|
||||
)
|
||||
return ApiResponse[QueryLeaseData | None](success=True, msg="success", data=data, code=0)
|
||||
|
||||
|
||||
@router.post(
|
||||
"/{query_id}/result",
|
||||
response_model=ApiResponse[QueryResultData],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def submit_query_result(
|
||||
query_id: str,
|
||||
payload: QueryResultRequest,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[QueryResultData]:
|
||||
"""本地回报查询结果
|
||||
|
||||
租约不是自己的(多半是超时后被重投给了下一轮)时报 6006 并丢弃本次结果,
|
||||
避免迟到的旧结果覆盖新结果。
|
||||
"""
|
||||
data = await container.query_queue.submit_result(
|
||||
query_id,
|
||||
worker_id=payload.worker_id,
|
||||
success=payload.success,
|
||||
result=payload.result,
|
||||
error_code=payload.error_code,
|
||||
error_message=payload.error_message,
|
||||
)
|
||||
return ApiResponse[QueryResultData](success=True, msg="success", data=data, code=0)
|
||||
|
||||
|
||||
@router.get(
|
||||
"/{query_id}",
|
||||
response_model=ApiResponse[QueryDetail],
|
||||
dependencies=[Depends(require_bearer_token)],
|
||||
)
|
||||
async def get_query(
|
||||
query_id: str,
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[QueryDetail]:
|
||||
"""取查询单状态与结果
|
||||
|
||||
status=succeeded 时 result 里是 worker 从站点读到的原文 + 规范化字段;
|
||||
failed/expired 时看 error。仍在 queued/leased 就是还没读到,稍后再来。
|
||||
"""
|
||||
data = await container.query_queue.get_detail(query_id)
|
||||
return ApiResponse[QueryDetail](success=True, msg="success", data=data, code=0)
|
||||
|
||||
|
||||
@router.get(
|
||||
"",
|
||||
response_model=ApiResponse[QueryListData],
|
||||
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),
|
||||
container: GatewayContainer = Depends(get_container),
|
||||
) -> ApiResponse[QueryListData]:
|
||||
"""查询单列表,按创建时间倒序(运维排查用)"""
|
||||
data = await container.query_queue.list_queries(
|
||||
status=status, kind=kind, limit=limit, offset=offset
|
||||
)
|
||||
return ApiResponse[QueryListData](success=True, msg="success", data=data, code=0)
|
||||
@@ -6,6 +6,7 @@ import logging
|
||||
from dataclasses import dataclass
|
||||
|
||||
from app.gateway.db import GatewayDB
|
||||
from app.gateway.query_queue import QueryQueue
|
||||
from app.gateway.task_queue import TaskQueue
|
||||
from app.shared.config import Settings
|
||||
|
||||
@@ -19,11 +20,17 @@ class GatewayContainer:
|
||||
与抓取/交易容器同样的依赖注入风格,但持有的是任务队列与 SQLite 连接,
|
||||
生命周期由 main.lifespan 管理:start 时打开 DB,close 时关闭。
|
||||
|
||||
`sweep_task` 是常驻后台扫描,把过期的 leased/running 推到 stale。
|
||||
即使没有 lease 请求,过期的任务也会被及时发现(规格 §5 的关键约束)。
|
||||
两条通道各有一个队列对象,共用同一个 GatewayDB:
|
||||
- `task_queue` 下单任务(不可逆写,租约过期只置 stale)
|
||||
- `query_queue` 账号只读查询(可安全重投,见 query_queue.py 顶部对照表)
|
||||
|
||||
`sweep_task` 是常驻后台扫描,把过期的 leased/running 推到 stale,并顺带扫
|
||||
一轮查询单(重投 / 过期 / 清理过保留期的结果)。即使没有 lease 请求,过期的
|
||||
任务也会被及时发现(规格 §5 的关键约束)。
|
||||
"""
|
||||
|
||||
settings: Settings
|
||||
db: GatewayDB
|
||||
task_queue: TaskQueue
|
||||
query_queue: QueryQueue
|
||||
sweep_task: asyncio.Task | None = None
|
||||
|
||||
+227
-4
@@ -1,10 +1,11 @@
|
||||
"""网关 SQLite 访问层:三张表 + 原生 SQL,无业务逻辑
|
||||
"""网关 SQLite 访问层:四张表 + 原生 SQL,无业务逻辑
|
||||
|
||||
业务规则(状态机、并发度、幂等)放在 `task_queue.py`,本模块只负责持久化与查询,
|
||||
业务规则(状态机、并发度、幂等)放在 `task_queue.py`(下单任务)与
|
||||
`query_queue.py`(账号只读查询),本模块只负责持久化与查询,
|
||||
返回 dataclass 行。所有时间戳以 ISO8601 UTC 字符串存储(以 `Z` 结尾),便于
|
||||
跨进程对账;时间戳运算在 task_queue 层完成。
|
||||
跨进程对账;时间戳运算在上层完成。
|
||||
|
||||
写操作由 task_queue 的 asyncio.Lock 串行化(详见该模块),本层不重复加锁,
|
||||
写操作由上层的 asyncio.Lock 串行化(详见对应模块),本层不重复加锁,
|
||||
因此**调用方必须确保写操作在外层锁的保护下进行**。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
@@ -49,6 +50,27 @@ CREATE TABLE IF NOT EXISTS workers (
|
||||
worker_id TEXT PRIMARY KEY,
|
||||
last_seen_at TEXT NOT NULL
|
||||
);
|
||||
|
||||
-- 账号只读查询单(见 docs/order-gateway.md §11)。与 tasks 分表,不是洁癖:
|
||||
-- 两者的安全约束相反(只读可重投 vs 写操作绝不重投),共表会让「租约过期怎么办」
|
||||
-- 这个最关键的分支变成一个 if,迟早被改错。
|
||||
CREATE TABLE IF NOT EXISTS account_queries (
|
||||
query_id TEXT PRIMARY KEY, -- 幂等键,语义同 task_id
|
||||
site TEXT NOT NULL,
|
||||
kind TEXT NOT NULL, -- order_list / order_detail
|
||||
params_json TEXT NOT NULL, -- 查询参数原文,网关不解释内容
|
||||
status TEXT NOT NULL, -- 见 shared.task_state.QueryStatus
|
||||
lease_owner TEXT,
|
||||
lease_expires_at TEXT,
|
||||
attempts INTEGER NOT NULL DEFAULT 0, -- 被领取过几次
|
||||
result_json TEXT, -- worker 回的结果原文
|
||||
error_code INTEGER,
|
||||
error_message TEXT,
|
||||
created_at TEXT NOT NULL,
|
||||
updated_at TEXT NOT NULL,
|
||||
completed_at TEXT
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_queries_pending ON account_queries(status, created_at);
|
||||
"""
|
||||
|
||||
|
||||
@@ -93,6 +115,34 @@ class WorkerRow:
|
||||
last_seen_at: str
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class QueryRow:
|
||||
"""account_queries 表一行的强类型视图"""
|
||||
|
||||
query_id: str
|
||||
site: str
|
||||
kind: str
|
||||
params_json: str
|
||||
status: str
|
||||
lease_owner: str | None
|
||||
lease_expires_at: str | None
|
||||
attempts: int
|
||||
result_json: str | None
|
||||
error_code: int | None
|
||||
error_message: str | None
|
||||
created_at: str
|
||||
updated_at: str
|
||||
completed_at: str | None
|
||||
|
||||
@property
|
||||
def params(self) -> dict[str, Any]:
|
||||
return json.loads(self.params_json)
|
||||
|
||||
@property
|
||||
def result(self) -> dict[str, Any] | None:
|
||||
return json.loads(self.result_json) if self.result_json else None
|
||||
|
||||
|
||||
def _row_to_task(row: aiosqlite.Row) -> TaskRow:
|
||||
return TaskRow(
|
||||
task_id=row["task_id"],
|
||||
@@ -124,6 +174,25 @@ def _row_to_worker(row: aiosqlite.Row) -> WorkerRow:
|
||||
return WorkerRow(worker_id=row["worker_id"], last_seen_at=row["last_seen_at"])
|
||||
|
||||
|
||||
def _row_to_query(row: aiosqlite.Row) -> QueryRow:
|
||||
return QueryRow(
|
||||
query_id=row["query_id"],
|
||||
site=row["site"],
|
||||
kind=row["kind"],
|
||||
params_json=row["params_json"],
|
||||
status=row["status"],
|
||||
lease_owner=row["lease_owner"],
|
||||
lease_expires_at=row["lease_expires_at"],
|
||||
attempts=row["attempts"],
|
||||
result_json=row["result_json"],
|
||||
error_code=row["error_code"],
|
||||
error_message=row["error_message"],
|
||||
created_at=row["created_at"],
|
||||
updated_at=row["updated_at"],
|
||||
completed_at=row["completed_at"],
|
||||
)
|
||||
|
||||
|
||||
class GatewayDB:
|
||||
"""网关 SQLite 访问对象
|
||||
|
||||
@@ -354,3 +423,157 @@ class GatewayDB:
|
||||
(status,),
|
||||
) as cur:
|
||||
return (await cur.fetchone())[0]
|
||||
|
||||
# ---- account_queries ----
|
||||
|
||||
async def get_query(self, query_id: str) -> QueryRow | None:
|
||||
async with self.conn.execute(
|
||||
"SELECT * FROM account_queries WHERE query_id = ?", (query_id,)
|
||||
) as cur:
|
||||
row = await cur.fetchone()
|
||||
return _row_to_query(row) if row else None
|
||||
|
||||
async def insert_query(self, row: QueryRow) -> bool:
|
||||
"""插入新查询单。返回 True=新建,False=query_id 已存在(幂等命中)"""
|
||||
try:
|
||||
await self.conn.execute(
|
||||
"INSERT INTO account_queries (query_id, site, kind, params_json, status, "
|
||||
"lease_owner, lease_expires_at, attempts, result_json, error_code, "
|
||||
"error_message, created_at, updated_at, completed_at) "
|
||||
"VALUES (?, ?, ?, ?, ?, NULL, NULL, 0, NULL, NULL, NULL, ?, ?, NULL)",
|
||||
(
|
||||
row.query_id,
|
||||
row.site,
|
||||
row.kind,
|
||||
row.params_json,
|
||||
row.status,
|
||||
row.created_at,
|
||||
row.updated_at,
|
||||
),
|
||||
)
|
||||
await self.conn.commit()
|
||||
return True
|
||||
except aiosqlite.IntegrityError:
|
||||
return False
|
||||
|
||||
async def update_query(
|
||||
self,
|
||||
query_id: str,
|
||||
*,
|
||||
status: str | None = None,
|
||||
lease_owner: str | None = None,
|
||||
lease_expires_at: str | None = None,
|
||||
attempts: int | None = None,
|
||||
result_json: str | None = None,
|
||||
error_code: int | None = None,
|
||||
error_message: str | None = None,
|
||||
updated_at: str | None = None,
|
||||
completed_at: str | None = None,
|
||||
clear_lease: bool = False,
|
||||
) -> None:
|
||||
"""更新查询单字段。clear_lease=True 时把 lease_owner/expires_at 置 NULL"""
|
||||
sets: list[str] = []
|
||||
params: list[Any] = []
|
||||
for column, value in (
|
||||
("status", status),
|
||||
("lease_owner", lease_owner),
|
||||
("lease_expires_at", lease_expires_at),
|
||||
("attempts", attempts),
|
||||
("result_json", result_json),
|
||||
("error_code", error_code),
|
||||
("error_message", error_message),
|
||||
("updated_at", updated_at),
|
||||
("completed_at", completed_at),
|
||||
):
|
||||
if value is not None:
|
||||
sets.append(f"{column} = ?")
|
||||
params.append(value)
|
||||
if clear_lease:
|
||||
sets.append("lease_owner = NULL")
|
||||
sets.append("lease_expires_at = NULL")
|
||||
if not sets:
|
||||
return
|
||||
params.append(query_id)
|
||||
await self.conn.execute(
|
||||
f"UPDATE account_queries SET {', '.join(sets)} WHERE query_id = ?",
|
||||
params,
|
||||
)
|
||||
await self.conn.commit()
|
||||
|
||||
async def pick_queued_query(self, site: str | None) -> QueryRow | None:
|
||||
"""取最早的 queued 查询单,可选按站点过滤"""
|
||||
if site:
|
||||
sql = (
|
||||
"SELECT * FROM account_queries WHERE status = ? AND site = ? "
|
||||
"ORDER BY created_at ASC LIMIT 1"
|
||||
)
|
||||
params: tuple[Any, ...] = ("queued", site)
|
||||
else:
|
||||
sql = "SELECT * FROM account_queries WHERE status = ? ORDER BY created_at ASC LIMIT 1"
|
||||
params = ("queued",)
|
||||
async with self.conn.execute(sql, params) as cur:
|
||||
row = await cur.fetchone()
|
||||
return _row_to_query(row) if row else None
|
||||
|
||||
async def list_queries(
|
||||
self,
|
||||
*,
|
||||
status: str | None = None,
|
||||
kind: str | None = None,
|
||||
limit: int = 50,
|
||||
offset: int = 0,
|
||||
) -> tuple[list[QueryRow], int]:
|
||||
"""分页列出查询单,按 created_at 降序(查询是即时行为,最近的先看)"""
|
||||
where = []
|
||||
params: list[Any] = []
|
||||
if status:
|
||||
where.append("status = ?")
|
||||
params.append(status)
|
||||
if kind:
|
||||
where.append("kind = ?")
|
||||
params.append(kind)
|
||||
clause = f"WHERE {' AND '.join(where)}" if where else ""
|
||||
|
||||
async with self.conn.execute(
|
||||
f"SELECT COUNT(*) FROM account_queries {clause}", params
|
||||
) as cur:
|
||||
total = (await cur.fetchone())[0]
|
||||
|
||||
sql = (
|
||||
f"SELECT * FROM account_queries {clause} "
|
||||
"ORDER BY created_at DESC LIMIT ? OFFSET ?"
|
||||
)
|
||||
async with self.conn.execute(sql, [*params, limit, offset]) as cur:
|
||||
rows = await cur.fetchall()
|
||||
return [_row_to_query(r) for r in rows], total
|
||||
|
||||
async def list_queries_in_statuses(self, statuses: tuple[str, ...]) -> list[QueryRow]:
|
||||
"""取处于给定状态集合的全部查询单(sweep 与健康检查用)"""
|
||||
if not statuses:
|
||||
return []
|
||||
placeholders = ", ".join("?" for _ in statuses)
|
||||
sql = (
|
||||
f"SELECT * FROM account_queries WHERE status IN ({placeholders}) "
|
||||
"ORDER BY created_at ASC"
|
||||
)
|
||||
async with self.conn.execute(sql, statuses) as cur:
|
||||
rows = await cur.fetchall()
|
||||
return [_row_to_query(r) for r in rows]
|
||||
|
||||
async def count_queries_by_status(self, status: str) -> int:
|
||||
async with self.conn.execute(
|
||||
"SELECT COUNT(*) FROM account_queries WHERE status = ?", (status,)
|
||||
) as cur:
|
||||
return (await cur.fetchone())[0]
|
||||
|
||||
async def delete_queries_completed_before(self, cutoff: str) -> int:
|
||||
"""清理 completed_at 早于 cutoff 的终态查询单,返回删除条数
|
||||
|
||||
结果里带站点原始 JSON,不清会一直涨。只删已完成的,进行中的一律不动。
|
||||
"""
|
||||
cursor = await self.conn.execute(
|
||||
"DELETE FROM account_queries WHERE completed_at IS NOT NULL AND completed_at < ?",
|
||||
(cutoff,),
|
||||
)
|
||||
await self.conn.commit()
|
||||
return cursor.rowcount or 0
|
||||
|
||||
+24
-3
@@ -16,8 +16,10 @@ from fastapi import FastAPI
|
||||
|
||||
from app.gateway.api.routes.health import router as health_router
|
||||
from app.gateway.api.routes.orders import router as orders_router
|
||||
from app.gateway.api.routes.queries import router as queries_router
|
||||
from app.gateway.container import GatewayContainer
|
||||
from app.gateway.db import GatewayDB
|
||||
from app.gateway.query_queue import QueryQueue
|
||||
from app.gateway.task_queue import TaskQueue
|
||||
from app.shared.api import register_exception_handlers
|
||||
from app.shared.config import get_settings
|
||||
@@ -29,7 +31,7 @@ SWEEP_INTERVAL_SECONDS = 60
|
||||
|
||||
|
||||
def build_container() -> GatewayContainer:
|
||||
"""构建网关容器:DB + 任务队列"""
|
||||
"""构建网关容器:DB + 下单任务队列 + 账号只读查询队列"""
|
||||
settings = get_settings()
|
||||
db = GatewayDB(settings.gateway_db_path_resolved)
|
||||
task_queue = TaskQueue(
|
||||
@@ -37,11 +39,23 @@ def build_container() -> GatewayContainer:
|
||||
lease_ttl_seconds=settings.lease_ttl_seconds,
|
||||
worker_offline_alert_seconds=settings.worker_offline_alert_seconds,
|
||||
)
|
||||
return GatewayContainer(settings=settings, db=db, task_queue=task_queue)
|
||||
query_queue = QueryQueue(
|
||||
db,
|
||||
lease_ttl_seconds=settings.query_lease_ttl_seconds,
|
||||
query_ttl_seconds=settings.query_ttl_seconds,
|
||||
max_attempts=settings.query_max_attempts,
|
||||
retention_seconds=settings.query_retention_seconds,
|
||||
)
|
||||
return GatewayContainer(
|
||||
settings=settings, db=db, task_queue=task_queue, query_queue=query_queue
|
||||
)
|
||||
|
||||
|
||||
async def _sweep_loop(container: GatewayContainer) -> None:
|
||||
"""常驻后台任务:每 60 秒扫一次过期租约,把 leased/running 推到 stale
|
||||
"""常驻后台任务:每 60 秒扫一次过期租约
|
||||
|
||||
下单侧把 leased/running 推到 stale;查询侧重投超时的、置 expired 超 TTL 的、
|
||||
清掉过保留期的结果。
|
||||
|
||||
lease 请求本身也会顺带扫一次,但 worker 退场后没人来 lease,必须靠这个
|
||||
兜底——否则过期的任务永远停在 active 状态,健康检查看不到,stale 任务也
|
||||
@@ -53,6 +67,9 @@ async def _sweep_loop(container: GatewayContainer) -> None:
|
||||
swept = await container.task_queue.sweep()
|
||||
if swept:
|
||||
logger.info("sweep 把 %s 个过期任务置为 stale", swept)
|
||||
query_swept = await container.query_queue.sweep()
|
||||
if query_swept:
|
||||
logger.info("sweep 处理了 %s 张超时/过期的查询单", query_swept)
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
except Exception: # noqa: BLE001
|
||||
@@ -78,6 +95,9 @@ async def lifespan(app: FastAPI):
|
||||
startup_swept = await container.task_queue.sweep()
|
||||
if startup_swept:
|
||||
logger.warning("启动时把 %s 个遗留过期任务置为 stale", startup_swept)
|
||||
startup_queries = await container.query_queue.sweep()
|
||||
if startup_queries:
|
||||
logger.warning("启动时处理了 %s 张遗留的超时/过期查询单", startup_queries)
|
||||
|
||||
container.sweep_task = asyncio.create_task(
|
||||
_sweep_loop(container), name="gateway-sweep"
|
||||
@@ -101,6 +121,7 @@ def create_app() -> FastAPI:
|
||||
app = FastAPI(title="Rakuten Order Gateway", lifespan=lifespan)
|
||||
app.include_router(health_router)
|
||||
app.include_router(orders_router)
|
||||
app.include_router(queries_router)
|
||||
register_exception_handlers(app)
|
||||
return app
|
||||
|
||||
|
||||
+105
-1
@@ -2,6 +2,7 @@
|
||||
|
||||
intent 字段刻意保留成 `dict[str, Any]`——网关不解释下单意图,结构由 trading 侧
|
||||
定义。网关只负责把它存下来、原样吐给 worker,避免业务规则悄悄渗进任务队列。
|
||||
账号只读查询的 `params` / `result` 同理。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
@@ -9,7 +10,7 @@ from typing import Any
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
from app.shared.task_state import OrderState, TaskStatus
|
||||
from app.shared.task_state import AccountQueryKind, OrderState, QueryStatus, TaskStatus
|
||||
|
||||
|
||||
# ---- POST /api/orders ----
|
||||
@@ -159,6 +160,106 @@ class TaskListData(BaseModel):
|
||||
offset: int
|
||||
|
||||
|
||||
# ---- 账号只读查询通道(见 docs/order-gateway.md §11)----
|
||||
|
||||
|
||||
class SubmitQueryRequest(BaseModel):
|
||||
"""上游提交一张账号只读查询单
|
||||
|
||||
query_id 语义同 task_id:上游自带的幂等键,重复提交返回既有单。
|
||||
params 由网关原样透传给 worker,结构由 trading 侧定义(见 §11.2):
|
||||
- order_list:`{"since": ISO8601?, "max_pages": int?}`
|
||||
- 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)
|
||||
|
||||
|
||||
class SubmitQueryData(BaseModel):
|
||||
"""提交响应。纯异步:这里只拿到单号,结果去 GET /api/account/queries/{id} 取"""
|
||||
|
||||
query_id: str
|
||||
status: QueryStatus
|
||||
created: bool # True=本次新建,False=命中既有单(幂等)
|
||||
|
||||
|
||||
class QueryLeaseData(BaseModel):
|
||||
"""GET /api/account/queries/lease 的响应
|
||||
|
||||
无可领查询单时整个 data 为 null(HTTP 仍 200)。没有 known_state 之类的字段——
|
||||
查询是无状态的一次性动作,attempt > 1 只是说明上一轮超时被重投了,worker
|
||||
不需要为此改变行为(只读,重跑安全)。
|
||||
"""
|
||||
|
||||
query_id: str
|
||||
site: str
|
||||
kind: AccountQueryKind
|
||||
params: dict[str, Any]
|
||||
lease_expires_at: str # ISO8601 UTC
|
||||
attempt: int
|
||||
|
||||
|
||||
class QueryResultRequest(BaseModel):
|
||||
"""worker 回报查询结果
|
||||
|
||||
success=true 时 result 必填;false 时 error_message 必填、error_code 可选
|
||||
(沿用 shared.errors 的错误码,便于上游按同一张表分支)。
|
||||
"""
|
||||
|
||||
worker_id: str
|
||||
success: bool
|
||||
result: dict[str, Any] | None = None
|
||||
error_code: int | None = None
|
||||
error_message: str = ""
|
||||
|
||||
|
||||
class QueryResultData(BaseModel):
|
||||
"""回报响应"""
|
||||
|
||||
query_id: str
|
||||
status: QueryStatus
|
||||
|
||||
|
||||
class QueryError(BaseModel):
|
||||
"""查询失败的原因"""
|
||||
|
||||
code: int | None = None
|
||||
message: str = ""
|
||||
|
||||
|
||||
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
|
||||
|
||||
|
||||
class QueryListData(BaseModel):
|
||||
"""查询单列表(运维排查用)"""
|
||||
|
||||
items: list[QueryDetail]
|
||||
total: int
|
||||
limit: int
|
||||
offset: int
|
||||
|
||||
|
||||
# ---- GET /health ----
|
||||
|
||||
|
||||
@@ -188,6 +289,9 @@ class GatewayHealthData(BaseModel):
|
||||
|
||||
status: str # "ok" 或 "degraded"
|
||||
queued_count: int
|
||||
# 等待本地 worker 领取的账号只读查询单数量。查询没有「无人领即告警」这条
|
||||
# 规则(它有自己的 TTL 会自动 expired),这里只是给运维一个可见的积压信号。
|
||||
queued_query_count: int = 0
|
||||
active_tasks: list[TaskDetail] = Field(default_factory=list)
|
||||
workers: list[WorkerHealthEntry] = Field(default_factory=list)
|
||||
offline_workers: list[WorkerHealthEntry] = Field(default_factory=list)
|
||||
|
||||
@@ -0,0 +1,357 @@
|
||||
"""账号只读查询通道的业务逻辑:状态机、租约、重投、过期与保留期清理
|
||||
|
||||
与 `task_queue.py` 是**刻意分开的两套语义**,不是重复代码(见 docs/order-gateway.md §11):
|
||||
|
||||
| | 下单任务队列 | 本模块(只读查询) |
|
||||
| --- | --- | --- |
|
||||
| 操作性质 | 不可逆写 | 只读 |
|
||||
| 全局并发度 1 | 是(同账号写操作必须串行) | 否(账号级串行由本地 SiteInteractor 的锁兜底) |
|
||||
| 租约过期 | 置 stale,**绝不自动重投**,等人工 reclaim | 回 queued **自动重投**,重跑无副作用 |
|
||||
| 续租 | 有(下单要几分钟) | 无(超时即重投,比续租简单且安全) |
|
||||
| 心跳 | lease 兼作心跳 | **不刷心跳**(见 `lease` 注释) |
|
||||
|
||||
这两张表共用一个 `GatewayDB` 连接,但各自持有自己的锁:查询的长轮询等待不该
|
||||
把下单任务的 lease/report 一起挂住。SQLite 单连接本身是串行的,两把锁只是各自
|
||||
保护「检查 + 写入」的原子序列。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import secrets
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
from app.gateway.db import GatewayDB, QueryRow
|
||||
from app.gateway.models import (
|
||||
QueryDetail,
|
||||
QueryError,
|
||||
QueryLeaseData,
|
||||
QueryListData,
|
||||
QueryResultData,
|
||||
SubmitQueryData,
|
||||
)
|
||||
from app.shared.errors import QueryLeaseInvalidError, QueryNotFoundError
|
||||
from app.shared.task_state import QUERY_TERMINAL_STATUSES, QueryStatus
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def utcnow() -> datetime:
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
def to_iso(dt: datetime) -> str:
|
||||
"""统一时间戳格式:ISO8601 UTC,秒精度,Z 结尾(与 task_queue 一致)"""
|
||||
return dt.astimezone(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z")
|
||||
|
||||
|
||||
def parse_iso(s: str) -> datetime:
|
||||
if s.endswith("Z"):
|
||||
s = s[:-1] + "+00:00"
|
||||
return datetime.fromisoformat(s)
|
||||
|
||||
|
||||
def _generate_query_id() -> str:
|
||||
"""服务端生成的 query_id:与 task_id 用不同前缀,日志里一眼能分清两条通道"""
|
||||
return f"q-{to_iso(utcnow())[:10].replace('-', '')}-{secrets.token_hex(4)}"
|
||||
|
||||
|
||||
def _row_to_detail(row: QueryRow) -> QueryDetail:
|
||||
error = None
|
||||
if row.error_code is not None or row.error_message:
|
||||
error = QueryError(code=row.error_code, message=row.error_message or "")
|
||||
return QueryDetail(
|
||||
query_id=row.query_id,
|
||||
site=row.site,
|
||||
kind=row.kind, # type: ignore[arg-type]
|
||||
params=row.params,
|
||||
status=QueryStatus(row.status),
|
||||
lease_owner=row.lease_owner,
|
||||
lease_expires_at=row.lease_expires_at,
|
||||
attempts=row.attempts,
|
||||
created_at=row.created_at,
|
||||
updated_at=row.updated_at,
|
||||
completed_at=row.completed_at,
|
||||
result=row.result,
|
||||
error=error,
|
||||
)
|
||||
|
||||
|
||||
class QueryQueue:
|
||||
"""账号只读查询单队列
|
||||
|
||||
生命周期与 TaskQueue 一致:由 GatewayContainer 持有,共用同一个 GatewayDB。
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
db: GatewayDB,
|
||||
*,
|
||||
lease_ttl_seconds: int,
|
||||
query_ttl_seconds: int,
|
||||
max_attempts: int,
|
||||
retention_seconds: int,
|
||||
):
|
||||
self._db = db
|
||||
self._lease_ttl = lease_ttl_seconds
|
||||
self._query_ttl = query_ttl_seconds
|
||||
self._max_attempts = max_attempts
|
||||
self._retention = retention_seconds
|
||||
self._lock = asyncio.Lock()
|
||||
self._cond = asyncio.Condition(self._lock)
|
||||
|
||||
# ---- 提交 ----
|
||||
|
||||
async def submit(
|
||||
self, *, query_id: str | None, site: str, kind: str, params: dict
|
||||
) -> SubmitQueryData:
|
||||
"""上游提交查询单。query_id 缺省时服务端生成;重复提交幂等"""
|
||||
qid = query_id or _generate_query_id()
|
||||
now = to_iso(utcnow())
|
||||
row = QueryRow(
|
||||
query_id=qid,
|
||||
site=site,
|
||||
kind=kind,
|
||||
params_json=json.dumps(params, ensure_ascii=False),
|
||||
status=QueryStatus.QUEUED.value,
|
||||
lease_owner=None,
|
||||
lease_expires_at=None,
|
||||
attempts=0,
|
||||
result_json=None,
|
||||
error_code=None,
|
||||
error_message=None,
|
||||
created_at=now,
|
||||
updated_at=now,
|
||||
completed_at=None,
|
||||
)
|
||||
async with self._lock:
|
||||
created = await self._db.insert_query(row)
|
||||
if created:
|
||||
self._cond.notify_all()
|
||||
else:
|
||||
existing = await self._db.get_query(qid)
|
||||
if existing is None:
|
||||
raise RuntimeError(f"查询单 {qid} 既未新建也无法读取,状态异常")
|
||||
row = existing
|
||||
return SubmitQueryData(
|
||||
query_id=row.query_id, status=QueryStatus(row.status), created=created
|
||||
)
|
||||
|
||||
# ---- 领取 ----
|
||||
|
||||
async def lease(
|
||||
self, *, worker_id: str, wait: int, site: str | None, max_wait: int
|
||||
) -> QueryLeaseData | None:
|
||||
"""本地长轮询领取一张查询单
|
||||
|
||||
与下单 lease 的三点不同:
|
||||
1. **没有全局并发度 1 的闸门**——只读操作允许多张单同时在飞;真正的账号级
|
||||
串行由本地 SiteInteractor 那把锁保证,网关不必也不该重复实现。
|
||||
2. **不刷 workers.last_seen_at**。心跳的语义是「下单 worker 还活着」,
|
||||
如果查询循环活着就刷心跳,下单主循环挂了也看不出来,告警会失真。
|
||||
3. 领走前先 sweep:把超时的重投回 queued、过期的置 expired。
|
||||
"""
|
||||
effective_wait = max(0, min(wait, max_wait))
|
||||
deadline = utcnow() + timedelta(seconds=effective_wait)
|
||||
|
||||
async with self._lock:
|
||||
await self._sweep_locked()
|
||||
|
||||
while True:
|
||||
row = await self._db.pick_queued_query(site)
|
||||
if row is not None:
|
||||
return await self._lease_to_locked(row, worker_id)
|
||||
|
||||
remaining = (deadline - utcnow()).total_seconds()
|
||||
if remaining <= 0:
|
||||
return None
|
||||
try:
|
||||
await asyncio.wait_for(self._cond.wait(), timeout=remaining)
|
||||
except asyncio.TimeoutError:
|
||||
return None
|
||||
|
||||
async def _lease_to_locked(self, row: QueryRow, worker_id: str) -> QueryLeaseData:
|
||||
"""把一张 queued 查询单租给 worker。调用方必须持有 self._lock"""
|
||||
now = utcnow()
|
||||
expires = to_iso(now + timedelta(seconds=self._lease_ttl))
|
||||
attempt = row.attempts + 1
|
||||
await self._db.update_query(
|
||||
row.query_id,
|
||||
status=QueryStatus.LEASED.value,
|
||||
lease_owner=worker_id,
|
||||
lease_expires_at=expires,
|
||||
attempts=attempt,
|
||||
updated_at=to_iso(now),
|
||||
)
|
||||
return QueryLeaseData(
|
||||
query_id=row.query_id,
|
||||
site=row.site,
|
||||
kind=row.kind, # type: ignore[arg-type]
|
||||
params=row.params,
|
||||
lease_expires_at=expires,
|
||||
attempt=attempt,
|
||||
)
|
||||
|
||||
# ---- 回结果 ----
|
||||
|
||||
async def submit_result(
|
||||
self,
|
||||
query_id: str,
|
||||
*,
|
||||
worker_id: str,
|
||||
success: bool,
|
||||
result: dict | None,
|
||||
error_code: int | None,
|
||||
error_message: str,
|
||||
) -> QueryResultData:
|
||||
"""worker 回报结果
|
||||
|
||||
校验 `worker_id == lease_owner` 且查询单未终结,否则 6006——最典型的场景是
|
||||
执行超时、单子已被重投给下一轮,迟到的结果必须丢弃,不能覆盖新结果。
|
||||
"""
|
||||
async with self._lock:
|
||||
row = await self._db.get_query(query_id)
|
||||
if row is None:
|
||||
raise QueryNotFoundError(query_id)
|
||||
status = QueryStatus(row.status)
|
||||
if status in QUERY_TERMINAL_STATUSES:
|
||||
raise QueryLeaseInvalidError(
|
||||
f"查询单 {query_id} 已是终态 {status.value},不再接受结果"
|
||||
)
|
||||
if row.lease_owner != worker_id:
|
||||
raise QueryLeaseInvalidError(
|
||||
f"worker {worker_id} 不是查询单 {query_id} 的租约持有者"
|
||||
f"(当前持有者 {row.lease_owner})"
|
||||
)
|
||||
|
||||
now = to_iso(utcnow())
|
||||
if success:
|
||||
await self._db.update_query(
|
||||
query_id,
|
||||
status=QueryStatus.SUCCEEDED.value,
|
||||
result_json=json.dumps(result or {}, ensure_ascii=False),
|
||||
updated_at=now,
|
||||
completed_at=now,
|
||||
clear_lease=True,
|
||||
)
|
||||
final = QueryStatus.SUCCEEDED
|
||||
else:
|
||||
await self._db.update_query(
|
||||
query_id,
|
||||
status=QueryStatus.FAILED.value,
|
||||
error_code=error_code,
|
||||
error_message=error_message or "worker 未给出失败原因",
|
||||
updated_at=now,
|
||||
completed_at=now,
|
||||
clear_lease=True,
|
||||
)
|
||||
final = QueryStatus.FAILED
|
||||
logger.warning(
|
||||
"查询单失败:query_id=%s kind=%s code=%s msg=%s",
|
||||
query_id, row.kind, error_code, error_message,
|
||||
)
|
||||
return QueryResultData(query_id=query_id, status=final)
|
||||
|
||||
# ---- 清扫:重投 / 过期 / 保留期 ----
|
||||
|
||||
async def sweep(self) -> int:
|
||||
"""扫一轮:租约超时重投、整体超时置 expired、过保留期的终态单清理
|
||||
|
||||
返回本轮发生状态变更的查询单条数(删除的不计入)。
|
||||
"""
|
||||
async with self._lock:
|
||||
changed = await self._sweep_locked()
|
||||
purged = await self._db.delete_queries_completed_before(
|
||||
to_iso(utcnow() - timedelta(seconds=self._retention))
|
||||
)
|
||||
if purged:
|
||||
logger.info("清理了 %s 张过保留期的查询单", purged)
|
||||
return changed
|
||||
|
||||
async def _sweep_locked(self) -> int:
|
||||
"""锁内清扫。调用方必须持有 self._lock"""
|
||||
now = utcnow()
|
||||
changed = 0
|
||||
rows = await self._db.list_queries_in_statuses(
|
||||
(QueryStatus.QUEUED.value, QueryStatus.LEASED.value)
|
||||
)
|
||||
for row in rows:
|
||||
# 整体 TTL 优先:单子已经没有意义了,再重投也只是白跑一趟
|
||||
if parse_iso(row.created_at) + timedelta(seconds=self._query_ttl) < now:
|
||||
await self._db.update_query(
|
||||
row.query_id,
|
||||
status=QueryStatus.EXPIRED.value,
|
||||
error_code=None,
|
||||
error_message=(
|
||||
f"查询单超过 {self._query_ttl} 秒仍未完成"
|
||||
f"(attempts={row.attempts},多半是本地 worker 不在线)"
|
||||
),
|
||||
updated_at=to_iso(now),
|
||||
completed_at=to_iso(now),
|
||||
clear_lease=True,
|
||||
)
|
||||
changed += 1
|
||||
logger.warning("查询单过期:query_id=%s kind=%s", row.query_id, row.kind)
|
||||
continue
|
||||
|
||||
if row.status != QueryStatus.LEASED.value:
|
||||
continue
|
||||
if not row.lease_expires_at or parse_iso(row.lease_expires_at) >= now:
|
||||
continue
|
||||
|
||||
if row.attempts >= self._max_attempts:
|
||||
await self._db.update_query(
|
||||
row.query_id,
|
||||
status=QueryStatus.FAILED.value,
|
||||
error_message=(
|
||||
f"已被领取 {row.attempts} 次仍未回结果,不再重投"
|
||||
"(可能是站点页面卡住或本地 worker 反复崩溃)"
|
||||
),
|
||||
updated_at=to_iso(now),
|
||||
completed_at=to_iso(now),
|
||||
clear_lease=True,
|
||||
)
|
||||
logger.warning(
|
||||
"查询单重投次数用尽,置 failed:query_id=%s attempts=%s",
|
||||
row.query_id, row.attempts,
|
||||
)
|
||||
else:
|
||||
# 只读操作可以安全重投——这正是查询与下单任务最大的语义差别
|
||||
await self._db.update_query(
|
||||
row.query_id,
|
||||
status=QueryStatus.QUEUED.value,
|
||||
updated_at=to_iso(now),
|
||||
clear_lease=True,
|
||||
)
|
||||
logger.info(
|
||||
"查询单租约过期,重投回队列:query_id=%s attempts=%s",
|
||||
row.query_id, row.attempts,
|
||||
)
|
||||
changed += 1
|
||||
|
||||
if changed:
|
||||
self._cond.notify_all()
|
||||
return changed
|
||||
|
||||
# ---- 查询 ----
|
||||
|
||||
async def get_detail(self, query_id: str) -> QueryDetail:
|
||||
row = await self._db.get_query(query_id)
|
||||
if row is None:
|
||||
raise QueryNotFoundError(query_id)
|
||||
return _row_to_detail(row)
|
||||
|
||||
async def list_queries(
|
||||
self, *, status: str | None, kind: str | None, limit: int, offset: int
|
||||
) -> QueryListData:
|
||||
rows, total = await self._db.list_queries(
|
||||
status=status, kind=kind, limit=limit, offset=offset
|
||||
)
|
||||
return QueryListData(
|
||||
items=[_row_to_detail(r) for r in rows], total=total, limit=limit, offset=offset
|
||||
)
|
||||
|
||||
async def queued_count(self) -> int:
|
||||
"""等待领取的查询单数量(/health 用)"""
|
||||
return await self._db.count_queries_by_status(QueryStatus.QUEUED.value)
|
||||
Reference in New Issue
Block a user