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]