diff --git a/shelfmark/download/clients/base_handler.py b/shelfmark/download/clients/base_handler.py index c4ca3cbc..eb8afc6b 100644 --- a/shelfmark/download/clients/base_handler.py +++ b/shelfmark/download/clients/base_handler.py @@ -15,6 +15,7 @@ from shelfmark.core.config import config from shelfmark.core.logger import setup_logger from shelfmark.core.request_helpers import normalize_optional_text from shelfmark.core.utils import is_audiobook +from shelfmark.download.activity import release_activity_grace, request_activity_grace from shelfmark.download.clients import ( DownloadClient, DownloadState, @@ -50,6 +51,20 @@ def _is_sabnzbd_like_client(candidate: DownloadClient) -> TypeGuard[_SabnzbdLike # How often to poll the download client for status (seconds) POLL_INTERVAL = 2 +# A torrent the client has queued is waiting for a free slot, not stalled. Its status never +# changes while it waits, so the orchestrator's stall timer would cancel it after +# STALL_TIMEOUT. Ask for a grace instead (the orchestrator caps one grace at 16 minutes), +# renew it while the torrent is still queued, and stop at a hard ceiling so a queue that +# never moves still ends. +QUEUE_GRACE_SECONDS = 900.0 +QUEUE_GRACE_RENEW_SECONDS = 600.0 +QUEUE_MAX_WAIT_SECONDS = 7200.0 +# Torrents are jerky: a swarm with one seed can sit still for several minutes and then carry +# on. The orchestrator's default window (STALL_TIMEOUT) is five minutes of no change. A torrent +# that has already moved gets this longer window, restarted every time its progress goes up, +# so a dead one still ends 15 minutes after its last movement. One that never moved (a magnet +# that cannot fetch metadata) keeps the short window and is cleaned up quickly. +MOVING_STALL_SECONDS = 900.0 WINDOWS_DRIVE_PREFIX_LENGTH = 2 SECONDS_PER_MINUTE = 60 SECONDS_PER_HOUR = 3600 @@ -927,6 +942,10 @@ class ExternalClientHandler(DownloadHandler, ABC): ) -> str | None: """Poll the download client for progress and handle completion.""" poll_interval = self._poll_interval() + queued_since: float | None = None + grace_requested_at = 0.0 + last_progress = 0.0 + start_window_given = False # Track consecutive "not found" errors - torrents may take time to appear in client not_found_count = 0 max_not_found_retries = 15 # 15 retries * poll interval ~= 30s grace period @@ -1012,6 +1031,31 @@ class ExternalClientHandler(DownloadHandler, ABC): # Reset not-found counter on successful status check not_found_count = 0 + if status.state == DownloadState.QUEUED: + now = time.monotonic() + if queued_since is None: + queued_since = now + if ( + now - queued_since < QUEUE_MAX_WAIT_SECONDS + and now - grace_requested_at >= QUEUE_GRACE_RENEW_SECONDS + ): + request_activity_grace(status_callback, QUEUE_GRACE_SECONDS) + grace_requested_at = now + elif queued_since is not None: + release_activity_grace(status_callback) + queued_since = None + grace_requested_at = 0.0 + + if status.progress > last_progress: + request_activity_grace(status_callback, MOVING_STALL_SECONDS) + start_window_given = True + elif not start_window_given and status.state != DownloadState.QUEUED: + # Metadata and the first peer can take longer than the orchestrator's + # stall timeout, so a torrent gets one full window to get going. + request_activity_grace(status_callback, MOVING_STALL_SECONDS) + start_window_given = True + last_progress = status.progress + # Build status message - use client message if provided, else build progress msg = status.message or self._build_progress_message(status) if status.state == DownloadState.PROCESSING: diff --git a/shelfmark/download/orchestrator.py b/shelfmark/download/orchestrator.py index 26201603..1f5af7d7 100644 --- a/shelfmark/download/orchestrator.py +++ b/shelfmark/download/orchestrator.py @@ -1097,11 +1097,13 @@ def _cancel_stalled_task(task_id: str) -> None: the progress lock across that blocks the hub and every other download worker. """ logger.warning("Download stalled for %s, cancelling", task_id) - book_queue.cancel_download(task_id) + # The message goes first: the terminal hook copies it into the history row, and a cancel + # that records the last poll's message ("Queued") hides that the timer ended the download. book_queue.update_status_message( task_id, f"Download stalled (no activity for {STALL_TIMEOUT}s)", ) + book_queue.cancel_download(task_id) def concurrent_download_loop() -> None: diff --git a/tests/download/test_orchestrator_stall.py b/tests/download/test_orchestrator_stall.py index adf1de24..8dbdd9b3 100644 --- a/tests/download/test_orchestrator_stall.py +++ b/tests/download/test_orchestrator_stall.py @@ -289,3 +289,19 @@ def test_cancel_stalled_task_does_not_hold_the_progress_lock(monkeypatch): mock_queue.update_status_message.assert_called_once_with( "book", f"Download stalled (no activity for {orchestrator.STALL_TIMEOUT}s)" ) + + +def test_the_stall_reason_is_set_before_the_cancel_finalises_history(monkeypatch): + """The terminal hook copies the task's message into the history row, so a stall that + sets its message after the cancel is recorded with whatever the last poll said ("Queued") + and cannot be told apart from a user cancelling it.""" + import shelfmark.download.orchestrator as orchestrator + + _reset(orchestrator) + mock_queue = MagicMock() + monkeypatch.setattr(orchestrator, "book_queue", mock_queue) + + orchestrator._cancel_stalled_task("book") + + names = [call[0] for call in mock_queue.method_calls] + assert names.index("update_status_message") < names.index("cancel_download") diff --git a/tests/download/test_torrent_queue_stall.py b/tests/download/test_torrent_queue_stall.py new file mode 100644 index 00000000..ae0258eb --- /dev/null +++ b/tests/download/test_torrent_queue_stall.py @@ -0,0 +1,169 @@ +"""A torrent the client has queued is waiting for a slot, not stalled.""" + +from __future__ import annotations + +import sys +from threading import Event +from unittest.mock import MagicMock, patch + +from shelfmark.core.models import DownloadTask +from shelfmark.download.activity import ACTIVITY_GRACE_STATUS +from shelfmark.download.clients import DownloadState, DownloadStatus +from shelfmark.release_sources.prowlarr.handler import ProwlarrHandler + +bh = sys.modules["shelfmark.download.clients.base_handler"] # loaded by the import above + + +def _status(state: DownloadState, progress: float = 0.0, **kw) -> DownloadStatus: + return DownloadStatus( + progress=progress, + state=state, + message="Queued" if state == DownloadState.QUEUED else None, + complete=kw.get("complete", False), + file_path=None, + ) + + +class _Recorder: + def __init__(self) -> None: + self.events: list[tuple[str, str | None]] = [] + + def status(self, status: str, message: str | None) -> None: + self.events.append((status, message)) + + @property + def graces(self) -> list[float]: + return [float(m) for s, m in self.events if s == ACTIVITY_GRACE_STATUS] + + +def _poll(statuses: list[DownloadStatus], clock: list[float] | None = None) -> _Recorder: + """Run the poll loop over a scripted list of client statuses, then cancel.""" + cancel = Event() + seen = iter(statuses) + client = MagicMock() + client.name = "qbittorrent" + + def get_status(_id: str) -> DownloadStatus: + try: + return next(seen) + except StopIteration: + cancel.set() + return statuses[-1] + + client.get_status.side_effect = get_status + rec = _Recorder() + ticks = iter(clock) if clock is not None else None + patches = [patch.object(ProwlarrHandler, "_poll_interval", return_value=0.0)] + if ticks is not None: + patches.append(patch.object(bh.time, "monotonic", side_effect=lambda: next(ticks))) + for p in patches: + p.start() + try: + ProwlarrHandler()._poll_and_complete( + client, + "hash", + "torrent", + DownloadTask(task_id="t", source="prowlarr", title="Book"), + cancel, + lambda _p: None, + rec.status, + ) + finally: + patch.stopall() + return rec + + +Q, D = DownloadState.QUEUED, DownloadState.DOWNLOADING + + +def test_a_queued_torrent_asks_for_a_grace_once() -> None: + rec = _poll([_status(Q)] * 4) + + assert rec.graces[0] == bh.QUEUE_GRACE_SECONDS + assert rec.graces.count(bh.QUEUE_GRACE_SECONDS) == 1 + + +def test_the_grace_is_released_when_the_torrent_starts() -> None: + rec = _poll([_status(Q), _status(Q), _status(D, 5.0), _status(D, 6.0)]) + + # queue grace, its release, then one movement grace per reading that went up + assert rec.graces == [ + bh.QUEUE_GRACE_SECONDS, + 0.0, + bh.MOVING_STALL_SECONDS, + bh.MOVING_STALL_SECONDS, + ] + + +def test_a_torrent_that_never_queued_only_gets_the_moving_window() -> None: + rec = _poll([_status(D, 1.0), _status(D, 2.0)]) + + assert rec.graces == [bh.MOVING_STALL_SECONDS] * 2 # no release, no queue renewal + + +# --- a torrent that is moving gets a longer window each time it moves --------------------- + + +def test_every_step_forward_pushes_the_stall_deadline_out() -> None: + rec = _poll([_status(D, 1.0), _status(D, 2.0), _status(D, 3.5)]) + + assert rec.graces == [bh.MOVING_STALL_SECONDS] * 3 + + +def test_a_torrent_that_stops_moving_stops_getting_more_time() -> None: + rec = _poll([_status(D, 5.0), _status(D, 5.0), _status(D, 5.0), _status(D, 5.0)]) + + assert rec.graces == [bh.MOVING_STALL_SECONDS] # only the first reading moved + + +def test_a_torrent_gets_a_full_window_to_start_even_before_any_data_arrives() -> None: + # Metadata and the first peer can take longer than the orchestrator's 5 minutes. + assert _poll([_status(D, 0.0)] * 4).graces == [bh.MOVING_STALL_SECONDS] + + +def test_the_start_window_is_given_once_not_renewed_while_stuck_at_zero() -> None: + graces = _poll([_status(D, 0.0)] * 20).graces + + assert graces.count(bh.MOVING_STALL_SECONDS) == 1 + + +def test_a_torrent_leaving_the_queue_at_zero_gets_a_fresh_start_window() -> None: + rec = _poll([_status(Q), _status(Q), _status(D, 0.0), _status(D, 0.0)]) + + # queue grace, its release, then the start window + assert rec.graces == [bh.QUEUE_GRACE_SECONDS, 0.0, bh.MOVING_STALL_SECONDS] + + +def test_a_queued_torrent_is_not_given_the_start_window_yet() -> None: + assert _poll([_status(Q)] * 4).graces == [bh.QUEUE_GRACE_SECONDS] # the queue grace only + + +def test_going_backwards_is_not_movement() -> None: + rec = _poll([_status(D, 50.0), _status(D, 40.0), _status(D, 40.0)]) + + assert rec.graces == [bh.MOVING_STALL_SECONDS] + + +def test_the_windows_fit_inside_what_the_orchestrator_allows() -> None: + from shelfmark.download import orchestrator + + cap = orchestrator._MAX_ACTIVITY_GRACE_SECONDS + assert bh.QUEUE_GRACE_SECONDS <= cap + assert bh.MOVING_STALL_SECONDS <= cap + assert bh.MOVING_STALL_SECONDS > orchestrator.STALL_TIMEOUT + + +def test_a_long_queue_renews_the_grace_but_only_until_the_ceiling() -> None: + # Each poll reads the clock once; times are seconds since the torrent was first seen queued. + times = [0, 100, 700, 800, 1400, 7000, 7300, 7400, 7500] + rec = _poll([_status(Q)] * len(times), clock=[1000 + t for t in times]) + + # first request at t=0, renewed after >= 600s (t=700, t=1400, t=7000), none past the 7200s ceiling + assert rec.graces.count(bh.QUEUE_GRACE_SECONDS) == 4 + + +def test_a_queue_wait_past_the_ceiling_gets_no_more_grace() -> None: + times = [0, 7300, 7400] + rec = _poll([_status(Q)] * 3, clock=[5000 + t for t in times]) + + assert rec.graces.count(bh.QUEUE_GRACE_SECONDS) == 1