"""Download engine using yt-dlp Python API with progress callbacks and sequential queue processing.""" import logging import os import re import threading import time from datetime import datetime, timezone from pathlib import Path from typing import Optional import yt_dlp from models import ArchiveItem, QueueItem from models.archive import ArchiveDB from models.queue_store import QueueStore logger = logging.getLogger(__name__) class DownloadEngine: """Handles video/playlist downloads with sequential queue processing.""" def __init__(self, queue_store: QueueStore, archive_db: ArchiveDB, yt_cli=None, socketio=None): self.queue_store = queue_store self.archive_db = archive_db self.yt_cli = yt_cli self.socketio = socketio self._active_download_id = None self._yt_dlp_instance = None self._queue_lock = threading.Lock() self._queue_processor_thread = None self._stop_event = threading.Event() self._recover_in_progress_downloads() self._start_queue_processor() def _recover_in_progress_downloads(self): """Recover downloads that were in progress when the server crashed.""" try: items = self.queue_store.get_all() recovered = 0 for item in items: if item.status == "downloading": logger.info(f"Recovering in-progress download: {item.id} ({item.title})") self.queue_store.update_status(item.id, "pending") self.queue_store.update_progress(item.id, 0.0) recovered += 1 if recovered > 0: logger.info(f"Recovered {recovered} in-progress download(s) from crash") except Exception as e: logger.error(f"Failed to recover in-progress downloads: {e}") def _start_queue_processor(self): """Start the background queue processor thread.""" self._stop_event.clear() self._queue_processor_thread = threading.Thread( target=self._queue_processor_loop, daemon=True ) self._queue_processor_thread.start() logger.info("Queue processor started") def _queue_processor_loop(self): """Main loop that processes one download at a time.""" while not self._stop_event.is_set(): with self._queue_lock: if self._active_download_id is not None: time.sleep(1) continue # Find next pending item all_items = self.queue_store.get_all() pending = [item for item in all_items if item.status == "pending"] if not pending: time.sleep(2) continue # Sort by added_at to process oldest first pending.sort(key=lambda x: x.added_at) next_item = pending[0] # Mark as downloading self._active_download_id = next_item.id self.queue_store.update_status(next_item.id, "downloading") self.queue_store.update_progress(next_item.id, 0.0) self._broadcast(next_item.id, "download:status", { "queueId": next_item.id, "status": "downloading", "progress": 0 }) # Run the actual download (blocks until done) if next_item.item_type == "playlist": self._run_playlist_download(next_item) else: self._run_video_download(next_item) # Release lock for next iteration with self._queue_lock: self._active_download_id = None self._yt_dlp_instance = None def _run_video_download(self, item: QueueItem): """Run a single video download synchronously.""" config = self.yt_cli.config url = item.url queue_id = item.id category = item.category network_folder = item.network_folder quality = item.quality base_dir = Path(config["download_dir"]) if not base_dir.exists(): base_dir.mkdir(parents=True, exist_ok=True) download_dir = base_dir / category if category else base_dir download_dir.mkdir(parents=True, exist_ok=True) # Build yt-dlp options with thumbnail support ytdlp_args = config.get("yt_dlp_args", {}) default_format = ytdlp_args.get("format", "bestvideo[height<=1080]+bestaudio/best") # Apply user quality preference fmt = self._build_format(quality, default_format) logger.info(f"Download {queue_id}: quality={quality}, format={fmt}") # Use %(title)s.%(ext)s template so thumbnail gets same base name output_template = str(download_dir / "%(title)s.%(ext)s") ydl_opts = { "format": fmt, "outtmpl": output_template, "write_thumbnail": True, "thumbnail_format": "jpg", "no_warnings": False, "restrict_filenames": True, "user_agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36", "retries": 5, "fragment_retries": 5, "extract_retries": 3, "concurrent_fragment_downloads": 4, "overwrites": True, "continuedl": True, "extractor_args": {"youtube": {"player_client": ["web", "ios", "android", "tv", "mediaconnect"]}}, } video_info = {} downloaded_filepath = None def progress_callback(d): nonlocal downloaded_filepath if d["status"] == "downloading": total = d.get("total_bytes") or 1 progress = d.get("downloaded_bytes", 0) / total * 100 speed = d.get("speed") speed_str = f"{speed / 1024 / 1024:.1f} MB/s" if speed else None eta = d.get("eta") eta_str = f"{int(eta)}s" if eta else None self.queue_store.update_progress(queue_id, progress, speed_str, eta_str) self._broadcast(queue_id, "download:progress", { "queueId": queue_id, "progress": round(progress, 1), "speed": speed_str, "eta": eta_str, }) elif d["status"] == "finished": downloaded_filepath = d.get("filename", "") self.queue_store.update_progress(queue_id, 100.0) self._broadcast(queue_id, "download:progress", { "queueId": queue_id, "progress": 100, "speed": None, "eta": None }) ydl_opts["progress_hooks"] = [progress_callback] try: ydl = yt_dlp.YoutubeDL(ydl_opts) self._yt_dlp_instance = ydl # Pre-fetch metadata try: info = ydl.extract_info(url, download=False) if info: video_info.update(info) except Exception as e: logger.warning(f"Failed to pre-fetch metadata: {e}") ydl.download([url]) # Find the actual downloaded file if not downloaded_filepath: downloaded_filepath = self._find_downloaded_file(download_dir, video_info) file_size = 0 if downloaded_filepath and os.path.exists(downloaded_filepath): file_size = os.path.exists(downloaded_filepath) and os.path.getsize(downloaded_filepath) or 0 # Check thumbnail was downloaded thumbnail_path = None if downloaded_filepath: base = os.path.splitext(downloaded_filepath)[0] for ext in ['.jpg', '.jpeg', '.webp', '.png']: tp = base + ext if os.path.exists(tp): thumbnail_path = tp break self.queue_store.update_status( queue_id, "completed", download_path=downloaded_filepath or "", file_size=str(file_size) ) # Archive vid = video_info.get("id", "") if vid: archive_item = ArchiveItem( video_id=vid, title=video_info.get("title", "Unknown Title"), url=video_info.get("webpage_url", url), description=video_info.get("description", ""), thumbnail=video_info.get("thumbnail", ""), channel=video_info.get("uploader", ""), views=video_info.get("view_count", 0) or 0, duration=self._format_duration(video_info.get("duration", 0)), category=category or "", download_path=downloaded_filepath or "", file_size=file_size, download_date=datetime.now(timezone.utc).isoformat(), ) self.archive_db.add_video(archive_item) # Network share copy if network_folder and config.get("network_share_path") and downloaded_filepath: self._copy_to_network_share(downloaded_filepath, config, network_folder) # Also copy thumbnail if it exists if thumbnail_path: self._copy_to_network_share(thumbnail_path, config, network_folder) self._broadcast(queue_id, "download:complete", { "queueId": queue_id, "downloadPath": downloaded_filepath or "", "fileSize": file_size, "thumbnailPath": thumbnail_path, }) except Exception as e: logger.error(f"Download failed for {queue_id}: {type(e).__name__}: {e}", exc_info=True) self.queue_store.update_status(queue_id, "failed", error_message=f"{type(e).__name__}: {e}") self._broadcast(queue_id, "download:failed", { "queueId": queue_id, "error": f"{type(e).__name__}: {e}" }) def _run_playlist_download(self, item: QueueItem): """Run a playlist download synchronously.""" config = self.yt_cli.config url = item.url queue_id = item.id category = item.category _network_folder = item.network_folder quality = item.quality base_dir = Path(config["download_dir"]) if not base_dir.exists(): base_dir.mkdir(parents=True, exist_ok=True) download_dir = base_dir / category if category else base_dir download_dir.mkdir(parents=True, exist_ok=True) # Get playlist title playlist_title = "Unknown Playlist" try: ydl_info = yt_dlp.YoutubeDL({ "flat_playlist": True, "no_warnings": True, }) info = ydl_info.extract_info(url, download=False) if info: playlist_title = info.get("title", "Unknown Playlist") except Exception: pass playlist_dir = download_dir / playlist_title playlist_dir.mkdir(parents=True, exist_ok=True) ytdlp_args = config.get("yt_dlp_args", {}) default_format = ytdlp_args.get("format", "bestvideo[height<=1080]+bestaudio/best") # Apply user quality preference fmt = self._build_format(quality, default_format) ydl_opts = { "format": fmt, "outtmpl": str(playlist_dir / "%(title)s.%(ext)s"), "write_thumbnail": True, "thumbnail_format": "jpg", "no_warnings": False, "restrict_filenames": True, "user_agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36", "retries": 5, "fragment_retries": 5, "extract_retries": 3, "concurrent_fragment_downloads": 4, "overwrites": True, "continuedl": True, "extractor_args": {"youtube": {"player_client": ["web", "ios", "android", "tv", "mediaconnect"]}}, } total_videos = None completed_videos = 0 def progress_callback(d): nonlocal completed_videos if d["status"] == "downloading": if total_videos: progress = (completed_videos / total_videos) * 100 else: total = d.get("total_bytes") or 1 progress = d.get("downloaded_bytes", 0) / total * 100 speed = d.get("speed") speed_str = f"{speed / 1024 / 1024:.1f} MB/s" if speed else None eta = d.get("eta") eta_str = f"{int(eta)}s" if eta else None self.queue_store.update_progress(queue_id, progress, speed_str, eta_str) self._broadcast(queue_id, "download:progress", { "queueId": queue_id, "progress": round(progress, 1), "speed": speed_str, "eta": eta_str }) elif d["status"] == "finished": completed_videos += 1 if total_videos: progress = (completed_videos / total_videos) * 100 else: progress = 100 self.queue_store.update_progress(queue_id, progress) self._broadcast(queue_id, "download:progress", { "queueId": queue_id, "progress": round(progress, 1) }) ydl_opts["progress_hooks"] = [progress_callback] try: ydl = yt_dlp.YoutubeDL(ydl_opts) self._yt_dlp_instance = ydl try: info = ydl.extract_info(url, download=False) if info and "entries" in info: total_videos = len(info["entries"]) except Exception: pass ydl.download([url]) self.queue_store.update_status(queue_id, "completed") self._broadcast(queue_id, "download:complete", { "queueId": queue_id, "videoCount": completed_videos }) # Archive playlist playlist_id = None id_match = re.search(r"(?:list=|\/)([0-9A-Za-z_-]{30,})", url) if id_match: playlist_id = id_match.group(1) archive_item = ArchiveItem( video_id=f"playlist_{playlist_id or 'unknown'}", title=f"Playlist: {playlist_title}", url=url, category=category or "", download_path=str(playlist_dir), download_date=datetime.now(timezone.utc).isoformat(), item_type="playlist", ) self.archive_db.add_video(archive_item) except Exception as e: logger.error(f"Playlist download failed for {queue_id}: {type(e).__name__}: {e}", exc_info=True) self.queue_store.update_status(queue_id, "failed", error_message=f"{type(e).__name__}: {e}") self._broadcast(queue_id, "download:failed", { "queueId": queue_id, "error": f"{type(e).__name__}: {e}" }) def enqueue_download(self, item: QueueItem): """Add a video to the queue (will be processed in order).""" self.queue_store.add_item(item) self._broadcast(item.id, "queue:enqueued", { "queueId": item.id, "status": "pending", "message": "Added to queue" }) def download_video(self, queue_id: str, url: str, config: dict, category: str = None, network_folder: str = None, quality: str = None): """Start a video download directly (bypasses queue).""" item = self.queue_store.get_item(queue_id) if not item or queue_id in (self._active_download_id,): return False # Mark immediately and run synchronously in a thread self.queue_store.update_status(queue_id, "downloading") self.queue_store.update_progress(queue_id, 0.0) self._broadcast(queue_id, "download:status", { "queueId": queue_id, "status": "downloading", "progress": 0 }) def _run_direct(): direct_item = self.queue_store.get_item(queue_id) if direct_item: self._run_video_download(direct_item) with self._queue_lock: if self._active_download_id == queue_id: self._active_download_id = None self._yt_dlp_instance = None threading.Thread(target=_run_direct, daemon=True).start() return True def download_playlist(self, queue_id: str, url: str, config: dict, category: str = None, network_folder: str = None, quality: str = None): """Start a playlist download directly (bypasses queue).""" item = self.queue_store.get_item(queue_id) if not item or queue_id == self._active_download_id: return False self.queue_store.update_status(queue_id, "downloading") self.queue_store.update_progress(queue_id, 0.0) self._broadcast(queue_id, "download:status", { "queueId": queue_id, "status": "downloading", "progress": 0 }) def _run_direct(): direct_item = self.queue_store.get_item(queue_id) if direct_item: self._run_playlist_download(direct_item) with self._queue_lock: if self._active_download_id == queue_id: self._active_download_id = None self._yt_dlp_instance = None threading.Thread(target=_run_direct, daemon=True).start() return True def cancel_download(self, queue_id: str) -> bool: """Cancel an active or pending download.""" item = self.queue_store.get_item(queue_id) if not item: return False # If it's the currently active download, try to cancel it if queue_id == self._active_download_id: ydl = self._yt_dlp_instance if ydl: try: ydl.quiet = True except Exception: pass self.queue_store.update_status(queue_id, "cancelled") with self._queue_lock: self._active_download_id = None self._yt_dlp_instance = None self._broadcast(queue_id, "download:status", { "queueId": queue_id, "status": "cancelled" }) return True # If it's pending in queue, just mark as cancelled if item.status == "pending": self.queue_store.update_status(queue_id, "cancelled") self._broadcast(queue_id, "download:status", { "queueId": queue_id, "status": "cancelled" }) return True return False def _find_downloaded_file(self, directory: Path, video_info: dict) -> Optional[str]: """Find the most recently downloaded video file in a directory.""" try: video_extensions = ['.mp4', '.mkv', '.webm', '.flv'] files = [] for f in directory.iterdir(): if f.is_file() and f.suffix.lower() in video_extensions: files.append(f) if files: latest = max(files, key=lambda f: f.stat().st_mtime) return str(latest) except Exception: pass return None def _broadcast(self, queue_id: str, event: str, data: dict): """Broadcast a WebSocket event.""" if self.socketio: try: self.socketio.emit(event, data) except Exception: pass def _copy_to_network_share(self, filepath: str, config: dict, network_folder: str): """Copy downloaded file to network share.""" try: import shutil network_path = Path(config["network_share_path"]) dest_dir = network_path / network_folder dest_dir.mkdir(parents=True, exist_ok=True) src = Path(filepath) if src.exists(): shutil.copy2(src, dest_dir / src.name) except Exception as e: logger.warning(f"Failed to copy to network share: {e}") def _build_format(self, quality: Optional[str], default_format: str) -> str: """Build yt-dlp format string based on quality setting.""" if quality and quality != "best": return f"bestvideo[height<={quality}]+bestaudio/best" return default_format def _format_duration(self, seconds): """Convert seconds to MM:SS or HH:MM:SS format.""" if not seconds: return "0:00" hours = int(seconds // 3600) minutes = int((seconds % 3600) // 60) secs = int(seconds % 60) if hours > 0: return f"{hours}:{minutes:02d}:{secs:02d}" return f"{minutes}:{secs:02d}"