Files
34047007@qq.com 73e25050df checkpoint: 调度任务存储落 PG——jobstore 归一 default + execute 写回真实触发 + 前端 store 字典更新
- ap_scheduler: 默认分组改 SQLAlchemyJobStore on PG,删 RedisJobStore 与 redis_instance;注释/告警口径同步
- node/service: jobstore 归一 default(兼容存量 sqlalchemy 名);cron/interval/date 写回 task_node.trigger/trigger_args(前端列表可见真实触发);run_job_now 用独立临时 id,不再覆盖节点既有排程
- enums: 删 APSCHEDULER_LOCK_KEY(jobstore 落 PG 后死锁)
- sys_dict_data: sys_job_store 字典归一 default,删 sqlalchemy 档
- 前端: 任务表单 store 选项与 default 归一
2026-08-07 21:47:38 +08:00

164 lines
6.2 KiB
Python

from collections.abc import AsyncGenerator
from typing import Any
from fastapi import FastAPI
from fastapi.concurrency import asynccontextmanager
from fastapi.openapi.docs import get_redoc_html, get_swagger_ui_html, get_swagger_ui_oauth2_redirect_html
from fastapi.responses import HTMLResponse
from fastapi.staticfiles import StaticFiles
from .common.enums import EnvironmentEnum
from .config import path_conf
from .config.setting import settings
from .core.exceptions import handle_exception
from .core.logger import logger
from .utils.common_util import import_module
from .utils.console import console_end, console_start
@asynccontextmanager
async def lifespan(app: FastAPI) -> AsyncGenerator[Any, Any]:
from app.api.v1.module_system.dict.service import DictDataService
from app.api.v1.module_system.params.service import ParamsService
from app.core.ap_scheduler import SchedulerUtil
from app.core.database import async_engine, redis_connect
from app.scripts.initialize import InitializeData
await InitializeData().init_db()
logger.info("✅ {}数据库初始化完成", settings.DATABASE_TYPE)
# Redis 启动不强依赖:连接失败降级启动(参数/字典读接口回源数据库,
# 登录会话/在线监控/AI配置/调度任务等存储类接口返回 503),不再中断应用。
await redis_connect(app, status=True)
redis_ready = app.state.redis is not None
if redis_ready:
logger.info("✅ Redis 连接初始化完成")
await ParamsService.init_cache(redis=app.state.redis)
logger.info("✅ Redis系统参数初始化完成")
await DictDataService.init_cache(redis=app.state.redis)
logger.info("✅ Redis数据字典初始化完成")
else:
logger.warning(
"⚠️ Redis 不可用,以降级模式启动:参数/字典读接口回源数据库;"
"登录会话/在线监控/AI配置接口将返回 503(调度任务 jobstore 已落 PG,不受影响)"
)
from app.utils.dict_util import preload_value2label_cache
await preload_value2label_cache()
logger.info("✅ 数据字典导出映射预载完成")
scheduler_ready = False
try:
await SchedulerUtil.init_scheduler()
scheduler_ready = SchedulerUtil.is_running()
logger.info("✅ 定时任务调度器初始化完成")
except Exception as e:
# 调度器初始化失败(jobstore 已落 PG,与主库同生共死):不阻断启动
SchedulerUtil.shutdown(wait=False)
logger.error("❌ 定时任务调度器初始化失败(继续启动,任务接口将返回 503): {}", e)
console_start(
host=settings.SERVER_HOST,
port=settings.SERVER_PORT,
reload=settings.DEBUG,
database_ready=True,
redis_ready=redis_ready,
scheduler_ready=scheduler_ready,
)
yield
try:
SchedulerUtil.shutdown(wait=True)
logger.info("✅ 定时任务调度器已关闭")
await redis_connect(app, status=False)
logger.info("✅ Redis 连接已关闭")
await async_engine.dispose()
logger.info("✅ 数据库引擎连接池已释放")
console_end()
except Exception as e:
logger.error("❌ 应用关闭过程中发生错误: {}", e)
raise SystemExit(1)
def register_middlewares(app: FastAPI) -> None:
for middleware in settings.MIDDLEWARE_LIST[::-1]:
if not middleware:
continue
middleware = import_module(middleware, desc="中间件")
app.add_middleware(middleware)
def register_exceptions(app: FastAPI) -> None:
handle_exception(app)
def register_routers(app: FastAPI) -> None:
from app.api.v1.module_ai import ai_router
from app.api.v1.module_bre import bre_router
from app.api.v1.module_bre.brapi.controller import brapi_router
from app.api.v1.module_common import common_router
from app.api.v1.module_generator import generator_router
from app.api.v1.module_monitor import monitor_router
from app.api.v1.module_system import system_router
from app.api.v1.module_task import task_router
app.include_router(common_router)
app.include_router(monitor_router)
app.include_router(system_router)
app.include_router(ai_router)
app.include_router(generator_router)
app.include_router(task_router)
app.include_router(bre_router)
app.include_router(brapi_router)
from app.core.discover import dynamic_router
dynamic_router.init_app(app)
def register_static(app: FastAPI) -> None:
"""注册静态文件路由。"""
path_conf.STATIC_DIR.mkdir(parents=True, exist_ok=True)
app.mount(path=settings.STATIC_URL, app=StaticFiles(directory=path_conf.STATIC_DIR), name=path_conf.STATIC_DIR.name)
def register_docs(app: FastAPI) -> None:
"""注册文档路由。
生产环境不注册 /docs /redoc:避免暴露全部接口定义(攻击者据此定位薄弱点)。
"""
if settings.ENVIRONMENT == EnvironmentEnum.PROD:
return
swagger_ui_redirect_url = str(app.swagger_ui_oauth2_redirect_url)
root_openapi_url = str(app.root_path) + str(app.openapi_url)
@app.get(swagger_ui_redirect_url, include_in_schema=False)
async def swagger_ui_redirect():
return get_swagger_ui_oauth2_redirect_html()
@app.get(settings.DOCS_URL, include_in_schema=False)
async def custom_swagger_ui_html() -> HTMLResponse:
return get_swagger_ui_html(
openapi_url=root_openapi_url,
title=app.title + " - Swagger UI",
oauth2_redirect_url=app.swagger_ui_oauth2_redirect_url,
swagger_js_url=settings.SWAGGER_JS_URL,
swagger_css_url=settings.SWAGGER_CSS_URL,
swagger_favicon_url=settings.FAVICON_URL,
)
@app.get(settings.REDOC_URL, include_in_schema=False)
async def custom_redoc_html():
return get_redoc_html(
openapi_url=root_openapi_url,
title=app.title + " - ReDoc",
redoc_js_url=settings.REDOC_JS_URL,
redoc_favicon_url=settings.FAVICON_URL,
)
def register_frontend(app: FastAPI) -> None:
if path_conf.FRONTEND_DIST_DIR.exists():
app.mount("/web", StaticFiles(directory=str(path_conf.FRONTEND_DIST_DIR), html=True), name="frontend")