diff --git a/shelfmark/bypass/internal_bypasser.py b/shelfmark/bypass/internal_bypasser.py index 445054b8..54123605 100644 --- a/shelfmark/bypass/internal_bypasser.py +++ b/shelfmark/bypass/internal_bypasser.py @@ -592,6 +592,8 @@ def _bypass(sb, max_retries: Optional[int] = None, cancel_flag: Optional[Event] last_challenge_type = None consecutive_same_challenge = 0 + # Allow at least one full pass through all bypass methods before aborting due to a "stuck" challenge. + min_same_challenge_before_abort = max(MAX_CONSECUTIVE_SAME_CHALLENGE, len(BYPASS_METHODS) + 1) for try_count in range(max_retries): _check_cancellation(cancel_flag, "Bypass cancelled by user") @@ -623,7 +625,7 @@ def _bypass(sb, max_retries: Optional[int] = None, cancel_flag: Optional[Event] if challenge_type == last_challenge_type: consecutive_same_challenge += 1 - if consecutive_same_challenge >= MAX_CONSECUTIVE_SAME_CHALLENGE: + if consecutive_same_challenge >= min_same_challenge_before_abort: logger.warning( f"Same challenge ({challenge_type}) detected {consecutive_same_challenge} times - aborting" ) diff --git a/shelfmark/download/http.py b/shelfmark/download/http.py index e3beb080..11d2caaf 100644 --- a/shelfmark/download/http.py +++ b/shelfmark/download/http.py @@ -3,7 +3,7 @@ import random import time from io import BytesIO -from threading import Event +from threading import Event, Thread from typing import Callable, Optional from urllib.parse import urlparse @@ -200,13 +200,32 @@ def html_get_page( if use_bypasser_now and _is_cf_bypass_enabled(): logger.debug(f"GET (bypasser): {current_url}") if status_callback: - status_callback("resolving", "Bypassing protection") + status_callback("resolving", "Bypassing protection...") + heartbeat_stop = Event() + heartbeat_thread: Optional[Thread] = None + if status_callback: + def _heartbeat() -> None: + # Keep the download "alive" during long bypass operations so the orchestrator + # doesn't flag it as stalled. + while not heartbeat_stop.wait(timeout=30): + if cancel_flag and cancel_flag.is_set(): + return + try: + status_callback("resolving", "Bypassing protection...") + except Exception: + return + heartbeat_thread = Thread(target=_heartbeat, daemon=True, name="BypassHeartbeat") + heartbeat_thread.start() try: result = get_bypassed_page(current_url, selector, cancel_flag) return result or "" except Exception as e: logger.warning(f"Bypasser error: {type(e).__name__}: {e}") return "" + finally: + heartbeat_stop.set() + if heartbeat_thread: + heartbeat_thread.join(timeout=1) logger.debug(f"GET: {current_url}") # Try with CF cookies/UA if available (from previous bypass) diff --git a/shelfmark/download/orchestrator.py b/shelfmark/download/orchestrator.py index df28f2a0..092f3520 100644 --- a/shelfmark/download/orchestrator.py +++ b/shelfmark/download/orchestrator.py @@ -50,6 +50,8 @@ _progress_lock = Lock() # Stall detection - track last activity time per download _last_activity: Dict[str, float] = {} +# De-duplicate status updates (keep-alive updates shouldn't spam clients) +_last_status_event: Dict[str, Tuple[str, Optional[str]]] = {} STALL_TIMEOUT = 300 # 5 minutes without progress/status update = stalled def search_books(query: str, filters: SearchFilters) -> List[Dict[str, Any]]: @@ -398,21 +400,29 @@ def update_download_status(book_id: str, status: str, message: Optional[str] = N 'cancelled': QueueStatus.CANCELLED, } - queue_status_enum = status_map.get(status.lower()) - if queue_status_enum: - book_queue.update_status(book_id, queue_status_enum) + status_key = status.lower() + queue_status_enum = status_map.get(status_key) + if not queue_status_enum: + return - # Track activity for stall detection - with _progress_lock: - _last_activity[book_id] = time.time() + # Always update activity timestamp (used by stall detection) even if the status + # event is a duplicate keep-alive update. + with _progress_lock: + _last_activity[book_id] = time.time() + status_event = (status_key, message) + if _last_status_event.get(book_id) == status_event: + return + _last_status_event[book_id] = status_event - # Update status message if provided (empty string clears the message) - if message is not None: - book_queue.update_status_message(book_id, message) + book_queue.update_status(book_id, queue_status_enum) - # Broadcast status update via WebSocket - if ws_manager: - ws_manager.broadcast_status_update(queue_status()) + # Update status message if provided (empty string clears the message) + if message is not None: + book_queue.update_status_message(book_id, message) + + # Broadcast status update via WebSocket + if ws_manager: + ws_manager.broadcast_status_update(queue_status()) def cancel_download(book_id: str) -> bool: """Cancel a download.""" @@ -450,6 +460,7 @@ def _cleanup_progress_tracking(task_id: str) -> None: _progress_last_broadcast.pop(task_id, None) _progress_last_broadcast.pop(f"{task_id}_progress", None) _last_activity.pop(task_id, None) + _last_status_event.pop(task_id, None) def _process_single_download(task_id: str, cancel_flag: Event) -> None: @@ -508,12 +519,14 @@ def concurrent_download_loop() -> None: with ThreadPoolExecutor(max_workers=max_workers, thread_name_prefix="Download") as executor: active_futures: Dict[Future, str] = {} # Track active download futures + stalled_tasks: set[str] = set() # Track tasks already cancelled due to stall while True: # Clean up completed futures completed_futures = [f for f in active_futures if f.done()] for future in completed_futures: task_id = active_futures.pop(future) + stalled_tasks.discard(task_id) try: future.result() # This will raise any exceptions from the worker except Exception as e: @@ -523,11 +536,14 @@ def concurrent_download_loop() -> None: current_time = time.time() with _progress_lock: for future, task_id in list(active_futures.items()): + if task_id in stalled_tasks: + continue last_active = _last_activity.get(task_id, current_time) if current_time - last_active > STALL_TIMEOUT: logger.warning(f"Download stalled for {task_id}, cancelling") book_queue.cancel_download(task_id) book_queue.update_status_message(task_id, f"Download stalled (no activity for {STALL_TIMEOUT}s)") + stalled_tasks.add(task_id) # Start new downloads if we have capacity while len(active_futures) < max_workers: