"""交易服务入口:FastAPI 应用创建与生命周期管理 与抓取服务(app/scraping/main.py)分成两个进程运行,理由不是「要不要登录」 这一条,而是运行特性根本不同: - 抓取无状态、可重试、可多开实例;这里的写操作**不可逆**,重复提交就是重复下单。 - 登录态 cookie 全局唯一,订单监控是常驻轮询;多开实例会让同一个账号被多个 进程并发操作,也会让轮询重复触发。 - 抓取被限速最多是慢,账号被风控是封号;两者不该共用出口 IP 与请求节奏。 因此本服务只能单实例运行(或按账号分片),扩容靠抓取服务那一侧。 SiteInteractor 始终启动:HTTP /api/auth/* 与 /api/cart/* 都要靠它持账号 cookie 做站点交互(加购、校验、清空、删除)。worker 启动条件:RAKUTEN_ORDER_GATEWAY_URL 非空。留空时不构造 worker 子组件(本地 DB / 证据目录 / 出站客户端都不创建), 但 SiteInteractor 与 Playwright 仍会启动。 """ 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.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 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 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, ) 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 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, )