Merge pull request #1134 from willmiao/perf/bulk-rename-apply

perf(rename): make bulk filename-template apply O(n) instead of O(n²)
This commit is contained in:
pixelpaws
2026-10-03 15:39:58 +08:00
committed by GitHub
8 changed files with 1080 additions and 63 deletions
+118 -18
View File
@@ -2,10 +2,12 @@
from __future__ import annotations
import asyncio
import json
import logging
import os
from typing import Any, Awaitable, Callable, Dict, Iterable, List, Mapping, Optional, TYPE_CHECKING, cast
from contextlib import asynccontextmanager
from typing import Any, AsyncIterator, Awaitable, Callable, Dict, Iterable, List, Mapping, Optional, TYPE_CHECKING, cast
from ..services.service_registry import ServiceRegistry
from ..services.pending_delete_service import get_pending_delete_service
@@ -107,6 +109,36 @@ def _require_path_in_library_roots(file_path: str, scanner, *, label: str = "pat
)
class BulkRenameContext:
"""Per-session state threaded through ``rename_model`` calls of a bulk rename.
Holds the lazily built recipe hash index so a bulk rename loop pays the
O(recipes) index build at most once (on the first recipe-touching rename)
instead of rescanning every recipe per renamed file. Also tracks whether
any recipe was re-pointed so the session finalizes recipe maintenance only
when needed.
"""
def __init__(self, recipe_scanner: Any) -> None:
self._recipe_scanner = recipe_scanner
self._recipe_hash_index: Optional[Dict[str, List[Dict[str, Any]]]] = None
self.recipes_touched = False
@property
def recipe_scanner(self) -> Any:
return self._recipe_scanner
async def get_recipe_hash_index(self) -> Optional[Dict[str, List[Dict[str, Any]]]]:
"""Return the lora-hash → recipes index, building it on first use."""
if self._recipe_scanner is None:
return None
if self._recipe_hash_index is None:
self._recipe_hash_index = (
await self._recipe_scanner.build_lora_hash_index()
)
return self._recipe_hash_index
class ModelLifecycleService:
"""Co-ordinate destructive and mutating model operations."""
@@ -365,10 +397,45 @@ class ModelLifecycleService:
return await self._scanner.bulk_delete_models(file_paths)
@asynccontextmanager
async def bulk_rename_session(self) -> AsyncIterator[BulkRenameContext]:
"""Context for bulk rename loops (filename-template "Apply to Library").
While active, the per-file ``update_single_model_cache`` resort/persist
chain and the per-file recipe folder-metadata refresh/resort are
deferred; both run exactly once when the outermost session exits — see
``ModelScanner.defer_cache_persist`` and
``RecipeScanner.finalize_bulk_filename_updates``. The finalize steps run
even on cancellation or mid-loop errors, because files are already
renamed on disk and the caches must not be left diverging.
Yields a :class:`BulkRenameContext` to pass as ``bulk_context`` into
each ``rename_model`` call of the loop.
"""
recipe_scanner = await self._recipe_scanner_factory()
context = BulkRenameContext(recipe_scanner)
async with self._scanner.defer_cache_persist():
try:
yield context
finally:
if recipe_scanner is not None and context.recipes_touched:
try:
await recipe_scanner.finalize_bulk_filename_updates()
except Exception as exc: # pragma: no cover - defensive logging
logger.error(
"Error finalizing bulk recipe updates: %s", exc
)
async def rename_model(
self, *, file_path: str, new_file_name: str
self, *, file_path: str, new_file_name: str, bulk_context: Optional[BulkRenameContext] = None
) -> Dict[str, object]:
"""Rename a model and its companion artefacts."""
"""Rename a model and its companion artefacts.
When ``bulk_context`` is given (bulk rename loop), the recipe
re-pointing uses the session's prebuilt hash index and defers recipe
maintenance to the session finalize; the scanner cache persist is
likewise deferred by the surrounding ``bulk_rename_session``.
"""
if not file_path or not new_file_name:
raise ValueError("File path and new file name are required")
@@ -419,20 +486,11 @@ class ModelLifecycleService:
raw_hash = metadata.get("sha256") if isinstance(metadata, dict) else None
hash_value = raw_hash if isinstance(raw_hash, str) else None
renamed_files: List[str] = []
new_metadata_path: Optional[str] = None
new_preview: Optional[str] = None
for old_path, pattern in existing_files:
ext = self._get_multipart_ext(pattern)
new_path = os.path.join(
os.path.dirname(old_path), f"{new_file_name}{ext}"
).replace(os.sep, "/")
os.rename(old_path, new_path)
renamed_files.append(new_path)
if ext == ".metadata.json":
new_metadata_path = new_path
renamed_files, new_metadata_path = await asyncio.to_thread(
self._rename_companion_files, existing_files, new_file_name
)
if metadata and new_metadata_path:
metadata["file_name"] = new_file_name
@@ -457,12 +515,26 @@ class ModelLifecycleService:
)
if hash_value and getattr(self._scanner, "model_type", "") == "lora":
recipe_scanner = await self._recipe_scanner_factory()
if bulk_context is not None:
recipe_scanner = bulk_context.recipe_scanner
hash_index = await bulk_context.get_recipe_hash_index()
defer_maintenance = True
else:
recipe_scanner = await self._recipe_scanner_factory()
hash_index = None
defer_maintenance = False
if recipe_scanner:
try:
await recipe_scanner.update_lora_filename_by_hash(
hash_value, new_file_name
file_count, cache_count = (
await recipe_scanner.update_lora_filename_by_hash(
hash_value,
new_file_name,
hash_index=hash_index,
defer_maintenance=defer_maintenance,
)
)
if bulk_context is not None and (file_count or cache_count):
bulk_context.recipes_touched = True
except Exception as exc: # pragma: no cover - defensive logging
logger.error(
"Error updating recipe references for %s: %s",
@@ -478,6 +550,34 @@ class ModelLifecycleService:
"reload_required": False,
}
def _rename_companion_files(
self,
existing_files: List[tuple[str, str]],
new_file_name: str,
) -> tuple[List[str], Optional[str]]:
"""Rename all companion files, off the event loop thread.
Runs the blocking ``os.rename`` sequence for one model in a worker
thread so a single file's HDD I/O does not stall the event loop.
Never parallelized across files: one model's renames stay sequential
and the helper holds no locks.
"""
renamed_files: List[str] = []
new_metadata_path: Optional[str] = None
for old_path, pattern in existing_files:
ext = self._get_multipart_ext(pattern)
new_path = os.path.join(
os.path.dirname(old_path), f"{new_file_name}{ext}"
).replace(os.sep, "/")
os.rename(old_path, new_path)
renamed_files.append(new_path)
if ext == ".metadata.json":
new_metadata_path = new_path
return renamed_files, new_metadata_path
@staticmethod
def _get_multipart_ext(filename: str) -> str:
"""Return the extension for files with compound suffixes."""
+115 -18
View File
@@ -4,6 +4,7 @@ import logging
import asyncio
import time
import shutil
from contextlib import asynccontextmanager
from dataclasses import dataclass
from typing import Any, Awaitable, Callable, Dict, List, Mapping, Optional, Sequence, Set, Tuple, Type, Union, cast
@@ -162,6 +163,13 @@ class ModelScanner:
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
# Bulk-operation deferral: while _defer_persist_depth > 0,
# update_single_model_cache() skips the per-call resort/persist and
# only marks _deferred_persist_pending; the exit of the outermost
# defer_cache_persist() context finalizes once (see
# _finalize_deferred_cache_persist).
self._defer_persist_depth = 0
self._deferred_persist_pending = False
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).
@@ -773,11 +781,11 @@ class ModelScanner:
except Exception as exc:
logger.warning("AutoV3 backfill failed: %s", exc)
async def _save_persistent_cache(self, scan_result: CacheBuildResult) -> None:
async def _save_persistent_cache(self, scan_result: CacheBuildResult, *, force: bool = False) -> None:
if not scan_result or not getattr(self, '_persistent_cache', None):
return
if self.is_cancelled():
if self.is_cancelled() and not force:
logger.info(
f"{self.model_type.capitalize()} Scanner: Skipping _save_persistent_cache "
"after cancellation"
@@ -836,7 +844,7 @@ class ModelScanner:
bucket.append(path)
return snapshot
async def _persist_current_cache(self) -> None:
async def _persist_current_cache(self, *, force: bool = False) -> None:
if self._cache is None or not getattr(self, '_persistent_cache', None):
return
@@ -851,7 +859,7 @@ class ModelScanner:
else None
),
)
await self._save_persistent_cache(snapshot)
await self._save_persistent_cache(snapshot, force=force)
await self._sync_download_history(snapshot.raw_data, source='scan')
def _count_model_files(self) -> int:
"""Count all model files with supported extensions in all roots
@@ -2646,11 +2654,85 @@ class ModelScanner:
logger.error(f"Error updating metadata paths: {e}", exc_info=True)
return None
@asynccontextmanager
async def defer_cache_persist(self):
"""Defer heavyweight cache maintenance for a bulk operation.
While at least one ``defer_cache_persist`` context is active,
:meth:`update_single_model_cache` performs only the in-memory entry
swap plus incremental index updates — it skips the full version-index
rebuild, the natsort resort, and the whole-table SQLite persist plus
download-history sync that normally run per call. When the outermost
context exits, the pending maintenance runs **once** (resort, persist,
download-history sync).
The final persist is forced: it runs even when the scanner's
cancellation flag is set or the wrapped block raised, because callers
use this around operations that already mutated files on disk and the
cache must not be left diverging from reality.
Intended for bulk rename/move loops (e.g. the filename-template "Apply
to Library" flow). Single-shot callers keep the immediate per-call
behavior by not entering this context.
"""
self._defer_persist_depth = getattr(self, "_defer_persist_depth", 0) + 1
try:
yield
finally:
self._defer_persist_depth -= 1
if self._defer_persist_depth == 0:
await self._finalize_deferred_cache_persist()
@property
def _cache_persist_deferred(self) -> bool:
"""True while cache resort/persist is deferred to a bulk finalize."""
return getattr(self, "_defer_persist_depth", 0) > 0
async def _finalize_deferred_cache_persist(self) -> None:
"""Run the resort + persist deferred by ``defer_cache_persist``.
Best-effort: failures are logged, never raised, so an error here
cannot mask the outcome of the bulk operation itself (including
cancellation).
"""
if not getattr(self, "_deferred_persist_pending", False):
return
self._deferred_persist_pending = False
if self._cache is None:
return
try:
# resort() rebuilds the version index and folder list, so the
# per-call rebuilds skipped during deferral are covered here.
await self._cache.resort()
await self._persist_current_cache(force=True)
self.bump_cache_version()
except Exception:
logger.error(
"%s Scanner: failed to finalize deferred cache persist",
self.model_type.capitalize(),
exc_info=True,
)
async def update_single_model_cache(self, original_path: str, new_path: str, metadata: Optional[Dict[str, Any]], recalculate_type: bool = False) -> Union[bool, Dict[str, Any]]:
"""Update cache after a model has been moved or modified"""
"""Update cache after a model has been moved or modified.
Performs the full maintenance chain (version-index rebuild, resort,
whole-table persist, download-history sync) unless the scanner is
inside a :meth:`defer_cache_persist` context, in which case only
the in-memory entry swap and incremental index updates run and the
heavy chain executes once at context exit.
"""
deferred = self._cache_persist_deferred
cache = await self.get_cached_data()
existing_item = next((item for item in cache.raw_data if item['file_path'] == original_path), None)
existing_index: Optional[int] = None
existing_item = None
for idx, item in enumerate(cache.raw_data):
if item['file_path'] == original_path:
existing_item = item
existing_index = idx
break
if existing_item:
cache.remove_from_version_index(existing_item)
@@ -2662,11 +2744,18 @@ class ModelScanner:
del self._tags_count[tag]
self._hash_index.remove_by_path(original_path)
cache.raw_data = [
item for item in cache.raw_data
if item['file_path'] != original_path
]
if deferred:
# In-place swap avoids the O(n) list rebuild per renamed file;
# indexes were already updated incrementally above/below, and the
# folder recompute happens in the single finalize resort().
if existing_index is not None:
cache.raw_data.pop(existing_index)
else:
cache.raw_data = [
item for item in cache.raw_data
if item['file_path'] != original_path
]
cache_modified = bool(existing_item) or bool(metadata)
cache_entry: Optional[Dict[str, Any]] = None
@@ -2707,8 +2796,11 @@ class ModelScanner:
cache_entry.get('autov3') or None,
)
all_folders = set(item['folder'] for item in cache.raw_data)
cache.folders = sorted(list(all_folders), key=lambda x: x.lower())
if not deferred:
# O(n) over raw_data; the finalize resort() recomputes the
# folder list once, so bulk callers skip it per file.
all_folders = set(item['folder'] for item in cache.raw_data)
cache.folders = sorted(list(all_folders), key=lambda x: x.lower())
# The move target may live in directories the last scan never saw;
# record the destination folder (and its parents) in the known
@@ -2723,13 +2815,18 @@ class ModelScanner:
for tag in cache_entry.get('tags', []):
self._tags_count[tag] = self._tags_count.get(tag, 0) + 1
cache.rebuild_version_index()
if deferred:
if cache_modified:
self._deferred_persist_pending = True
self.bump_cache_version()
else:
cache.rebuild_version_index()
await cache.resort()
await cache.resort()
if cache_modified:
await self._persist_current_cache()
self.bump_cache_version()
if cache_modified:
await self._persist_current_cache()
self.bump_cache_version()
if metadata and cache_entry is not None:
return cache_entry
+64 -8
View File
@@ -4580,14 +4580,64 @@ class RecipeScanner:
return syntax_parts
async def build_lora_hash_index(self) -> Dict[str, List[Dict[str, Any]]]:
"""Build a one-shot lowercase-LoRA-hash → recipes index.
Scans the recipe cache exactly once (O(recipes × loras)) and returns
a mapping of lowercase lora ``hash`` to the list of recipe dicts
containing it. Bulk rename loops pass this index to
:meth:`update_lora_filename_by_hash` so per-file lookups are O(1)
instead of rescanning every recipe for each renamed LoRA.
"""
cache = await self.get_cached_data()
index: Dict[str, List[Dict[str, Any]]] = {}
if not cache or not cache.raw_data:
return index
for recipe in cache.raw_data:
loras = recipe.get("loras", [])
if not isinstance(loras, list):
continue
for lora in loras:
if not isinstance(lora, dict):
continue
hash_value = (lora.get("hash") or "").lower()
if hash_value:
index.setdefault(hash_value, []).append(recipe)
return index
async def finalize_bulk_filename_updates(self) -> None:
"""Run once after a bulk rename session that deferred maintenance.
Refreshes folder metadata and schedules a single re-sort. Filename-only
renames never change recipe folders, so the deferred refresh is
redundant but cheap; skipping it per file is what makes bulk renames
O(1)-per-file.
"""
if self._cache is None:
return
self._schedule_resort()
async def update_lora_filename_by_hash(
self, hash_value: str, new_file_name: str
self,
hash_value: str,
new_file_name: str,
*,
hash_index: Optional[Dict[str, List[Dict[str, Any]]]] = None,
defer_maintenance: bool = False,
) -> Tuple[int, int]:
"""Update file_name in all recipes that contain a LoRA with the specified hash.
Args:
hash_value: The SHA256 hash value of the LoRA
new_file_name: The new file_name to set
hash_index: Optional prebuilt index from
:meth:`build_lora_hash_index`. When given, the O(recipes)
cache scan (and its folder-metadata walk) is skipped and the
affected recipes are looked up directly — the bulk rename path.
defer_maintenance: When True, skip the folder-metadata refresh and
resort scheduling. The caller MUST run
:meth:`finalize_bulk_filename_updates` exactly once afterwards.
Returns:
Tuple[int, int]: (number of recipes updated in files, number of recipes updated in cache)
@@ -4598,17 +4648,21 @@ class RecipeScanner:
# Always use lowercase hash for consistency
hash_value = hash_value.lower()
# Get cache
cache = await self.get_cached_data()
if not cache or not cache.raw_data:
return 0, 0
if hash_index is not None:
candidate_recipes = hash_index.get(hash_value, [])
else:
# Get cache
cache = await self.get_cached_data()
if not cache or not cache.raw_data:
return 0, 0
candidate_recipes = cache.raw_data
file_updated_count = 0
cache_updated_count = 0
# Find recipes that need updating from the cache
# Find recipes that need updating
recipes_to_update = []
for recipe in cache.raw_data:
for recipe in candidate_recipes:
loras = recipe.get("loras", [])
if not isinstance(loras, list):
continue
@@ -4654,7 +4708,9 @@ class RecipeScanner:
# We don't necessarily need to resort because LoRA file_name isn't a sort key,
# but we might want to schedule a resort if we're paranoid or if searching relies on sorted state.
# Given it's a rename of a dependency, search results might change if searching by LoRA name.
self._schedule_resort()
# Bulk callers defer this to a single finalize_bulk_filename_updates() call.
if not defer_maintenance:
self._schedule_resort()
return file_updated_count, cache_updated_count
@@ -33,7 +33,9 @@ class FilenameTemplateUseCase:
An empty template restores the recorded original filename instead of
rendering a template. Shares the auto-organize lock (and its in-progress
error) so a bulk rename never runs concurrently with an auto-organize
operation.
operation. The whole loop runs inside a bulk rename session so cache
persist/resort and recipe maintenance happen once at the end instead of
per renamed file.
"""
def __init__(
@@ -106,23 +108,24 @@ class FilenameTemplateUseCase:
await self._emit_progress(progress_callback, result, "started")
for index in range(0, result.total, AUTO_ORGANIZE_BATCH_SIZE):
if self._scanner.is_cancelled():
logger.info(
"Filename template apply cancelled for %s", self._model_type
)
break
batch = models[index : index + AUTO_ORGANIZE_BATCH_SIZE]
for model in batch:
async with self._lifecycle_service.bulk_rename_session() as bulk_context:
for index in range(0, result.total, AUTO_ORGANIZE_BATCH_SIZE):
if self._scanner.is_cancelled():
logger.info(
"Filename template apply cancelled for %s", self._model_type
)
break
await self._process_model(model, template, result)
result.processed += 1
await self._emit_progress(progress_callback, result, "processing")
# Yield between batches so the server stays responsive.
await asyncio.sleep(0.1)
batch = models[index : index + AUTO_ORGANIZE_BATCH_SIZE]
for model in batch:
if self._scanner.is_cancelled():
break
await self._process_model(model, template, result, bulk_context)
result.processed += 1
await self._emit_progress(progress_callback, result, "processing")
# Yield between batches so the server stays responsive.
await asyncio.sleep(0.1)
if self._scanner.is_cancelled():
result.status = "cancelled"
@@ -150,6 +153,7 @@ class FilenameTemplateUseCase:
model: Dict[str, Any],
template: str,
result: AutoOrganizeResult,
bulk_context: Any = None,
) -> None:
model_name = model.get("model_name", "Unknown")
try:
@@ -177,7 +181,7 @@ class FilenameTemplateUseCase:
return
await self._lifecycle_service.rename_model(
file_path=file_path, new_file_name=new_stem
file_path=file_path, new_file_name=new_stem, bulk_context=bulk_context
)
result.success_count += 1