fix(download): do not cancel a queued torrent as stalled (#1420)

qBittorrent only runs a few downloads at once (`max_active_downloads`
defaults to 3) and holds the rest as `queuedDL`. Shelfmark's stall timer
cancels a download after five minutes without a changed status or
progress. A queued torrent reports "Queued" on every poll, so it looks
idle, and Shelfmark cancels a download that was only waiting for a slot.

I run Shelfmark against the same qBittorrent that Sonarr, Radarr and
Lidarr use, which I assume is a common setup. Their torrents and
Shelfmark's share one limit, so anything that fills the slots makes
Shelfmark's torrent wait. In my case three torrents that never fetched
metadata held every slot. I'd expect a busy Sonarr queue to hold them
the same way, but I haven't reproduced that. Once the stall timer had
cancelled the waiting downloads they stayed in the client, still queued
behind the same blockers. On my instance 10 AudioBookBay downloads were
cancelled while qBittorrent still had them queued.

This changes the torrent poll loop. When the client reports a queued
torrent, it asks for an activity grace (the mechanism `html_get_page`
already uses for protection bypasses) and releases it as soon as the
torrent leaves the queue. The grace is renewed every ten minutes and
stops at two hours, so a queue that never moves still ends. A torrent
that is downloading and moving is timed as before until the new window
below applies.

The second commit sets the stall message before the cancel runs. The
terminal hook copies the message into the history row, and it was set
afterwards, so a download the timer ended was recorded with the last
poll's "Queued". That left no way to tell it apart from a person
cancelling.

The last two commits are about slow torrents. The timer only counts a
change in progress, so a torrent that needs more than five minutes to
fetch metadata or find its first peer is cancelled before it has any
progress to count, and one that moves a few megabytes at a time can be
cancelled between bursts. I hit this with a torrent that had a single
peer. A torrent now gets one 15 minute window the first time it is seen
outside the queue, and each step forward in progress renews it. One that
makes no progress for 15 minutes is still cancelled, so a dead torrent
now takes about 15 minutes to end instead of five. That trade is the
part I'm least sure about, so the number is easy to change.

Tests cover the grace being requested once, released when the torrent
starts, renewed, and capped, the order of the message and the cancel,
and the start and movement windows. The full suite passes, and I broke
each rule on purpose to check a test catches it.

The two hour ceiling is a guess. If you'd rather have a setting for it,
or a different number, I'm happy to change it.
This commit is contained in:
splitsec2
2026-10-02 22:12:25 -04:00
committed by GitHub
parent bac0b93566
commit 88f73e3bc3
4 changed files with 232 additions and 1 deletions
@@ -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:
+3 -1
View File
@@ -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:
+16
View File
@@ -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")
+169
View File
@@ -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