mirror of
https://github.com/calibrain/shelfmark.git
synced 2026-10-06 09:24:40 +01:00
Fix: Magnet and hash handling + Various bug fixes (#511)
This commit is contained in:
@@ -1231,6 +1231,16 @@ def advanced_settings():
|
||||
description="Path to a script to run after each successful download. Must be executable.",
|
||||
placeholder="/path/to/script.sh",
|
||||
),
|
||||
SelectField(
|
||||
key="CUSTOM_SCRIPT_PATH_MODE",
|
||||
label="Custom Script Path Mode",
|
||||
description="Pass the path to the custom script as an absolute path or relative to the destination folder.",
|
||||
options=[
|
||||
{"value": "absolute", "label": "Absolute", "description": "Pass the full destination path (default)."},
|
||||
{"value": "relative", "label": "Relative", "description": "Pass the path relative to the destination folder."},
|
||||
],
|
||||
default="absolute",
|
||||
),
|
||||
CheckboxField(
|
||||
key="DEBUG",
|
||||
label="Debug Mode",
|
||||
|
||||
@@ -141,6 +141,20 @@ def _perform_nfs_fallback(source: Path, dest: Path, is_move: bool) -> None:
|
||||
raise
|
||||
|
||||
|
||||
def _claim_destination(path: Path) -> bool:
|
||||
"""Atomically claim a destination path by creating a placeholder file.
|
||||
|
||||
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)
|
||||
except FileExistsError:
|
||||
return False
|
||||
else:
|
||||
os.close(fd)
|
||||
return True
|
||||
|
||||
|
||||
def atomic_move(source_path: Path, dest_path: Path, max_attempts: int = 100) -> Path:
|
||||
"""Move a file with collision detection.
|
||||
|
||||
@@ -170,29 +184,42 @@ def atomic_move(source_path: Path, dest_path: Path, max_attempts: int = 100) ->
|
||||
try_path = dest_path if attempt == 0 else parent / f"{base}_{attempt}{ext}"
|
||||
|
||||
# Check for existing file (os.rename would overwrite on Unix)
|
||||
claimed = False
|
||||
if try_path.exists():
|
||||
continue
|
||||
# Some filesystems can report false positives for exists() with
|
||||
# special characters. Probe with O_EXCL to confirm.
|
||||
claimed = _claim_destination(try_path)
|
||||
if not claimed:
|
||||
continue
|
||||
|
||||
try:
|
||||
# os.rename is atomic on same filesystem and triggers inotify events
|
||||
os.rename(str(source_path), str(try_path))
|
||||
if claimed:
|
||||
os.replace(str(source_path), str(try_path))
|
||||
else:
|
||||
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)
|
||||
continue
|
||||
except OSError as e:
|
||||
# Cross-filesystem - fall back to exclusive create + verified copy + delete.
|
||||
if e.errno != errno.EXDEV:
|
||||
if claimed:
|
||||
try_path.unlink(missing_ok=True)
|
||||
raise
|
||||
|
||||
expected_size = source_path.stat().st_size
|
||||
|
||||
try:
|
||||
# Claim destination path atomically.
|
||||
fd = os.open(str(try_path), os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o666)
|
||||
os.close(fd)
|
||||
if not claimed:
|
||||
# Claim destination path atomically.
|
||||
fd = os.open(str(try_path), os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o666)
|
||||
os.close(fd)
|
||||
|
||||
# Copy to a temp file first, then replace to avoid partial files.
|
||||
temp_path = try_path.parent / f".{try_path.name}.tmp"
|
||||
|
||||
@@ -20,6 +20,19 @@ logger = setup_logger(__name__)
|
||||
FOLDER_OUTPUT_MODE = "folder"
|
||||
|
||||
|
||||
def _resolve_custom_script_target(target_path: Path, destination: Path, path_mode: str) -> Path:
|
||||
mode = (path_mode or "absolute").strip().lower()
|
||||
if mode != "relative":
|
||||
return target_path
|
||||
|
||||
try:
|
||||
return target_path.relative_to(destination)
|
||||
except ValueError:
|
||||
if target_path.is_absolute():
|
||||
return Path(target_path.name)
|
||||
return target_path
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class _ProcessingPlan:
|
||||
destination: Path
|
||||
@@ -129,22 +142,41 @@ def process_folder_output(
|
||||
record_step(steps, step_name, source=str(temp_file), dest=str(prepared.output_plan.staging_dir))
|
||||
|
||||
def run_custom_script(script_path: str, target_path: Path, phase: str) -> bool:
|
||||
record_step(steps, "custom_script", script=str(script_path), target=str(target_path), phase=phase)
|
||||
path_mode = core_config.config.get("CUSTOM_SCRIPT_PATH_MODE", "absolute")
|
||||
script_target = _resolve_custom_script_target(target_path, plan.destination, path_mode)
|
||||
env = {
|
||||
**os.environ,
|
||||
"SHELFMARK_CUSTOM_SCRIPT_TARGET": str(target_path),
|
||||
"SHELFMARK_CUSTOM_SCRIPT_RELATIVE": str(_resolve_custom_script_target(target_path, plan.destination, "relative")),
|
||||
"SHELFMARK_CUSTOM_SCRIPT_DESTINATION": str(plan.destination),
|
||||
"SHELFMARK_CUSTOM_SCRIPT_MODE": str(path_mode),
|
||||
"SHELFMARK_CUSTOM_SCRIPT_PHASE": phase,
|
||||
}
|
||||
record_step(
|
||||
steps,
|
||||
"custom_script",
|
||||
script=str(script_path),
|
||||
target=str(script_target),
|
||||
target_abs=str(target_path),
|
||||
mode=str(path_mode),
|
||||
phase=phase,
|
||||
)
|
||||
log_plan_steps(task.task_id, steps)
|
||||
logger.info(
|
||||
"Task %s: running custom script %s on %s (%s)",
|
||||
task.task_id,
|
||||
script_path,
|
||||
target_path,
|
||||
script_target,
|
||||
phase,
|
||||
)
|
||||
try:
|
||||
result = subprocess.run(
|
||||
[script_path, str(target_path)],
|
||||
[script_path, str(script_target)],
|
||||
check=True,
|
||||
timeout=300, # 5 minute timeout
|
||||
capture_output=True,
|
||||
text=True,
|
||||
env=env,
|
||||
)
|
||||
if result.stdout:
|
||||
logger.debug("Task %s: custom script stdout: %s", task.task_id, result.stdout.strip())
|
||||
|
||||
@@ -105,15 +105,22 @@ def is_torrent_source(source_path: Path, task: DownloadTask) -> bool:
|
||||
return False
|
||||
|
||||
|
||||
def _max_attempts_for_batch(file_count: int, default: int = 100) -> int:
|
||||
if file_count <= 1:
|
||||
return default
|
||||
return max(default, file_count + default)
|
||||
|
||||
|
||||
def _transfer_single_file(
|
||||
source_path: Path,
|
||||
dest_path: Path,
|
||||
use_hardlink: bool,
|
||||
is_torrent: bool,
|
||||
preserve_source: bool = False,
|
||||
max_attempts: int = 100,
|
||||
) -> Tuple[Path, str]:
|
||||
if use_hardlink:
|
||||
final_path = atomic_hardlink(source_path, dest_path)
|
||||
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:
|
||||
return final_path, "hardlink"
|
||||
@@ -122,9 +129,9 @@ def _transfer_single_file(
|
||||
return final_path, "copy"
|
||||
|
||||
if is_torrent or preserve_source:
|
||||
return atomic_copy(source_path, dest_path), "copy"
|
||||
return atomic_copy(source_path, dest_path, max_attempts=max_attempts), "copy"
|
||||
|
||||
return atomic_move(source_path, dest_path), "move"
|
||||
return atomic_move(source_path, dest_path, max_attempts=max_attempts), "move"
|
||||
|
||||
|
||||
def transfer_book_files(
|
||||
@@ -141,6 +148,7 @@ def transfer_book_files(
|
||||
|
||||
is_audiobook = check_audiobook(task.content_type)
|
||||
organization_mode = organization_mode or get_file_organization(is_audiobook)
|
||||
max_attempts = _max_attempts_for_batch(len(book_files))
|
||||
|
||||
final_paths: List[Path] = []
|
||||
|
||||
@@ -160,6 +168,7 @@ def transfer_book_files(
|
||||
use_hardlink,
|
||||
is_torrent,
|
||||
preserve_source=preserve_source,
|
||||
max_attempts=max_attempts,
|
||||
)
|
||||
final_paths.append(final_path)
|
||||
logger.debug(f"{op.capitalize()} to destination: {final_path.name}")
|
||||
@@ -179,6 +188,7 @@ def transfer_book_files(
|
||||
use_hardlink,
|
||||
is_torrent,
|
||||
preserve_source=preserve_source,
|
||||
max_attempts=max_attempts,
|
||||
)
|
||||
final_paths.append(final_path)
|
||||
logger.debug(f"{op.capitalize()} to destination: {final_path.name}")
|
||||
@@ -210,6 +220,7 @@ def transfer_book_files(
|
||||
use_hardlink,
|
||||
is_torrent,
|
||||
preserve_source=preserve_source,
|
||||
max_attempts=max_attempts,
|
||||
)
|
||||
final_paths.append(final_path)
|
||||
logger.debug(f"{op.capitalize()} to destination: {final_path.name}")
|
||||
@@ -286,7 +297,13 @@ def transfer_file_to_library(
|
||||
dest_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
is_torrent = is_torrent_source(source_path, task)
|
||||
final_path, op = _transfer_single_file(source_path, dest_path, use_hardlink, is_torrent)
|
||||
final_path, op = _transfer_single_file(
|
||||
source_path,
|
||||
dest_path,
|
||||
use_hardlink,
|
||||
is_torrent,
|
||||
max_attempts=_max_attempts_for_batch(1),
|
||||
)
|
||||
logger.info(f"Library {op}: {final_path}")
|
||||
|
||||
if use_hardlink and temp_file and not is_torrent_source(temp_file, task):
|
||||
@@ -327,12 +344,19 @@ def transfer_directory_to_library(
|
||||
|
||||
is_torrent = is_torrent_source(source_dir, task)
|
||||
transferred_paths: List[Path] = []
|
||||
max_attempts = _max_attempts_for_batch(len(source_files))
|
||||
|
||||
if len(source_files) == 1:
|
||||
source_file = source_files[0]
|
||||
ext = source_file.suffix.lstrip(".")
|
||||
dest_path = base_library_path.with_suffix(f".{ext}")
|
||||
final_path, op = _transfer_single_file(source_file, dest_path, use_hardlink, is_torrent)
|
||||
final_path, op = _transfer_single_file(
|
||||
source_file,
|
||||
dest_path,
|
||||
use_hardlink,
|
||||
is_torrent,
|
||||
max_attempts=max_attempts,
|
||||
)
|
||||
logger.debug(f"Library {op}: {source_file.name} -> {final_path}")
|
||||
transferred_paths.append(final_path)
|
||||
else:
|
||||
@@ -345,7 +369,13 @@ def transfer_directory_to_library(
|
||||
file_path = build_library_path(library_base, template, file_metadata, extension=ext)
|
||||
file_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
final_path, op = _transfer_single_file(source_file, file_path, use_hardlink, is_torrent)
|
||||
final_path, op = _transfer_single_file(
|
||||
source_file,
|
||||
file_path,
|
||||
use_hardlink,
|
||||
is_torrent,
|
||||
max_attempts=max_attempts,
|
||||
)
|
||||
logger.debug(f"Library {op}: {source_file.name} -> {final_path}")
|
||||
transferred_paths.append(final_path)
|
||||
|
||||
|
||||
@@ -323,7 +323,9 @@ class DownloadClient(ABC):
|
||||
"""
|
||||
pass
|
||||
|
||||
def find_existing(self, url: str) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
def find_existing(
|
||||
self, url: str, category: Optional[str] = None
|
||||
) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
"""
|
||||
Check if a download for this URL already exists in the client.
|
||||
|
||||
@@ -332,6 +334,7 @@ class DownloadClient(ABC):
|
||||
|
||||
Args:
|
||||
url: Download URL (magnet link, .torrent URL, or NZB URL)
|
||||
category: Category to filter by (usenet clients only)
|
||||
|
||||
Returns:
|
||||
Tuple of (download_id, status) if found, None if not found.
|
||||
|
||||
@@ -369,7 +369,9 @@ class DelugeClient(DownloadClient):
|
||||
self._log_error("get_download_path", e, level="debug")
|
||||
return None
|
||||
|
||||
def find_existing(self, url: str) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
def find_existing(
|
||||
self, url: str, category: Optional[str] = None
|
||||
) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
try:
|
||||
self._ensure_connected()
|
||||
|
||||
|
||||
@@ -540,7 +540,9 @@ class QBittorrentClient(DownloadClient):
|
||||
logger.debug(f"qBittorrent could not derive path from files: {type(e).__name__}: {e}")
|
||||
return None
|
||||
|
||||
def find_existing(self, url: str) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
def find_existing(
|
||||
self, url: str, category: Optional[str] = None
|
||||
) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
"""Check if a torrent for this URL already exists in qBittorrent."""
|
||||
try:
|
||||
torrent_info = extract_torrent_info(url)
|
||||
|
||||
@@ -274,7 +274,9 @@ class RTorrentClient(DownloadClient):
|
||||
logger.debug(f"rTorrent get_download_path failed ({error_type}): {e}")
|
||||
return None
|
||||
|
||||
def find_existing(self, url: str) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
def find_existing(
|
||||
self, url: str, category: Optional[str] = None
|
||||
) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
"""Check if a torrent for this URL already exists in rTorrent."""
|
||||
try:
|
||||
torrent_info = extract_torrent_info(url)
|
||||
|
||||
@@ -482,7 +482,9 @@ class SABnzbdClient(DownloadClient):
|
||||
status = self.get_status(download_id)
|
||||
return status.file_path
|
||||
|
||||
def find_existing(self, url: str) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
def find_existing(
|
||||
self, url: str, category: Optional[str] = None
|
||||
) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
"""
|
||||
Check if an NZB for this URL already exists in SABnzbd.
|
||||
|
||||
@@ -493,6 +495,7 @@ class SABnzbdClient(DownloadClient):
|
||||
|
||||
Args:
|
||||
url: NZB URL
|
||||
category: Category to filter by (defaults to configured category)
|
||||
|
||||
Returns:
|
||||
Tuple of (nzo_id, status) if found, None if not found.
|
||||
@@ -518,10 +521,15 @@ class SABnzbdClient(DownloadClient):
|
||||
if not filename:
|
||||
return None
|
||||
|
||||
# Search queue
|
||||
# Use provided category or fall back to configured default
|
||||
search_category = category or self._category
|
||||
|
||||
# Search queue (SABnzbd uses "cat" field for category in queue)
|
||||
queue_result = self._api_call("queue")
|
||||
queue = queue_result.get("queue", {})
|
||||
for slot in queue.get("slots", []):
|
||||
if slot.get("cat", "") != search_category:
|
||||
continue
|
||||
slot_name = slot.get("filename", "")
|
||||
if filename.lower() in slot_name.lower():
|
||||
nzo_id = slot.get("nzo_id")
|
||||
@@ -530,10 +538,12 @@ class SABnzbdClient(DownloadClient):
|
||||
logger.debug(f"Found existing NZB in SABnzbd queue: {nzo_id}")
|
||||
return (nzo_id, status)
|
||||
|
||||
# Search history
|
||||
# Search history (SABnzbd uses "category" field in history)
|
||||
history_result = self._api_call("history", {"limit": 100})
|
||||
history = history_result.get("history", {})
|
||||
for slot in history.get("slots", []):
|
||||
if slot.get("category", "") != search_category:
|
||||
continue
|
||||
slot_name = slot.get("name", "")
|
||||
if filename.lower() in slot_name.lower():
|
||||
nzo_id = slot.get("nzo_id")
|
||||
|
||||
@@ -14,50 +14,6 @@ from shelfmark.core.logger import setup_logger
|
||||
|
||||
logger = setup_logger(__name__)
|
||||
|
||||
_PROWLARR_DOWNLOAD_PATH = re.compile(r"(?:/api/v1/indexer)?/\d+/download$")
|
||||
|
||||
|
||||
def _decode_prowlarr_link(link_value: str) -> Optional[str]:
|
||||
"""Decode Prowlarr's link param into a usable URL, if possible."""
|
||||
if not link_value:
|
||||
return None
|
||||
|
||||
value = link_value.strip()
|
||||
if not value:
|
||||
return None
|
||||
|
||||
if value.startswith(("http://", "https://", "magnet:")):
|
||||
return value
|
||||
|
||||
# Try urlsafe + standard base64 decoding with padding.
|
||||
padded = value + "=" * (-len(value) % 4)
|
||||
for decoder in (base64.urlsafe_b64decode, base64.b64decode):
|
||||
try:
|
||||
decoded = decoder(padded).decode("utf-8", errors="ignore").strip()
|
||||
except Exception:
|
||||
continue
|
||||
if decoded.startswith(("http://", "https://", "magnet:")):
|
||||
return decoded
|
||||
|
||||
return None
|
||||
|
||||
|
||||
def _get_prowlarr_fallback_url(url: str) -> Optional[str]:
|
||||
"""Try to extract the original download URL from a Prowlarr download proxy URL."""
|
||||
try:
|
||||
parsed = urlparse(url)
|
||||
if not _PROWLARR_DOWNLOAD_PATH.search(parsed.path):
|
||||
return None
|
||||
|
||||
params = parse_qs(parsed.query)
|
||||
link_value = (params.get("link") or [None])[0]
|
||||
if not link_value:
|
||||
return None
|
||||
|
||||
return _decode_prowlarr_link(link_value)
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
@dataclass
|
||||
class TorrentInfo:
|
||||
@@ -91,7 +47,6 @@ def extract_torrent_info(
|
||||
url: str,
|
||||
fetch_torrent: bool = True,
|
||||
expected_hash: Optional[str] = None,
|
||||
allow_prowlarr_fallback: bool = True,
|
||||
) -> TorrentInfo:
|
||||
"""Extract info_hash from magnet link or .torrent URL.
|
||||
|
||||
@@ -100,17 +55,9 @@ def extract_torrent_info(
|
||||
requires the `X-Api-Key` header. If `PROWLARR_API_KEY` is configured,
|
||||
include it for the torrent fetch request.
|
||||
|
||||
This mirrors how Sonarr builds an authenticated download request via the
|
||||
indexer when grabbing torrent files.
|
||||
Redirects to magnet links are handled explicitly so we can extract a
|
||||
hash from the magnet when available.
|
||||
"""
|
||||
fallback_url: Optional[str] = None
|
||||
if allow_prowlarr_fallback:
|
||||
decoded_url = _get_prowlarr_fallback_url(url)
|
||||
if decoded_url and decoded_url != url:
|
||||
logger.debug(f"Decoded Prowlarr link, using direct URL: {decoded_url[:80]}...")
|
||||
fallback_url = url
|
||||
url = decoded_url
|
||||
|
||||
is_magnet = url.startswith("magnet:")
|
||||
|
||||
# Try to extract hash from magnet URL
|
||||
@@ -121,13 +68,10 @@ def extract_torrent_info(
|
||||
return TorrentInfo(info_hash=info_hash, torrent_data=None, is_magnet=True, magnet_url=url)
|
||||
|
||||
# Not a magnet - try to fetch and parse the .torrent file
|
||||
if expected_hash:
|
||||
if not fetch_torrent:
|
||||
return TorrentInfo(info_hash=expected_hash, torrent_data=None, is_magnet=False)
|
||||
|
||||
if not fetch_torrent:
|
||||
return TorrentInfo(info_hash=None, torrent_data=None, is_magnet=False)
|
||||
|
||||
headers: dict[str, str] = {}
|
||||
headers: dict[str, str] = {"Accept": "application/x-bittorrent"}
|
||||
api_key = str(config.get("PROWLARR_API_KEY", "") or "").strip()
|
||||
if api_key:
|
||||
headers["X-Api-Key"] = api_key
|
||||
@@ -179,7 +123,7 @@ def extract_torrent_info(
|
||||
except Exception:
|
||||
pass # Not text, continue with torrent parsing
|
||||
|
||||
info_hash = extract_info_hash_from_torrent(torrent_data)
|
||||
info_hash = extract_info_hash_from_torrent(torrent_data) or expected_hash
|
||||
if info_hash:
|
||||
logger.debug(f"Extracted hash from torrent file: {info_hash}")
|
||||
else:
|
||||
@@ -187,74 +131,7 @@ def extract_torrent_info(
|
||||
return TorrentInfo(info_hash=info_hash, torrent_data=torrent_data, is_magnet=False)
|
||||
except Exception as e:
|
||||
logger.debug(f"Could not fetch torrent file: {e}")
|
||||
if allow_prowlarr_fallback and fallback_url:
|
||||
logger.debug(f"Retrying torrent fetch via Prowlarr proxy: {fallback_url[:80]}...")
|
||||
return extract_torrent_info(
|
||||
fallback_url,
|
||||
fetch_torrent=fetch_torrent,
|
||||
expected_hash=expected_hash,
|
||||
allow_prowlarr_fallback=False,
|
||||
)
|
||||
return TorrentInfo(info_hash=None, torrent_data=None, is_magnet=False)
|
||||
|
||||
|
||||
headers: dict[str, str] = {}
|
||||
api_key = str(config.get("PROWLARR_API_KEY", "") or "").strip()
|
||||
if api_key:
|
||||
headers["X-Api-Key"] = api_key
|
||||
|
||||
def resolve_url(current: str, location: str) -> str:
|
||||
if not location:
|
||||
return current
|
||||
# Support relative redirect locations
|
||||
return urljoin(current, location)
|
||||
|
||||
try:
|
||||
logger.debug(f"Fetching torrent file from: {url[:80]}...")
|
||||
|
||||
# Use allow_redirects=False to handle magnet link redirects manually
|
||||
# Some indexers redirect download URLs to magnet links
|
||||
resp = requests.get(url, timeout=30, allow_redirects=False, headers=headers)
|
||||
|
||||
# Check if this is a redirect to a magnet link
|
||||
if resp.status_code in (301, 302, 303, 307, 308):
|
||||
redirect_url = resolve_url(url, resp.headers.get("Location", ""))
|
||||
if redirect_url.startswith("magnet:"):
|
||||
logger.debug("Download URL redirected to magnet link")
|
||||
info_hash = extract_hash_from_magnet(redirect_url)
|
||||
return TorrentInfo(
|
||||
info_hash=info_hash, torrent_data=None, is_magnet=True, magnet_url=redirect_url
|
||||
)
|
||||
# Not a magnet redirect, follow it manually
|
||||
logger.debug(f"Following redirect to: {redirect_url[:80]}...")
|
||||
resp = requests.get(redirect_url, timeout=30, headers=headers)
|
||||
|
||||
resp.raise_for_status()
|
||||
torrent_data = resp.content
|
||||
|
||||
# Check if response is actually a magnet link (text response)
|
||||
# Some indexers return magnet links as plain text instead of redirecting
|
||||
if len(torrent_data) < 2000: # Magnet links are typically short
|
||||
try:
|
||||
text_content = torrent_data.decode("utf-8", errors="ignore").strip()
|
||||
if text_content.startswith("magnet:"):
|
||||
logger.debug("Download URL returned magnet link as response body")
|
||||
info_hash = extract_hash_from_magnet(text_content)
|
||||
return TorrentInfo(
|
||||
info_hash=info_hash, torrent_data=None, is_magnet=True, magnet_url=text_content
|
||||
)
|
||||
except Exception:
|
||||
pass # Not text, continue with torrent parsing
|
||||
|
||||
info_hash = extract_info_hash_from_torrent(torrent_data)
|
||||
if info_hash:
|
||||
logger.debug(f"Extracted hash from torrent file: {info_hash}")
|
||||
else:
|
||||
logger.warning("Could not extract hash from torrent file")
|
||||
return TorrentInfo(info_hash=info_hash, torrent_data=torrent_data, is_magnet=False)
|
||||
except Exception as e:
|
||||
logger.debug(f"Could not fetch torrent file: {e}")
|
||||
return TorrentInfo(info_hash=None, torrent_data=None, is_magnet=False)
|
||||
return TorrentInfo(info_hash=expected_hash, torrent_data=None, is_magnet=False)
|
||||
|
||||
|
||||
def parse_transmission_url(url: str) -> Tuple[str, int, str]:
|
||||
@@ -346,7 +223,10 @@ def extract_info_hash_from_torrent(torrent_data: bytes) -> Optional[str]:
|
||||
return None
|
||||
|
||||
info_bencoded = bencode_encode(decoded[b'info'])
|
||||
return hashlib.sha1(info_bencoded).hexdigest().lower()
|
||||
info_dict = decoded[b'info']
|
||||
if isinstance(info_dict, dict) and b'pieces' in info_dict:
|
||||
return hashlib.sha1(info_bencoded).hexdigest().lower()
|
||||
return hashlib.sha256(info_bencoded).hexdigest().lower()
|
||||
except Exception as e:
|
||||
logger.debug(f"Failed to parse torrent file: {e}")
|
||||
return None
|
||||
@@ -360,7 +240,42 @@ def extract_hash_from_magnet(magnet_url: str) -> Optional[str]:
|
||||
parsed = urlparse(magnet_url)
|
||||
params = parse_qs(parsed.query)
|
||||
|
||||
for xt in params.get("xt", []):
|
||||
def extract_btmh(value: str) -> Optional[str]:
|
||||
raw_value = value.strip()
|
||||
if not raw_value:
|
||||
return None
|
||||
|
||||
data: Optional[bytes] = None
|
||||
if re.fullmatch(r"[a-fA-F0-9]+", raw_value):
|
||||
if len(raw_value) % 2 != 0:
|
||||
return None
|
||||
try:
|
||||
data = bytes.fromhex(raw_value)
|
||||
except ValueError:
|
||||
return None
|
||||
else:
|
||||
padded = raw_value.upper() + "=" * (-len(raw_value) % 8)
|
||||
try:
|
||||
data = base64.b32decode(padded, casefold=True)
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
if not data:
|
||||
return None
|
||||
|
||||
if len(data) >= 34 and data[0] == 0x12 and data[1] == 0x20:
|
||||
digest = data[2:34]
|
||||
if len(digest) == 32:
|
||||
return digest.hex().lower()
|
||||
|
||||
if len(data) == 32:
|
||||
return data.hex().lower()
|
||||
|
||||
return None
|
||||
|
||||
xt_values = params.get("xt", [])
|
||||
|
||||
for xt in xt_values:
|
||||
# Format: urn:btih:<hash> (32 or 40 chars)
|
||||
match = re.match(r"urn:btih:([a-fA-F0-9]{40}|[a-zA-Z0-9]{32})", xt)
|
||||
if match:
|
||||
@@ -380,4 +295,11 @@ def extract_hash_from_magnet(magnet_url: str) -> Optional[str]:
|
||||
# Fallback: return as-is
|
||||
return hash_value.lower()
|
||||
|
||||
for xt in xt_values:
|
||||
if xt.startswith("urn:btmh:"):
|
||||
btmh_value = xt[len("urn:btmh:"):]
|
||||
btmh_hash = extract_btmh(btmh_value)
|
||||
if btmh_hash:
|
||||
return btmh_hash
|
||||
|
||||
return None
|
||||
|
||||
@@ -249,7 +249,9 @@ class TransmissionClient(DownloadClient):
|
||||
self._log_error("get_download_path", e, level="debug")
|
||||
return None
|
||||
|
||||
def find_existing(self, url: str) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
def find_existing(
|
||||
self, url: str, category: Optional[str] = None
|
||||
) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
"""Check if a torrent for this URL already exists in Transmission."""
|
||||
try:
|
||||
torrent_info = extract_torrent_info(url)
|
||||
|
||||
@@ -188,7 +188,8 @@ class ProwlarrHandler(DownloadHandler):
|
||||
|
||||
# Check if this download already exists in the client
|
||||
status_callback("resolving", f"Checking {client.name}")
|
||||
existing = client.find_existing(download_url)
|
||||
category = self._get_category_for_task(client, task)
|
||||
existing = client.find_existing(download_url, category=category)
|
||||
|
||||
if existing:
|
||||
download_id, existing_status = existing
|
||||
|
||||
@@ -712,6 +712,47 @@ class TestCustomScriptExecution:
|
||||
result_path = Path(result)
|
||||
assert call_args[0][0] == ["/path/to/script.sh", str(result_path)]
|
||||
|
||||
def test_runs_custom_script_with_relative_path_mode(self, temp_dirs, sample_direct_task):
|
||||
"""Runs custom script with a destination-relative path when configured."""
|
||||
from shelfmark.download.postprocess.router import post_process_download as _post_process_download
|
||||
import subprocess
|
||||
|
||||
temp_file = temp_dirs["staging"] / "book.epub"
|
||||
temp_file.write_bytes(b"content")
|
||||
|
||||
status_cb = MagicMock()
|
||||
cancel_flag = Event()
|
||||
|
||||
with patch('shelfmark.core.config.config') as mock_config, \
|
||||
patch('shelfmark.config.env.TMP_DIR', temp_dirs["staging"]), \
|
||||
patch('subprocess.run') as mock_run:
|
||||
|
||||
mock_config.USE_BOOK_TITLE = False
|
||||
mock_config.CUSTOM_SCRIPT = "/path/to/script.sh"
|
||||
_sync_core_config(mock_config, mock_config)
|
||||
mock_config.get = _mock_destination_config(
|
||||
temp_dirs["ingest"],
|
||||
{"CUSTOM_SCRIPT_PATH_MODE": "relative"},
|
||||
)
|
||||
_sync_core_config(mock_config, mock_config)
|
||||
|
||||
mock_run.return_value = MagicMock(stdout="", returncode=0)
|
||||
|
||||
result = _post_process_download(
|
||||
temp_file=temp_file,
|
||||
task=sample_direct_task,
|
||||
cancel_flag=cancel_flag,
|
||||
status_callback=status_cb,
|
||||
)
|
||||
|
||||
assert result is not None
|
||||
result_path = Path(result)
|
||||
expected_relative = result_path.relative_to(temp_dirs["ingest"])
|
||||
script_args = mock_run.call_args[0][0]
|
||||
assert script_args[0] == "/path/to/script.sh"
|
||||
assert script_args[1] == str(expected_relative)
|
||||
assert not Path(script_args[1]).is_absolute()
|
||||
|
||||
def test_runs_custom_script_for_directory_download_once(self, temp_dirs):
|
||||
"""Runs custom script once after transferring a directory download."""
|
||||
from shelfmark.download.postprocess.router import post_process_download as _post_process_download
|
||||
|
||||
@@ -275,6 +275,29 @@ class TestAtomicMove:
|
||||
assert dest.read_text() == "existing"
|
||||
assert result.read_text() == "new"
|
||||
|
||||
def test_false_positive_exists_probe(self, tmp_path, monkeypatch):
|
||||
"""Moves file even if exists() falsely reports a collision."""
|
||||
from shelfmark.download.fs import atomic_move as _atomic_move
|
||||
|
||||
source = tmp_path / "source.txt"
|
||||
source.write_text("content")
|
||||
dest = tmp_path / "dest.txt"
|
||||
|
||||
original_exists = Path.exists
|
||||
|
||||
def _fake_exists(self):
|
||||
if self == dest:
|
||||
return True
|
||||
return original_exists(self)
|
||||
|
||||
monkeypatch.setattr(Path, "exists", _fake_exists)
|
||||
|
||||
result = _atomic_move(source, dest)
|
||||
|
||||
assert result == dest
|
||||
assert not source.exists()
|
||||
assert result.read_text() == "content"
|
||||
|
||||
def test_cross_filesystem_fallback(self):
|
||||
"""Falls back to copy when cross-filesystem."""
|
||||
from shelfmark.download.fs import atomic_move as _atomic_move
|
||||
@@ -865,8 +888,8 @@ class TestTorrentSourceCleanupProtection:
|
||||
"HARDLINK_TORRENTS": hardlink,
|
||||
"HARDLINK_TORRENTS_AUDIOBOOK": hardlink,
|
||||
# Supported formats
|
||||
"SUPPORTED_FORMATS": ["epub", "mobi"],
|
||||
"SUPPORTED_AUDIOBOOK_FORMATS": ["mp3"],
|
||||
"SUPPORTED_FORMATS": ["epub", "mobi", "cbz", "cbr", "azw3", "fb2", "djvu", "pdf"],
|
||||
"SUPPORTED_AUDIOBOOK_FORMATS": ["mp3", "m4a", "m4b", "flac"],
|
||||
}.get(key, default))
|
||||
|
||||
# ==================== EPUB EBOOK TESTS ====================
|
||||
|
||||
@@ -160,15 +160,15 @@ class TestExtractInfoHash:
|
||||
"""Tests for extracting info hash from torrent files."""
|
||||
|
||||
def test_extract_hash_from_simple_torrent(self):
|
||||
"""Test extracting hash from a simple torrent structure."""
|
||||
# Create a minimal valid torrent structure
|
||||
info_dict = {b"name": b"test.txt", b"length": 100}
|
||||
"""Test extracting hash from a simple v1 torrent structure."""
|
||||
# Create a minimal valid v1 torrent structure (has 'pieces' key)
|
||||
info_dict = {b"name": b"test.txt", b"length": 100, b"pieces": b"\x00" * 20}
|
||||
torrent = {b"info": info_dict}
|
||||
torrent_bytes = _bencode_encode(torrent)
|
||||
|
||||
result = _extract_info_hash_from_torrent(torrent_bytes)
|
||||
|
||||
# Should return a 40-character hex string
|
||||
# V1 torrents return SHA-1 hash (40-character hex string)
|
||||
assert result is not None
|
||||
assert len(result) == 40
|
||||
assert all(c in "0123456789abcdef" for c in result)
|
||||
|
||||
@@ -125,7 +125,9 @@ class MockClient(DownloadClient):
|
||||
def get_download_path(self, download_id: str) -> Optional[str]:
|
||||
return "/downloads/test-file.epub"
|
||||
|
||||
def find_existing(self, url: str) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
def find_existing(
|
||||
self, url: str, category: Optional[str] = None
|
||||
) -> Optional[Tuple[str, DownloadStatus]]:
|
||||
return None
|
||||
|
||||
|
||||
|
||||
@@ -667,6 +667,7 @@ class TestSABnzbdClientFindExisting:
|
||||
{
|
||||
"nzo_id": "SABnzbd_nzo_found",
|
||||
"filename": "Test_Book.nzb",
|
||||
"cat": "cwabd",
|
||||
"status": "Downloading",
|
||||
"percentage": "50",
|
||||
"timeleft": "",
|
||||
@@ -716,6 +717,7 @@ class TestSABnzbdClientFindExisting:
|
||||
{
|
||||
"nzo_id": "SABnzbd_nzo_history",
|
||||
"name": "Test Book",
|
||||
"category": "cwabd",
|
||||
"status": "Completed",
|
||||
"storage": "/downloads/Test Book",
|
||||
}
|
||||
@@ -774,3 +776,69 @@ class TestSABnzbdClientFindExisting:
|
||||
result = client.find_existing("https://example.com/unknown.nzb")
|
||||
|
||||
assert result is None
|
||||
|
||||
def test_find_existing_ignores_different_category(self, monkeypatch):
|
||||
"""Test that find_existing ignores downloads in different categories.
|
||||
|
||||
This tests the fix for issue #508 where SABnzbd's test download
|
||||
in the 'default' category was incorrectly matched.
|
||||
"""
|
||||
config_values = {
|
||||
"SABNZBD_URL": "http://localhost:8080",
|
||||
"SABNZBD_API_KEY": "abc123",
|
||||
"SABNZBD_CATEGORY": "books",
|
||||
}
|
||||
monkeypatch.setattr(
|
||||
"shelfmark.release_sources.prowlarr.clients.sabnzbd.config.get",
|
||||
lambda key, default="": config_values.get(key, default),
|
||||
)
|
||||
|
||||
def mock_api_call(mode, params=None):
|
||||
if mode == "queue":
|
||||
return {
|
||||
"queue": {
|
||||
"slots": [
|
||||
{
|
||||
"nzo_id": "SABnzbd_nzo_test",
|
||||
"filename": "test_download_1000MB",
|
||||
"cat": "default", # Different category
|
||||
"status": "Downloading",
|
||||
"percentage": "50",
|
||||
"timeleft": "",
|
||||
"kbpersec": "",
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
if mode == "history":
|
||||
return {
|
||||
"history": {
|
||||
"slots": [
|
||||
{
|
||||
"nzo_id": "SABnzbd_nzo_old_test",
|
||||
"name": "test_download_1000MB",
|
||||
"category": "default", # Different category
|
||||
"status": "Completed",
|
||||
"storage": "/downloads/default/test_download_1000MB",
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
return {}
|
||||
|
||||
from shelfmark.release_sources.prowlarr.clients.sabnzbd import (
|
||||
SABnzbdClient,
|
||||
)
|
||||
|
||||
with patch.object(SABnzbdClient, "__init__", lambda x: None):
|
||||
client = SABnzbdClient()
|
||||
client.url = "http://localhost:8080"
|
||||
client.api_key = "abc123"
|
||||
client._category = "books"
|
||||
client._api_call = mock_api_call
|
||||
|
||||
# Even though "download" might match "test_download_1000MB",
|
||||
# it should be ignored because it's in the "default" category
|
||||
result = client.find_existing("https://example.com/download/Book.nzb")
|
||||
|
||||
assert result is None
|
||||
|
||||
@@ -8,6 +8,9 @@ Tests:
|
||||
- extract_hash_from_magnet
|
||||
"""
|
||||
|
||||
import base64
|
||||
import hashlib
|
||||
|
||||
import pytest
|
||||
|
||||
from shelfmark.release_sources.prowlarr.clients.torrent_utils import (
|
||||
@@ -218,6 +221,8 @@ class TestBencodeEncode:
|
||||
result = bencode_encode(data)
|
||||
assert result == b"d4:listli1ei2ei3ee3:numi42ee"
|
||||
|
||||
|
||||
|
||||
def test_encode_invalid_type_raises(self):
|
||||
"""Test that invalid types raise ValueError."""
|
||||
with pytest.raises(ValueError):
|
||||
@@ -276,7 +281,12 @@ class TestExtractInfoHash:
|
||||
|
||||
def test_extract_hash_from_simple_torrent(self):
|
||||
"""Test extracting hash from a simple torrent structure."""
|
||||
info_dict = {b"name": b"test.txt", b"length": 100}
|
||||
info_dict = {
|
||||
b"name": b"test.txt",
|
||||
b"length": 100,
|
||||
b"piece length": 16384,
|
||||
b"pieces": b"\x00" * 20,
|
||||
}
|
||||
torrent = {b"info": info_dict}
|
||||
torrent_bytes = bencode_encode(torrent)
|
||||
|
||||
@@ -321,6 +331,19 @@ class TestExtractInfoHash:
|
||||
|
||||
assert hash1 != hash2
|
||||
|
||||
def test_extract_hash_v2_without_pieces(self):
|
||||
"""Use SHA-256 when torrent lacks v1 pieces."""
|
||||
info_dict = {
|
||||
b"meta version": 2,
|
||||
b"file tree": {b"test.txt": {b"": {b"length": 123}}},
|
||||
b"piece length": 16384,
|
||||
}
|
||||
torrent = {b"info": info_dict}
|
||||
torrent_bytes = bencode_encode(torrent)
|
||||
|
||||
expected = hashlib.sha256(bencode_encode(info_dict)).hexdigest().lower()
|
||||
assert extract_info_hash_from_torrent(torrent_bytes) == expected
|
||||
|
||||
|
||||
class TestExtractHashFromMagnet:
|
||||
"""Tests for extracting hash from magnet links."""
|
||||
@@ -380,3 +403,20 @@ class TestExtractHashFromMagnet:
|
||||
)
|
||||
result = extract_hash_from_magnet(magnet)
|
||||
assert result == "3b245504cf5f11bbdbe1201cea6a6bf45aee1bc0"
|
||||
|
||||
def test_extract_hash_from_btmh_hex(self):
|
||||
"""Test extracting v2 hash from btmh (hex multihash)."""
|
||||
digest = bytes(range(1, 33))
|
||||
multihash = b"\x12\x20" + digest
|
||||
magnet = f"magnet:?xt=urn:btmh:{multihash.hex()}&dn=test"
|
||||
result = extract_hash_from_magnet(magnet)
|
||||
assert result == digest.hex()
|
||||
|
||||
def test_extract_hash_from_btmh_base32(self):
|
||||
"""Test extracting v2 hash from btmh (base32 multihash)."""
|
||||
digest = bytes(range(1, 33))
|
||||
multihash = b"\x12\x20" + digest
|
||||
b32 = base64.b32encode(multihash).decode("ascii").rstrip("=")
|
||||
magnet = f"magnet:?xt=urn:btmh:{b32}&dn=test"
|
||||
result = extract_hash_from_magnet(magnet)
|
||||
assert result == digest.hex()
|
||||
|
||||
Reference in New Issue
Block a user