新增 §12 定时下派通道:网关自己定期派 order_list / order_detail 查询,把账号
真实订单沉淀进新表 account_orders(回写网关),上游可直接查 GET /api/account/orders
拿到账号里实际有哪些订单,不必自己记 order_number。
- collector.py: OrderDiscoveryCollector 常驻后台任务(与 sweep 并列)。每隔
account_discovery_interval_seconds 派 order_list,收割结果后把「编目里没有的
订单」upsert 进编目,再逐笔派 order_detail 沉淀 delivery_status/order_state。
状态无痕:不新增编排跟踪表,仅内存 _pending + _detail_dispatched_this_day;
detail query_id 按天分段(discover-detail-<order>-<日期>),同天幂等防重复派。
边界:collector 只消费 worker 结果的规范化字段,绝不解析 raw/raw_pages。
- db.py: account_orders 表 + AccountOrderRow + upsert/single/list/stale 访问;
upsert 以 order_number 为主键,重复采集只刷新,列表重扫不抹详情节点。
- 路由: GET /api/account/orders(列编目)、GET /api/account/orders/{n}(6005)、
POST /api/account/discovery/trigger(手动立即派一轮 list)。
- config: RAKUTEN_ACCOUNT_DISCOVERY_ENABLED / _INTERVAL_SECONDS / _MAX_PAGES /
RAKUTEN_ACCOUNT_DETAIL_REFRESH_SECONDS。
- 既有网关测试夹具统一关 discovery(会启动即派扫描污染查询队列语义),
新表测试在 tests/test_gateway_discovery.py(13 用例,真实 QueryQueue+DB 装配)。
472 tests passed.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
354 lines
17 KiB
Python
354 lines
17 KiB
Python
"""定时下派通道:周期性发现账号真实订单并沉淀到网关编目(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-<order_number>-<日期>`
|
|
**按天分段**:同一天内对同一笔订单重复 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")
|