tasks manager
This commit is contained in:
@@ -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()
|
||||
|
||||
|
||||
Binary file not shown.
@@ -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()
|
||||
@@ -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
|
||||
@@ -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")
|
||||
@@ -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()
|
||||
@@ -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)
|
||||
@@ -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()
|
||||
@@ -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
|
||||
@@ -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")
|
||||
|
||||
@@ -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()
|
||||
+8
-6
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user