Files
rakuten-api/app/trading/worker/runner.py
T
q792602257andClaude Opus 5 3284e086ad feat(trading): 浏览器掉线兜底——任务边界自愈重建,中途掉线转 needs_human
SiteInteractor 的 Chromium 是进程级单例,此前启动后即假定永远活着:全仓唯一的
is_connected() 探活在 scraping 侧,交易侧既不探活也不重启。容器里 Chromium 崩溃
是有真实前提的(/dev/shm 不足、OOM kill、seccomp 挡 sandbox,docker-compose.yml
里已有相关注释),一旦发生,进程还活着但之后每一单都会失败,且 /health 恒返回
ok,restart: unless-stopped 永远不会被触发。

更隐蔽的一条:clear_cart / enter_checkout / pay 等处的 new_page()、context.request
写在 try 之外,掉线抛的 TargetClosedError 不是 AppError,会穿过 runner 的
except AppError 落到主循环那个只记日志的兜底里——任务一次都不上报,网关侧要干等
整个 lease_ttl(默认 300s)才被 sweep 置 stale。clear_cart 是 execute() 的 step 0,
浏览器死时最先撞上的正是它。

分三层处理:

- 任务边界自愈。_launch() 从 start() 抽出,_context_options() 统一 context 参数
  (重建必须与启动完全一致,指纹漂移就是一次风控事件);_refresh_context_if_stale
  改名 _ensure_context_ready(),先探活重建再做原有的 storage_state mtime 检查
  (顺序不可换,mtime 重建要用 self._browser)。重建前丢弃 _checkout_pages 里的
  残留确认页并记 warning——那些 Page 已随浏览器一起没了。

- 中途掉线不重建,抛 BrowserDeadError(新增,5006)。新增 _new_page() 与
  _request() 两个壳收口裸异常;_request() 只在确认浏览器真死了时才改写异常,
  站点 5xx 这类正常业务失败原样抛出。submit_order / pay 入口用
  _require_live_browser() 直接拒绝:这两步复用 enter_checkout 留存的 Page,
  重建救不回服务端订单草稿,而 pay 跑的时候订单已经真的提交了。顺带修掉一个
  误诊——浏览器死时 submit_order 原先报「未找到确认按钮」,把「浏览器崩了」
  说成「站点改版了」,两者的处置方式完全不同。

- runner 把 BrowserDeadError 转 needs_human 而非 failed,except 分支排在
  except AppError 之前(子类,顺序反了就报 failed)。掉线发生在动作中途,
  站点侧生效与否无从判断,不能给上游「明确失败」的结论。

/health 暴露 browser 状态,掉线时 degraded + HTTP 503 + code 5006,让
Dockerfile.trading 的 HEALTHCHECK 探到并重启容器。重启不会导致重复下单:网关侧
任务绝不自动重投,租约过期只置 stale 等人工 reclaim(docs/order-gateway.md §5),
重启只是恢复领新任务的能力。启动窗口期 started=False 不算掉线。

README 错误码表补 5005(此前遗漏)与 5006。

新增 15 个用例覆盖探活三态、is_connected() 自身抛错、边界重建/不重建/丢弃残留页/
重建失败、两个包装壳的分支、submit/pay 拒绝、runner 转 needs_human、/health 503。
真实 Chromium 崩溃无法在离线测试里制造,用例模拟的是 is_connected() 返回 False
这个唯一可观测信号,覆盖的是代码对该信号的反应而非崩溃本身;容器 HEALTHCHECK
真的触发重启这条链路尚未实跑验证。

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-20 16:01:47 +08:00

619 lines
26 KiB
Python

"""worker 主循环与每任务执行流
主循环(规格 §6):
while running:
task = await gateway.lease(worker_id, wait=30)
if task is None:
continue
if local_db.has_finished(task.task_id): # 本地幂等闸门
await gateway.report(...); continue
if task.lease_count > 1: # 恢复领取
verdict = await verify_on_site(task)
if verdict.already_ordered: await report(...); continue
if verdict.unknown: await report(needs_human); continue
async with renew_lease_every(60):
await execute(task)
每个动作的顺序固定为:动作 → 落证据 → 写本地 SQLite → 回报 gateway。
顺序不能颠倒——先落证据再回报,保证服务器上看到的状态一定有本地证据可查。
"""
from __future__ import annotations
import asyncio
import contextlib
import logging
from typing import TYPE_CHECKING
from app.shared.errors import (
AppError,
BrowserDeadError,
CheckoutBlockedError,
OrderGuardError,
)
from app.shared.task_state import OrderState, TaskStatus
from app.trading.worker import verify
from app.trading.worker.client import GatewayClient
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.site_interact import PageSnapshot, SiteInteractor
if TYPE_CHECKING:
from app.shared.config import Settings
logger = logging.getLogger(__name__)
def _coerce_state(value: str | None) -> OrderState:
"""把网关返回的 state 字符串安全地包成 OrderState;None 或非法值回落到 CREATED"""
if not value:
return OrderState.CREATED
try:
return OrderState(value)
except ValueError:
logger.warning("未知的 order state 字符串,回落到 CREATED:%s", value)
return OrderState.CREATED
class WorkerRunner:
"""worker 主循环与执行流
持有 gateway 客户端、本地 DB、证据存储。`run()` 是常驻 asyncio 任务,
`handle()` 处理单个任务,`execute()` 是站点交互的实际调度(当前多数步骤未实现)。
"""
def __init__(
self,
*,
settings: "Settings",
gateway_client: GatewayClient,
local_db: LocalDB,
evidence: EvidenceStore,
site: SiteInteractor,
):
self._settings = settings
self._gateway = gateway_client
self._db = local_db
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
# ---- 主循环 ----
async def run(self) -> None:
"""常驻循环:长轮询领任务,按 §6 调度"""
self._running = True
logger.info("worker 启动:worker_id=%s", self.worker_id)
while self._running:
try:
task = await self._gateway.lease(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 task is None:
continue
try:
await self.handle(task)
except Exception: # noqa: BLE001
# 任何未预期的异常都不应让 worker 退出
logger.exception("handle 任务异常:task_id=%s", task.task_id)
# ---- 单任务调度 ----
async def handle(self, task: LeaseTask) -> None:
"""单任务调度入口:本地幂等闸门 → 恢复核对 → 执行"""
# 本地幂等闸门:之前已完成过的任务不再执行
if await self._db.has_finished(task.task_id):
final_state = await self._db.final_state(task.task_id)
logger.info(
"本地已完成,补报终态:task_id=%s state=%s", task.task_id, final_state
)
await self._report_safe(
task,
state=_coerce_state(final_state),
detail="本地已完成,补报避免重复执行",
terminal=True,
)
return
# 恢复领取:lease_count > 1 表示 stale → reclaim,必须先核对站点订单
if task.lease_count > 1:
await self._handle_recovery(task)
return
# 常规执行
await self._execute_with_renewal(task)
async def _handle_recovery(self, task: LeaseTask) -> None:
"""lease_count > 1 的恢复路径:先核对,不直接执行"""
logger.warning(
"恢复领取(lease_count=%s),先核对站点订单:task_id=%s",
task.lease_count,
task.task_id,
)
verdict = await verify.verify_on_site(task, gateway=self._gateway, site=self._site)
if verdict.verdict == verify.VerifyVerdict.ALREADY_ORDERED:
await self._report_safe(
task,
state=OrderState.ORDERED,
site_order_id=verdict.site_order_id,
detail=f"核对结论:已下单({verdict.detail})",
terminal=True,
terminal_status=TaskStatus.SUCCEEDED,
)
await self._db.ensure_started(task.task_id, task.site, task.intent)
await self._db.mark_finished(task.task_id, OrderState.ORDERED.value)
return
# NOT_ORDERED 也走 needs_human(2026-08-13 verify_on_site 已实现,会真的
# 返回这个值,但刻意仍不自动重新执行):核对逻辑目前只有「账号 1 笔订单」
# 的真实数据支撑 ALREADY_ORDERED 分支,NOT_ORDERED 分支只验证过逻辑本身,
# 没有被真实多单场景跑过——在这条判断被更多真实数据验证之前,即使确认
# 未下单也交人工决定是否重新提交,不自动触发新的下单动作。这是刻意的
# 保守选择,不是遗漏;要不要放开需要显式决定,不在这里静默改。
await self._report_safe(
task,
state=_coerce_state(task.known_state),
detail=f"恢复核对结论({verdict.verdict.value}):{verdict.detail}",
terminal=True,
terminal_status=TaskStatus.NEEDS_HUMAN,
)
# ---- 常规执行 ----
async def _execute_with_renewal(self, task: LeaseTask) -> None:
"""在租约自动续期的上下文里执行任务"""
async with self._renew_lease_every(task, interval=60):
try:
await self.execute(task)
except NotImplementedError as exc:
# 站点交互未实现(规格 §10):上报 needs_human,不视为 worker 失败
logger.warning(
"站点交互未实现,转 needs_human:task_id=%s what=%s",
task.task_id,
exc,
)
await self._report_safe(
task,
state=_coerce_state(task.known_state),
detail=f"站点交互未实现:{exc}",
terminal=True,
terminal_status=TaskStatus.NEEDS_HUMAN,
)
except OrderGuardError as exc:
# 金额守卫等闸门拦截:上报 needs_human,站点侧无提交动作
logger.warning(
"下单闸门拦截,转 needs_human:task_id=%s msg=%s",
task.task_id,
exc.message,
)
await self._report_safe(
task,
state=_coerce_state(task.known_state),
detail=exc.message,
terminal=True,
terminal_status=TaskStatus.NEEDS_HUMAN,
)
except CheckoutBlockedError as exc:
# 站点风控拦截(session upgrade / 3DS 等):转 needs_human 交人工接管,
# 不当普通失败重试(规格 §10.1)
logger.warning(
"结算被站点风控拦截,转 needs_human:task_id=%s msg=%s",
task.task_id,
exc.message,
)
await self._report_safe(
task,
state=_coerce_state(task.known_state),
detail=exc.message,
terminal=True,
terminal_status=TaskStatus.NEEDS_HUMAN,
)
except BrowserDeadError as exc:
# 浏览器在执行途中没了:站点侧到底生效没有无从判断(尤其掉线发生在
# submit_order / pay 前后),按 needs_human 交人工核对订单列表。
# 必须排在 except AppError 前面——BrowserDeadError 是 AppError 的
# 子类,顺序反了就会被当成普通失败报 failed。
logger.error(
"浏览器掉线导致任务中断,转 needs_human:task_id=%s msg=%s",
task.task_id,
exc.message,
)
await self._report_safe(
task,
state=_coerce_state(task.known_state),
detail=exc.message,
terminal=True,
terminal_status=TaskStatus.NEEDS_HUMAN,
)
except AppError as exc:
logger.warning(
"执行失败:task_id=%s code=%s msg=%s",
task.task_id,
exc.err_code,
exc.message,
)
await self._report_safe(
task,
state=_coerce_state(task.known_state),
detail=exc.message,
terminal=True,
terminal_status=TaskStatus.FAILED,
)
async def execute(self, task: LeaseTask) -> None:
"""站点交互的实际调度:清购物车 → 加购 → 校验 → 确认页 → 金额守卫 → 提交 → 付款
每一步的顺序:动作 → 落证据 → 写本地 SQLite → 回报 gateway。
开单前的「清购物车」是本机侧卫生步骤(step 0,见下方代码注释):落本地
步骤证据(页面 + 截图 + meta)但不上报 gateway、不记状态事件;站点交互
失败会转 _execute_with_renewal 的 except 分支上报 needs_human / failed。
"""
await self._db.ensure_started(task.task_id, task.site, task.intent)
# 交易服务只覆盖乐天市场。其他站点直接转人工
if task.site != "rakuten":
raise NotImplementedError(f"site={task.site} 暂不在交易服务范围内(仅 rakuten)")
# 步骤 0:开单前清理购物车。上一单若在提交(submit_order)之前失败——比如
# 加购成功但 enter_checkout 被 session upgrade 拦截、或金额守卫拦下——残留的
# 商品不会自己消失,会一直留在购物车里;下一次 add_to_cart 把新商品叠加在旧
# 商品上,结算时会把上一单的一起买走(已实测踩到过)。所以每单开跑前先清空
# 购物车,保证「一单 = 只买本单商品」的不变式——无论上一单是怎么失败的
# (包括进程中途崩溃,残留都没人清),这里都从空车起步。
#
# 这条不在 gateway 上报、不记 order_events:它是本机侧的卫生操作,不是订单
# 进度的状态迁移。但清理后的页面快照、整页截图与 meta 要落本地证据并登记
# evidence_index——清理没跑干净被下面的闸门拦下转 needs_human 时,这份现场
# 就是排查依据,所以证据必须先于闸门判断落盘。清理后 count 非 0(含
# count API 拿不到结果返回 -1)说明清理没跑干净,**不能**带着残留往下
# 加购——有把上一单买走的真实风险,按闸门语义拦截转 needs_human,
# 宁可卡住等人核对,也不赌「残留不会被买走」。
cleared = await self._site.clear_cart()
evidence_ref = self._evidence.write_step(
task.task_id, 0, "cart-clear",
html=cleared.get("html") or None,
png=cleared.get("screenshot") or None,
meta={
"step": "cart-clear",
"removed_count": cleared.get("removed_count"),
"cart_count": cleared.get("cart_count"),
},
)
await self._db.index_evidence(task.task_id, 0, "cart-clear", evidence_ref)
if cleared.get("cart_count", -1) != 0:
raise OrderGuardError(
"开单前清理购物车后仍未清空"
f"(cart_count={cleared.get('cart_count')},removed_count="
f"{cleared.get('removed_count')}):残留商品可能随本次下单一起被买走,"
"中止转人工核对"
)
# 步骤 1:加购
await self._run_step(
task, step_no=1, step_name="cart-add",
action=lambda: self._site.add_to_cart(task),
state=OrderState.IN_CART,
detail="已加入购物车",
)
# 步骤 2:校验购物车
await self._run_step(
task, step_no=2, step_name="cart-check",
action=lambda: self._site.verify_cart(task),
state=OrderState.IN_CART,
detail="购物车已校验",
)
# 步骤 3:进入下单确认页 + 金额守卫
checkout = await self._site.enter_checkout(task)
summary = await self._site.parse_checkout(checkout.html)
self._enforce_amount_guard(task, summary.payable_yen)
await self._run_step(
task, step_no=3, step_name="order-confirm",
action=self._noop(),
state=OrderState.CREATED,
detail=f"下单确认页已解析:应付 {summary.payable_yen} 円",
html=checkout.html or None,
png=checkout.screenshot or None,
evidence_meta={
"payable_yen": summary.payable_yen,
"site_order_id": summary.site_order_id,
"pay_deadline": summary.pay_deadline,
},
)
# 步骤 4:提交下单
submit = await self._site.submit_order(task)
site_order_id = submit.site_order_id
await self._run_step(
task, step_no=4, step_name="order-submit",
action=self._noop(),
state=OrderState.ORDERED,
site_order_id=site_order_id,
detail=f"已提交下单,站点订单号 {site_order_id}",
html=submit.evidence.html or None,
png=submit.evidence.screenshot or None,
)
# 步骤 5:付款
await self._run_step(
task, step_no=5, step_name="payment",
action=lambda: self._site.pay(task, site_order_id),
state=OrderState.AWAITING_PAYMENT,
site_order_id=site_order_id,
payable_yen=summary.payable_yen,
pay_deadline=summary.pay_deadline,
detail="已进入付款流程",
)
await self._report_safe(
task,
state=OrderState.PAID,
site_order_id=site_order_id,
payable_yen=summary.payable_yen,
pay_deadline=summary.pay_deadline,
detail="付款完成",
terminal=True,
terminal_status=TaskStatus.SUCCEEDED,
)
await self._db.mark_finished(task.task_id, OrderState.PAID.value)
# 步骤 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,
png=snapshot.screenshot or None,
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 时拦截"""
intent = task.intent or {}
limit = intent.get("max_total_yen")
if limit is None or limit <= 0:
limit = self._settings.order_max_total_yen
if limit > 0 and payable_yen > limit:
raise OrderGuardError(
f"实际应付 {payable_yen} 円超过上限 {limit} 円"
f"(task_id={task.task_id})"
)
@staticmethod
def _noop() -> "callable":
async def _f() -> None:
return None
return _f
async def _run_step(
self,
task: LeaseTask,
*,
step_no: int,
step_name: str,
action: "callable",
state: OrderState,
detail: str,
site_order_id: str | None = None,
payable_yen: int | None = None,
pay_deadline: str | None = None,
evidence_meta: dict | None = None,
html: str | None = None,
png: bytes | None = None,
) -> None:
"""单步执行:动作 → 落证据 → 写本地 → 回报 gateway
`action` 返回 PageSnapshot(add_to_cart / verify_cart / pay 等站点方法的
证据载体)时,其 html / screenshot 即本步骤要落盘的页面与整页截图;
显式传入的 `html` / `png` 优先级更高,供 step3/step4 等已经单独拿到
页面的调用点使用。
"""
result = await action()
if isinstance(result, PageSnapshot):
if html is None:
html = result.html or None
if png is None:
png = result.screenshot or None
elif html is None and isinstance(result, str):
# 旧契约兼容:站点方法返回字符串时视为页面 HTML
html = result
meta = {
"step": step_name,
"state": state.value,
**(evidence_meta or {}),
}
evidence_ref = self._evidence.write_step(
task.task_id, step_no, step_name, html=html, png=png, meta=meta
)
await self._db.index_evidence(task.task_id, step_no, step_name, evidence_ref)
await self._db.record_event(
task.task_id, state.value, detail=detail, evidence_ref=evidence_ref
)
await self._gateway.report(
task.task_id,
self.worker_id,
state=state,
payable_yen=payable_yen,
pay_deadline=pay_deadline,
site_order_id=site_order_id,
evidence_ref=evidence_ref,
detail=detail,
)
async def _report_safe(
self,
task: LeaseTask,
*,
state: OrderState,
detail: str,
terminal: bool = False,
terminal_status: TaskStatus | None = None,
site_order_id: str | None = None,
payable_yen: int | None = None,
pay_deadline: str | None = None,
) -> None:
"""回报 gateway,失败只记日志不抛——主循环不能因为回报失败退出"""
try:
await self._gateway.report(
task.task_id,
self.worker_id,
state=state,
detail=detail,
terminal=terminal,
terminal_status=terminal_status,
site_order_id=site_order_id,
payable_yen=payable_yen,
pay_deadline=pay_deadline,
)
except Exception: # noqa: BLE001
logger.exception("回报 gateway 失败:task_id=%s state=%s", task.task_id, state)
@contextlib.asynccontextmanager
async def _renew_lease_every(self, task: LeaseTask, *, interval: int):
"""后台任务:每隔 interval 秒续租一次。退出时停掉"""
stop_event = asyncio.Event()
async def _loop() -> None:
while not stop_event.is_set():
try:
await asyncio.wait_for(stop_event.wait(), timeout=interval)
except asyncio.TimeoutError:
pass
if stop_event.is_set():
return
try:
await self._gateway.renew(task.task_id, self.worker_id)
except AppError as exc:
logger.warning(
"续租失败:task_id=%s code=%s msg=%s",
task.task_id,
exc.err_code,
exc.message,
)
return
except Exception: # noqa: BLE001
logger.exception("续租异常:task_id=%s", task.task_id)
return
renewal = asyncio.create_task(_loop(), name=f"renew-{task.task_id}")
try:
yield
finally:
stop_event.set()
with contextlib.suppress(asyncio.CancelledError):
await renewal