"""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, 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 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 def stop(self) -> None: self._running = False @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) 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:当前 verify 桩不会返回这个值,但留接口给 # 未来真正能可靠核对时——那时再决定 NOT_ORDERED 是否直接重新执行 await self._report_safe( task, state=_coerce_state(task.known_state), detail=f"恢复核对无法定论:{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 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。 站点交互当前未实现(site_interact 抛 NotImplementedError),第一步就会 转到 _execute_with_renewal 的 except 分支上报 needs_human。 """ await self._db.ensure_started(task.task_id, task.site, task.intent) # ラクマ 加购契约未实现,首版只对乐天。其他站点直接转人工 if task.site != "rakuten": raise NotImplementedError(f"site={task.site} 暂未实现下单(首版仅 rakuten)") # 步骤 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_html = 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} 円", evidence_meta={ "payable_yen": summary.payable_yen, "site_order_id": summary.site_order_id, "pay_deadline": summary.pay_deadline, }, ) # 步骤 4:提交下单 site_order_id = await self._site.submit_order(task) 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}", ) # 步骤 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:付款后监控(非阻塞,常驻轮询;当前未实现) try: await self._site.monitor(task, site_order_id) except NotImplementedError: logger.info("付款后监控未实现,跳过:task_id=%s", 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, ) -> None: """单步执行:动作 → 落证据 → 写本地 → 回报 gateway""" await action() meta = { "step": step_name, "state": state.value, **(evidence_meta or {}), } evidence_ref = self._evidence.write_step( task.task_id, step_no, step_name, 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