diff --git a/core/tasks/periodic_manager.py b/core/tasks/periodic_manager.py index faff4dec..6e444d0c 100644 --- a/core/tasks/periodic_manager.py +++ b/core/tasks/periodic_manager.py @@ -68,12 +68,14 @@ def _run_process_loop_task(task_id: str, runner: LoopRunner) -> None: async def _run_process_loop_task_async(task_id: str, runner: LoopRunner) -> None: + from database.db import reset_async_db_engine from bot import bot from core.bootstrap import bootstrap from database import async_session_maker, init_db logger.info("[PeriodicManager] Process-loop задача {} запущена, PID={}", task_id, os.getpid()) try: + reset_async_db_engine() await init_db() await bootstrap() await runner(bot, async_session_maker) diff --git a/database/__init__.py b/database/__init__.py index 77523c17..d9ec33ca 100644 --- a/database/__init__.py +++ b/database/__init__.py @@ -1,7 +1,7 @@ from .audit import * from .bans import * from .coupons import * -from .db import async_session_maker +from .db import Base, async_session_maker, engine, reset_async_db_engine from .gifts import * from . import identities from .hot_leads import * diff --git a/database/db.py b/database/db.py index 6aaa5640..d1264343 100644 --- a/database/db.py +++ b/database/db.py @@ -19,17 +19,23 @@ if USE_PGBOUNCER and "+asyncpg" in DATABASE_URL: _db_url = f"{_db_url}{sep}prepared_statement_cache_size=0" _pool_recycle = 60 if USE_PGBOUNCER else 300 -engine = create_async_engine( - _db_url, - echo=False, - future=True, - pool_size=DB_POOL_SIZE, - max_overflow=DB_MAX_OVERFLOW, - pool_timeout=60, - pool_pre_ping=True, - pool_recycle=_pool_recycle, - connect_args=_connect_args, -) + + +def _create_engine(): + return create_async_engine( + _db_url, + echo=False, + future=True, + pool_size=DB_POOL_SIZE, + max_overflow=DB_MAX_OVERFLOW, + pool_timeout=60, + pool_pre_ping=True, + pool_recycle=_pool_recycle, + connect_args=_connect_args, + ) + + +engine = _create_engine() async_session_maker = async_sessionmaker( bind=engine, @@ -42,6 +48,12 @@ Base = declarative_base() WARM_POOL_COUNT = 10 +def reset_async_db_engine() -> None: + global engine + engine = _create_engine() + async_session_maker.configure(bind=engine) + + async def warm_pool() -> None: """ Прогревает пул соединений при старте. diff --git a/database/init_db.py b/database/init_db.py index c278d755..e64284a1 100644 --- a/database/init_db.py +++ b/database/init_db.py @@ -3,15 +3,15 @@ from datetime import datetime from sqlalchemy import select from config import ADMIN_ID -from database.db import async_session_maker, engine +from database import db from database.models import Admin, Base, User async def init_db(): - async with engine.begin() as conn: + async with db.engine.begin() as conn: await conn.run_sync(Base.metadata.create_all) - async with async_session_maker() as session: + async with db.async_session_maker() as session: result = await session.execute(select(User).where(User.tg_id == 0)) if not result.scalar_one_or_none(): session.add(