diff --git a/py/services/model_scanner.py b/py/services/model_scanner.py index c884c9de..bf5d430e 100644 --- a/py/services/model_scanner.py +++ b/py/services/model_scanner.py @@ -161,6 +161,7 @@ class ModelScanner: self._persistent_cache = get_persistent_cache() self._name_display_mode = self._resolve_name_display_mode() self._cancel_requested = False # Flag for cancellation + self._move_locks: Dict[str, asyncio.Lock] = {} # Per-source-file move locks self._autov3_backfill_scheduled = False # One-time AutoV3 backfill trigger per process # Guard against concurrent all-folders backfill walks (cold fallback # for persisted snapshots that predate folder recording). @@ -2408,18 +2409,31 @@ class ModelScanner: async def move_model(self, source_path: str, target_path: str) -> Optional[Dict[str, Any]]: """Move a model and its associated files to a new location - + Args: source_path: Original file path target_path: Target directory path - + Returns: Optional[str]: New file path if successful, None if failed """ + source_path = source_path.replace(os.sep, '/') + target_path = target_path.replace(os.sep, '/') + + # Serialize moves per source file: concurrent requests for the same + # model (auto-organize racing a manual move, duplicate clicks) must + # not interleave, or the second mover hits a missing source file. + lock_key = os.path.normcase(os.path.abspath(os.path.normpath(source_path))) + lock = self._move_locks.setdefault(lock_key, asyncio.Lock()) + try: + async with lock: + return await self._move_model_locked(source_path, target_path) + finally: + if not lock.locked(): + self._move_locks.pop(lock_key, None) + + async def _move_model_locked(self, source_path: str, target_path: str) -> Optional[Dict[str, Any]]: try: - source_path = source_path.replace(os.sep, '/') - target_path = target_path.replace(os.sep, '/') - file_ext = os.path.splitext(source_path)[1] if not file_ext or file_ext.lower() not in self.file_extensions: @@ -2450,10 +2464,20 @@ class ModelScanner: if final_filename != f"{base_name}{file_ext}": logger.info(f"Renamed {base_name}{file_ext} to {final_filename} to avoid filename conflict") - real_source = os.path.realpath(source_path) - real_target = os.path.realpath(target_file) - - shutil.move(real_source, real_target) + # Business paths (abspath, symlinks NOT resolved) per project + # convention: file mutations must operate on the paths as they + # appear under the configured model roots. + move_source = os.path.abspath(source_path) + move_target = os.path.abspath(target_file) + + if not os.path.exists(move_source): + # The source is gone — typically a previous move already + # succeeded but the cache/metadata were left pointing at the + # old path. Repair that state instead of failing. + natural_target = os.path.join(target_path, f"{base_name}{file_ext}").replace(os.sep, '/') + return await self._reconcile_already_moved(source_path, [target_file, natural_target]) + + shutil.move(move_source, move_target) # Move all associated files with the same base name source_metadata = None @@ -2532,7 +2556,70 @@ class ModelScanner: except Exception as e: logger.error(f"Error moving model: {e}", exc_info=True) return None - + + async def _reconcile_already_moved(self, source_path: str, target_candidates: List[str]) -> Optional[Dict[str, Any]]: + """Repair state when a move's source file is already gone. + + A previous move may have relocated the file while the cache/metadata + still point at the old path (crash mid-move, concurrent request, or + external tools). If the model is found at its new location, update + the cache and metadata to match reality instead of failing. + """ + candidates: List[str] = [] + source_hash = self.get_hash_by_path(source_path) + if source_hash: + indexed_path = self.get_path_by_hash(source_hash) + if indexed_path: + candidates.append(indexed_path) + candidates.extend(target_candidates) + + for candidate in candidates: + if not candidate or os.path.normpath(candidate) == os.path.normpath(source_path): + continue + if not os.path.exists(os.path.abspath(candidate)): + continue + + new_path = candidate.replace(os.sep, '/') + logger.info( + f"Move source {source_path} no longer exists; the model is already " + f"at {new_path}. Reconciling cache and metadata." + ) + + cache = await self.get_cached_data() + existing_at_target = next((item for item in cache.raw_data if item['file_path'] == new_path), None) + if existing_at_target is not None: + # Cache already tracks the moved file (a previous move updated + # it); just drop the stale source entry without appending a + # duplicate. + await self.update_single_model_cache(source_path, new_path, None) + return {"new_path": new_path, "cache_entry": existing_at_target} + + metadata = None + metadata_path = get_metadata_path(new_path) + if os.path.exists(metadata_path): + metadata = await self._update_metadata_paths(metadata_path, new_path) + + if metadata is None: + # No sidecar at the new location — reuse the stale cache entry + # so the model card keeps its data under the corrected path. + existing_item = next((item for item in cache.raw_data if item['file_path'] == source_path), None) + if existing_item: + metadata = dict(existing_item) + metadata['file_path'] = new_path + metadata['file_name'] = os.path.splitext(os.path.basename(new_path))[0] + + update_result = await self.update_single_model_cache(source_path, new_path, metadata, recalculate_type=True) + return { + "new_path": new_path, + "cache_entry": update_result if isinstance(update_result, dict) else None, + } + + logger.error( + f"Cannot move model: source file not found: {source_path} " + f"(already moved or deleted outside LoRA Manager?)" + ) + return None + async def _update_metadata_paths(self, metadata_path: str, model_path: str) -> Optional[Dict[str, Any]]: """Update file paths in metadata file""" try: diff --git a/tests/services/test_checkpoint_scanner_sub_type.py b/tests/services/test_checkpoint_scanner_sub_type.py index b28aceaf..19914782 100644 --- a/tests/services/test_checkpoint_scanner_sub_type.py +++ b/tests/services/test_checkpoint_scanner_sub_type.py @@ -164,6 +164,7 @@ def _make_move_scanner(ckpt_root: Path, unet_root: Path) -> CheckpointScanner: scanner._persistent_cache = MagicMock() scanner._name_display_mode = "model_name" scanner._cancel_requested = False + scanner._move_locks = {} scanner._all_folders_backfill_running = False roots = [str(ckpt_root), str(unet_root)] scanner.get_model_roots = lambda: roots diff --git a/tests/services/test_model_scanner_move.py b/tests/services/test_model_scanner_move.py new file mode 100644 index 00000000..410873ee --- /dev/null +++ b/tests/services/test_model_scanner_move.py @@ -0,0 +1,229 @@ +"""Tests for ModelScanner.move_model robustness. + +Covers the failure mode from issue #1126: a move whose source file is +already gone (previous move succeeded but cache/metadata were left stale, +or a duplicate/concurrent move request arrived) must not fail with a raw +FileNotFoundError and leave the model card pointing at empty paths. +""" + +from __future__ import annotations + +import asyncio +import json +import os +from pathlib import Path +from unittest.mock import MagicMock + +import pytest + +from py import config as config_module +from py.services.checkpoint_scanner import CheckpointScanner +from py.services.model_cache import ModelCache +from py.services.model_hash_index import ModelHashIndex +from py.utils.models import CheckpointMetadata + + +def _normalize(path) -> str: + return str(path).replace(os.sep, "/") + + +def _make_scanner(roots) -> CheckpointScanner: + """Create a CheckpointScanner wired for move tests without async init.""" + scanner = object.__new__(CheckpointScanner) + scanner.model_type = "checkpoint" + scanner.model_class = CheckpointMetadata + scanner.file_extensions = {".safetensors"} + scanner._cache = None + scanner._cache_version = 0 + scanner._hash_index = ModelHashIndex() + scanner._tags_count = {} + scanner._excluded_models = [] + scanner._is_initializing = False + scanner._persistent_cache = MagicMock() + scanner._name_display_mode = "model_name" + scanner._cancel_requested = False + scanner._move_locks = {} + scanner._all_folders_backfill_running = False + scanner.get_model_roots = lambda: [_normalize(r) for r in roots] + return scanner + + +@pytest.fixture +def library(tmp_path, monkeypatch): + root = tmp_path / "checkpoints" + root.mkdir() + monkeypatch.setattr(config_module.config, "checkpoints_roots", [str(root)]) + monkeypatch.setattr(config_module.config, "unet_roots", []) + monkeypatch.setattr(config_module.config, "extra_checkpoints_roots", []) + monkeypatch.setattr(config_module.config, "extra_unet_roots", []) + return root + + +def _cache_entry(file_path: str, name: str) -> dict: + return { + "file_path": file_path, + "file_name": name, + "model_name": name, + "folder": "", + "sha256": "abc123", + "sub_type": "checkpoint", + "tags": [], + } + + +def _write_metadata(sidecar: Path, file_path: str, name: str) -> None: + sidecar.write_text( + json.dumps( + { + "file_path": file_path, + "file_name": name, + "model_name": name, + "sha256": "abc123", + "sub_type": "checkpoint", + "hash_status": "completed", + "tags": [], + } + ) + ) + + +@pytest.mark.asyncio +async def test_move_reconciles_when_source_already_moved(library: Path): + """Source missing but the file sits at the natural target: repair cache + and metadata instead of failing with WinError 2.""" + old_dir = library / "old" + old_dir.mkdir() + new_dir = library / "new" + new_dir.mkdir() + + old_path = _normalize(old_dir / "model.safetensors") + moved = new_dir / "model.safetensors" + moved.write_bytes(b"weights") + _write_metadata(new_dir / "model.metadata.json", old_path, "model") + + scanner = _make_scanner([library]) + scanner._cache = ModelCache(raw_data=[_cache_entry(old_path, "model")], folders=[""]) + + result = await scanner.move_model(old_path, _normalize(new_dir)) + + assert result is not None + assert result["new_path"] == _normalize(moved) + + cache = await scanner.get_cached_data() + assert [item["file_path"] for item in cache.raw_data] == [_normalize(moved)] + + saved = json.loads((new_dir / "model.metadata.json").read_text()) + assert saved["file_path"] == _normalize(moved) + + +@pytest.mark.asyncio +async def test_move_reconciles_via_hash_index_when_file_elsewhere(library: Path): + """Source missing and the file is NOT at the requested target (a previous + move took it elsewhere): the hash index locates it, and the stale cache + entry is reused when no sidecar exists at the new location.""" + old_dir = library / "old" + old_dir.mkdir() + elsewhere = library / "elsewhere" + elsewhere.mkdir() + moved = elsewhere / "model.safetensors" + moved.write_bytes(b"weights") + + old_path = _normalize(old_dir / "model.safetensors") + + scanner = _make_scanner([library]) + scanner._cache = ModelCache(raw_data=[_cache_entry(old_path, "model")], folders=[""]) + scanner._hash_index.add_entry("abc123", _normalize(moved)) + + result = await scanner.move_model(old_path, _normalize(library / "target")) + + assert result is not None + assert result["new_path"] == _normalize(moved) + + cache = await scanner.get_cached_data() + entries = [item for item in cache.raw_data] + assert [item["file_path"] for item in entries] == [_normalize(moved)] + # Card data preserved from the stale cache entry + assert entries[0]["model_name"] == "model" + assert entries[0]["sha256"] == "abc123" + + +@pytest.mark.asyncio +async def test_move_returns_none_when_source_missing_and_nowhere_found(library: Path): + """Source gone and no trace of the file anywhere: fail with a clear + error, leaving the cache untouched (a rescan will clean it up).""" + old_path = _normalize(library / "ghost.safetensors") + + scanner = _make_scanner([library]) + scanner._cache = ModelCache(raw_data=[_cache_entry(old_path, "ghost")], folders=[""]) + + result = await scanner.move_model(old_path, _normalize(library / "target")) + + assert result is None + cache = await scanner.get_cached_data() + assert [item["file_path"] for item in cache.raw_data] == [old_path] + + +@pytest.mark.asyncio +async def test_concurrent_moves_of_same_source_are_serialized(library: Path): + """Two simultaneous move requests for the same file: one performs the + move, the other reconciles — no FileNotFoundError, no duplicate cache + entries.""" + source_file = library / "model.safetensors" + source_file.write_bytes(b"weights") + source = _normalize(source_file) + _write_metadata(library / "model.metadata.json", source, "model") + + target_dir = library / "target" + target_file = target_dir / "model.safetensors" + + scanner = _make_scanner([library]) + scanner._cache = ModelCache(raw_data=[_cache_entry(source, "model")], folders=[""]) + + results = await asyncio.gather( + scanner.move_model(source, _normalize(target_dir)), + scanner.move_model(source, _normalize(target_dir)), + ) + + assert all(r is not None for r in results) + assert target_file.exists() + assert not source_file.exists() + + cache = await scanner.get_cached_data() + paths = [item["file_path"] for item in cache.raw_data] + assert paths == [_normalize(target_file)] + + saved = json.loads((target_dir / "model.metadata.json").read_text()) + assert saved["file_path"] == _normalize(target_file) + + +@pytest.mark.asyncio +async def test_move_through_symlinked_directory(library: Path, tmp_path: Path): + """Moving a model that lives under a symlinked directory uses the + business path: the file leaves the physical directory, the symlink + itself stays intact, and the cache records the unresolved path.""" + real_dir = tmp_path / "real_root" + real_dir.mkdir() + link_dir = library / "linked" + link_dir.symlink_to(real_dir, target_is_directory=True) + + model = real_dir / "model.safetensors" + model.write_bytes(b"weights") + source = _normalize(link_dir / "model.safetensors") + _write_metadata(real_dir / "model.metadata.json", source, "model") + + target_dir = library / "target" + target_file = target_dir / "model.safetensors" + + scanner = _make_scanner([library]) + scanner._cache = ModelCache(raw_data=[_cache_entry(source, "model")], folders=[""]) + + result = await scanner.move_model(source, _normalize(target_dir)) + + assert result is not None + assert result["new_path"] == _normalize(target_file) + assert target_file.exists() + assert not model.exists() + assert link_dir.is_symlink() + + cache = await scanner.get_cached_data() + assert [item["file_path"] for item in cache.raw_data] == [_normalize(target_file)]