From 022e50a0ba1858ba1eb4a9e302fd9fe628650e2c Mon Sep 17 00:00:00 2001 From: Alex <25013571+alexhb1@users.noreply.github.com> Date: Mon, 9 Feb 2026 18:04:10 +0000 Subject: [PATCH] Add threading to file system operations (#602) --- shelfmark/core/logger.py | 42 ++++--- shelfmark/download/fs.py | 109 +++++++++++------- shelfmark/download/orchestrator.py | 5 +- shelfmark/download/permissions_debug.py | 31 +++-- shelfmark/download/postprocess/destination.py | 12 +- shelfmark/download/postprocess/scan.py | 12 +- shelfmark/download/postprocess/transfer.py | 49 +++++--- shelfmark/download/postprocess/workspace.py | 17 ++- shelfmark/download/staging.py | 15 +-- .../prowlarr/clients/qbittorrent.py | 24 ++++ shelfmark/release_sources/prowlarr/handler.py | 18 +-- tests/prowlarr/test_qbittorrent_client.py | 72 ++++++++++++ 12 files changed, 303 insertions(+), 103 deletions(-) diff --git a/shelfmark/core/logger.py b/shelfmark/core/logger.py index eb0cf985..e12b828d 100644 --- a/shelfmark/core/logger.py +++ b/shelfmark/core/logger.py @@ -39,22 +39,38 @@ class CustomLogger(logging.Logger): self.debug(msg, *args, exc_info=has_exception, **kwargs) def log_resource_usage(self): - import psutil + # Best-effort only; this should never raise during exception logging. + try: + import psutil - # Sum RSS of all processes for actual app memory - app_memory_mb = 0 - for proc in psutil.process_iter(['memory_info']): + # Sum RSS of all processes for actual app memory (container-friendly), + # but fall back gracefully on platforms that restrict process enumeration. + app_memory_mb = 0.0 try: - if proc.info['memory_info']: - app_memory_mb += proc.info['memory_info'].rss / (1024 * 1024) - except (psutil.NoSuchProcess, psutil.AccessDenied): - continue + for proc in psutil.process_iter(['memory_info']): + try: + mem = proc.info.get('memory_info') + if mem: + app_memory_mb += mem.rss / (1024 * 1024) + except (psutil.NoSuchProcess, psutil.AccessDenied, KeyError, AttributeError): + continue + except (PermissionError, psutil.AccessDenied, OSError): + try: + app_memory_mb = psutil.Process().memory_info().rss / (1024 * 1024) + except Exception: + app_memory_mb = 0.0 - memory = psutil.virtual_memory() - system_used_mb = memory.used / (1024 * 1024) - available_mb = memory.available / (1024 * 1024) - cpu_percent = psutil.cpu_percent() - self.debug(f"Container Memory: App={app_memory_mb:.2f} MB, System={system_used_mb:.2f} MB, Available={available_mb:.2f} MB, CPU: {cpu_percent:.2f}%") + memory = psutil.virtual_memory() + system_used_mb = memory.used / (1024 * 1024) + available_mb = memory.available / (1024 * 1024) + cpu_percent = psutil.cpu_percent() + self.debug( + f"Container Memory: App={app_memory_mb:.2f} MB, System={system_used_mb:.2f} MB, " + f"Available={available_mb:.2f} MB, CPU: {cpu_percent:.2f}%" + ) + except Exception: + # Avoid breaking the original log call if psutil is missing or restricted. + return def setup_logger(name: str, log_file: Path = LOG_FILE) -> CustomLogger: diff --git a/shelfmark/download/fs.py b/shelfmark/download/fs.py index 33454f93..9075b5bf 100644 --- a/shelfmark/download/fs.py +++ b/shelfmark/download/fs.py @@ -11,7 +11,7 @@ import subprocess import tempfile import time from pathlib import Path -from typing import Any, Callable, Optional, TypeVar +from typing import Any, Callable, Optional, TypeVar, cast from shelfmark.core.logger import setup_logger from shelfmark.download.permissions_debug import log_transfer_permission_context @@ -45,10 +45,27 @@ def _get_io_threadpool() -> "_GeventThreadPool": return _IO_THREADPOOL +def _call_and_capture(func: Callable[..., T], args: tuple[Any, ...], kwargs: dict[str, Any]) -> tuple[bool, T | Exception]: + try: + return True, func(*args, **kwargs) + except Exception as exc: + return False, exc + + def run_blocking_io(func: Callable[..., T], *args: Any, **kwargs: Any) -> T: - """Run blocking I/O in a native thread when under gevent.""" + """Run blocking I/O in a native thread when under gevent. + + gevent's threadpool will eagerly log exceptions raised inside worker threads, + even when the caller expects and handles those errors (e.g. FileExistsError for + collision retries, EXDEV for cross-device moves). Capture and re-raise in the + caller to avoid noisy, misleading tracebacks. + """ if _use_gevent_threadpool(): - return _get_io_threadpool().apply(func, args, kwds=kwargs) + ok, result = _get_io_threadpool().apply(_call_and_capture, (func, args, kwargs)) + if ok: + return cast(T, result) + exc = cast(Exception, result) + raise exc return func(*args, **kwargs) @@ -66,7 +83,8 @@ def _verify_transfer_size( Some filesystems (especially remote NAS/CIFS/NFS) can report stale sizes briefly after large writes. Do a second stat after a short delay before declaring failure. """ - actual_size = dest.stat().st_size + # On network filesystems, `stat()` can block long enough to starve the gevent hub. + actual_size = run_blocking_io(dest.stat).st_size if actual_size == expected_size: return @@ -76,7 +94,7 @@ def _verify_transfer_size( ) time.sleep(_VERIFY_IO_WAIT_SECONDS) - actual_size = dest.stat().st_size + actual_size = run_blocking_io(dest.stat).st_size if actual_size != expected_size: raise IOError( f"File {action} incomplete, data loss may have occurred. " @@ -109,11 +127,16 @@ def atomic_write(dest_path: Path, data: bytes, max_attempts: int = 100) -> Path: try_path = dest_path if attempt == 0 else parent / f"{base}_{attempt}{ext}" try: # O_CREAT | O_EXCL fails atomically if file exists - fd = os.open(str(try_path), os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o666) + fd = run_blocking_io( + os.open, + str(try_path), + os.O_CREAT | os.O_EXCL | os.O_WRONLY, + 0o666, + ) try: - os.write(fd, data) + run_blocking_io(os.write, fd, data) finally: - os.close(fd) + run_blocking_io(os.close, fd) if attempt > 0: logger.info(f"File collision resolved: {try_path.name}") return try_path @@ -142,7 +165,7 @@ def _system_op(op: str, source: Path, dest: Path) -> None: def _perform_nfs_fallback(source: Path, dest: Path, is_move: bool) -> None: """Handle NFS/SMB permission errors by falling back to copyfile -> system op.""" - expected_size = source.stat().st_size + expected_size = run_blocking_io(source.stat).st_size try: # Fallback 1: copy content only @@ -150,12 +173,12 @@ def _perform_nfs_fallback(source: Path, dest: Path, is_move: bool) -> None: _verify_transfer_size(dest, expected_size, "copy") if is_move: - source.unlink() + run_blocking_io(source.unlink) return except Exception as copy_error: # Clean up failed copy attempt if it exists - dest.unlink(missing_ok=True) + run_blocking_io(dest.unlink, missing_ok=True) if _is_permission_error(copy_error): log_transfer_permission_context("nfs_fallback_copyfile", source=source, dest=dest, error=copy_error) @@ -166,14 +189,14 @@ def _perform_nfs_fallback(source: Path, dest: Path, is_move: bool) -> None: try: _system_op(op, source, dest) # Best-effort verify after external command. - if dest.exists(): + if run_blocking_io(dest.exists): _verify_transfer_size(dest, expected_size, op) if is_move: - source.unlink(missing_ok=True) + run_blocking_io(source.unlink, missing_ok=True) except subprocess.CalledProcessError as sys_error: log_transfer_permission_context("nfs_fallback_system", source=source, dest=dest, error=sys_error) logger.error("System %s failed (%s -> %s): %s", op, source, dest, sys_error.stderr) - dest.unlink(missing_ok=True) + run_blocking_io(dest.unlink, missing_ok=True) raise @@ -183,11 +206,16 @@ def _claim_destination(path: Path) -> bool: Returns True if the placeholder was created. Caller must replace or unlink it. """ try: - fd = os.open(str(path), os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o666) + fd = run_blocking_io( + os.open, + str(path), + os.O_CREAT | os.O_EXCL | os.O_WRONLY, + 0o666, + ) except FileExistsError: return False else: - os.close(fd) + run_blocking_io(os.close, fd) return True @@ -205,12 +233,13 @@ def _hardlink_not_supported(error: OSError) -> bool: def _create_temp_path(dest_path: Path) -> Path: - fd, temp_path = tempfile.mkstemp( + fd, temp_path = run_blocking_io( + tempfile.mkstemp, prefix=f".{dest_path.name}.", suffix=".tmp", dir=str(dest_path.parent), ) - os.close(fd) + run_blocking_io(os.close, fd) return Path(temp_path) @@ -220,8 +249,8 @@ def _publish_temp_file(temp_path: Path, dest_path: Path) -> bool: Returns True on success, False if the destination already exists. """ try: - os.link(str(temp_path), str(dest_path)) - temp_path.unlink(missing_ok=True) + run_blocking_io(os.link, str(temp_path), str(dest_path)) + run_blocking_io(temp_path.unlink, missing_ok=True) return True except FileExistsError: return False @@ -244,9 +273,9 @@ def _publish_temp_file(temp_path: Path, dest_path: Path) -> bool: if not claimed: return False try: - os.replace(str(temp_path), str(dest_path)) + run_blocking_io(os.replace, str(temp_path), str(dest_path)) except Exception: - dest_path.unlink(missing_ok=True) + run_blocking_io(dest_path.unlink, missing_ok=True) raise return True raise @@ -282,7 +311,7 @@ def atomic_move(source_path: Path, dest_path: Path, max_attempts: int = 100) -> # Check for existing file (os.rename would overwrite on Unix) claimed = False - if try_path.exists(): + if run_blocking_io(try_path.exists): # Some filesystems can report false positives for exists() with # special characters. Probe with O_EXCL to confirm. claimed = _claim_destination(try_path) @@ -292,27 +321,27 @@ def atomic_move(source_path: Path, dest_path: Path, max_attempts: int = 100) -> try: # os.rename is atomic on same filesystem and triggers inotify events if claimed: - os.replace(str(source_path), str(try_path)) + run_blocking_io(os.replace, str(source_path), str(try_path)) else: - os.rename(str(source_path), str(try_path)) + run_blocking_io(os.rename, str(source_path), str(try_path)) if attempt > 0: logger.info(f"File collision resolved: {try_path.name}") return try_path except FileExistsError: # Race condition: file created between exists() check and rename() if claimed: - try_path.unlink(missing_ok=True) + run_blocking_io(try_path.unlink, missing_ok=True) continue except OSError as e: # Cross-filesystem - copy to temp and publish atomically. if e.errno != errno.EXDEV: if claimed: - try_path.unlink(missing_ok=True) + run_blocking_io(try_path.unlink, missing_ok=True) raise - expected_size = source_path.stat().st_size + expected_size = run_blocking_io(source_path.stat).st_size if claimed: - try_path.unlink(missing_ok=True) + run_blocking_io(try_path.unlink, missing_ok=True) claimed = False temp_path: Optional[Path] = None @@ -336,16 +365,16 @@ def atomic_move(source_path: Path, dest_path: Path, max_attempts: int = 100) -> _verify_transfer_size(temp_path, expected_size, "move") published = _publish_temp_file(temp_path, try_path) if not published: - temp_path.unlink(missing_ok=True) + run_blocking_io(temp_path.unlink, missing_ok=True) continue try: _verify_transfer_size(try_path, expected_size, "move") except Exception: - try_path.unlink(missing_ok=True) + run_blocking_io(try_path.unlink, missing_ok=True) raise - source_path.unlink() + run_blocking_io(source_path.unlink) if attempt > 0: logger.info(f"File collision resolved: {try_path.name}") @@ -353,11 +382,11 @@ def atomic_move(source_path: Path, dest_path: Path, max_attempts: int = 100) -> except FileExistsError: if temp_path: - temp_path.unlink(missing_ok=True) + run_blocking_io(temp_path.unlink, missing_ok=True) continue except Exception: if temp_path: - temp_path.unlink(missing_ok=True) + run_blocking_io(temp_path.unlink, missing_ok=True) raise except (PermissionError, OSError) as e: @@ -413,7 +442,7 @@ def atomic_hardlink(source_path: Path, dest_path: Path, max_attempts: int = 100) for attempt in range(max_attempts): try_path = dest_path if attempt == 0 else parent / f"{base}_{attempt}{ext}" try: - os.link(str(source_path), str(try_path)) + run_blocking_io(os.link, str(source_path), str(try_path)) if attempt > 0: logger.info(f"File collision resolved: {try_path.name}") return try_path @@ -460,11 +489,11 @@ def atomic_copy(source_path: Path, dest_path: Path, max_attempts: int = 100) -> base = dest_path.stem ext = dest_path.suffix parent = dest_path.parent - expected_size = source_path.stat().st_size + expected_size = run_blocking_io(source_path.stat).st_size for attempt in range(max_attempts): try_path = dest_path if attempt == 0 else parent / f"{base}_{attempt}{ext}" - if try_path.exists(): + if run_blocking_io(try_path.exists): continue temp_path: Optional[Path] = None try: @@ -502,13 +531,13 @@ def atomic_copy(source_path: Path, dest_path: Path, max_attempts: int = 100) -> _verify_transfer_size(temp_path, expected_size, "copy") published = _publish_temp_file(temp_path, try_path) if not published: - temp_path.unlink(missing_ok=True) + run_blocking_io(temp_path.unlink, missing_ok=True) continue try: _verify_transfer_size(try_path, expected_size, "copy") except Exception: - try_path.unlink(missing_ok=True) + run_blocking_io(try_path.unlink, missing_ok=True) raise if attempt > 0: @@ -516,7 +545,7 @@ def atomic_copy(source_path: Path, dest_path: Path, max_attempts: int = 100) -> return try_path except Exception: if temp_path: - temp_path.unlink(missing_ok=True) + run_blocking_io(temp_path.unlink, missing_ok=True) raise raise RuntimeError(f"Could not copy file after {max_attempts} attempts: {dest_path}") diff --git a/shelfmark/download/orchestrator.py b/shelfmark/download/orchestrator.py index dc6b7d4c..b6fa0c82 100644 --- a/shelfmark/download/orchestrator.py +++ b/shelfmark/download/orchestrator.py @@ -18,6 +18,7 @@ from shelfmark.core.logger import setup_logger from shelfmark.core.models import BookInfo, DownloadTask, QueueStatus, SearchFilters, SearchMode from shelfmark.core.queue import book_queue from shelfmark.core.utils import transform_cover_url +from shelfmark.download.fs import run_blocking_io from shelfmark.download.postprocess.pipeline import is_torrent_source, safe_cleanup_path from shelfmark.download.postprocess.router import post_process_download from shelfmark.release_sources import direct_download, get_handler, get_source_display_name @@ -184,7 +185,7 @@ def queue_status() -> Dict[str, Dict[str, Any]]: status = book_queue.get_status() for _, tasks in status.items(): for _, task in tasks.items(): - if task.download_path and not os.path.exists(task.download_path): + if task.download_path and not run_blocking_io(os.path.exists, task.download_path): task.download_path = None # Convert Enum keys to strings and DownloadTask objects to dicts for JSON serialization @@ -295,7 +296,7 @@ def _download_task(task_id: str, cancel_flag: Event) -> Optional[str]: return None temp_file = Path(temp_path) - if not temp_file.exists(): + if not run_blocking_io(temp_file.exists): logger.error(f"Handler returned non-existent path: {temp_path}") return None diff --git a/shelfmark/download/permissions_debug.py b/shelfmark/download/permissions_debug.py index 6e5ecca0..4da53980 100644 --- a/shelfmark/download/permissions_debug.py +++ b/shelfmark/download/permissions_debug.py @@ -16,6 +16,23 @@ from shelfmark.core.logger import setup_logger logger = setup_logger(__name__) +def _run_io(func, *args, **kwargs): + """Best-effort offload for potentially blocking filesystem calls. + + Keep this module import-cycle safe: `shelfmark.download.fs` imports this module, + so we only import `run_blocking_io` lazily at call-time. + """ + try: + from shelfmark.download.fs import run_blocking_io as _run_blocking_io + except Exception: + return func(*args, **kwargs) + + try: + return _run_blocking_io(func, *args, **kwargs) + except Exception: + # Fall back to direct call if threadpool offload is unavailable. + return func(*args, **kwargs) + def _format_uid(uid: int) -> str: try: @@ -59,12 +76,12 @@ def log_path_permission_context(label: str, path: Path) -> None: for probe in [path, path.parent]: try: - resolved = probe.resolve() + resolved = _run_io(probe.resolve) except Exception: resolved = probe try: - st = probe.stat() + st = _run_io(probe.stat) logger.debug( "Path permissions (%s): path=%s resolved=%s mode=%s owner=%s(%d) group=%s(%d) dir=%s symlink=%s", label, @@ -75,8 +92,8 @@ def log_path_permission_context(label: str, path: Path) -> None: st.st_uid, _format_gid(st.st_gid), st.st_gid, - probe.is_dir(), - probe.is_symlink(), + _run_io(probe.is_dir), + _run_io(probe.is_symlink), ) except Exception as stat_error: logger.debug("Path permissions (%s): stat failed for %s: %s", label, probe, stat_error) @@ -106,7 +123,7 @@ def log_transfer_permission_context(label: str, source: Path, dest: Path, error: for probe in [source, dest, dest.parent]: try: - st = probe.stat() + st = _run_io(probe.stat) logger.debug( "Path permissions (%s): path=%s mode=%s owner=%s(%d) group=%s(%d) exists=%s dir=%s", label, @@ -116,8 +133,8 @@ def log_transfer_permission_context(label: str, source: Path, dest: Path, error: st.st_uid, _format_gid(st.st_gid), st.st_gid, - probe.exists(), - probe.is_dir(), + _run_io(probe.exists), + _run_io(probe.is_dir), ) except Exception as stat_error: logger.debug("Path permissions (%s): stat failed for %s: %s", label, probe, stat_error) diff --git a/shelfmark/download/postprocess/destination.py b/shelfmark/download/postprocess/destination.py index e08f6df7..deae8b95 100644 --- a/shelfmark/download/postprocess/destination.py +++ b/shelfmark/download/postprocess/destination.py @@ -10,6 +10,7 @@ from shelfmark.core.utils import ( get_destination, is_audiobook as check_audiobook, ) +from shelfmark.download.fs import run_blocking_io from shelfmark.download.permissions_debug import log_path_permission_context logger = setup_logger("shelfmark.download.postprocess.pipeline") @@ -23,14 +24,15 @@ def validate_destination(destination: Path, status_callback) -> bool: status_callback("error", f"Destination must be absolute: {destination}") return False - if destination.exists() and not destination.is_dir(): + destination_exists = run_blocking_io(destination.exists) + if destination_exists and not run_blocking_io(destination.is_dir): logger.warning(f"Destination is not a directory: {destination}") status_callback("error", f"Destination is not a directory: {destination}") return False - if not destination.exists(): + if not destination_exists: try: - destination.mkdir(parents=True, exist_ok=True) + run_blocking_io(destination.mkdir, parents=True, exist_ok=True) except (OSError, PermissionError) as exc: log_path_permission_context("destination_create", destination) logger.warning(f"Cannot create destination: {destination} ({exc})") @@ -44,8 +46,8 @@ def validate_destination(destination: Path, status_callback) -> bool: f"This file was created to verify if '{destination}' is writable. " "It should've been automatically deleted. Feel free to delete it.\n" ) - test_path.write_text(test_content) - test_path.unlink(missing_ok=True) + run_blocking_io(test_path.write_text, test_content) + run_blocking_io(test_path.unlink, missing_ok=True) except Exception as exc: logger.debug("Destination write probe path: %s", test_path) log_path_permission_context("destination_write_probe", destination) diff --git a/shelfmark/download/postprocess/scan.py b/shelfmark/download/postprocess/scan.py index a5a45c89..18ef8419 100644 --- a/shelfmark/download/postprocess/scan.py +++ b/shelfmark/download/postprocess/scan.py @@ -80,7 +80,7 @@ def extract_archive_files( ) if cleanup_archive: - archive_path.unlink(missing_ok=True) + run_blocking_io(archive_path.unlink, missing_ok=True) cleanup_paths = [output_dir] @@ -107,8 +107,12 @@ def scan_directory_tree( """Scan a directory tree for book files, trackable-but-unsupported files, and archives.""" try: - with os.scandir(directory) as it: - next(it, None) + def _probe_dir() -> None: + # Force a fast error if the dir is missing/inaccessible. + with os.scandir(directory) as it: + next(it, None) + + run_blocking_io(_probe_dir) except PermissionError as exc: log_path_permission_context("scan_directory", directory) logger.warning(f"Permission denied scanning directory: {directory} ({exc})") @@ -280,7 +284,7 @@ def collect_staged_files( status_callback, cleanup_archives: bool, ) -> Tuple[List[Path], List[Path], List[Path], Optional[str]]: - if working_path.is_dir(): + if run_blocking_io(working_path.is_dir): if status_callback: status_callback("resolving", "Processing download folder") return collect_directory_files( diff --git a/shelfmark/download/postprocess/transfer.py b/shelfmark/download/postprocess/transfer.py index 37da2a39..4db11869 100644 --- a/shelfmark/download/postprocess/transfer.py +++ b/shelfmark/download/postprocess/transfer.py @@ -15,7 +15,7 @@ from shelfmark.core.naming import ( sanitize_filename, ) from shelfmark.core.utils import is_audiobook as check_audiobook -from shelfmark.download.fs import atomic_copy, atomic_hardlink, atomic_move +from shelfmark.download.fs import atomic_copy, atomic_hardlink, atomic_move, run_blocking_io from shelfmark.download.postprocess.policy import get_file_organization, get_template from .scan import collect_directory_files, scan_directory_tree @@ -70,10 +70,11 @@ def resolve_hardlink_source( if hardlink_enabled and task.original_download_path: hardlink_source = Path(task.original_download_path) - if destination and hardlink_source.exists() and same_filesystem(hardlink_source, destination): + hardlink_source_exists = run_blocking_io(hardlink_source.exists) + if destination and hardlink_source_exists and run_blocking_io(same_filesystem, hardlink_source, destination): use_hardlink = True source_path = hardlink_source - elif hardlink_source.exists(): + elif hardlink_source_exists: logger.warning( f"Cannot hardlink: {hardlink_source} and {destination} are on different filesystems. " "Falling back to copy. To fix: ensure torrent client downloads to same filesystem as destination." @@ -97,7 +98,7 @@ def is_torrent_source(source_path: Path, task: DownloadTask) -> bool: original_path = Path(task.original_download_path) try: - return source_path.resolve() == original_path.resolve() + return run_blocking_io(source_path.resolve) == run_blocking_io(original_path.resolve) except (OSError, ValueError): try: return os.path.normpath(str(source_path)) == os.path.normpath(str(original_path)) @@ -122,7 +123,7 @@ def _transfer_single_file( if use_hardlink: final_path = atomic_hardlink(source_path, dest_path, max_attempts=max_attempts) try: - if os.stat(source_path).st_ino == os.stat(final_path).st_ino: + if run_blocking_io(source_path.stat).st_ino == run_blocking_io(final_path.stat).st_ino: return final_path, "hardlink" except OSError: return final_path, "hardlink" @@ -160,8 +161,14 @@ def transfer_book_files( if len(book_files) == 1: source_file = book_files[0] ext = source_file.suffix.lstrip(".") or task.format or "" - dest_path = build_library_path(str(destination), template, metadata, extension=ext or None) - dest_path.parent.mkdir(parents=True, exist_ok=True) + dest_path = run_blocking_io( + build_library_path, + str(destination), + template, + metadata, + extension=ext or None, + ) + run_blocking_io(dest_path.parent.mkdir, parents=True, exist_ok=True) final_path, op = _transfer_single_file( source_file, @@ -181,8 +188,14 @@ def transfer_book_files( for source_file, part_number in files_with_parts: ext = source_file.suffix.lstrip(".") or task.format or "" file_metadata = {**metadata, "PartNumber": part_number} - dest_path = build_library_path(str(destination), template, file_metadata, extension=ext or None) - dest_path.parent.mkdir(parents=True, exist_ok=True) + dest_path = run_blocking_io( + build_library_path, + str(destination), + template, + file_metadata, + extension=ext or None, + ) + run_blocking_io(dest_path.parent.mkdir, parents=True, exist_ok=True) final_path, op = _transfer_single_file( source_file, @@ -297,8 +310,8 @@ def transfer_file_to_library( use_hardlink: bool, ) -> Optional[str]: extension = source_path.suffix.lstrip(".") or task.format - dest_path = build_library_path(library_base, template, metadata, extension) - dest_path.parent.mkdir(parents=True, exist_ok=True) + dest_path = run_blocking_io(build_library_path, library_base, template, metadata, extension) + run_blocking_io(dest_path.parent.mkdir, parents=True, exist_ok=True) is_torrent = is_torrent_source(source_path, task) final_path, op = _transfer_single_file( @@ -349,8 +362,14 @@ def transfer_directory_to_library( safe_cleanup_path(temp_file, task) return None - base_library_path = build_library_path(library_base, template, metadata, extension=None) - base_library_path.parent.mkdir(parents=True, exist_ok=True) + base_library_path = run_blocking_io( + build_library_path, + library_base, + template, + metadata, + extension=None, + ) + run_blocking_io(base_library_path.parent.mkdir, parents=True, exist_ok=True) is_torrent = is_torrent_source(source_dir, task) transferred_paths: List[Path] = [] @@ -378,8 +397,8 @@ def transfer_directory_to_library( for source_file, part_number in files_with_parts: ext = source_file.suffix.lstrip(".") file_metadata = {**metadata, "PartNumber": part_number} - file_path = build_library_path(library_base, template, file_metadata, extension=ext) - file_path.parent.mkdir(parents=True, exist_ok=True) + file_path = run_blocking_io(build_library_path, library_base, template, file_metadata, extension=ext) + run_blocking_io(file_path.parent.mkdir, parents=True, exist_ok=True) final_path, op = _transfer_single_file( source_file, diff --git a/shelfmark/download/postprocess/workspace.py b/shelfmark/download/postprocess/workspace.py index a1da29f7..43198905 100644 --- a/shelfmark/download/postprocess/workspace.py +++ b/shelfmark/download/postprocess/workspace.py @@ -22,8 +22,20 @@ def _tmp_dir() -> Path: def is_within_tmp_dir(path: Path) -> bool: """Legacy helper: True if path is inside TMP_DIR.""" + # Fast path: avoid `resolve()` (can block on NFS) for obviously-non-TMP paths. + # This is a *negative* check only; for potential TMP paths we still resolve to + # prevent symlink escapes from being treated as managed. + tmp_dir = _tmp_dir() try: - path.resolve().relative_to(_tmp_dir().resolve()) + if path.is_absolute() and tmp_dir.is_absolute(): + if path != tmp_dir and tmp_dir not in path.parents: + return False + except Exception: + # Fall back to the slower resolve-based check below. + pass + + try: + run_blocking_io(path.resolve).relative_to(run_blocking_io(tmp_dir.resolve)) return True except (OSError, ValueError): return False @@ -43,7 +55,8 @@ def _is_original_download(path: Optional[Path], task: DownloadTask) -> bool: if not path or not task.original_download_path: return False try: - return path.resolve() == Path(task.original_download_path).resolve() + original = Path(task.original_download_path) + return run_blocking_io(path.resolve) == run_blocking_io(original.resolve) except (OSError, ValueError): return False diff --git a/shelfmark/download/staging.py b/shelfmark/download/staging.py index da6075eb..f566458e 100644 --- a/shelfmark/download/staging.py +++ b/shelfmark/download/staging.py @@ -20,7 +20,7 @@ STAGE_MOVE: StageAction = "move" def get_staging_dir() -> Path: """Get the staging directory for downloads.""" tmp_dir = env_config.TMP_DIR - tmp_dir.mkdir(parents=True, exist_ok=True) + run_blocking_io(tmp_dir.mkdir, parents=True, exist_ok=True) return tmp_dir @@ -41,11 +41,11 @@ def build_staging_dir(prefix: str | None, task_id: str) -> Path: staging_dir = base_dir / f"{prefix}_{safe_id}" counter = 1 - while staging_dir.exists(): + while run_blocking_io(staging_dir.exists): staging_dir = base_dir / f"{prefix}_{safe_id}_{counter}" counter += 1 - staging_dir.mkdir(parents=True, exist_ok=True) + run_blocking_io(staging_dir.mkdir, parents=True, exist_ok=True) return staging_dir @@ -63,8 +63,9 @@ def stage_path(source: Path, staging_dir: Path, action: StageAction) -> Path: staged_path = staging_dir / source.name counter = 1 - if source.is_dir(): - while staged_path.exists(): + source_is_dir = run_blocking_io(source.is_dir) + if source_is_dir: + while run_blocking_io(staged_path.exists): staged_path = staging_dir / f"{source.name}_{counter}" counter += 1 if action == STAGE_COPY: @@ -72,7 +73,7 @@ def stage_path(source: Path, staging_dir: Path, action: StageAction) -> Path: else: run_blocking_io(shutil.move, str(source), str(staged_path)) else: - while staged_path.exists(): + while run_blocking_io(staged_path.exists): staged_path = staging_dir / f"{source.stem}_{counter}{source.suffix}" counter += 1 if action == STAGE_COPY: @@ -80,6 +81,6 @@ def stage_path(source: Path, staging_dir: Path, action: StageAction) -> Path: else: run_blocking_io(shutil.move, str(source), str(staged_path)) - staged_kind = "directory" if source.is_dir() else "file" + staged_kind = "directory" if source_is_dir else "file" logger.debug("Staged %s via %s: %s -> %s", staged_kind, action, source, staged_path) return staged_path diff --git a/shelfmark/release_sources/prowlarr/clients/qbittorrent.py b/shelfmark/release_sources/prowlarr/clients/qbittorrent.py index 2f48cbd4..ffaaaf10 100644 --- a/shelfmark/release_sources/prowlarr/clients/qbittorrent.py +++ b/shelfmark/release_sources/prowlarr/clients/qbittorrent.py @@ -1,6 +1,7 @@ """qBittorrent download client for Prowlarr integration.""" import time +from pathlib import Path from types import SimpleNamespace from typing import Optional, Tuple @@ -463,14 +464,37 @@ class QBittorrentClient(DownloadClient): Centralizes the logic shared by `get_status()` and `get_download_path()`: - accept `content_path` only when it's not equal to `save_path` + - when the torrent is complete and both `content_path` and `save_path` are present, + prefer a path rooted at `save_path` to avoid races where qBittorrent briefly reports + a temp/incomplete `content_path` and then moves the payload - otherwise derive via properties+files - finally fall back to `save_path + name` """ + torrent_progress = getattr(torrent, "progress", 0.0) + try: + progress = float(torrent_progress) + except (TypeError, ValueError): + progress = 0.0 + # Prefer content_path, but treat content_path == save_path as invalid. content_path = getattr(torrent, "content_path", "") save_path = getattr(torrent, "save_path", "") if content_path and (not save_path or str(content_path) != str(save_path)): + # When using a temp/incomplete directory, qBittorrent can briefly keep reporting + # `content_path` under that temp path right at completion, then move the files + # into `save_path`. Returning the temp path can race with that move. + if save_path and progress >= 1.0: + # Use the basename of content_path under save_path (works for single-file + # torrents and multi-file torrents where content_path is a top-level dir). + try: + content_basename = str(Path(str(content_path)).name) + except Exception: + content_basename = "" + rooted = self._build_path(str(save_path), content_basename) + if rooted: + return rooted + return str(content_path) download_id = getattr(torrent, "hash", "") diff --git a/shelfmark/release_sources/prowlarr/handler.py b/shelfmark/release_sources/prowlarr/handler.py index cb67f515..c3570d25 100644 --- a/shelfmark/release_sources/prowlarr/handler.py +++ b/shelfmark/release_sources/prowlarr/handler.py @@ -10,6 +10,7 @@ from shelfmark.core.config import config from shelfmark.core.logger import setup_logger from shelfmark.core.models import DownloadTask from shelfmark.core.utils import is_audiobook +from shelfmark.download.fs import run_blocking_io from shelfmark.release_sources import DownloadHandler, register_handler from shelfmark.release_sources.prowlarr.cache import get_release, remove_release from shelfmark.release_sources.prowlarr.clients import ( @@ -159,15 +160,15 @@ class ProwlarrHandler(DownloadHandler): logger.warning(f"Refusing to delete unsafe path for {client.name} {download_id}: {delete_path}") return - if not delete_path.exists(): + if not run_blocking_io(delete_path.exists): logger.debug(f"Local download path does not exist for cleanup: {delete_path}") return try: - if delete_path.is_dir(): - shutil.rmtree(delete_path) + if run_blocking_io(delete_path.is_dir): + run_blocking_io(shutil.rmtree, delete_path) else: - delete_path.unlink() + run_blocking_io(delete_path.unlink) logger.info(f"Deleted local download data for {client.name} {download_id}: {delete_path}") except Exception as e: logger.warning(f"Failed to delete local download data for {client.name} {download_id}: {e}") @@ -284,17 +285,18 @@ class ProwlarrHandler(DownloadHandler): ) if log_details: + remapped_exists = run_blocking_io(remapped.exists) logger.debug( "Remap result: %s -> %s (exists=%s, changed=%s, matched=%s)", source_path_obj, remapped, - remapped.exists(), + remapped_exists, remapped != source_path_obj, matched_mapping, ) if matched_mapping: - if remapped.exists(): + if run_blocking_io(remapped.exists): logger.info( "Remapped download path for %s (%s): %s -> %s", client.name, @@ -321,7 +323,7 @@ class ProwlarrHandler(DownloadHandler): return None, message if mappings: - if source_path_obj.exists(): + if run_blocking_io(source_path_obj.exists): logger.info( "No remote path mapping matched for %s (%s); using client path: %s", client.name, @@ -344,7 +346,7 @@ class ProwlarrHandler(DownloadHandler): ) return None, message - if not source_path_obj.exists(): + if not run_blocking_io(source_path_obj.exists): hint = _diagnose_path_issue(raw_path) message = hint if log_details: diff --git a/tests/prowlarr/test_qbittorrent_client.py b/tests/prowlarr/test_qbittorrent_client.py index bc15af76..02cf9294 100644 --- a/tests/prowlarr/test_qbittorrent_client.py +++ b/tests/prowlarr/test_qbittorrent_client.py @@ -244,6 +244,43 @@ class TestQBittorrentClientGetStatus: assert status.complete is True assert status.file_path == "/downloads/completed.epub" + def test_get_status_complete_roots_content_path_at_save_path(self, monkeypatch): + """Prefer a save_path-rooted path when qBittorrent reports a temp/incomplete content_path.""" + config_values = { + "QBITTORRENT_URL": "http://localhost:8080", + "QBITTORRENT_USERNAME": "admin", + "QBITTORRENT_PASSWORD": "password", + "QBITTORRENT_CATEGORY": "test", + } + monkeypatch.setattr( + "shelfmark.release_sources.prowlarr.clients.qbittorrent.config.get", + lambda key, default="": config_values.get(key, default), + ) + + mock_torrent = MockTorrent( + hash_val="abc123", + progress=1.0, + state="uploading", + content_path="/media/incomplete/book.m4b", + ) + + mock_client_instance = MagicMock() + # Include save_path in the info payload to simulate a temp/incomplete directory config. + info_payload = mock_torrent.to_dict() | {"save_path": "/media"} + mock_client_instance._session.get.return_value = create_mock_session_response([info_payload], status_code=200) + mock_client_class = MagicMock(return_value=mock_client_instance) + + with patch.dict('sys.modules', {'qbittorrentapi': MagicMock(Client=mock_client_class)}): + import importlib + import shelfmark.release_sources.prowlarr.clients.qbittorrent as qb_module + importlib.reload(qb_module) + + client = qb_module.QBittorrentClient() + status = client.get_status("abc123") + + assert status.complete is True + assert status.file_path == "/media/book.m4b" + def test_get_status_complete_derives_when_content_path_equals_save_path(self, monkeypatch): """Keep get_status() and get_download_path() consistent.""" config_values = { @@ -644,6 +681,41 @@ class TestQBittorrentClientGetDownloadPath: assert path == "/downloads/some/book.epub" + def test_get_download_path_roots_content_path_at_save_path_when_complete(self, monkeypatch): + """Mirror get_status(): completed torrents should return the save_path-rooted path.""" + config_values = { + "QBITTORRENT_URL": "http://localhost:8080", + "QBITTORRENT_USERNAME": "admin", + "QBITTORRENT_PASSWORD": "password", + "QBITTORRENT_CATEGORY": "test", + } + monkeypatch.setattr( + "shelfmark.release_sources.prowlarr.clients.qbittorrent.config.get", + lambda key, default="": config_values.get(key, default), + ) + + mock_torrent = MockTorrent( + hash_val="abc123", + progress=1.0, + state="uploading", + content_path="/media/incomplete/book.m4b", + ) + + mock_client_instance = MagicMock() + info_payload = mock_torrent.to_dict() | {"save_path": "/media"} + mock_client_instance._session.get.return_value = create_mock_session_response([info_payload], status_code=200) + mock_client_class = MagicMock(return_value=mock_client_instance) + + with patch.dict('sys.modules', {'qbittorrentapi': MagicMock(Client=mock_client_class)}): + import importlib + import shelfmark.release_sources.prowlarr.clients.qbittorrent as qb_module + importlib.reload(qb_module) + + client = qb_module.QBittorrentClient() + path = client.get_download_path("abc123") + + assert path == "/media/book.m4b" + def test_get_download_path_does_not_accept_content_path_equal_save_path(self, monkeypatch): """content_path == save_path indicates a path error.""" config_values = {