fix(scanner): make move_model resilient to missing source files (#1126)

Concurrent or repeated move requests for the same model raced each other:
the first move succeeded, the rest failed with FileNotFoundError, leaving
the model card pointing at stale/empty paths.

- Serialize moves per source file with an asyncio.Lock keyed on the
  normalized source path
- When the source file is already gone, reconcile instead of failing:
  locate the model via the hash index or the expected target paths, repair
  the metadata sidecar and cache entry, and reuse the stale cache entry
  when no sidecar exists at the new location
- Avoid duplicate cache entries when the cache already tracks the moved
  file; only drop the stale source entry
- Move via business paths (abspath) instead of realpath, matching every
  other file mutation and the containment check; realpath stays reserved
  for scanner dedup per project convention
This commit is contained in:
Will Miao
2026-10-01 08:55:47 +08:00
parent 42fa8294df
commit eeb9270827
3 changed files with 327 additions and 10 deletions
+97 -10
View File
@@ -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:
@@ -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
+229
View File
@@ -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)]