import os import re import logging import queue import subprocess import threading import time import uuid import typing from typing import Any, Dict, Optional, List, Union from pathlib import Path # Configure logging logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s") logger = logging.getLogger(__name__) class DownloadService: """Service for managing downloads using spotdl.""" def emit_progress(self, download_id: str, progress: float, stage: str, status: str) -> None: """Emit download progress to clients via SocketIO and update history. Args: download_id: Unique identifier for the download (usually the URL) progress: Progress percentage (0-100) stage: Current stage of the download (e.g., 'queued', 'downloading', 'completed', 'error') status: Status message to display """ if download_id in self.download_history: self.download_history[download_id].update({ "progress": progress, "stage": stage, "status": status, "timestamp": int(time.time()) }) # Build progress payload with full metadata for the UI info = self.download_history.get(download_id, {}) progress_payload = { "id": download_id, "url": info.get("url", download_id), "name": info.get("name", ""), "artist": info.get("artist", ""), "type": info.get("type", "track"), "progress": progress, "stage": stage, "status": status } self.socketio.emit("download_progress", progress_payload) # Also push aggregated status history self.socketio.emit("status", {"downloads": list(self.download_history.values())}) def __init__(self, config, playlist_manager, socketio, download_queue, download_history): """Initialize the download service. Args: config: Configuration service instance playlist_manager: Playlist manager instance socketio: SocketIO instance for real-time updates download_queue: Shared queue for download items download_history: Shared dictionary for download history """ self.config = config self.playlist_manager = playlist_manager self.socketio = socketio self.download_queue = download_queue self.download_history = download_history self.active_downloads = {} self.spodtdl_subprocess = None self._stop_event = threading.Event() self.worker_thread = threading.Thread( target=self.process_downloads, # fixed method name daemon=True ) self.active_process = None def _start_download_thread(self) -> None: """Start the background thread for processing downloads.""" self._stop_event.clear() self.worker_thread.start() def add_item_to_queue(self, data: Dict[str, Any]) -> None: """Add a download item to the queue. Args: data: Dictionary containing download details - type: 'track', 'album', 'playlist', or 'artist' - name: Name of the item - artist: Artist name (for tracks/albums) - url: Spotify URL """ spotify_url = data.get("url") item_type = data.get("type", "track") item_name = data.get("name", "Unknown") item_artist = data.get("artist", "Various Artists") # Always generate a unique id for each download (not just url) download_id = data.get("id") or str(uuid.uuid4()) download_info = { "id": download_id, "name": item_name, "type": item_type, "artist": item_artist, "url": spotify_url, "status": "Pending...", "progress": 0, "stage": "queued", "timestamp": int(time.time()) } logger.info(f"Adding to download queue: {download_info}") # Add to queue and history self.download_queue.put((spotify_url, download_info)) self.download_history[download_id] = download_info # Emit initial status self.socketio.emit("status", {"downloads": list(self.download_history.values())}) self.socketio.emit("download_progress", { "id": download_id, "url": spotify_url, "name": item_name, "artist": item_artist, "type": item_type, "progress": 0, "stage": "queued", "status": "Pending..." }) logger.info(f"Added to queue: {item_name} by {item_artist} (id={download_id})") def get_status(self, download_id: Optional[str] = None) -> Union[Dict[str, Any], List[Dict[str, Any]]]: """Get the status of a specific download or all downloads.""" if download_id: return self.download_history.get(download_id, {}) return list(self.download_history.values()) def cancel_download(self, download_id: str) -> bool: """Cancel a queued or active download. Args: download_id: ID of the download to cancel Returns: bool: True if cancelled, False if not found """ # Cancel active download if download_id in self.active_downloads: if self.active_process and self.active_process.poll() is None: try: self.active_process.terminate() except Exception as e: logger.warning(f"Error terminating active process for {download_id}: {e}") self.active_process = None self.active_downloads[download_id]["status"] = "cancelled" self.active_downloads[download_id]["stage"] = "cancelled" self.active_downloads[download_id]["message"] = "Download cancelled" self.active_downloads[download_id]["progress"] = 0 # Move to history and clean up self.download_history[download_id] = self.active_downloads.pop(download_id) self._broadcast_update(self.download_history[download_id]) logger.info(f"Cancelled active download: {download_id}") return True # Cancel queued download with self.download_queue.mutex: items = list(self.download_queue.queue) found = False for i, (url, item) in enumerate(items): if item.get("id") == download_id: del items[i] found = True break if found: self.download_queue.queue.clear() for entry in items: self.download_queue.put(entry) if download_id in self.download_history: self.download_history[download_id].update({ "status": "cancelled", "stage": "cancelled", "message": "Download cancelled", "progress": 0 }) self._broadcast_update(self.download_history[download_id]) logger.info(f"Cancelled queued download: {download_id}") return True logger.info(f"Cancel requested but not found: {download_id}") return False def process_downloads(self) -> None: """Process downloads from the queue in a background thread. This method runs in a loop, processing download items from the queue one at a time. It handles the entire download lifecycle including progress tracking, error handling, and cleanup. """ while not self._stop_event.is_set(): try: try: url, download_info = self.download_queue.get(timeout=1) except queue.Empty: continue download_id = download_info.get('id') logger.info(f"Starting download process for {download_id}") if download_info.get("status") == "cancelled": logger.info(f"Skipping cancelled download: {download_id}") self.download_queue.task_done() continue process = None try: self.active_downloads[download_id] = download_info self.emit_progress(download_id, 5, "preparing", "Preparing download...") try: cmd = self._build_spotdl_command(url, download_info) logger.info(f"Executing spotdl command: {' '.join(cmd)}") logger.info(f"Download type: {download_info.get('type', 'unknown')}, URL: {url}") except Exception as e: error_msg = f"Failed to build download command: {str(e)}" logger.error(error_msg, exc_info=True) raise RuntimeError(error_msg) from e self.emit_progress(download_id, 10, "starting", "Starting download...") try: output_dir = self._get_output_path(download_info) # Ensure output directory exists os.makedirs(output_dir, exist_ok=True) process = subprocess.Popen( cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, bufsize=1, env=os.environ.copy() # Don't set cwd since we're using absolute paths in --output ) self.active_process = process logger.info(f"Started download process (PID: {process.pid}) for {download_id}") # Enhanced detailed progress tracking with comprehensive stage detection last_progress = 0 current_stage = "initializing" current_track = "" tracks_processed = 0 total_tracks = 1 # Default for single track # Extract expected track count for albums if download_info.get('type') == 'album': total_tracks = download_info.get('track_count', 1) while True: if self._stop_event.is_set(): logger.warning("Shutdown requested, terminating download process") process.terminate() raise KeyboardInterrupt("Shutdown requested") return_code = process.poll() if return_code is not None: break stdout_line = process.stdout.readline() if process.stdout else '' if stdout_line: stdout_line = stdout_line.strip() logger.info(f"spotdl: {stdout_line}") # Changed to info level for visibility # Detect stage changes and update progress with detailed information new_progress = None stage_changed = False # Stage 1: Initialization and metadata if any(word in stdout_line.lower() for word in ['initializing', 'starting', 'begin']): if current_stage != "initializing": current_stage = "initializing" stage_changed = True new_progress = max(5, last_progress) # Stage 2: Searching for tracks elif any(word in stdout_line.lower() for word in ['searching', 'finding', 'looking', 'query']): if current_stage != "searching": current_stage = "searching" stage_changed = True new_progress = max(15, last_progress) # Stage 3: Metadata retrieval elif any(word in stdout_line.lower() for word in ['metadata', 'info', 'details', 'fetching']): if current_stage != "metadata": current_stage = "metadata" stage_changed = True new_progress = max(25, last_progress) # Stage 4: Audio downloading elif any(word in stdout_line.lower() for word in ['downloading', 'download', 'fetching audio']): if current_stage != "downloading": current_stage = "downloading" stage_changed = True # Stage 5: Audio processing/conversion elif any(word in stdout_line.lower() for word in ['processing', 'converting', 'encoding', 'ffmpeg']): if current_stage != "processing": current_stage = "processing" stage_changed = True new_progress = max(80, last_progress) # Stage 6: File operations elif any(word in stdout_line.lower() for word in ['saving', 'writing', 'moving', 'organizing']): if current_stage != "saving": current_stage = "saving" stage_changed = True new_progress = max(90, last_progress) # Track completion detection if any(phrase in stdout_line.lower() for phrase in ['saved', 'complete', 'finished', 'done']): tracks_processed += 1 if download_info.get('type') == 'album': new_progress = min(95, (tracks_processed / total_tracks) * 90 + 10) # Extract current track name for display track_match = re.search(r'(?:downloading|processing|saving).*?["\']([^"\']+)["\']', stdout_line, re.IGNORECASE) if track_match: current_track = track_match.group(1)[:50] # Limit length # Parse actual progress from line parsed_progress = self._parse_progress(stdout_line) if parsed_progress is not None: new_progress = parsed_progress # Generate detailed status messages if new_progress is not None and (new_progress > last_progress or stage_changed): last_progress = max(new_progress, last_progress) # Create detailed status message based on current stage if current_stage == "initializing": status_msg = f"🚀 Initializing download... ({new_progress:.0f}%)" stage_label = "initializing" elif current_stage == "searching": item_name = download_info.get('name', 'content') status_msg = f"🔍 Searching for {item_name}... ({new_progress:.0f}%)" stage_label = "searching" elif current_stage == "metadata": status_msg = f"📋 Fetching track metadata... ({new_progress:.0f}%)" stage_label = "metadata" elif current_stage == "downloading": if current_track: if download_info.get('type') == 'album': status_msg = f"⬇️ Downloading: {current_track} ({tracks_processed}/{total_tracks}) - {new_progress:.0f}%" else: status_msg = f"⬇️ Downloading: {current_track} - {new_progress:.0f}%" else: status_msg = f"⬇️ Downloading audio... ({new_progress:.0f}%)" stage_label = "downloading" elif current_stage == "processing": if current_track: status_msg = f"🎵 Processing: {current_track} - {new_progress:.0f}%" else: status_msg = f"🎵 Processing audio files... ({new_progress:.0f}%)" stage_label = "converting" elif current_stage == "saving": status_msg = f"💾 Saving to library... ({new_progress:.0f}%)" stage_label = "saving" else: status_msg = f"📈 Progress: {new_progress:.0f}%" stage_label = "downloading" # Add album progress indicator if download_info.get('type') == 'album' and total_tracks > 1: status_msg += f" | Album: {tracks_processed}/{total_tracks} tracks" self.emit_progress(download_id, new_progress, stage_label, status_msg) logger.info(f"Progress update: {status_msg}") stderr_line = process.stderr.readline() if process.stderr else '' if stderr_line: stderr_line = stderr_line.strip() logger.warning(f"spotdl stderr: {stderr_line}") # Enhanced error detection and reporting if any(err in stderr_line.lower() for err in ['error', 'failed', 'exception', 'unable', 'cannot']): error_msg = f"⚠️ {stderr_line[:150]}..." self.emit_progress(download_id, last_progress, "error", error_msg) logger.error(f"Download error for {download_id}: {stderr_line}") elif any(warn in stderr_line.lower() for warn in ['warning', 'skip', 'retry']): warn_msg = f"⚠️ {stderr_line[:100]}..." self.emit_progress(download_id, last_progress, current_stage, warn_msg) time.sleep(0.05) # Reduced sleep for more responsive updates if return_code == 0: # Show final completion stages self.emit_progress(download_id, 95, "finalizing", "🔄 Finalizing download...") logger.info(f"Download process completed successfully: {download_id}") # Check if files were actually created try: files_created = [] if os.path.exists(output_dir): for root, dirs, files in os.walk(output_dir): for file in files: if file.endswith(('.mp3', '.flac', '.m4a', '.ogg')): files_created.append(os.path.join(root, file)) if files_created: file_count = len(files_created) if download_info.get('type') == 'album': self.emit_progress(download_id, 98, "finalizing", f"✅ Album downloaded: {file_count} tracks saved to library") else: self.emit_progress(download_id, 98, "finalizing", f"✅ Track downloaded and saved to library") logger.info(f"Successfully created {file_count} audio files for {download_id}") else: logger.warning(f"Download completed but no audio files found in {output_dir}") self.emit_progress(download_id, 90, "warning", "⚠️ Download completed but no audio files found") except Exception as e: logger.error(f"Error checking downloaded files: {e}") # Trigger Jellyfin library refresh if enabled if self.config.get("trigger_jellyfin_scan", False): try: self.emit_progress(download_id, 99, "scanning", "📚 Updating Jellyfin library...") # Use the download output directory for a targeted refresh self.playlist_manager.trigger_media_refresh([output_dir]) logger.info("Jellyfin library refresh triggered after download") self.emit_progress(download_id, 100, "complete", "🎶 Download complete! Added to Jellyfin library") except Exception as e: logger.error(f"Failed to trigger Jellyfin refresh: {e}", exc_info=True) self.emit_progress(download_id, 100, "complete", "🎶 Download complete! (Jellyfin scan failed)") else: self.emit_progress(download_id, 100, "complete", "🎶 Download completed successfully! Enjoy your music!") logger.info(f"Download fully completed: {download_id}") else: try: stdout, stderr = process.communicate(timeout=5) except Exception: stdout, stderr = '', '' error_msg = f"Download failed with return code {return_code}" logger.error(f"Command that failed: {' '.join(cmd)}") if stderr: logger.error(f"Stderr output: {stderr.strip()}") error_msg = f"{error_msg}: {stderr.strip()}" elif stdout: logger.error(f"Stdout output: {stdout.strip()}") error_msg = f"{error_msg}: {stdout.strip()}" logger.error(error_msg) self.emit_progress(download_id, 0, "failed", error_msg[:500]) except subprocess.TimeoutExpired: if process: process.kill() stdout, stderr = process.communicate() else: stdout, stderr = '', '' error_msg = "Download timed out" if stderr: error_msg = f"{error_msg}: {stderr.strip()}" logger.error(error_msg) self.emit_progress(download_id, 0, "failed", error_msg) except Exception as e: error_msg = f"Download error: {str(e)}" logger.error(error_msg, exc_info=True) self.emit_progress(download_id, 0, "error", error_msg) except Exception as e: error_msg = f"Unexpected error during download: {str(e)}" logger.error(error_msg, exc_info=True) self.emit_progress(download_id, 0, "error", error_msg) finally: if process and process.poll() is None: try: process.terminate() process.wait(timeout=5) except Exception as e: logger.warning(f"Error terminating process: {e}") download_info['timestamp'] = self._get_timestamp() self.download_history[download_id] = download_info self.active_downloads.pop(download_id, None) self.download_queue.task_done() if hasattr(self, 'active_process') and self.active_process: if self.active_process.poll() is None: self.active_process.terminate() self.active_process = None logger.info(f"Completed processing for {download_id}") except Exception as e: logger.critical(f"Critical error in download worker: {e}", exc_info=True) time.sleep(5) # Prevent tight loop on repeated errors def _build_spotdl_command(self, url: str, download_info: Dict[str, Any]) -> List[str]: """Build the spotdl command for a download. This method constructs the command-line arguments for spotdl based on the download type and configuration. It handles different types of downloads (track, album, playlist, etc.) and applies appropriate output templates and options. Args: url: URL to download (Spotify URL or search query) download_info: Dictionary containing download details including: - type: Type of download ('track', 'album', 'playlist', 'artist') - output_path: Base output directory (optional, defaults to /data) - format: Output format (defaults to 'mp3') - bitrate: Audio bitrate (defaults to '320k') - additional_args: List of additional spotdl arguments (optional) Returns: List[str]: Command and arguments for subprocess.Popen Raises: ValueError: If required parameters are missing or invalid """ if not url or not isinstance(url, str): raise ValueError("Invalid or missing URL") # Get download type and validate download_type = download_info.get('type', 'track') if download_type not in ['track', 'album', 'playlist', 'artist']: logger.warning(f"Unknown download type '{download_type}', defaulting to 'track'") download_type = 'track' # Get output path and ensure it exists output_path = self._get_output_path(download_info) os.makedirs(output_path, exist_ok=True) # Determine audio format early (used below) audio_format = str(download_info.get('format', 'mp3')).lower() if audio_format not in ['mp3', 'flac', 'ogg', 'opus', 'm4a', 'wav']: audio_format = 'mp3' # sanitize invalid format # Build spotdl CLI command (correct format: spotdl download [URL] [options]) cmd = [ 'spotdl', 'download', url, # URL must come immediately after download subcommand '--format', audio_format, '--bitrate', str(download_info.get('bitrate', '320k')), '--overwrite', 'skip', '--print-errors', '--sponsor-block', '--threads', str(min(4, (os.cpu_count() or 1))), '--max-retries', '3', ] # Add output template based on download type for proper folder structure if download_type == 'track': # For single tracks: let spotdl handle album detection naturally cmd += ['--output', '/data/{artist}/{album}/{title}.{output-ext}'] elif download_type == 'album': # For albums: Artist/Album - (Year)/Track Number Artist - Track Name.ext cmd += ['--output', '/data/{artist}/{album} - ({year})/{track-number:02d} {artist} - {title}.{output-ext}'] elif download_type == 'playlist': # For playlists: keep playlist structure cmd += ['--output', '/data/Playlists/{list-name}/{track-number:02d} {artist} - {title}.{output-ext}'] elif download_type == 'artist': # For artist downloads: Artist/Album - (Year)/Track Number Artist - Track Name.ext cmd += ['--output', '/data/{artist}/{album} - ({year})/{track-number:02d} {artist} - {title}.{output-ext}'] else: # Default fallback format - single track cmd += ['--output', '/data/{artist}/{title}.{output-ext}'] # Type-specific flags if download_type == 'album': cmd += ['--album-type', 'album'] elif download_type == 'playlist': cmd += ['--playlist-numbering'] # artist / track need no extra flags # Additional CLI args from request additional_args = download_info.get('additional_args', []) if additional_args and isinstance(additional_args, list): cmd += [str(arg) for arg in additional_args] logger.debug(f"Built spotdl command: {' '.join(cmd)}") return cmd def _get_output_path(self, download_item: Dict[str, Any]) -> str: """Get the output directory for a download item. Args: download_item: Download item details with 'type' and other metadata Returns: str: Output directory path (using container paths) Raises: ValueError: If download_item is missing required fields """ if not download_item or 'type' not in download_item: raise ValueError("Invalid download item: missing 'type' field") # Use the container path for music files base_dir = '/data' download_type = download_item['type'] try: if download_type == 'track': artist = self._sanitize_filename(download_item.get('artist', 'Unknown')) album = self._sanitize_filename(download_item.get('album', 'Single')) year = str(download_item.get('year', '')).strip() album_dir = f"{album} - ({year})" if year and year.lower() != 'none' and year != '0' else album output_path = os.path.join(base_dir, artist, album_dir) elif download_type == 'album': artist = self._sanitize_filename(download_item.get('artist', 'Unknown')) album = self._sanitize_filename(download_item.get('name', 'Unknown Album')) year = str(download_item.get('year', '')).strip() album_dir = f"{album} - ({year})" if year and year.lower() != 'none' and year != '0' else album output_path = os.path.join(base_dir, artist, album_dir) elif download_type == 'playlist': playlist_name = download_item.get('name', 'Unnamed Playlist') if not playlist_name or playlist_name.lower() == 'none': playlist_name = 'Unnamed Playlist' playlist_dir = self._sanitize_filename(playlist_name) output_path = os.path.join(base_dir, 'Playlists', playlist_dir) elif download_type == 'artist': artist = download_item.get('name', 'Unknown Artist') if not artist or artist.lower() == 'none': artist = 'Unknown Artist' output_path = os.path.join(base_dir, self._sanitize_filename(artist)) else: output_path = os.path.join(base_dir, 'Downloads') # Ensure the directory exists and has correct permissions os.makedirs(output_path, exist_ok=True) os.chmod(output_path, 0o755) # Ensure proper permissions logger.debug(f"Output path for {download_type} download: {output_path}") return output_path except Exception as e: logger.error(f"Error determining output path: {e}") # Fallback to a safe directory if there's any error fallback_path = os.path.join(base_dir, 'Downloads') os.makedirs(fallback_path, exist_ok=True) return fallback_path @staticmethod def _parse_progress(line: str) -> Optional[float]: """Parse download progress from spotdl output with enhanced pattern recognition. Args: line: Output line from spotdl output Returns: float: Progress percentage (0-100) or None if no progress found Example formats: " 98%|█████████▊| 49.0/50.0 [00:10<00:00, 4.90s/s]" "[download] 99.9% of 5.00MiB at 1.2MiB/s ETA 00:01" "Processing: 75% complete" "Downloaded 3/10 tracks (30%)" """ try: # Pattern 1: Standard progress bar format with percentage and bar # Example: " 98%|█████████▊| 49.0/50.0 [00:10<00:00, 4.90s/s]" progress_match = re.search(r'(\d+(?:\.\d+)?)%\s*\|', line) if progress_match: progress = float(progress_match.group(1)) return min(100.0, max(0.0, progress)) # Pattern 2: Download progress format # Example: "[download] 99.9% of 5.00MiB at 1.2MiB/s ETA 00:01" alt_match = re.search(r'\[(?:download|fetching)\]\s*(\d+\.?\d*)%', line, re.IGNORECASE) if alt_match: progress = float(alt_match.group(1)) return min(100.0, max(0.0, progress)) # Pattern 3: Simple percentage in text # Example: "Processing: 75% complete" simple_match = re.search(r'(\d+(?:\.\d+)?)%\s*(?:complete|done|finished)', line, re.IGNORECASE) if simple_match: progress = float(simple_match.group(1)) return min(100.0, max(0.0, progress)) # Pattern 4: Track counting format for albums # Example: "Downloaded 3/10 tracks (30%)" or "Track 3 of 10 (30%)" track_match = re.search(r'(?:downloaded|track)\s*(\d+)(?:\s*of\s*|\s*/\s*)(\d+)', line, re.IGNORECASE) if track_match: current = float(track_match.group(1)) total = float(track_match.group(2)) if total > 0: progress = (current / total) * 100 return min(100.0, max(0.0, progress)) # Pattern 5: YouTube-dl style progress # Example: "[download] 45.2% of 12.34MiB at 1.23MiB/s ETA 00:08" ytdl_match = re.search(r'\[download\]\s*(\d+\.?\d*)%', line, re.IGNORECASE) if ytdl_match: progress = float(ytdl_match.group(1)) return min(100.0, max(0.0, progress)) # Pattern 6: Generic percentage anywhere in the line (as fallback) # Example: "Converting audio: 87%" generic_match = re.search(r'(\d+(?:\.\d+)?)%', line) if generic_match: progress = float(generic_match.group(1)) # Only return if it's a reasonable progress value if 0 <= progress <= 100: return progress # Look for downloading indicators and return incremental progress if 'downloading' in line.lower() or 'fetching' in line.lower(): # Return a small progress increment if we detect downloading activity return 5.0 # Look for processing/converting indicators if any(word in line.lower() for word in ['processing', 'converting', 'encoding', 'finalizing']): return 85.0 # Try to extract progress from ffmpeg output ffmpeg_match = re.search(r'size=\s*\d+\w+\s+time=(\d+):(\d+):(\d+\.\d+)', line) if ffmpeg_match: # This is a simplified example - in a real implementation, you'd need to know # the total duration to calculate percentage logger.debug(f"FFmpeg progress detected: {line.strip()}") return 90.0 except (ValueError, IndexError) as e: logger.warning(f"Error parsing progress from line: {line.strip()}. Error: {e}") return None @staticmethod def _sanitize_filename(filename: str, max_length: int = 255, replace: str = "_") -> str: """Sanitize a string to be used as a filename. Args: filename: Input string to be sanitized max_length: Maximum length of the resulting filename (default: 255) replace: Replacement string for invalid characters (default: "_") Returns: str: Sanitized filename that is safe for all filesystems Example: Input: "My/File:Name?.txt" Output: "My_File_Name_.txt" """ if not isinstance(filename, str): filename = str(filename) # Replace invalid characters invalid_chars = '<>:"/\\|?*\x00-\x1F\x7F' # Includes control characters for char in invalid_chars: filename = filename.replace(char, replace) # Handle reserved filenames on Windows reserved_names = { 'CON', 'PRN', 'AUX', 'NUL', 'COM1', 'COM2', 'COM3', 'COM4', 'COM5', 'COM6', 'COM7', 'COM8', 'COM9', 'LPT1', 'LPT2', 'LPT3', 'LPT4', 'LPT5', 'LPT6', 'LPT7', 'LPT8', 'LPT9' } # Check if filename is a reserved name (case-insensitive) base_name = os.path.splitext(filename)[0].upper() if base_name in reserved_names: filename = f"{filename}{replace}" # Append underscore to make it non-reserved # Remove leading/trailing whitespace and dots (Windows doesn't like these) filename = re.sub(r'^[\s.]+|[\s.]+$', '', filename) # Ensure filename is not empty after sanitization if not filename: filename = f"unnamed_file_{int(time.time())}" # Truncate if too long (reserve space for extension) name, ext = os.path.splitext(filename) if len(name) > max_length - len(ext): name = name[:max_length - len(ext) - 1] # Leave room for a single character filename = f"{name}{ext}" # Remove any remaining control characters filename = re.sub(r'[\x00-\x1f\x7f]', replace, filename) return filename def _broadcast_update(self, download_item: Dict[str, Any]) -> None: """Broadcast download status update to all connected clients. This method ensures that all status updates are properly formatted and broadcast to connected clients with appropriate error handling. Args: download_item: Dictionary containing download details with at least 'id' and 'status' keys Raises: ValueError: If download_item is missing required fields """ if not download_item or 'id' not in download_item or 'status' not in download_item: logger.error("Invalid download_item in _broadcast_update") return if not self.socketio: logger.warning("SocketIO not initialized, cannot broadcast update") return download_id = download_item['id'] status = download_item.get('status', 'unknown') try: # Ensure we have required fields with defaults download_item.setdefault('timestamp', self._get_timestamp()) download_item.setdefault('progress', 0) download_item.setdefault('message', '') # Log the update for debugging logger.debug(f"Broadcasting update for {download_id}: {status} - {download_item.get('message')}") # Emit specific events based on status if status == 'completed': self.socketio.emit('download_complete', download_item) logger.info(f"Download completed: {download_id}") elif status in ['failed', 'error', 'cancelled']: error_data = { 'id': download_id, 'error': download_item.get('message', 'Unknown error'), 'status': status, 'timestamp': download_item.get('timestamp') } self.socketio.emit('download_error', error_data) logger.warning(f"Download {status}: {download_id} - {error_data['error']}") else: # For all other statuses (queued, downloading, etc.) self.socketio.emit('download_progress', download_item) # Also emit a generic update for any listeners that might be interested update_data = { 'downloads': [download_item], 'timestamp': download_item['timestamp'] } self.socketio.emit('status', update_data) except Exception as e: logger.error(f"Error broadcasting update for {download_id}: {e}", exc_info=True) # Try to at least emit a basic error if we can try: self.socketio.emit('download_error', { 'id': download_id, 'error': f'Error processing update: {str(e)}', 'status': 'error', 'timestamp': self._get_timestamp() }) except Exception as inner_e: logger.critical(f"Critical error in error handling: {inner_e}", exc_info=True) @staticmethod def _get_timestamp() -> str: """Get current timestamp in ISO format.""" from datetime import datetime return datetime.utcnow().isoformat() + 'Z' def stop(self) -> None: """Stop the download service and clean up resources.""" self._stop_event.set() # Terminate active process if any if self.active_process and self.active_process.poll() is None: self.active_process.terminate() # Wait for thread to finish if self.worker_thread and self.worker_thread.is_alive(): self.worker_thread.join(timeout=5.0) # Clear queue while not self.download_queue.empty(): try: self.download_queue.get_nowait() self.download_queue.task_done() except queue.Empty: break def cleanup_old_downloads(self, days_old: int = 30) -> int: """Clean up download history older than specified days. Args: days_old: Number of days after which to clean up downloads Returns: int: Number of items cleaned up """ cutoff_time = int(time.time()) - (days_old * 24 * 60 * 60) cleaned_count = 0 # Create a list of keys to remove to avoid changing dict during iteration keys_to_remove = [] for download_id, download_info in self.download_history.items(): if download_info.get('timestamp', 0) < cutoff_time: keys_to_remove.append(download_id) # Remove old entries for key in keys_to_remove: del self.download_history[key] cleaned_count += 1 if cleaned_count > 0: logger.info(f"Cleaned up {cleaned_count} old download entries") # Emit updated status self.socketio.emit("status", {"downloads": list(self.download_history.values())}) return cleaned_count def clear_completed_downloads(self) -> int: """Clear all completed downloads from history. Returns: int: Number of items cleared """ keys_to_remove = [] for download_id, download_info in self.download_history.items(): if download_info.get('stage') in ['completed', 'error']: keys_to_remove.append(download_id) # Remove completed entries for key in keys_to_remove: del self.download_history[key] cleared_count = len(keys_to_remove) if cleared_count > 0: logger.info(f"Cleared {cleared_count} completed download entries") # Emit updated status self.socketio.emit("status", {"downloads": list(self.download_history.values())}) return cleared_count def get_download_stats(self) -> Dict[str, Any]: """Get download statistics. Returns: Dict with download statistics """ total = len(self.download_history) completed = sum(1 for d in self.download_history.values() if d.get('stage') == 'completed') failed = sum(1 for d in self.download_history.values() if d.get('stage') == 'error') pending = sum(1 for d in self.download_history.values() if d.get('stage') == 'queued') downloading = sum(1 for d in self.download_history.values() if d.get('stage') == 'downloading') return { 'total': total, 'completed': completed, 'failed': failed, 'pending': pending, 'downloading': downloading, 'success_rate': (completed / total * 100) if total > 0 else 0 } def get_stats(self) -> Dict[str, Any]: """Alias for get_download_stats for compatibility. Returns: Dict with download statistics """ return self.get_download_stats() def cancel_active_download(self) -> bool: """Cancel the currently active download. Returns: True if a download was cancelled, False if none was active """ active_downloads = [d for d in self.download_history.values() if d.get('stage') == 'downloading'] if active_downloads: for download in active_downloads: download_id = download.get('id', '') if download_id: return self.cancel_download(download_id) return False def cancel_pending_downloads(self) -> int: """Cancel all pending downloads. Returns: Number of downloads cancelled """ pending_downloads = [d for d in self.download_history.values() if d.get('stage') == 'queued'] cancelled_count = 0 for download in pending_downloads: download_id = download.get('id', '') if download_id and self.cancel_download(download_id): cancelled_count += 1 return cancelled_count def cancel_specific_item(self, item_id: str) -> bool: """Cancel a specific download item. Args: item_id: The ID of the item to cancel Returns: True if the item was cancelled successfully """ return self.cancel_download(item_id)