# main.py
import os
import asyncio
import threading
import logging
import time
from typing import Optional
from aiogram import Bot, Dispatcher, F
from aiogram.types import Message, CallbackQuery, FSInputFile
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,
)
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_playlist_videos,
optimize_for_telegram,
DownloadCancelled,
)
from middleware import PrivateMiddleware
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__)
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()
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] = {}
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
eta = d.get("eta")
eta_str = str(int(float(eta))) if eta and eta != "?" else "?"
bar = render_bar(percent)
text = f"⏬ Загрузка\n{bar} {percent:.0f}%\n⏱ Осталось: {eta_str} сек"
await message.edit_text(text, reply_markup=cancel_keyboard() if "playlist" in message.text.lower() else None, parse_mode="HTML")
except Exception as e:
logger.error(f"Progress error: {e}")
def cb(d):
asyncio.run_coroutine_threadsafe(update(d), loop)
return cb
# ---------------- TikTok + Instagram (автозагрузка) ----------------
async def process_tiktok_auto(message: Message, user_id: int, url: str):
status = await message.answer("🎵 Загружаю TikTok (видео + аудио)...", 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")
tmp_video = cache_path(TMP_DIR, key_video, "mp4")
os.makedirs(os.path.dirname(tmp_video), exist_ok=True)
audio_cache = cache_path(CACHE_DIR, key_audio, "mp3")
tmp_audio = cache_path(TMP_DIR, key_audio, "mp3")
os.makedirs(os.path.dirname(tmp_audio), exist_ok=True)
if os.path.exists(video_cache) and os.path.exists(audio_cache):
await status.edit_text("📤 Отправляю из кэша...", 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"✅ Готово! (TikTok)\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
# Уже готовы Telegram-safe mp4/m3, просто перемещаем
if os.path.exists(video_cache): os.remove(video_cache)
os.rename(tmp_video, video_cache)
if os.path.exists(audio_cache): os.remove(audio_cache)
os.rename(tmp_audio, audio_cache)
await status.edit_text("📤 Отправляю видео + аудио...", 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"✅ Готово! (TikTok)\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):
status = await message.answer("📸 Загружаю Instagram (1080p)...", parse_mode="HTML")
key = cache_key(url, "instagram_video", audio=False)
final_cache = cache_path(CACHE_DIR, key, "mp4")
tmp_path = cache_path(TMP_DIR, key, "mp4")
os.makedirs(os.path.dirname(tmp_path), exist_ok=True)
if os.path.exists(final_cache):
await status.edit_text("📤 Отправляю из кэша...", 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"✅ Готово! (Instagram 1080p)\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
if os.path.exists(final_cache): os.remove(final_cache)
os.rename(tmp_path, final_cache)
await status.edit_text("📤 Отправляю видео...", 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"✅ Готово! (Instagram 1080p)\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(
"👋 Привет!\n\n"
"📥 Скачиваю видео и аудио по ссылке.\n\n"
"✨ Новые возможности:\n"
"• 🎬 TikTok: автоматически видео + аудио (Telegram-safe)\n"
"• 📸 Instagram: автоматически видео в 1080p (Telegram-safe)\n"
"• 🎵 Аудио из любого видео\n"
"• 📁 Плейлисты YouTube\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)
if platform_info == "tiktok":
await process_tiktok_auto(message, user_id, url)
return
if platform_info == "instagram":
await process_instagram_auto(message, user_id, url)
return
if await asyncio.to_thread(is_playlist, url):
await message.answer(
"📁 Обнаружен плейлист!\n\n"
"Выберите действие:",
reply_markup=playlist_keyboard(),
parse_mode="HTML"
)
return
await message.answer(
"🔽 Выбери качество:",
reply_markup=youtube_quality_keyboard(),
parse_mode="HTML"
)
@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("🔍 Анализирую ссылку…", parse_mode="HTML")
key = cache_key(url, quality, audio=False)
final_path = cache_path(CACHE_DIR, key, "mp4")
tmp_path = cache_path(TMP_DIR, key, "mp4")
os.makedirs(os.path.dirname(tmp_path), exist_ok=True)
if os.path.exists(final_path):
await status.edit_text("📤 Отправляю файл из кэша…", 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"✅ Готово!\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():
if os.path.exists(tmp_path): os.remove(tmp_path)
await status.edit_text("⛔ Загрузка отменена")
return
if os.path.exists(final_path): os.remove(final_path)
os.rename(tmp_path, final_path)
except DownloadCancelled:
if os.path.exists(tmp_path): os.remove(tmp_path)
await status.edit_text("⛔ Загрузка отменена")
return
except Exception as e:
logger.error(f"Error downloading video: {e}")
if os.path.exists(tmp_path): os.remove(tmp_path)
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"✅ Готово!\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"✅ Отправлено как документ\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("🎧 Подготовка аудио…", parse_mode="HTML")
key = cache_key(url, "audio", audio=True)
final_path = cache_path(CACHE_DIR, key, "mp3")
tmp_path = cache_path(TMP_DIR, key, "mp3")
os.makedirs(os.path.dirname(tmp_path), exist_ok=True)
if os.path.exists(final_path):
await status.edit_text("📤 Отправляю аудио из кэша…", 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"✅ Готово!\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("Аудио файл не был создан")
if os.path.getsize(tmp_path) == 0:
os.remove(tmp_path)
raise Exception("Создан пустой аудио файл")
if os.path.exists(final_path): os.remove(final_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("📤 Отправляю аудио…", 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"✅ Готово!\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("📁 Анализирую плейлист…", 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"⚠️ Внимание!\n\n"
f"Плейлист содержит {video_count} видео.\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):
import uuid
import shutil
video_count = len(playlist_info['entries'])
playlist_title = playlist_info.get('title', 'Плейлист')
await status.edit_text(
f"📁 Начинаю загрузку плейлиста\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)
try:
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"📤 Отправляю {len(downloaded_files)} видео…", 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"✅ Плейлист загружен!\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(
"🔽 Выбери формат загрузки для первого видео:",
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)
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}")