Files
downloads-all-bot/bot/main.py
2026-02-18 14:05:36 +07:00

742 lines
27 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# main.py
import os
import asyncio
import threading
import logging
import subprocess
import time
from datetime import datetime
from typing import Optional
from aiogram import Bot, Dispatcher, F
from aiogram.types import Message, CallbackQuery, FSInputFile, InputMediaDocument
from aiogram.client.session.aiohttp import AiohttpSession
from aiogram.client.telegram import TelegramAPIServer
from config import (
BOT_TOKEN,
LOCAL_API_URL,
CACHE_DIR,
TMP_DIR,
COOKIES_FILE,
RATE_LIMIT_SECONDS,
CACHE_MAX_AGE_DAYS,
CACHE_MAX_SIZE_MB,
)
from keyboards import youtube_quality_keyboard, cancel_keyboard, playlist_keyboard
from downloader import (
download_tiktok_video_and_audio,
download_instagram_video,
download_video,
download_audio,
download_original_quality,
download_playlist_videos,
DownloadCancelled,
)
from middleware import PrivateMiddleware
from rate_limit import check_rate_limit
from info import extract_info, is_playlist, get_platform_info
from cache import cache_key, cache_path
from cleanup import cleanup_tmp
# Настройка логирования
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
# Создаем директории
os.makedirs(CACHE_DIR, exist_ok=True)
os.makedirs(TMP_DIR, exist_ok=True)
# Инициализация бота
if LOCAL_API_URL:
api_server = TelegramAPIServer.from_base(LOCAL_API_URL)
session = AiohttpSession(api=api_server)
bot = Bot(token=BOT_TOKEN, session=session)
else:
bot = Bot(token=BOT_TOKEN)
dp = Dispatcher()
# Регистрация middleware
private_middleware = PrivateMiddleware()
dp.message.middleware(private_middleware)
dp.callback_query.middleware(private_middleware)
USER_URLS: dict[int, str] = {}
USER_DATA: dict[int, dict] = {}
ACTIVE_DOWNLOADS: dict[int, dict] = {}
# -------------------- HELPERS --------------
def render_bar(percent: float, size: int = 10) -> str:
filled = int(size * percent / 100)
return "" * filled + "" * (size - filled)
def make_progress_cb(loop, message):
last_percent = {"value": 0}
last_update = {"time": 0}
async def update(d):
try:
downloaded = d.get("downloaded_bytes", 0)
total = d.get("total_bytes") or d.get("total_bytes_estimate") or 1
if total <= 0:
return
percent = min(100, downloaded * 100 / total)
current_time = time.time()
if percent - last_percent["value"] < 2 and current_time - last_update["time"] < 2:
return
last_percent["value"] = percent
last_update["time"] = current_time
bar = render_bar(percent)
eta = d.get("eta")
eta_str = str(int(float(eta))) if eta and eta != "?" else "?"
text = (
"⏬ <b>Загрузка</b>\n"
f"<code>{bar}</code> {percent:.0f}%\n"
f"⏱ Осталось: {eta_str} сек"
)
await message.edit_text(
text,
reply_markup=cancel_keyboard() if "playlist" not in message.text.lower() else None,
parse_mode="HTML"
)
except Exception as e:
logger.error(f"Error updating progress: {e}")
def cb(d):
asyncio.run_coroutine_threadsafe(update(d), loop)
return cb
async def process_tiktok_auto(message: Message, user_id: int, url: str):
"""Автоматически загружает TikTok: видео (1080p) + аудио (MP3)"""
status = await message.answer("🎵 <b>Загружаю TikTok (видео + аудио)...</b>", parse_mode="HTML")
key_video = cache_key(url, "tiktok_video", audio=False)
key_audio = cache_key(url, "tiktok_audio", audio=True)
video_cache = cache_path(CACHE_DIR, key_video, "mp4")
audio_cache = cache_path(CACHE_DIR, key_audio, "mp3")
tmp_video = os.path.join(TMP_DIR, f"{key_video}.mp4")
tmp_audio = os.path.join(TMP_DIR, f"{key_audio}.mp3")
# Проверка кэша
if os.path.exists(video_cache) and os.path.exists(audio_cache):
await status.edit_text("📤 <b>Отправляю из кэша...</b>", parse_mode="HTML")
try:
await message.answer_video(FSInputFile(video_cache), supports_streaming=True)
await message.answer_audio(FSInputFile(audio_cache))
size_mb = (os.path.getsize(video_cache) + os.path.getsize(audio_cache)) / (1024 * 1024)
await message.answer(f"✅ <b>Готово! (TikTok)</b>\n📦 Размер: {size_mb:.1f} МБ", parse_mode="HTML")
except Exception as e:
logger.error(f"TikTok cache error: {e}")
await status.edit_text("❌ Ошибка при отправке")
finally:
cleanup_tmp(TMP_DIR)
return
# Загрузка
cancel_event = threading.Event()
ACTIVE_DOWNLOADS[user_id] = {"cancel": cancel_event}
loop = asyncio.get_running_loop()
progress_cb = make_progress_cb(loop, status)
try:
await asyncio.to_thread(
download_tiktok_video_and_audio,
url,
tmp_video,
tmp_audio,
COOKIES_FILE,
cancel_event,
progress_cb,
)
if cancel_event.is_set():
await status.edit_text("⛔ Загрузка отменена")
for p in [tmp_video, tmp_audio]:
if os.path.exists(p): os.remove(p)
return
# Перемещаем в кэш
os.makedirs(os.path.dirname(video_cache), exist_ok=True)
for path in [video_cache, audio_cache]:
if os.path.exists(path): os.remove(path)
os.rename(tmp_video, video_cache)
os.rename(tmp_audio, audio_cache)
# Отправка
await status.edit_text("📤 <b>Отправляю видео + аудио...</b>", parse_mode="HTML")
await message.answer_video(FSInputFile(video_cache), supports_streaming=True)
await message.answer_audio(FSInputFile(audio_cache))
size_mb = (os.path.getsize(video_cache) + os.path.getsize(audio_cache)) / (1024 * 1024)
await message.answer(f"✅ <b>Готово! (TikTok)</b>\n📦 Размер: {size_mb:.1f} МБ", parse_mode="HTML")
except Exception as e:
logger.error(f"TikTok download error: {e}")
await status.edit_text(f"❌ Ошибка: {str(e)[:120]}")
finally:
ACTIVE_DOWNLOADS.pop(user_id, None)
cleanup_tmp(TMP_DIR)
async def process_instagram_auto(message: Message, user_id: int, url: str):
"""Автоматически загружает Instagram: видео в 1080p"""
status = await message.answer("📸 <b>Загружаю Instagram (1080p)...</b>", parse_mode="HTML")
key = cache_key(url, "instagram_video", audio=False)
final_cache = cache_path(CACHE_DIR, key, "mp4")
tmp_path = os.path.join(TMP_DIR, f"{key}.mp4")
# Проверка кэша
if os.path.exists(final_cache):
await status.edit_text("📤 <b>Отправляю из кэша...</b>", parse_mode="HTML")
try:
await message.answer_video(FSInputFile(final_cache), supports_streaming=True)
size_mb = os.path.getsize(final_cache) / (1024 * 1024)
await message.answer(f"✅ <b>Готово! (Instagram 1080p)</b>\n📦 Размер: {size_mb:.1f} МБ", parse_mode="HTML")
except Exception as e:
logger.error(f"Instagram cache error: {e}")
await status.edit_text("❌ Ошибка при отправке")
finally:
cleanup_tmp(TMP_DIR)
return
cancel_event = threading.Event()
ACTIVE_DOWNLOADS[user_id] = {"cancel": cancel_event}
loop = asyncio.get_running_loop()
progress_cb = make_progress_cb(loop, status)
try:
await asyncio.to_thread(
download_instagram_video,
url,
tmp_path,
COOKIES_FILE,
cancel_event,
progress_cb,
)
if cancel_event.is_set():
await status.edit_text("⛔ Загрузка отменена")
if os.path.exists(tmp_path): os.remove(tmp_path)
return
# Перемещаем в кэш
os.makedirs(os.path.dirname(final_cache), exist_ok=True)
if os.path.exists(final_cache): os.remove(final_cache)
os.rename(tmp_path, final_cache)
await status.edit_text("📤 <b>Отправляю видео...</b>", parse_mode="HTML")
await message.answer_video(FSInputFile(final_cache), supports_streaming=True)
size_mb = os.path.getsize(final_cache) / (1024 * 1024)
await message.answer(f"✅ <b>Готово! (Instagram 1080p)</b>\n📦 Размер: {size_mb:.1f} МБ", parse_mode="HTML")
except Exception as e:
logger.error(f"Instagram download error: {e}")
await status.edit_text(f"❌ Ошибка: {str(e)[:120]}")
finally:
ACTIVE_DOWNLOADS.pop(user_id, None)
cleanup_tmp(TMP_DIR)
# -------------------- HANDLERS --------------
@dp.message(F.text == "/start")
async def start(message: Message):
await message.answer(
"👋 <b>Привет!</b>\n\n"
"📥 Я скачиваю <b>видео</b> и <b>звук из видео</b> по ссылке.\n\n"
"✨ <b>Новые возможности:</b>\n"
"• 🎬 TikTok: автоматически видео + аудио\n"
"• 📸 Instagram: автоматически видео в 1080p\n"
"• 🎵 Отдельный звук из видео\n"
"• 📁 Плейлисты YouTube\n"
"• 🔄 Автоочистка кэша\n\n"
"👉 Просто отправь ссылку.",
parse_mode="HTML"
)
@dp.message(F.text.startswith("http"))
async def handle_link(message: Message):
url = message.text.strip()
user_id = message.from_user.id
USER_URLS[user_id] = url
# Определяем платформу
platform_info = await asyncio.to_thread(get_platform_info, url)
# TikTok — автоматически
if platform_info == "tiktok":
await process_tiktok_auto(message, user_id, url)
return
# Instagram — автоматически
if platform_info == "instagram":
await process_instagram_auto(message, user_id, url)
return
# Плейлисты
if await asyncio.to_thread(is_playlist, url):
await message.answer(
"📁 <b>Обнаружен плейлист!</b>\n\n"
"Выберите действие:",
reply_markup=playlist_keyboard(),
parse_mode="HTML"
)
return
# YouTube/VK — выбор качества
await message.answer(
"🔽 <b>Выбери качество:</b>",
reply_markup=youtube_quality_keyboard(),
parse_mode="HTML"
)
# ------------------------ YouTube/VK handlers ------------------------
@dp.callback_query(F.data.startswith("q:"))
async def handle_video(callback: CallbackQuery):
await callback.answer()
user_id = callback.from_user.id
url = USER_URLS.get(user_id)
quality = callback.data.split(":", 1)[1]
if not url:
await callback.message.answer("❌ Ссылка не найдена")
return
await callback.message.edit_reply_markup(reply_markup=None)
status = await callback.message.answer("🔍 <b>Анализирую ссылку…</b>", parse_mode="HTML")
key = cache_key(url, quality, audio=False)
final_path = cache_path(CACHE_DIR, key, "mp4")
tmp_path = os.path.join(TMP_DIR, f"{key}.mp4")
optimized_path = os.path.join(TMP_DIR, f"{key}_optimized.mp4")
# Проверяем кэш
if os.path.exists(final_path):
await status.edit_text("📤 <b>Отправляю файл из кэша…</b>", parse_mode="HTML")
try:
await callback.message.answer_video(FSInputFile(final_path))
size_mb = os.path.getsize(final_path) / 1024 / 1024
await callback.message.answer(
f"✅ <b>Готово!</b>\n📦 Размер: {size_mb:.1f} МБ",
parse_mode="HTML"
)
except Exception as e:
logger.error(f"Error sending cached file: {e}")
await status.edit_text("❌ Ошибка при отправке файла")
return
cancel_event = threading.Event()
ACTIVE_DOWNLOADS[user_id] = {"cancel": cancel_event}
loop = asyncio.get_running_loop()
progress_cb = make_progress_cb(loop, status)
try:
await asyncio.to_thread(
download_video,
url,
quality,
tmp_path,
COOKIES_FILE,
cancel_event,
progress_cb,
)
if cancel_event.is_set():
for p in [tmp_path, optimized_path]:
if os.path.exists(p): os.remove(p)
await status.edit_text("⛔ Загрузка отменена")
return
# Оптимизируем видео для Telegram (кроме оригинального качества)
if quality != "original":
await status.edit_text("⚙️ <b>Оптимизирую видео для телеграма…</b>", parse_mode="HTML")
await asyncio.to_thread(optimize_for_telegram, tmp_path, optimized_path)
for p in [tmp_path, optimized_path]:
if os.path.exists(p): os.rename(p, optimized_path)
else:
os.rename(tmp_path, final_path)
os.rename(optimized_path, final_path)
except DownloadCancelled:
for p in [tmp_path, optimized_path]:
if os.path.exists(p): os.remove(p)
await status.edit_text("⛔ Загрузка отменена")
return
except Exception as e:
logger.error(f"Error downloading video: {e}")
for p in [tmp_path, optimized_path]:
if os.path.exists(p): os.remove(p)
await status.edit_text("Не удалось скачать видео")
return
finally:
ACTIVE_DOWNLOADS.pop(user_id, None)
try:
await callback.message.answer_video(FSInputFile(final_path), supports_streaming=True)
size_mb = os.path.getsize(final_path) / (1024 * 1024)
await callback.message.answer(
f"✅ <b>Готово!</b>\n📦 Размер: {size_mb:.1f} МБ",
parse_mode="HTML"
)
except Exception as e:
logger.error(f"Error sending video: {e}")
try:
await callback.message.answer_document(FSInputFile(final_path))
size_mb = os.path.getsize(final_path) / (1024 * 1024)
await callback.message.answer(
f"✅ <b>Отправлено как документ</b>\n📦 Размер: {size_mb:.1f} МБ",
parse_mode="HTML"
)
except Exception as e2:
logger.error(f"Error sending as document: {e2}")
await status.edit_text("❌ Ошибка при отправке файла")
@dp.callback_query(F.data == "audio")
async def handle_audio(callback: CallbackQuery):
await callback.answer()
user_id = callback.from_user.id
url = USER_URLS.get(user_id)
if not url:
await callback.message.answer("❌ Ссылка не найдена")
return
await callback.message.edit_reply_markup(reply_markup=None)
status = await callback.message.answer("🎧 <b>Подготовка аудио…</b>", parse_mode="HTML")
key = cache_key(url, "audio", audio=True)
final_path = cache_path(CACHE_DIR, key, "mp3")
tmp_path = os.path.join(TMP_DIR, f"{key}.mp3")
# Создаем директории
os.makedirs(TMP_DIR, exist_ok=True)
os.makedirs(os.path.dirname(final_path), exist_ok=True)
# Проверяем кэш
if os.path.exists(final_path):
await status.edit_text("📤 <b>Отправляю аудио из кэша…</b>", parse_mode="HTML")
try:
await callback.message.answer_audio(FSInputFile(final_path))
size_mb = os.path.getsize(final_path) / (1024 * 1024)
await callback.message.answer(
f"✅ <b>Готово!</b>\n📦 Размер: {size_mb:.1f} МБ",
parse_mode="HTML"
)
except Exception as e:
logger.error(f"Error sending cached audio: {e}")
await status.edit_text("❌ Ошибка при отправке аудио")
return
cancel_event = threading.Event()
ACTIVE_DOWNLOADS[user_id] = {"cancel": cancel_event}
loop = asyncio.get_running_loop()
progress_cb = make_progress_cb(loop, status)
try:
await asyncio.to_thread(
download_audio,
url,
tmp_path,
COOKIES_FILE,
cancel_event,
progress_cb,
)
if cancel_event.is_set():
await status.edit_text("⛔ Загрузка отменена")
if os.path.exists(tmp_path): os.remove(tmp_path)
return
if not os.path.exists(tmp_path):
raise Exception("Аудио файл не был создан")
file_size = os.path.getsize(tmp_path)
if file_size == 0:
os.remove(tmp_path)
raise Exception("Создан пустой аудио файл")
# Добавляем метаданные
video_info = None
try:
video_info = await asyncio.to_thread(extract_info, url, COOKIES_FILE)
except:
pass
if video_info:
metadata = {
'title': video_info.get('title', ''),
'artist': video_info.get('uploader', ''),
'album': video_info.get('title', '')[:50],
}
await asyncio.to_thread(
lambda: add_metadata_to_audio(tmp_path, tmp_path + "_meta.mp3", metadata)
)
if os.path.exists(tmp_path + "_meta.mp3"):
os.remove(tmp_path)
os.rename(tmp_path + "_meta.mp3", tmp_path)
os.rename(tmp_path, final_path)
except DownloadCancelled:
await status.edit_text("⛔ Загрузка отменена")
if os.path.exists(tmp_path): os.remove(tmp_path)
return
except Exception as e:
logger.error(f"Error downloading audio: {str(e)}")
await status.edit_text(f"❌ Ошибка: {str(e)[:100]}")
if os.path.exists(tmp_path): os.remove(tmp_path)
return
finally:
ACTIVE_DOWNLOADS.pop(user_id, None)
await status.edit_text("📤 <b>Отправляю аудио…</b>", parse_mode="HTML")
try:
await callback.message.answer_audio(FSInputFile(final_path))
size_mb = os.path.getsize(final_path) / (1024 * 1024)
await callback.message.answer(
f"✅ <b>Готово!</b>\n📦 Размер: {size_mb:.1f} МБ",
parse_mode="HTML"
)
except Exception as e:
logger.error(f"Error sending audio: {e}")
await status.edit_text("❌ Ошибка при отправке аудио")
# -------------------- PLAYLIST HANDLERS --------------
@dp.callback_query(F.data == "playlist_all")
async def handle_playlist_all(callback: CallbackQuery):
await callback.answer()
user_id = callback.from_user.id
url = USER_URLS.get(user_id)
if not url:
await callback.message.answer("❌ Ссылка не найдена")
return
await callback.message.edit_reply_markup(reply_markup=None)
status = await callback.message.answer("📁 <b>Анализирую плейлист…</b>", parse_mode="HTML")
try:
from info import get_playlist_info
playlist_info = await asyncio.to_thread(get_playlist_info, url)
if not playlist_info or 'entries' not in playlist_info:
await status.edit_text("Не удалось получить информацию о плейлисте")
return
video_count = len(playlist_info['entries'])
if video_count == 0:
await status.edit_text("❌ Плейлист пуст")
return
if video_count > 10:
await callback.message.answer(
f"⚠️ <b>Внимание!</b>\n\n"
f"Плейлист содержит <b>{video_count}</b> видео.\n"
f"Это может занять много времени и места.\n\n"
f"Продолжить загрузку?",
reply_markup=playlist_keyboard(confirm=True),
parse_mode="HTML"
)
USER_DATA[user_id] = {"playlist_info": playlist_info, "status_message": status}
return
await download_playlist_confirm(callback, user_id, playlist_info, status)
except Exception as e:
logger.error(f"Error analyzing playlist: {e}")
await status.edit_text("❌ Ошибка при анализе плейлиста")
@dp.callback_query(F.data == "playlist_confirm_yes")
async def handle_playlist_confirm(callback: CallbackQuery):
await callback.answer()
user_id = callback.from_user.id
data = USER_DATA.get(user_id, {})
playlist_info = data.get("playlist_info")
status = data.get("status_message")
if not playlist_info or not status:
await callback.message.answer("❌ Данные плейлиста не найдены")
return
await callback.message.edit_reply_markup(reply_markup=None)
await download_playlist_confirm(callback, user_id, playlist_info, status)
async def download_playlist_confirm(callback, user_id, playlist_info, status):
try:
import uuid
import shutil
video_count = len(playlist_info['entries'])
playlist_title = playlist_info.get('title', 'Плейлист')
await status.edit_text(
f"📁 <b>Начинаю загрузку плейлиста</b>\n\n"
f"🎬 Название: {playlist_title}\n"
f"📹 Видео: {video_count}\n"
f"⏳ Подготовка...",
parse_mode="HTML"
)
cancel_event = threading.Event()
ACTIVE_DOWNLOADS[user_id] = {"cancel": cancel_event}
loop = asyncio.get_running_loop()
progress_cb = make_progress_cb(loop, status)
playlist_dir = os.path.join(TMP_DIR, f"playlist_{uuid.uuid4().hex[:8]}")
os.makedirs(playlist_dir, exist_ok=True)
downloaded_files = await asyncio.to_thread(
download_playlist_videos,
playlist_info,
playlist_dir,
COOKIES_FILE,
cancel_event,
progress_cb
)
if cancel_event.is_set():
await status.edit_text("⛔ Загрузка плейлиста отменена")
shutil.rmtree(playlist_dir, ignore_errors=True)
return
if not downloaded_files:
await status.edit_text("Не удалось загрузить видео из плейлиста")
shutil.rmtree(playlist_dir, ignore_errors=True)
return
await status.edit_text(f"📤 <b>Отправляю {len(downloaded_files)} видео…</b>", parse_mode="HTML")
downloaded_files.sort(key=lambda x: os.path.getsize(x))
sent_count = 0
for i, file_path in enumerate(downloaded_files, 1):
if cancel_event.is_set():
break
try:
file_name = os.path.basename(file_path)
display_name = os.path.splitext(file_name)[0]
await callback.message.answer_document(
FSInputFile(file_path),
caption=f"🎬 Видео {i}/{len(downloaded_files)}\n📁 {display_name[:50]}"
)
sent_count += 1
await asyncio.sleep(1)
except Exception as e:
logger.error(f"Error sending file {file_path}: {e}")
continue
total_size = sum(os.path.getsize(f) for f in downloaded_files)
total_size_mb = total_size / (1024 * 1024)
await callback.message.answer(
f"✅ <b>Плейлист загружен!</b>\n\n"
f"📁 Видео в плейлисте: {video_count}\n"
f"📤 Отправлено: {sent_count}\n"
f"💾 Общий размер: {total_size_mb:.1f} МБ\n"
f"🎬 Название: {playlist_title}",
parse_mode="HTML"
)
shutil.rmtree(playlist_dir, ignore_errors=True)
except DownloadCancelled:
await status.edit_text("⛔ Загрузка плейлиста отменена")
except Exception as e:
logger.error(f"Error downloading playlist: {e}")
await status.edit_text(f"❌ Ошибка: {str(e)[:100]}")
finally:
ACTIVE_DOWNLOADS.pop(user_id, None)
cleanup_tmp(TMP_DIR)
@dp.callback_query(F.data == "playlist_confirm_no")
async def handle_playlist_cancel(callback: CallbackQuery):
await callback.answer("Отменено", show_alert=True)
@dp.callback_query(F.data == "playlist_first")
async def handle_playlist_first(callback: CallbackQuery):
await callback.answer()
user_id = callback.from_user.id
url = USER_URLS.get(user_id)
if not url:
await callback.message.answer("❌ Ссылка не найдена")
return
try:
from info import get_first_video_from_playlist
video_url = await asyncio.to_thread(get_first_video_from_playlist, url)
if not video_url:
await callback.message.answer("Не удалось получить видео из плейлиста")
return
USER_URLS[user_id] = video_url
await callback.message.edit_reply_markup(reply_markup=None)
await callback.message.answer(
"🔽 <b>Выбери формат загрузки для первого видео:</b>",
reply_markup=youtube_quality_keyboard(),
parse_mode="HTML"
)
except Exception as e:
logger.error(f"Error getting first video: {e}")
await callback.message.answer("❌ Ошибка при получении видео из плейлиста")
@dp.callback_query(F.data == "cancel")
async def cancel_download(callback: CallbackQuery):
user_id = callback.from_user.id
data = ACTIVE_DOWNLOADS.get(user_id)
if data:
data["cancel"].set()
ACTIVE_DOWNLOADS.pop(user_id, None)
await callback.answer("⛔ Загрузка отменена", show_alert=True)
else:
await callback.answer("❌ Нет активной загрузки", show_alert=True)
# -------------------- ENTRYPOINT --------------
async def main():
# Очистка временных файлов при старте
cleanup_tmp(TMP_DIR)
try:
await dp.start_polling(bot)
except KeyboardInterrupt:
logger.info("Bot stopped by user")
if __name__ == "__main__":
try:
asyncio.run(main())
except Exception as e:
logger.error(f"Fatal error: {e}")