mirror of
https://github.com/calibrain/shelfmark.git
synced 2026-10-04 19:31:13 +01:00
Add threading to file system operations (#602)
This commit is contained in:
+29
-13
@@ -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:
|
||||
|
||||
+69
-40
@@ -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}")
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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", "")
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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 = {
|
||||
|
||||
Reference in New Issue
Block a user