feat(scanner): report and parallelize the reconcile walk

A regular Refresh walked every configured root sequentially, so a full
25 TB drive delayed the roots behind it, and the dialog sat on "Checking
for changes..." at 0 % for the whole walk with no way to tell it was
working. On a cold external drive the walk itself dominates the cost, so
the fix is to overlap the drives and to show what the walk is doing.

Backend (py/services/model_scanner.py):
* Extract the per-root walk into the synchronous _walk_root_for_reconcile()
  worker and merge its results on the event loop afterwards, in configured
  root order: which business path wins a file reachable through several
  roots must not depend on the order the workers happened to finish in.
* Group roots by device (_root_device_key: drive letter on Windows, st_dev
  on POSIX) and run one worker per device. Roots sharing a device stay
  sequential, so directory claims and the overlap dedup (#871, #1041) keep
  their configured-order semantics; different devices run in parallel.
* Track walk progress per root (_ReconcileWalkTracker), weighted by each
  root's cached entry count because the real file count is only known once
  the walk ends. The bar splits into walk (0-50 %) and new-file (50-99 %)
  phases so it never jumps backwards, and the walk broadcasts files seen,
  active roots and the ETA counters.
* Replace the Windows case-insensitive fallback -- a scan of every cached
  path per miss, i.e. O(files x cached) -- with a lazily built lower-cased
  index (_CachedPathLookups, also now guarding the realpath alias index for
  worker threads), and lift the os.name == "nt" gate into the module-level
  _CASE_INSENSITIVE_PATHS so the branch is testable off Windows.
* Excluded-model membership is a set lookup instead of a list scan.

Frontend:
* render the walk phase as "Checking for changes... <roots> (N files)" with
  the ETA, and reset the ETA tracker when the stage changes: the per-file
  rate of counting files says nothing about processing them.
* add common.scanProgress.walkFiles; the other locales keep the sanctioned
  [TODO: Translate] placeholder until the translation pass.

Verified: no-change reconcile over 10k files/500 dirs 101 ms and 10k/5000
dirs 230 ms (was 93/229 ms, within noise); a two-device sandbox walk runs
both roots concurrently, names them in the progress messages and finishes
with added=15, removed=0; 3665 passed, 7 skipped; frontend 1456 passed,
vue widgets 96 passed.
This commit is contained in:
Will Miao
2026-10-06 12:40:57 +08:00
parent 8cd53c20f8
commit 14e660e876
16 changed files with 1059 additions and 187 deletions
+558 -118
View File
@@ -2,10 +2,11 @@ import json
import os
import logging
import asyncio
import threading
import time
import shutil
from contextlib import asynccontextmanager
from dataclasses import dataclass
from contextlib import asynccontextmanager, suppress
from dataclasses import dataclass, field
from typing import Any, Awaitable, Callable, Dict, List, Mapping, Optional, Sequence, Set, Tuple, Type, Union, cast
from ..utils.models import BaseModelMetadata, autov3_from_civitai_files
@@ -89,6 +90,388 @@ PAGE_TYPE_MAP = {
'other': 'other',
}
# Case-insensitive cache matching is a Windows-only safety net (NTFS treats
# "Model.txt" and "model.TXT" as the same file). Module-level so the branch can
# be exercised on case-sensitive filesystems in tests.
_CASE_INSENSITIVE_PATHS: bool = os.name == "nt"
# Cadence of walk-phase progress broadcasts during a reconcile.
_WALK_PROGRESS_INTERVAL_SECONDS = 0.5
# Share of the progress bar owned by the reconcile walk. The new-file pass
# takes the rest (50-99) so the bar never jumps backwards between phases.
_WALK_PROGRESS_SHARE = 50
# How many files a walk worker may scan between two progress reports.
_WALK_PROGRESS_REPORT_EVERY = 256
def _new_file_pass_progress(processed: int, total: int) -> int:
"""Bar percentage for the new-file pass (walk share .. 99)."""
if total <= 0:
return _WALK_PROGRESS_SHARE
ratio = min(max(processed / total, 0.0), 1.0)
return _WALK_PROGRESS_SHARE + min(49, int(ratio * 49))
def _root_display_label(root: str) -> str:
"""Short label identifying a model root in progress messages."""
drive, _tail = os.path.splitdrive(root)
if drive:
return drive
normalized = root.rstrip("/\\")
return os.path.basename(normalized) or root
def _normalized_root_prefix(root: str) -> str:
"""Forward-slash root prefix used to attribute cached paths to a root."""
prefix = os.path.normpath(root).replace(os.sep, "/")
if not prefix.endswith("/"):
prefix += "/"
return prefix.lower() if _CASE_INSENSITIVE_PATHS else prefix
def _count_cached_entries_per_root(
cached_paths: Set[str], roots: Sequence[str]
) -> Dict[str, int]:
"""Attribute cached entries to model roots (longest prefix wins).
Used as the walk-workload weight for progress reporting: the walk itself is
what discovers the real file count, so the cached entry count is the only
estimate available up front.
"""
counts: Dict[str, int] = {root: 0 for root in roots}
if not roots:
return counts
prefixes = sorted(
((_normalized_root_prefix(root), root) for root in roots),
key=lambda item: len(item[0]),
reverse=True,
)
case_insensitive = _CASE_INSENSITIVE_PATHS
for path in cached_paths:
candidate = path.lower() if case_insensitive else path
for prefix, root in prefixes:
if candidate.startswith(prefix):
counts[root] += 1
break
return counts
def _root_device_key(root: str) -> str:
"""Group key identifying the storage device that holds a root.
Roots sharing a device must be walked sequentially so directory claims stay
deterministic (configured order wins), while roots on different devices are
independent and can be walked by parallel workers.
"""
drive, _tail = os.path.splitdrive(root)
if drive:
return drive.lower()
try:
return f"dev:{os.stat(root).st_dev}"
except OSError:
return f"path:{os.path.normpath(root)}"
def _group_roots_by_device(roots: Sequence[str]) -> List[List[str]]:
"""Group roots by device, preserving the configured root order.
Roots on distinct devices cannot alias unless a symlink crosses devices, so
walking them in parallel cannot disturb the directory-claim order in any
realistic layout (and even then the realpath alias index keeps the cached
business path stable).
"""
groups: Dict[str, List[str]] = {}
for root in roots:
groups.setdefault(_root_device_key(root), []).append(root)
return list(groups.values())
def _build_realpath_index(cached_paths: Set[str]) -> Dict[str, str]:
"""Map realpath -> cached business path (overlapping roots / symlinks)."""
index: Dict[str, str] = {}
for cached_path in cached_paths:
try:
index.setdefault(os.path.realpath(cached_path), cached_path)
except Exception:
continue
return index
def _build_casefold_index(cached_paths: Set[str]) -> Dict[str, str]:
"""Map lower-cased path -> cached business path (Windows case fallback)."""
index: Dict[str, str] = {}
for cached_path in cached_paths:
index.setdefault(cached_path.lower(), cached_path)
return index
class _CachedPathLookups:
"""Lazily built path indexes over the cached paths, safe for workers.
Both indexes are only ever consulted for walk *misses*, so building them on
first use keeps a no-change refresh free of per-cached-entry realpath work.
The case-fold index replaces the former linear scan over every cached path,
which was O(files x cached entries) whenever it was reached.
"""
def __init__(self, cached_paths: Set[str]) -> None:
self.cached_paths = cached_paths
self._lock = threading.Lock()
self._realpath_index: Optional[Dict[str, str]] = None
self._casefold_index: Optional[Dict[str, str]] = None
def match_real_path(self, real_path: str) -> Optional[str]:
index = self._realpath_index
if index is None:
with self._lock:
if self._realpath_index is None:
self._realpath_index = _build_realpath_index(self.cached_paths)
index = self._realpath_index
return index.get(real_path)
def match_casefold_path(self, file_path: str) -> Optional[str]:
index = self._casefold_index
if index is None:
with self._lock:
if self._casefold_index is None:
self._casefold_index = _build_casefold_index(self.cached_paths)
index = self._casefold_index
return index.get(file_path.lower())
class _RealDirClaims:
"""Thread-safe claim set for real directory paths.
Mirrors the historical global ``visited_real_paths`` set: a directory
reachable through several roots (overlapping roots, symlinked aliases) must
only be walked once, otherwise the same physical file could end up in the
library under two business paths.
"""
def __init__(self) -> None:
self._lock = threading.Lock()
self._claimed: Set[str] = set()
def claim(self, real_path: str) -> bool:
with self._lock:
if real_path in self._claimed:
return False
self._claimed.add(real_path)
return True
class _ReconcileWalkTracker:
"""Thread-safe progress bookkeeping for the reconcile walk.
Per-root progress is weighted by the number of cached entries under that
root: the walk is what discovers the real file count, so the cached count
is the only workload estimate available before it runs. Roots with no
cached entries contribute a single weight unit and are reported only once
they finish.
"""
def __init__(self, roots: Sequence[str], expected: Mapping[str, int]) -> None:
self._lock = threading.Lock()
self._roots: List[str] = list(roots)
self._labels = {root: _root_display_label(root) for root in self._roots}
self._expected = {
root: max(int(expected.get(root, 0) or 0), 0) for root in self._roots
}
self._weight = {
root: max(self._expected[root], 1) for root in self._roots
}
self._seen = {root: 0 for root in self._roots}
self._done = {root: False for root in self._roots}
self._active: Set[str] = set()
def start_root(self, root: str) -> None:
with self._lock:
if root in self._seen:
self._active.add(root)
def report(self, root: str, files_seen: int) -> None:
with self._lock:
if root in self._seen:
self._seen[root] = max(int(files_seen), 0)
def finish_root(self, root: str) -> None:
with self._lock:
self._done[root] = True
self._active.discard(root)
def snapshot(self) -> Optional[Dict[str, Any]]:
"""Return the current payload extras plus the bar percentage."""
with self._lock:
if not self._roots:
return None
total_weight = sum(self._weight.values()) or 1
completed = 0.0
files_seen = 0
roots_done = 0
for root in self._roots:
files_seen += self._seen[root]
if self._done[root]:
completed += self._weight[root]
roots_done += 1
continue
expected = self._expected[root]
if expected > 0:
completed += self._weight[root] * min(
self._seen[root] / expected, 0.99
)
active = [root for root in self._roots if root in self._active]
if active:
current = max(active, key=lambda root: self._seen[root])
else:
current = self._roots[-1]
expected_total = sum(self._expected.values())
ratio = min(completed / total_weight, 1.0)
return {
'progress': min(_WALK_PROGRESS_SHARE, int(ratio * _WALK_PROGRESS_SHARE)),
# processed/total feed the client-side ETA; the total is a lower
# bound and never below the files already seen, so the estimate
# stays coherent when the library grew since the last scan.
'processed': files_seen,
'total': max(expected_total, files_seen),
'files_seen': files_seen,
'roots_total': len(self._roots),
'roots_done': roots_done,
'active_roots': [self._labels[root] for root in active][:4],
'current_name': self._labels[current],
}
@dataclass
class _RootWalkResult:
"""Outcome of walking a single model root (no scanner state touched)."""
root_path: str
found_paths: Set[str] = field(default_factory=set)
stale_paths: List[str] = field(default_factory=list)
# (business path, real path) pairs for files that are neither cached nor
# excluded; the caller claims real paths in configured root order.
new_candidates: List[Tuple[str, str]] = field(default_factory=list)
discovered_folders: Set[str] = field(default_factory=set)
files_seen: int = 0
cancelled: bool = False
def _walk_root_for_reconcile(
*,
root_path: str,
file_extensions: Set[str],
cached_paths: Set[str],
path_to_item: Mapping[str, Dict[str, Any]],
lookups: _CachedPathLookups,
dir_claims: _RealDirClaims,
excluded_models: Set[str],
is_cancelled: Callable[[], bool],
report_progress: Callable[[int], None],
) -> _RootWalkResult:
"""Walk one model root and classify every model file found.
Runs synchronously (worker thread) and never mutates scanner state: all
cache updates happen on the event loop once every root has been walked.
"""
result = _RootWalkResult(root_path=root_path)
stale_seen: Set[str] = set()
files_since_report = 0
def mark_stale_if_needed(cached_path: str) -> None:
"""Queue a cached path for file_name repair when it drifted."""
if cached_path in stale_seen:
return
item = path_to_item.get(cached_path)
if item is None:
return
if item.get("file_name") == _file_name_stem(cached_path):
return
stale_seen.add(cached_path)
result.stale_paths.append(cached_path)
for root, dirnames, files in os.walk(root_path, followlinks=True):
dirnames[:] = [d for d in dirnames if not _is_excluded_dir(d)]
real_root = os.path.realpath(root)
if not dir_claims.claim(real_root):
continue
# Record every visited directory (including empty ones) so the folder
# tree stays accurate without a live walk.
rel_dir = os.path.relpath(
os.path.abspath(root), os.path.abspath(root_path)
).replace(os.path.sep, "/")
if rel_dir != "." and not _is_hidden_relative_path(rel_dir):
result.discovered_folders.add(rel_dir)
for file in files:
ext = os.path.splitext(file)[1].lower()
if ext not in file_extensions:
continue
result.files_seen += 1
files_since_report += 1
if files_since_report >= _WALK_PROGRESS_REPORT_EVERY:
# A flat library (thousands of files in one directory) would
# otherwise report nothing until that directory ends.
files_since_report = 0
report_progress(result.files_seen)
# Construct paths exactly as they would be in cache
file_path = os.path.join(root, file).replace(os.sep, '/')
if file_path in cached_paths:
result.found_paths.add(file_path)
mark_stale_if_needed(file_path)
continue
# Only a cache miss needs the physical path, so the realpath
# syscalls are paid per changed file rather than per file in the
# library.
real_file_path = os.path.realpath(os.path.join(root, file))
cached_real_match = lookups.match_real_path(real_file_path)
if cached_real_match:
result.found_paths.add(cached_real_match)
mark_stale_if_needed(cached_real_match)
continue
if file_path in excluded_models:
continue
# Windows: a cached path may differ from the walk result only by
# case. O(1) via the lower-cased index (was a full scan of the
# cached paths per miss).
if _CASE_INSENSITIVE_PATHS:
cached_case_match = lookups.match_casefold_path(file_path)
if cached_case_match:
result.found_paths.add(cached_case_match)
mark_stale_if_needed(cached_case_match)
continue
# Not cached yet; the caller claims the real path in root order so
# a file reachable through several roots is only added once.
result.new_candidates.append((file_path, real_file_path))
if files_since_report:
# Publish per directory so a library spread over many directories
# reports as it progresses, not only at the end of the walk.
files_since_report = 0
report_progress(result.files_seen)
if is_cancelled():
result.cancelled = True
break
report_progress(result.files_seen)
return result
def _is_pending_delete_path(path: str) -> bool:
"""Return True when any path component is the pending-delete staging dir."""
@@ -1082,129 +1465,76 @@ class ModelScanner:
cached_paths = {item['file_path'] for item in self._cache.raw_data}
path_to_item = {item['file_path']: item for item in self._cache.raw_data}
# physical path -> cached business path, for the alias case where the
# same file is reachable under a different path than the cached one
# (overlapping roots / symlink layout changes): keep the existing
# entry instead of delete + re-add (which would re-read metadata and
# re-hash every file). Built lazily on the first miss, because a
# realpath per cached entry is ~half the cost of a no-change
# reconcile and the map is only ever consulted for misses.
cached_real_paths: Optional[Dict[str, str]] = None
# Every configured root that is currently reachable. A root that is
# missing (drive switched off, unmounted share) is skipped, so its
# cached entries are reported as missing below.
roots: List[str] = []
seen_roots: Set[str] = set()
for root_path in self.get_model_roots():
if not root_path or root_path in seen_roots:
continue
if not os.path.exists(root_path):
continue
seen_roots.add(root_path)
roots.append(root_path)
def lookup_cached_real_path(real_path: str) -> Optional[str]:
nonlocal cached_real_paths
if cached_real_paths is None:
cached_real_paths = {}
for cached_path in cached_paths:
try:
cached_real_paths.setdefault(os.path.realpath(cached_path), cached_path)
except Exception:
continue
return cached_real_paths.get(real_path)
# Track found files and new files
found_paths = set()
new_files = []
# Roots on different devices are walked by parallel workers (a cold
# or slow drive then no longer serializes the others); roots sharing
# a device stay sequential so directory claims remain deterministic.
tracker = _ReconcileWalkTracker(
roots, _count_cached_entries_per_root(cached_paths, roots)
)
walk_results = await self._walk_roots_for_reconcile(
roots=roots,
tracker=tracker,
cached_paths=cached_paths,
path_to_item=path_to_item,
)
# Final walk snapshot: the bar reaches the walk share, then the
# new-file pass continues from there.
await self._broadcast_walk_progress(tracker)
if self.is_cancelled():
logger.info(f"{self.model_type.capitalize()} Scanner: Reconcile scan cancelled")
await self._broadcast_scan_progress(
'cancelled', 'reconcile_scan', 0, False,
elapsed_seconds=time.time() - start_time,
)
return
# Merge the per-root results in configured root order: which
# business path wins a file reachable through several roots (and
# which cached entry counts as found) must not depend on the order
# the workers happened to finish in.
found_paths: Set[str] = set()
discovered_folders: Set[str] = set()
# Cached entries whose stored file_name no longer matches the file
# on disk (e.g. dotted stems truncated by the legacy .civitai.info
# on disk (e.g. dotted stems truncated by the legacy .civitai_info
# migration, issue #1112). Repaired in place after the walk; the
# list stays empty on a clean library, so a no-change reconcile
# only pays one string compare per cached file.
stale_paths: List[str] = []
stale_seen: Set[str] = set()
new_files: List[str] = []
discovered_real_files: Set[str] = set()
def mark_stale_if_needed(cached_path: str) -> None:
"""Queue a cached path for file_name repair when it drifted."""
if cached_path in stale_seen:
return
item = path_to_item.get(cached_path)
if item is None:
return
if item.get("file_name") == _file_name_stem(cached_path):
return
stale_seen.add(cached_path)
stale_paths.append(cached_path)
visited_real_paths = set()
discovered_real_files = set()
discovered_folders: Set[str] = set()
# Scan all model roots
for root_path in self.get_model_roots():
if not os.path.exists(root_path):
continue
# Recursively scan directory
for root, dirnames, files in os.walk(root_path, followlinks=True):
dirnames[:] = [d for d in dirnames if not _is_excluded_dir(d)]
real_root = os.path.realpath(root)
if real_root in visited_real_paths:
for result in walk_results:
found_paths |= result.found_paths
discovered_folders |= result.discovered_folders
for cached_path in result.stale_paths:
if cached_path in stale_seen:
continue
visited_real_paths.add(real_root)
stale_seen.add(cached_path)
stale_paths.append(cached_path)
# Record every visited directory (including empty ones) so
# the folder tree stays accurate without a live walk.
rel_dir = os.path.relpath(
os.path.abspath(root), os.path.abspath(root_path)
).replace(os.path.sep, "/")
if rel_dir != "." and not _is_hidden_relative_path(rel_dir):
discovered_folders.add(rel_dir)
for file in files:
ext = os.path.splitext(file)[1].lower()
if ext in self.file_extensions:
# Construct paths exactly as they would be in cache
file_path = os.path.join(root, file).replace(os.sep, '/')
# Check if this file is already in cache
if file_path in cached_paths:
found_paths.add(file_path)
mark_stale_if_needed(file_path)
continue
# Only a cache miss needs the physical path, so the
# realpath syscalls are paid per changed file rather
# than per file in the library.
real_file_path = os.path.realpath(os.path.join(root, file))
cached_real_match = lookup_cached_real_path(real_file_path)
if cached_real_match:
found_paths.add(cached_real_match)
mark_stale_if_needed(cached_real_match)
continue
if file_path in self._excluded_models:
continue
# Try case-insensitive match on Windows
if os.name == 'nt':
lower_path = file_path.lower()
matched = False
for cached_path in cached_paths:
if cached_path.lower() == lower_path:
found_paths.add(cached_path)
mark_stale_if_needed(cached_path)
matched = True
break
if matched:
continue
if real_file_path in discovered_real_files:
continue
discovered_real_files.add(real_file_path)
# This is a new file to process
new_files.append(file_path)
# Yield control periodically
await asyncio.sleep(0)
if self.is_cancelled():
logger.info(f"{self.model_type.capitalize()} Scanner: Reconcile scan cancelled")
await self._broadcast_scan_progress(
'cancelled', 'reconcile_scan', 0, False,
elapsed_seconds=time.time() - start_time,
)
return
for result in walk_results:
for file_path, real_file_path in result.new_candidates:
if real_file_path in discovered_real_files:
continue
discovered_real_files.add(real_file_path)
# This is a new file to process
new_files.append(file_path)
# Process new files in batches
total_added = 0
@@ -1281,7 +1611,7 @@ class ModelScanner:
last_progress_time = current_time
await self._broadcast_scan_progress(
'processing', 'process_new',
min(99, int(1 + (processed_new / total_new) * 98)), False,
_new_file_pass_progress(processed_new, total_new), False,
processed=processed_new, total=total_new,
current_name=os.path.basename(path),
)
@@ -1290,7 +1620,7 @@ class ModelScanner:
logger.info(f"{self.model_type.capitalize()} Scanner: Reconcile processing cancelled")
await self._broadcast_scan_progress(
'cancelled', 'process_new',
min(99, int(1 + (processed_new / total_new) * 98)), False,
_new_file_pass_progress(processed_new, total_new), False,
elapsed_seconds=time.time() - start_time,
)
return
@@ -1444,7 +1774,117 @@ class ModelScanner:
finally:
self._is_initializing = False # Unset flag
self.bump_cache_version()
async def _walk_roots_for_reconcile(
self,
*,
roots: List[str],
tracker: _ReconcileWalkTracker,
cached_paths: Set[str],
path_to_item: Mapping[str, Dict[str, Any]],
) -> List[_RootWalkResult]:
"""Walk every root off the event loop and return results in root order.
One worker per device: cold or slow drives no longer serialize the
others, and the event loop stays free for the whole walk (previously it
was re-entered once per directory). Results are re-ordered to the
configured root order so the merge below is deterministic.
"""
if not roots:
return []
lookups = _CachedPathLookups(cached_paths)
dir_claims = _RealDirClaims()
excluded_models = set(self._excluded_models)
loop = asyncio.get_running_loop()
monitor = asyncio.create_task(self._monitor_walk_progress(tracker))
try:
futures = [
loop.run_in_executor(
None,
self._walk_root_group_sync,
group,
tracker,
lookups,
dir_claims,
excluded_models,
path_to_item,
)
for group in _group_roots_by_device(roots)
]
outcomes = await asyncio.gather(*futures, return_exceptions=True)
finally:
monitor.cancel()
with suppress(asyncio.CancelledError, Exception):
await monitor
for outcome in outcomes:
if isinstance(outcome, BaseException):
raise outcome
by_root: Dict[str, _RootWalkResult] = {}
for outcome in outcomes:
for result in cast(List[_RootWalkResult], outcome):
by_root[result.root_path] = result
return [by_root[root] for root in roots if root in by_root]
def _walk_root_group_sync(
self,
roots: List[str],
tracker: _ReconcileWalkTracker,
lookups: _CachedPathLookups,
dir_claims: _RealDirClaims,
excluded_models: Set[str],
path_to_item: Mapping[str, Dict[str, Any]],
) -> List[_RootWalkResult]:
"""Walk the roots of one device sequentially (worker-thread entry point)."""
results: List[_RootWalkResult] = []
for root_path in roots:
if self.is_cancelled():
break
tracker.start_root(root_path)
def report_progress(files_seen: int, _root: str = root_path) -> None:
tracker.report(_root, files_seen)
try:
result = _walk_root_for_reconcile(
root_path=root_path,
file_extensions=self.file_extensions,
cached_paths=lookups.cached_paths,
path_to_item=path_to_item,
lookups=lookups,
dir_claims=dir_claims,
excluded_models=excluded_models,
is_cancelled=self.is_cancelled,
report_progress=report_progress,
)
finally:
tracker.finish_root(root_path)
results.append(result)
if result.cancelled:
break
return results
async def _monitor_walk_progress(self, tracker: _ReconcileWalkTracker) -> None:
"""Broadcast walk progress while the root workers are running."""
while True:
await asyncio.sleep(_WALK_PROGRESS_INTERVAL_SECONDS)
await self._broadcast_walk_progress(tracker)
async def _broadcast_walk_progress(self, tracker: _ReconcileWalkTracker) -> None:
"""Emit one walk-phase progress message (best effort)."""
snapshot = tracker.snapshot()
if not snapshot:
return
progress = int(snapshot.pop('progress', 0))
await self._broadcast_scan_progress(
'processing', 'reconcile_scan', progress, False, **snapshot
)
def is_initializing(self) -> bool:
"""Check if the scanner is currently initializing"""
return self._is_initializing