main.py 3.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140
  1. """FastAPI 应用入口。"""
  2. from __future__ import annotations
  3. import asyncio
  4. import logging
  5. from contextlib import asynccontextmanager
  6. from typing import AsyncIterator
  7. from fastapi import FastAPI
  8. from app.exception_handlers import register_exception_handlers
  9. from app.response import ApiResponse
  10. from app.routers import search
  11. from app.schemas import CacheRefreshResult, CacheStatus
  12. from app.services.cache import (
  13. WORKER_ID,
  14. get_snapshot,
  15. leader_heartbeat,
  16. refresh_cache,
  17. scheduler_loop,
  18. try_become_leader,
  19. wait_cache_ready,
  20. )
  21. from app.services.redis_client import (
  22. KEY_LOCK_LEADER,
  23. close_redis,
  24. release_lock,
  25. )
  26. logging.basicConfig(
  27. level=logging.INFO,
  28. format="%(asctime)s | %(levelname)s | %(name)s | %(message)s",
  29. )
  30. logger = logging.getLogger(__name__)
  31. @asynccontextmanager
  32. async def lifespan(app: FastAPI) -> AsyncIterator[None]:
  33. """启动:选举 leader → 仅 leader 刷新缓存与跑定时任务;关闭:释放资源。"""
  34. is_leader = await try_become_leader()
  35. tasks: list[asyncio.Task[None]] = []
  36. if is_leader:
  37. logger.info("[%s] 当选 leader,开始首次刷新", WORKER_ID)
  38. await refresh_cache()
  39. tasks.append(
  40. asyncio.create_task(scheduler_loop(), name="faq-scheduler")
  41. )
  42. tasks.append(
  43. asyncio.create_task(leader_heartbeat(), name="leader-heartbeat")
  44. )
  45. logger.info("[%s] 定时任务与 leader 心跳已启动", WORKER_ID)
  46. else:
  47. logger.info("[%s] 不是 leader,等待缓存就绪…", WORKER_ID)
  48. ready = await wait_cache_ready()
  49. if ready:
  50. logger.info("[%s] 共享缓存已就绪", WORKER_ID)
  51. else:
  52. logger.warning(
  53. "[%s] 等待缓存就绪超时,搜索将以空 sec_ids 降级运行",
  54. WORKER_ID,
  55. )
  56. try:
  57. yield
  58. finally:
  59. logger.info("[%s] 关闭:取消后台任务…", WORKER_ID)
  60. for task in tasks:
  61. task.cancel()
  62. for task in tasks:
  63. try:
  64. await task
  65. except asyncio.CancelledError:
  66. pass
  67. if is_leader:
  68. # 主动释放 leader 锁,让其他 worker 可以更快接管
  69. await release_lock(KEY_LOCK_LEADER, WORKER_ID)
  70. await close_redis()
  71. app = FastAPI(
  72. title="Zendesk FAQ Keyword Search",
  73. description="封装 Zendesk Help Center 搜索;FAQ sec_ids 共享缓存(Redis),每日定时刷新。",
  74. version="1.0.0",
  75. lifespan=lifespan,
  76. )
  77. # 注册全局异常处理器(统一错误响应格式 {code, data, msg})
  78. register_exception_handlers(app)
  79. # 路由
  80. app.include_router(search.router)
  81. @app.get(
  82. "/health",
  83. response_model=ApiResponse[dict[str, str]],
  84. tags=["meta"],
  85. )
  86. async def health() -> ApiResponse[dict[str, str]]:
  87. """健康检查。"""
  88. return ApiResponse.ok(data={"status": "ok"})
  89. @app.get(
  90. "/health/cache",
  91. response_model=ApiResponse[CacheStatus],
  92. tags=["meta"],
  93. )
  94. async def cache_status() -> ApiResponse[CacheStatus]:
  95. """缓存状态查看。"""
  96. snap = await get_snapshot()
  97. return ApiResponse.ok(
  98. data=CacheStatus(
  99. sec_ids_count=len(snap["sec_ids"]),
  100. last_updated_at=snap["last_updated_at"],
  101. next_refresh_at=snap["next_refresh_at"],
  102. )
  103. )
  104. @app.post(
  105. "/admin/cache/refresh",
  106. response_model=ApiResponse[CacheRefreshResult],
  107. tags=["meta"],
  108. )
  109. async def refresh_cache_now() -> ApiResponse[CacheRefreshResult]:
  110. """手动触发刷新(运维 / 测试用途)。
  111. 使用 force=True 跳过去重锁,确保管理员请求立即生效。
  112. """
  113. await refresh_cache(force=True)
  114. snap = await get_snapshot()
  115. return ApiResponse.ok(
  116. data=CacheRefreshResult(
  117. sec_ids_count=len(snap["sec_ids"]),
  118. last_updated_at=snap["last_updated_at"],
  119. ),
  120. msg="缓存刷新成功",
  121. )