diff --git a/api/v1/routes/management.py b/api/v1/routes/management.py index 4bf26774..83a15929 100644 --- a/api/v1/routes/management.py +++ b/api/v1/routes/management.py @@ -18,7 +18,7 @@ from api.depends import get_session, verify_admin_token, verify_admin_token_shor from database import async_session_maker from config import API_TOKEN, BOT_SERVICE from core.bootstrap import MANAGEMENT_CONFIG -from core.executor import get_thread_pool +from core.executor import run_io from core.settings.management_config import update_management_config from database.models import Key, User from database.models import Server @@ -65,14 +65,7 @@ async def _restart_bot() -> None: is_systemd = parent and "systemd" in parent.name().lower() if is_systemd: - loop = asyncio.get_running_loop() - await loop.run_in_executor( - get_thread_pool(), - lambda: subprocess.run( - ["sudo", "systemctl", "restart", BOT_SERVICE], - check=True, - ), - ) + await run_io(lambda: subprocess.run(["sudo", "systemctl", "restart", BOT_SERVICE], check=True)) else: python_exe = sys.executable script_path = os.path.abspath(sys.argv[0]) diff --git a/api/v1/routes/modules.py b/api/v1/routes/modules.py index 4836dcb8..11e0fd44 100644 --- a/api/v1/routes/modules.py +++ b/api/v1/routes/modules.py @@ -6,6 +6,7 @@ from fastapi import APIRouter, Depends, HTTPException from pydantic import BaseModel from api.depends import verify_admin_token +from core.executor import run_io from utils.modules_loader import _is_safe_module_name from utils.modules_manager import manager @@ -80,24 +81,23 @@ def _read_local_module_version(name: str) -> str | None: return None -@router.get("/") -async def list_modules(admin=Depends(verify_admin_token)): - refresh = getattr(manager, "refresh_state", None) +def sync_list_modules() -> list: + """Вся синхронная работа со списком модулей (файлы, состояние). Вызывать через run_io().""" + refresh = getattr(manager, "refresh_state", None) or getattr(manager, "_load_state", None) if callable(refresh): refresh() - else: - legacy_refresh = getattr(manager, "_load_state", None) - if callable(legacy_refresh): - legacy_refresh() module_names = _available_module_names() _prune_missing_state(set(module_names)) modules = [_module_state(name) for name in module_names] - for item in modules: name = str(item.get("name") or "").strip() - local_version = _read_local_module_version(name) - item["local_version"] = local_version + item["local_version"] = _read_local_module_version(name) + return modules + +@router.get("/") +async def list_modules(admin=Depends(verify_admin_token)): + modules = await run_io(sync_list_modules) return {"items": modules} diff --git a/api/v2/routes/management.py b/api/v2/routes/management.py index 71c36517..a2459b34 100644 --- a/api/v2/routes/management.py +++ b/api/v2/routes/management.py @@ -18,7 +18,7 @@ from api.depends import get_session, verify_identity_admin, verify_identity_admi from database import async_session_maker from config import API_TOKEN, BOT_SERVICE from core.bootstrap import MANAGEMENT_CONFIG -from core.executor import get_thread_pool +from core.executor import run_io from core.settings.management_config import update_management_config from database.models import Key, Server, User from handlers.admin.sender.sender_service import BroadcastService @@ -64,14 +64,7 @@ async def _restart_bot() -> None: parent = psutil.Process(os.getpid()).parent() is_systemd = parent and "systemd" in parent.name().lower() if is_systemd: - loop = asyncio.get_running_loop() - await loop.run_in_executor( - get_thread_pool(), - lambda: subprocess.run( - ["sudo", "systemctl", "restart", BOT_SERVICE], - check=True, - ), - ) + await run_io(lambda: subprocess.run(["sudo", "systemctl", "restart", BOT_SERVICE], check=True)) else: python_exe = sys.executable script_path = os.path.abspath(sys.argv[0]) diff --git a/api/v2/routes/modules.py b/api/v2/routes/modules.py index 43ed47df..5f24035f 100644 --- a/api/v2/routes/modules.py +++ b/api/v2/routes/modules.py @@ -6,6 +6,7 @@ from fastapi import APIRouter, Depends, HTTPException from pydantic import BaseModel from api.depends import verify_identity_admin +from core.executor import run_io from utils.modules_loader import _is_safe_module_name from utils.modules_manager import manager @@ -77,23 +78,24 @@ def _read_local_module_version(name: str) -> str | None: return None -@router.get("/") -async def list_modules(identity=Depends(verify_identity_admin)): - """Список модулей с состоянием и локальной версией.""" - refresh = getattr(manager, "refresh_state", None) +def sync_list_modules() -> list: + """Вся синхронная работа со списком модулей (файлы, состояние). Вызывать через run_io().""" + refresh = getattr(manager, "refresh_state", None) or getattr(manager, "_load_state", None) if callable(refresh): refresh() - else: - legacy_refresh = getattr(manager, "_load_state", None) - if callable(legacy_refresh): - legacy_refresh() module_names = _available_module_names() _prune_missing_state(set(module_names)) modules = [_module_state(name) for name in module_names] for item in modules: name = str(item.get("name") or "").strip() - local_version = _read_local_module_version(name) - item["local_version"] = local_version + item["local_version"] = _read_local_module_version(name) + return modules + + +@router.get("/") +async def list_modules(identity=Depends(verify_identity_admin)): + """Список модулей с состоянием и локальной версией.""" + modules = await run_io(sync_list_modules) return {"items": modules} diff --git a/core/app.cpython-312-x86_64-linux-gnu.so b/core/app.cpython-312-x86_64-linux-gnu.so index f9e8a091..b1061e6c 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/executor.py b/core/executor.py index 34230ef8..9c7e881c 100644 --- a/core/executor.py +++ b/core/executor.py @@ -1,10 +1,14 @@ +import asyncio import atexit import signal import multiprocessing from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor +from typing import Callable, TypeVar from logger import logger +T = TypeVar("T") + _thread_pool: ThreadPoolExecutor | None = None _process_pool: ProcessPoolExecutor | None = None @@ -30,7 +34,7 @@ def get_thread_pool() -> ThreadPoolExecutor: from config import EXECUTOR_POOL_SIZE size = max(1, int(EXECUTOR_POOL_SIZE)) _thread_pool = ThreadPoolExecutor(max_workers=size, thread_name_prefix="bot-thread") - logger.debug("Thread pool started (workers=%s)", size) + logger.debug("[Executor] Пул потоков: {} воркеров", size) return _thread_pool @@ -40,7 +44,7 @@ def shutdown_thread_pool() -> None: if _thread_pool is not None: _thread_pool.shutdown(wait=True) _thread_pool = None - logger.debug("Thread pool shut down") + logger.debug("[Executor] Пул потоков остановлен") def get_process_pool() -> ProcessPoolExecutor: @@ -56,7 +60,7 @@ def get_process_pool() -> ProcessPoolExecutor: ctx.Process = _IgnoreSIGINTProcess _process_pool = ProcessPoolExecutor(max_workers=size, mp_context=ctx) atexit.register(_atexit_shutdown_pools) - logger.debug("Process pool started (workers=%s)", size) + logger.debug("[Executor] Пул процессов: {} воркеров", size) return _process_pool @@ -70,4 +74,16 @@ def shutdown_process_pool() -> None: pass _process_pool.shutdown(wait=True) _process_pool = None - logger.debug("Process pool shut down") + logger.debug("[Executor] Пул процессов остановлен") + + +async def run_io(fn: Callable[..., T], *args: object) -> T: + """Выполняет fn(*args) в пуле потоков (I/O). Один вызов для всех блокирующих операций.""" + loop = asyncio.get_running_loop() + return await loop.run_in_executor(get_thread_pool(), lambda: fn(*args)) + + +async def run_cpu(fn: Callable[..., T], *args: object) -> T: + """Выполняет fn(*args) в пуле процессов (CPU). fn — функция уровня модуля (для pickle).""" + loop = asyncio.get_running_loop() + return await loop.run_in_executor(get_process_pool(), fn, *args) diff --git a/database/identities.py b/database/identities.py index 91e6873d..dbec7a97 100644 --- a/database/identities.py +++ b/database/identities.py @@ -7,6 +7,7 @@ from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from config import API_TOKEN_TTL_DAYS +from core.executor import run_cpu, run_io from database.models import Admin, Identity, User @@ -90,7 +91,7 @@ async def get_identity_by_token_hash(session: AsyncSession, token_hash: str) -> async def issue_token_for_identity(session: AsyncSession, identity: Identity) -> str: """Генерирует токен, сохраняет хеш и token_issued_at в identity, возвращает токен (показать один раз).""" token = generate_token() - identity.api_token_hash = hash_token(token) + identity.api_token_hash = await run_io(hash_token, token) identity.token_issued_at = datetime.utcnow() await session.commit() await session.refresh(identity) @@ -114,7 +115,7 @@ async def create_identity_with_token( """Создаёт идентичность и выдаёт API-токен. При регистрации по почте передать email и password.""" identity = await create_identity(session, email=email, tg_id=tg_id) if password: - identity.password_hash = hash_password(password) + identity.password_hash = await run_cpu(hash_password, password) await session.commit() await session.refresh(identity) token = await issue_token_for_identity(session, identity) @@ -126,7 +127,8 @@ async def verify_identity_token(session: AsyncSession, identity_id: str, token: identity = await get_identity_by_id(session, identity_id) if not identity or not identity.api_token_hash: return None - if hash_token(token) != identity.api_token_hash: + token_hash = await run_io(hash_token, token) + if token_hash != identity.api_token_hash: return None if _is_token_expired(identity): return None @@ -136,7 +138,9 @@ async def verify_identity_token(session: AsyncSession, identity_id: str, token: async def login_by_email(session: AsyncSession, email: str, password: str) -> tuple[Identity, str] | None: """Вход по email и паролю: проверяет пароль, выдаёт новый токен; возвращает (identity, token) или None.""" identity = await get_identity_by_email(session, email) - if not identity or not check_password(password, identity.password_hash): + if not identity: + return None + if not await run_cpu(check_password, password, identity.password_hash): return None token = await issue_token_for_identity(session, identity) return identity, token diff --git a/handlers/admin/management/database.py b/handlers/admin/management/database.py index a918daf6..122d39cd 100644 --- a/handlers/admin/management/database.py +++ b/handlers/admin/management/database.py @@ -11,6 +11,7 @@ from aiogram.fsm.state import State, StatesGroup from aiogram.types import CallbackQuery, Message from config import DB_NAME, DB_PASSWORD, DB_USER, PG_HOST, PG_PORT +from core.executor import run_io from filters.admin import IsAdminFilter from logger import logger @@ -18,6 +19,62 @@ from . import router from .keyboard import AdminPanelCallback, build_back_to_db_menu, build_database_kb, build_export_db_sources_kb +def sync_restore_database( + tmp_path: str, + db_name: str, + db_user: str, + db_password: str, + pg_host: str, + pg_port: str, +) -> tuple[bool, str]: + """Восстановление БД из файла. Вызывать через run_io().""" + is_custom_dump = False + with open(tmp_path, "rb") as f: + if f.read(5) == b"PGDMP": + is_custom_dump = True + + try: + subprocess.run( + [ + "sudo", "-u", "postgres", "psql", "-d", "postgres", + "-c", + f"SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = '{db_name}' AND pid <> pg_backend_pid();", + ], + check=True, + ) + subprocess.run( + ["sudo", "-u", "postgres", "psql", "-d", "postgres", "-c", f"DROP DATABASE IF EXISTS {db_name};"], + check=True, + ) + subprocess.run( + ["sudo", "-u", "postgres", "psql", "-d", "postgres", "-c", f"CREATE DATABASE {db_name} OWNER {db_user};"], + check=True, + ) + except subprocess.CalledProcessError as e: + return False, (e.stderr or e.stdout or str(e)) + + os.environ["PGPASSWORD"] = db_password + try: + if is_custom_dump: + result = subprocess.run( + [ + "pg_restore", f"--dbname={db_name}", "-U", db_user, + "-h", pg_host, "-p", pg_port, "--no-owner", "--exit-on-error", tmp_path, + ], + capture_output=True, + text=True, + ) + else: + result = subprocess.run( + ["psql", "-U", db_user, "-h", pg_host, "-p", pg_port, "-d", db_name, "-f", tmp_path], + capture_output=True, + text=True, + ) + return result.returncode == 0, result.stderr or "" + finally: + del os.environ["PGPASSWORD"] + + class DatabaseState(StatesGroup): waiting_for_backup_file = State() @@ -53,93 +110,31 @@ async def restore_database(message: Message, state: FSMContext, bot: Bot): tmp_path = tmp_file.name await bot.download(document, destination=tmp_path) - logger.info(f"[Restore] Файл получен и сохранён: {tmp_path}") + logger.info("[Restore] Файл получен: {}", tmp_path) - is_custom_dump = False - with open(tmp_path, "rb") as f: - signature = f.read(5) - if signature == b"PGDMP": - is_custom_dump = True - - subprocess.run( - [ - "sudo", - "-u", - "postgres", - "psql", - "-d", - "postgres", - "-c", - f"SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = '{DB_NAME}' AND pid <> pg_backend_pid();", - ], - check=True, + success, err_msg = await run_io( + sync_restore_database, + tmp_path, + DB_NAME, + DB_USER, + DB_PASSWORD, + PG_HOST, + PG_PORT, ) - subprocess.run( - ["sudo", "-u", "postgres", "psql", "-d", "postgres", "-c", f"DROP DATABASE IF EXISTS {DB_NAME};"], - check=True, - ) - - subprocess.run( - ["sudo", "-u", "postgres", "psql", "-d", "postgres", "-c", f"CREATE DATABASE {DB_NAME} OWNER {DB_USER};"], - check=True, - ) - - logger.info("[Restore] База данных пересоздана") - - os.environ["PGPASSWORD"] = DB_PASSWORD - - if is_custom_dump: - result = subprocess.run( - [ - "pg_restore", - f"--dbname={DB_NAME}", - "-U", - DB_USER, - "-h", - PG_HOST, - "-p", - PG_PORT, - "--no-owner", - "--exit-on-error", - tmp_path, - ], - capture_output=True, - text=True, - ) - else: - result = subprocess.run( - [ - "psql", - "-U", - DB_USER, - "-h", - PG_HOST, - "-p", - PG_PORT, - "-d", - DB_NAME, - "-f", - tmp_path, - ], - capture_output=True, - text=True, - ) - - del os.environ["PGPASSWORD"] - - if result.returncode != 0: - logger.error(f"[Restore] Ошибка восстановления: {result.stderr}") + if not success: + logger.error("[Restore] Ошибка: {}", err_msg) await message.answer( - f"❌ Ошибка при восстановлении базы данных:\n
{result.stderr}
", + f"❌ Ошибка при восстановлении базы данных:\n
{err_msg}
", ) return + logger.info("[Restore] База восстановлена") await message.answer( "✅ База данных восстановлена.", reply_markup=build_back_to_db_menu(), ) - logger.info("[Restore] Успешно восстановлено. Завершаем процесс для перезапуска.") + logger.info("[Restore] Завершение для перезапуска") await state.clear() sys.exit(0) diff --git a/handlers/admin/management/file_upload.py b/handlers/admin/management/file_upload.py index a65d73bc..1d97d3c8 100644 --- a/handlers/admin/management/file_upload.py +++ b/handlers/admin/management/file_upload.py @@ -6,6 +6,7 @@ from aiogram.fsm.state import State, StatesGroup from aiogram.types import CallbackQuery, Message from aiogram.utils.keyboard import InlineKeyboardBuilder +from core.executor import run_io from filters.admin import IsAdminFilter from logger import logger @@ -79,7 +80,7 @@ async def handle_admin_file_upload(message: Message, state: FSMContext): else: base_dir = os.path.abspath(".") - os.makedirs(base_dir, exist_ok=True) + await run_io(lambda: os.makedirs(base_dir, exist_ok=True)) dest_path = os.path.join(base_dir, file_name) try: diff --git a/handlers/admin/module/module_handler.py b/handlers/admin/module/module_handler.py index 7abb009a..f0941b9f 100644 --- a/handlers/admin/module/module_handler.py +++ b/handlers/admin/module/module_handler.py @@ -6,6 +6,7 @@ from aiogram.fsm.context import FSMContext from aiogram.types import CallbackQuery from sqlalchemy.ext.asyncio import AsyncSession +from core.executor import run_io from filters.admin import IsAdminFilter from handlers.admin.panel.keyboard import AdminPanelCallback from utils.modules_manager import manager @@ -50,7 +51,7 @@ async def handle_modules(callback_query: CallbackQuery, state: FSMContext, sessi packed = AdminPanelCallback.unpack(callback_query.data) page = max(1, packed.page or 1) - all_items = list_installed_modules() + all_items = await run_io(list_installed_modules) items = [(n, v) for n, v in all_items if n != "web_admin_panel"] per_page = 12 @@ -100,7 +101,7 @@ async def handle_module_restart(callback_query: CallbackQuery, state: FSMContext except Exception as e: result = f"❌ Ошибка перезапуска: {e}" - items = dict(list_installed_modules()) + items = dict(await run_io(list_installed_modules)) ver = items.get(name) title = f"{name} v{ver}" if ver else name text = f"🧩 Модуль: {title}\n\n{result}" @@ -129,7 +130,7 @@ async def handle_module_stop(callback_query: CallbackQuery, state: FSMContext, s except Exception as e: result = f"❌ Ошибка остановки: {e}" - items = dict(list_installed_modules()) + items = dict(await run_io(list_installed_modules)) ver = items.get(name) title = f"{name} v{ver}" if ver else name text = f"🧩 Модуль: {title}\n\n{result}" @@ -158,7 +159,7 @@ async def handle_module_start(callback_query: CallbackQuery, state: FSMContext, except Exception as e: result = f"❌ Ошибка запуска: {e}" - items = dict(list_installed_modules()) + items = dict(await run_io(list_installed_modules)) ver = items.get(name) title = f"{name} v{ver}" if ver else name text = f"🧩 Модуль: {title}\n\n{result}" diff --git a/handlers/admin/panel/panel_handler.py b/handlers/admin/panel/panel_handler.py index 883fc457..d2c7149a 100644 --- a/handlers/admin/panel/panel_handler.py +++ b/handlers/admin/panel/panel_handler.py @@ -9,6 +9,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from database.models import Admin from filters.admin import IsAdminFilter from logger import logger +from core.executor import run_io from utils.versioning import get_version from .keyboard import AdminPanelCallback, build_panel_kb @@ -19,7 +20,7 @@ router = Router() @router.callback_query(AdminPanelCallback.filter(F.action == "admin"), IsAdminFilter()) async def handle_admin_callback_query(callback_query: CallbackQuery, state: FSMContext, session: AsyncSession): - text = f"🤖 Панель администратора\n\nВерсия бота:\n
{get_version()}
" + text = f"🤖 Панель администратора\n\nВерсия бота:\n
{await run_io(get_version)}
" await state.clear() @@ -60,7 +61,7 @@ async def handle_admin_callback_query_simple(callback_query: CallbackQuery, stat @router.message(Command("admin"), IsAdminFilter()) async def handle_admin_message(message: Message, state: FSMContext, session: AsyncSession): - text = f"🤖 Панель администратора\n\nВерсия бота:\n
{get_version()}
" + text = f"🤖 Панель администратора\n\nВерсия бота:\n
{await run_io(get_version)}
" await state.clear() diff --git a/handlers/admin/restart/restart_handler.py b/handlers/admin/restart/restart_handler.py index 88b79175..42b16442 100644 --- a/handlers/admin/restart/restart_handler.py +++ b/handlers/admin/restart/restart_handler.py @@ -6,7 +6,7 @@ import sys import psutil from aiogram import F, Router -from core.executor import get_thread_pool +from core.executor import run_io from aiogram.types import CallbackQuery from filters.admin import IsAdminFilter @@ -33,14 +33,7 @@ async def restart_bot(): is_systemd = parent and "systemd" in parent.name().lower() if is_systemd: - loop = asyncio.get_running_loop() - await loop.run_in_executor( - get_thread_pool(), - lambda: subprocess.run( - ["sudo", "systemctl", "restart", "bot.service"], - check=True, - ), - ) + await run_io(lambda: subprocess.run(["sudo", "systemctl", "restart", "bot.service"], check=True)) else: python_exe = sys.executable script_path = os.path.abspath(sys.argv[0]) diff --git a/handlers/keys/key_connect.py b/handlers/keys/key_connect.py index defa7339..3deeb3d9 100644 --- a/handlers/keys/key_connect.py +++ b/handlers/keys/key_connect.py @@ -48,6 +48,21 @@ from logger import logger router = Router() +def generate_key_qr_file(qr_data: str, email: str) -> str: + """Генерация QR в файл. Вызывать через run_cpu(). Возвращает путь к файлу.""" + qr = qrcode.QRCode(version=1, box_size=10, border=4) + qr.add_data(qr_data) + qr.make(fit=True) + img = qr.make_image(fill_color="black", back_color="white") + buffer = BytesIO() + img.save(buffer, format="PNG") + buffer.seek(0) + qr_path = f"/tmp/qrcode_{email}.png" + with open(qr_path, "wb") as f: + f.write(buffer.read()) + return qr_path + + @router.callback_query(F.data.startswith("connect_device|")) async def handle_connect_device(callback_query: CallbackQuery, session: AsyncSession): try: @@ -226,18 +241,9 @@ async def show_qr_code(callback_query: types.CallbackQuery, session: AsyncSessio await callback_query.message.answer("❌ У этой подписки отсутствует ссылка для подключения.") return - qr = qrcode.QRCode(version=1, box_size=10, border=4) - qr.add_data(qr_data) - qr.make(fit=True) + from core.executor import run_cpu - img = qr.make_image(fill_color="black", back_color="white") - buffer = BytesIO() - img.save(buffer, format="PNG") - buffer.seek(0) - - qr_path = f"/tmp/qrcode_{record.email}.png" - with open(qr_path, "wb") as f: - f.write(buffer.read()) + qr_path = await run_cpu(generate_key_qr_file, qr_data, record.email) builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text=BACK, callback_data=f"view_key|{record.email}")) diff --git a/handlers/notifications/general_notifications.py b/handlers/notifications/general_notifications.py index 8f2aafba..d2524288 100644 --- a/handlers/notifications/general_notifications.py +++ b/handlers/notifications/general_notifications.py @@ -527,7 +527,7 @@ async def notify_expiring_keys( results = await asyncio.gather(*tasks, return_exceptions=True) for r in results: if isinstance(r, Exception): - logger.error("Ошибка в задаче продления: %s", r) + logger.error("Ошибка в задаче продления: {}", r) continue renew_results.append(r) else: diff --git a/handlers/refferal.py b/handlers/refferal.py index 6b625578..cd01eb4f 100644 --- a/handlers/refferal.py +++ b/handlers/refferal.py @@ -49,6 +49,21 @@ from .utils import edit_or_send_message, format_days router = Router() +def generate_referral_qr_file(referral_link: str, chat_id: str) -> str: + """Генерация QR в файл. Вызывать через run_cpu(). Возвращает путь к файлу.""" + qr = qrcode.QRCode(version=1, box_size=10, border=4) + qr.add_data(referral_link) + qr.make(fit=True) + image = qr.make_image(fill_color="black", back_color="white") + buffer = BytesIO() + image.save(buffer, format="PNG") + buffer.seek(0) + qr_path = f"/tmp/qrcode_referral_{chat_id}.png" + with open(qr_path, "wb") as file: + file.write(buffer.read()) + return qr_path + + @router.callback_query(F.data == "invite") @router.message(F.text == "/invite") async def invite_handler(callback_query_or_message: Message | CallbackQuery, session: AsyncSession): @@ -159,21 +174,11 @@ async def inline_referral_handler(inline_query: InlineQuery, session: AsyncSessi @router.callback_query(F.data.startswith("show_referral_qr|")) async def show_referral_qr(callback_query: CallbackQuery): try: + from core.executor import run_cpu + chat_id = callback_query.data.split("|")[1] referral_link = get_referral_link(chat_id) - - qr = qrcode.QRCode(version=1, box_size=10, border=4) - qr.add_data(referral_link) - qr.make(fit=True) - - image = qr.make_image(fill_color="black", back_color="white") - buffer = BytesIO() - image.save(buffer, format="PNG") - buffer.seek(0) - - qr_path = f"/tmp/qrcode_referral_{chat_id}.png" - with open(qr_path, "wb") as file: - file.write(buffer.read()) + qr_path = await run_cpu(generate_referral_qr_file, referral_link, chat_id) builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text=BACK, callback_data="invite")) diff --git a/handlers/tariffs/buy/key_tariffs.py b/handlers/tariffs/buy/key_tariffs.py index dfd670aa..06391421 100644 --- a/handlers/tariffs/buy/key_tariffs.py +++ b/handlers/tariffs/buy/key_tariffs.py @@ -737,7 +737,7 @@ async def select_tariff_plan(callback_query: CallbackQuery, session: Any, state: tg_id = callback_query.from_user.id tariff_id = int(callback_query.data.split("|")[1]) - logger.info("[TARIFF_CFG] select_tariff_plan: tg_id=%s tariff_id=%s", tg_id, tariff_id) + logger.info("[TARIFF_CFG] select_tariff_plan: tg_id={} tariff_id={}", tg_id, tariff_id) tariff = await get_tariff_by_id(session, tariff_id) if not tariff: diff --git a/hooks/hooks.py b/hooks/hooks.py index 9aff042d..ee112df3 100644 --- a/hooks/hooks.py +++ b/hooks/hooks.py @@ -25,12 +25,12 @@ def register_hook(name: str, func: Callable[..., Any] | None = None): def deco(f: Callable[..., Any]): _hooks.setdefault(name, []).append((f, owner(f))) - logger.info(f"[Hook] Зарегистрирован хук '{name}': {f.__name__}") + logger.info("[Hook] {} -> {}", name, f.__name__) return f return deco _hooks.setdefault(name, []).append((func, owner(func))) - logger.info(f"[Hook] Зарегистрирован хук '{name}': {func.__name__}") + logger.info("[Hook] {} -> {}", name, func.__name__) def unregister_module_hooks(module_name: str): @@ -58,23 +58,26 @@ async def run_hooks(name: str, require_enabled: bool = True, **kwargs) -> list[A if inspect.iscoroutinefunction(func): coro = func(**kwargs) else: - - async def _run_sync(): - return func(**kwargs) - - coro = _run_sync() + from core.executor import run_io + coro = run_io(lambda: func(**kwargs)) result = await asyncio.wait_for(coro, timeout=DEFAULT_HOOK_TIMEOUT) if result: results.append(result) except TimeoutError: logger.error( - f"[HOOK:{name}] Таймаут в {getattr(func, '__name__', func)} при timeout={DEFAULT_HOOK_TIMEOUT}", + "[Hook:{}] Таймаут {} с в {}", + name, + DEFAULT_HOOK_TIMEOUT, + getattr(func, "__name__", func), exc_info=True, ) except Exception as e: logger.error( - f"[HOOK:{name}] Ошибка в {getattr(func, '__name__', func)}: {e}", + "[Hook:{}] Ошибка в {}: {}", + name, + getattr(func, "__name__", func), + e, exc_info=True, ) return results diff --git a/servers.py b/servers.py index 5b312ba7..bbc57d23 100644 --- a/servers.py +++ b/servers.py @@ -10,7 +10,7 @@ from ping3 import ping from bot import bot from config import ADMIN_ID, PING_TIME -from core.executor import get_thread_pool +from core.executor import run_io from database import async_session_maker, get_servers from handlers.admin.servers.keyboard import AdminServerCallback from logger import logger @@ -30,11 +30,7 @@ def _sync_ping(server_ip: str, timeout: float = 3): async def ping_server(server_ip: str) -> bool: async with PING_SEMAPHORE: try: - loop = asyncio.get_running_loop() - response = await loop.run_in_executor( - get_thread_pool(), - lambda: _sync_ping(server_ip, 3), - ) + response = await run_io(_sync_ping, server_ip, 3) if response is not None and response is not False: return True return await check_tcp_connection(server_ip, 443) diff --git a/utils/backup.py b/utils/backup.py index 005fb8a3..585f576e 100644 --- a/utils/backup.py +++ b/utils/backup.py @@ -40,34 +40,35 @@ async def backup_database() -> Exception | None: Returns: Optional[Exception]: Исключение в случае ошибки или None при успешном выполнении """ - import asyncio - from core.executor import get_process_pool + from core.executor import run_io - loop = asyncio.get_event_loop() - pool = get_process_pool() + # Создаём бэкап в пуле потоков (run_io), а не процессов (run_cpu), чтобы файл + # создавался в том же процессе, что и отправка — иначе путь может быть недоступен + # (воркер уведомлений и воркер пула процессов могут иметь разный cwd/окружение). if BACKUP_CREATE_ARCHIVE: if not any([BACKUP_INCLUDE_DB, BACKUP_INCLUDE_CONFIG, BACKUP_INCLUDE_TEXTS, BACKUP_INCLUDE_IMG]): - backup_file_path, exception = await loop.run_in_executor(pool, _create_database_backup) + backup_file_path, exception = await run_io(_create_database_backup) else: - backup_file_path, exception = await loop.run_in_executor(pool, _create_backup_archive) + backup_file_path, exception = await run_io(_create_backup_archive) else: - backup_file_path, exception = await loop.run_in_executor(pool, _create_database_backup) + backup_file_path, exception = await run_io(_create_database_backup) if exception: - logger.error(f"Ошибка при создании бэкапа: {exception}") + logger.error("[Backup] Ошибка при создании: {}", exception) return exception + logger.info("[Backup] Файл создан: {}", backup_file_path) try: await _send_backup_to_admins(backup_file_path) - exception = _cleanup_old_backups() + exception = await run_io(_cleanup_old_backups) if exception: - logger.error(f"Ошибка при удалении старых бэкапов: {exception}") + logger.error("[Backup] Ошибка при очистке старых: {}", exception) return exception return None except Exception as e: - logger.error(f"Ошибка при отправке бэкапа: {e}") + logger.error("[Backup] Ошибка при отправке: {}", e) return e @@ -108,13 +109,13 @@ def _create_database_backup() -> tuple[str | None, Exception | None]: capture_output=True, text=True, ) - logger.info(f"Бэкап базы данных создан: {filename}") + logger.info("[Backup] БД создана: {}", filename) return str(filename), None except subprocess.CalledProcessError as e: - logger.error(f"Ошибка при выполнении pg_dump: {e.stderr}") + logger.error("[Backup] pg_dump: {}", e.stderr) return None, e except Exception as e: - logger.error(f"Непредвиденная ошибка при создании бэкапа: {e}") + logger.error("[Backup] Непредвиденная ошибка: {}", e) return None, e finally: if "PGPASSWORD" in os.environ: @@ -143,26 +144,26 @@ def _create_backup_archive() -> tuple[str | None, Exception | None]: if BACKUP_INCLUDE_DB: db_backup_path, db_exception = _create_database_backup() if db_exception: - logger.warning(f"Не удалось создать бекап БД для архива: {db_exception}") + logger.warning("[Backup] БД для архива не создана: {}", db_exception) elif db_backup_path and os.path.exists(db_backup_path): tar.add(db_backup_path, arcname=f"{archive_folder}/database.sql") - logger.info("База данных добавлена в архив") + logger.info("[Backup] БД добавлена в архив") if BACKUP_INCLUDE_CONFIG: config_path = project_root / "config.py" if config_path.exists(): tar.add(config_path, arcname=f"{archive_folder}/config.py") - logger.info("config.py добавлен в архив") + logger.info("[Backup] config.py в архив") else: - logger.warning("config.py не найден, пропущен") + logger.warning("[Backup] config.py не найден") if BACKUP_INCLUDE_TEXTS: texts_path = project_root / "handlers" / "texts.py" if texts_path.exists(): tar.add(texts_path, arcname=f"{archive_folder}/texts.py") - logger.info("texts.py добавлен в архив") + logger.info("[Backup] texts.py в архив") else: - logger.warning("handlers/texts.py не найден, пропущен") + logger.warning("[Backup] handlers/texts.py не найден") if BACKUP_INCLUDE_IMG: img_dir = project_root / "img" @@ -170,23 +171,23 @@ def _create_backup_archive() -> tuple[str | None, Exception | None]: img_files = [f for f in img_dir.iterdir() if f.is_file()] for img_file in img_files: tar.add(img_file, arcname=f"{archive_folder}/img/{img_file.name}") - logger.info(f"Папка img/ добавлена в архив ({len(img_files)} файлов)") + logger.info("[Backup] img/ в архив ({} файлов)", len(img_files)) else: - logger.warning("Папка img/ не найдена, пропущена") + logger.warning("[Backup] img/ не найдена") - logger.info(f"Архив бекапа создан: {archive_path}") + logger.info("[Backup] Архив создан: {}", archive_path) if db_backup_path and os.path.exists(db_backup_path) and db_backup_path != str(archive_path): try: os.unlink(db_backup_path) - logger.info(f"Временный файл БД удален: {db_backup_path}") + logger.info("[Backup] Временный файл БД удалён: {}", db_backup_path) except Exception as e: - logger.warning(f"Не удалось удалить временный файл БД: {e}") + logger.warning("[Backup] Не удалось удалить временный файл БД: {}", e) return str(archive_path), None except Exception as e: - logger.error(f"Непредвиденная ошибка при создании архива бекапа: {e}") + logger.error("[Backup] Ошибка создания архива: {}", e) return None, e @@ -209,19 +210,19 @@ def _cleanup_old_backups() -> Exception | None: file_mtime = datetime.fromtimestamp(backup_file.stat().st_mtime) if file_mtime < cutoff_date: backup_file.unlink() - logger.info(f"Удален старый бэкап: {backup_file}") + logger.info("[Backup] Удалён старый: {}", backup_file) for archive_file in backup_dir.glob("*.tar.gz"): if archive_file.is_file(): file_mtime = datetime.fromtimestamp(archive_file.stat().st_mtime) if file_mtime < cutoff_date: archive_file.unlink() - logger.info(f"Удален старый архив: {archive_file}") + logger.info("[Backup] Удалён старый архив: {}", archive_file) - logger.info("Очистка старых бэкапов завершена") + logger.info("[Backup] Очистка старых завершена") return None except Exception as e: - logger.error(f"Ошибка при удалении старых бэкапов: {e}") + logger.error("[Backup] Ошибка при очистке: {}", e) return e @@ -255,9 +256,9 @@ async def _send_backup_to_admins(backup_file_path: str) -> None: for admin_id in ADMIN_ID: try: await bot.send_document(chat_id=admin_id, document=backup_input_file) - logger.info(f"Бэкап базы данных отправлен админу: {admin_id}") + logger.info("[Backup] Отправлено админу: {}", admin_id) except Exception as e: - logger.error(f"Не удалось отправить бэкап админу {admin_id}: {e}") + logger.error("[Backup] Не отправлено админу {}: {}", admin_id, e) try: async with aiofiles.open(backup_file_path, "rb") as backup_file: @@ -267,12 +268,13 @@ async def _send_backup_to_admins(backup_file_path: str) -> None: if BACKUP_SEND_MODE == "default": await send_default() + logger.info("[Backup] Отправлено всем админам") elif BACKUP_SEND_MODE == "channel": channel_id = BACKUP_CHANNEL_ID.strip() thread_id = BACKUP_CHANNEL_THREAD_ID.strip() if not channel_id: - logger.error("BACKUP_CHANNEL_ID не задан для режима 'channel', fallback на default") + logger.error("[Backup] BACKUP_CHANNEL_ID не задан, fallback на default") await send_default() return send_kwargs = {"chat_id": channel_id, "document": backup_input_file} @@ -282,14 +284,14 @@ async def _send_backup_to_admins(backup_file_path: str) -> None: send_kwargs["caption"] = BACKUP_CAPTION try: await bot.send_document(**send_kwargs) - logger.info(f"Бэкап базы данных отправлен в канал: {channel_id} (топик: {thread_id})") + logger.info("[Backup] Отправлено в канал: {} (топик: {})", channel_id, thread_id) except Exception as e: - logger.error(f"Не удалось отправить бэкап в канал {channel_id}: {e}, fallback на default") + logger.error("[Backup] Не отправлено в канал {}: {}, fallback", channel_id, e) await send_default() elif BACKUP_SEND_MODE == "bot": if not BACKUP_OTHER_BOT_TOKEN: - logger.error("BACKUP_OTHER_BOT_TOKEN не задан для режима 'bot', fallback на default") + logger.error("[Backup] BACKUP_OTHER_BOT_TOKEN не задан, fallback") await send_default() return other_bot = Bot(token=BACKUP_OTHER_BOT_TOKEN) @@ -300,16 +302,16 @@ async def _send_backup_to_admins(backup_file_path: str) -> None: if BACKUP_CAPTION: send_kwargs["caption"] = BACKUP_CAPTION await other_bot.send_document(**send_kwargs) - logger.info(f"Бэкап базы данных отправлен админу через другого бота: {admin_id}") + logger.info("[Backup] Отправлено через другого бота админу: {}", admin_id) except Exception as e: - logger.error(f"Не удалось отправить бэкап админу {admin_id} через другого бота: {e}") + logger.error("[Backup] Не отправлено админу {} через другого бота: {}", admin_id, e) await other_bot.session.close() except Exception as e: - logger.error(f"Ошибка при отправке через другого бота: {e}, fallback на default") + logger.error("[Backup] Ошибка через другого бота: {}, fallback", e) await send_default() else: - logger.error(f"Неизвестный BACKUP_SEND_MODE: {BACKUP_SEND_MODE}, fallback на default") + logger.error("[Backup] Неизвестный BACKUP_SEND_MODE: {}, fallback", BACKUP_SEND_MODE) await send_default() except Exception as e: - logger.error(f"Ошибка при отправке бэкапа: {e}") + logger.error("[Backup] Ошибка при отправке: {}", e) raise diff --git a/utils/modules_manager.py b/utils/modules_manager.py index 8e7b8053..11cc7826 100644 --- a/utils/modules_manager.py +++ b/utils/modules_manager.py @@ -5,6 +5,7 @@ import sys from aiogram import Router +from core.executor import run_io from hooks.hooks import unregister_module_hooks from logger import logger @@ -97,7 +98,7 @@ class ModulesManager: if name in self.disabled: self.disabled.discard(name) - self._save_state() + await run_io(self._save_state) logger.info(f"[Modules] {name} запущен.") @@ -108,7 +109,7 @@ class ModulesManager: logger.info(f"[Modules] {name} уже остановлен или не найден.") if name not in self.disabled: self.disabled.add(name) - self._save_state() + await run_io(self._save_state) return try: @@ -127,7 +128,7 @@ class ModulesManager: if name not in self.disabled: self.disabled.add(name) - self._save_state() + await run_io(self._save_state) logger.info(f"[Modules] {name} остановлен.") diff --git a/utils/versioning.py b/utils/versioning.py index b0cfcca3..001cbcce 100644 --- a/utils/versioning.py +++ b/utils/versioning.py @@ -17,7 +17,7 @@ def _get_git_commit_number_uncached() -> str: if not os.path.isdir(os.path.join(cwd, ".git")): cwd = "/root/Solo_bot" - logger.info(f"[Git] .git не найден в текущем каталоге, используем {cwd}") + logger.info("[Git] .git не найден, используем {}", cwd) env = os.environ.copy() env["GIT_DIR"] = os.path.join(cwd, ".git") @@ -48,7 +48,7 @@ def _get_git_commit_number_uncached() -> str: branch = "dev" except Exception as e: - logger.error(f"[Git] Ошибка при получении локального коммита: {e}") + logger.error("[Git] Ошибка локального коммита: {}", e) return f"\n(Требуется обновление через CLI (команда sudo solobot): {e})" try: @@ -72,7 +72,7 @@ def _get_git_commit_number_uncached() -> str: f'#{remote_number})' ) except Exception as e: - logger.error(f"[Git] Ошибка при получении удалённого коммита: {e}") + logger.error("[Git] Ошибка удалённого коммита: {}", e) return "\n(Требуется обновление через CLI, команда sudo solobot)"