From 125f7eb800494a4603fb5b5fc79472ff50e29de6 Mon Sep 17 00:00:00 2001 From: MpakobecEx Date: Wed, 27 May 2026 21:54:15 +0500 Subject: [PATCH] =?UTF-8?q?=D0=9F=D0=B5=D1=80=D0=B2=D1=8B=D0=B9=20=D0=BA?= =?UTF-8?q?=D0=BE=D0=BC=D0=B8=D1=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env | 6 + .gitignore | 9 + bot.py | 34 +++ config.py | 16 + database.py | 223 ++++++++++++++ email_sender.py | 148 +++++++++ gettoken.py | 11 + gmail_sender.py | 103 +++++++ handlers.py | 755 ++++++++++++++++++++++++++++++++++++++++++++++ keyboards.py | 127 ++++++++ queue_manager.py | 73 +++++ scheduler.py | 57 ++++ timezone_utils.py | 40 +++ 13 files changed, 1602 insertions(+) create mode 100644 .env create mode 100644 .gitignore create mode 100644 bot.py create mode 100644 config.py create mode 100644 database.py create mode 100644 email_sender.py create mode 100644 gettoken.py create mode 100644 gmail_sender.py create mode 100644 handlers.py create mode 100644 keyboards.py create mode 100644 queue_manager.py create mode 100644 scheduler.py create mode 100644 timezone_utils.py diff --git a/.env b/.env new file mode 100644 index 0000000..d222e4f --- /dev/null +++ b/.env @@ -0,0 +1,6 @@ +BOT_TOKEN=<> +ADMIN_ID=<> + +# Gmail настройки +GMAIL_USER=<> +GMAIL_APP_PASSWORD=<> \ No newline at end of file diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..c4e9573 --- /dev/null +++ b/.gitignore @@ -0,0 +1,9 @@ +*.pickle +credentials.json +bot_database.db +__pycache__/*.pyc +.idea/EMailBot.iml +.idea/inspectionProfiles/profiles_settings.xml +.idea/misc.xml +.idea/modules.xml +EmailBot/ diff --git a/bot.py b/bot.py new file mode 100644 index 0000000..4347e6c --- /dev/null +++ b/bot.py @@ -0,0 +1,34 @@ +# bot.py +import asyncio +import logging +from aiogram import Bot, Dispatcher +from config import BOT_TOKEN +from database import init_db +from handlers import router +from scheduler import batch_sender + +logging.basicConfig(level=logging.INFO) +logger = logging.getLogger(__name__) + + +async def main(): + # Инициализируем БД + await init_db() + + # Создаём бота и диспетчер + bot = Bot(token=BOT_TOKEN) + dp = Dispatcher() + + # Подключаем обработчики + dp.include_router(router) + + # Запускаем фоновую задачу для отправки писем + asyncio.create_task(batch_sender(bot)) + logger.info("Бот запущен, фоновый отправитель активен") + + # Запускаем поллинг + await dp.start_polling(bot) + + +if __name__ == "__main__": + asyncio.run(main()) \ No newline at end of file diff --git a/config.py b/config.py new file mode 100644 index 0000000..0abc847 --- /dev/null +++ b/config.py @@ -0,0 +1,16 @@ +# config.py +import os +from dotenv import load_dotenv + +load_dotenv() + +BOT_TOKEN = os.getenv("BOT_TOKEN") + +# Настройки Gmail +SMTP_HOST = "smtp.gmail.com" +SMTP_PORT = 587 +SMTP_USER = os.getenv("GMAIL_USER") +SMTP_PASSWORD = os.getenv("GMAIL_APP_PASSWORD") + +# ID админа по умолчанию (первый запуск) +ADMIN_ID = int(os.getenv("ADMIN_ID", 0)) \ No newline at end of file diff --git a/database.py b/database.py new file mode 100644 index 0000000..63971dd --- /dev/null +++ b/database.py @@ -0,0 +1,223 @@ +# database.py +from typing import Any, Coroutine + +import aiosqlite +from datetime import timedelta +import pytz + +DB_NAME = "bot_database.db" + + +async def init_db(): + async with aiosqlite.connect(DB_NAME) as db: + await db.execute(""" + CREATE TABLE IF NOT EXISTS users ( + user_id INTEGER PRIMARY KEY, + username TEXT, + full_name TEXT, + is_allowed BOOLEAN DEFAULT 0, + last_notification_time TIMESTAMP DEFAULT NULL + ) + """) + + # Добавляем колонку если её нет + try: + await db.execute("ALTER TABLE users ADD COLUMN last_notification_time TIMESTAMP DEFAULT NULL") + except aiosqlite.OperationalError: + pass + + await db.execute(""" + CREATE TABLE IF NOT EXISTS settings ( + key TEXT PRIMARY KEY, + value TEXT + ) + """) + + await db.execute( + "INSERT OR IGNORE INTO settings (key, value) VALUES (?, ?)", + ("target_email", "") + ) + await db.execute( + "INSERT OR IGNORE INTO settings (key, value) VALUES (?, ?)", + ("admin_id", "") + ) + await db.commit() + + +async def add_user(user_id: int, username: str, full_name: str): + async with aiosqlite.connect(DB_NAME) as db: + await db.execute( + "INSERT OR IGNORE INTO users (user_id, username, full_name, is_allowed, last_notification_time) VALUES (?, ?, ?, 0, NULL)", + (user_id, username, full_name) + ) + await db.commit() + + +async def set_user_allowed(user_id: int, allowed: bool): + async with aiosqlite.connect(DB_NAME) as db: + await db.execute( + "UPDATE users SET is_allowed = ? WHERE user_id = ?", + (1 if allowed else 0, user_id) + ) + await db.commit() + + +async def is_user_allowed(user_id: int) -> bool: + async with aiosqlite.connect(DB_NAME) as db: + async with db.execute("SELECT is_allowed FROM users WHERE user_id = ?", (user_id,)) as cursor: + row = await cursor.fetchone() + return row is not None and row[0] == 1 + + +async def get_target_email() -> str: + async with aiosqlite.connect(DB_NAME) as db: + async with db.execute("SELECT value FROM settings WHERE key = 'target_email'") as cursor: + row = await cursor.fetchone() + return row[0] if row else "" + + +async def set_target_email(email: str): + async with aiosqlite.connect(DB_NAME) as db: + await db.execute( + "UPDATE settings SET value = ? WHERE key = 'target_email'", + (email,) + ) + await db.commit() + + +async def get_admin_id() -> int: + async with aiosqlite.connect(DB_NAME) as db: + async with db.execute("SELECT value FROM settings WHERE key = 'admin_id'") as cursor: + row = await cursor.fetchone() + if row and row[0]: + return int(row[0]) + return None + + +async def set_admin_id(admin_id: int): + async with aiosqlite.connect(DB_NAME) as db: + await db.execute( + "UPDATE settings SET value = ? WHERE key = 'admin_id'", + (str(admin_id),) + ) + await db.commit() + + +async def get_all_users(): + async with aiosqlite.connect(DB_NAME) as db: + async with db.execute("SELECT user_id, username, full_name, is_allowed FROM users ORDER BY user_id") as cursor: + return await cursor.fetchall() + + +async def delete_user(user_id: int): + async with aiosqlite.connect(DB_NAME) as db: + await db.execute("DELETE FROM users WHERE user_id = ?", (user_id,)) + await db.commit() + return True + + +async def get_user_by_id(user_id: int): + async with aiosqlite.connect(DB_NAME) as db: + async with db.execute("SELECT user_id, username, full_name, is_allowed FROM users WHERE user_id = ?", + (user_id,)) as cursor: + return await cursor.fetchone() + + +async def can_send_notification(user_id: int) -> bool: + """Проверяет, можно ли отправить уведомление (с учётом часового пояса Екатеринбурга)""" + async with aiosqlite.connect(DB_NAME) as db: + async with db.execute( + "SELECT is_allowed, last_notification_time FROM users WHERE user_id = ?", + (user_id,) + ) as cursor: + row = await cursor.fetchone() + + if not row: + return True + + is_allowed, last_time = row + + # Если пользователь уже имеет доступ, не отправляем уведомления + if is_allowed == 1: + return False + + # Если нет доступа и не было уведомлений + if last_time is None: + return True + + from timezone_utils import parse_time_from_utc, get_current_time + + # Преобразуем время из БД в Екатеринбург + last_notification = parse_time_from_utc(last_time) + + if last_notification is None: + return True + + # Проверяем, прошёл ли час + now = get_current_time() + return (now - last_notification) >= timedelta(hours=1) + + +async def update_notification_time(user_id: int): + """Обновляет время последнего уведомления (храним в UTC)""" + from timezone_utils import get_current_time, parse_time_to_utc + + local_time = get_current_time() + utc_time = parse_time_to_utc(local_time) + + async with aiosqlite.connect(DB_NAME) as db: + await db.execute( + "UPDATE users SET last_notification_time = ? WHERE user_id = ?", + (utc_time.isoformat(), user_id) + ) + await db.commit() + + +async def reset_notification_timer(user_id: int): + """Сбрасывает таймер уведомлений""" + async with aiosqlite.connect(DB_NAME) as db: + await db.execute( + "UPDATE users SET last_notification_time = NULL WHERE user_id = ?", + (user_id,) + ) + await db.commit() + + +async def get_last_notification_time(user_id: int) -> Any | None: + """Получить время последнего уведомления в формате Екатеринбурга""" + async with aiosqlite.connect(DB_NAME) as db: + async with db.execute( + "SELECT last_notification_time FROM users WHERE user_id = ?", + (user_id,) + ) as cursor: + row = await cursor.fetchone() + if row and row[0]: + from timezone_utils import parse_time_from_utc, format_time + local_time = parse_time_from_utc(row[0]) + if local_time: + return format_time(local_time) + return None + + +async def get_next_allowed_time(user_id: int) -> Any | None: + """ + Возвращает время (в Екатеринбурге), когда пользователь сможет отправить следующий запрос + """ + async with aiosqlite.connect(DB_NAME) as db: + async with db.execute( + "SELECT last_notification_time FROM users WHERE user_id = ?", + (user_id,) + ) as cursor: + row = await cursor.fetchone() + + if not row or row[0] is None: + return None + + from timezone_utils import parse_time_from_utc, format_time + from datetime import timedelta + + last_time = parse_time_from_utc(row[0]) + if last_time: + next_time = last_time + timedelta(hours=1) + return format_time(next_time) + return None \ No newline at end of file diff --git a/email_sender.py b/email_sender.py new file mode 100644 index 0000000..fb5eb7f --- /dev/null +++ b/email_sender.py @@ -0,0 +1,148 @@ +# email_sender.py +import os +import pickle +import base64 +from email.mime.multipart import MIMEMultipart +from email.mime.text import MIMEText +from email.mime.image import MIMEImage +from email.mime.application import MIMEApplication +from typing import List +from queue_manager import QueuedFile +from timezone_utils import format_time +import logging +import mimetypes + +logger = logging.getLogger(__name__) + +# Пути к файлам аутентификации +TOKEN_FILE = 'token.pickle' +CREDENTIALS_FILE = 'credentials.json' # нужен для обновления токена +SCOPES = ['https://www.googleapis.com/auth/gmail.send'] + + +def get_gmail_service(): + """ + Загружает сохранённый токен и возвращает сервис Gmail API + Если токен протух и есть credentials.json — обновляет его + """ + creds = None + + # Проверяем наличие token.pickle + if not os.path.exists(TOKEN_FILE): + raise Exception( + f"Файл {TOKEN_FILE} не найден! " + "Загрузи его на сервер или получи заново через get_token.py" + ) + + # Загружаем токен + with open(TOKEN_FILE, 'rb') as token: + creds = pickle.load(token) + + # Если токен протух — обновляем + from google.auth.transport.requests import Request + if creds and creds.expired and creds.refresh_token: + logger.info("Токен протух, обновляем...") + creds.refresh(Request()) + # Сохраняем обновлённый токен + with open(TOKEN_FILE, 'wb') as token: + pickle.dump(creds, token) + logger.info("Токен обновлён и сохранён") + elif creds and creds.expired and not creds.refresh_token: + # Нет refresh_token — нужна повторная авторизация + if os.path.exists(CREDENTIALS_FILE): + logger.warning("Нет refresh_token, выполняем повторную авторизацию...") + from google_auth_oauthlib.flow import InstalledAppFlow + flow = InstalledAppFlow.from_client_secrets_file(CREDENTIALS_FILE, SCOPES) + creds = flow.run_local_server(port=0) + with open(TOKEN_FILE, 'wb') as token: + pickle.dump(creds, token) + logger.info("Новый токен получен и сохранён") + else: + raise Exception( + f"Токен протух, а {CREDENTIALS_FILE} не найден. " + "Загрузи credentials.json для автообновления или получи новый token.pickle" + ) + + from googleapiclient.discovery import build + return build('gmail', 'v1', credentials=creds) + + +async def send_batch_to_email(target_email: str, files: List[QueuedFile]): + """ + Отправляет письмо с вложениями через Gmail API + Работает без SMTP портов (использует HTTPS 443) + + Args: + target_email: Кому отправить + files: Список файлов из очереди + """ + if not target_email: + raise ValueError("Целевая почта не настроена") + + if not files: + logger.warning("Нет файлов для отправки") + return + + logger.info(f"Начинаем отправку {len(files)} файлов через Gmail API") + + # Создаём письмо + msg = MIMEMultipart() + msg['To'] = target_email + msg['From'] = 'me' # 'me' = аккаунт, который авторизован в Gmail API + msg['Subject'] = f"📎 {len(files)} новых файлов" + + # Формируем текст письма со списком файлов + total_size = sum(len(f.file_bytes) for f in files) + total_size_mb = total_size / (1024 * 1024) + + body = f"Получено {len(files)} файлов.\n" + body += f"Общий размер: {total_size_mb:.2f} MB\n\n" + body += "Список вложений:\n" + body += "-" * 30 + "\n" + + for idx, file in enumerate(files, 1): + file_size_kb = len(file.file_bytes) / 1024 + if file_size_kb > 1024: + size_str = f"{file_size_kb / 1024:.2f} MB" + else: + size_str = f"{file_size_kb:.2f} KB" + body += f"{idx}. {file.filename} ({size_str})\n" + + body += "\n---\n" + body += f"Время отправки: {format_time(files[0].timestamp)}\n" + body += f"Отправлено через Gmail API" + + msg.attach(MIMEText(body, 'plain')) + + # Добавляем вложения + for file in files: + try: + mime_type, encoding = mimetypes.guess_type(file.filename) + + if mime_type and mime_type.startswith('image/'): + img_type = mime_type.split('/')[1] + mime_file = MIMEImage(file.file_bytes, _subtype=img_type) + else: + mime_file = MIMEApplication(file.file_bytes, _subtype='octet-stream') + + mime_file.add_header('Content-Disposition', 'attachment', filename=file.filename) + msg.attach(mime_file) + logger.debug(f"Добавлено вложение: {file.filename}") + except Exception as e: + logger.error(f"Ошибка при добавлении вложения {file.filename}: {e}") + + # Кодируем письмо в base64 для Gmail API + try: + raw_message = base64.urlsafe_b64encode(msg.as_bytes()).decode('utf-8') + except Exception as e: + raise Exception(f"Ошибка кодирования письма: {e}") + + body_message = {'raw': raw_message} + + # Отправляем через Gmail API + try: + service = get_gmail_service() + result = service.users().messages().send(userId='me', body=body_message).execute() + logger.info(f"✅ Письмо отправлено через Gmail API, ID: {result['id']}, файлов: {len(files)}") + except Exception as e: + raise Exception(f"Gmail API ошибка: {str(e)}") \ No newline at end of file diff --git a/gettoken.py b/gettoken.py new file mode 100644 index 0000000..e197478 --- /dev/null +++ b/gettoken.py @@ -0,0 +1,11 @@ +# get_token.py +from google_auth_oauthlib.flow import InstalledAppFlow + +SCOPES = ['https://www.googleapis.com/auth/gmail.send'] +flow = InstalledAppFlow.from_client_secrets_file('credentials.json', SCOPES) +creds = flow.run_local_server(port=0) + +import pickle +with open('token.pickle', 'wb') as token: + pickle.dump(creds, token) +print("✅ token.pickle создан") \ No newline at end of file diff --git a/gmail_sender.py b/gmail_sender.py new file mode 100644 index 0000000..644b481 --- /dev/null +++ b/gmail_sender.py @@ -0,0 +1,103 @@ +# gmail_sender.py +import os +import pickle +import base64 +from email.mime.multipart import MIMEMultipart +from email.mime.text import MIMEText +from email.mime.image import MIMEImage +from email.mime.application import MIMEApplication +from typing import List +from queue_manager import QueuedFile +from timezone_utils import format_time +import logging +import mimetypes + +logger = logging.getLogger(__name__) + +# Пути к файлам аутентификации +CREDENTIALS_FILE = 'credentials.json' # если нужен для обновления токена +TOKEN_FILE = 'token.pickle' +SCOPES = ['https://www.googleapis.com/auth/gmail.send'] + + +def get_gmail_service(): + """Загружает сохранённый токен и возвращает сервис Gmail API""" + creds = None + + if not os.path.exists(TOKEN_FILE): + raise Exception(f"Файл {TOKEN_FILE} не найден. Загрузи его на сервер.") + + with open(TOKEN_FILE, 'rb') as token: + creds = pickle.load(token) + + # Если токен протух — обновляем (потребуется credentials.json) + from google.auth.transport.requests import Request + if creds and creds.expired and creds.refresh_token: + creds.refresh(Request()) + # Сохраняем обновлённый токен + with open(TOKEN_FILE, 'wb') as token: + pickle.dump(creds, token) + + from googleapiclient.discovery import build + return build('gmail', 'v1', credentials=creds) + + +async def send_batch_to_email(target_email: str, files: List[QueuedFile]): + """Отправка письма через Gmail API (без SMTP)""" + if not target_email: + raise ValueError("Целевая почта не настроена") + + if not files: + return + + # Создаём письмо + msg = MIMEMultipart() + msg['To'] = target_email + msg['From'] = 'me' # 'me' = авторизованный аккаунт + msg['Subject'] = f"📎 {len(files)} новых файлов" + + # Текст письма + body = f"Получено {len(files)} файлов.\n\n" + body += "Список вложений:\n" + body += "-" * 30 + "\n" + + for idx, file in enumerate(files, 1): + file_size_kb = len(file.file_bytes) / 1024 + if file_size_kb > 1024: + size_str = f"{file_size_kb / 1024:.2f} MB" + else: + size_str = f"{file_size_kb:.2f} KB" + body += f"{idx}. {file.filename} ({size_str})\n" + + body += "\n---\n" + body += f"Время отправки: {format_time(files[0].timestamp)}" + + msg.attach(MIMEText(body, 'plain')) + + # Добавляем вложения + for file in files: + try: + mime_type, encoding = mimetypes.guess_type(file.filename) + + if mime_type and mime_type.startswith('image/'): + img_type = mime_type.split('/')[1] + mime_file = MIMEImage(file.file_bytes, _subtype=img_type) + else: + mime_file = MIMEApplication(file.file_bytes, _subtype='octet-stream') + + mime_file.add_header('Content-Disposition', 'attachment', filename=file.filename) + msg.attach(mime_file) + except Exception as e: + logger.error(f"Ошибка при добавлении вложения {file.filename}: {e}") + + # Кодируем в base64 для API + raw_message = base64.urlsafe_b64encode(msg.as_bytes()).decode() + body_message = {'raw': raw_message} + + # Отправляем через Gmail API + try: + service = get_gmail_service() + result = service.users().messages().send(userId='me', body=body_message).execute() + logger.info(f"✅ Письмо отправлено через Gmail API, ID: {result['id']}") + except Exception as e: + raise Exception(f"Gmail API ошибка: {str(e)}") \ No newline at end of file diff --git a/handlers.py b/handlers.py new file mode 100644 index 0000000..51ebcc6 --- /dev/null +++ b/handlers.py @@ -0,0 +1,755 @@ +# handlers.py - добавляем обработчики передачи прав +import io +from aiogram import Router, F +from aiogram.filters import Command, StateFilter +from aiogram.fsm.context import FSMContext +from aiogram.fsm.state import StatesGroup, State +from aiogram.types import Message, CallbackQuery, Document + +from database import ( + add_user, set_user_allowed, is_user_allowed, + get_target_email, set_target_email, get_all_users, + delete_user, get_user_by_id, get_admin_id, set_admin_id, + can_send_notification, update_notification_time, get_last_notification_time, + reset_notification_timer, get_next_allowed_time # новые импорты +) +from email_sender import send_batch_to_email +from queue_manager import file_queue +from keyboards import ( + admin_decision_keyboard, admin_panel_keyboard, + users_list_keyboard, confirm_delete_keyboard, + confirm_transfer_keyboard, quick_decision_keyboard +) +from config import ADMIN_ID as DEFAULT_ADMIN_ID +import logging + +logger = logging.getLogger(__name__) +router = Router() + + +class ChangeEmailState(StatesGroup): + waiting_for_email = State() + + +async def is_admin(user_id: int) -> bool: + """Проверка, является ли пользователь администратором""" + current_admin_id = await get_admin_id() + if current_admin_id: + return user_id == current_admin_id + return user_id == DEFAULT_ADMIN_ID + + +@router.message(Command("start")) +async def cmd_start(message: Message): + user_id = message.from_user.id + username = message.from_user.username or "без юзернейма" + full_name = message.from_user.full_name + + await add_user(user_id, username, full_name) + + # Проверяем, есть ли в БД админ, если нет - устанавливаем + current_admin = await get_admin_id() + if not current_admin: + await set_admin_id(DEFAULT_ADMIN_ID) + current_admin = DEFAULT_ADMIN_ID + + if await is_admin(user_id): + await message.answer( + "👑 Вы администратор.\n" + "Вот ваша панель управления:", + reply_markup=admin_panel_keyboard() + ) + return + + if await is_user_allowed(user_id): + await message.answer( + "✅ Доступ подтверждён!\n" + "Отправляйте мне любые файлы (фото, документы, архивы и т.д.).\n" + "Они будут накапливаться и отправляться пачкой раз в минуту на почту." + ) + else: + # Проверяем, можно ли отправлять уведомление админу + if await can_send_notification(user_id): + # Обновляем время уведомления + await update_notification_time(user_id) + + # Отправляем уведомление админу + admin_id = await get_admin_id() + await message.bot.send_message( + admin_id, + f"🆕 Запрос доступа от пользователя:\n" + f"ID: {user_id}\n" + f"Имя: {message.from_user.full_name}\n" + f"Username: @{message.from_user.username or 'нет'}", + reply_markup=admin_decision_keyboard(user_id) + ) + await message.answer( + "⏳ Ваш запрос отправлен администратору.\n" + "Ожидайте подтверждения. Новый запрос можно будет отправить через час." + ) + else: + # Получаем время следующего доступного запроса + next_allowed = await get_next_allowed_time(user_id) + + if next_allowed: + await message.answer( + f"⏳ Вы уже отправляли запрос администратору.\n\n" + f"🕐 Следующий запрос можно отправить после: **{next_allowed}** (по Екатеринбургу)\n\n" + f"Если доступ нужен срочно, свяжитесь с администратором напрямую.", + parse_mode="Markdown" + ) + else: + await message.answer( + "⏳ Ваш запрос уже был отправлен администратору.\n\n" + "Пожалуйста, ожидайте. Новый запрос можно отправить через час.\n\n" + "Если доступ нужен срочно, свяжитесь с администратором напрямую." + ) + + +@router.callback_query(F.data.startswith("approve_")) +async def approve_user(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + user_id = int(callback.data.split("_")[1]) + await set_user_allowed(user_id, True) + + # Сбрасываем таймер уведомлений для этого пользователя + await reset_notification_timer(user_id) + + await callback.bot.send_message( + user_id, + "🎉 Ваш доступ подтверждён! Теперь вы можете отправлять любые файлы.\n" + "Они будут отправляться пачкой раз в минуту." + ) + await callback.answer("Пользователь добавлен в список разрешённых") + await callback.message.edit_reply_markup() + + +@router.callback_query(F.data.startswith("deny_")) +async def deny_user(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + user_id = int(callback.data.split("_")[1]) + await set_user_allowed(user_id, False) + + # Сбрасываем таймер уведомлений для этого пользователя + await reset_notification_timer(user_id) + + # Получаем текущее время Екатеринбурга для сообщения + from timezone_utils import format_time, get_current_time + current_time = format_time(get_current_time()) + + await callback.bot.send_message( + user_id, + f"❌ Ваш запрос доступа отклонён администратором.\n\n" + f"Вы можете отправить новый запрос командой /start в любое время.\n" + f"🕐 Текущее время: {current_time} (Екатеринбург)" + ) + await callback.answer("Доступ отклонён") + await callback.message.edit_reply_markup() + + +@router.callback_query(F.data == "change_email") +async def change_email_request(callback: CallbackQuery, state: FSMContext): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + await callback.message.answer("📧 Введите новый email для получения файлов:") + await state.set_state(ChangeEmailState.waiting_for_email) + await callback.answer() + + +@router.message(ChangeEmailState.waiting_for_email) +async def set_new_email(message: Message, state: FSMContext): + if not await is_admin(message.from_user.id): + await message.answer("Нет прав") + return + + new_email = message.text.strip() + if "@" not in new_email or "." not in new_email: + await message.answer("❌ Некорректный email. Попробуйте ещё раз:") + return + await set_target_email(new_email) + await message.answer(f"✅ Почта для отправки изменена на: {new_email}") + await state.clear() + + +@router.callback_query(F.data == "list_users") +async def list_users(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + users = await get_all_users() + if not users: + await callback.message.answer("📋 Список пользователей пуст") + await callback.answer() + return + + # Разделяем пользователей по статусу для наглядности + pending = [u for u in users if u[3] == 0] + allowed = [u for u in users if u[3] == 1] + + text = "👥 **Управление пользователями**\n\n" + + if pending: + text += f"⏳ **Ожидают доступа ({len(pending)}):**\n" + text += "Нажмите на пользователя для быстрых действий\n\n" + else: + text += "⏳ Нет пользователей в ожидании\n\n" + + if allowed: + text += f"✅ **Подтверждённые ({len(allowed)}):**\n" + text += "Нажмите на пользователя для удаления\n\n" + + await callback.message.answer( + text, + reply_markup=users_list_keyboard(users, 0, for_transfer=False), + parse_mode="Markdown" + ) + await callback.answer() + + +@router.callback_query(F.data.startswith("users_page_")) +async def users_page(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + parts = callback.data.split("_") + page = int(parts[2]) + for_transfer = parts[3] == "True" if len(parts) > 3 else False + + users = await get_all_users() + + if not users: + await callback.message.edit_text("📋 Список пользователей пуст") + await callback.answer() + return + + await callback.message.edit_reply_markup(reply_markup=users_list_keyboard(users, page, for_transfer)) + await callback.answer() + + +@router.callback_query(F.data.startswith("user_")) +async def user_selected(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + user_id = int(callback.data.split("_")[1]) + user = await get_user_by_id(user_id) + + if not user: + await callback.answer("Пользователь не найден") + await callback.message.edit_text("❌ Пользователь уже был удалён") + return + + user_id, username, full_name, is_allowed = user + + text = f"🗑️ **Удаление пользователя**\n\n" + text += f"ID: `{user_id}`\n" + text += f"Имя: {full_name}\n" + text += f"Username: @{username if username and username != 'None' else 'нет'}\n" + text += f"Статус: {'✅ Доступ разрешён' if is_allowed else '⏳ Ожидает'}\n\n" + text += "Вы уверены, что хотите удалить этого пользователя?" + + await callback.message.edit_text( + text, + reply_markup=confirm_delete_keyboard(user_id, username), + parse_mode="Markdown" + ) + await callback.answer() + + +@router.callback_query(F.data.startswith("confirm_delete_")) +async def confirm_delete(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + user_id = int(callback.data.split("_")[2]) + user = await get_user_by_id(user_id) + + if not user: + await callback.message.edit_text("❌ Пользователь уже был удалён") + await callback.answer() + return + + user_id, username, full_name, is_allowed = user + + await delete_user(user_id) + + try: + await callback.bot.send_message( + user_id, + "❌ Ваш доступ был отозван администратором. Вы больше не можете отправлять файлы.\n" + "Для повторного доступа нужно снова запросить разрешение через /start" + ) + except Exception: + pass + + await callback.message.edit_text( + f"✅ Пользователь {full_name} (ID: {user_id}) удалён из базы данных.\n\n" + f"Если нужно восстановить доступ, пользователь должен снова написать /start" + ) + await callback.answer("Пользователь удалён") + + +@router.callback_query(F.data == "cancel_delete") +async def cancel_delete(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + users = await get_all_users() + await callback.message.edit_text( + "👥 Управление пользователями\n\nНажмите на пользователя, чтобы удалить его.", + reply_markup=users_list_keyboard(users, 0, for_transfer=False) + ) + await callback.answer() + + +@router.callback_query(F.data == "close_users") +async def close_users(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + await callback.message.delete() + await callback.answer() + + +@router.callback_query(F.data == "queue_status") +async def queue_status(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + size = await file_queue.get_queue_size() + await callback.message.answer(f"📊 В очереди {size} файлов") + await callback.answer() + + +# ========== НОВЫЕ ОБРАБОТЧИКИ ДЛЯ ПЕРЕДАЧИ ПРАВ АДМИНА ========== + +@router.callback_query(F.data == "transfer_admin") +async def transfer_admin_request(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + users = await get_all_users() + allowed_users = [u for u in users if u[3] == 1] # Только пользователи с доступом + + if not allowed_users: + await callback.message.answer("❌ Нет пользователей с подтверждённым доступом для передачи прав") + await callback.answer() + return + + text = "👑 **Передача прав администратора**\n\n" + text += "Выберите пользователя, которому хотите передать права:\n" + text += "(только пользователи с подтверждённым доступом ✅)\n\n" + text += "⚠️ **Внимание!** После передачи вы потеряете права администратора!" + + await callback.message.answer( + text, + reply_markup=users_list_keyboard(allowed_users, 0, for_transfer=True), + parse_mode="Markdown" + ) + await callback.answer() + + +@router.callback_query(F.data.startswith("select_new_admin_")) +async def select_new_admin(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + new_admin_id = int(callback.data.split("_")[3]) + user = await get_user_by_id(new_admin_id) + + if not user: + await callback.answer("Пользователь не найден") + return + + user_id, username, full_name, is_allowed = user + + if not is_allowed: + await callback.answer("У этого пользователя нет доступа", show_alert=True) + return + + text = f"⚠️ **Подтверждение передачи прав**\n\n" + text += f"Вы собираетесь передать права администратора пользователю:\n" + text += f"👤 {full_name}\n" + text += f"📱 @{username if username and username != 'None' else 'нет'}\n" + text += f"🆔 ID: `{user_id}`\n\n" + text += f"После этого действия **ВЫ ПОТЕРЯЕТЕ ПРАВА АДМИНИСТРАТОРА**.\n\n" + text += f"Вы уверены?" + + await callback.message.edit_text( + text, + reply_markup=confirm_transfer_keyboard(user_id, username, full_name), + parse_mode="Markdown" + ) + await callback.answer() + + +@router.callback_query(F.data.startswith("confirm_transfer_")) +async def confirm_transfer(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + new_admin_id = int(callback.data.split("_")[2]) + old_admin_id = callback.from_user.id + + user = await get_user_by_id(new_admin_id) + + if not user: + await callback.message.edit_text("❌ Пользователь не найден") + await callback.answer() + return + + user_id, username, full_name, is_allowed = user + + if not is_allowed: + await callback.message.edit_text("❌ У этого пользователя нет доступа") + await callback.answer() + return + + # Передаём права админа + await set_admin_id(new_admin_id) + + # Уведомляем нового админа + try: + await callback.bot.send_message( + new_admin_id, + f"👑 **Поздравляем!**\n\n" + f"Пользователь {callback.from_user.full_name} передал вам права администратора.\n\n" + f"Теперь вы можете управлять ботом: изменять почту, одобрять пользователей и т.д.\n\n" + f"Используйте /start для доступа к панели управления.", + parse_mode="Markdown" + ) + except Exception: + pass + + # Уведомляем старого админа + await callback.message.edit_text( + f"✅ **Права администратора переданы!**\n\n" + f"Новый администратор: {full_name} (ID: {new_admin_id})\n\n" + f"Вы больше не являетесь администратором. Для управления ботом обратитесь к новому админу." + ) + + # Отправляем старому админу сообщение в личку (если не в чате с ботом) + try: + await callback.bot.send_message( + old_admin_id, + f"👑 Вы передали права администратора пользователю {full_name}.\n" + f"Вы больше не являетесь администратором бота." + ) + except Exception: + pass + + await callback.answer("Права переданы") + + +@router.callback_query(F.data == "cancel_transfer") +async def cancel_transfer(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + await callback.message.edit_text("❌ Передача прав отменена") + await callback.answer() + + +# Обработчики файлов +@router.message(F.document) +async def handle_document(message: Message): + user_id = message.from_user.id + + if not await is_admin(user_id) and not await is_user_allowed(user_id): + await message.answer("⛔ У вас нет доступа. Запросите его через /start") + return + + document = message.document + file = await message.bot.get_file(document.file_id) + + from io import BytesIO + file_content = await message.bot.download_file(file.file_path) + + if isinstance(file_content, BytesIO): + file_bytes = file_content.getvalue() + elif isinstance(file_content, bytes): + file_bytes = file_content + else: + file_bytes = file_content.read() + + filename = document.file_name + user_info = f"{message.from_user.full_name}" + + await file_queue.add_file(file_bytes, filename, user_info, user_id, file_type='document') + + queue_size = await file_queue.get_queue_size() + await message.answer( + f"✅ Файл '{filename}' добавлен в очередь (всего: {queue_size})\n" + f"📧 Отправка произойдёт в течение минуты" + ) + + +@router.message(F.photo) +async def handle_photo(message: Message): + user_id = message.from_user.id + + if not await is_admin(user_id) and not await is_user_allowed(user_id): + await message.answer("⛔ У вас нет доступа. Запросите его через /start") + return + + photo = message.photo[-1] + file = await message.bot.get_file(photo.file_id) + + from io import BytesIO + file_content = await message.bot.download_file(file.file_path) + + if isinstance(file_content, BytesIO): + file_bytes = file_content.getvalue() + elif isinstance(file_content, bytes): + file_bytes = file_content + else: + file_bytes = file_content.read() + + filename = f"photo_{user_id}_{photo.file_unique_id}.jpg" + user_info = f"{message.from_user.full_name}" + + await file_queue.add_file(file_bytes, filename, user_info, user_id, file_type='photo') + + queue_size = await file_queue.get_queue_size() + await message.answer( + f"✅ Фото добавлено в очередь (всего: {queue_size})\n" + f"📧 Отправка произойдёт в течение минуты" + ) + +@router.message(Command("send_now")) +async def force_send(message: Message): + """Принудительная отправка всех накопленных файлов""" + if not await is_admin(message.from_user.id): + await message.answer("⛔ Нет прав. Только администратор может использовать эту команду.") + return + + # Получаем все файлы из очереди + files = await file_queue.get_all_files() + + if not files: + await message.answer("📭 Очередь пуста. Нечего отправлять.") + return + + target_email = await get_target_email() + if not target_email: + await message.answer("⚠️ Целевая почта не настроена. Используйте кнопку 'Сменить почту'.") + return + + # Отправляем уведомление о начале отправки + status_msg = await message.answer(f"📤 Начинаю отправку {len(files)} файлов...") + + try: + # Отправляем пачку файлов + await send_batch_to_email(target_email, files) + + # Успешно + await status_msg.edit_text( + f"✅ Принудительно отправлено {len(files)} файлов на почту {target_email}\n" + f"📊 Размер: {sum(len(f.file_bytes) for f in files) / (1024 * 1024):.2f} MB" + ) + + # Логируем действие + logger.info(f"Админ {message.from_user.id} принудительно отправил {len(files)} файлов") + + except Exception as e: + # Ошибка при отправке + await status_msg.edit_text( + f"❌ Ошибка при отправке: {str(e)[:200]}\n\n" + f"Файлы остались в очереди и будут отправлены при следующем автоматическом запуске." + ) + # Возвращаем файлы обратно в очередь + for file in files: + await file_queue.add_file( + file.file_bytes, + file.filename, + file.user_info, + file.user_id, + file.file_type + ) + logger.error(f"Ошибка при принудительной отправке: {e}") + + +@router.message(Command("queue")) +async def check_queue(message: Message): + """Показать содержимое очереди""" + if not await is_admin(message.from_user.id): + await message.answer("⛔ Нет прав") + return + + size = await file_queue.get_queue_size() + + if size == 0: + await message.answer("📭 Очередь пуста") + return + + # Получаем первые 10 файлов для просмотра + peek_files = await file_queue.peek_queue(10) + + text = f"📊 **Очередь отправки**\n\n" + text += f"📎 Всего файлов: {size}\n\n" + + if peek_files: + text += "**Первые файлы в очереди:**\n" + for i, f in enumerate(peek_files, 1): + text += f"{i}. {f['filename']}\n" + text += f" └ от {f['user_info']} ({f['size_kb']:.1f} KB) в {f['timestamp']}\n" + + if size > 10: + text += f"\n... и ещё {size - 10} файлов\n" + + text += f"\n⏱️ Автоотправка: каждую минуту\n" + text += f"🚀 Принудительно: /send_now" + + await message.answer(text, parse_mode="Markdown") + + @router.message(Command("stats")) + async def bot_stats(message: Message): + """Показать статистику бота""" + if not await is_admin(message.from_user.id): + await message.answer("⛔ Нет прав") + return + + users = await get_all_users() + allowed_users = [u for u in users if u[3] == 1] + pending_users = [u for u in users if u[3] == 0] + + target_email = await get_target_email() + queue_size = await file_queue.get_queue_size() + + text = f"📊 **Статистика бота**\n\n" + text += f"👥 Всего пользователей: {len(users)}\n" + text += f"✅ С доступом: {len(allowed_users)}\n" + text += f"⏳ Ожидают: {len(pending_users)}\n" + text += f"📎 В очереди: {queue_size}\n" + text += f"📧 Почта для отправки: {target_email or 'не настроена'}\n" + + await message.answer(text, parse_mode="Markdown") + + +@router.callback_query(F.data.startswith("quick_approve_")) +async def quick_approve(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + user_id = int(callback.data.split("_")[2]) + user = await get_user_by_id(user_id) + + if not user: + await callback.answer("Пользователь не найден") + await callback.message.edit_text("❌ Пользователь уже был удалён") + return + + await set_user_allowed(user_id, True) + + # Сбрасываем таймер уведомлений для этого пользователя + await reset_notification_timer(user_id) + + await callback.bot.send_message( + user_id, + "🎉 Ваш доступ подтверждён! Теперь вы можете отправлять любые файлы.\n" + "Они будут отправляться пачкой раз в минуту." + ) + + # Обновляем сообщение + users = await get_all_users() + await callback.message.edit_text( + "👥 Управление пользователями\n\n✅ Доступ подтверждён", + reply_markup=users_list_keyboard(users, 0, for_transfer=False) + ) + await callback.answer("Доступ подтверждён") + + +@router.callback_query(F.data.startswith("quick_deny_")) +async def quick_deny(callback: CallbackQuery): + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + user_id = int(callback.data.split("_")[2]) + user = await get_user_by_id(user_id) + + if not user: + await callback.answer("Пользователь не найден") + await callback.message.edit_text("❌ Пользователь уже был удалён") + return + + await set_user_allowed(user_id, False) + + # Сбрасываем таймер уведомлений для этого пользователя + await reset_notification_timer(user_id) + + # Получаем текущее время Екатеринбурга для сообщения + from timezone_utils import format_time, get_current_time + current_time = format_time(get_current_time()) + + await callback.bot.send_message( + user_id, + f"❌ Ваш запрос доступа отклонён администратором.\n\n" + f"Вы можете отправить новый запрос командой /start в любое время.\n" + f"🕐 Текущее время: {current_time} (Екатеринбург)" + ) + + # Обновляем сообщение админу + users = await get_all_users() + await callback.message.edit_text( + "👥 Управление пользователями\n\n❌ Доступ отклонён", + reply_markup=users_list_keyboard(users, 0, for_transfer=False) + ) + await callback.answer("Доступ отклонён") + + +@router.callback_query(F.data.startswith("pending_user_")) +async def pending_user_info(callback: CallbackQuery): + """Показывает информацию об ожидающем пользователе""" + if not await is_admin(callback.from_user.id): + await callback.answer("Нет прав", show_alert=True) + return + + user_id = int(callback.data.split("_")[2]) + user = await get_user_by_id(user_id) + + if not user: + await callback.answer("Пользователь не найден") + return + + user_id, username, full_name, is_allowed = user + + from database import get_last_notification_time + last_notification = await get_last_notification_time(user_id) + + text = f"👤 **Информация о пользователе**\n\n" + text += f"ID: `{user_id}`\n" + text += f"Имя: {full_name}\n" + text += f"Username: @{username if username and username != 'None' else 'нет'}\n" + text += f"Статус: ⏳ Ожидает подтверждения\n" + if last_notification: + text += f"Последний запрос: {last_notification}\n\n" + + text += "Выберите действие:" + + await callback.message.edit_text( + text, + reply_markup=quick_decision_keyboard(user_id), + parse_mode="Markdown" + ) + await callback.answer() diff --git a/keyboards.py b/keyboards.py new file mode 100644 index 0000000..c1cc117 --- /dev/null +++ b/keyboards.py @@ -0,0 +1,127 @@ +# keyboards.py +from aiogram.types import InlineKeyboardMarkup, InlineKeyboardButton + + +def admin_decision_keyboard(user_id: int): + return InlineKeyboardMarkup(inline_keyboard=[ + [ + InlineKeyboardButton(text="✅ Да", callback_data=f"approve_{user_id}"), + InlineKeyboardButton(text="❌ Нет", callback_data=f"deny_{user_id}") + ] + ]) + + +def admin_panel_keyboard(): + return InlineKeyboardMarkup(inline_keyboard=[ + [InlineKeyboardButton(text="✏️ Сменить почту для отправки", callback_data="change_email")], + [InlineKeyboardButton(text="📋 Список пользователей", callback_data="list_users")], + [InlineKeyboardButton(text="📊 Статус очереди", callback_data="queue_status")], + [InlineKeyboardButton(text="👑 Передать права админа", callback_data="transfer_admin")] + ]) + + +def users_list_keyboard(users: list, page: int = 0, for_transfer: bool = False): + """ + Создаёт клавиатуру со списком пользователей + Для ожидающих пользователей (is_allowed=0) показывает кнопки "Дать доступ/Отклонить" + Для подтверждённых (is_allowed=1) показывает кнопку удаления + """ + from aiogram.types import InlineKeyboardButton, InlineKeyboardMarkup + + ITEMS_PER_PAGE = 3 # Уменьшил до 3, чтобы влезали кнопки действий + start_idx = page * ITEMS_PER_PAGE + end_idx = start_idx + ITEMS_PER_PAGE + page_users = users[start_idx:end_idx] + + keyboard = [] + + for user in page_users: + user_id, username, full_name, is_allowed = user + status = "✅" if is_allowed else "⏳" + name = full_name[:20] if full_name else f"ID:{user_id}" + if username and username != "None": + name = f"@{username}"[:15] + + if for_transfer: + # Для передачи прав - только выбор пользователя + keyboard.append([ + InlineKeyboardButton( + text=f"{status} {name}", + callback_data=f"select_new_admin_{user_id}" + ) + ]) + elif is_allowed == 0: + # Пользователь ожидает доступа - показываем кнопки "Дать/Отклонить" + keyboard.append([ + InlineKeyboardButton( + text=f"⏳ {name} (ожидает)", + callback_data=f"pending_user_{user_id}" + ) + ]) + keyboard.append([ + InlineKeyboardButton(text="✅ Дать доступ", callback_data=f"quick_approve_{user_id}"), + InlineKeyboardButton(text="❌ Отклонить", callback_data=f"quick_deny_{user_id}") + ]) + keyboard.append([InlineKeyboardButton(text="─" * 20, callback_data="noop")]) + else: + # Подтверждённый пользователь - кнопка удаления + keyboard.append([ + InlineKeyboardButton( + text=f"✅ {name}", + callback_data=f"user_{user_id}" + ) + ]) + + # Кнопки навигации + nav_buttons = [] + if page > 0: + nav_buttons.append(InlineKeyboardButton(text="◀️ Назад", callback_data=f"users_page_{page - 1}_{for_transfer}")) + if end_idx < len(users): + nav_buttons.append( + InlineKeyboardButton(text="Вперёд ▶️", callback_data=f"users_page_{page + 1}_{for_transfer}")) + + if nav_buttons: + keyboard.append(nav_buttons) + + # Кнопка отмены/закрытия + if for_transfer: + keyboard.append([ + InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_transfer") + ]) + else: + keyboard.append([ + InlineKeyboardButton(text="🔄 Обновить", callback_data="list_users"), + InlineKeyboardButton(text="❌ Закрыть", callback_data="close_users") + ]) + + return InlineKeyboardMarkup(inline_keyboard=keyboard) + + +def confirm_delete_keyboard(user_id: int, username: str): + """Клавиатура подтверждения удаления""" + return InlineKeyboardMarkup(inline_keyboard=[ + [ + InlineKeyboardButton(text="✅ Да, удалить", callback_data=f"confirm_delete_{user_id}"), + InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_delete") + ] + ]) + + +def confirm_transfer_keyboard(user_id: int, username: str, full_name: str): + """Клавиатура подтверждения передачи прав админа""" + return InlineKeyboardMarkup(inline_keyboard=[ + [ + InlineKeyboardButton(text="✅ Да, передать права", callback_data=f"confirm_transfer_{user_id}"), + InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_transfer") + ] + ]) + + +def quick_decision_keyboard(user_id: int): + """Быстрая клавиатура для принятия решения по ожидающему пользователю""" + return InlineKeyboardMarkup(inline_keyboard=[ + [ + InlineKeyboardButton(text="✅ Дать доступ", callback_data=f"quick_approve_{user_id}"), + InlineKeyboardButton(text="❌ Отклонить", callback_data=f"quick_deny_{user_id}") + ] + ]) \ No newline at end of file diff --git a/queue_manager.py b/queue_manager.py new file mode 100644 index 0000000..5857fd0 --- /dev/null +++ b/queue_manager.py @@ -0,0 +1,73 @@ +# queue_manager.py - обновляем timestamp +import asyncio +from datetime import datetime +from typing import List, Dict +from dataclasses import dataclass +import logging +from timezone_utils import get_current_time + +logging.basicConfig(level=logging.INFO) +logger = logging.getLogger(__name__) + + +@dataclass +class QueuedFile: + """Класс для хранения любого файла в очереди""" + file_bytes: bytes + filename: str + user_info: str + user_id: int + file_type: str + timestamp: datetime # Теперь будет в часовом поясе Екатеринбурга + + +class FileQueue: + """Управление очередью файлов""" + + def __init__(self): + self.queue: List[QueuedFile] = [] + self.lock = asyncio.Lock() + + async def add_file(self, file_bytes: bytes, filename: str, user_info: str, user_id: int, file_type: str): + """Добавить файл в очередь с временем Екатеринбурга""" + async with self.lock: + self.queue.append(QueuedFile( + file_bytes=file_bytes, + filename=filename, + user_info=user_info, + user_id=user_id, + file_type=file_type, + timestamp=get_current_time() # Время Екатеринбурга + )) + logger.info(f"Файл '{filename}' добавлен в очередь. Всего в очереди: {len(self.queue)}") + + async def get_all_files(self) -> List[QueuedFile]: + """Получить все файлы и очистить очередь""" + async with self.lock: + files = self.queue.copy() + self.queue.clear() + return files + + async def get_queue_size(self) -> int: + """Получить размер очереди""" + async with self.lock: + return len(self.queue) + + async def peek_queue(self, limit: int = 10) -> List[dict]: + """Посмотреть первые N файлов в очереди (не забирая их)""" + async with self.lock: + from timezone_utils import format_time + result = [] + for file in self.queue[:limit]: + result.append({ + 'filename': file.filename, + 'user_info': file.user_info, + 'size_kb': len(file.file_bytes) / 1024, + 'timestamp': format_time(file.timestamp) + }) + return result + + +# Глобальный экземпляр очереди +file_queue = FileQueue() +image_queue = file_queue # Для обратной совместимости \ No newline at end of file diff --git a/scheduler.py b/scheduler.py new file mode 100644 index 0000000..43dbbc1 --- /dev/null +++ b/scheduler.py @@ -0,0 +1,57 @@ +# scheduler.py +import asyncio +import logging +from database import get_target_email +from email_sender import send_batch_to_email +from queue_manager import file_queue +from timezone_utils import get_current_time, format_time + +logging.basicConfig(level=logging.INFO) +logger = logging.getLogger(__name__) + + +async def batch_sender(bot): + """ + Фоновая задача: каждую минуту отправляет накопленные файлы + """ + while True: + try: + await asyncio.sleep(60) + + queue_size = await file_queue.get_queue_size() + + if queue_size == 0: + logger.info("Очередь пуста, пропускаем отправку") + continue + + logger.info(f"Начинаем отправку {queue_size} файлов...") + + target_email = await get_target_email() + + if not target_email: + logger.warning("Целевая почта не настроена! Пропускаем отправку.") + continue + + files = await file_queue.get_all_files() + + if files: + # Отправляем пачку + await send_batch_to_email(target_email, files) + + # Используем время Екатеринбурга для логов + current_time = format_time() + logger.info(f"✅ Отправлено {len(files)} файлов в {current_time}") + + try: + from config import ADMIN_ID + await bot.send_message( + ADMIN_ID, + f"📧 Отправлен отчёт с {len(files)} файлами\n" + f"Время: {current_time}" + ) + except Exception as e: + logger.error(f"Не удалось уведомить админа: {e}") + + except Exception as e: + logger.error(f"Ошибка в batch_sender: {e}") + await asyncio.sleep(10) \ No newline at end of file diff --git a/timezone_utils.py b/timezone_utils.py new file mode 100644 index 0000000..772320f --- /dev/null +++ b/timezone_utils.py @@ -0,0 +1,40 @@ +# timezone_utils.py +from datetime import datetime, timedelta +import pytz + +# Часовой пояс Екатеринбурга (UTC+5) +TARGET_TIMEZONE = "Asia/Yekaterinburg" + +def get_current_time(): + """Возвращает текущее время в Екатеринбурге""" + tz = pytz.timezone(TARGET_TIMEZONE) + return datetime.now(tz) + +def format_time(dt: datetime = None, format_str: str = "%Y-%m-%d %H:%M:%S"): + """Форматирует время для вывода""" + if dt is None: + dt = get_current_time() + return dt.strftime(format_str) + +def parse_time_to_utc(dt: datetime): + """Преобразует локальное время (Екатеринбург) в UTC для хранения в БД""" + if dt.tzinfo is None: + # Если время без часового пояса, добавляем Екатеринбург + tz = pytz.timezone(TARGET_TIMEZONE) + dt = tz.localize(dt) + return dt.astimezone(pytz.UTC) + +def parse_time_from_utc(utc_time_str: str): + """Преобразует UTC строку из БД в время Екатеринбурга""" + if not utc_time_str: + return None + try: + # Парсим UTC время + utc_time = datetime.fromisoformat(utc_time_str) + if utc_time.tzinfo is None: + utc_time = pytz.UTC.localize(utc_time) + # Конвертируем в Екатеринбург + tz = pytz.timezone(TARGET_TIMEZONE) + return utc_time.astimezone(tz) + except Exception: + return None