Make rescans robust: DB-aware mtime gating, wipe-proof discovery
- RootScanner takes indexed_paths (store torrent keys); a candidate is skipped only when its mtime is unchanged AND it is present in the index, so an emptied/partial index self-heals on the next pass (non-indexable OTHER candidates still skip on mtime alone) - discovery: a directory that cannot be listed (transient SMB failure) now carries its previously known paths into the Sync set instead of dropping them, and an unreadable media root fails the whole scan — previously both produced a partial/empty known_paths that made sync_torrent_paths delete items that still exist on disk - worker: broad per-directory exception guard (carry known paths, stay hot, keep other workers running); node completion hardened against double-decrement; backoff state not updated from cancelled walks
This commit is contained in:
@@ -82,10 +82,14 @@ class RootScanner:
|
|||||||
root_id: str,
|
root_id: str,
|
||||||
media_root: Path,
|
media_root: Path,
|
||||||
send: Send,
|
send: Send,
|
||||||
|
indexed_paths: Callable[[], set[str]] | None = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
self.root_id = root_id
|
self.root_id = root_id
|
||||||
self.media_root = media_root
|
self.media_root = media_root
|
||||||
self._send = send
|
self._send = send
|
||||||
|
# Returns the torrent paths currently present in the index; used to
|
||||||
|
# reprocess items missing from the database despite unchanged mtimes.
|
||||||
|
self._indexed_paths = indexed_paths
|
||||||
self._output_dir = media_root / DEFAULT_OUTPUT_FOLDER
|
self._output_dir = media_root / DEFAULT_OUTPUT_FOLDER
|
||||||
self._scanignore = ScanIgnore(media_root)
|
self._scanignore = ScanIgnore(media_root)
|
||||||
|
|
||||||
@@ -391,6 +395,9 @@ class RootScanner:
|
|||||||
media_root_str = self.media_root.as_posix()
|
media_root_str = self.media_root.as_posix()
|
||||||
now = time.time()
|
now = time.time()
|
||||||
dirs_visited = 0
|
dirs_visited = 0
|
||||||
|
# Torrent paths currently in the index; a matching mtime alone is not
|
||||||
|
# enough to skip an item that the database does not actually have.
|
||||||
|
indexed = self._indexed_paths() if self._indexed_paths else None
|
||||||
|
|
||||||
media_container_dirs = {"BDMV", "VIDEO_TS", "HVDVD_TS"}
|
media_container_dirs = {"BDMV", "VIDEO_TS", "HVDVD_TS"}
|
||||||
video_extensions = {
|
video_extensions = {
|
||||||
@@ -469,7 +476,16 @@ class RootScanner:
|
|||||||
except OSError, ValueError:
|
except OSError, ValueError:
|
||||||
return False
|
return False
|
||||||
if self._seen_mtimes.get(relpath) == mtime:
|
if self._seen_mtimes.get(relpath) == mtime:
|
||||||
return False
|
if indexed is None or relpath in indexed:
|
||||||
|
return False
|
||||||
|
# Unchanged on disk but missing from the index (snapshot
|
||||||
|
# wiped, upsert lost, ...): reprocess it unless it is not
|
||||||
|
# indexable content anyway.
|
||||||
|
parsed = await parse_download(path)
|
||||||
|
if parsed.content_type is ContentType.OTHER:
|
||||||
|
return False
|
||||||
|
downloads.append(parsed)
|
||||||
|
return True
|
||||||
found_mtimes[relpath] = mtime
|
found_mtimes[relpath] = mtime
|
||||||
downloads.append(await parse_download(path))
|
downloads.append(await parse_download(path))
|
||||||
return True
|
return True
|
||||||
@@ -484,7 +500,7 @@ class RootScanner:
|
|||||||
def _complete_node(node: dict) -> None:
|
def _complete_node(node: dict) -> None:
|
||||||
"""Fold a finished subtree into the backoff state and its parent."""
|
"""Fold a finished subtree into the backoff state and its parent."""
|
||||||
rel = node["rel"]
|
rel = node["rel"]
|
||||||
if rel is not None:
|
if rel is not None and not stop_event.is_set():
|
||||||
prev = self._dir_state.get(rel)
|
prev = self._dir_state.get(rel)
|
||||||
if prev and prev.get("items") and not node["items"]:
|
if prev and prev.get("items") and not node["items"]:
|
||||||
# Items vanished from this subtree — stay hot so the
|
# Items vanished from this subtree — stay hot so the
|
||||||
@@ -553,6 +569,7 @@ class RootScanner:
|
|||||||
nonlocal dirs_visited
|
nonlocal dirs_visited
|
||||||
while True:
|
while True:
|
||||||
_, _, path, mtime, node = await queue.get()
|
_, _, path, mtime, node = await queue.get()
|
||||||
|
rel = make_relative_path(str(path), media_root_str)
|
||||||
try:
|
try:
|
||||||
if stop_event.is_set():
|
if stop_event.is_set():
|
||||||
continue
|
continue
|
||||||
@@ -565,7 +582,11 @@ class RootScanner:
|
|||||||
) = await asyncio.to_thread(_scan_dir, path, stop_event)
|
) = await asyncio.to_thread(_scan_dir, path, stop_event)
|
||||||
except OSError, PermissionError:
|
except OSError, PermissionError:
|
||||||
logger.debug("Cannot list directory: %s", path)
|
logger.debug("Cannot list directory: %s", path)
|
||||||
# Unreadable — don't let it earn backoff.
|
# Unreadable (e.g. a transient network-mount failure):
|
||||||
|
# keep its previously known paths so the Sync event
|
||||||
|
# cannot delete the subtree's items, and don't let it
|
||||||
|
# earn backoff.
|
||||||
|
_carry_known(rel)
|
||||||
node["changed"] = True
|
node["changed"] = True
|
||||||
child_dirs, video_files, is_media_container = [], [], False
|
child_dirs, video_files, is_media_container = [], [], False
|
||||||
if self._scanignore:
|
if self._scanignore:
|
||||||
@@ -580,7 +601,6 @@ class RootScanner:
|
|||||||
if not self._scanignore.is_excluded(f[0])
|
if not self._scanignore.is_excluded(f[0])
|
||||||
]
|
]
|
||||||
if dirs_visited % 8 == 1:
|
if dirs_visited % 8 == 1:
|
||||||
rel = make_relative_path(str(path), media_root_str) or str(path)
|
|
||||||
await _report(f"Scanning: {rel} ({len(downloads)} found)")
|
await _report(f"Scanning: {rel} ({len(downloads)} found)")
|
||||||
|
|
||||||
if is_media_container or not child_dirs:
|
if is_media_container or not child_dirs:
|
||||||
@@ -593,7 +613,6 @@ class RootScanner:
|
|||||||
# (SMB), so fold in the newest video file mtime —
|
# (SMB), so fold in the newest video file mtime —
|
||||||
# file mtimes are reliable, and this is what makes a
|
# file mtimes are reliable, and this is what makes a
|
||||||
# still-growing download show up immediately.
|
# still-growing download show up immediately.
|
||||||
rel = make_relative_path(str(path), media_root_str)
|
|
||||||
if is_media_container or video_files:
|
if is_media_container or video_files:
|
||||||
node["items"] = True
|
node["items"] = True
|
||||||
if video_files:
|
if video_files:
|
||||||
@@ -604,16 +623,27 @@ class RootScanner:
|
|||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
for file_path, file_mtime in video_files:
|
for file_path, file_mtime in video_files:
|
||||||
rel = make_relative_path(str(file_path), media_root_str)
|
file_rel = make_relative_path(
|
||||||
|
str(file_path), media_root_str
|
||||||
|
)
|
||||||
node["items"] = True
|
node["items"] = True
|
||||||
node["changed"] = node[
|
node["changed"] = node[
|
||||||
"changed"
|
"changed"
|
||||||
] or await _register_candidate(file_path, rel, file_mtime)
|
] or await _register_candidate(
|
||||||
|
file_path, file_rel, file_mtime
|
||||||
|
)
|
||||||
await _enqueue_children(node, child_dirs)
|
await _enqueue_children(node, child_dirs)
|
||||||
|
except Exception:
|
||||||
|
# Never let one bad directory kill the worker or truncate
|
||||||
|
# the walk — an incomplete known_paths set would make the
|
||||||
|
# Sync event delete items that still exist on disk.
|
||||||
|
logger.exception("Discovery failed for directory: %s", path)
|
||||||
|
_carry_known(rel)
|
||||||
|
node["changed"] = True
|
||||||
|
finally:
|
||||||
node["pending"] -= 1
|
node["pending"] -= 1
|
||||||
if node["pending"] == 0:
|
if node["pending"] == 0:
|
||||||
_complete_node(node)
|
_complete_node(node)
|
||||||
finally:
|
|
||||||
queue.task_done()
|
queue.task_done()
|
||||||
|
|
||||||
logger.info("Starting filesystem discovery at %s", self.media_root)
|
logger.info("Starting filesystem discovery at %s", self.media_root)
|
||||||
@@ -628,8 +658,10 @@ class RootScanner:
|
|||||||
stop_event.set()
|
stop_event.set()
|
||||||
raise
|
raise
|
||||||
except OSError, PermissionError:
|
except OSError, PermissionError:
|
||||||
logger.exception("Cannot list media root: %s", self.media_root)
|
# Failing to list the root must fail the whole scan: returning an
|
||||||
return downloads, found_mtimes, known_paths
|
# empty known_paths here would make the Sync event delete every
|
||||||
|
# item in the index.
|
||||||
|
raise RuntimeError(f"Cannot list media root: {self.media_root}")
|
||||||
|
|
||||||
if self._scanignore:
|
if self._scanignore:
|
||||||
child_dirs = [
|
child_dirs = [
|
||||||
|
|||||||
@@ -150,6 +150,22 @@ class IndexStore:
|
|||||||
return None
|
return None
|
||||||
return item.info.tmdb_id
|
return item.info.tmdb_id
|
||||||
|
|
||||||
|
def torrent_paths(self) -> set[str]:
|
||||||
|
"""All media-root-relative torrent paths currently in the index.
|
||||||
|
|
||||||
|
The scanner uses this to reprocess items that are missing from the
|
||||||
|
index even though their mtime is unchanged (e.g. after the snapshot
|
||||||
|
was wiped or an upsert never landed).
|
||||||
|
"""
|
||||||
|
paths: set[str] = set()
|
||||||
|
for movie in self.movies.values():
|
||||||
|
paths.update(movie.files)
|
||||||
|
for show in self.series.values():
|
||||||
|
for season in show.seasons:
|
||||||
|
for episode in season.episodes:
|
||||||
|
paths.update(episode.files)
|
||||||
|
return paths
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def _newest_from_files(files: dict[str, Torrent]) -> int | None:
|
def _newest_from_files(files: dict[str, Torrent]) -> int | None:
|
||||||
timestamps = [t.added_at for t in files.values() if t.added_at]
|
timestamps = [t.added_at for t in files.values() if t.added_at]
|
||||||
|
|||||||
+6
-1
@@ -652,7 +652,12 @@ async def _attach_scanners() -> None:
|
|||||||
for ctx in supervisor.all_contexts().values():
|
for ctx in supervisor.all_contexts().values():
|
||||||
if ctx.scanner is None and ctx.status == "ready":
|
if ctx.scanner is None and ctx.status == "ready":
|
||||||
try:
|
try:
|
||||||
scanner = RootScanner(ctx.root_id, ctx.root_path, ctx.send_event)
|
scanner = RootScanner(
|
||||||
|
ctx.root_id,
|
||||||
|
ctx.root_path,
|
||||||
|
ctx.send_event,
|
||||||
|
ctx.store.torrent_paths,
|
||||||
|
)
|
||||||
await scanner.start()
|
await scanner.start()
|
||||||
ctx.scanner = scanner
|
ctx.scanner = scanner
|
||||||
except Exception:
|
except Exception:
|
||||||
|
|||||||
Reference in New Issue
Block a user