"""账号只读查询 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