From 51e253843889b29062c87e6b41410b1b5cdeadb5 Mon Sep 17 00:00:00 2001 From: Jerry Yan <792602257@qq.com> Date: Thu, 13 Aug 2026 23:01:35 +0800 Subject: [PATCH] =?UTF-8?q?=E5=AE=9E=E7=8E=B0=E4=BB=98=E6=AC=BE=E5=90=8E?= =?UTF-8?q?=E8=AE=A2=E5=8D=95=E7=9B=91=E6=8E=A7=EF=BC=9A=E7=9C=9F=E5=AE=9E?= =?UTF-8?q?=E6=8E=A2=E6=B5=8B=20order.my.rakuten.co.jp=20=E9=85=8D?= =?UTF-8?q?=E9=80=81=E9=98=B6=E6=AE=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 用真实订单号 306087-20260813-0863947697 探测 order.my.rakuten.co.jp(订单列表/ 详情页),拿到真实 DOM 结构后实现 SiteInteractor.check_order_status: - 详情页 URL 可直接从 site_order_id 构造(shop_id 是订单号第一段) - 配送阶段用「进度条」组件的 4 个固定阶段(ショップ/出荷/配達店/配達完了), 当前阶段的 class 带 -active-- 中缀,映射到 OrderState.SHIPPED/DELIVERED - 查不到订单号、进度条解析不出新阶段都不算错误,交给轮询循环继续重试 runner.WorkerRunner 新增 _monitor_order 后台轮询:付款成功上报后以 asyncio.create_task 起后台任务(不阻塞主循环领下一单,因为 SiteInteractor 的 Playwright 操作全程持锁串行化),状态变化时用 terminal=False 追加 report; 新增 cancel_monitors() 在服务关闭时于 SiteInteractor.close() 之前收尾。 Co-Authored-By: Claude Sonnet 5 --- .env.example | 4 + app/shared/config.py | 7 ++ app/trading/main.py | 3 + app/trading/worker/runner.py | 107 ++++++++++++++++++++- app/trading/worker/site_interact.py | 123 ++++++++++++++++++++++-- docs/order-gateway.md | 7 +- tests/test_site_interact.py | 78 +++++++++++++++- tests/test_worker_runner.py | 139 +++++++++++++++++++++++++++- 8 files changed, 450 insertions(+), 18 deletions(-) diff --git a/.env.example b/.env.example index bd89819..eef3bcc 100644 --- a/.env.example +++ b/.env.example @@ -81,6 +81,10 @@ RAKUTEN_AUTH_STATE_DIR=.auth # 下单金额上限(日元):实际应付超过该值直接拒绝提交,防止解析出错或页面改版 # 导致买到远超预期的订单。设为 0 表示不设上限(不建议)。 RAKUTEN_ORDER_MAX_TOTAL_YEN=30000 +# 付款后监控(订单列表页轮询)间隔(秒)与最大轮询次数。间隔不宜太短——同一账号 +# 频繁访问订单页有被风控盯上的风险。默认 3 小时一次,最多 80 次(约 10 天)。 +RAKUTEN_ORDER_MONITOR_POLL_INTERVAL_SECONDS=10800 +RAKUTEN_ORDER_MONITOR_MAX_CHECKS=80 # ---- 以下仅下单任务网关使用 ---- # 任务队列 SQLite 文件路径(相对项目根目录)。务必放在持久化卷上,丢了等于 diff --git a/app/shared/config.py b/app/shared/config.py index 6c6628d..d17adaf 100644 --- a/app/shared/config.py +++ b/app/shared/config.py @@ -111,6 +111,13 @@ class Settings(BaseSettings): # 自动下单时在支付方式页选择的选项文案(按可见文案匹配,非 value/id)。 # 2026-08-11 未经真实确认页验证:见 site_interact.py::_select_payment_method。 order_payment_method: str = "クレジットカード" + # 付款后监控(订单列表页轮询):间隔与最大轮询次数。2026-08-13 用真实订单号 + # 实测过 order.my.rakuten.co.jp 的结构,见 site_interact.py::check_order_status。 + # 间隔不宜太短——同一账号频繁访问订单页有被风控盯上的风险;默认 3 小时一次, + # 最多轮询 80 次(约 10 天,覆盖绝大多数国内配送时长),到期仍未到「配達完了」 + # 就停止(不是失败,只是不再继续追踪,详见 runner.py::_monitor_order)。 + order_monitor_poll_interval_seconds: int = 3 * 3600 + order_monitor_max_checks: int = 80 # ---- 自动重登(仅交易服务,需要 account.yaml)---- # 检测到登录态失效时,是否在 require_logged_in 内自动触发重登。需要项目根 diff --git a/app/trading/main.py b/app/trading/main.py index 43b47b9..a9c28b4 100644 --- a/app/trading/main.py +++ b/app/trading/main.py @@ -130,6 +130,9 @@ async def lifespan(app: FastAPI): await worker_task except asyncio.CancelledError: pass + # 付款后监控是独立于主循环的后台任务(见 runner.py::_spawn_monitor), + # 取消 worker_task 不会连带取消它们,必须在关 SiteInteractor 前单独收尾。 + await container.worker_runner.cancel_monitors() # type: ignore[union-attr] 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/runner.py b/app/trading/worker/runner.py index 5b27d3c..d000dbb 100644 --- a/app/trading/worker/runner.py +++ b/app/trading/worker/runner.py @@ -76,10 +76,28 @@ class WorkerRunner: self._evidence = evidence self._site = site self._running = False + # 付款后监控后台任务:task_id → asyncio.Task。与主循环解耦(不阻塞领下一单), + # 详见 _spawn_monitor / _monitor_order。用 set 而非 list 是因为只需要成员管理, + # 不需要顺序。 + self._monitor_tasks: set[asyncio.Task] = set() def stop(self) -> None: self._running = False + async def cancel_monitors(self) -> None: + """服务关闭时调用:取消所有还在跑的付款后监控后台任务并等它们收尾 + + 必须在 SiteInteractor.close() 之前调用——监控任务还在用同一个 + Playwright context,先关 context 再取消会导致监控任务在 await 到一半时 + 踩上已关闭的资源,报一堆无意义的异常。 + """ + tasks = list(self._monitor_tasks) + for t in tasks: + t.cancel() + for t in tasks: + with contextlib.suppress(asyncio.CancelledError): + await t + @property def worker_id(self) -> str: return self._settings.worker_id_effective @@ -311,11 +329,90 @@ class WorkerRunner: ) await self._db.mark_finished(task.task_id, OrderState.PAID.value) - # 步骤 6:付款后监控(非阻塞,常驻轮询;当前未实现) - try: - await self._site.monitor(task, site_order_id) - except NotImplementedError: - logger.info("付款后监控未实现,跳过:task_id=%s", task.task_id) + # 步骤 6:付款后监控——后台常驻轮询,不阻塞主循环领下一单(规格 §6) + self._spawn_monitor(task, site_order_id) + + def _spawn_monitor(self, task: LeaseTask, site_order_id: str) -> None: + """把 _monitor_order 起成独立的后台任务,登记到 _monitor_tasks 便于关服时收尾""" + monitor_task = asyncio.create_task( + self._monitor_order(task, site_order_id), name=f"monitor-{task.task_id}" + ) + self._monitor_tasks.add(monitor_task) + monitor_task.add_done_callback(self._monitor_tasks.discard) + + async def _monitor_order(self, task: LeaseTask, site_order_id: str) -> None: + """付款后台轮询:订单配送阶段变化时继续 report(规格 §6「付款后监控」) + + 任务本身已经在 execute() 里上报过 succeeded 终态——这里的 report 都是 + `terminal=False` 的追加上报(规格 §4.5:任务已 terminal 的仍可上报, + gateway 追加到 task_reports),监控本身的成败不影响已经完成的下单结果。 + + 轮询间隔/次数上限见 settings.order_monitor_poll_interval_seconds / + order_monitor_max_checks;达到上限仍未看到「配達完了」不算失败,只是 + 停止追踪(记一条 warning)。单次探测异常(登录态失效、页面打不开)只记 + 日志、下一轮重试,不能把异常抛出去——这是后台任务,没有人等着 catch 它, + 未捕获异常会被 asyncio 直接吞掉且只在垃圾回收时打一条难查的警告。 + """ + last_state: OrderState | None = None + step_no = 6 + interval = self._settings.order_monitor_poll_interval_seconds + max_checks = self._settings.order_monitor_max_checks + for attempt in range(1, max_checks + 1): + await asyncio.sleep(interval) + try: + snapshot = await self._site.check_order_status(site_order_id) + except Exception: + logger.warning( + "订单监控探测失败(第 %s/%s 次),下一轮继续重试:" + "task_id=%s site_order_id=%s", + attempt, max_checks, task.task_id, site_order_id, + exc_info=True, + ) + continue + + if not snapshot.found: + logger.info( + "订单监控:第 %s/%s 次仍未在详情页找到订单号,继续轮询:task_id=%s", + attempt, max_checks, task.task_id, + ) + continue + + if snapshot.order_state is None or snapshot.order_state == last_state: + continue + + last_state = snapshot.order_state + step_name = f"monitor-{snapshot.order_state.value}" + detail = f"订单监控:进度「{snapshot.stage_label}」→ {snapshot.order_state.value}" + evidence_ref = self._evidence.write_step( + task.task_id, step_no, step_name, + html=snapshot.html, + meta={ + "step": step_name, + "state": snapshot.order_state.value, + "stage_label": snapshot.stage_label, + }, + ) + await self._db.index_evidence(task.task_id, step_no, step_name, evidence_ref) + await self._db.record_event( + task.task_id, snapshot.order_state.value, + detail=detail, evidence_ref=evidence_ref, + ) + await self._report_safe( + task, + state=snapshot.order_state, + detail=detail, + site_order_id=site_order_id, + ) + step_no += 1 + + if snapshot.order_state == OrderState.DELIVERED: + logger.info("订单监控:已送达,停止轮询:task_id=%s", task.task_id) + return + + logger.warning( + "订单监控达到最大轮询次数 %s 仍未看到「配達完了」,停止追踪:task_id=%s", + max_checks, task.task_id, + ) def _enforce_amount_guard(self, task: LeaseTask, payable_yen: int) -> None: """金额守卫:实际应付超过 intent.max_total_yen 或 RAKUTEN_ORDER_MAX_TOTAL_YEN 时拦截""" diff --git a/app/trading/worker/site_interact.py b/app/trading/worker/site_interact.py index 8042c83..c251b91 100644 --- a/app/trading/worker/site_interact.py +++ b/app/trading/worker/site_interact.py @@ -48,6 +48,14 @@ data/evidence/checkout-research-20260811/NOTES.md): 这次真实提交**只此一单**,其余账号/商品/金额组合下的中间步骤分支(电话补录、 非默认地址、3DS/OTP 拦截等)仍未被真实数据验证过,继续按原有猜测处理并保留 CheckoutBlockedError 兜底,不做无限重试。 +- check_order_status(付款后监控的单次探测,循环轮询在 runner.py)**2026-08-13 + 用真实订单号 306087-20260813-0863947697 实测过** order.my.rakuten.co.jp 的 + 订单列表页与详情页:两者共用同一套「配送阶段进度条」组件,固定 4 阶段 + (ショップ→出荷→配達店→配達完了),当前阶段的 class 里带 `-active--` 中缀, + 详见 _ORDER_STEPPER_ITEM_PATTERN 上方注释。这次真实订单当时还停在「ショップ」 + (刚接单,未发货)这一阶段,「出荷」「配達完了」两个状态转换点没有被真实数据 + 验证过,只是按进度条文案直译映射(_ORDER_STAGE_TO_STATE);取消/退款没有 + 找到可靠信号,检测不到,遇到需要人工核对订单列表页。 **httpx 不能用于带账号的写操作**:Rakuten 对账号操作有 TLS/HTTP2 指纹校验, 同一份 cookie Playwright 能用、httpx 不能。所以本模块全程使用 Playwright @@ -80,6 +88,7 @@ from app.shared.purchase_contract import ( basket_domain_of, inventory_flag_for, ) +from app.shared.task_state import OrderState from app.trading.core import auth_site from app.trading.worker.models import LeaseTask @@ -259,10 +268,32 @@ _PAYMENT_BLOCK_INDICATORS = ( 'input[autocomplete="one-time-code"]', ) -_NOT_IMPLEMENTED_MSG = ( - "站点交互未实现:见 docs/order-gateway.md §10 与 " - "project://jp-rakuten/checkout-flow-probe-findings。" +# ---- 付款后监控:订单列表/详情页轮询(2026-08-13 用真实订单号 +# 306087-20260813-0863947697 实测过,见模块顶部说明)---- +# 详情页 URL 可以直接从 site_order_id 构造,不需要额外查询:shop_id 就是订单号 +# 第一段(本单验证 306087-20260813-... 对应 shop_id=306087)。 +_ORDER_DETAIL_URL_TEMPLATE = ( + "https://order.my.rakuten.co.jp/purchase-history/" + "?order_number={order_number}&shop_id={shop_id}&act=detail_page_view" ) +# 订单列表页与详情页共用同一套「配送阶段进度条」组件:4 个固定阶段的
  • , +# 当前阶段比其余几个多一段形如 `item-shipping-active--{hash}` 的 class +# (hash 是 CSS modules 编译产物,逐次构建会变;`-active--` 这个中缀是本模块 +# 唯一依赖的稳定信号)。用 re.findall 顺序返回 4 个 (class, 阶段文案) 二元组。 +_ORDER_STEPPER_ITEM_PATTERN = re.compile( + r'
  • .*?
    ([^<]*)
  • ', re.DOTALL, +) +_ORDER_STEPPER_ACTIVE_MARKER = "-active--" +# 「ショップ」(店铺已接单,未发货)阶段没有对应的 OrderState 取值——报告过 +# ORDERED/PAID 了,不需要 monitor 再报一次;「配達店」(配送中转)也没有单独 +# 状态,归入 SHIPPED。只有「出荷」→SHIPPED、「配達完了」→DELIVERED 两个转换点 +# 会触发 runner._monitor_order 的 report。真实订单当时仍停在「ショップ」阶段, +# 这两个映射本身未经真实状态转换验证,只是按进度条文案直译。 +_ORDER_STAGE_TO_STATE: dict[str, OrderState] = { + "出荷": OrderState.SHIPPED, + "配達店": OrderState.SHIPPED, + "配達完了": OrderState.DELIVERED, +} @dataclass(slots=True) @@ -277,6 +308,25 @@ class CheckoutSummary: pay_deadline: str | None = None +@dataclass(slots=True) +class OrderStatusSnapshot: + """订单列表/详情页单次探测结果(check_order_status 的返回值) + + found=False 表示页面上没找到这个订单号——**不当错误处理**:站点自己说 + 「ご注文の反映に10分ほどかかります」,下单后短时间内查不到是正常的,调用方 + (runner._monitor_order)应该继续下一轮轮询,不是放弃。 + stage_label 是真实站点用的阶段原文(如「出荷」),供日志/证据留原始信号; + order_state 是映射到 OrderState 的结果,映射不到(比如仍在「ショップ」这个 + 起始阶段,或者进度条没解析出来)时为 None——同样不当错误,只是「这次没有 + 新状态可报」。html 是抓取到的整页内容,供调用方按需落证据。 + """ + + found: bool + stage_label: str | None = None + order_state: OrderState | None = None + html: str = "" + + class SiteInteractor: """Rakuten 站点交互器:持有 Playwright 浏览器 context,复用账号 cookie @@ -1381,9 +1431,48 @@ class SiteInteractor: finally: await page.close() - async def monitor(self, task: LeaseTask, site_order_id: str) -> None: - """付款后监控。**未实现**:订单列表页的真实结构没有实测过""" - raise NotImplementedError(_NOT_IMPLEMENTED_MSG) + async def check_order_status(self, site_order_id: str) -> OrderStatusSnapshot: + """付款后监控的单次探测:查一次订单详情页的配送阶段,不循环 + + 循环轮询(间隔、次数上限、状态变化时 report)由 + `runner.WorkerRunner._monitor_order` 负责——本方法只做一次 + 「导航 + 解析」,站点交互与轮询节奏解耦,也方便离线单测轮询逻辑 + (用桩替换本方法)与解析逻辑(`_parse_order_status`,纯函数)。 + + 与 enter_checkout 不同,本方法开自己的临时 Page 并在返回前关闭—— + submit_order/pay 已经处理完并关闭了 `_checkout_pages` 里留存的会话, + 监控阶段没有需要跨调用复用的页面状态。 + + Returns: + OrderStatusSnapshot;订单号暂时查不到、或进度条解析不出新阶段都 + **不算错误**(详见该 dataclass 文档),由调用方决定是否继续轮询。 + + Raises: + NotLoggedInError: 登录态失效 + OrderOperationError: 订单详情页打开/渲染失败 + """ + shop_id = site_order_id.split("-", 1)[0] + url = _ORDER_DETAIL_URL_TEMPLATE.format(order_number=site_order_id, shop_id=shop_id) + + async with self._lock: + await self._auth_session.require_logged_in("rakuten") + await self._refresh_context_if_stale() + + page = await self._context.new_page() + try: + try: + await page.goto(url, wait_until="domcontentloaded", timeout=30_000) + await page.wait_for_timeout(2_000) + except Exception as exc: + raise OrderOperationError( + f"订单详情页打开失败:site_order_id={site_order_id} " + f"{type(exc).__name__}: {exc}" + ) from exc + html = await page.content() + finally: + await page.close() + + return _parse_order_status(html, site_order_id) # ---- 模块级辅助函数(纯函数,便于单测)---- @@ -1531,3 +1620,25 @@ def _parse_checkout_summary(html: str) -> CheckoutSummary: site_order_id=site_order_id, pay_deadline=pay_deadline, ) + + +def _parse_order_status(html: str, site_order_id: str) -> OrderStatusSnapshot: + """订单列表/详情页解析核心逻辑(纯函数,供 check_order_status 调用,便于离线单测) + + 2026-08-13 用真实订单号 306087-20260813-0863947697 的订单详情页 HTML 验证过: + 见 _ORDER_STEPPER_ITEM_PATTERN 上方注释。找不到订单号、或进度条没解析出 + 「当前阶段」都返回 found/order_state 相应地为 False/None,不抛错——这是 + 「暂时没有新信息」而不是「站点交互失败」,抛错的语义留给导航/渲染失败。 + """ + if site_order_id not in html: + return OrderStatusSnapshot(found=False, html=html) + + for cls, label in _ORDER_STEPPER_ITEM_PATTERN.findall(html): + if _ORDER_STEPPER_ACTIVE_MARKER in cls: + return OrderStatusSnapshot( + found=True, + stage_label=label, + order_state=_ORDER_STAGE_TO_STATE.get(label), + html=html, + ) + return OrderStatusSnapshot(found=True, html=html) diff --git a/docs/order-gateway.md b/docs/order-gateway.md index 7927f31..ba36bc7 100644 --- a/docs/order-gateway.md +++ b/docs/order-gateway.md @@ -324,7 +324,12 @@ trading 侧新增: 1. 自动付款是否触发 3D Secure 或短信验证。若触发,这条路走不通,付款环节改为 「下单到 `awaiting_payment` + 上报 `needs_human` 交人工」,其余环节不变。 2. 下单确认页的实际应付金额、付款方式、付款期限、站点订单号各自在哪个字段。 -3. 订单列表页能否按商品 + 时间窗口可靠地反查出「这单下没下」(§5 的恢复核对依赖它)。 +3. 订单列表页能否按商品 + 时间窗口可靠地反查出「这单下没下」(§5 的恢复核对依赖它, + `verify.verify_on_site` 仍是恒返回 unknown 的桩)。2026-08-13 已经拿到过 + `order.my.rakuten.co.jp` 的真实 HTML(用于实现付款后监控,见 + `site_interact.py::check_order_status` / `_parse_order_status`),页面按 + 订单号能查到「注文番号」「配送阶段进度条」,但没有验证过按商品名/时间窗口 + 反查、也没有验证过多页/分页场景——对实现 §5 恢复核对有参考价值,不能直接复用。 > **交易范围**:交易服务只覆盖乐天市场(rakuten)。ラクマ 的抓取仍在抓取服务里提供, > 但不进入交易链路——没有加购契约,也永不实现下单/付款/订单监控。 diff --git a/tests/test_site_interact.py b/tests/test_site_interact.py index 03ee93f..65e45a4 100644 --- a/tests/test_site_interact.py +++ b/tests/test_site_interact.py @@ -17,13 +17,16 @@ import pytest from app.shared.errors import CartOperationError, InvalidRequestError, OrderOperationError from app.trading.worker.models import LeaseTask +from app.shared.task_state import OrderState from app.trading.worker.site_interact import ( CheckoutSummary, + OrderStatusSnapshot, SiteInteractor, _extract_error_message, _extract_purchase_fields, _parse_checkout_summary, _parse_initial_state, + _parse_order_status, ) FIXTURES = Path(__file__).parent / "fixtures" @@ -250,13 +253,80 @@ def test_site_interactor_construction_does_not_require_playwright(): assert site._per_task_state == {} -# ---- monitor 仍未实现,抛 NotImplementedError ---- +# ---- _parse_order_status:付款后监控的解析核心,2026-08-13 用真实订单号 +# 306087-20260813-0863947697 的详情页 HTML 验证过(见 _ORDER_STEPPER_ITEM_PATTERN +# 上方注释)。下面 fixture 里的 4 个
  • 是从真实页面摘录的进度条结构原样保留, +# 只是把每个 class 里易变的 CSS modules hash 换成了固定占位,不影响所依赖的 +# `-active--` 中缀。 + +_REAL_ORDER_ID = "306087-20260813-0863947697" -async def test_monitor_unimplemented(): +def _stepper_html(*, active_stage: str | None, order_id: str = _REAL_ORDER_ID) -> str: + """构造「订单号 + 4 阶段进度条」的最小 fixture,active_stage 指定哪一阶带 -active-- class""" + stages = ["ショップ", "出荷", "配達店", "配達完了"] + items = [] + for i, stage in enumerate(stages): + active_cls = " item-shipping-active--1mu0i" if stage == active_stage else "" + items.append( + f'
  • ' + f'
    {stage}
  • ' + ) + return f"
    注文番号:{order_id}
    " + + +def test_parse_order_status_not_found_when_order_id_missing(): + snapshot = _parse_order_status(_stepper_html(active_stage="ショップ", order_id="999999-x"), _REAL_ORDER_ID) + assert snapshot.found is False + assert snapshot.order_state is None + + +def test_parse_order_status_shop_stage_has_no_order_state_mapping(): + """「ショップ」(已接单未发货)阶段没有对应的 OrderState——已经在 ORDERED/PAID 报过了""" + snapshot = _parse_order_status(_stepper_html(active_stage="ショップ"), _REAL_ORDER_ID) + assert snapshot.found is True + assert snapshot.stage_label == "ショップ" + assert snapshot.order_state is None + + +def test_parse_order_status_shipped_stage_maps_to_shipped(): + snapshot = _parse_order_status(_stepper_html(active_stage="出荷"), _REAL_ORDER_ID) + assert snapshot.order_state == OrderState.SHIPPED + + +def test_parse_order_status_depot_stage_also_maps_to_shipped(): + """「配達店」(配送网点中转)没有单独状态,归入 SHIPPED""" + snapshot = _parse_order_status(_stepper_html(active_stage="配達店"), _REAL_ORDER_ID) + assert snapshot.order_state == OrderState.SHIPPED + + +def test_parse_order_status_delivered_stage_maps_to_delivered(): + snapshot = _parse_order_status(_stepper_html(active_stage="配達完了"), _REAL_ORDER_ID) + assert snapshot.order_state == OrderState.DELIVERED + + +def test_parse_order_status_no_active_marker_returns_found_without_state(): + """进度条 4 项都没有 -active-- class(页面结构变了/解析不出):found 但 order_state=None""" + snapshot = _parse_order_status(_stepper_html(active_stage=None), _REAL_ORDER_ID) + assert snapshot.found is True + assert snapshot.order_state is None + assert snapshot.stage_label is None + + +def test_order_status_snapshot_defaults(): + s = OrderStatusSnapshot(found=False) + assert s.stage_label is None + assert s.order_state is None + assert s.html == "" + + +# ---- check_order_status 在没启动 Playwright 时应失败 ---- + + +async def test_check_order_status_without_start_raises(): site = SiteInteractor(auth_session=None, settings=None) # type: ignore[arg-type] - with pytest.raises(NotImplementedError): - await site.monitor(_make_task(), "ord-1") + with pytest.raises((AttributeError, TypeError)): + await site.check_order_status(_REAL_ORDER_ID) # ---- submit_order / pay:没有 enter_checkout 留存的确认页会话时应报错,不静默成功 ---- diff --git a/tests/test_worker_runner.py b/tests/test_worker_runner.py index 55a2487..b6fed54 100644 --- a/tests/test_worker_runner.py +++ b/tests/test_worker_runner.py @@ -12,6 +12,7 @@ evidence store / site,覆盖 §6 主循环的分支: """ from __future__ import annotations +import asyncio from dataclasses import dataclass, field from pathlib import Path from typing import Any @@ -25,7 +26,7 @@ from app.trading.worker.evidence import EvidenceStore from app.trading.worker.local_db import LocalDB from app.trading.worker.models import LeaseTask from app.trading.worker.runner import WorkerRunner -from app.trading.worker.site_interact import CheckoutSummary, SiteInteractor +from app.trading.worker.site_interact import CheckoutSummary, OrderStatusSnapshot, SiteInteractor # ---- 桩:网关客户端 ---- @@ -123,6 +124,10 @@ def runner(tmp_path, local_db, evidence, site) -> WorkerRunner: class _FakeSettings: worker_id = "test" order_max_total_yen = 30000 + # 测试环境不真的等 3 小时:间隔设 0,次数设小,_monitor_order 的循环 + # 靠 asyncio.sleep(0) 立刻推进,不拖慢测试 + order_monitor_poll_interval_seconds = 0 + order_monitor_max_checks = 5 @property def worker_id_effective(self) -> str: @@ -234,7 +239,7 @@ async def test_unimplemented_site_interaction_becomes_needs_human( """站点交互抛 NotImplementedError → runner 转 needs_human add_to_cart / verify_cart / enter_checkout / parse_checkout / submit_order / - pay 现在均已实现(monitor 仍未实现),这里用桩显式模拟「某一步没实现」, + pay / check_order_status 现在均已实现,这里用桩显式模拟「某一步没实现」, 覆盖 runner._execute_with_renewal 里 NotImplementedError → needs_human 的分支, 与具体哪个方法真的未实现解耦。 """ @@ -432,3 +437,133 @@ async def test_evidence_files_exist_before_each_report( # 至少有一个 step 调了 report,且每次 report 之前证据都在 assert seen_evidence_at_report, "应当至少有一次带 evidence_ref 的 report" assert all(seen_evidence_at_report), "某次 report 之前证据文件未落盘" + + +# ---- 付款后监控(规格 §6「付款后监控」)---- + + +async def test_successful_order_spawns_monitor_task_and_cancel_monitors_cleans_up( + runner: WorkerRunner, local_db: LocalDB +): + """付款成功后 execute() 应起一个后台监控任务、不阻塞返回;cancel_monitors 能干净收尾 + + check_order_status 桩故意永不返回(挂在一个手动控制的 Event 上),模拟「监控任务 + 还在跑」这个可观察状态——用来确认 execute() 没有同步 await 监控循环 + (不然 handle() 早该被这个永不返回的调用卡死,测试会超时而不是正常结束)。 + """ + + async def _noop(task): # noqa: ANN001 + return None + + async def _checkout_html(task): # noqa: ANN001 + return "checkout" + + async def _parse(html: str): + return CheckoutSummary(payable_yen=297) + + async def _submit(task): # noqa: ANN001 + return "306087-20260813-0863947697" + + async def _pay(task, site_order_id): # noqa: ANN001 + return None + + never_resolves = asyncio.Event() + + async def _hang_check(order_id: str): # noqa: ANN001 + await never_resolves.wait() + + runner._site.add_to_cart = _noop # type: ignore[assignment] + runner._site.verify_cart = _noop # type: ignore[assignment] + runner._site.enter_checkout = _checkout_html # type: ignore[assignment] + runner._site.parse_checkout = _parse # type: ignore[assignment] + runner._site.submit_order = _submit # type: ignore[assignment] + runner._site.pay = _pay # type: ignore[assignment] + runner._site.check_order_status = _hang_check # type: ignore[assignment] + + await runner.handle(_make_task(task_id="t1")) + + gateway: FakeGateway = runner._gateway_for_test # type: ignore[attr-defined] + terminal = gateway.last_terminal_report() + assert terminal["state"] == OrderState.PAID + assert terminal["terminal_status"] == TaskStatus.SUCCEEDED + + assert len(runner._monitor_tasks) == 1 + + await runner.cancel_monitors() + assert len(runner._monitor_tasks) == 0 + + +async def test_monitor_order_reports_state_changes_and_stops_at_delivered( + runner: WorkerRunner, evidence: EvidenceStore +): + """状态没变化不重复 report;到「配達完了」立刻停止轮询,不再多探测一次""" + snapshots = [ + OrderStatusSnapshot(found=False), + OrderStatusSnapshot(found=True, stage_label="ショップ", order_state=None, html="shop"), + OrderStatusSnapshot( + found=True, stage_label="出荷", order_state=OrderState.SHIPPED, html="shipped" + ), + OrderStatusSnapshot( + found=True, stage_label="配達完了", order_state=OrderState.DELIVERED, html="delivered" + ), + ] + calls: list[OrderStatusSnapshot] = [] + + async def _check(order_id: str): # noqa: ANN001 + calls.append(snapshots[len(calls)]) + return calls[-1] + + runner._site.check_order_status = _check # type: ignore[assignment] + + await runner._monitor_order(_make_task(task_id="t1"), "306087-20260813-0863947697") + + assert len(calls) == 4 # 配達完了那次之后立刻返回,不会再多轮询一次 + + gateway: FakeGateway = runner._gateway_for_test # type: ignore[attr-defined] + reported_states = [r["state"] for r in gateway.reports] + assert reported_states == [OrderState.SHIPPED, OrderState.DELIVERED] + assert all(r["terminal"] is False for r in gateway.reports) # 追加上报,不重新终结任务 + + base = evidence.step_dir("t1") + assert (base / "06-monitor-shipped.html").exists() + assert (base / "07-monitor-delivered.html").exists() + + +async def test_monitor_order_stops_after_max_checks_without_finding_order( + runner: WorkerRunner, +): + """一直查不到订单(found=False):轮询到上限后正常退出,不报错、不 report""" + + async def _check(order_id: str): # noqa: ANN001 + return OrderStatusSnapshot(found=False) + + runner._site.check_order_status = _check # type: ignore[assignment] + + await runner._monitor_order(_make_task(task_id="t1"), "ord-1") + + gateway: FakeGateway = runner._gateway_for_test # type: ignore[attr-defined] + assert gateway.reports == [] + + +async def test_monitor_order_swallows_check_errors_and_keeps_polling( + runner: WorkerRunner, +): + """单次探测异常(如登录态失效)只记日志继续重试,不会让后台任务崩掉""" + call_count = 0 + + async def _check(order_id: str): # noqa: ANN001 + nonlocal call_count + call_count += 1 + if call_count <= 2: + raise RuntimeError("模拟登录态失效") + return OrderStatusSnapshot( + found=True, stage_label="配達完了", order_state=OrderState.DELIVERED, html="ok" + ) + + runner._site.check_order_status = _check # type: ignore[assignment] + + await runner._monitor_order(_make_task(task_id="t1"), "ord-1") + + assert call_count == 3 # 前两次异常被吞掉继续重试,第三次成功拿到 DELIVERED 后停止 + gateway: FakeGateway = runner._gateway_for_test # type: ignore[attr-defined] + assert [r["state"] for r in gateway.reports] == [OrderState.DELIVERED]