"""账号只读查询 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 opentelemetry import trace from opentelemetry.trace import SpanKind from app.shared.errors import AppError, InvalidRequestError from app.shared.task_state import AccountQueryKind from app.shared.telemetry import record_error, set_attributes, suppressed 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__) tracer = trace.get_tracer(__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: # 与下单主循环同样:空转的长轮询不埋点,见 runner.py::run with suppressed(): 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: """执行一张查询单并回报结果。任何失败都转成一次「失败回报」,不抛出去 这里开的 span 是一张查询单的根:底下挂着站点读取(`site.*`)与回报网关的 HTTP 调用,一个 query_id 对应一条 trace。查询失败不抛出去(都转成失败 回报),所以失败信息由各分支显式记到 span 上——否则 trace 里会显示成功。 """ with tracer.start_as_current_span("account_query", kind=SpanKind.CONSUMER) as span: set_attributes( span, { "query.query_id": query.query_id, "query.kind": query.kind, "query.site": query.site, "query.attempt": query.attempt, "query.worker_id": self.worker_id, }, ) await self._handle_traced(query, span) async def _handle_traced(self, query: QueryTask, span: "trace.Span") -> None: """handle() 的实际执行体,拆出来只为让根 span 的 with 块保持一层缩进""" 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 as exc: logger.warning( "查询超时(%s 秒,多半是下单任务正占着账号锁):query_id=%s", self._settings.account_query_timeout_seconds, query.query_id, ) span.set_attribute("query.outcome", "timeout") record_error(span, exc) 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, ) span.set_attribute("query.outcome", "failed") record_error(span, exc) 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) span.set_attribute("query.outcome", "unexpected_error") record_error(span, exc) 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: span.set_attribute("query.outcome", "oversized") await self._report_safe( query, success=False, error_code=1003, error_message=oversized ) return span.set_attribute("query.outcome", "succeeded") 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, # 站点侧结构化配送状态码(如 "CHECKING_ORDER");列表页/stepper 路径 # 拿不到时为 None。这是 raw 里 orderData 的规范化摘要,供上游快速对账。 "delivery_status": status.delivery_status, # 详情页额外带出的站点原始 JSON 原样透传,由上游自担结构变动风险 # (见 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)