diff --git a/app/gateway/api/routes/account.py b/app/gateway/api/routes/account.py new file mode 100644 index 0000000..bd05800 --- /dev/null +++ b/app/gateway/api/routes/account.py @@ -0,0 +1,100 @@ +"""账号订单编目与定时下派触发路由(docs/order-gateway.md §12) + +collector 周期性把账号真实订单沉淀进 `account_orders` 编目,这里是让上游与运维 +读取编目、以及手动触发一轮扫描的入口: + +- GET /api/account/orders 列出已发现的账号订单(分页、可按状态过滤) +- GET /api/account/orders/{order_number} 单笔编目订单 +- POST /api/account/discovery/trigger 立即派一轮 order_list 扫描 + +与 `/api/account/queries` 是两回事:那边是上游主动提交的一次性查询单;这里是 +网关自己周期盘点产生的持久编目。编目行只是规范化字段(订单号/店铺/日期/配送 +状态),不含站点原始 JSON——要看细节去拿对应订单的 order_detail 查询单。 +""" +from __future__ import annotations + +from fastapi import APIRouter, Depends, Query + +from app.gateway.container import GatewayContainer +from app.gateway.models import ( + CatalogOrderListData, + CatalogedOrder, + TriggerDiscoveryData, +) +from app.shared.api import ApiResponse, get_container, require_bearer_token +from app.shared.errors import CatalogOrderNotFoundError + +router = APIRouter(prefix="/api/account", tags=["account-orders"]) + + +@router.get( + "/orders", + response_model=ApiResponse[CatalogOrderListData], + dependencies=[Depends(require_bearer_token)], +) +async def list_cataloged_orders( + state: str | None = Query(default=None, description="按映射后的 OrderState 过滤"), + limit: int = Query(default=50, ge=1, le=500), + offset: int = Query(default=0, ge=0), + container: GatewayContainer = Depends(get_container), +) -> ApiResponse[CatalogOrderListData]: + """列出网关编目里已发现的账号订单,按最近出现时间倒序""" + rows, total = await container.db.list_account_orders( + state=state, limit=limit, offset=offset + ) + items = [_row_to_model(r) for r in rows] + return ApiResponse[CatalogOrderListData]( + success=True, msg="success", + data=CatalogOrderListData(items=items, total=total, limit=limit, offset=offset), + code=0, + ) + + +@router.get( + "/orders/{order_number}", + response_model=ApiResponse[CatalogedOrder], + dependencies=[Depends(require_bearer_token)], +) +async def get_cataloged_order( + order_number: str, + container: GatewayContainer = Depends(get_container), +) -> ApiResponse[CatalogedOrder]: + """取单笔编目订单;编目里没有则 6005(与查询单同号语义)""" + row = await container.db.get_account_order(order_number) + if row is None: + raise CatalogOrderNotFoundError(order_number) + return ApiResponse[CatalogedOrder]( + success=True, msg="success", data=_row_to_model(row), code=0 + ) + + +@router.post( + "/discovery/trigger", + response_model=ApiResponse[TriggerDiscoveryData], + dependencies=[Depends(require_bearer_token)], +) +async def trigger_discovery( + container: GatewayContainer = Depends(get_container), +) -> ApiResponse[TriggerDiscoveryData]: + """手动立即派一轮 order_list 扫描,不等定时器到点 + + 返回本轮下派的单号,结果稍后看对应订单的 order_detail 查询单或编目。 + """ + data = await container.collector.trigger_list_sweep() + return ApiResponse[TriggerDiscoveryData]( + success=True, msg="success", data=TriggerDiscoveryData(**data), code=0 + ) + + +def _row_to_model(row) -> CatalogedOrder: + return CatalogedOrder( + order_number=row.order_number, + shop_id=row.shop_id, + shop_name=row.shop_name, + order_date=row.order_date, + delivery_status=row.delivery_status, + order_state=row.order_state, + discovered_at=row.discovered_at, + detail_fetched_at=row.detail_fetched_at, + last_seen_at=row.last_seen_at, + ) diff --git a/app/gateway/collector.py b/app/gateway/collector.py new file mode 100644 index 0000000..8b70f4a --- /dev/null +++ b/app/gateway/collector.py @@ -0,0 +1,353 @@ +"""定时下派通道:周期性发现账号真实订单并沉淀到网关编目(docs/order-gateway.md §12) + +**它解决什么问题**:上游经「账号只读查询通道」(§11)要订单详情,得**自己先知道 +order_number 再下派**。可账号在站点上可能走了别的渠道下单、被商家改状态,网关的 +状态镜像根本看不到。本模块让网关**自己定期去盘点账号**:派一张 order_list 查询 +(走 §11 那条队列,本地 worker 出站真读),拿到列表后把「网关编目里还没有的订单」 +回写进 `account_orders` 表,再逐笔派 order_detail 查询把配送阶段沉淀下来。之后 +上游直接查 `/api/account/orders` 就能拿到账号里真实有哪些订单,不必再自己记。 + +**实现里最要紧的一条边界**:collector 只读 `query_queue` 查询单结果里的 +**规范化字段**(`orders[].order_number/order_date/shop_id/shop_name` 与 +`delivery_status/order_state`——这些是 worker 在 query_runner 里专门为对外契约 +抽好的稳定字段,见 §11.3),**绝不解析结果里的 `raw` / `raw_pages` 站点原始 JSON**。 +「网关不解释 query 的 params/result」这条原则照旧适用于上游自提的查询单;collector +自主派发的这批单是网关自己的编排数据,消费自己的规范化契约不违反它。 + +**状态模型刻意保持无痕(不新增跟踪表)**: +- collector 只在内存里记「本轮派出去但还没处理结果」的 query_id(`_pending`), + 进程重启即清空。重启后既有的 discover-* 查询单要么被 worker 消费完留在 DB 里 + (无主结果,7 天保留期后被 sweep 清掉),要么过期/失败;下一轮 tick 会派一张 + 全新的 list 扫描,结果幂等 upsert 进编目,不会重复。 +- `order_detail` 的 query_id 用 `discover-detail--<日期>` + **按天分段**:同一天内对同一笔订单重复 submit 会命中幂等返回同一张单(不会 + 因每轮 tick 都满足「详情已过期」就反复创建);跨天自然生成新单刷新详情。 +- 编目 upsert 以 `order_number` 为主键,重复采集只刷新,不产生重复行。 + +**谁负责执行**:本模块只负责派单与沉淀结果;实际的订单列表/详情抓取仍然由本地 +worker 的 `query_runner` 经 §11 队列执行,本模块不碰浏览器。 +""" +from __future__ import annotations + +import asyncio +import logging +from datetime import datetime, timedelta, timezone +from typing import TYPE_CHECKING, Any + +from app.gateway.db import AccountOrderRow, GatewayDB +from app.gateway.query_queue import QueryQueue +from app.shared.task_state import ( + AccountQueryKind, + QueryStatus, + TERMINAL_STATUSES, +) + +if TYPE_CHECKING: + from app.shared.config import Settings + +logger = logging.getLogger(__name__) + +# collector 主循环的节拍(秒):常驻时定时去收割结果与检查排队状态。派单间隔 +# 由 settings.account_discovery_interval_seconds / account_detail_refresh_seconds +# 决定,不靠这个节拍驱动。 +_TICK_SECONDS = 60 + +# 查询单终态集合(字符串形式):SUCCEEDED/FAILED/EXPIRED。收割时据此判断「单子 +# 出结果了没」——仍在 queued/leased 的一律不碰,留到下一轮。 +_TERMINAL = tuple(s.value for s in TERMINAL_STATUSES) + + +def utcnow() -> datetime: + return datetime.now(timezone.utc) + + +def _day_key(dt: datetime) -> str: + """给 order_detail 的 query_id 按天分桶用:如 20260816""" + return dt.astimezone(timezone.utc).strftime("%Y%m%d") + + +def _coerce_str(value: Any) -> str | None: + """把 result 里的值安全地转成字符串;缺失/非法返回 None,不报错""" + if value is None: + return None + if isinstance(value, bool): + return None + return str(value) + + +def _parse_list_result(result: dict[str, Any] | None) -> list[dict[str, Any]]: + """从 order_list 查询结果里抽规范化订单行,形状不符则保守地返回空 + + 只取 `orders` 数组里每个元素的 order_number 这几个规范化字段;`raw_pages` + 原地透传的站点原始 JSON **一律不碰**。这里不抛错——collector 是后台任务, + 某个字段缺了只会导致这笔订单索引不到,不该让整轮扫描崩溃。 + """ + if not result: + return [] + orders = result.get("orders") + if not isinstance(orders, list): + logger.warning("order_list 结果里没有 orders 数组,本轮不沉淀:result_keys=%s", list(result)) + return [] + return [o for o in orders if isinstance(o, dict) and o.get("order_number")] + + +class OrderDiscoveryCollector: + """网关侧的周期性订单发现器 + + `run()` 是常驻后台循环(与 `_sweep_loop` 并列),每 `_TICK_SECONDS` 跑一次 + `tick()`。tick 做三件事: + + 1. 收割:把 `_pending` 里已到终态的查询单结果消化掉,写进 `account_orders`; + 2. 派列表:距上次 list 扫描超过 `list_interval_seconds` 且没有在飞的 list 单时, + 派一张 order_list 查询(整段账号历史,`max_pages` 由设置给); + 3. 派详情:把「还没成功取过详情」或「详情已过 detail_refresh_seconds」的编目 + 订单各派一张 order_detail 查询。 + + 生命周期由 `GatewayContainer` 持有、`main.lifespan` 起停。被禁用时( + `settings.account_discovery_enabled=False`)`run()` 直接返回,不做事。 + """ + + def __init__( + self, + *, + settings: "Settings", + db: GatewayDB, + query_queue: QueryQueue, + ) -> None: + self._settings = settings + self._db = db + self._queue = query_queue + self._running = False + # 本进程派出去、还没处理结果的 discover-* 单:query_id → kind。 + # 用 dict 而不只存 id,是为了在「是否有一张 list 单还在飞」的判断上不再 + # 靠前缀猜——list 单的 id 是自动生成的 q-...,detail 单是 discover-detail-..., + # 混在一个 set 里光看前缀容易错。 + self._pending: dict[str, str] = {} + # 本进程已为哪些订单派过当天的详情查询(order_number 集合)。防止「详情 + # found=False 导致订单一直算过期、每轮 tick 都重发一张同天详情单」的重复 + # 提交——同日内每笔订单只派一张详情,跨天由《日期分段》的 query_id 自然刷新。 + self._detail_dispatched_this_day: set[str] = set() + self._last_list_submitted: datetime | None = None + + def stop(self) -> None: + self._running = False + + @property + def enabled(self) -> bool: + return self._settings.account_discovery_enabled + + # ---- 主循环 ---- + + async def run(self) -> None: + """常驻循环:被启用时每 _TICK_SECONDS 跑一次 tick""" + if not self.enabled: + logger.info("定时下派通道未启用(RAKUTEN_ACCOUNT_DISCOVERY_ENABLED=false),不启动") + return + self._running = True + logger.info( + "定时下派通道启动:list 扫描间隔 %ss,详情刷新间隔 %ss", + self._settings.account_discovery_interval_seconds, + self._settings.account_detail_refresh_seconds, + ) + while self._running: + try: + await self.tick() + except asyncio.CancelledError: + raise + except Exception: # noqa: BLE001 + # 后台任务不能因为偶发错误退出,否则编目再也不会更新 + logger.exception("定时下派 tick 异常,继续下一轮") + await asyncio.sleep(_TICK_SECONDS) + + # ---- 单轮 ---- + + async def tick(self) -> None: + await self._collect_results() + await self._dispatch_list_if_due() + await self._dispatch_stale_details() + + async def _collect_results(self) -> None: + """把 _pending 里到终态的查询单结果消化进编目""" + pending = list(self._pending) + for query_id in pending: + try: + detail = await self._queue.get_detail(query_id) + except Exception: # noqa: BLE001 (6005:查询单已过保留期被清) + self._pending.pop(query_id, None) + logger.warning("discover 查询单已不存在(可能被清理),丢弃:query_id=%s", query_id) + continue + if detail.status not in _TERMINAL: + continue # 还在飞,下一轮再看 + kind = self._pending.pop(query_id, None) + try: + if detail.status == QueryStatus.SUCCEEDED: + if kind == AccountQueryKind.ORDER_LIST.value: + await self._ingest_list_result(detail.result, query_id=query_id) + elif kind == AccountQueryKind.ORDER_DETAIL.value: + await self._ingest_detail_result(detail.result) + else: + logger.warning( + "discover 查询单未成功:query_id=%s kind=%s status=%s error=%s", + query_id, kind, detail.status, + (detail.error.message if detail.error else ""), + ) + except Exception: # noqa: BLE001 + logger.exception("消化 discover 查询单结果失败:query_id=%s", query_id) + + async def _ingest_list_result(self, result: dict[str, Any] | None, *, query_id: str) -> None: + """把 order_list 结果里每笔订单 upsert 进编目(回写网关) + + 这是「采集到网关中不存在的订单就写回」的核心。只写规范化字段;订单是否 + 真的覆盖完整窗口已由 worker 的 `window_fully_covered` 标明,否则这里到的 + 只是部分窗口——仍然先沉淀拿到的那部分,日志里注明。 + """ + orders = _parse_list_result(result) + seen = 0 + for entry in orders: + order_number = _coerce_str(entry.get("order_number")) or "" + existing = await self._db.get_account_order(order_number) if order_number else None + order = AccountOrderRow( + order_number=order_number, + shop_id=_coerce_str(entry.get("shop_id")), + shop_name=_coerce_str(entry.get("shop_name")) or "", + order_date=_coerce_str(entry.get("order_date")), + # 列表结果不带配送阶段/详情时间。订单若已编目过(详情取过),这里 + # 保留既有数据,别拿 None 把详情阶段悄悄抹掉;新订单则为 None,等 + # 后面的 detail 下派来填。 + delivery_status=existing.delivery_status if existing else None, + order_state=existing.order_state if existing else None, + discovered_at=existing.discovered_at if existing else _iso(utcnow()), + detail_fetched_at=existing.detail_fetched_at if existing else None, + last_seen_at=_iso(utcnow()), + ) + inserted = await self._db.upsert_account_order(order) + seen += 1 + if inserted: + logger.info( + "定时下派:发现网关编目中不存在的新订单并回写:order_number=%s shop=%s", + order.order_number, order.shop_name, + ) + covered = bool(result and result.get("window_fully_covered")) + logger.info( + "discover list 收割完成:query_id=%s orders=%s window_fully_covered=%s", + query_id, seen, covered, + ) + + async def _ingest_detail_result(self, result: dict[str, Any] | None) -> None: + """把 order_detail 结果里规范化出的配送阶段写进编目对应订单 + + 只更新 order_number 命中编目的行:detail 查询是我们为编目订单派出的, + 正常情况下一定命中。found=False 是「站点说这单还没反映出来」,不覆盖已有 + 数据,只刷新 detail_fetched_at 免得每轮都重派。 + """ + if not result: + return + order_number = _coerce_str(result.get("order_number")) + if not order_number: + logger.warning("order_detail 结果缺 order_number,跳过:result_keys=%s", list(result)) + return + row = await self._db.get_account_order(order_number) + if row is None: + # 理论上不该发生:detail 单都是为编目订单派出的。可能编目被删或并发 + # 重置,保守处理:只记日志。 + logger.warning("discover detail 拿到了编目不存在的订单,忽略:order_number=%s", order_number) + return + now = _iso(utcnow()) + found = bool(result.get("found")) + # found=False 是「站点说这单还没反映出来」:不动已有的配送阶段/详情时间, + # 免得一次瞬时抖动把编目里已沉淀的状态抹掉(结果里这些字段是 null)。 + upsert = AccountOrderRow( + order_number=row.order_number, + shop_id=row.shop_id, + shop_name=row.shop_name, + order_date=row.order_date, + delivery_status=( + _coerce_str(result.get("delivery_status")) if found else row.delivery_status + ), + order_state=_coerce_str(result.get("order_state")) if found else row.order_state, + discovered_at=row.discovered_at, + detail_fetched_at=now if found else row.detail_fetched_at, + last_seen_at=row.last_seen_at, + ) + await self._db.upsert_account_order(upsert) + if found: + logger.info( + "discover detail 已沉淀:order_number=%s delivery_status=%s order_state=%s", + order_number, upsert.delivery_status, upsert.order_state, + ) + + # ---- 派单 ---- + + async def _dispatch_list_if_due(self) -> None: + """距上次 list 扫描够久且没有在飞的 list 单时,派一张 order_list 查询""" + interval = self._settings.account_discovery_interval_seconds + due = self._last_list_submitted is None or ( + utcnow() - self._last_list_submitted >= timedelta(seconds=interval) + ) + has_inflight_list = any( + kind == AccountQueryKind.ORDER_LIST.value for kind in self._pending.values() + ) + if not due or has_inflight_list: + return + params = {"max_pages": self._settings.account_discovery_max_pages} + data = await self._queue.submit( + query_id=None, site="rakuten", + kind=AccountQueryKind.ORDER_LIST.value, params=params, + ) + self._last_list_submitted = utcnow() + self._pending[data.query_id] = AccountQueryKind.ORDER_LIST.value + logger.info("定时下派:已派 order_list 扫描,query_id=%s", data.query_id) + + async def _dispatch_stale_details(self) -> None: + """把「没取过详情」或「详情已过刷新期」的编目订单各派一张 order_detail 查询 + + 当天已派过详情的订单跳过(`_detail_dispatched_this_day`),避免因 + found=False 一直算「过期」而每轮 tick 重复提交同一张同天详情单。 + """ + cutoff = _iso(utcnow() - timedelta(seconds=self._settings.account_detail_refresh_seconds)) + stale = await self._db.list_account_orders_needing_detail(cutoff) + for order in stale: + if order.order_number in self._detail_dispatched_this_day: + continue + await self._dispatch_one_detail(order.order_number) + self._detail_dispatched_this_day.add(order.order_number) + + async def _dispatch_one_detail(self, order_number: str) -> None: + """派一张 order_detail 查询;同一天内重复派命中幂等返回既有单 + + 返回的 query_id 无论新旧都放进 _pending,保证它到终态时会被收割—— + 旧单若早已终态,下一轮 _collect_results 收割即丢弃,无副作用。 + """ + # 按天分段,跨天自然刷新,同一天内不重复派 + qid = f"discover-detail-{order_number}-{_day_key(utcnow())}" + data = await self._queue.submit( + query_id=qid, site="rakuten", + kind=AccountQueryKind.ORDER_DETAIL.value, + params={"order_number": order_number}, + ) + self._pending[data.query_id] = AccountQueryKind.ORDER_DETAIL.value + logger.info("定时下派:已派 order_detail,order_number=%s query_id=%s", order_number, data.query_id) + + # ---- 手动触发 ---- + + async def trigger_list_sweep(self) -> dict[str, Any]: + """立即派一轮 order_list 扫描(POST /api/account/discovery/trigger 用) + + 返回 submit 的 data 以表示本轮是否新建。因为不重记 last_list_submitted, + 触发后到下一个整间隔前,定时器不会重复派(仍受「有在飞 list 单就不派」保护)。 + """ + params = {"max_pages": self._settings.account_discovery_max_pages} + data = await self._queue.submit( + query_id=None, site="rakuten", + kind=AccountQueryKind.ORDER_LIST.value, params=params, + ) + self._pending[data.query_id] = AccountQueryKind.ORDER_LIST.value + return { + "query_id": data.query_id, + "status": data.status.value, + "created": data.created, + } + + +def _iso(dt: datetime) -> str: + return dt.astimezone(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z") diff --git a/app/gateway/container.py b/app/gateway/container.py index 7cd0e45..e387705 100644 --- a/app/gateway/container.py +++ b/app/gateway/container.py @@ -3,7 +3,7 @@ from __future__ import annotations import asyncio import logging -from dataclasses import dataclass +from dataclasses import dataclass, field from app.gateway.db import GatewayDB from app.gateway.query_queue import QueryQueue @@ -27,10 +27,16 @@ class GatewayContainer: `sweep_task` 是常驻后台扫描,把过期的 leased/running 推到 stale,并顺带扫 一轮查询单(重投 / 过期 / 清理过保留期的结果)。即使没有 lease 请求,过期的 任务也会被及时发现(规格 §5 的关键约束)。 + + `collector` 是定时下派通道(§12):周期性派 order_list / order_detail 查询、 + 把账户真实订单沉淀进 `account_orders` 编目。`collector_task` 是它的后台任务, + 在 lifespan 里与 `sweep_task` 一起启动/取消。 """ settings: Settings db: GatewayDB task_queue: TaskQueue query_queue: QueryQueue + collector: object | None = field(default=None) # app.gateway.collector.OrderDiscoveryCollector sweep_task: asyncio.Task | None = None + collector_task: asyncio.Task | None = None diff --git a/app/gateway/db.py b/app/gateway/db.py index 1e8c5a5..d35ea1c 100644 --- a/app/gateway/db.py +++ b/app/gateway/db.py @@ -71,6 +71,23 @@ CREATE TABLE IF NOT EXISTS account_queries ( completed_at TEXT ); CREATE INDEX IF NOT EXISTS idx_queries_pending ON account_queries(status, created_at); + +-- 定时下派通道(docs/order-gateway.md §12):账号里真实订单的「已发现」目录。 +-- collector 周期性派 order_list / order_detail 查询,把 worker 回结果里的 +-- **规范化字段**(订单号/店铺/日期/配送状态)沉淀到这里——这是网关自己的编目, +-- 不是上游查询单的原样透传。只消费规范化契约,绝不碰结果里的站点原始 JSON。 +CREATE TABLE IF NOT EXISTS account_orders ( + order_number TEXT PRIMARY KEY, + shop_id TEXT, + shop_name TEXT, + order_date TEXT, -- 站点 ISO8601 字符串 + delivery_status TEXT, -- 最近一次 detail 查询带回的站点码(如 CHECKING_ORDER) + order_state TEXT, -- 映射后的 OrderState(如 delivered) + discovered_at TEXT NOT NULL, -- 首次经 account 采集到该订单的时间 + detail_fetched_at TEXT, -- 最近一次成功取到详情页的时间 + last_seen_at TEXT NOT NULL -- 最近一次出现在订单列表里的时间 +); +CREATE INDEX IF NOT EXISTS idx_account_orders_last_seen ON account_orders(last_seen_at); """ @@ -143,6 +160,21 @@ class QueryRow: return json.loads(self.result_json) if self.result_json else None +@dataclass(slots=True) +class AccountOrderRow: + """account_orders 表一行的强类型视图(定时下派通道的编目,见 §12)""" + + order_number: str + shop_id: str | None + shop_name: str + order_date: str | None + delivery_status: str | None + order_state: str | None + discovered_at: str + detail_fetched_at: str | None + last_seen_at: str + + def _row_to_task(row: aiosqlite.Row) -> TaskRow: return TaskRow( task_id=row["task_id"], @@ -193,6 +225,20 @@ def _row_to_query(row: aiosqlite.Row) -> QueryRow: ) +def _row_to_account_order(row: aiosqlite.Row) -> AccountOrderRow: + return AccountOrderRow( + order_number=row["order_number"], + shop_id=row["shop_id"], + shop_name=row["shop_name"] or "", + order_date=row["order_date"], + delivery_status=row["delivery_status"], + order_state=row["order_state"], + discovered_at=row["discovered_at"], + detail_fetched_at=row["detail_fetched_at"], + last_seen_at=row["last_seen_at"], + ) + + class GatewayDB: """网关 SQLite 访问对象 @@ -577,3 +623,89 @@ class GatewayDB: ) await self.conn.commit() return cursor.rowcount or 0 + + # ---- account_orders(定时下派通道的编目,见 docs/order-gateway.md §12)---- + + async def upsert_account_order(self, row: AccountOrderRow) -> bool: + """把一笔账号订单写进编目(回写)。存在则刷新,返回 True=新写入 + + 这正是「采集到网关中不存在的订单就写回」的落点:主键是 order_number, + 重复采集同一笔订单只是刷新 last_seen 等字段,不会产生重复行。 + """ + async with self.conn.execute( + "SELECT 1 FROM account_orders WHERE order_number = ?", (row.order_number,) + ) as cur: + existed = await cur.fetchone() is not None + await self.conn.execute( + "INSERT INTO account_orders (order_number, shop_id, shop_name, order_date, " + "delivery_status, order_state, discovered_at, detail_fetched_at, last_seen_at) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) " + "ON CONFLICT(order_number) DO UPDATE SET " + "shop_id = excluded.shop_id, shop_name = excluded.shop_name, " + "order_date = excluded.order_date, delivery_status = excluded.delivery_status, " + "order_state = excluded.order_state, detail_fetched_at = excluded.detail_fetched_at, " + "last_seen_at = excluded.last_seen_at", + ( + row.order_number, + row.shop_id, + row.shop_name, + row.order_date, + row.delivery_status, + row.order_state, + row.discovered_at, + row.detail_fetched_at, + row.last_seen_at, + ), + ) + await self.conn.commit() + return not existed + + async def get_account_order(self, order_number: str) -> AccountOrderRow | None: + async with self.conn.execute( + "SELECT * FROM account_orders WHERE order_number = ?", (order_number,) + ) as cur: + row = await cur.fetchone() + return _row_to_account_order(row) if row else None + + async def list_account_orders( + self, + *, + state: str | None = None, + limit: int = 50, + offset: int = 0, + ) -> tuple[list[AccountOrderRow], int]: + """分页列出编目订单,按 last_seen(最近出现)倒序""" + where = [] + params: list[Any] = [] + if state: + where.append("order_state = ?") + params.append(state) + clause = f"WHERE {' AND '.join(where)}" if where else "" + + async with self.conn.execute( + f"SELECT COUNT(*) FROM account_orders {clause}", params + ) as cur: + total = (await cur.fetchone())[0] + + sql = ( + f"SELECT * FROM account_orders {clause} " + "ORDER BY last_seen_at DESC LIMIT ? OFFSET ?" + ) + async with self.conn.execute(sql, [*params, limit, offset]) as cur: + rows = await cur.fetchall() + return [_row_to_account_order(r) for r in rows], total + + async def list_account_orders_needing_detail(self, cutoff_iso: str) -> list[AccountOrderRow]: + """取「还没成功取过详情」或「详情已过刷新期(detail_fetched_at < cutoff)」的编目订单 + + collector 据此派新一批 order_detail 查询。校验用 detail_fetched_at + IS NULL OR < cutoff,保证不会对同一笔订单反复派(详情查询结果回来后才 + 落 detail_fetched_at)。 + """ + async with self.conn.execute( + "SELECT * FROM account_orders WHERE detail_fetched_at IS NULL " + "OR detail_fetched_at < ? ORDER BY last_seen_at ASC", + (cutoff_iso,), + ) as cur: + rows = await cur.fetchall() + return [_row_to_account_order(r) for r in rows] diff --git a/app/gateway/main.py b/app/gateway/main.py index d262379..0273429 100644 --- a/app/gateway/main.py +++ b/app/gateway/main.py @@ -14,9 +14,11 @@ from contextlib import asynccontextmanager from fastapi import FastAPI +from app.gateway.api.routes.account import router as account_router 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.collector import OrderDiscoveryCollector from app.gateway.container import GatewayContainer from app.gateway.db import GatewayDB from app.gateway.query_queue import QueryQueue @@ -46,8 +48,12 @@ def build_container() -> GatewayContainer: max_attempts=settings.query_max_attempts, retention_seconds=settings.query_retention_seconds, ) + collector = OrderDiscoveryCollector( + settings=settings, db=db, query_queue=query_queue + ) return GatewayContainer( - settings=settings, db=db, task_queue=task_queue, query_queue=query_queue + settings=settings, db=db, task_queue=task_queue, + query_queue=query_queue, collector=collector, ) @@ -102,10 +108,22 @@ async def lifespan(app: FastAPI): container.sweep_task = asyncio.create_task( _sweep_loop(container), name="gateway-sweep" ) + # 定时下派通道(§12):被启用时 collector.run() 里自带常驻循环,禁用则直接返回 + assert container.collector is not None + container.collector_task = asyncio.create_task( + container.collector.run(), name="gateway-collector" + ) try: yield finally: + if container.collector_task is not None: + container.collector_task.cancel() + try: + await container.collector_task + except asyncio.CancelledError: + pass + container.collector_task = None if container.sweep_task is not None: container.sweep_task.cancel() try: @@ -122,6 +140,7 @@ def create_app() -> FastAPI: app.include_router(health_router) app.include_router(orders_router) app.include_router(queries_router) + app.include_router(account_router) register_exception_handlers(app) return app diff --git a/app/gateway/models.py b/app/gateway/models.py index c8a5958..b26a902 100644 --- a/app/gateway/models.py +++ b/app/gateway/models.py @@ -260,6 +260,46 @@ class QueryListData(BaseModel): offset: int +# ---- 定时下派通道:账号订单编目(docs/order-gateway.md §12)---- + + +class CatalogedOrder(BaseModel): + """account_orders 编目的单笔订单 + + 这是网关 collector 收到 worker 回结果后**自己缩写出来的目录行**(只读规范化 + 字段:订单号/店铺/日期/配送状态),不是查询单 result 的原样透传——`params` + 与 `result` 仍保持「网关不解释内容」,唯独 collector 消费的是 worker 那边 + 已经规范化了的稳定契约字段(见 §12 说明)。 + """ + + order_number: str + shop_id: str | None = None + shop_name: str = "" + order_date: str | None = None + delivery_status: str | None = None + order_state: str | None = None + discovered_at: str + detail_fetched_at: str | None = None + last_seen_at: str + + +class CatalogOrderListData(BaseModel): + """编目订单列表(GET /api/account/orders)""" + + items: list[CatalogedOrder] + total: int + limit: int + offset: int + + +class TriggerDiscoveryData(BaseModel): + """手动触发一轮下派的响应""" + + query_id: str # 本轮 order_list 下派的单号 + status: QueryStatus + created: bool # 本轮扫描是否新建了一张 order_list 单(False=已有一张在飞/已入队) + + # ---- GET /health ---- diff --git a/app/shared/config.py b/app/shared/config.py index 5523b2a..1bb33df 100644 --- a/app/shared/config.py +++ b/app/shared/config.py @@ -174,6 +174,19 @@ class Settings(BaseSettings): # 规范化字段并标记 raw_omitted——避免把几 MB 的页面状态塞进网关 SQLite。 query_result_max_bytes: int = 1024 * 1024 + # ---- 定时下派通道(仅网关进程 app.gateway.main 使用,见 docs/order-gateway.md §12)---- + # 网关周期性派 order_list / order_detail 查询、把账号真实订单沉淀进编目 + # account_orders 表(回写网关)。关闭则 collector 完全不启动。 + account_discovery_enabled: bool = True + # 两轮 order_list 扫描的间隔(秒)。间隔也是「该订单详情何时到期需刷新」的 + # 参照之一,默认 6 小时一次列表、24 小时刷新一次单笔详情,符合「周期性盘点 + # 账号」的节奏,也不至于频繁戳订单页触发风控。 + account_discovery_interval_seconds: int = 6 * 3600 + account_discovery_max_pages: int = 3 + # 已编目订单的详情刷新间隔(秒)。detail_fetched_at 距今超过这个值就再派一张 + # order_detail 查询更新配送阶段。 + account_detail_refresh_seconds: int = 24 * 3600 + # ---- 本地下单 worker(仅交易服务内的 worker 子模块使用)---- # 网关 URL。**留空则不启动 worker**,交易服务只跑登录态接口。 # 部署形态:本地机(NAT 后无公网入口)通过出站长轮询领任务,详见 diff --git a/app/shared/errors.py b/app/shared/errors.py index 821e5ea..ed537db 100644 --- a/app/shared/errors.py +++ b/app/shared/errors.py @@ -277,3 +277,21 @@ class QueryLeaseInvalidError(AppError): retryable=True, status_code=409, ) + + +class CatalogOrderNotFoundError(AppError): + """编目订单不存在:定时下派通道的 account_orders 里没有这笔订单 + + 沿用 6005(与查询单同号的「不存在」语义),区别只在对象是编目里的订单而非 + 查询单,便于上游按同一张错误码表分支。 + """ + + def __init__(self, order_number: str): + super().__init__( + message=f"编目订单不存在:{order_number}", + code="CATALOG_ORDER_NOT_FOUND", + err_code=6005, + retryable=False, + status_code=404, + ) + self.order_number = order_number diff --git a/docs/order-gateway.md b/docs/order-gateway.md index 4f5348f..108b0a8 100644 --- a/docs/order-gateway.md +++ b/docs/order-gateway.md @@ -503,3 +503,69 @@ worker 回结果撞上 6006 属于正常情况(上一轮超时、单子已重 - [x] 执行超时按失败回报,不把查询单挂到租约过期 - [x] worker 任何异常都不逃出主循环,每张单都有一次回报 +## 12. 定时下派通道:周期性发现账号订单并回写网关(2026-08-16) + +### 12.1 解决什么 + +§11 的查询通道要拿订单详情,得**上游自己已经知道 order_number 再下派**。可账号在 +站点上可能走了别的渠道下单、被商家改状态,网关的状态镜像根本看不到。本通道让 +网关**自己定期盘点账号**:派一张 `order_list` 查询(走 §11 那条队列,本地 worker +出站真读),把「网关编目里还没有的订单」**回写进新的 `account_orders` 表**,再逐笔 +派 `order_detail` 查询把配送阶段沉淀进去。之后上游直接查 `GET /api/account/orders` +就能拿到账号里真实有哪些订单,不必自己记单号。 + +一句话:**gateway 从「等上游告诉它账号里有什么」变成「自己定时去账号里看,并把 +新订单写回自己」**。 + +### 12.2 实现与边界(最重要的两条) + +**1. collector 只消费「规范化字段」,绝不解析站点原始 JSON。**「网关不解释 +query 的 params/result」这条原则照旧适用于上游自提的查询单;collector 自主派发的 +这批 `discover-*` 单是**网关自己的编排数据**,它们的结果里 worker 已经把 +`orders[].order_number/order_date/shop_id/shop_name` 与 `delivery_status/order_state` +抽成了稳定契约字段(§11.3 明说这是「规范化字段是稳定契约」),collector 把它 +沉淀进编目不算解释站点内容。代码见 `app/gateway/collector.py::_parse_list_result`( +只读这几个键)。 + +**2. 状态无痕(不新增任何编排跟踪表)。** collector 只在内存记「本进程派出去 +但还没处理结果」的 `_pending`;进程重启即清空,重启后下一轮 tick 派一张全新 list +扫描,结果幂等 upsert 进编目,没有重复行。`order_detail` 的 query_id 用 +`discover-detail--<日期>` **按天分段**:同一天内重复 submit 命中幂等 +返回既有单(不会因每轮 tick 都满足「详情过期」而反复创建),跨天自然生成新单刷新 +详情;同一天内已派过详情的订单由 `_detail_dispatched_this_day` 在内存里挡掉,避免 +「found=False 订单一直算过期、每轮都重发」。已消费的 `discover-*` 查询单留在 +`account_queries`,由既有 sweep 按保留期(默认 7 天)清理。 + +### 12.3 接口 + +- `GET /api/account/orders?state=&limit=&offset=` — 列出编目订单,按最近出现倒序。 + 编目行只含规范化字段(订单号/店铺/日期/配送状态),要看细节仍去拿对应订单的 + `order_detail` 查询单。 +- `GET /api/account/orders/{order_number}` — 单笔编目订单;编目里没有返回 6005。 +- `POST /api/account/discovery/trigger` — 手动立即派一轮 `order_list` 扫描 + (不等定时器到点),返回本轮下派的单号与 `created`。 + +collector 本身不回写「已经存在于下单任务表 `tasks` 的订单」——`account_orders` 是 +**账号真实订单的编目**,与 `tasks`(下单意图及其镜像)是两回事,两者通过 +`order_number` 天然对账(上游要确认「这单网上下了没」就是查 `tasks`,要「账号里 +实际有什么」就是查 `account_orders`)。 + +### 12.4 配置项 + +| 配置 | 默认 | 说明 | +| --- | --- | --- | +| `RAKUTEN_ACCOUNT_DISCOVERY_ENABLED` | true | 关闭则 collector 完全不启动 | +| `RAKUTEN_ACCOUNT_DISCOVERY_INTERVAL_SECONDS` | 21600 | 两轮 `order_list` 扫描间隔(6h) | +| `RAKUTEN_ACCOUNT_DISCOVERY_MAX_PAGES` | 3 | list 扫描翻页上限 | +| `RAKUTEN_ACCOUNT_DETAIL_REFRESH_SECONDS` | 86400 | 单笔详情刷新间隔(24h) | + +### 12.5 验收清单 + +- [x] 定时派 `order_list` / `order_detail`,结果落到 `account_orders` 编目 +- [x] 采集到编目中不存在的订单 → 回写新行;重复采集只刷新不重复 +- [x] 列表扫描不抹掉已取过的详情阶段(配送状态/详情时间保留) +- [x] 同一天内同一订单的详情单只派一张;跨天自然刷新详情 +- [x] 手动 trigger 立即派一轮 list,返回的单在查询队列里可见 +- [x] `GET /api/account/orders` 列编目;单笔不存在返回 6005 +- [x] collector 只读规范化字段,不解析 raw/raw_pages 站点原始 JSON + diff --git a/tests/test_gateway_api.py b/tests/test_gateway_api.py index eaf4cf7..1ed43f1 100644 --- a/tests/test_gateway_api.py +++ b/tests/test_gateway_api.py @@ -29,6 +29,9 @@ def gateway_client(tmp_path: Path, monkeypatch): """ db_path = tmp_path / "gw.db" monkeypatch.setenv("RAKUTEN_GATEWAY_DB_PATH", str(db_path)) + # 定时下派通道(§12)会在一启动就派 order_list 扫描,污染查询队列语义测试, + # 这类「队列/租赁/鉴权」用例统一关掉它(它有自己的专门测试文件)。 + monkeypatch.setenv("RAKUTEN_ACCOUNT_DISCOVERY_ENABLED", "false") get_settings.cache_clear() try: from app.gateway.main import create_app diff --git a/tests/test_gateway_discovery.py b/tests/test_gateway_discovery.py new file mode 100644 index 0000000..39eea97 --- /dev/null +++ b/tests/test_gateway_discovery.py @@ -0,0 +1,290 @@ +"""定时下派通道测试(docs/order-gateway.md §12) + +collector 靠 QueryQueue + GatewayDB 真实装配驱动(与 test_gateway_queries.py 同一套 +思路),用「真实队列 + 假 worker 回结果」来模拟完整闭环: + + 派 order_list → lease → worker 回带订单的 result → collector 收割 → 写进 + account_orders 编目(回写网关)→ 派 order_detail → 回 detail → 更新配送阶段 + +跑测试时关掉 collector 的常驻循环,只驱动单次 tick()/各方法,避免引入 sleep。 +""" +from __future__ import annotations + +from pathlib import Path + +import pytest + +from app.gateway.collector import OrderDiscoveryCollector, _parse_list_result +from app.gateway.db import GatewayDB +from app.gateway.query_queue import QueryQueue +from app.shared.config import get_settings + +TOKEN = get_settings().bearer_token +AUTH = {"Authorization": f"Bearer {TOKEN}"} + +ORDER = { + "order_number": "306087-20260813-0863947697", + "order_date": "2026-08-13T09:47:28.000Z", + "shop_id": 306087, + "shop_name": "テスト店舗", + "items": [], +} + + +@pytest.fixture +async def collector(tmp_path: Path, monkeypatch): + """真实装配的 (collector, db, queue) 三元组,collector 的常驻循环不启动""" + monkeypatch.setenv("RAKUTEN_ACCOUNT_DISCOVERY_ENABLED", "true") + get_settings.cache_clear() + db = GatewayDB(tmp_path / "c.db") + await db.start() + queue = QueryQueue( + db, lease_ttl_seconds=60, query_ttl_seconds=600, + max_attempts=3, retention_seconds=7 * 24 * 3600, + ) + c = OrderDiscoveryCollector( + settings=get_settings(), db=db, query_queue=queue, + ) + try: + yield c, db, queue + finally: + await db.close() + get_settings.cache_clear() + + +async def _complete_query(queue, query_id: str, result: dict, *, success: bool = True): + """模拟 worker:lease 这张单并回一个结果""" + leased = await queue.lease(worker_id="w-fake", wait=0, site=None, max_wait=60) + assert leased.query_id == query_id + await queue.submit_result( + query_id, worker_id="w-fake", success=success, result=result, + error_code=None if success else 1234, + error_message="" if success else "boom", + ) + + +# ---- 纯函数:从列表结果里抽规范化订单 ---- + +def test_parse_list_result_extracts_valid_rows(): + result = { + "orders": [ + {**ORDER}, + {"order_number": "x", "order_date": "t"}, + {}, # 缺 order_number,应被丢弃 + "not-a-dict", # 非 dict,应被丢弃 + ] + } + rows = _parse_list_result(result) + assert len(rows) == 2 + assert rows[0]["order_number"] == ORDER["order_number"] + + +def test_parse_list_result_empty_on_no_orders(): + assert _parse_list_result(None) == [] + assert _parse_list_result({}) == [] + assert _parse_list_result({"orders": "oops"}) == [] + + +# ---- 回写:采集到的订单写进编目 ---- + +async def test_ingest_list_result_writes_back_new_order(collector): + c, db, queue = collector + await c._queue.submit(query_id="q1", site="rakuten", kind="order_list", params={}) + # 模拟 collector 已派这张单并在等结果(_collect_results 只处理 _pending 里的单) + c._pending["q1"] = "order_list" + await _complete_query(queue, "q1", {"orders": [ORDER], "window_fully_covered": True}) + + await c._collect_results() # 收割 q1 的结果,回写编目 + + row = await db.get_account_order(ORDER["order_number"]) + assert row is not None + assert row.shop_name == "テスト店舗" + assert row.delivery_status is None # 列表结果不带配送阶段,回到编目后仍待详情 + assert row.detail_fetched_at is None + + +async def test_reingest_same_order_no_duplicate(collector): + """同一订单反复被采集,编目里只有一行(upsert 以 order_number 为主键)""" + c, db, _ = collector + for _ in range(2): + await c._ingest_list_result({"orders": [ORDER]}, query_id="q") + rows, total = await db.list_account_orders() + assert total == 1 + assert rows[0].order_number == ORDER["order_number"] + + +async def test_list_reingest_preserves_detail_data(collector): + """订单已取过详情后,再来一轮 list 扫描不得把配送阶段/详情时间抹掉""" + c, db, _ = collector + # 先经 detail 沉淀配送阶段 + await c._ingest_list_result({"orders": [ORDER]}, query_id="q") + await c._ingest_detail_result( + { + "order_number": ORDER["order_number"], "found": True, + "delivery_status": "CHECKING_ORDER", "order_state": "created", + } + ) + before = await db.get_account_order(ORDER["order_number"]) + assert before.detail_fetched_at is not None + + # 再来一轮列表扫描 + await c._ingest_list_result({"orders": [ORDER]}, query_id="q2") + + after = await db.get_account_order(ORDER["order_number"]) + assert after.delivery_status == "CHECKING_ORDER" + assert after.detail_fetched_at == before.detail_fetched_at + + +async def test_found_false_detail_does_not_wipe_status(collector): + """详情 found=False(订单还没反映出来)时,不覆盖已沉淀的配送阶段""" + c, db, _ = collector + await c._ingest_list_result({"orders": [ORDER]}, query_id="q") + await c._ingest_detail_result( + { + "order_number": ORDER["order_number"], "found": True, + "delivery_status": "CHECKING_ORDER", "order_state": "created", + } + ) + + # 一瞬抖动:这单暂时查不到(found=False,结果里配送字段为 null) + await c._ingest_detail_result( + {"order_number": ORDER["order_number"], "found": False, + "delivery_status": None, "order_state": None} + ) + + after = await db.get_account_order(ORDER["order_number"]) + assert after.delivery_status == "CHECKING_ORDER" # 没被抹掉 + assert after.detail_fetched_at is not None + + +# ---- 完整闭环:派 list → 消化 → 派 detail → 更新 ---- + +async def test_full_discovery_loop_populates_catalog(collector): + """一次 tick 派 list,回来后被写进编目,且触发 detail 下派与更新""" + c, db, queue = collector + + # 第一次 tick:派出一张 order_list 扫描 + await c._dispatch_list_if_due() + assert c._pending # 有一张在飞 + + # writer worker 领走并回一个有订单的列表 + (qid,) = c._pending + await _complete_query(queue, qid, {"orders": [ORDER], "window_fully_covered": True}) + + # 收割 + 派 stale detail + await c.tick() + + # 编目里有了订单 + row = await db.get_account_order(ORDER["order_number"]) + assert row is not None and row.detail_fetched_at is None + + # detail 下派已产生一张 order_detail 单 + pending_detail = [q for q, k in c._pending.items() if k == "order_detail"] + assert pending_detail + + # worker 回 detail 结果 + dqid = pending_detail[0] + await _complete_query( + queue, dqid, + { + "order_number": ORDER["order_number"], + "found": True, + "delivery_status": "CHECKING_ORDER", + "order_state": "created", + }, + ) + await c._collect_results() + + after = await db.get_account_order(ORDER["order_number"]) + assert after.delivery_status == "CHECKING_ORDER" + assert after.detail_fetched_at is not None + + +async def test_detail_dispatch_is_same_day_idempotent(collector): + """同一天内轮询到同一订单,detail 单只派一张(幂等 submit 命中既有单)""" + c, db, _ = collector + await c._ingest_list_result({"orders": [ORDER]}, query_id="q") + + await c._dispatch_one_detail(ORDER["order_number"]) + first = set(c._pending) + await c._dispatch_one_detail(ORDER["order_number"]) + second = set(c._pending) + + assert first == second + assert len(first) == 1 + + +async def test_stale_detail_not_redispatched_within_day(collector): + """详情 found=False(订单还没反映)时,当天不再重复派同一订单的详情""" + c, db, _ = collector + await c._ingest_list_result({"orders": [ORDER]}, query_id="q") + + # 第一轮:派详情 → 回 found=False → 订单仍算「未取过详情」 + await c._dispatch_stale_details() + (dqid,) = [q for q, k in c._pending.items() if k == "order_detail"] + await _complete_query( + c._queue, dqid, + {"order_number": ORDER["order_number"], "found": False}, + ) + await c._collect_results() + assert (await db.get_account_order(ORDER["order_number"])).detail_fetched_at is None + + # 第二轮:不应再为同一订单派详情(当天已派过) + before = set(c._pending) + await c._dispatch_stale_details() + assert set(c._pending) == before + + +async def test_no_list_redispatch_while_pending(collector): + """有在飞的 list 单时,即使间隔到了也不重复派""" + c, _, _ = collector + await c._dispatch_list_if_due() + assert len(c._pending) == 1 + # 模拟时间流逝:把 last_list_submitted 拨回到很久以前 + c._last_list_submitted = None + await c._dispatch_list_if_due() + # 幂等?list 用的是自动生成的 q- 前缀 id,重新 submit 会再新建—— + # 但「有在飞 list」会挡住。此时 pending 里那张 q- 单还没到终态,不应再派。 + assert len([k for k in c._pending.values() if k == "order_list"]) == 1 + + +# ---- HTTP 层:编目读取与手动触发 ---- + +@pytest.fixture +def gateway_client(tmp_path: Path, monkeypatch): + monkeypatch.setenv("RAKUTEN_GATEWAY_DB_PATH", str(tmp_path / "gw.db")) + monkeypatch.setenv("RAKUTEN_ACCOUNT_DISCOVERY_ENABLED", "false") + get_settings.cache_clear() + try: + from app.gateway.main import create_app + from fastapi.testclient import TestClient + + app = create_app() + with TestClient(app) as client: + yield client + finally: + get_settings.cache_clear() + + +def test_get_cataloged_order_not_found_6005(gateway_client): + response = gateway_client.get("/api/account/orders/nope", headers=AUTH) + assert response.status_code == 404 + assert response.json()["code"] == 6005 + + +def test_trigger_discovery_dispatch_list_sweep(gateway_client): + """手动触发应派一张 order_list 单(走真实队列)""" + # 触发里 collector 真实存在;手动派单不应因 discovery disabled 而失效—— + # 但这个夹具关了 discovery,trigger 只是 submit 一张单,仍应正常。 + body = gateway_client.post("/api/account/discovery/trigger", headers=AUTH).json() + assert body["success"] is True + assert body["data"]["query_id"].startswith("q-") + # 这张单真的进了查询队列 + detail = gateway_client.get(f"/api/account/queries/{body['data']['query_id']}", headers=AUTH).json() + assert detail["data"]["kind"] == "order_list" + + +def test_catalog_endpoints_reject_missing_token(gateway_client): + assert gateway_client.get("/api/account/orders").status_code == 401 + assert gateway_client.get("/api/account/orders/x").status_code == 401 + assert gateway_client.post("/api/account/discovery/trigger").status_code == 401 diff --git a/tests/test_gateway_leasing.py b/tests/test_gateway_leasing.py index d912aae..df5c8c4 100644 --- a/tests/test_gateway_leasing.py +++ b/tests/test_gateway_leasing.py @@ -26,6 +26,7 @@ AUTH = {"Authorization": f"Bearer {TOKEN}"} def gateway_client(tmp_path: Path, monkeypatch): db_path = tmp_path / "gw.db" monkeypatch.setenv("RAKUTEN_GATEWAY_DB_PATH", str(db_path)) + monkeypatch.setenv("RAKUTEN_ACCOUNT_DISCOVERY_ENABLED", "false") # worker_offline_alert_seconds 设小一点,方便测 /health 失联判定 monkeypatch.setenv("RAKUTEN_WORKER_OFFLINE_ALERT_SECONDS", "1") get_settings.cache_clear() diff --git a/tests/test_gateway_queries.py b/tests/test_gateway_queries.py index 6d74922..66c4c68 100644 --- a/tests/test_gateway_queries.py +++ b/tests/test_gateway_queries.py @@ -32,6 +32,7 @@ QUERIES = "/api/account/queries" def gateway_client(tmp_path: Path, monkeypatch): """起一个独立 DB 的网关应用(与 test_gateway_api.py 同一套装配方式)""" monkeypatch.setenv("RAKUTEN_GATEWAY_DB_PATH", str(tmp_path / "gw.db")) + monkeypatch.setenv("RAKUTEN_ACCOUNT_DISCOVERY_ENABLED", "false") get_settings.cache_clear() try: from app.gateway.main import create_app diff --git a/tests/test_gateway_recover.py b/tests/test_gateway_recover.py index 27ab285..c2d73dc 100644 --- a/tests/test_gateway_recover.py +++ b/tests/test_gateway_recover.py @@ -21,6 +21,7 @@ AUTH = {"Authorization": f"Bearer {TOKEN}"} def gateway_client(tmp_path: Path, monkeypatch): db_path = tmp_path / "gw.db" monkeypatch.setenv("RAKUTEN_GATEWAY_DB_PATH", str(db_path)) + monkeypatch.setenv("RAKUTEN_ACCOUNT_DISCOVERY_ENABLED", "false") # TTL 设小一点,让 lease 一过期就能被 sweep 转 stale monkeypatch.setenv("RAKUTEN_LEASE_TTL_SECONDS", "1") get_settings.cache_clear()