Files
q792602257andClaude Opus 5 c03158488b feat(gateway): 账号只读查询通道——从已登录账号取真实订单
上游要的不只是网关记的任务状态镜像,还有「已登录账号在站点上的真实订单」,
但账号只在 NAT 后本地机上,只能经网关队列走。新增独立查询通道(§11):

- gateway 单开 account_queries 表 + QueryStatus 状态机,接口
  POST /api/account/queries(幂等)/ lease / {id}/result / {id}
- 不复用下单任务队列:查询是只读,租约过期可安全重投(与下单「绝不自动
  重投」相反),且不该被全局并发度 1 堵死、task_reports 是订单镜像不能污染
- 本地交易服务起第二条常驻循环 query_runner,领到即调 SiteInteractor 真读:
  order_list 复用已实测的 list_recent_orders(规范化字段 + 站点
  orderListData 原文),order_detail 复用 fetch_order_detail(配送阶段 +
  页面 __INITIAL_STATE__ 原样透传,结构未经真实样本,不抽字段)
- 账号级串行仍由 SiteInteractor 的锁保证;每次执行套超时按失败回报
- 错误码 6005/6006(查询通道,可重试只读区别于 6001-6004);/health 暴露
  queued_query_count;结果体积上限先丢原始 JSON

openapi.json 重导,docs/order-gateway.md §11、README、.env.example 补全
配置与实测边界。全量测试 404→454 通过。

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

242 lines
10 KiB
Python

"""交易服务入口:FastAPI 应用创建与生命周期管理
与抓取服务(app/scraping/main.py)分成两个进程运行,理由不是「要不要登录」
这一条,而是运行特性根本不同:
- 抓取无状态、可重试、可多开实例;这里的写操作**不可逆**,重复提交就是重复下单。
- 登录态 cookie 全局唯一,订单监控是常驻轮询;多开实例会让同一个账号被多个
进程并发操作,也会让轮询重复触发。
- 抓取被限速最多是慢,账号被风控是封号;两者不该共用出口 IP 与请求节奏。
因此本服务只能单实例运行(或按账号分片),扩容靠抓取服务那一侧。
SiteInteractor 始终启动:HTTP /api/auth/* 与 /api/cart/* 都要靠它持账号 cookie
做站点交互(加购、校验、清空、删除)。worker 启动条件:RAKUTEN_ORDER_GATEWAY_URL
非空。留空时不构造 worker 子组件(本地 DB / 证据目录 / 出站客户端都不创建),
但 SiteInteractor 与 Playwright 仍会启动。配了网关 URL 时会起**两条**常驻循环:
下单 worker 与账号只读查询 worker(后者见 docs/order-gateway.md §11)。
"""
from __future__ import annotations
import asyncio
import logging
from contextlib import asynccontextmanager
from fastapi import FastAPI
from app.shared.api import register_exception_handlers
from app.shared.config import get_settings
from app.shared.logging_setup import configure_logging
from app.shared.telemetry import instrument_app, setup_telemetry, shutdown_telemetry
from app.trading.api.routes.auth import router as auth_router
from app.trading.api.routes.cart import router as cart_router
from app.trading.api.routes.health import router as health_router
from app.trading.container import TradingContainer
from app.trading.services.auth_session import AuthSession
logger = logging.getLogger(__name__)
def build_container() -> TradingContainer:
"""构建交易服务容器
- 登录态会话与站点交互器始终构造(HTTP /api/cart/* 与 /api/auth/* 都依赖)
- worker 仅在配置了网关 URL 时构造(依赖 httpx 客户端、本地 DB、证据目录)
"""
settings = get_settings()
container = TradingContainer(settings=settings, auth_session=AuthSession(settings))
# SiteInteractor 总是构造:HTTP 接口需要它,worker(如果配置了)也需要。
# 延迟 import 让没装 Playwright 的环境能起服务做基本配置检查。
from app.trading.worker.site_interact import SiteInteractor
container.site = SiteInteractor(auth_session=container.auth_session, settings=settings)
if settings.order_gateway_url:
# 延迟 import:未配置网关 URL 时不加载 worker 模块(也就不会拉起 aiosqlite 等)
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.query_runner import QueryRunner
from app.trading.worker.runner import WorkerRunner
client = GatewayClient(
settings.order_gateway_url,
settings.bearer_token,
settings=settings,
timeout=max(60.0, settings.lease_max_wait_seconds + 10),
)
local_db = LocalDB(settings.trading_db_path_resolved)
evidence = EvidenceStore(settings.evidence_path)
runner = WorkerRunner(
settings=settings,
gateway_client=client,
local_db=local_db,
evidence=evidence,
site=container.site,
)
container.worker_client = client
container.worker_local_db = local_db
container.worker_evidence = evidence
container.worker_runner = runner
# 只读查询循环与下单 worker 共用同一个网关客户端与 SiteInteractor:
# 前者是同一份连接池,后者那把锁正是账号级串行的保证。
container.query_runner = QueryRunner(
settings=settings, gateway_client=client, site=container.site
)
return container
async def auto_login_on_start(container: TradingContainer) -> None:
"""启动后自动把登录态准备好(RAKUTEN_AUTO_LOGIN_ON_START=true 时)
为容器部署而存在:镜像里没有落盘的 storage_state,不希望每次新起容器都得先
在宿主机跑一遍 scripts/login.py。逐站点探测,未登录就按 account.yaml 登一次。
刻意做成后台任务而不是启动阻塞:登录最长要等 relogin_timeout_seconds(默认
300s,撞验证码时在等人工),阻塞会让 /health 在这段时间里连端口都不通。
任何失败都只记日志——服务照常提供 /health 与 /api/auth/*,人工接管后调
/api/auth/login 重试即可。
"""
for site in container.auth_session.sites:
try:
status = await container.auth_session.check(site)
if status.logged_in:
logger.info("启动自动登录跳过:site=%s 已是登录态", site)
continue
logger.info("启动自动登录:site=%s 当前未登录(%s)", site, status.detail)
if await container.auth_session.try_relogin(site):
status = await container.auth_session.check(site)
logger.info(
"启动自动登录结束:site=%s logged_in=%s detail=%s",
site, status.logged_in, status.detail,
)
except Exception:
logger.exception("启动自动登录异常:site=%s(服务继续运行)", site)
@asynccontextmanager
async def lifespan(app: FastAPI):
"""应用生命周期管理:登录态会话 + 站点交互器(必起)+ 可选 worker"""
container = build_container()
app.state.container = container
configure_logging(container.settings)
# 与抓取侧同样:必须在创建 httpx 客户端之前 setup。
setup_telemetry(container.settings, service_name="rakuten-trading")
logger.info(
"交易服务启动:%s:%s",
container.settings.trading_host,
container.settings.trading_port,
)
logger.info("当前环境:%s", container.settings.app_env)
await container.auth_session.start()
# SiteInteractor 总是启动:HTTP /api/cart/* 接口要靠它发请求
assert container.site is not None
await container.site.start() # type: ignore[union-attr]
login_task: asyncio.Task | None = None
if container.settings.auto_login_on_start:
login_task = asyncio.create_task(
auto_login_on_start(container), name="trading-auto-login"
)
else:
logger.info("未开启 RAKUTEN_AUTO_LOGIN_ON_START,启动时不自动登录")
worker_task: asyncio.Task | None = None
query_task: asyncio.Task | None = None
if container.worker_runner is not None:
# 顺序:本地 DB 在 worker 写证据前先就绪;SiteInteractor 已在外层启动
assert container.worker_local_db is not None
assert container.worker_client is not None
await container.worker_local_db.start()
worker_task = asyncio.create_task(
container.worker_runner.run(), name="trading-worker"
)
logger.info(
"下单 worker 已启动:worker_id=%s gateway=%s",
container.settings.worker_id_effective,
container.settings.order_gateway_url,
)
# 第二条循环:账号只读查询。独立于下单主循环,否则一笔下单跑几分钟,
# 期间上游一句「账号里那笔订单什么状态」就得干等(docs/order-gateway.md §11)。
assert container.query_runner is not None
query_task = asyncio.create_task(
container.query_runner.run(), name="trading-query-worker"
)
logger.info("账号只读查询 worker 已启动")
else:
logger.info(
"未配置 RAKUTEN_ORDER_GATEWAY_URL,下单 worker 不启动(仅登录态与购物车接口)"
)
try:
yield
finally:
if login_task is not None and not login_task.done():
# 关服务时正在登的那次不要了:浏览器由 login_one 自己的 async with 收尾
login_task.cancel()
try:
await login_task
except asyncio.CancelledError:
pass
if worker_task is not None:
container.worker_runner.stop() # type: ignore[union-attr]
worker_task.cancel()
try:
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 query_task is not None:
# 与下单 worker 同样:先 stop 再 cancel,正在跑的那次读取直接放弃
# (只读,放弃没有副作用;网关侧那张单会超时重投)。
container.query_runner.stop() # type: ignore[union-attr]
query_task.cancel()
try:
await query_task
except asyncio.CancelledError:
pass
if container.site is not None:
await container.site.close() # type: ignore[union-attr]
if container.worker_local_db is not None:
await container.worker_local_db.close()
if container.worker_client is not None:
await container.worker_client.aclose()
await container.auth_session.close()
shutdown_telemetry()
def create_app() -> FastAPI:
"""创建 FastAPI 应用实例,注册路由和异常处理器"""
app = FastAPI(title="Rakuten Trading Service", lifespan=lifespan)
app.include_router(health_router)
app.include_router(auth_router)
app.include_router(cart_router)
register_exception_handlers(app)
instrument_app(app)
return app
app = create_app()
if __name__ == "__main__":
import uvicorn
settings = get_settings()
configure_logging(settings)
uvicorn.run(
"app.trading.main:app",
host=settings.trading_host,
port=settings.trading_port,
log_config=None,
timeout_keep_alive=120,
# 单进程:登录态与后续的订单监控都不能有第二份
workers=1,
)