import asyncio import fnmatch import os import logging import shutil from typing import Any, Dict, List, Optional, Sequence, Set from abc import ABC, abstractmethod from ..utils.utils import calculate_relative_path_for_model, remove_empty_dirs from ..utils.constants import AUTO_ORGANIZE_BATCH_SIZE, MODEL_FILE_EXTENSIONS from ..services.settings_manager import get_settings_manager from ..services.model_lifecycle_service import _require_path_in_library_roots from ..services.pending_delete_service import PENDING_DELETE_DIR_NAME logger = logging.getLogger(__name__) class ProgressCallback(ABC): """Abstract callback interface for progress reporting""" @abstractmethod async def on_progress(self, progress_data: Dict[str, Any]) -> None: """Called when progress is updated""" pass class AutoOrganizeResult: """Result object for auto-organize operations""" def __init__(self): self.total: int = 0 self.processed: int = 0 self.success_count: int = 0 self.failure_count: int = 0 self.skipped_count: int = 0 self.operation_type: str = 'unknown' self.cleanup_counts: Dict[str, int] = {} self.results: List[Dict[str, Any]] = [] self.results_truncated: bool = False self.sample_results: List[Dict[str, Any]] = [] self.is_flat_structure: bool = False self.status: str = 'success' def to_dict(self) -> Dict[str, Any]: """Convert result to dictionary""" result: Dict[str, Any] = { 'success': self.status != 'error', 'status': self.status, 'message': f'Auto-organize {self.operation_type} completed: {self.success_count} moved, {self.skipped_count} skipped, {self.failure_count} failed out of {self.total} total', 'summary': { 'total': self.total, 'success': self.success_count, 'skipped': self.skipped_count, 'failures': self.failure_count, 'organization_type': 'flat' if self.is_flat_structure else 'structured', 'cleaned_dirs': self.cleanup_counts, 'operation_type': self.operation_type } } if self.results_truncated: result['results_truncated'] = True result['sample_results'] = self.sample_results else: result['results'] = self.results return result class ModelFileService: """Service for handling model file operations and organization""" def __init__(self, scanner, model_type: str): """Initialize the service Args: scanner: Model scanner instance model_type: Type of model (e.g., 'lora', 'checkpoint') """ self.scanner = scanner self.model_type = model_type def get_model_roots(self) -> List[str]: """Get model root directories""" return self.scanner.get_model_roots() async def auto_organize_models( self, file_paths: Optional[List[str]] = None, progress_callback: Optional[ProgressCallback] = None, exclusion_patterns: Optional[Sequence[str]] = None, ) -> AutoOrganizeResult: """Auto-organize models based on current settings Args: file_paths: Optional list of specific file paths to organize. If None, organizes all models. progress_callback: Optional callback for progress updates Returns: AutoOrganizeResult object with operation results """ result = AutoOrganizeResult() source_directories: Set[str] = set() self.scanner.reset_cancellation() try: # Get all models from cache cache = await self.scanner.get_cached_data() all_models = cache.raw_data settings_manager = get_settings_manager() normalized_exclusions = settings_manager.normalize_auto_organize_exclusions( exclusion_patterns if exclusion_patterns is not None else settings_manager.get_auto_organize_exclusions() ) # Filter models if specific file paths are provided if file_paths: all_models = [model for model in all_models if model.get('file_path') in file_paths] result.operation_type = 'bulk' else: result.operation_type = 'all' model_roots = self.get_model_roots() if not model_roots: raise ValueError('No model roots configured') if normalized_exclusions: all_models = [ model for model in all_models if not self._should_exclude_model( model.get('file_path'), normalized_exclusions, model_roots ) ] # Check if flat structure is configured for this model type settings_manager = get_settings_manager() path_template = settings_manager.get_download_path_template(self.model_type) result.is_flat_structure = not path_template # Initialize tracking result.total = len(all_models) # Send initial progress if progress_callback: await progress_callback.on_progress({ 'type': 'auto_organize_progress', 'status': 'started', 'total': result.total, 'processed': 0, 'success': 0, 'failures': 0, 'skipped': 0, 'operation_type': result.operation_type }) if result.total == 0: if progress_callback: await asyncio.sleep(0.1) payload = { 'type': 'auto_organize_progress', 'total': 0, 'processed': 0, 'success': 0, 'failures': 0, 'skipped': 0, 'operation_type': result.operation_type } await progress_callback.on_progress({**payload, 'status': 'processing'}) await progress_callback.on_progress({ **payload, 'status': 'cleaning', 'message': 'Cleaning up empty directories...' }) result.cleanup_counts = {} await progress_callback.on_progress({ **payload, 'status': 'completed', 'cleanup': result.cleanup_counts }) return result # Process models in batches await self._process_models_in_batches( all_models, model_roots, result, progress_callback, source_directories # Pass the set to track source directories ) if self.scanner.is_cancelled(): result.status = 'cancelled' if progress_callback: await progress_callback.on_progress({ 'type': 'auto_organize_progress', 'status': 'cancelled', 'total': result.total, 'processed': result.processed, 'success': result.success_count, 'failures': result.failure_count, 'skipped': result.skipped_count, 'operation_type': result.operation_type }) return result # Send cleanup progress if progress_callback: await progress_callback.on_progress({ 'type': 'auto_organize_progress', 'status': 'cleaning', 'total': result.total, 'processed': result.processed, 'success': result.success_count, 'failures': result.failure_count, 'skipped': result.skipped_count, 'message': 'Cleaning up empty directories...', 'operation_type': result.operation_type }) # Clean up empty directories - only in affected directories for bulk operations cleanup_paths = list(source_directories) if result.operation_type == 'bulk' else model_roots result.cleanup_counts = await self._cleanup_empty_directories(cleanup_paths) # Send completion message if progress_callback: await progress_callback.on_progress({ 'type': 'auto_organize_progress', 'status': 'completed', 'total': result.total, 'processed': result.processed, 'success': result.success_count, 'failures': result.failure_count, 'skipped': result.skipped_count, 'cleanup': result.cleanup_counts, 'operation_type': result.operation_type }) return result except Exception as e: logger.error(f"Error in auto_organize_models: {e}", exc_info=True) # Send error message if progress_callback: await progress_callback.on_progress({ 'type': 'auto_organize_progress', 'status': 'error', 'error': str(e), 'operation_type': result.operation_type }) raise e async def _process_models_in_batches( self, all_models: List[Dict[str, Any]], model_roots: List[str], result: AutoOrganizeResult, progress_callback: Optional[ProgressCallback], source_directories: Optional[Set[str]] = None ) -> None: """Process models in batches to avoid overwhelming the system""" for i in range(0, result.total, AUTO_ORGANIZE_BATCH_SIZE): if self.scanner.is_cancelled(): logger.info(f"{self.model_type.capitalize()} File Service: Auto-organize cancelled by user") break batch = all_models[i:i + AUTO_ORGANIZE_BATCH_SIZE] for model in batch: if self.scanner.is_cancelled(): break await self._process_single_model(model, model_roots, result, source_directories) result.processed += 1 # Send progress update after each batch if progress_callback: await progress_callback.on_progress({ 'type': 'auto_organize_progress', 'status': 'processing', 'total': result.total, 'processed': result.processed, 'success': result.success_count, 'failures': result.failure_count, 'skipped': result.skipped_count, 'operation_type': result.operation_type }) # Small delay between batches await asyncio.sleep(0.1) async def _process_single_model( self, model: Dict[str, Any], model_roots: List[str], result: AutoOrganizeResult, source_directories: Optional[Set[str]] = None ) -> None: """Process a single model for organization""" try: file_path = model.get('file_path') model_name = model.get('model_name', 'Unknown') if not file_path: self._add_result(result, model_name, False, "No file path found") result.failure_count += 1 return # Find which model root this file belongs to current_root = self._find_model_root(file_path, model_roots) if not current_root: self._add_result(result, model_name, False, "Model file not found in any configured root directory") result.failure_count += 1 return # Determine target directory target_dir = await self._calculate_target_directory( model, current_root, result.is_flat_structure ) if target_dir is None: self._add_result(result, model_name, False, "Skipped - insufficient metadata for organization") result.skipped_count += 1 return current_dir = os.path.dirname(file_path) # Skip if already in correct location if current_dir.replace(os.sep, '/') == target_dir.replace(os.sep, '/'): result.skipped_count += 1 return # Check for conflicts file_name = os.path.basename(file_path) target_file_path = os.path.join(target_dir, file_name) if os.path.exists(target_file_path): self._add_result(result, model_name, False, f"Target file already exists: {target_file_path}") result.failure_count += 1 return # Store the source directory for potential cleanup if source_directories is not None: source_directories.add(current_dir) # Perform the move success = await self.scanner.move_model(file_path, target_dir) if success: result.success_count += 1 else: self._add_result(result, model_name, False, "Failed to move model") result.failure_count += 1 except Exception as e: logger.error(f"Error processing model {model.get('model_name', 'Unknown')}: {e}", exc_info=True) self._add_result(result, model.get('model_name', 'Unknown'), False, f"Error: {str(e)}") result.failure_count += 1 def _find_model_root(self, file_path: str, model_roots: List[str]) -> Optional[str]: """Find which model root the file belongs to""" for root in model_roots: # Normalize paths for comparison normalized_root = os.path.normpath(root).replace(os.sep, '/') normalized_file = os.path.normpath(file_path).replace(os.sep, '/') if normalized_file.startswith(normalized_root): return root return None def _should_exclude_model( self, file_path: Optional[str], patterns: Sequence[str], model_roots: Sequence[str], ) -> bool: if not file_path or not patterns: return False normalized_path = os.path.normpath(file_path).replace(os.sep, '/') filename = os.path.basename(normalized_path) relative_path = None if model_roots: root = self._find_model_root(file_path, list(model_roots)) if root: normalized_root = os.path.normpath(root) try: relative = os.path.relpath(file_path, normalized_root) except ValueError: relative = None if relative is not None: relative_path = relative.replace(os.sep, '/') for pattern in patterns: if fnmatch.fnmatch(filename, pattern): return True if relative_path and fnmatch.fnmatch(relative_path, pattern): return True if fnmatch.fnmatch(normalized_path, pattern): return True return False async def _calculate_target_directory( self, model: Dict[str, Any], current_root: str, is_flat_structure: bool ) -> Optional[str]: """Calculate the target directory for a model""" if is_flat_structure: file_path = model.get('file_path') if not isinstance(file_path, str): return None current_dir = os.path.dirname(file_path) # Check if already in root directory if os.path.normpath(current_dir) == os.path.normpath(current_root): return None # Signal to skip return current_root else: # Calculate new relative path based on settings new_relative_path = calculate_relative_path_for_model(model, self.model_type) if not new_relative_path: return None # Signal to skip return os.path.join(current_root, new_relative_path).replace(os.sep, '/') def _add_result( self, result: AutoOrganizeResult, model_name: str, success: bool, message: str ) -> None: """Add a result entry if under the limit""" if len(result.results) < 100: # Limit detailed results result.results.append({ "model": model_name, "success": success, "message": message }) elif len(result.results) == 100: # Mark as truncated and save sample result.results_truncated = True result.sample_results = result.results[:50] async def _cleanup_empty_directories(self, paths: List[str]) -> Dict[str, int]: """Clean up empty directories after organizing Args: paths: List of paths to check for empty directories Returns: Dictionary with counts of removed directories by root path """ cleanup_counts = {} for path in paths: removed = remove_empty_dirs(path) cleanup_counts[path] = removed return cleanup_counts class ModelMoveService: """Service for handling individual model moves""" def __init__(self, scanner, model_type: str): """Initialize the service Args: scanner: Model scanner instance model_type: Type of model (e.g., 'lora', 'checkpoint') """ self.scanner = scanner self.model_type = model_type async def create_folder(self, folder_path: str) -> Dict[str, Any]: """Create a directory inside the model library roots. Args: folder_path: Absolute path of the directory to create (business path — symlinks are not resolved) Returns: Dictionary with success flag, the created path and the library-relative folder name used by folder trees. """ try: if not folder_path or not str(folder_path).strip(): return {"success": False, "error": "Folder path is required"} _require_path_in_library_roots(folder_path, self.scanner, label="Folder path") absolute_path = os.path.abspath(folder_path) already_exists = os.path.isdir(absolute_path) os.makedirs(absolute_path, exist_ok=True) relative_folder = self._calculate_relative_folder(absolute_path) if relative_folder: await self.scanner.add_known_folder(relative_folder) return { "success": True, "folder_path": absolute_path.replace(os.sep, "/"), "folder": relative_folder, "created": not already_exists, } except ValueError as exc: return {"success": False, "error": str(exc)} except Exception as exc: logger.error(f"Error creating folder: {exc}", exc_info=True) return {"success": False, "error": str(exc)} def _calculate_relative_folder(self, absolute_path: str) -> str: """Return the library-relative folder for an absolute directory path.""" normalized = os.path.abspath(absolute_path) for root in self.scanner.get_model_roots(): abs_root = os.path.abspath(root) try: rel = os.path.relpath(normalized, abs_root) except ValueError: continue if rel == ".": return "" if not rel.startswith(".."): return rel.replace(os.sep, "/") return "" async def delete_folder(self, folder_path: str, dry_run: bool = False) -> Dict[str, Any]: """Delete a model-free directory inside the model library roots. Only directories whose subtree holds no model weight files can be removed: a folder-level cascade would bypass the per-model lifecycle bookkeeping (metadata sidecars, previews, cache entries, pending-delete staging and recipe references), so it is deliberately refused. Leftover non-model files (stray previews, sidecars, ``.bak`` files) are reported in the manifest before they are removed. Args: folder_path: Absolute path of the directory to remove (business path — symlinks are not resolved) dry_run: When true, only report what would be removed Returns: Dictionary with the success flag plus a removal manifest (``model_count``/``file_count``/``dir_count``/``symlink_count``/ ``total_bytes``/``restorable``) on success. """ try: if not folder_path or not str(folder_path).strip(): return {"success": False, "error": "Folder path is required"} _require_path_in_library_roots(folder_path, self.scanner, label="Folder path") absolute_path = os.path.abspath(folder_path) if os.path.islink(absolute_path): # shutil.rmtree refuses symlinked roots, and silently deleting # the link (leaving the real directory behind) is a separate # decision we do not make here. return { "success": False, "error": "Symlinked folders cannot be deleted", } if not os.path.isdir(absolute_path): return {"success": False, "error": "Folder no longer exists"} if self._is_model_root(absolute_path): return { "success": False, "error": "The library root itself cannot be deleted", } manifest = self._collect_folder_manifest(absolute_path) if manifest["pending_delete_job"]: return { "success": False, "code": "busy", "error": ( "A staged delete is still pending inside this folder; " "wait for the undo window to expire" ), "manifest": manifest, } if manifest["model_count"] > 0: return { "success": False, "code": "not_empty", "error": ( f"Folder still contains {manifest['model_count']} model " "file(s); delete or move them first" ), "manifest": manifest, } relative_folder = self._calculate_relative_folder(absolute_path) if dry_run: return { "success": True, "dry_run": True, "folder_path": absolute_path.replace(os.sep, "/"), "folder": relative_folder, **manifest, } shutil.rmtree(absolute_path) await self._forget_folder(relative_folder) return { "success": True, "dry_run": False, "folder_path": absolute_path.replace(os.sep, "/"), "folder": relative_folder, **manifest, } except ValueError as exc: return {"success": False, "error": str(exc)} except Exception as exc: logger.error(f"Error deleting folder: {exc}", exc_info=True) return {"success": False, "error": str(exc)} def _is_model_root(self, absolute_path: str) -> bool: """Return True when the path *is* one of the configured library roots.""" normalized = os.path.normpath(absolute_path) for root in self.scanner.get_model_roots(): if os.path.normpath(os.path.abspath(root)) == normalized: return True return False @staticmethod def _is_model_file(file_name: str) -> bool: """Return True when the file name carries a model weight extension.""" return os.path.splitext(file_name)[1].lower() in MODEL_FILE_EXTENSIONS def _collect_folder_manifest(self, absolute_path: str) -> Dict[str, Any]: """Describe everything a recursive delete of *absolute_path* removes. Walking is intentional: the scanner cache can be stale, and a model file that appeared on disk since the last scan must still block the delete. Symbolic links are never followed (``os.walk`` default) and are counted separately — ``shutil.rmtree`` unlinks them without touching their targets. """ model_count = 0 file_count = 0 dir_count = 0 symlink_count = 0 total_bytes = 0 pending_delete_job = False for dirpath, dirnames, filenames in os.walk(absolute_path): if PENDING_DELETE_DIR_NAME in dirnames: pending_delete_job = True for name in dirnames: if os.path.islink(os.path.join(dirpath, name)): symlink_count += 1 else: dir_count += 1 for name in filenames: full_path = os.path.join(dirpath, name) if os.path.islink(full_path): symlink_count += 1 continue if self._is_model_file(name): model_count += 1 else: file_count += 1 try: total_bytes += os.path.getsize(full_path) except OSError: # pragma: no cover - defensive pass return { "model_count": model_count, "file_count": file_count, "dir_count": dir_count, "symlink_count": symlink_count, "total_bytes": total_bytes, "pending_delete_job": pending_delete_job, # A truly empty directory is the only case an "undo" can restore by # simply recreating it; a folder holding stray files is gone for good. "restorable": ( model_count == 0 and file_count == 0 and dir_count == 0 and symlink_count == 0 ), } async def _forget_folder(self, relative_folder: str) -> None: """Drop a removed directory from the scanner's folder/cache records.""" if not relative_folder: return remove_known_folder = getattr(self.scanner, "remove_known_folder", None) if callable(remove_known_folder): await remove_known_folder(relative_folder) async def rename_folder(self, folder_path: str, new_name: str) -> Dict[str, Any]: """Rename a directory inside the model library roots. Unlike :meth:`delete_folder` this works on folders that hold models. A rename keeps every file, so no per-model lifecycle step is bypassed: the directory is renamed on disk and the affected folder, cache, hash index and metadata-sidecar records are re-keyed onto the new prefix by the scanner. Args: folder_path: Absolute path of the directory to rename (business path — symlinks are not resolved) new_name: New leaf name; a single path segment, not a path Returns: Dictionary with the success flag, the previous/next library-relative folder names and whether the directory actually moved. """ try: if not folder_path or not str(folder_path).strip(): return {"success": False, "error": "Folder path is required"} new_name = str(new_name or "").strip() if not new_name: return {"success": False, "error": "New folder name is required"} if new_name in (".", "..") or any( char in new_name for char in '/\\:*?"<>|' ): return {"success": False, "error": "Invalid characters in folder name"} _require_path_in_library_roots(folder_path, self.scanner, label="Folder path") absolute_path = os.path.abspath(folder_path) if os.path.islink(absolute_path): return { "success": False, "error": "Symlinked folders cannot be renamed", } if not os.path.isdir(absolute_path): return {"success": False, "error": "Folder no longer exists"} if self._is_model_root(absolute_path): return { "success": False, "error": "The library root itself cannot be renamed", } previous_relative = self._calculate_relative_folder(absolute_path) target = os.path.join(os.path.dirname(absolute_path), new_name) if os.path.normpath(target) == os.path.normpath(absolute_path): return { "success": True, "renamed": False, "folder": previous_relative, "previous_folder": previous_relative, "folder_path": absolute_path.replace(os.sep, "/"), } if os.path.exists(target): return { "success": False, "code": "target_exists", "error": f"A folder named \"{new_name}\" already exists here", } # A staging manifest records absolute original/staged paths, so # moving a folder that holds one would break its undo and purge. if self._has_pending_delete_job(absolute_path): return { "success": False, "code": "busy", "error": ( "A staged delete is still pending inside this folder; " "wait for the undo window to expire" ), } os.rename(absolute_path, target) new_relative = self._calculate_relative_folder(target) await self._rename_folder_records( previous_relative, new_relative, absolute_path, target ) return { "success": True, "renamed": True, "folder": new_relative, "previous_folder": previous_relative, "folder_path": target.replace(os.sep, "/"), } except ValueError as exc: return {"success": False, "error": str(exc)} except Exception as exc: logger.error(f"Error renaming folder: {exc}", exc_info=True) return {"success": False, "error": str(exc)} @staticmethod def _has_pending_delete_job(absolute_path: str) -> bool: """Return True when a staged-delete batch lives inside the subtree.""" for _dirpath, dirnames, _filenames in os.walk(absolute_path): if PENDING_DELETE_DIR_NAME in dirnames: return True return False async def _rename_folder_records( self, previous_relative: str, new_relative: str, previous_path: str, new_path: str, ) -> None: """Hand the rename to the scanner so folder/cache records follow it.""" if not previous_relative or not new_relative: return rename_known_folder = getattr(self.scanner, "rename_known_folder", None) if callable(rename_known_folder): await rename_known_folder( previous_relative, new_relative, previous_path=previous_path, new_path=new_path, ) async def move_model(self, file_path: str, target_path: str, use_default_paths: bool = False) -> Dict[str, Any]: """Move a single model file Args: file_path: Source file path target_path: Target directory path (used as root if use_default_paths is True) use_default_paths: Whether to use default path template for organization Returns: Dictionary with move result """ try: _require_path_in_library_roots(file_path, self.scanner, label="Source path") _require_path_in_library_roots(target_path, self.scanner, label="Target path") if use_default_paths: # Find the model in cache to get metadata cache = await self.scanner.get_cached_data() model_data = next((m for m in cache.raw_data if m.get('file_path') == file_path), None) if model_data: from ..utils.utils import calculate_relative_path_for_model relative_path = calculate_relative_path_for_model(model_data, self.model_type) if relative_path: target_path = os.path.join(target_path, relative_path).replace(os.sep, '/') elif not get_settings_manager().get_download_path_template(self.model_type): # Flat structure, target_path remains the root pass else: # Could not calculate relative path (e.g. missing metadata) # Fallback to manual target_path or skip? pass source_dir = os.path.dirname(file_path) if os.path.normpath(source_dir) == os.path.normpath(target_path): logger.info(f"Source and target directories are the same: {source_dir}") return { 'success': True, 'message': 'Source and target directories are the same', 'original_file_path': file_path, 'new_file_path': file_path } move_result = await self.scanner.move_model(file_path, target_path) if move_result: new_file_path = move_result.get("new_path") cache_entry = move_result.get("cache_entry") return { 'success': True, 'original_file_path': file_path, 'new_file_path': new_file_path, 'cache_entry': cache_entry } else: return { 'success': False, 'error': 'Failed to move model', 'original_file_path': file_path, 'new_file_path': None } except Exception as e: logger.error(f"Error moving model: {e}", exc_info=True) return { 'success': False, 'error': str(e), 'original_file_path': file_path, 'new_file_path': None } async def move_models_bulk(self, file_paths: List[str], target_path: str, use_default_paths: bool = False) -> Dict[str, Any]: """Move multiple model files Args: file_paths: List of source file paths target_path: Target directory path (used as root if use_default_paths is True) use_default_paths: Whether to use default path template for organization Returns: Dictionary with bulk move results """ try: results = [] self.scanner.reset_cancellation() for file_path in file_paths: if self.scanner.is_cancelled(): logger.info(f"{self.model_type.capitalize()} Move Service: Bulk move cancelled by user") break result = await self.move_model(file_path, target_path, use_default_paths=use_default_paths) results.append({ "original_file_path": file_path, "new_file_path": result.get('new_file_path'), "success": result['success'], "message": result.get('message', result.get('error', 'Unknown')), "cache_entry": result.get('cache_entry') }) success_count = sum(1 for r in results if r["success"]) failure_count = len(results) - success_count return { 'success': True, 'message': f'Moved {success_count} of {len(file_paths)} models', 'results': results, 'success_count': success_count, 'failure_count': failure_count } except Exception as e: logger.error(f"Error moving models in bulk: {e}", exc_info=True) return { 'success': False, 'error': str(e), 'results': [], 'success_count': 0, 'failure_count': len(file_paths) }