diff --git a/bot.py b/bot.py index b15ea11d..df8c61d5 100644 --- a/bot.py +++ b/bot.py @@ -36,13 +36,13 @@ dp.message.filter(IsPrivateFilter()) dp.callback_query.filter(IsPrivateFilter()) async def _on_dispatcher_startup(*_args, **_kwargs): - from handlers.notifications.task_manager import ensure_periodic_task_manager_started + from core.tasks import ensure_periodic_task_manager_started await ensure_periodic_task_manager_started(bot, async_session_maker) async def _on_dispatcher_shutdown(*_args, **_kwargs): - from handlers.notifications.task_manager import ensure_periodic_task_manager_stopped + from core.tasks import ensure_periodic_task_manager_stopped await ensure_periodic_task_manager_stopped() diff --git a/core/app.cpython-312-x86_64-linux-gnu.so b/core/app.cpython-312-x86_64-linux-gnu.so index 46277b43..d6ddfc24 100644 Binary files a/core/app.cpython-312-x86_64-linux-gnu.so and b/core/app.cpython-312-x86_64-linux-gnu.so differ diff --git a/core/periodic_manager.py b/core/periodic_manager.py deleted file mode 100644 index 09eefe6d..00000000 --- a/core/periodic_manager.py +++ /dev/null @@ -1,129 +0,0 @@ -import asyncio -import fcntl -import os - -from collections.abc import Awaitable, Callable -from dataclasses import dataclass - -from aiogram import Bot -from apscheduler.schedulers.asyncio import AsyncIOScheduler -from apscheduler.triggers.base import BaseTrigger -from sqlalchemy.ext.asyncio import async_sessionmaker - -from logger import logger - - -LoopRunner = Callable[[Bot, async_sessionmaker], Awaitable[None]] -CronRunner = Callable[[], Awaitable[None]] - - -@dataclass -class ManagedLoopTask: - task_id: str - runner: LoopRunner - - -@dataclass -class ManagedCronTask: - task_id: str - runner: CronRunner - trigger: BaseTrigger - - -class PeriodicTaskManager: - def __init__(self, timezone_name: str = "Europe/Moscow") -> None: - self.timezone_name = timezone_name - self._loop_tasks: dict[str, ManagedLoopTask] = {} - self._cron_tasks: dict[str, ManagedCronTask] = {} - self._running_tasks: dict[str, asyncio.Task] = {} - self._scheduler: AsyncIOScheduler | None = None - self._started = False - self._lock = asyncio.Lock() - self._process_lock_file = None - self._process_lock_path = "/tmp/solo_bot_periodic_manager.lock" - - def register_loop_task(self, task_id: str, runner: LoopRunner) -> None: - self._loop_tasks[task_id] = ManagedLoopTask(task_id=task_id, runner=runner) - - def register_cron_task(self, task_id: str, runner: CronRunner, trigger: BaseTrigger) -> None: - self._cron_tasks[task_id] = ManagedCronTask(task_id=task_id, runner=runner, trigger=trigger) - - def _acquire_process_lock(self) -> bool: - if self._process_lock_file is not None: - return True - lock_file = open(self._process_lock_path, "a+", encoding="utf-8") - try: - fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) - lock_file.seek(0) - lock_file.truncate() - lock_file.write(str(os.getpid())) - lock_file.flush() - self._process_lock_file = lock_file - return True - except OSError: - lock_file.close() - return False - - def _release_process_lock(self) -> None: - if self._process_lock_file is None: - return - try: - fcntl.flock(self._process_lock_file.fileno(), fcntl.LOCK_UN) - except OSError: - pass - try: - self._process_lock_file.close() - except OSError: - pass - self._process_lock_file = None - - async def start(self, bot: Bot, sessionmaker: async_sessionmaker) -> None: - async with self._lock: - if self._started: - return - if not self._acquire_process_lock(): - logger.info("[PeriodicManager] Уже запущен в другом процессе, текущий запуск пропущен") - return - scheduler = AsyncIOScheduler(timezone=self.timezone_name) - for cron_task in self._cron_tasks.values(): - scheduler.add_job( - cron_task.runner, - cron_task.trigger, - id=cron_task.task_id, - replace_existing=True, - max_instances=1, - coalesce=True, - ) - scheduler.start() - self._scheduler = scheduler - for loop_task in self._loop_tasks.values(): - self._running_tasks[loop_task.task_id] = asyncio.create_task(loop_task.runner(bot, sessionmaker)) - self._started = True - logger.info( - "[PeriodicManager] Запущен: loop=%s cron=%s", - len(self._loop_tasks), - len(self._cron_tasks), - ) - - async def stop(self) -> None: - async with self._lock: - if not self._started: - return - if self._scheduler is not None: - try: - self._scheduler.shutdown(wait=False) - except Exception: - pass - self._scheduler = None - tasks = list(self._running_tasks.values()) - self._running_tasks.clear() - for task in tasks: - task.cancel() - if tasks: - await asyncio.gather(*tasks, return_exceptions=True) - self._started = False - self._release_process_lock() - logger.info("[PeriodicManager] Остановлен") - - -periodic_task_manager = PeriodicTaskManager() diff --git a/core/tasks/__init__.py b/core/tasks/__init__.py new file mode 100644 index 00000000..fe62e969 --- /dev/null +++ b/core/tasks/__init__.py @@ -0,0 +1,3 @@ +from .periodic_manager import PeriodicTaskManager, periodic_task_manager +from .registry import register_periodic_tasks +from .lifecycle import ensure_periodic_task_manager_started, ensure_periodic_task_manager_stopped diff --git a/core/tasks/cron_tasks.py b/core/tasks/cron_tasks.py new file mode 100644 index 00000000..cf9da1b5 --- /dev/null +++ b/core/tasks/cron_tasks.py @@ -0,0 +1,45 @@ +import asyncio + +from apscheduler.triggers.cron import CronTrigger + +from database import async_session_maker, cancel_expired_pending_payments +from logger import logger + + +async def scheduled_audit_drain() -> None: + from audit import drain_audit_redis_to_db + + try: + drained = await drain_audit_redis_to_db(async_session_maker) + logger.info("[AuditDrain] Ночной drain завершён, записано событий: {}", drained) + except Exception as error: + logger.error("[AuditDrain] Ошибка ночного drain: {}", error) + + +async def scheduled_stats_report() -> None: + from handlers.admin.stats.stats_handler import send_daily_stats_report + + async with async_session_maker() as session: + await send_daily_stats_report(session) + + +async def sweep_stale_payments_job() -> None: + async with async_session_maker() as session: + await cancel_expired_pending_payments(session) + + +def scheduled_audit_drain_process_runner() -> None: + asyncio.run(scheduled_audit_drain()) + + +def scheduled_stats_report_process_runner() -> None: + asyncio.run(scheduled_stats_report()) + + +def sweep_stale_payments_process_runner() -> None: + asyncio.run(sweep_stale_payments_job()) + + +AUDIT_DRAIN_TRIGGER = CronTrigger(hour=0, minute=0, timezone="Europe/Moscow") +DAILY_STATS_REPORT_TRIGGER = CronTrigger(hour=0, minute=1, timezone="Europe/Moscow") +STALE_PAYMENTS_SWEEP_TRIGGER = CronTrigger(minute=0, timezone="Europe/Moscow") diff --git a/core/tasks/lifecycle.py b/core/tasks/lifecycle.py new file mode 100644 index 00000000..09976494 --- /dev/null +++ b/core/tasks/lifecycle.py @@ -0,0 +1,36 @@ +import os + +from core.tasks.periodic_manager import periodic_task_manager +from core.tasks.registry import register_periodic_tasks +from hooks.hooks import register_hook +from logger import logger + + +def _should_start_manager() -> bool: + if os.environ.get("NOTIFICATION_WORKER", "").strip() == "1": + return True + if os.environ.get("NOTIFICATION_WORKER_SEPARATE", "").strip() == "1": + return False + return True + + +async def ensure_periodic_task_manager_started(bot, sessionmaker) -> None: + register_periodic_tasks() + if not _should_start_manager(): + logger.info("[PeriodicManager] Пропуск старта в текущем процессе") + return + await periodic_task_manager.start(bot, sessionmaker) + + +async def ensure_periodic_task_manager_stopped() -> None: + await periodic_task_manager.stop() + + +@register_hook("startup") +async def start_periodic_task_manager(bot, sessionmaker, **_kwargs): + await ensure_periodic_task_manager_started(bot, sessionmaker) + + +@register_hook("shutdown") +async def stop_periodic_task_manager(**_kwargs): + await ensure_periodic_task_manager_stopped() diff --git a/core/tasks/loop_tasks.py b/core/tasks/loop_tasks.py new file mode 100644 index 00000000..4c67081a --- /dev/null +++ b/core/tasks/loop_tasks.py @@ -0,0 +1,64 @@ +import asyncio + +from logger import logger + + +async def notifications_loop(bot, sessionmaker) -> None: + from handlers.notifications.general_notifications import periodic_notifications + + await periodic_notifications(bot, sessionmaker=sessionmaker) + + +async def scheduled_broadcasts_loop_task(bot, _sessionmaker) -> None: + from handlers.admin.sender.scheduled_service import scheduled_broadcasts_loop + + await scheduled_broadcasts_loop(bot) + + +async def backup_loop(bot, _sessionmaker) -> None: + from config import BACKUP_TIME + from utils.backup import backup_database + + if BACKUP_TIME <= 0: + await asyncio.Event().wait() + return + while True: + error = await backup_database(bot_instance=bot) + if error: + logger.error("[Backup] Ошибка: {}", error) + await asyncio.sleep(BACKUP_TIME) + + +def backup_thread_loop(stop_event, _bot, _sessionmaker) -> None: + from aiogram import Bot + from aiogram.client.default import DefaultBotProperties + from aiogram.enums import ParseMode + from config import API_TOKEN, BACKUP_TIME + from utils.backup import backup_database + + if BACKUP_TIME <= 0: + stop_event.wait() + return + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + backup_bot = Bot(token=API_TOKEN, default=DefaultBotProperties(parse_mode=ParseMode.HTML)) + try: + while not stop_event.is_set(): + error = loop.run_until_complete(backup_database(bot_instance=backup_bot)) + if error: + logger.error("[Backup] Ошибка: {}", error) + if stop_event.wait(BACKUP_TIME): + break + finally: + loop.run_until_complete(backup_bot.session.close()) + loop.close() + + +async def server_checks_loop(_bot, sessionmaker) -> None: + from config import PING_TIME + from servers import check_servers + + if PING_TIME <= 0: + await asyncio.Event().wait() + return + await check_servers(sessionmaker=sessionmaker) diff --git a/core/tasks/periodic_manager.py b/core/tasks/periodic_manager.py new file mode 100644 index 00000000..faff4dec --- /dev/null +++ b/core/tasks/periodic_manager.py @@ -0,0 +1,313 @@ +import asyncio +import fcntl +import inspect +import multiprocessing +import os +import threading + +from collections.abc import Awaitable, Callable +from dataclasses import dataclass +from typing import Literal + +from aiogram import Bot +from apscheduler.executors.asyncio import AsyncIOExecutor +from apscheduler.executors.pool import ProcessPoolExecutor as APSchedulerProcessPoolExecutor +from apscheduler.executors.pool import ThreadPoolExecutor as APSchedulerThreadPoolExecutor +from apscheduler.schedulers.asyncio import AsyncIOScheduler +from apscheduler.triggers.base import BaseTrigger +from sqlalchemy.ext.asyncio import async_sessionmaker + +from logger import logger + + +LoopRunner = Callable[[Bot, async_sessionmaker], Awaitable[None]] +ThreadLoopRunner = Callable[[threading.Event, Bot, async_sessionmaker], None] +CronRunner = Callable[[], Awaitable[None]] | Callable[[], None] +CronExecutionMode = Literal["async", "thread", "process"] + + +@dataclass +class ManagedLoopTask: + task_id: str + runner: LoopRunner + + +@dataclass +class ManagedThreadLoopTask: + task_id: str + runner: ThreadLoopRunner + + +@dataclass +class ManagedProcessLoopTask: + task_id: str + runner: LoopRunner + + +@dataclass +class ManagedCronTask: + task_id: str + runner: CronRunner + trigger: BaseTrigger + execution_mode: CronExecutionMode + + +@dataclass +class RunningThreadLoopTask: + thread: threading.Thread + stop_event: threading.Event + + +@dataclass +class RunningProcessLoopTask: + process: multiprocessing.Process + + +def _run_process_loop_task(task_id: str, runner: LoopRunner) -> None: + asyncio.run(_run_process_loop_task_async(task_id, runner)) + + +async def _run_process_loop_task_async(task_id: str, runner: LoopRunner) -> None: + 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: + await init_db() + await bootstrap() + await runner(bot, async_session_maker) + except Exception as error: + logger.error("[PeriodicManager] Ошибка process-loop задачи {}: {}", task_id, error) + raise + finally: + try: + await bot.session.close() + except Exception: + pass + + +class PeriodicTaskManager: + def __init__(self, timezone_name: str = "Europe/Moscow") -> None: + self.timezone_name = timezone_name + self._loop_tasks: dict[str, ManagedLoopTask] = {} + self._thread_loop_tasks: dict[str, ManagedThreadLoopTask] = {} + self._process_loop_tasks: dict[str, ManagedProcessLoopTask] = {} + self._cron_tasks: dict[str, ManagedCronTask] = {} + self._running_tasks: dict[str, asyncio.Task] = {} + self._running_thread_tasks: dict[str, RunningThreadLoopTask] = {} + self._running_process_tasks: dict[str, RunningProcessLoopTask] = {} + self._scheduler: AsyncIOScheduler | None = None + self._scheduler_process_workers: int | None = None + self._started = False + self._lock = asyncio.Lock() + self._process_lock_file = None + self._process_lock_path = "/tmp/solo_bot_periodic_manager.lock" + + def register_loop_task(self, task_id: str, runner: LoopRunner) -> None: + self._loop_tasks[task_id] = ManagedLoopTask(task_id=task_id, runner=runner) + + def register_thread_loop_task(self, task_id: str, runner: ThreadLoopRunner) -> None: + self._thread_loop_tasks[task_id] = ManagedThreadLoopTask(task_id=task_id, runner=runner) + + def register_process_loop_task(self, task_id: str, runner: LoopRunner) -> None: + self._process_loop_tasks[task_id] = ManagedProcessLoopTask(task_id=task_id, runner=runner) + + def set_scheduler_process_workers(self, workers: int | None) -> None: + self._scheduler_process_workers = None if workers is None else max(0, int(workers)) + + def register_cron_task( + self, + task_id: str, + runner: CronRunner, + trigger: BaseTrigger, + execution_mode: CronExecutionMode = "async", + ) -> None: + if execution_mode not in {"async", "thread", "process"}: + raise ValueError(f"Unsupported execution_mode: {execution_mode}") + if execution_mode != "async" and inspect.iscoroutinefunction(runner): + raise ValueError( + f"Cron task {task_id} with execution_mode={execution_mode} must be a sync function" + ) + self._cron_tasks[task_id] = ManagedCronTask( + task_id=task_id, + runner=runner, + trigger=trigger, + execution_mode=execution_mode, + ) + + def _acquire_process_lock(self) -> bool: + if self._process_lock_file is not None: + return True + lock_file = open(self._process_lock_path, "a+", encoding="utf-8") + try: + fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + lock_file.seek(0) + lock_file.truncate() + lock_file.write(str(os.getpid())) + lock_file.flush() + self._process_lock_file = lock_file + return True + except OSError: + lock_file.close() + return False + + def _release_process_lock(self) -> None: + if self._process_lock_file is None: + return + try: + fcntl.flock(self._process_lock_file.fileno(), fcntl.LOCK_UN) + except OSError: + pass + try: + self._process_lock_file.close() + except OSError: + pass + self._process_lock_file = None + + def _build_scheduler(self) -> AsyncIOScheduler: + from config import EXECUTOR_POOL_SIZE, PROCESS_POOL_SIZE + + thread_workers = max(1, int(EXECUTOR_POOL_SIZE)) + configured_process_workers = self._scheduler_process_workers + if configured_process_workers is None: + configured_process_workers = int(PROCESS_POOL_SIZE) + process_workers = max(0, min(configured_process_workers, multiprocessing.cpu_count() or 1)) + executors = { + "default": AsyncIOExecutor(), + "threadpool": APSchedulerThreadPoolExecutor(max_workers=thread_workers), + } + if process_workers > 0: + executors["processpool"] = APSchedulerProcessPoolExecutor(max_workers=process_workers) + return AsyncIOScheduler(timezone=self.timezone_name, executors=executors) + + @staticmethod + def _cron_executor_name(execution_mode: CronExecutionMode) -> str: + if execution_mode == "thread": + return "threadpool" + if execution_mode == "process": + return "processpool" + return "default" + + @staticmethod + def _thread_loop_entry( + task_id: str, + runner: ThreadLoopRunner, + stop_event: threading.Event, + bot: Bot, + sessionmaker: async_sessionmaker, + ) -> None: + try: + runner(stop_event, bot, sessionmaker) + except Exception as error: + logger.error("[PeriodicManager] Ошибка thread-loop задачи {}: {}", task_id, error) + + async def _join_thread_task(self, task_id: str, running_task: RunningThreadLoopTask) -> None: + running_task.stop_event.set() + await asyncio.to_thread(running_task.thread.join, 5) + if running_task.thread.is_alive(): + logger.warning("[PeriodicManager] Thread-loop задача {} не завершилась вовремя", task_id) + + async def _stop_process_task(self, task_id: str, running_task: RunningProcessLoopTask) -> None: + process = running_task.process + if not process.is_alive(): + await asyncio.to_thread(process.join, 1) + return + process.terminate() + await asyncio.to_thread(process.join, 5) + if process.is_alive(): + process.kill() + await asyncio.to_thread(process.join, 5) + if process.is_alive(): + logger.warning("[PeriodicManager] Process-loop задача {} не завершилась вовремя", task_id) + + async def start(self, bot: Bot, sessionmaker: async_sessionmaker) -> None: + async with self._lock: + if self._started: + return + if not self._acquire_process_lock(): + logger.info("[PeriodicManager] Уже запущен в другом процессе, текущий запуск пропущен") + return + scheduler = self._build_scheduler() + for cron_task in self._cron_tasks.values(): + scheduler.add_job( + cron_task.runner, + cron_task.trigger, + id=cron_task.task_id, + executor=self._cron_executor_name(cron_task.execution_mode), + replace_existing=True, + max_instances=1, + coalesce=True, + ) + scheduler.start() + self._scheduler = scheduler + for loop_task in self._loop_tasks.values(): + self._running_tasks[loop_task.task_id] = asyncio.create_task(loop_task.runner(bot, sessionmaker)) + for loop_task in self._thread_loop_tasks.values(): + stop_event = threading.Event() + thread = threading.Thread( + target=self._thread_loop_entry, + args=(loop_task.task_id, loop_task.runner, stop_event, bot, sessionmaker), + name=f"periodic-{loop_task.task_id}", + daemon=True, + ) + thread.start() + self._running_thread_tasks[loop_task.task_id] = RunningThreadLoopTask( + thread=thread, + stop_event=stop_event, + ) + for loop_task in self._process_loop_tasks.values(): + ctx = multiprocessing.get_context("spawn") + process = ctx.Process( + target=_run_process_loop_task, + args=(loop_task.task_id, loop_task.runner), + name=f"periodic-{loop_task.task_id}", + daemon=True, + ) + process.start() + self._running_process_tasks[loop_task.task_id] = RunningProcessLoopTask(process=process) + self._started = True + logger.info( + "[PeriodicManager] Запущен: async-loop=%s thread-loop=%s process-loop=%s cron=%s", + len(self._loop_tasks), + len(self._thread_loop_tasks), + len(self._process_loop_tasks), + len(self._cron_tasks), + ) + + async def stop(self) -> None: + async with self._lock: + if not self._started: + return + if self._scheduler is not None: + try: + self._scheduler.shutdown(wait=False) + except Exception: + pass + self._scheduler = None + tasks = list(self._running_tasks.values()) + self._running_tasks.clear() + for task in tasks: + task.cancel() + if tasks: + await asyncio.gather(*tasks, return_exceptions=True) + thread_tasks = list(self._running_thread_tasks.items()) + self._running_thread_tasks.clear() + if thread_tasks: + await asyncio.gather( + *(self._join_thread_task(task_id, running_task) for task_id, running_task in thread_tasks), + return_exceptions=True, + ) + process_tasks = list(self._running_process_tasks.items()) + self._running_process_tasks.clear() + if process_tasks: + await asyncio.gather( + *(self._stop_process_task(task_id, running_task) for task_id, running_task in process_tasks), + return_exceptions=True, + ) + self._started = False + self._release_process_lock() + logger.info("[PeriodicManager] Остановлен") + + +periodic_task_manager = PeriodicTaskManager() diff --git a/core/tasks/registry.py b/core/tasks/registry.py new file mode 100644 index 00000000..f2f3b710 --- /dev/null +++ b/core/tasks/registry.py @@ -0,0 +1,101 @@ +from core.tasks.cron_tasks import ( + AUDIT_DRAIN_TRIGGER, + DAILY_STATS_REPORT_TRIGGER, + STALE_PAYMENTS_SWEEP_TRIGGER, + scheduled_audit_drain, + scheduled_audit_drain_process_runner, + scheduled_stats_report, + scheduled_stats_report_process_runner, + sweep_stale_payments_job, + sweep_stale_payments_process_runner, +) +from core.tasks.loop_tasks import ( + backup_loop, + backup_thread_loop, + notifications_loop, + scheduled_broadcasts_loop_task, + server_checks_loop, +) +from core.tasks.periodic_manager import periodic_task_manager + + +_TASKS_REGISTERED = False + + +def register_periodic_tasks() -> None: + global _TASKS_REGISTERED + if _TASKS_REGISTERED: + return + from config import PROCESS_POOL_SIZE + + process_budget = max(0, int(PROCESS_POOL_SIZE) if int(PROCESS_POOL_SIZE) > 1 else 0) + + if process_budget > 0: + periodic_task_manager.register_process_loop_task("notifications", notifications_loop) + process_budget -= 1 + else: + periodic_task_manager.register_loop_task("notifications", notifications_loop) + + if process_budget > 0: + periodic_task_manager.register_process_loop_task("scheduled_broadcasts", scheduled_broadcasts_loop_task) + process_budget -= 1 + else: + periodic_task_manager.register_loop_task("scheduled_broadcasts", scheduled_broadcasts_loop_task) + + if process_budget > 0: + periodic_task_manager.register_process_loop_task("backup", backup_loop) + process_budget -= 1 + else: + periodic_task_manager.register_thread_loop_task("backup", backup_thread_loop) + + if process_budget > 0: + periodic_task_manager.register_process_loop_task("server_checks", server_checks_loop) + process_budget -= 1 + else: + periodic_task_manager.register_loop_task("server_checks", server_checks_loop) + + periodic_task_manager.set_scheduler_process_workers(process_budget) + + if process_budget > 0: + periodic_task_manager.register_cron_task( + "audit_drain_midnight", + scheduled_audit_drain_process_runner, + AUDIT_DRAIN_TRIGGER, + execution_mode="process", + ) + else: + periodic_task_manager.register_cron_task( + "audit_drain_midnight", + scheduled_audit_drain, + AUDIT_DRAIN_TRIGGER, + ) + + if process_budget > 0: + periodic_task_manager.register_cron_task( + "daily_stats_report", + scheduled_stats_report_process_runner, + DAILY_STATS_REPORT_TRIGGER, + execution_mode="process", + ) + else: + periodic_task_manager.register_cron_task( + "daily_stats_report", + scheduled_stats_report, + DAILY_STATS_REPORT_TRIGGER, + ) + + if process_budget > 0: + periodic_task_manager.register_cron_task( + "sweep_stale_payments", + sweep_stale_payments_process_runner, + STALE_PAYMENTS_SWEEP_TRIGGER, + execution_mode="process", + ) + else: + periodic_task_manager.register_cron_task( + "sweep_stale_payments", + sweep_stale_payments_job, + STALE_PAYMENTS_SWEEP_TRIGGER, + ) + + _TASKS_REGISTERED = True diff --git a/handlers/notifications/__init__.py b/handlers/notifications/__init__.py index 584275a4..600caeb0 100644 --- a/handlers/notifications/__init__.py +++ b/handlers/notifications/__init__.py @@ -2,9 +2,9 @@ __all__ = ("router",) from aiogram import Router +from core.tasks import lifecycle as _task_lifecycle from .general_notifications import router as general_notifications_router from .special_notifications import router as special_notifications_router -from . import task_manager as _task_manager router = Router(name="notifications_main_router") diff --git a/handlers/notifications/task_manager.py b/handlers/notifications/task_manager.py deleted file mode 100644 index 6341b4b8..00000000 --- a/handlers/notifications/task_manager.py +++ /dev/null @@ -1,120 +0,0 @@ -import asyncio -import os - -from apscheduler.triggers.cron import CronTrigger - -from core.periodic_manager import periodic_task_manager -from database import async_session_maker, cancel_expired_pending_payments -from hooks.hooks import register_hook -from logger import logger - - -_TASKS_REGISTERED = False - - -async def _backup_loop(_bot, _sessionmaker) -> None: - from config import BACKUP_TIME - from utils.backup import backup_database - - if BACKUP_TIME <= 0: - await asyncio.Event().wait() - return - while True: - error = await backup_database() - if error: - logger.error("[Backup] Ошибка: {}", error) - await asyncio.sleep(BACKUP_TIME) - - -async def _server_checks_loop(_bot, sessionmaker) -> None: - from config import PING_TIME - from servers import check_servers - - if PING_TIME <= 0: - await asyncio.Event().wait() - return - await check_servers(sessionmaker=sessionmaker) - - -async def _scheduled_audit_drain() -> None: - from audit import drain_audit_redis_to_db - - try: - drained = await drain_audit_redis_to_db(async_session_maker) - logger.info("[AuditDrain] Ночной drain завершён, записано событий: {}", drained) - except Exception as error: - logger.error("[AuditDrain] Ошибка ночного drain: {}", error) - - -async def _scheduled_stats_report() -> None: - from handlers.admin.stats.stats_handler import send_daily_stats_report - - async with async_session_maker() as session: - await send_daily_stats_report(session) - - -async def _sweep_stale_payments_job() -> None: - async with async_session_maker() as session: - await cancel_expired_pending_payments(session) - - -def _register_periodic_tasks() -> None: - global _TASKS_REGISTERED - if _TASKS_REGISTERED: - return - from handlers.admin.sender.scheduled_service import scheduled_broadcasts_loop - from handlers.notifications.general_notifications import periodic_notifications - - periodic_task_manager.register_loop_task( - "notifications", - lambda bot, sessionmaker: periodic_notifications(bot, sessionmaker=sessionmaker), - ) - periodic_task_manager.register_loop_task("scheduled_broadcasts", lambda bot, sessionmaker: scheduled_broadcasts_loop(bot)) - periodic_task_manager.register_loop_task("backup", _backup_loop) - periodic_task_manager.register_loop_task("server_checks", _server_checks_loop) - periodic_task_manager.register_cron_task( - "audit_drain_midnight", - _scheduled_audit_drain, - CronTrigger(hour=0, minute=0, timezone="Europe/Moscow"), - ) - periodic_task_manager.register_cron_task( - "daily_stats_report", - _scheduled_stats_report, - CronTrigger(hour=0, minute=1, timezone="Europe/Moscow"), - ) - periodic_task_manager.register_cron_task( - "sweep_stale_payments", - _sweep_stale_payments_job, - CronTrigger(minute=0, timezone="Europe/Moscow"), - ) - _TASKS_REGISTERED = True - - -def _should_start_manager() -> bool: - if os.environ.get("NOTIFICATION_WORKER", "").strip() == "1": - return True - if os.environ.get("NOTIFICATION_WORKER_SEPARATE", "").strip() == "1": - return False - return True - - -async def ensure_periodic_task_manager_started(bot, sessionmaker) -> None: - _register_periodic_tasks() - if not _should_start_manager(): - logger.info("[PeriodicManager] Пропуск старта в текущем процессе") - return - await periodic_task_manager.start(bot, sessionmaker) - - -async def ensure_periodic_task_manager_stopped() -> None: - await periodic_task_manager.stop() - - -@register_hook("startup") -async def start_periodic_task_manager(bot, sessionmaker, **_kwargs): - await ensure_periodic_task_manager_started(bot, sessionmaker) - - -@register_hook("shutdown") -async def stop_periodic_task_manager(**_kwargs): - await ensure_periodic_task_manager_stopped() diff --git a/utils/backup.py b/utils/backup.py index b25a7c5a..e2e1552b 100644 --- a/utils/backup.py +++ b/utils/backup.py @@ -32,7 +32,7 @@ from config import ( from logger import logger -async def backup_database() -> Exception | None: +async def backup_database(bot_instance: Bot | None = None) -> Exception | None: """ Создает резервную копию базы данных (или полный архив) и отправляет его администраторам. Блокирующие операции (pg_dump и т.д.) выполняются в пуле процессов, не блокируя event loop и используя другие ядра CPU. @@ -56,7 +56,7 @@ async def backup_database() -> Exception | None: logger.info("[Backup] Файл создан: {}", backup_file_path) try: - await _send_backup_to_admins(backup_file_path) + await _send_backup_to_admins(backup_file_path, bot_instance=bot_instance) exception = await run_io(_cleanup_old_backups) if exception: @@ -234,7 +234,7 @@ async def create_backup_and_send_to_admins(client) -> None: await client.database.export() -async def _send_backup_to_admins(backup_file_path: str) -> None: +async def _send_backup_to_admins(backup_file_path: str, bot_instance: Bot | None = None) -> None: """ Отправляет файл бэкапа всем администраторам через Telegram. @@ -247,12 +247,14 @@ async def _send_backup_to_admins(backup_file_path: str) -> None: if not backup_file_path or not os.path.exists(backup_file_path): raise FileNotFoundError(f"Файл бэкапа не найден: {backup_file_path}") - from bot import bot + active_bot = bot_instance + if active_bot is None: + from bot import bot as active_bot async def send_default(): for admin_id in ADMIN_ID: try: - await bot.send_document(chat_id=admin_id, document=backup_input_file) + await active_bot.send_document(chat_id=admin_id, document=backup_input_file) logger.info("[Backup] Отправлено админу: {}", admin_id) except Exception as e: logger.error("[Backup] Не отправлено админу {}: {}", admin_id, e) @@ -280,7 +282,7 @@ async def _send_backup_to_admins(backup_file_path: str) -> None: if BACKUP_CAPTION: send_kwargs["caption"] = BACKUP_CAPTION try: - await bot.send_document(**send_kwargs) + await active_bot.send_document(**send_kwargs) logger.info("[Backup] Отправлено в канал: {} (топик: {})", channel_id, thread_id) except Exception as e: logger.error("[Backup] Не отправлено в канал {}: {}, fallback", channel_id, e)