import os import sys import logging import subprocess import threading import time from typing import Dict, Any, Optional logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s") logger = logging.getLogger(__name__) class DownloadService: """Handles downloading of tracks, albums, and playlists from Spotify.""" def __init__(self, config, playlist_manager, socketio, download_queue, download_history): """Initialize the download service. Args: config: Configuration object playlist_manager: Playlist manager instance socketio: SocketIO instance for real-time updates download_queue: Queue for managing downloads download_history: Dictionary to track 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_process = None self.is_running = True # Start the download processor in a separate thread self.processor_thread = threading.Thread(target=self._process_downloads, daemon=True) self.processor_thread.start() def add_item_to_queue(self, item_data: Dict[str, Any]) -> None: """Add an item to the download queue. Args: item_data: Dictionary containing item details (url, type, name, artist, etc.) """ try: item_id = item_data.get('url', str(time.time())) download_info = { 'id': item_id, 'status': 'queued', 'progress': 0, 'name': item_data.get('name', 'Unknown'), 'artist': item_data.get('artist', 'Unknown'), 'type': item_data.get('type', 'track'), 'url': item_data.get('url'), 'added_at': time.time(), 'started_at': None, 'completed_at': None, 'error': None } self.download_queue.put(download_info) self.download_history[item_id] = download_info # Notify clients of the update self._notify_status_update() logger.info(f"Added to queue: {item_data.get('name')}") except Exception as e: logger.error(f"Error adding to queue: {str(e)}") raise def _process_downloads(self) -> None: """Process items from the download queue.""" while self.is_running: try: if self.download_queue.empty(): time.sleep(1) continue # Get the next item from the queue item = self.download_queue.get() item_id = item['id'] # Update status item['status'] = 'downloading' item['started_at'] = time.time() self.download_history[item_id] = item self._notify_status_update() try: # Determine output directory based on item type # Always use main download directory; templates are handled by spotdl internally output_dir = self.config.download_dir os.makedirs(output_dir, exist_ok=True) # Build the spotdl command cmd = [ 'spotdl', '--output', output_dir, '--temp-dir', self.config.temp_dir, '--format', 'mp3', '--bitrate', '320k', '--overwrite', 'skip', item['url'] ] logger.info(f"Starting download: {' '.join(cmd)}") # Start the download process self.active_process = subprocess.Popen( cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, bufsize=1, universal_newlines=True ) # Process output in real-time for line in iter(self.active_process.stdout.readline, ''): if not line.strip(): continue # Parse progress from spotdl output if '%' in line: try: progress = int(line.split('%')[0].strip()) item['progress'] = min(progress, 100) self.download_history[item_id] = item self._notify_status_update() except (ValueError, IndexError): pass logger.debug(f"spotdl: {line.strip()}") # Wait for the process to complete return_code = self.active_process.wait() if return_code == 0: item['status'] = 'completed' item['progress'] = 100 logger.info(f"Download completed: {item['name']}") # Add to playlist if configured if hasattr(self.config, 'auto_add_to_playlist') and self.config.auto_add_to_playlist: self.playlist_manager.add_to_playlist( self.config.auto_add_to_playlist, { 'name': item['name'], 'artist': item['artist'], 'url': item['url'], 'downloaded_at': time.time(), 'path': os.path.join( output_dir, f"{item['artist']} - {item['name']}.mp3" ) } ) # Trigger Jellyfin scan if configured if hasattr(self.config, 'jellyfin_address') and hasattr(self.config, 'jellyfin_api_key'): self._trigger_jellyfin_scan() else: item['status'] = 'failed' item['error'] = f"Process returned non-zero exit code: {return_code}" logger.error(f"Download failed: {item['name']} - {item['error']}") except Exception as e: item['status'] = 'failed' item['error'] = str(e) logger.error(f"Error during download: {str(e)}") # Update completion time and status item['completed_at'] = time.time() self.download_history[item_id] = item self._notify_status_update() # Mark task as done self.download_queue.task_done() except Exception as e: logger.error(f"Error in download processor: {str(e)}") time.sleep(5) # Prevent tight loop on errors def _notify_status_update(self) -> None: """Notify clients of a status update.""" try: self.socketio.emit('update_status', { 'queue': list(self.download_queue.queue), 'history': list(self.download_history.values()), 'active': [item for item in self.download_history.values() if item.get('status') == 'downloading'] }) except Exception as e: logger.error(f"Error notifying status update: {str(e)}") def _trigger_jellyfin_scan(self) -> None: """Trigger a Jellyfin library scan.""" if not hasattr(self.config, 'jellyfin_address') or not hasattr(self.config, 'jellyfin_api_key'): return try: import requests headers = { 'X-Emby-Token': self.config.jellyfin_api_key, 'Accept': 'application/json' } response = requests.post( f"{self.config.jellyfin_address.rstrip('/')}/Library/Refresh", headers=headers, timeout=10 ) response.raise_for_status() logger.info("Triggered Jellyfin library scan") except Exception as e: logger.error(f"Failed to trigger Jellyfin scan: {str(e)}") def cancel_download(self, item_id: str) -> bool: """Cancel a queued or active download. Args: item_id: ID of the item to cancel Returns: bool: True if cancelled, False if not found """ # Check if it's the currently downloading item if self.active_process and self.active_process.poll() is None: current_item = next((item for item in self.download_history.values() if item.get('status') == 'downloading'), None) if current_item and current_item.get('id') == item_id: try: self.active_process.terminate() current_item['status'] = 'cancelled' current_item['completed_at'] = time.time() self._notify_status_update() return True except Exception as e: logger.error(f"Error cancelling download: {str(e)}") return False # Check if it's in the queue with self.download_queue.mutex: for i, item in enumerate(list(self.download_queue.queue)): if item.get('id') == item_id: # Remove from queue self.download_queue.queue.remove(item) # Update status in history if item_id in self.download_history: self.download_history[item_id]['status'] = 'cancelled' self.download_history[item_id]['completed_at'] = time.time() self._notify_status_update() return True return False def get_status(self) -> Dict[str, Any]: """Get the current download status. Returns: Dict containing queue, history, and active downloads """ return { 'queue': list(self.download_queue.queue), 'history': list(self.download_history.values()), 'active': [item for item in self.download_history.values() if item.get('status') == 'downloading'] } def shutdown(self) -> None: """Shut down the download service.""" self.is_running = False if self.active_process and self.active_process.poll() is None: self.active_process.terminate() if self.processor_thread.is_alive(): self.processor_thread.join(timeout=5) logger.info("Download service shutdown complete")