diff --git a/.env.example b/.env.example index 3673d64..2ac80de 100644 --- a/.env.example +++ b/.env.example @@ -105,7 +105,8 @@ RAKUTEN_RELOGIN_TIMEOUT_SECONDS=300 # ---- 以下仅下单任务网关使用 ---- # 任务队列 SQLite 文件路径(相对项目根目录)。务必放在持久化卷上,丢了等于 -# 丢了一批下单任务。详见 docs/order-gateway.md。 +# 丢了一批下单任务。查询通道(account_queries 表)在同一份库文件里。 +# 详见 docs/order-gateway.md。 RAKUTEN_GATEWAY_DB_PATH=data/gateway.db # 任务租约 TTL(秒)。worker 领取后必须在此时间内首次 report 或 renew,否则 # 任务被置为 stale(**绝不自动重投**,需要人工 reclaim)。 @@ -115,6 +116,24 @@ RAKUTEN_LEASE_MAX_WAIT_SECONDS=60 # worker 心跳超时阈值(秒)。超过即视为失联,/health 报 degraded。 RAKUTEN_WORKER_OFFLINE_ALERT_SECONDS=300 +# ---- 账号只读查询通道(网关与本地 worker 共享),详见 docs/order-gateway.md §11 ---- +# 查询单租约 TTL(秒)。worker 领走后需在此期间回结果,否则自动重投回丢列。 +# 与下单任务的关键区别:**只读操作可以安全重投**。 +RAKUTEN_QUERY_LEASE_TTL_SECONDS=180 +# 查询单整体存活上限(秒)。从创建算起,超过仍未完成即置 expired(多半本地 worker 不在线)。 +RAKUTEN_QUERY_TTL_SECONDS=900 +# 同一张查询单最多被领取几次。重投累计到此仍无结果即置 failed。 +RAKUTEN_QUERY_MAX_ATTEMPTS=3 +# 终态查询单保留时长(秒)。结果里带站点原始 JSON,超期由 sweep 清理。 +RAKUTEN_QUERY_RETENTION_SECONDS=604800 +# worker 侧单次站点读取上限(秒)。查询在等账号锁(下单正跑)时也会超时失败, +# 上游重发即可——只读,重发没有副作用。 +RAKUTEN_ACCOUNT_QUERY_TIMEOUT_SECONDS=120 +# order_list 查询缺省翻页上限。上游可用 params.max_pages 覆盖(1..20)。 +RAKUTEN_ACCOUNT_QUERY_DEFAULT_MAX_PAGES=3 +# 回报结果 JSON 体积上限(字节)。超限先丢站点原始 JSON,仍超限判失败让上游缩小窗口。 +RAKUTEN_QUERY_RESULT_MAX_BYTES=1048576 + # ---- 以下仅交易服务内的下单 worker 使用 ---- # 网关 URL。**留空则不启动 worker**,交易服务只跑登录态接口。 # 部署形态:本地机(NAT 后无公网入口)通过出站长轮询从这里领任务。 diff --git a/README.md b/README.md index 3efc1e3..8ba2e39 100644 --- a/README.md +++ b/README.md @@ -146,8 +146,15 @@ PC UA 在搜索页、详情页、店铺页上都能拿到完整模板。因此 | `POST /api/orders/{id}/reclaim` | 把 stale 任务重新租给 worker(**绝不自动重投**) | | `GET /api/orders/{id}` | 任务详情 + 完整状态历史 | | `GET /api/orders` | 任务列表(运维与上游对账用) | +| `POST /api/account/queries` | 提交**账号只读查询**(订单列表 / 单笔详情,幂等) | +| `GET /api/account/queries/lease` | 本地 worker 领取查询单(可安全重投,与下单任务相反) | +| `POST /api/account/queries/{id}/result` | 本地回报查询结果(站点真实订单数据) | +| `GET /api/account/queries/{id}` | 取查询单状态与结果 | +| `GET /api/account/queries` | 查询单列表(运维排查) | + +网关的契约与状态机详见 [docs/order-gateway.md](docs/order-gateway.md) +(账号只读查询通道见该文档 §11,用于「从已登录账号拉取真实订单」)。 -网关的契约与状态机详见 [docs/order-gateway.md](docs/order-gateway.md)。 三个服务共用同一个 Bearer Token,错误码表也是同一份。 启动后分别在 `http://127.0.0.1:31107/docs`、`:31108/docs`、`:31109/docs` 查看 OpenAPI 文档。 @@ -635,12 +642,16 @@ trading 加购时的字段选择策略:多规格挑第一个非售罄的 varia | 6002 | 租约无效:不是持有者、已过期或任务已终结(网关) | 409 | | 6003 | 任务状态不允许该操作(如对已终结任务 reclaim)(网关) | 409 | | 6004 | 已有任务在执行中,本次不发放(正常返回空,仅诊断用) | 200 | +| 6005 | 查询单不存在(网关,账号只读查询通道) | 404 | +| 6006 | 查询单租约无效(网关,账号只读查询通道) | 409 | 错误码在两站、三个服务之间通用。ラクマ 链路不会出现 `3002`(无反爬拦截行为) 与 `4002`(无子站跳转);`5xxx` 只会来自交易服务——抓取服务全程匿名,不会有登录态问题。 `5001` 与 `5004` 都标记为不可重试:前者要人工重新登录,后者要调用方改入参。 -`6xxx` 只会来自网关,全部标记为不可重试——任务编排侧重试无意义,部分场景 -(如租约过期)重试可能变成重复下单。 +`6xxx` 只会来自网关:`6001`–`6004`(下单任务通道)全部标记为不可重试——任务编排侧 +重试无意义,部分场景(如租约过期)重试可能变成重复下单;`6005`–`6006` +(账号只读查询通道)**可以重试**——查询只读,重发没有副作用 +(见 [docs/order-gateway.md §11](docs/order-gateway.md#11-账号只读查询通道))。 ## 常用环境变量 @@ -651,6 +662,7 @@ trading 加购时的字段选择策略:多规格挑第一个非售罄的 varia - 抓取服务:`RAKUTEN_APP_HOST`、`RAKUTEN_APP_PORT`(默认 `31107`)、`RAKUTEN_APP_ENV` - 交易服务:`RAKUTEN_TRADING_HOST`、`RAKUTEN_TRADING_PORT`(默认 `31108`)、`RAKUTEN_AUTH_STATE_DIR`(默认 `.auth`)、`RAKUTEN_ORDER_MAX_TOTAL_YEN`(默认 `30000`) - 下单任务网关:`RAKUTEN_GATEWAY_HOST`、`RAKUTEN_GATEWAY_PORT`(默认 `31109`)、`RAKUTEN_GATEWAY_DB_PATH`(默认 `data/gateway.db`)、`RAKUTEN_LEASE_TTL_SECONDS`(默认 `300`)、`RAKUTEN_WORKER_OFFLINE_ALERT_SECONDS`(默认 `300`) +- 账号只读查询通道(网关 + 本地 worker 两侧共用):`RAKUTEN_QUERY_LEASE_TTL_SECONDS`(默认 `180`)、`RAKUTEN_QUERY_TTL_SECONDS`(默认 `900`)、`RAKUTEN_QUERY_MAX_ATTEMPTS`(默认 `3`)、`RAKUTEN_QUERY_RETENTION_SECONDS`(默认 `604800`)、`RAKUTEN_ACCOUNT_QUERY_TIMEOUT_SECONDS`(默认 `120`)、`RAKUTEN_ACCOUNT_QUERY_DEFAULT_MAX_PAGES`(默认 `3`)、`RAKUTEN_QUERY_RESULT_MAX_BYTES`(默认 `1048576`) - 本地下单 worker(在交易服务内,按 `RAKUTEN_ORDER_GATEWAY_URL` 是否配置决定是否启动):`RAKUTEN_ORDER_GATEWAY_URL`、`RAKUTEN_WORKER_ID`、`RAKUTEN_TRADING_DB_PATH`(默认 `data/trading.db`)、`RAKUTEN_EVIDENCE_DIR`(默认 `data/evidence`)、`RAKUTEN_SCRAPER_BASE_URL` - 鉴权:`RAKUTEN_BEARER_TOKEN`(三个服务共用) - 抓取:`RAKUTEN_MAX_SITE_CONCURRENCY`(默认 `8`,两站各自独立计数)、`RAKUTEN_HTTP_MAX_ATTEMPTS`(默认 `3`)、`RAKUTEN_SESSION_TTL_SECONDS`(默认 `1800`,仅乐天) diff --git a/app/gateway/api/routes/health.py b/app/gateway/api/routes/health.py index da9b321..7478b3e 100644 --- a/app/gateway/api/routes/health.py +++ b/app/gateway/api/routes/health.py @@ -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 ) diff --git a/app/gateway/api/routes/queries.py b/app/gateway/api/routes/queries.py new file mode 100644 index 0000000..b1ea47f --- /dev/null +++ b/app/gateway/api/routes/queries.py @@ -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) diff --git a/app/gateway/container.py b/app/gateway/container.py index edb65ac..7cd0e45 100644 --- a/app/gateway/container.py +++ b/app/gateway/container.py @@ -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 diff --git a/app/gateway/db.py b/app/gateway/db.py index 25163b8..1e8c5a5 100644 --- a/app/gateway/db.py +++ b/app/gateway/db.py @@ -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 diff --git a/app/gateway/main.py b/app/gateway/main.py index 7b2d680..d262379 100644 --- a/app/gateway/main.py +++ b/app/gateway/main.py @@ -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 diff --git a/app/gateway/models.py b/app/gateway/models.py index 8bcb3bb..c8a5958 100644 --- a/app/gateway/models.py +++ b/app/gateway/models.py @@ -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) diff --git a/app/gateway/query_queue.py b/app/gateway/query_queue.py new file mode 100644 index 0000000..d93ac8a --- /dev/null +++ b/app/gateway/query_queue.py @@ -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) diff --git a/app/shared/config.py b/app/shared/config.py index bb10064..5523b2a 100644 --- a/app/shared/config.py +++ b/app/shared/config.py @@ -151,6 +151,29 @@ class Settings(BaseSettings): # 正常 worker 每 30 秒来一次 lease,超过该阈值未来 lease 即视为异常。 worker_offline_alert_seconds: int = 300 + # ---- 账号只读查询通道(网关 + worker 两侧共用,见 docs/order-gateway.md §11)---- + # 查询单租约 TTL(秒)。worker 领走后必须在此时间内回结果,否则网关把它 + # **重投回 queued**——只读操作重复执行没有副作用,与下单任务的 stale 语义相反。 + query_lease_ttl_seconds: int = 180 + # 查询单整体存活上限(秒)。从创建算起,超过仍未拿到结果即置 expired + # (多半是本地 worker 不在线)。上游据此判断「这张单不用再等了」。 + query_ttl_seconds: int = 900 + # 同一张查询单最多被领取几次。租约过期重投累计到这个次数仍无结果即置 failed, + # 避免某条查询让 worker 反复卡死在同一个页面上。 + query_max_attempts: int = 3 + # 终态查询单的保留时长(秒)。结果里带站点原始 JSON,体积不小,超期由 sweep + # 清掉;默认 7 天,够上游对账与排查。 + query_retention_seconds: int = 7 * 24 * 3600 + # worker 侧单次站点读取的耗时上限(秒)。查询与下单共用同一把账号锁,正在 + # 跑的下单会让查询排队等待,这个上限保证查询不会一直挂着——超时按失败回报, + # 上游重发即可(只读,重发无副作用)。 + account_query_timeout_seconds: int = 120 + # 订单列表查询缺省翻页上限。上游可在 params.max_pages 覆盖(1..20)。 + account_query_default_max_pages: int = 3 + # 回报给网关的结果 JSON 体积上限(字节)。超过时丢掉站点原始 JSON、只回 + # 规范化字段并标记 raw_omitted——避免把几 MB 的页面状态塞进网关 SQLite。 + query_result_max_bytes: int = 1024 * 1024 + # ---- 本地下单 worker(仅交易服务内的 worker 子模块使用)---- # 网关 URL。**留空则不启动 worker**,交易服务只跑登录态接口。 # 部署形态:本地机(NAT 后无公网入口)通过出站长轮询领任务,详见 diff --git a/app/shared/errors.py b/app/shared/errors.py index aea793d..821e5ea 100644 --- a/app/shared/errors.py +++ b/app/shared/errors.py @@ -241,3 +241,39 @@ class InvalidTaskStateError(AppError): retryable=False, status_code=409, ) + + +# ---- 账号只读查询通道(仅网关进程使用,见 docs/order-gateway.md §11)---- +# 与下单任务不同:查询是只读的,**可以安全重投**,因此这两类错误标 retryable=True, +# 上游重发一张查询单不会有任何副作用。 + + +class QueryNotFoundError(AppError): + """查询单不存在(或已过保留期被清理)""" + + def __init__(self, query_id: str): + super().__init__( + message=f"查询单不存在:{query_id}", + code="QUERY_NOT_FOUND", + err_code=6005, + retryable=False, + status_code=404, + ) + self.query_id = query_id + + +class QueryLeaseInvalidError(AppError): + """查询单租约无效:不是持有者、已被重投给别人或已终结 + + 最常见的触发场景是 worker 执行超时、查询单被 sweep 重投后,原 worker 才姗姗 + 来迟地回结果——这时结果必须被拒绝,否则会覆盖掉新一轮的执行结果。 + """ + + def __init__(self, message: str = "查询单租约无效"): + super().__init__( + message=message, + code="QUERY_LEASE_INVALID", + err_code=6006, + retryable=True, + status_code=409, + ) diff --git a/app/shared/task_state.py b/app/shared/task_state.py index e46695a..74f6308 100644 --- a/app/shared/task_state.py +++ b/app/shared/task_state.py @@ -49,6 +49,43 @@ LEASABLE_STATUSES: frozenset[TaskStatus] = frozenset({TaskStatus.QUEUED}) RECLAIMABLE_STATUSES: frozenset[TaskStatus] = frozenset({TaskStatus.STALE}) +class QueryStatus(StrEnum): + """账号只读查询单状态(网关权威,见 docs/order-gateway.md §11) + + 与 TaskStatus 刻意分成两套词汇表,因为安全约束正好相反: + + - 下单是不可逆写操作 → 租约过期只能置 stale 等人工 reclaim,绝不自动重投 + - 查询是只读操作 → 租约过期直接回 queued 自动重投,重复执行没有副作用 + + 没有 running:查询是「领走 → 一次性回结果」,中间没有需要单独表达的进行态, + 也不需要续租(超时就重投)。 + """ + + QUEUED = "queued" # 已入队,等待 worker 领取 + LEASED = "leased" # 已被 worker 领走,等待回结果 + SUCCEEDED = "succeeded" # 终态:拿到结果 + FAILED = "failed" # 终态:worker 明确失败,或重投次数用尽 + EXPIRED = "expired" # 终态:超过查询单 TTL 仍未完成(多半是 worker 不在线) + + +# 查询单终态集合:到达后 result 一律报 6006,sweep 也不再动它 +QUERY_TERMINAL_STATUSES: frozenset[QueryStatus] = frozenset( + {QueryStatus.SUCCEEDED, QueryStatus.FAILED, QueryStatus.EXPIRED} +) + + +class AccountQueryKind(StrEnum): + """账号只读查询的种类 + + 放在 shared 是因为它是网关与 worker 之间的 HTTP 契约的一部分:网关不解释 + `params` 的内容(与 intent 同样原样透传),但**校验 kind 合法**——不然上游 + 拼错一个字符,要等 worker 领走、执行、回报失败才知道,比当场 422 差得多。 + """ + + ORDER_LIST = "order_list" # 订单列表(可带时间窗与翻页上限) + ORDER_DETAIL = "order_detail" # 单笔订单详情(按注文番号) + + class OrderState(StrEnum): """订单状态(本地权威,网关只存镜像) diff --git a/app/trading/container.py b/app/trading/container.py index dea92e3..964811b 100644 --- a/app/trading/container.py +++ b/app/trading/container.py @@ -41,3 +41,6 @@ class TradingContainer: worker_local_db: object | None = None # app.trading.worker.local_db.LocalDB worker_evidence: object | None = None # app.trading.worker.evidence.EvidenceStore worker_runner: object | None = None # app.trading.worker.runner.WorkerRunner + # 账号只读查询的第二条循环(docs/order-gateway.md §11)。与下单 worker 同生共死, + # 但跑在各自的 asyncio 任务里:下单一跑几分钟,查询不该排在它后面。 + query_runner: object | None = None # app.trading.worker.query_runner.QueryRunner diff --git a/app/trading/main.py b/app/trading/main.py index b9fd480..c5daf15 100644 --- a/app/trading/main.py +++ b/app/trading/main.py @@ -13,7 +13,8 @@ SiteInteractor 始终启动:HTTP /api/auth/* 与 /api/cart/* 都要靠它持账号 cookie 做站点交互(加购、校验、清空、删除)。worker 启动条件:RAKUTEN_ORDER_GATEWAY_URL 非空。留空时不构造 worker 子组件(本地 DB / 证据目录 / 出站客户端都不创建), -但 SiteInteractor 与 Playwright 仍会启动。 +但 SiteInteractor 与 Playwright 仍会启动。配了网关 URL 时会起**两条**常驻循环: +下单 worker 与账号只读查询 worker(后者见 docs/order-gateway.md §11)。 """ from __future__ import annotations @@ -56,6 +57,7 @@ def build_container() -> TradingContainer: from app.trading.worker.client import GatewayClient from app.trading.worker.evidence import EvidenceStore from app.trading.worker.local_db import LocalDB + from app.trading.worker.query_runner import QueryRunner from app.trading.worker.runner import WorkerRunner client = GatewayClient( @@ -77,6 +79,11 @@ def build_container() -> TradingContainer: container.worker_local_db = local_db container.worker_evidence = evidence container.worker_runner = runner + # 只读查询循环与下单 worker 共用同一个网关客户端与 SiteInteractor: + # 前者是同一份连接池,后者那把锁正是账号级串行的保证。 + container.query_runner = QueryRunner( + settings=settings, gateway_client=client, site=container.site + ) return container @@ -139,6 +146,7 @@ async def lifespan(app: FastAPI): logger.info("未开启 RAKUTEN_AUTO_LOGIN_ON_START,启动时不自动登录") worker_task: asyncio.Task | None = None + query_task: asyncio.Task | None = None if container.worker_runner is not None: # 顺序:本地 DB 在 worker 写证据前先就绪;SiteInteractor 已在外层启动 assert container.worker_local_db is not None @@ -152,6 +160,13 @@ async def lifespan(app: FastAPI): container.settings.worker_id_effective, container.settings.order_gateway_url, ) + # 第二条循环:账号只读查询。独立于下单主循环,否则一笔下单跑几分钟, + # 期间上游一句「账号里那笔订单什么状态」就得干等(docs/order-gateway.md §11)。 + assert container.query_runner is not None + query_task = asyncio.create_task( + container.query_runner.run(), name="trading-query-worker" + ) + logger.info("账号只读查询 worker 已启动") else: logger.info( "未配置 RAKUTEN_ORDER_GATEWAY_URL,下单 worker 不启动(仅登录态与购物车接口)" @@ -177,6 +192,15 @@ async def lifespan(app: FastAPI): # 付款后监控是独立于主循环的后台任务(见 runner.py::_spawn_monitor), # 取消 worker_task 不会连带取消它们,必须在关 SiteInteractor 前单独收尾。 await container.worker_runner.cancel_monitors() # type: ignore[union-attr] + if query_task is not None: + # 与下单 worker 同样:先 stop 再 cancel,正在跑的那次读取直接放弃 + # (只读,放弃没有副作用;网关侧那张单会超时重投)。 + container.query_runner.stop() # type: ignore[union-attr] + query_task.cancel() + try: + await query_task + except asyncio.CancelledError: + pass if container.site is not None: await container.site.close() # type: ignore[union-attr] if container.worker_local_db is not None: diff --git a/app/trading/worker/client.py b/app/trading/worker/client.py index 1c6ade9..e48880a 100644 --- a/app/trading/worker/client.py +++ b/app/trading/worker/client.py @@ -20,7 +20,7 @@ from app.shared.config import Settings from app.shared.errors import AppError from app.shared.proxy import httpx_client_options from app.shared.task_state import OrderState, TaskStatus -from app.trading.worker.models import LeaseTask +from app.trading.worker.models import LeaseTask, QueryTask logger = logging.getLogger(__name__) @@ -158,3 +158,55 @@ class GatewayClient: lease_count=data.get("lease_count", 1), known_state=data.get("known_state"), ) + + # ---- 账号只读查询通道(docs/order-gateway.md §11)---- + + async def lease_query(self, worker_id: str, *, wait: int = 30) -> QueryTask | None: + """长轮询领取一张只读查询单。无单可领时返回 None + + 与 lease() 走的是两条互不干扰的通道:下单任务在跑的时候,查询照样能领到 + (账号级串行由 SiteInteractor 那把锁保证,不靠网关排队)。 + """ + body = await self._request( + "GET", + "/api/account/queries/lease", + params={"worker_id": worker_id, "wait": wait}, + ) + data = body.get("data") + if not data: + return None + return QueryTask( + query_id=data["query_id"], + site=data["site"], + kind=data["kind"], + params=data.get("params") or {}, + lease_expires_at=data.get("lease_expires_at", ""), + attempt=data.get("attempt", 1), + ) + + async def report_query_result( + self, + query_id: str, + worker_id: str, + *, + success: bool, + result: dict[str, Any] | None = None, + error_code: int | None = None, + error_message: str = "", + ) -> dict[str, Any]: + """回报查询结果,返回网关响应里的 data 字段 + + 网关可能回 6006(租约已被重投给下一轮)——那说明本次结果已经作废, + 调用方按「丢弃」处理即可,不要重发。 + """ + payload: dict[str, Any] = { + "worker_id": worker_id, + "success": success, + "result": result, + "error_code": error_code, + "error_message": error_message, + } + body = await self._request( + "POST", f"/api/account/queries/{query_id}/result", json=payload + ) + return body["data"] diff --git a/app/trading/worker/models.py b/app/trading/worker/models.py index b6ea96d..b4b6018 100644 --- a/app/trading/worker/models.py +++ b/app/trading/worker/models.py @@ -26,3 +26,20 @@ class LeaseTask: lease_expires_at: str = "" lease_count: int = 1 known_state: str | None = None + + +@dataclass(slots=True) +class QueryTask: + """worker 从网关领到的账号只读查询单 + + 对应网关 GET /api/account/queries/lease 的 QueryLeaseData。没有 known_state + 之类的字段——查询是一次性只读动作,`attempt > 1` 只说明上一轮超时被重投了, + worker 不需要因此改变行为(重跑安全,这正是它与下单任务的根本差别)。 + """ + + query_id: str + site: str + kind: str + params: dict[str, Any] = field(default_factory=dict) + lease_expires_at: str = "" + attempt: int = 1 diff --git a/app/trading/worker/query_runner.py b/app/trading/worker/query_runner.py new file mode 100644 index 0000000..6e7d377 --- /dev/null +++ b/app/trading/worker/query_runner.py @@ -0,0 +1,307 @@ +"""账号只读查询 worker:第二条常驻循环,与下单主循环并行 + +为什么单独一条循环而不是塞进 `runner.WorkerRunner`(docs/order-gateway.md §11): +下单主循环执行一笔任务要几分钟(加购 → 确认页 → 提交 → 付款),期间它整个人 +都占在那笔任务上;查询要是排在同一条循环里,上游问一句「账号里那笔订单现在 +什么状态」就得等下单跑完。两条循环各自长轮询各自的通道,互不阻塞。 + +那「同一账号并发操作」的硬约束靠什么保证?靠 `SiteInteractor` 里那把 +`asyncio.Lock`——所有站点交互(下单写操作与这里的只读操作)都要先拿它,天然串行。 +网关侧因此**不需要**为查询再实现一遍「全局并发度 1」。 + +代价是:正在下单时,查询会卡在那把锁上等。所以每次执行都套 +`account_query_timeout_seconds` 超时,超时按失败回报,上游重发即可 +(只读,重发没有副作用)——总比让查询单一直挂到租约过期强。 +""" +from __future__ import annotations + +import asyncio +import json +import logging +from datetime import datetime, timezone +from typing import TYPE_CHECKING, Any + +from app.shared.errors import AppError, InvalidRequestError +from app.shared.task_state import AccountQueryKind +from app.trading.worker.client import GatewayClient +from app.trading.worker.models import QueryTask +from app.trading.worker.site_interact import SiteInteractor + +if TYPE_CHECKING: + from app.shared.config import Settings + +logger = logging.getLogger(__name__) + +# 没给 since 时的窗口下界。用 epoch 而不是 None,是为了让翻页终止条件 +# (`order_date < since`)与「不设下界」共用同一条代码路径,不多一个分支。 +_EPOCH = datetime(1970, 1, 1, tzinfo=timezone.utc) + +# 超时时回报的错误码:沿用 2002(资源繁忙 / 等槽位超时),这里等的是账号锁。 +_TIMEOUT_ERR_CODE = 2002 +# 未预期异常回报的错误码,与 HTTP 层兜底处理器保持一致 +_UNEXPECTED_ERR_CODE = 1500 + + +def _parse_since(raw: Any) -> datetime: + """解析 params.since。缺省用 epoch(不设下界);给了但解析不出直接报错,不猜""" + if raw in (None, ""): + return _EPOCH + if not isinstance(raw, str): + raise InvalidRequestError(f"params.since 必须是 ISO8601 字符串,收到 {type(raw).__name__}") + try: + parsed = datetime.fromisoformat(raw.replace("Z", "+00:00")) + except ValueError as exc: + raise InvalidRequestError(f"params.since 不是合法的 ISO8601 时间:{raw}") from exc + # 站点返回的下单时间带时区,比较两侧必须都 aware,否则直接 TypeError + return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc) + + +def _parse_max_pages(raw: Any, default: int) -> int: + if raw in (None, ""): + return default + if not isinstance(raw, int) or isinstance(raw, bool): + raise InvalidRequestError(f"params.max_pages 必须是整数,收到 {raw!r}") + if raw < 1: + raise InvalidRequestError(f"params.max_pages 必须 ≥ 1,收到 {raw}") + return raw + + +class QueryRunner: + """账号只读查询的常驻循环 + + 只依赖网关客户端与 SiteInteractor:查询不落证据、不写本地 SQLite——它没有 + 「执行事实」需要留痕,读到什么就回什么。 + """ + + def __init__( + self, + *, + settings: "Settings", + gateway_client: GatewayClient, + site: SiteInteractor, + ): + self._settings = settings + self._gateway = gateway_client + self._site = site + self._running = False + + def stop(self) -> None: + self._running = False + + @property + def worker_id(self) -> str: + return self._settings.worker_id_effective + + # ---- 主循环 ---- + + async def run(self) -> None: + """常驻循环:长轮询领查询单 → 执行 → 回结果""" + self._running = True + logger.info("账号查询 worker 启动:worker_id=%s", self.worker_id) + while self._running: + try: + query = await self._gateway.lease_query(self.worker_id, wait=30) + except AppError as exc: + logger.warning("查询 lease 失败:%s (err=%s)", exc.message, exc.err_code) + await asyncio.sleep(5) + continue + except Exception: # noqa: BLE001 + logger.exception("查询 lease 异常") + await asyncio.sleep(5) + continue + + if query is None: + continue + + try: + await self.handle(query) + except Exception: # noqa: BLE001 + # 任何未预期的异常都不应让这条循环退出 + logger.exception("处理查询单异常:query_id=%s", query.query_id) + + # ---- 单张查询单 ---- + + async def handle(self, query: QueryTask) -> None: + """执行一张查询单并回报结果。任何失败都转成一次「失败回报」,不抛出去""" + logger.info( + "领到查询单:query_id=%s kind=%s attempt=%s", + query.query_id, query.kind, query.attempt, + ) + try: + result = await asyncio.wait_for( + self.execute(query), + timeout=self._settings.account_query_timeout_seconds, + ) + except asyncio.TimeoutError: + logger.warning( + "查询超时(%s 秒,多半是下单任务正占着账号锁):query_id=%s", + self._settings.account_query_timeout_seconds, query.query_id, + ) + await self._report_safe( + query, + success=False, + error_code=_TIMEOUT_ERR_CODE, + error_message=( + f"查询超过 {self._settings.account_query_timeout_seconds} 秒未完成" + "(账号锁被下单任务占用或站点页面无响应),可稍后重发" + ), + ) + return + except AppError as exc: + logger.warning( + "查询失败:query_id=%s code=%s msg=%s", + query.query_id, exc.err_code, exc.message, + ) + await self._report_safe( + query, success=False, error_code=exc.err_code, error_message=exc.message + ) + return + except Exception as exc: # noqa: BLE001 + logger.exception("查询未预期异常:query_id=%s", query.query_id) + await self._report_safe( + query, + success=False, + error_code=_UNEXPECTED_ERR_CODE, + error_message=f"{type(exc).__name__}: {exc}", + ) + return + + payload, oversized = self._enforce_result_size(result) + if oversized is not None: + await self._report_safe( + query, success=False, error_code=1003, error_message=oversized + ) + return + await self._report_safe(query, success=True, result=payload) + + async def execute(self, query: QueryTask) -> dict[str, Any]: + """按 kind 分发到 SiteInteractor 的只读方法,返回要回给上游的结果字典""" + if query.site != "rakuten": + raise InvalidRequestError( + f"site={query.site} 不在交易服务范围内(账号查询仅支持 rakuten)" + ) + + if query.kind == AccountQueryKind.ORDER_LIST.value: + return await self._execute_order_list(query) + if query.kind == AccountQueryKind.ORDER_DETAIL.value: + return await self._execute_order_detail(query) + raise InvalidRequestError(f"未知的查询种类:{query.kind}") + + async def _execute_order_list(self, query: QueryTask) -> dict[str, Any]: + """订单列表:规范化条目 + 站点原始 orderListData(逐页)""" + since = _parse_since(query.params.get("since")) + max_pages = _parse_max_pages( + query.params.get("max_pages"), self._settings.account_query_default_max_pages + ) + + window = await self._site.list_recent_orders(since=since, max_pages=max_pages) + return { + "kind": AccountQueryKind.ORDER_LIST.value, + "since": since.isoformat().replace("+00:00", "Z"), + "max_pages": max_pages, + # False=翻页没能覆盖完窗口(命中页数上限 / 页面结构不对 / 掉登录), + # 此时「列表里没有某笔订单」不等于「账号里没有这笔订单」。 + "window_fully_covered": window.window_fully_covered, + "orders": [ + { + "order_number": entry.order_number, + "order_date": entry.order_date, + "shop_id": entry.shop_id, + "shop_name": entry.shop_name, + "items": [ + { + "item_id": item.item_id, + "item_name": item.item_name, + "item_url": item.item_url, + } + for item in entry.items + ], + } + for entry in window.entries + ], + # 站点 __INITIAL_STATE__.orderListData 原文,按翻页顺序。规范化字段 + # 只挑了下单/核对用得上的那几个,别的字段(金额、配送、状态文案)在这里。 + "raw_pages": window.raw_pages, + } + + async def _execute_order_detail(self, query: QueryTask) -> dict[str, Any]: + """单笔订单详情:已实测的配送阶段 + 未实测结构的页面原始状态""" + order_number = query.params.get("order_number") + if not order_number or not isinstance(order_number, str): + raise InvalidRequestError("params.order_number 必填(站点注文番号)") + + detail = await self._site.fetch_order_detail(order_number) + status = detail.status + return { + "kind": AccountQueryKind.ORDER_DETAIL.value, + "order_number": order_number, + # found=False 不是错误:站点自己说下单后订单要 10 分钟左右才反映到 + # 订单页,刚下的单查不到属正常。 + "found": status.found, + "stage_label": status.stage_label, + "order_state": status.order_state.value if status.order_state else None, + # 详情页的 __INITIAL_STATE__ 结构从未拿到真实样本,本服务不做字段抽取, + # 原样透传由上游自担结构变动风险(见 site_interact.OrderDetailSnapshot)。 + "raw": detail.raw, + "raw_available": detail.raw is not None, + } + + # ---- 结果体积与回报 ---- + + def _enforce_result_size(self, result: dict[str, Any]) -> tuple[dict[str, Any], str | None]: + """体积闸门:超限先丢站点原始 JSON,仍超限则判失败 + + 返回 (要回报的结果, 超限说明)。超限说明非 None 时表示这次要按失败回报—— + 与其把几 MB 页面状态塞进网关 SQLite,不如让上游把窗口调小重来。 + """ + limit = self._settings.query_result_max_bytes + if self._size_of(result) <= limit: + return result, None + + trimmed = dict(result) + trimmed["raw_pages"] = [] + trimmed["raw"] = None + trimmed["raw_omitted"] = f"站点原始 JSON 超过 {limit} 字节上限,已丢弃,仅保留规范化字段" + if self._size_of(trimmed) <= limit: + logger.warning("查询结果超限,已丢弃站点原始 JSON:limit=%s", limit) + return trimmed, None + + return result, ( + f"查询结果超过 {limit} 字节上限,丢弃站点原始 JSON 后仍超限;" + "请缩小查询窗口(减小 params.max_pages 或收紧 params.since)后重试" + ) + + @staticmethod + def _size_of(payload: dict[str, Any]) -> int: + return len(json.dumps(payload, ensure_ascii=False).encode("utf-8")) + + async def _report_safe( + self, + query: QueryTask, + *, + success: bool, + result: dict[str, Any] | None = None, + error_code: int | None = None, + error_message: str = "", + ) -> None: + """回报结果,失败只记日志不抛——这条循环不能因为回报失败退出 + + 网关回 6006 属于正常情况(本次租约已被重投给下一轮,结果作废), + 同样只记日志,不重发。 + """ + try: + await self._gateway.report_query_result( + query.query_id, + self.worker_id, + success=success, + result=result, + error_code=error_code, + error_message=error_message, + ) + except AppError as exc: + logger.warning( + "回报查询结果被网关拒绝:query_id=%s code=%s msg=%s", + query.query_id, exc.err_code, exc.message, + ) + except Exception: # noqa: BLE001 + logger.exception("回报查询结果失败:query_id=%s", query.query_id) diff --git a/app/trading/worker/site_interact.py b/app/trading/worker/site_interact.py index 9f85b6d..04cf853 100644 --- a/app/trading/worker/site_interact.py +++ b/app/trading/worker/site_interact.py @@ -64,6 +64,11 @@ data/evidence/checkout-research-20260811/NOTES.md): 不需要正则抠 DOM。当时账号只有 1 笔订单、1 页,验证了单页解析与 `?page=2` 超出范围返回空列表;多页翻页时「新订单排在前面」的排序假设、以及 分页游标本身,都没有被真实多页数据验证过。 +- fetch_order_detail / list_recent_orders 同时供「账号只读查询通道」使用 + (docs/order-gateway.md §11,上游经网关问「账号里真实的订单长什么样」)。 + 查询返回的规范化字段全部来自上面这两条已实测路径;额外带出的站点原始 JSON 里, + **只有订单列表页的 orderListData 是实测过的结构**,详情页的 __INITIAL_STATE__ + 从未拿到过真实样本,只做「解析得动就原样透传」,本模块不猜它的字段。 **httpx 不能用于带账号的写操作**:Rakuten 对账号操作有 TLS/HTTP2 指纹校验, 同一份 cookie Playwright 能用、httpx 不能。所以本模块全程使用 Playwright @@ -376,11 +381,17 @@ class OrderListEntry: @dataclass(slots=True) class OrderListPage: - """order-list 单页解析结果(_parse_order_list 的返回值)""" + """order-list 单页解析结果(_parse_order_list 的返回值) + + `raw` 是站点 `__INITIAL_STATE__.orderListData` 的原文,供只读查询接口原样 + 透传给上游(见 docs/order-gateway.md §11)——站点比我们的 dataclass 多给的 + 字段(金额、配送、状态文案等)不该在这一层被悄悄丢掉。解析不到时为 None。 + """ entries: list[OrderListEntry] orders_found: int | None = None page_size: int | None = None + raw: dict | None = None @dataclass(slots=True) @@ -389,13 +400,16 @@ class OrderListWindow: window_fully_covered=True 才代表「窗口内的订单已经看全」——可能是因为翻到了 比 since 更早的订单、可能是列表本身翻完了、也可能是 ordersFound 已经对上。 - False 表示翻页在覆盖完窗口前就停了(命中 _ORDER_LIST_MAX_PAGES,或页面结构 + False 表示翻页在覆盖完窗口前就停了(命中 max_pages 上限,或页面结构 解析不出 ordersFound 之类的异常),此时 entries 里「没有匹配」不能当作 「确实没下单」——调用方必须转 unknown,不能默认 NOT_ORDERED。 + + raw_pages 按翻页顺序保存每页的站点原始 orderListData,只读查询接口用。 """ entries: list[OrderListEntry] window_fully_covered: bool + raw_pages: list[dict] = field(default_factory=list) def parse_order_datetime(raw: str | None) -> datetime | None: @@ -416,6 +430,7 @@ class _OrderListAccumulator: total_found: int | None = None is_first_page: bool = True gave_up: bool = False # True=第一页就拿不到结构化数据,翻页无意义,直接放弃 + raw_pages: list[dict] = field(default_factory=list) def _accumulate_order_list_page( @@ -431,21 +446,39 @@ def _accumulate_order_list_page( 结构化列表数据(页面结构变了/act 分支不对/未登录跳转),此时也会停止翻页, 但调用方必须按「没覆盖」处理,不能当成真的翻完了。 """ + raw_pages = acc.raw_pages + [page.raw] if page.raw is not None else acc.raw_pages + if acc.is_first_page and page.orders_found is None: return ( - _OrderListAccumulator(entries=acc.entries, total_found=acc.total_found, is_first_page=False, gave_up=True), + _OrderListAccumulator( + entries=acc.entries, + total_found=acc.total_found, + is_first_page=False, + gave_up=True, + raw_pages=raw_pages, + ), True, ) total_found = acc.total_found if acc.total_found is not None else page.orders_found if not page.entries: return ( - _OrderListAccumulator(entries=acc.entries, total_found=total_found, is_first_page=False), + _OrderListAccumulator( + entries=acc.entries, + total_found=total_found, + is_first_page=False, + raw_pages=raw_pages, + ), True, ) new_entries = acc.entries + page.entries - new_acc = _OrderListAccumulator(entries=new_entries, total_found=total_found, is_first_page=False) + new_acc = _OrderListAccumulator( + entries=new_entries, + total_found=total_found, + is_first_page=False, + raw_pages=raw_pages, + ) oldest_dt = parse_order_datetime(page.entries[-1].order_date) if oldest_dt is not None and oldest_dt < since: @@ -486,6 +519,23 @@ class OrderStatusSnapshot: html: str = "" +@dataclass(slots=True) +class OrderDetailSnapshot: + """订单详情页的一次完整读取结果(fetch_order_detail 的返回值) + + `status` 是已实测的那部分(配送阶段进度条,见 _parse_order_status); + `raw` 是整页 `window.__INITIAL_STATE__` 的原文——**详情页的这份结构从未被 + 真实数据验证过**(2026-08-13 那次实测只验了进度条组件),因此这里刻意 + 只做「能解析成 JSON 就原样带出去」,不写任何字段抽取逻辑。要金额、收货 + 地址、付款方式这些字段的调用方,自己从 raw 里取并自担结构变动风险;等拿到 + 真实详情页样本后再在本模块补规范化解析,不要在没有样本的情况下先猜着写。 + 解析不出(页面没有内联状态、或不是 JSON)时为 None。 + """ + + status: OrderStatusSnapshot + raw: dict | None = None + + class SiteInteractor: """Rakuten 站点交互器:持有 Playwright 浏览器 context,复用账号 cookie @@ -1772,23 +1822,36 @@ class SiteInteractor: async def check_order_status(self, site_order_id: str) -> OrderStatusSnapshot: """付款后监控的单次探测:查一次订单详情页的配送阶段,不循环 + `fetch_order_detail` 的薄封装——监控只关心配送阶段,不需要页面原始状态。 循环轮询(间隔、次数上限、状态变化时 report)由 - `runner.WorkerRunner._monitor_order` 负责——本方法只做一次 - 「导航 + 解析」,站点交互与轮询节奏解耦,也方便离线单测轮询逻辑 - (用桩替换本方法)与解析逻辑(`_parse_order_status`,纯函数)。 - - 与 enter_checkout 不同,本方法开自己的临时 Page 并在返回前关闭—— - submit_order/pay 已经处理完并关闭了 `_checkout_pages` 里留存的会话, - 监控阶段没有需要跨调用复用的页面状态。 - - 执行中掉登录会自动重登一次并重跑(详见 `_read_with_relogin_retry`): - 本方法是纯读,重跑没有副作用;不加这层的话订单页被踢到 SSO 会静默走进 - 解析逻辑、被当成「订单号还没出现」,轮询白转几个小时也看不出原因。 + `runner.WorkerRunner._monitor_order` 负责。 Returns: OrderStatusSnapshot;订单号暂时查不到、或进度条解析不出新阶段都 **不算错误**(详见该 dataclass 文档),由调用方决定是否继续轮询。 + Raises: + NotLoggedInError: 登录态失效且自动重登没能恢复 + OrderOperationError: 订单详情页打开/渲染失败 + """ + return (await self.fetch_order_detail(site_order_id)).status + + async def fetch_order_detail(self, site_order_id: str) -> OrderDetailSnapshot: + """读一次订单详情页:配送阶段 + 页面原始 __INITIAL_STATE__ + + 两个调用方:付款后监控(`check_order_status`,只要配送阶段)与账号只读 + 查询通道的 order_detail(规格 §11,还要 raw 原文)。站点交互与轮询节奏 + 解耦,也方便离线单测轮询逻辑(用桩替换本方法)与解析逻辑 + (`_parse_order_status`,纯函数)。 + + 与 enter_checkout 不同,本方法开自己的临时 Page 并在返回前关闭—— + submit_order/pay 已经处理完并关闭了 `_checkout_pages` 里留存的会话, + 监控阶段没有需要跨调用复用的页面状态。 + + 执行中掉登录会自动重登一次并重跑(详见 `_read_with_relogin_retry`): + 本方法是纯读,重跑没有副作用;不加这层的话订单页被踢到 SSO 会静默走进 + 解析逻辑、被当成「订单号还没出现」,轮询白转几个小时也看不出原因。 + Raises: NotLoggedInError: 登录态失效且自动重登没能恢复 OrderOperationError: 订单详情页打开/渲染失败 @@ -1796,7 +1859,7 @@ class SiteInteractor: shop_id = site_order_id.split("-", 1)[0] url = _ORDER_DETAIL_URL_TEMPLATE.format(order_number=site_order_id, shop_id=shop_id) - async def read() -> OrderStatusSnapshot: + async def read() -> OrderDetailSnapshot: page = await self._context.new_page() try: try: @@ -1820,20 +1883,27 @@ class SiteInteractor: "rakuten", final_url=final_url, body=html ): raise _LoggedOutMidRead(f"订单详情页落地 {final_url}") - return snapshot + return OrderDetailSnapshot(status=snapshot, raw=_parse_initial_state(html)) return await self._read_with_relogin_retry( f"check_order_status site_order_id={site_order_id}", read ) - async def list_recent_orders(self, *, since: datetime) -> OrderListWindow: - """恢复核对用:拉取「任务创建时间之后」的订单列表(规格 §5 依赖它) + async def list_recent_orders( + self, *, since: datetime, max_pages: int | None = None + ) -> OrderListWindow: + """拉取「since 之后」的订单列表 - 由 verify.verify_on_site 调用,不由 worker 主循环直接调。翻页直到看到 - order_date 早于 since 的订单(说明窗口内的都已经看过一遍)、或 - ordersFound 已经全部翻完、或到达 _ORDER_LIST_MAX_PAGES 上限。命中上限仍 - 没能确认覆盖完整窗口时,`OrderListWindow.window_fully_covered=False`—— - 调用方据此转 unknown,绝不能把「没翻完」当成「翻完了但没有」。 + 两个调用方: + - `verify.verify_on_site` 恢复核对(规格 §5 依赖它) + - 账号只读查询通道的 order_list(规格 §11),此时 `max_pages` 由上游给, + `raw_pages` 会被原样透传出去 + + 翻页直到看到 order_date 早于 since 的订单(说明窗口内的都已经看过一遍)、或 + ordersFound 已经全部翻完、或到达 max_pages 上限(缺省 + `_ORDER_LIST_MAX_PAGES`)。命中上限仍没能确认覆盖完整窗口时, + `OrderListWindow.window_fully_covered=False`——调用方据此转 unknown,绝不能 + 把「没翻完」当成「翻完了但没有」。 2026-08-13 只用「账号只有 1 笔订单、1 页」的真实数据验证过单页解析与 page=2 返回空列表这两点;多页翻页的排序假设(新订单在前)未经真实数据 @@ -1848,13 +1918,15 @@ class SiteInteractor: NotLoggedInError: 登录态失效且自动重登没能恢复 OrderOperationError: 订单列表页打开/渲染失败 """ + page_limit = max(1, min(max_pages or _ORDER_LIST_MAX_PAGES, _ORDER_LIST_MAX_PAGES)) + async def read() -> OrderListWindow: acc = _OrderListAccumulator() stop = False page = await self._context.new_page() try: - for page_num in range(1, _ORDER_LIST_MAX_PAGES + 1): + for page_num in range(1, page_limit + 1): url = _ORDER_LIST_URL if page_num == 1 else f"{_ORDER_LIST_URL}?page={page_num}" try: await page.goto(url, wait_until="domcontentloaded", timeout=30_000) @@ -1880,7 +1952,9 @@ class SiteInteractor: await page.close() return OrderListWindow( - entries=acc.entries, window_fully_covered=stop and not acc.gave_up + entries=acc.entries, + window_fully_covered=stop and not acc.gave_up, + raw_pages=acc.raw_pages, ) return await self._read_with_relogin_retry("list_recent_orders", read) @@ -2093,4 +2167,5 @@ def _parse_order_list(html: str) -> OrderListPage: entries=entries, orders_found=data.get("ordersFound"), page_size=data.get("pageSize"), + raw=data or None, ) diff --git a/docs/order-gateway.md b/docs/order-gateway.md index ba36bc7..0e77375 100644 --- a/docs/order-gateway.md +++ b/docs/order-gateway.md @@ -54,6 +54,11 @@ queued ──lease──> leased ──首次 report──> running ──termin └── 只有人工介入才能从 stale 回到 queued(见 §5) ``` +> 账号**只读**查询(上游问「账号里真实的订单长什么样」, +> 见 [§11](#11-账号只读查询通道))有一套独立的查询单状态机与通道,不复用这张任务表—— +> 查询是可安全重投的,与「下单绝不重投」是两套相反的安全语义,分表就是要防止 +> 这两个 if 被改错。 + ## 4. order-gateway 技术栈与现有仓库保持一致:Python 3.13 + FastAPI + pydantic-settings + SQLite, @@ -296,6 +301,8 @@ trading 侧新增: | 6002 | 租约无效:不是持有者、已过期或任务已终结 | 409 | | 6003 | 任务状态不允许该操作(如对已终结任务 reclaim) | 409 | | 6004 | 已有任务在执行中,本次不发放(正常返回空即可,仅诊断用) | 200 | +| 6005 | 查询单不存在(账号只读查询通道,见 §11.4) | 404 | +| 6006 | 查询单租约无效(账号只读查询通道,见 §11.4) | 409 | ## 9. 验收清单 @@ -333,3 +340,153 @@ trading 侧新增: > **交易范围**:交易服务只覆盖乐天市场(rakuten)。ラクマ 的抓取仍在抓取服务里提供, > 但不进入交易链路——没有加购契约,也永不实现下单/付款/订单监控。 + +## 11. 账号只读查询通道 + +上游有时要的不是「网关记的这笔任务走到哪一步」,而是「**已登录账号在站点上真实的 +订单是什么样**」——两者会不一致:站点侧被商家取消、金额调整、或者根本是走别的 +渠道下的单,网关的状态镜像都看不见。 + +但账号只在本地机上(NAT 后无公网入口),上游只能打网关。所以这条通道与下单任务 +同构:上游把查询意图放进队列,本地 worker 出站长轮询领走、去站点上真读一次、回结果。 + +### 11.1 为什么不复用 `/api/orders` 那条队列 + +三条理由,任何一条单独成立就足够: + +1. **全局并发度 1 是给写操作设的**。查询排在下单后面会被堵死,而只读本来可以并行。 +2. **「租约过期绝不自动重投」是写操作的安全约束**。只读操作重投是安全的,套上 + §5 那套只会制造一堆需要人工 reclaim 的 stale 噪音。 +3. **`task_reports` 的语义是订单状态镜像**(主键 `(task_id, state)`),塞查询结果 + 会污染对账口径。 + +因此单开一张 `account_queries` 表、一套 `QueryStatus` 词汇表、一条 worker 循环。 +两条通道的语义对照写在 `app/gateway/query_queue.py` 顶部,改任一边前先看它。 + +``` +queued ──lease──> leased ──result──> succeeded + ▲ │ └─> failed + │ │ + └─ 租约过期自动重投 ┘ (attempts 用尽 → failed;超过查询单 TTL → expired) +``` + +**账号级串行靠什么保证**:本地 `SiteInteractor` 里那把 `asyncio.Lock`——下单的写 +操作与这里的只读操作都要先拿它。网关侧不重复实现一遍并发闸门。代价是下单跑着时 +查询要排队等锁,因此 worker 侧对每次执行套 `RAKUTEN_ACCOUNT_QUERY_TIMEOUT_SECONDS`, +超时按失败回报,上游重发即可(只读,重发无副作用)。 + +### 11.2 接口 + +**纯异步**:提交只拿单号,结果去 GET 取。网关不提供「挂起等结果」的接口——下单 +可能占着账号锁跑几分钟,同步等待会把上游的连接一起卡住。 + +- `POST /api/account/queries` — 提交查询单(幂等,同 `query_id` 返回既有单) + + ```jsonc + { + "query_id": "aq-20260816-0001", // 可选,不传服务端生成 + "site": "rakuten", + "kind": "order_list", // order_list | order_detail + "params": { // 网关原样透传,结构由 trading 侧定义 + "since": "2026-08-01T00:00:00Z", // order_list 可选,缺省不设下界 + "max_pages": 3 // order_list 可选,1..20 + // order_detail 必填:{"order_number": "306087-20260813-0863947697"} + } + } + ``` + + 响应 `data`:`{"query_id": "...", "status": "queued", "created": true}` + +- `GET /api/account/queries/lease` — 本地长轮询领取(参数 `worker_id`、`wait`、`site`)。 + **不刷 worker 心跳**:心跳的语义是「下单 worker 还活着」,查询循环活着不代表它活着。 +- `POST /api/account/queries/{id}/result` — 本地回结果 + (`{worker_id, success, result?, error_code?, error_message?}`) +- `GET /api/account/queries/{id}` — 取状态与结果 +- `GET /api/account/queries?status=&kind=&limit=&offset=` — 列表(运维排查) + +### 11.3 结果结构:原样透传 + 规范化字段 + +两者都给。规范化字段是稳定契约,站点原始 JSON 是逃生舱——站点比我们的模型多给的 +字段(金额、配送、状态文案)不该在中间层被悄悄丢掉。 + +`kind=order_list`(复用 `list_recent_orders`,2026-08-13 实测过的路径): + +```jsonc +{ + "kind": "order_list", + "since": "2026-08-01T00:00:00Z", + "max_pages": 3, + "window_fully_covered": true, // ← 见下,false 时结论不可当「没有」用 + "orders": [{"order_number": "...", "order_date": "...", "shop_id": 306087, + "shop_name": "...", "items": [{"item_id": 0, "item_name": "", "item_url": ""}]}], + "raw_pages": [ /* 站点 __INITIAL_STATE__.orderListData 原文,按翻页顺序 */ ] +} +``` + +**`window_fully_covered=false` 时,「列表里没有某笔订单」不等于「账号里没有」**—— +它表示翻页在覆盖完窗口前就停了(命中 `max_pages`、页面改版、掉登录)。上游据此 +决定是缩小窗口重查还是交人工,不要把它当成确定性的否定结论。这与 §5 恢复核对里 +「核对不出结论就报 needs_human」是同一条原则。 + +`kind=order_detail`(复用 `fetch_order_detail`): + +```jsonc +{ + "kind": "order_detail", + "order_number": "306087-20260813-0863947697", + "found": true, // false 不是错误:站点自己说下单后要 ~10 分钟才反映 + "stage_label": "出荷", // 站点原文(已实测的配送阶段进度条) + "order_state": "shipped", // 映射到 OrderState;映射不到为 null + "raw": { /* 页面 __INITIAL_STATE__ 原文 */ }, + "raw_available": true +} +``` + +> **实测边界(重要)**:`stage_label` / `order_state` 走的是 2026-08-13 用真实订单号 +> 验证过的进度条解析;但**详情页的 `__INITIAL_STATE__` 结构从未拿到真实样本**。 +> 因此本服务对它只做「解析得动就原样透传」,不抽任何字段——想要金额、收货地址、 +> 付款方式的调用方自己从 `raw` 里取并自担结构变动风险。等拿到真实样本后再在 +> `site_interact.py` 里补规范化解析,**不要在没有样本的情况下先猜着写**。 + +结果体积上限 `RAKUTEN_QUERY_RESULT_MAX_BYTES`:超限先丢 `raw_pages`/`raw` 并置 +`raw_omitted`;丢完仍超限则整单判失败,让上游缩小窗口重来——不把几 MB 页面状态 +塞进网关 SQLite。 + +### 11.4 错误码(新增 `600x`) + +| code | 含义 | HTTP | +| --- | --- | --- | +| 6005 | 查询单不存在(或已过保留期被清理) | 404 | +| 6006 | 查询单租约无效:不是持有者、已被重投给别人或已终结 | 409 | + +worker 回结果撞上 6006 属于正常情况(上一轮超时、单子已重投),**丢弃本次结果即可, +不要重发**——否则会覆盖新一轮的结果。 + +### 11.5 配置项 + +| 配置 | 默认 | 说明 | +| --- | --- | --- | +| `RAKUTEN_QUERY_LEASE_TTL_SECONDS` | 180 | 查询单租约 TTL,过期自动重投 | +| `RAKUTEN_QUERY_TTL_SECONDS` | 900 | 查询单整体存活上限,超时置 expired | +| `RAKUTEN_QUERY_MAX_ATTEMPTS` | 3 | 最多被领取几次,用尽置 failed | +| `RAKUTEN_QUERY_RETENTION_SECONDS` | 604800 | 终态查询单保留时长,过期由 sweep 清理 | +| `RAKUTEN_ACCOUNT_QUERY_TIMEOUT_SECONDS` | 120 | worker 单次站点读取上限(含等账号锁) | +| `RAKUTEN_ACCOUNT_QUERY_DEFAULT_MAX_PAGES` | 3 | order_list 缺省翻页上限 | +| `RAKUTEN_QUERY_RESULT_MAX_BYTES` | 1048576 | 结果 JSON 体积上限 | + +### 11.6 验收清单 + +- [x] 同 `query_id` 提交两次只产生一张单 +- [x] `kind` 非法在提交时就 422,不等 worker 领走才失败 +- [x] 有查询在飞时,下单任务照样能被 lease(两条通道互不阻塞) +- [x] 多张查询单可以同时在飞(没有并发度 1 闸门) +- [x] 租约过期 → 回 `queued` 且能被再次领取,`attempt` 递增 +- [x] 重投次数用尽 → `failed`,不再无限重投 +- [x] 超过查询单 TTL → `expired`,且不会再被领取 +- [x] 重投后旧 worker 的迟到结果被拒(6006),不覆盖新结果 +- [x] 终态单过保留期被清理,进行中的单不受影响 +- [x] 站点原始 JSON 完整出现在结果里;翻页上限生效且如实标 `window_fully_covered` +- [x] 站点报错(掉登录 5001 等)原样带着错误码回给上游 +- [x] 执行超时按失败回报,不把查询单挂到租约过期 +- [x] worker 任何异常都不逃出主循环,每张单都有一次回报 + diff --git a/tests/test_gateway_queries.py b/tests/test_gateway_queries.py new file mode 100644 index 0000000..6d74922 --- /dev/null +++ b/tests/test_gateway_queries.py @@ -0,0 +1,390 @@ +"""账号只读查询通道测试(docs/order-gateway.md §11) + +分两层: +- HTTP 层:鉴权、幂等、6005/6006、lease → result 全链路、列表与 /health +- 队列语义层:直接驱动 QueryQueue,测那些靠 HTTP 不好造的时序——租约超时重投、 + 重投次数用尽、整体 TTL 过期、迟到结果被拒、过保留期清理 + +这里最关键的一条是「只读可重投」:它与下单任务的「绝不自动重投」正好相反, +两条通道的语义分歧全在这几个用例里,改任何一边前先看它们。 +""" +from __future__ import annotations + +from datetime import timedelta +from pathlib import Path + +import pytest +from fastapi.testclient import TestClient + +from app.gateway.db import GatewayDB, QueryRow +from app.gateway.query_queue import QueryQueue, to_iso, utcnow +from app.shared.config import get_settings +from app.shared.errors import QueryLeaseInvalidError +from app.shared.task_state import AccountQueryKind, QueryStatus + +TOKEN = get_settings().bearer_token +AUTH = {"Authorization": f"Bearer {TOKEN}"} + +QUERIES = "/api/account/queries" + + +@pytest.fixture +def gateway_client(tmp_path: Path, monkeypatch): + """起一个独立 DB 的网关应用(与 test_gateway_api.py 同一套装配方式)""" + monkeypatch.setenv("RAKUTEN_GATEWAY_DB_PATH", str(tmp_path / "gw.db")) + get_settings.cache_clear() + try: + from app.gateway.main import create_app + + app = create_app() + with TestClient(app) as client: + yield client + finally: + get_settings.cache_clear() + + +@pytest.fixture +async def queue(tmp_path: Path): + """直接可用的 (QueryQueue, GatewayDB) 对 + + 队列语义用例要模拟「租约已过期」「结果早就写完了」这类时间流逝,直接改库里的 + 时间戳比 sleep 快得多,所以把 db 一并交出去。 + """ + db = GatewayDB(tmp_path / "q.db") + await db.start() + try: + yield ( + QueryQueue( + db, + lease_ttl_seconds=60, + query_ttl_seconds=600, + max_attempts=3, + retention_seconds=7 * 24 * 3600, + ), + db, + ) + finally: + await db.close() + + +def _submit(client: TestClient, **overrides) -> dict: + payload = { + "query_id": "q1", + "site": "rakuten", + "kind": AccountQueryKind.ORDER_LIST.value, + "params": {"max_pages": 1}, + } + payload.update(overrides) + return client.post(QUERIES, json=payload, headers=AUTH).json() + + +# ---- HTTP:鉴权与参数校验 ---- + + +@pytest.mark.parametrize( + "path,method", + [ + (QUERIES, "POST"), + (f"{QUERIES}/lease", "GET"), + (f"{QUERIES}/q1", "GET"), + ], +) +def test_query_endpoints_reject_missing_token(gateway_client, path, method): + response = gateway_client.request(method, path) + assert response.status_code == 401 + assert response.json()["code"] == 1001 + + +def test_unknown_kind_is_rejected_at_submit(gateway_client): + """kind 拼错当场 422,而不是等 worker 领走、执行、回报失败才知道""" + response = gateway_client.post( + QUERIES, + json={"site": "rakuten", "kind": "order_lst", "params": {}}, + headers=AUTH, + ) + assert response.status_code == 422 + assert response.json()["code"] == 1002 + + +# ---- HTTP:提交、幂等、领取、回结果 ---- + + +def test_submit_returns_query_id_and_queued(gateway_client): + body = _submit(gateway_client) + assert body["success"] is True + assert body["data"]["query_id"] == "q1" + assert body["data"]["status"] == QueryStatus.QUEUED.value + assert body["data"]["created"] is True + + +def test_submit_is_idempotent_on_same_query_id(gateway_client): + assert _submit(gateway_client)["data"]["created"] is True + assert _submit(gateway_client)["data"]["created"] is False + + +def test_submit_generates_query_id_when_absent(gateway_client): + body = gateway_client.post( + QUERIES, + json={"site": "rakuten", "kind": AccountQueryKind.ORDER_DETAIL.value, "params": {}}, + headers=AUTH, + ).json() + assert body["data"]["query_id"].startswith("q-") + + +def test_lease_returns_null_when_no_query(gateway_client): + response = gateway_client.get(f"{QUERIES}/lease?worker_id=w1&wait=0", headers=AUTH) + assert response.status_code == 200 + assert response.json()["data"] is None + + +def test_lease_then_result_completes_the_query(gateway_client): + _submit(gateway_client) + + leased = gateway_client.get(f"{QUERIES}/lease?worker_id=w1&wait=0", headers=AUTH).json() + assert leased["data"]["query_id"] == "q1" + assert leased["data"]["kind"] == AccountQueryKind.ORDER_LIST.value + assert leased["data"]["params"] == {"max_pages": 1} + assert leased["data"]["attempt"] == 1 + + result = {"kind": "order_list", "orders": [{"order_number": "306087-20260813-0863947697"}]} + reported = gateway_client.post( + f"{QUERIES}/q1/result", + json={"worker_id": "w1", "success": True, "result": result}, + headers=AUTH, + ).json() + assert reported["data"]["status"] == QueryStatus.SUCCEEDED.value + + detail = gateway_client.get(f"{QUERIES}/q1", headers=AUTH).json()["data"] + assert detail["status"] == QueryStatus.SUCCEEDED.value + assert detail["result"] == result + assert detail["error"] is None + assert detail["lease_owner"] is None + + +def test_failed_result_carries_error_code_and_message(gateway_client): + _submit(gateway_client) + gateway_client.get(f"{QUERIES}/lease?worker_id=w1&wait=0", headers=AUTH) + gateway_client.post( + f"{QUERIES}/q1/result", + json={ + "worker_id": "w1", + "success": False, + "error_code": 5001, + "error_message": "账号未登录", + }, + headers=AUTH, + ) + detail = gateway_client.get(f"{QUERIES}/q1", headers=AUTH).json()["data"] + assert detail["status"] == QueryStatus.FAILED.value + assert detail["error"] == {"code": 5001, "message": "账号未登录"} + + +def test_two_queries_can_be_in_flight_at_once(gateway_client): + """与下单任务不同:只读查询没有「全局并发度 1」的闸门""" + _submit(gateway_client, query_id="q1") + _submit(gateway_client, query_id="q2") + + first = gateway_client.get(f"{QUERIES}/lease?worker_id=w1&wait=0", headers=AUTH).json() + second = gateway_client.get(f"{QUERIES}/lease?worker_id=w2&wait=0", headers=AUTH).json() + + assert first["data"]["query_id"] == "q1" + assert second["data"]["query_id"] == "q2" + + +def test_order_lease_is_not_disturbed_by_pending_query(gateway_client): + """两条通道互不干扰:有查询在飞,下单任务照领""" + _submit(gateway_client) + gateway_client.get(f"{QUERIES}/lease?worker_id=w1&wait=0", headers=AUTH) + + gateway_client.post( + "/api/orders", + json={"task_id": "t1", "site": "rakuten", "intent": {}}, + headers=AUTH, + ) + leased = gateway_client.get("/api/orders/lease?worker_id=w1&wait=0", headers=AUTH).json() + assert leased["data"]["task_id"] == "t1" + + +# ---- HTTP:错误码 ---- + + +def test_get_unknown_query_returns_6005(gateway_client): + response = gateway_client.get(f"{QUERIES}/no-such", headers=AUTH) + assert response.status_code == 404 + assert response.json()["code"] == 6005 + + +def test_result_from_non_owner_returns_6006(gateway_client): + _submit(gateway_client) + gateway_client.get(f"{QUERIES}/lease?worker_id=w1&wait=0", headers=AUTH) + + response = gateway_client.post( + f"{QUERIES}/q1/result", + json={"worker_id": "w2", "success": True, "result": {}}, + headers=AUTH, + ) + assert response.status_code == 409 + assert response.json()["code"] == 6006 + + +def test_result_on_terminal_query_returns_6006(gateway_client): + _submit(gateway_client) + gateway_client.get(f"{QUERIES}/lease?worker_id=w1&wait=0", headers=AUTH) + gateway_client.post( + f"{QUERIES}/q1/result", + json={"worker_id": "w1", "success": True, "result": {"a": 1}}, + headers=AUTH, + ) + + response = gateway_client.post( + f"{QUERIES}/q1/result", + json={"worker_id": "w1", "success": True, "result": {"a": 2}}, + headers=AUTH, + ) + assert response.status_code == 409 + assert response.json()["code"] == 6006 + # 结果没有被第二次回报覆盖 + detail = gateway_client.get(f"{QUERIES}/q1", headers=AUTH).json()["data"] + assert detail["result"] == {"a": 1} + + +# ---- HTTP:列表与健康检查 ---- + + +def test_list_queries_filters_by_kind(gateway_client): + _submit(gateway_client, query_id="q1", kind=AccountQueryKind.ORDER_LIST.value) + _submit(gateway_client, query_id="q2", kind=AccountQueryKind.ORDER_DETAIL.value) + + body = gateway_client.get( + f"{QUERIES}?kind={AccountQueryKind.ORDER_DETAIL.value}", headers=AUTH + ).json() + assert body["data"]["total"] == 1 + assert body["data"]["items"][0]["query_id"] == "q2" + + +def test_health_reports_queued_query_count(gateway_client): + _submit(gateway_client) + body = gateway_client.get("/health").json() + assert body["data"]["queued_query_count"] == 1 + + +# ---- 队列语义:只读可重投(与下单任务的核心分歧)---- + + +async def test_expired_lease_goes_back_to_queued(queue): + """租约过期→回 queued 自动重投。下单任务在这里是 stale 且绝不重投""" + q, db = queue + await q.submit(query_id="q1", site="rakuten", kind="order_list", params={}) + lease = await q.lease(worker_id="w1", wait=0, site=None, max_wait=60) + assert lease is not None + + await db.update_query("q1", lease_expires_at=to_iso(utcnow() - timedelta(seconds=1))) + assert await q.sweep() == 1 + + detail = await q.get_detail("q1") + assert detail.status == QueryStatus.QUEUED + assert detail.lease_owner is None + + again = await q.lease(worker_id="w2", wait=0, site=None, max_wait=60) + assert again is not None + assert again.attempt == 2 + + +async def test_late_result_after_reinvest_is_rejected(queue): + """超时重投后,上一轮 worker 迟到的结果必须被拒,不能覆盖新一轮""" + q, db = queue + await q.submit(query_id="q1", site="rakuten", kind="order_list", params={}) + await q.lease(worker_id="w1", wait=0, site=None, max_wait=60) + await db.update_query("q1", lease_expires_at=to_iso(utcnow() - timedelta(seconds=1))) + await q.sweep() + await q.lease(worker_id="w2", wait=0, site=None, max_wait=60) + + with pytest.raises(QueryLeaseInvalidError): + await q.submit_result( + "q1", worker_id="w1", success=True, result={"stale": True}, + error_code=None, error_message="", + ) + + +async def test_reinvest_stops_at_max_attempts(queue): + """重投次数用尽 → failed,不再无限重投同一张卡死的单""" + q, db = queue + await q.submit(query_id="q1", site="rakuten", kind="order_list", params={}) + for _ in range(3): + await q.lease(worker_id="w1", wait=0, site=None, max_wait=60) + await db.update_query("q1", lease_expires_at=to_iso(utcnow() - timedelta(seconds=1))) + await q.sweep() + + detail = await q.get_detail("q1") + assert detail.status == QueryStatus.FAILED + assert detail.attempts == 3 + assert "不再重投" in detail.error.message + + +async def test_query_expires_after_ttl(tmp_path): + """整体 TTL 到点 → expired(本地 worker 不在线的典型表现)""" + db = GatewayDB(tmp_path / "q.db") + await db.start() + try: + queue = QueryQueue( + db, lease_ttl_seconds=60, query_ttl_seconds=60, max_attempts=3, + retention_seconds=3600, + ) + # 直接插一张「10 分钟前创建」的单,比 sleep 现实 + old = to_iso(utcnow() - timedelta(minutes=10)) + await db.insert_query( + QueryRow( + query_id="q1", site="rakuten", kind="order_list", params_json="{}", + status=QueryStatus.QUEUED.value, lease_owner=None, lease_expires_at=None, + attempts=0, result_json=None, error_code=None, error_message=None, + created_at=old, updated_at=old, completed_at=None, + ) + ) + + assert await queue.sweep() == 1 + detail = await queue.get_detail("q1") + assert detail.status == QueryStatus.EXPIRED + assert "worker 不在线" in detail.error.message + # 过期的单不会再被领走 + assert await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60) is None + finally: + await db.close() + + +async def test_sweep_purges_queries_past_retention(tmp_path): + """终态查询单过保留期被清掉——结果里带站点原始 JSON,不清会一直涨""" + db = GatewayDB(tmp_path / "q.db") + await db.start() + try: + queue = QueryQueue( + db, lease_ttl_seconds=60, query_ttl_seconds=600, max_attempts=3, + retention_seconds=1, + ) + await queue.submit(query_id="q1", site="rakuten", kind="order_list", params={}) + await queue.lease(worker_id="w1", wait=0, site=None, max_wait=60) + await queue.submit_result( + "q1", worker_id="w1", success=True, result={"a": 1}, + error_code=None, error_message="", + ) + await db.update_query("q1", completed_at=to_iso(utcnow() - timedelta(seconds=60))) + + await queue.sweep() + assert await db.get_query("q1") is None + finally: + await db.close() + + +async def test_sweep_keeps_unfinished_queries(tmp_path): + """保留期清理只动终态单,进行中的一律不碰""" + db = GatewayDB(tmp_path / "q.db") + await db.start() + try: + queue = QueryQueue( + db, lease_ttl_seconds=60, query_ttl_seconds=600, max_attempts=3, + retention_seconds=1, + ) + await queue.submit(query_id="q1", site="rakuten", kind="order_list", params={}) + await queue.sweep() + assert await db.get_query("q1") is not None + finally: + await db.close() \ No newline at end of file diff --git a/tests/test_query_runner.py b/tests/test_query_runner.py new file mode 100644 index 0000000..8b8e46e --- /dev/null +++ b/tests/test_query_runner.py @@ -0,0 +1,417 @@ +"""账号只读查询 worker 测试(app/trading/worker/query_runner.py) + +站点交互与网关都用桩替换,测的是 QueryRunner 自己的编排: +- 两种 kind 各自把站点结果整理成什么样的 result(原始 JSON 有没有原样带出来) +- 失败怎么转成一次「失败回报」——超时、站点 AppError、参数不合法、未知 kind +- 结果体积闸门:先丢站点原始 JSON,仍超限就判失败让上游缩小窗口 + +一条贯穿始终的约定:**worker 永远不把异常抛回主循环**,每张查询单都要有一次 +回报(成功或失败),否则网关侧那张单只能干等到租约超时。 +""" +from __future__ import annotations + +import asyncio +from datetime import datetime, timezone +from typing import Any + +import pytest + +from app.shared.config import Settings +from app.shared.errors import AppError, NotLoggedInError +from app.shared.task_state import AccountQueryKind, OrderState +from app.trading.worker.models import QueryTask +from app.trading.worker.query_runner import QueryRunner +from app.trading.worker.site_interact import ( + OrderDetailSnapshot, + OrderListEntry, + OrderListItem, + OrderListWindow, + OrderStatusSnapshot, +) + + +class _FakeGateway: + """记录回报内容的网关客户端替身 + + `lease_query` 里那句 `sleep(0)` 是必需的:真实实现是 HTTP 长轮询,一定会把 + 控制权交回事件循环;纯内存替身不 await 任何东西的话,run() 会一直霸占循环, + 同一个 event loop 里的停机任务永远排不上,测试直接挂死。 + """ + + def __init__(self, queries: list[QueryTask | None] | None = None): + self._queries = list(queries or []) + self.reports: list[dict[str, Any]] = [] + self.lease_calls = 0 + + async def lease_query(self, worker_id: str, *, wait: int = 30) -> QueryTask | None: + self.lease_calls += 1 + await asyncio.sleep(0) + if not self._queries: + return None + return self._queries.pop(0) + + async def report_query_result( + self, + query_id: str, + worker_id: str, + *, + success: bool, + result: dict[str, Any] | None = None, + error_code: int | None = None, + error_message: str = "", + ) -> dict[str, Any]: + self.reports.append( + { + "query_id": query_id, + "worker_id": worker_id, + "success": success, + "result": result, + "error_code": error_code, + "error_message": error_message, + } + ) + return {"query_id": query_id, "status": "succeeded" if success else "failed"} + + +class _FakeSite: + """SiteInteractor 替身:按脚本返回窗口/详情,或抛出指定异常""" + + def __init__(self, *, window=None, detail=None, error: Exception | None = None, delay: float = 0): + self._window = window + self._detail = detail + self._error = error + self._delay = delay + self.list_calls: list[dict[str, Any]] = [] + self.detail_calls: list[str] = [] + + async def list_recent_orders(self, *, since, max_pages=None): + self.list_calls.append({"since": since, "max_pages": max_pages}) + await asyncio.sleep(self._delay) + if self._error: + raise self._error + return self._window + + async def fetch_order_detail(self, site_order_id: str): + self.detail_calls.append(site_order_id) + await asyncio.sleep(self._delay) + if self._error: + raise self._error + return self._detail + + +def _runner(site: _FakeSite, gateway: _FakeGateway, **settings_kwargs) -> QueryRunner: + settings = Settings(worker_id="w1", **settings_kwargs) + return QueryRunner(settings=settings, gateway_client=gateway, site=site) # type: ignore[arg-type] + + +def _query(kind: str, params: dict[str, Any] | None = None, **kwargs) -> QueryTask: + return QueryTask( + query_id=kwargs.get("query_id", "q1"), + site=kwargs.get("site", "rakuten"), + kind=kind, + params=params or {}, + attempt=kwargs.get("attempt", 1), + ) + + +def _window(*, covered: bool = True) -> OrderListWindow: + return OrderListWindow( + entries=[ + OrderListEntry( + order_number="306087-20260813-0863947697", + order_date="2026-08-13T09:47:28.000Z", + shop_id=306087, + shop_name="BACKYARD FAMILY インテリアタウン", + items=[ + OrderListItem( + item_url="https://item.rakuten.co.jp/backyard/x/", + item_name="壁掛けフック", + item_id=10012345, + ) + ], + ) + ], + window_fully_covered=covered, + raw_pages=[{"ordersFound": 1, "orderList": [{"orderNumber": "306087-20260813-0863947697"}]}], + ) + + +# ---- order_list ---- + + +async def test_order_list_returns_normalized_and_raw(): + site = _FakeSite(window=_window()) + gateway = _FakeGateway() + runner = _runner(site, gateway) + + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value, {"max_pages": 2})) + + report = gateway.reports[0] + assert report["success"] is True + result = report["result"] + assert result["kind"] == "order_list" + assert result["window_fully_covered"] is True + assert result["orders"][0]["order_number"] == "306087-20260813-0863947697" + assert result["orders"][0]["items"][0]["item_id"] == 10012345 + # 站点原文原样带出,规范化模型没覆盖的字段上游还能自己取 + assert result["raw_pages"][0]["ordersFound"] == 1 + assert site.list_calls[0]["max_pages"] == 2 + + +async def test_order_list_since_is_passed_through(): + site = _FakeSite(window=_window()) + runner = _runner(site, _FakeGateway()) + + await runner.handle( + _query(AccountQueryKind.ORDER_LIST.value, {"since": "2026-08-01T00:00:00Z"}) + ) + + assert site.list_calls[0]["since"] == datetime(2026, 8, 1, tzinfo=timezone.utc) + + +async def test_order_list_without_since_uses_epoch_and_default_max_pages(): + site = _FakeSite(window=_window()) + runner = _runner(site, _FakeGateway(), account_query_default_max_pages=5) + + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value)) + + assert site.list_calls[0]["since"] == datetime(1970, 1, 1, tzinfo=timezone.utc) + assert site.list_calls[0]["max_pages"] == 5 + + +async def test_order_list_reports_partial_coverage_honestly(): + """翻页没覆盖完窗口时必须如实标出来——上游据此知道「没查到」不等于「没有」""" + site = _FakeSite(window=_window(covered=False)) + gateway = _FakeGateway() + runner = _runner(site, gateway) + + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value)) + + assert gateway.reports[0]["result"]["window_fully_covered"] is False + + +async def test_invalid_since_is_reported_as_failure(): + site = _FakeSite(window=_window()) + gateway = _FakeGateway() + runner = _runner(site, gateway) + + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value, {"since": "上周"})) + + assert gateway.reports[0]["success"] is False + assert gateway.reports[0]["error_code"] == 1003 + assert site.list_calls == [] # 参数不合法就没去碰站点 + + +async def test_invalid_max_pages_is_reported_as_failure(): + gateway = _FakeGateway() + runner = _runner(_FakeSite(window=_window()), gateway) + + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value, {"max_pages": 0})) + + assert gateway.reports[0]["success"] is False + assert gateway.reports[0]["error_code"] == 1003 + + +# ---- order_detail ---- + + +async def test_order_detail_returns_stage_and_raw(): + detail = OrderDetailSnapshot( + status=OrderStatusSnapshot( + found=True, stage_label="出荷", order_state=OrderState.SHIPPED, html="" + ), + raw={"pageType": "ph-detail"}, + ) + site = _FakeSite(detail=detail) + gateway = _FakeGateway() + runner = _runner(site, gateway) + + await runner.handle( + _query(AccountQueryKind.ORDER_DETAIL.value, {"order_number": "306087-20260813-0863947697"}) + ) + + result = gateway.reports[0]["result"] + assert result["found"] is True + assert result["stage_label"] == "出荷" + assert result["order_state"] == OrderState.SHIPPED.value + assert result["raw"] == {"pageType": "ph-detail"} + assert result["raw_available"] is True + assert site.detail_calls == ["306087-20260813-0863947697"] + + +async def test_order_detail_not_found_is_success_not_failure(): + """站点自己说订单要 10 分钟才反映出来:查不到是正常结果,不是查询失败""" + site = _FakeSite(detail=OrderDetailSnapshot(status=OrderStatusSnapshot(found=False), raw=None)) + gateway = _FakeGateway() + runner = _runner(site, gateway) + + await runner.handle(_query(AccountQueryKind.ORDER_DETAIL.value, {"order_number": "x-1-2"})) + + assert gateway.reports[0]["success"] is True + assert gateway.reports[0]["result"]["found"] is False + assert gateway.reports[0]["result"]["raw_available"] is False + + +async def test_order_detail_requires_order_number(): + gateway = _FakeGateway() + runner = _runner(_FakeSite(), gateway) + + await runner.handle(_query(AccountQueryKind.ORDER_DETAIL.value, {})) + + assert gateway.reports[0]["success"] is False + assert gateway.reports[0]["error_code"] == 1003 + + +# ---- 失败路径 ---- + + +async def test_unknown_kind_is_reported_as_failure(): + gateway = _FakeGateway() + runner = _runner(_FakeSite(), gateway) + + await runner.handle(_query("order_everything")) + + assert gateway.reports[0]["success"] is False + assert "未知的查询种类" in gateway.reports[0]["error_message"] + + +async def test_non_rakuten_site_is_rejected(): + """交易服务只覆盖乐天市场,ラクマ 永不进入账号链路""" + gateway = _FakeGateway() + runner = _runner(_FakeSite(), gateway) + + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value, site="rakuma")) + + assert gateway.reports[0]["success"] is False + assert "rakuma" in gateway.reports[0]["error_message"] + + +async def test_site_app_error_keeps_its_error_code(): + """站点侧错误码原样回给上游(掉登录=5001),不被翻译成一个笼统的失败""" + site = _FakeSite(error=NotLoggedInError(site="rakuten")) + gateway = _FakeGateway() + runner = _runner(site, gateway) + + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value)) + + assert gateway.reports[0]["success"] is False + assert gateway.reports[0]["error_code"] == 5001 + + +async def test_unexpected_exception_is_reported_not_raised(): + site = _FakeSite(error=RuntimeError("boom")) + gateway = _FakeGateway() + runner = _runner(site, gateway) + + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value)) + + assert gateway.reports[0]["success"] is False + assert gateway.reports[0]["error_code"] == 1500 + assert "RuntimeError" in gateway.reports[0]["error_message"] + + +async def test_timeout_is_reported_as_retryable_failure(): + """账号锁被下单占着 → 超时按失败回报(上游重发即可),不挂到租约过期""" + site = _FakeSite(window=_window(), delay=0.2) + gateway = _FakeGateway() + runner = _runner(site, gateway, account_query_timeout_seconds=0) + + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value)) + + assert gateway.reports[0]["success"] is False + assert gateway.reports[0]["error_code"] == 2002 + assert "可稍后重发" in gateway.reports[0]["error_message"] + + +async def test_report_failure_does_not_escape(): + """回报本身失败(比如租约已被重投)只记日志,不能把循环带崩""" + + class _RejectingGateway(_FakeGateway): + async def report_query_result(self, *args: Any, **kwargs: Any) -> dict[str, Any]: + raise AppError(message="lease invalid", code="X", err_code=6006, status_code=409) + + runner = _runner(_FakeSite(window=_window()), _RejectingGateway()) + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value)) # 不抛异常即通过 + + +# ---- 结果体积闸门 ---- + + +async def test_oversized_result_drops_raw_first(): + big_window = _window() + big_window.raw_pages = [{"padding": "x" * 5000}] + gateway = _FakeGateway() + runner = _runner(_FakeSite(window=big_window), gateway, query_result_max_bytes=2000) + + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value)) + + report = gateway.reports[0] + assert report["success"] is True + assert report["result"]["raw_pages"] == [] + assert "raw_omitted" in report["result"] + # 规范化字段还在——上游至少拿得到订单号 + assert report["result"]["orders"][0]["order_number"] == "306087-20260813-0863947697" + + +async def test_result_still_oversized_after_dropping_raw_fails(): + """丢掉原始 JSON 仍超限:判失败并告诉上游缩小窗口,不硬塞进网关库""" + window = _window() + window.entries = [ + OrderListEntry(order_number=f"o{i}", order_date="2026-08-13T09:47:28.000Z") + for i in range(200) + ] + gateway = _FakeGateway() + runner = _runner(_FakeSite(window=window), gateway, query_result_max_bytes=500) + + await runner.handle(_query(AccountQueryKind.ORDER_LIST.value)) + + assert gateway.reports[0]["success"] is False + assert "缩小查询窗口" in gateway.reports[0]["error_message"] + + +# ---- 主循环 ---- + + +async def test_run_loop_handles_then_stops(): + """领到单就处理;stop() 后循环退出,不吃掉后续任务""" + gateway = _FakeGateway([_query(AccountQueryKind.ORDER_LIST.value), None]) + runner = _runner(_FakeSite(window=_window()), gateway) + + async def stop_soon(): + # 让 run() 先跑几轮(领单 → 处理 → 再领一次拿到 None),再停 + for _ in range(20): + await asyncio.sleep(0) + runner.stop() + + await asyncio.wait_for(asyncio.gather(runner.run(), stop_soon()), timeout=5) + + assert [r["query_id"] for r in gateway.reports] == ["q1"] + + +async def test_run_loop_survives_lease_error(monkeypatch): + """lease 报错不能让循环退出(网关重启、网络抖动都算正常)""" + real_sleep = asyncio.sleep + + class _FlakyGateway(_FakeGateway): + """前两次 lease 都报错,第二次顺便把循环停掉,避免测试无限转""" + + def __init__(self, stopper): + super().__init__() + self._stopper = stopper + + async def lease_query(self, worker_id: str, *, wait: int = 30): + self.lease_calls += 1 + if self.lease_calls >= 2: + self._stopper() + raise AppError(message="gateway down", code="X", err_code=3001) + + holder: dict[str, QueryRunner] = {} + gateway = _FlakyGateway(lambda: holder["runner"].stop()) + runner = _runner(_FakeSite(), gateway) + holder["runner"] = runner + # 出错后那句 sleep(5) 在测试里没必要真等;捕获原函数再替换,否则递归自调 + monkeypatch.setattr(asyncio, "sleep", lambda *_: real_sleep(0)) + + await runner.run() + + assert gateway.lease_calls == 2 diff --git a/tests/test_site_interact.py b/tests/test_site_interact.py index 92e7e3b..64ef005 100644 --- a/tests/test_site_interact.py +++ b/tests/test_site_interact.py @@ -754,6 +754,109 @@ async def test_list_recent_orders_without_start_raises(): await site.list_recent_orders(since=_SINCE) +# ---- 账号只读查询通道用到的两条只读路径(docs/order-gateway.md §11)---- +# +# 查询接口不新写解析,复用的就是上面这两条已实测路径;这里补的是「站点原始 +# JSON 有没有被完整带出来」和「翻页上限有没有生效」——这两点是查询接口独有的, +# 恢复核对那条老路径不关心。 + + +def test_parse_order_list_keeps_raw_order_list_data(): + """规范化字段之外,站点 orderListData 原文要原样留着供上游取用""" + html = _wrap_state(_real_order_list_state()) + page = _parse_order_list(html) + assert page.raw is not None + assert page.raw["ordersFound"] == 1 + # 规范化模型里没有的字段也在(这正是「原样透传」的意义) + assert page.raw["orderList"][0]["shopName"] == "BACKYARD FAMILY インテリアタウン" + + +def test_parse_order_list_raw_is_none_when_page_type_wrong(): + """不是 ph-list(改版 / 掉登录 / act 分支不对):raw 为 None,不给上游半截数据""" + html = _wrap_state(_real_order_list_state(page_type="ph-detail")) + assert _parse_order_list(html).raw is None + + +def test_accumulate_collects_raw_pages_in_order(): + """多页翻页时,每页的原始 JSON 按顺序累积""" + acc = _OrderListAccumulator() + page1 = OrderListPage( + entries=[_entry("o1", "2026-08-10T00:00:00Z")], orders_found=3, raw={"page": 1} + ) + acc, stop = _accumulate_order_list_page(acc, page1, since=_SINCE) + assert stop is False + page2 = OrderListPage( + entries=[_entry("o2", "2026-07-01T00:00:00Z")], orders_found=3, raw={"page": 2} + ) + acc, stop = _accumulate_order_list_page(acc, page2, since=_SINCE) + assert stop is True + assert acc.raw_pages == [{"page": 1}, {"page": 2}] + + +async def test_list_recent_orders_respects_max_pages(tmp_path): + """max_pages 是硬上限:翻到上限就停,且必须如实报 window_fully_covered=False""" + # 每页都还有更多订单(orders_found 远大于已取回数),正常会一直翻下去 + state = _real_order_list_state(orders_found=99) + page = _FakeOrderPage([(_ORDER_LIST_LANDED_URL, _wrap_state(state))]) + site = _build_site(tmp_path, _FakeContext([page]), _FakeAuthSession()) + + window = await site.list_recent_orders(since=_SINCE, max_pages=2) + + assert len(page.goto_urls) == 2 + assert page.goto_urls[1].endswith("?page=2") + assert window.window_fully_covered is False + assert len(window.raw_pages) == 2 + + +async def test_list_recent_orders_window_carries_raw_pages(tmp_path): + html = _wrap_state(_real_order_list_state()) + page = _FakeOrderPage([(_ORDER_LIST_LANDED_URL, html)]) + site = _build_site(tmp_path, _FakeContext([page]), _FakeAuthSession()) + + window = await site.list_recent_orders(since=_SINCE) + + assert window.window_fully_covered is True + assert [raw["ordersFound"] for raw in window.raw_pages] == [1] + + +async def test_fetch_order_detail_returns_status_and_raw_state(tmp_path): + """详情页:已实测的配送阶段照常解析,页面原始状态原样带出""" + html = _wrap_state({"pageType": "ph-detail", "whatever": {"a": 1}}) + _stepper_html( + active_stage="出荷" + ) + page = _FakeOrderPage([(_ORDER_DETAIL_LANDED_URL, html)]) + site = _build_site(tmp_path, _FakeContext([page]), _FakeAuthSession()) + + detail = await site.fetch_order_detail(_REAL_ORDER_ID) + + assert detail.status.found is True + assert detail.status.order_state == OrderState.SHIPPED + assert detail.raw == {"pageType": "ph-detail", "whatever": {"a": 1}} + + +async def test_fetch_order_detail_raw_is_none_without_inline_state(tmp_path): + """页面没有内联状态时 raw=None——不编造,也不因此把这次读取判成失败""" + page = _FakeOrderPage([(_ORDER_DETAIL_LANDED_URL, _stepper_html(active_stage="出荷"))]) + site = _build_site(tmp_path, _FakeContext([page]), _FakeAuthSession()) + + detail = await site.fetch_order_detail(_REAL_ORDER_ID) + + assert detail.raw is None + assert detail.status.found is True + + +async def test_check_order_status_still_returns_only_status(tmp_path): + """付款后监控那条老路径不受影响:仍然拿到 OrderStatusSnapshot""" + page = _FakeOrderPage([(_ORDER_DETAIL_LANDED_URL, _stepper_html(active_stage="配達完了"))]) + site = _build_site(tmp_path, _FakeContext([page]), _FakeAuthSession()) + + snapshot = await site.check_order_status(_REAL_ORDER_ID) + + assert isinstance(snapshot, OrderStatusSnapshot) + assert snapshot.order_state == OrderState.DELIVERED + + + # ---- submit_order / pay:没有 enter_checkout 留存的确认页会话时应报错,不静默成功 ----