Files
shelfmark/tests/download/test_orchestrator_lifecycle.py
T
CaliBrain cc1a95f965 Fix protection bypass cancelled by stall detection at exactly 300s (#1184)
A download that hits Cloudflare hung on "Bypassing protection..." for
five
minutes and then died, regardless of which bypasser was configured.

html_get_page() started a BypassHeartbeat thread to keep the download
marked
alive during a bypass, but the thread had no loop: it fired one status
event
and returned. Even with the loop restored it could not have worked,
because
update_download_status() dedupes identical (status, message) tuples and
returns before refreshing _last_activity, and the heartbeat re-sent the
byte-identical payload already emitted just above it.

So _last_activity was frozen for the whole bypass, while both bypassers
are
allowed to run longer than STALL_TIMEOUT (external FlareSolverr ~394s at
default settings, internal 420s per get() call). The watchdog always
won.
From a reporter's log: 403 at 07:04:33.390, cancelled at 07:09:33.987 -
exactly 300.000s, and 41s before the bypasser would have finished and
reported the real error, an HTTP 500 from FlareSolverr the user never
saw.

The regression is not one commit. 1f093de (#536) added the heartbeat and
the
dedup together and refreshed activity before the dedup return, so it
worked.
ff094be (#832) moved the refresh below that return while tightening
stall
detection for #823. 3a3a3ce (#845) then deleted the heartbeat's while
loop
to silence a B023 lint, removing the last evidence of intent.

The dedup itself is correct and stays: a keep-alive that ticks on a
timer
proves nothing about whether an operation is progressing, so letting it
refresh the stall clock would make a wedged download immortal. Split the
two
concerns instead.

Add shelfmark/download/activity.py. A long single-shot operation
declares its
own upper bound once, over a sentinel status carried on the existing
status_callback channel - so no new parameter has to be threaded through
every
handler, post-processor and output module. The orchestrator intercepts
the
sentinel in its per-task closure and records an absolute deadline in
_activity_grace, which stall detection honours alongside STALL_TIMEOUT.
The
grace never extends itself and is clamped to
_MAX_ACTIVITY_GRACE_SECONDS, so
an operation that overruns its own declared budget is still cancelled.

Each bypasser now reports max_duration_seconds() derived from its own
retry
and timeout settings, and http.py asks whichever is active, plus 30s of
slack
so the bypasser's own deadline expires first and the user sees its real
failure. On that path html_get_page() also emits
status_callback("error", ...)
rather than silently returning an empty page.

Three further fixes on the same code path:

- Extract the watchdog into _find_stalled_tasks() and
_cancel_stalled_task().
It was the only place holding _progress_lock across a call into
book_queue,
whose terminal-status hooks reach a sqlite write that gevent does not
patch,
blocking the hub and every download worker. It now holds the lock for
dict
  reads only.
- Bound _CDP_WORKER.run(), which waited with timeout=None while holding
the
module-wide LOCKED, so a single wedged in-process CDP session blocked
every
  subsequent bypass forever on non-Docker installs.
- Broaden the coordinator loop's except clause back to Exception, with
  escalating backoff. 8d98e12 (#868) narrowed it to a six-type tuple to
silence BLE001, which let gevent's LoopExit and similar kill the only
thread
driving the download queue - undoing #832's fix for #823 and resurfacing
it
  as #1166. GreenletExit and gevent.Timeout still propagate.

Fixes #1001
Refs #1166, #823
2026-08-11 02:25:28 -04:00

228 lines
7.6 KiB
Python

from __future__ import annotations
from unittest.mock import ANY, MagicMock
import pytest
class _StopLoop(BaseException):
"""Sentinel used to stop the infinite coordinator loop during tests."""
class _UnexpectedCall(BaseException):
"""Guard-rail failure that the coordinator's `except Exception` must not swallow."""
class _FakeExecutor:
def __init__(self, *args, **kwargs) -> None:
self.args = args
self.kwargs = kwargs
def __enter__(self) -> _FakeExecutor:
return self
def __exit__(self, exc_type, exc, tb) -> bool:
return False
def submit(self, *args, **kwargs): # pragma: no cover - not expected in these tests
raise _UnexpectedCall("submit() should not be called in this test")
class _StopCoordinator(BaseException):
"""Sentinel used to stop a real coordinator thread cleanly in tests."""
def test_concurrent_download_loop_logs_and_recovers_after_loop_error(monkeypatch):
import shelfmark.download.orchestrator as orchestrator
call_count = 0
def fake_get_next():
nonlocal call_count
call_count += 1
if call_count == 1:
raise RuntimeError("boom")
return None
sleep_delays: list[float] = []
def fake_sleep(delay: float) -> None:
sleep_delays.append(delay)
if len(sleep_delays) >= 2:
raise _StopLoop()
mock_queue = MagicMock()
mock_queue.get_next.side_effect = fake_get_next
error_trace = MagicMock()
monkeypatch.setattr(orchestrator, "book_queue", mock_queue)
monkeypatch.setattr(orchestrator, "ThreadPoolExecutor", _FakeExecutor)
monkeypatch.setattr(orchestrator.time, "sleep", fake_sleep)
monkeypatch.setattr(orchestrator.logger, "error_trace", error_trace)
with pytest.raises(_StopLoop):
orchestrator.concurrent_download_loop()
assert mock_queue.get_next.call_count == 2
error_trace.assert_called_once_with("Download coordinator loop error: %s", ANY)
assert sleep_delays == [
orchestrator.COORDINATOR_LOOP_ERROR_RETRY_DELAY,
orchestrator.config.MAIN_LOOP_SLEEP_TIME,
]
def test_concurrent_download_loop_survives_exceptions_outside_the_legacy_list(monkeypatch):
"""Regression for #823/#1166: the coordinator must not die on an unlisted exception.
`gevent.exceptions.LoopExit` and `StopIteration` are Exceptions that the old narrow
except tuple let through, silently killing the only thread that drives the queue.
"""
import shelfmark.download.orchestrator as orchestrator
raised: list[type[BaseException]] = []
escapes = [StopIteration, ZeroDivisionError, KeyboardInterrupt]
def fake_get_next():
if escapes:
exc = escapes.pop(0)
if exc is KeyboardInterrupt:
# BaseException: must still propagate and stop the loop.
raise KeyboardInterrupt
raised.append(exc)
raise exc("boom")
return None
mock_queue = MagicMock()
mock_queue.get_next.side_effect = fake_get_next
monkeypatch.setattr(orchestrator, "book_queue", mock_queue)
monkeypatch.setattr(orchestrator, "ThreadPoolExecutor", _FakeExecutor)
monkeypatch.setattr(orchestrator.time, "sleep", lambda _delay: None)
monkeypatch.setattr(orchestrator.logger, "error_trace", MagicMock())
with pytest.raises(KeyboardInterrupt):
orchestrator.concurrent_download_loop()
assert raised == [StopIteration, ZeroDivisionError]
def test_concurrent_download_loop_backs_off_on_repeated_errors(monkeypatch):
"""A persistent failure must not spin the loop at 1Hz forever."""
import shelfmark.download.orchestrator as orchestrator
sleep_delays: list[float] = []
def fake_sleep(delay: float) -> None:
sleep_delays.append(delay)
if len(sleep_delays) >= 4:
raise _StopLoop()
mock_queue = MagicMock()
mock_queue.get_next.side_effect = RuntimeError("persistent boom")
monkeypatch.setattr(orchestrator, "book_queue", mock_queue)
monkeypatch.setattr(orchestrator, "ThreadPoolExecutor", _FakeExecutor)
monkeypatch.setattr(orchestrator.time, "sleep", fake_sleep)
monkeypatch.setattr(orchestrator.logger, "error_trace", MagicMock())
with pytest.raises(_StopLoop):
orchestrator.concurrent_download_loop()
base = orchestrator.COORDINATOR_LOOP_ERROR_RETRY_DELAY
# First delay stays unchanged for a normal transient blip, then doubles.
assert sleep_delays == [base, base * 2, base * 4, base * 8]
assert max(sleep_delays) <= orchestrator._COORDINATOR_LOOP_ERROR_MAX_DELAY
def test_concurrent_download_loop_recovers_and_processes_task_after_transient_loop_error(
monkeypatch,
):
import threading
import shelfmark.download.orchestrator as orchestrator
processed = threading.Event()
class FlakyQueue:
def __init__(self) -> None:
self.calls = 0
def get_next(self):
self.calls += 1
if self.calls == 1:
raise RuntimeError("boom")
if self.calls == 2:
return ("task-1", threading.Event())
if processed.is_set():
raise _StopCoordinator()
return None
# These raise _UnexpectedCall rather than AssertionError because the coordinator
# loop catches Exception: an AssertionError here would be swallowed and logged,
# letting the test pass vacuously instead of reporting the unexpected call.
def cancel_download(self, task_id: str) -> None: # pragma: no cover - unused
raise _UnexpectedCall(f"cancel_download unexpectedly called for {task_id}")
def update_status_message(
self, task_id: str, message: str
) -> None: # pragma: no cover - unused
raise _UnexpectedCall(
f"update_status_message unexpectedly called for {task_id}: {message}"
)
queue = FlakyQueue()
error_trace = MagicMock()
monkeypatch.setattr(orchestrator, "book_queue", queue)
monkeypatch.setattr(
orchestrator,
"_process_single_download",
lambda task_id, cancel_flag: processed.set(),
)
monkeypatch.setattr(orchestrator.logger, "error_trace", error_trace)
monkeypatch.setattr(orchestrator, "COORDINATOR_LOOP_ERROR_RETRY_DELAY", 0.01)
monkeypatch.setattr(orchestrator.config, "MAX_CONCURRENT_DOWNLOADS", 1, raising=False)
monkeypatch.setattr(orchestrator.config, "MAIN_LOOP_SLEEP_TIME", 0.01, raising=False)
def run_loop() -> None:
try:
orchestrator.concurrent_download_loop()
except _StopCoordinator:
pass
thread = threading.Thread(target=run_loop, daemon=True, name="TestDownloadCoordinator")
thread.start()
assert processed.wait(timeout=1.0) is True
thread.join(timeout=1.0)
assert thread.is_alive() is False
assert queue.calls >= 3
error_trace.assert_called_once_with("Download coordinator loop error: %s", ANY)
def test_start_replaces_dead_coordinator_thread(monkeypatch):
import shelfmark.download.orchestrator as orchestrator
dead_thread = MagicMock()
dead_thread.is_alive.return_value = False
new_thread = MagicMock()
new_thread.is_alive.return_value = True
thread_factory = MagicMock(return_value=new_thread)
monkeypatch.setattr(orchestrator, "_coordinator_thread", dead_thread)
monkeypatch.setattr(orchestrator.threading, "Thread", thread_factory)
orchestrator.start()
thread_factory.assert_called_once_with(
target=orchestrator.concurrent_download_loop,
daemon=True,
name="DownloadCoordinator",
)
new_thread.start.assert_called_once_with()
assert orchestrator._coordinator_thread is new_thread