# 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}")