Первый комит
This commit is contained in:
commit
125f7eb800
|
|
@ -0,0 +1,6 @@
|
||||||
|
BOT_TOKEN=<>
|
||||||
|
ADMIN_ID=<>
|
||||||
|
|
||||||
|
# Gmail настройки
|
||||||
|
GMAIL_USER=<>
|
||||||
|
GMAIL_APP_PASSWORD=<>
|
||||||
|
|
@ -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/
|
||||||
|
|
@ -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())
|
||||||
|
|
@ -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))
|
||||||
|
|
@ -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
|
||||||
|
|
@ -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)}")
|
||||||
|
|
@ -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 создан")
|
||||||
|
|
@ -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)}")
|
||||||
|
|
@ -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()
|
||||||
|
|
@ -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}")
|
||||||
|
]
|
||||||
|
])
|
||||||
|
|
@ -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 # Для обратной совместимости
|
||||||
|
|
@ -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)
|
||||||
|
|
@ -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
|
||||||
Loading…
Reference in New Issue