"""FastAPI 应用工厂。""" from __future__ import annotations import logging import os import time from collections.abc import AsyncIterator from contextlib import asynccontextmanager from fastapi import FastAPI from fastapi.middleware.cors import CORSMiddleware from src.core.config import settings logger = logging.getLogger(__name__) # 进程启动时间(monotonic,不受系统时钟跳变影响):/health 的 uptime 来源 _PROCESS_STARTED_MONOTONIC = time.monotonic() async def _fail_stale_ingest_jobs() -> None: """P1-E: 启动时将上次遗留的 pending/running ingest_jobs 标 failed。 进程异常退出(重启/OOM)会导致 ingest_jobs 残留为 pending/running, 这些任务实际已不在执行,启动时一次性标 failed 避免永久"执行中"。 尽力而为:失败只记 warning,不阻断启动。 """ from datetime import datetime, timezone from sqlalchemy import update from src.db.base import AsyncSessionLocal from src.db.models import IngestJob try: async with AsyncSessionLocal() as session: stmt = ( update(IngestJob) .where(IngestJob.status.in_(["pending", "running"])) .values(status="failed", error="进程重启:任务被终止", finished_at=datetime.now(timezone.utc)) ) result = await session.execute(stmt) await session.commit() if result.rowcount: logger.info("P1-E: 已将 %d 条残留 pending/running ingest_jobs 标 failed", result.rowcount) except Exception: logger.warning("P1-E: 清理残留 ingest_jobs 失败,不影响启动", exc_info=True) @asynccontextmanager async def lifespan(app: FastAPI) -> AsyncIterator[None]: from src.db.base import init_db from src.core.http_client import close_client from src.core.runtime_config import ( ensure_admin_password_hashed, migrate_plaintext_sensitive_settings, ) from src.core.security_check import assert_security_on_startup from src.core.scheduler import scheduler, quality_scheduler from src.api.routes.schedules import _run_scheduled_task from src.data.config import BZZOIRO_LEAGUE_IDS await init_db() # 验证连接,不建表 await migrate_plaintext_sensitive_settings() # 明文敏感配置 → 加密(幂等) await ensure_admin_password_hashed() # .env 明文密码 → scrypt 哈希(幂等) await assert_security_on_startup() # 启动安全校验(生产拒绝/开发警告) await _fail_stale_ingest_jobs() # P1-E: 上次遗留的 pending/running 标 failed # D7(工程债): 进程内限流(_RateLimiter)与 KeyRing 均为单进程状态; # 多 worker 部署时各进程独立计数,限流阈值会按 worker 数放大、KeyRing 不共享。 # 生产环境应将限流前置到 Nginx/网关,或以单 worker 运行(见 src/api/deps.py 注释)。 # 此处仅提醒一次,不阻断启动,也不引入 Redis 等外部依赖。 if settings.APP_ENV == "production": logger.warning( "APP_ENV=production: 进程内 rate-limit 与 KeyRing 仅单进程有效;" "多 worker 部署请将限流前置到 Nginx/网关,或以单 worker 运行" ) # P3-3:STRICT_SINGLE_WORKER 启动期强制校验,拒绝多 worker 静默配额漂移。 # uvicorn 通过 --workers 传入;此处以环境变量 UVICORN_WORKERS 或启动参数判定。 # 为避免耦合 uvicorn 内部,仅校验一个显式传入的标记:当 STRICT_SINGLE_WORKER=True 时, # 要求环境变量 UVICORN_WORKERS 不为空且 <=1,否则拒绝启动。 if settings.STRICT_SINGLE_WORKER: workers = os.environ.get("UVICORN_WORKERS", "1") try: n_workers = int(workers) except ValueError: n_workers = 1 if n_workers > 1: raise RuntimeError( f"STRICT_SINGLE_WORKER=True 但以 {n_workers} worker 启动会被拒绝 " f"(应用内限流/KeyRing 多 worker 下各自独立计数,配额放大 {n_workers} 倍)。" f"请前置 Nginx/网关全局限流后再启用多 worker,或保持单 worker。" ) logger.info("STRICT_SINGLE_WORKER=True:已确认单 worker 启动,限流配额不会漂移") # 注册默认定时任务(如果数据库中没有) from src.db.base import AsyncSessionLocal from sqlalchemy import select from src.db.models import Schedule async with AsyncSessionLocal() as session: existing = (await session.execute(select(Schedule.id))).scalars().all() if "daily-events" not in existing: session.add(Schedule(id="daily-events", task="events", cron="0 8 * * *", enabled=False)) if "daily-standings" not in existing: session.add(Schedule(id="daily-standings", task="standings", cron="0 9 * * *", enabled=False)) if "daily-stats" not in existing: session.add(Schedule(id="daily-stats", task="stats", cron="*/30 * * * *", enabled=False)) await session.commit() # 从数据库加载所有启用的定时任务 async with AsyncSessionLocal() as session: schedules = (await session.execute(select(Schedule).where(Schedule.enabled))).scalars().all() for s in schedules: leagues = s.leagues.split(",") if s.leagues else list(BZZOIRO_LEAGUE_IDS.keys()) scheduler.register( s.id, s.cron, lambda sid=s.id: _run_scheduled_task(sid), enabled=s.enabled, ) await scheduler.start() await quality_scheduler.start() logger.info("应用启动完成") yield await scheduler.stop() await quality_scheduler.stop() await close_client() def create_app() -> FastAPI: from src.core.log_buffer import setup_logging setup_logging(settings.LOG_LEVEL, settings.LOG_FILE) # 生产环境不暴露 OpenAPI 文档(避免向访客泄露接口结构) openapi_url = "/openapi.json" if settings.APP_ENV != "production" else None app = FastAPI( title="Profeto API", description="足球数据 + LLM 预测服务", version="0.1.0", lifespan=lifespan, openapi_url=openapi_url, ) origins = [o.strip() for o in settings.CORS_ORIGINS.split(",") if o.strip()] methods = [m.strip() for m in settings.CORS_METHODS.split(",") if m.strip()] headers = [h.strip() for h in settings.CORS_HEADERS.split(",") if h.strip()] app.add_middleware( CORSMiddleware, allow_origins=origins, allow_credentials=True, allow_methods=methods, allow_headers=headers, ) from src.api.routes.matches import router as matches_router from src.api.routes.predict import router as predict_router from src.api.routes.ingest import router as ingest_router from src.api.routes.eval import router as eval_router from src.api.routes.backtest import router as backtest_router from src.api.routes.auth import router as auth_router from src.api.routes.admin_settings import router as admin_settings_router from src.api.routes.schedules import router as schedules_router from src.api.routes.admin_monitoring import router as admin_monitoring_router app.include_router(matches_router) app.include_router(predict_router) app.include_router(ingest_router) app.include_router(eval_router) app.include_router(backtest_router) app.include_router(auth_router) app.include_router(admin_settings_router) app.include_router(schedules_router) app.include_router(admin_monitoring_router) @app.get("/health") async def health(): """存活检查。version/uptime_seconds 供管理端监控页展示。""" version = None try: from importlib.metadata import version as _pkg_version version = _pkg_version("profeto") except Exception: # 包元数据缺失时返回 None,前端降级隐藏版本卡片 version = None return { "status": "healthy", "service": "profeto", "version": version, "uptime_seconds": round(time.monotonic() - _PROCESS_STARTED_MONOTONIC), } @app.get("/health/ready") async def health_ready(): """就绪检查:验证数据库连接。 数据库不可达时返回 HTTP 503,而非 200 + not_ready —— 这样 K8s/Compose 的 readinessProbe 才能正确判定「未就绪」并停止流量。 """ from src.db.base import engine from fastapi.responses import JSONResponse try: async with engine.begin() as conn: await conn.run_sync(lambda conn: None) return {"status": "ready"} except Exception as e: logger.warning("就绪检查失败(数据库不可达): %s", e) return JSONResponse( status_code=503, content={"status": "not_ready", "reason": "database_unreachable"}, ) return app app = create_app()