mirror of
https://github.com/calibrain/shelfmark.git
synced 2026-09-24 13:30:47 +01:00
Fix: Bypass activation on non-fast downloads (#536)
This commit is contained in:
@@ -592,6 +592,8 @@ def _bypass(sb, max_retries: Optional[int] = None, cancel_flag: Optional[Event]
|
|||||||
|
|
||||||
last_challenge_type = None
|
last_challenge_type = None
|
||||||
consecutive_same_challenge = 0
|
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):
|
for try_count in range(max_retries):
|
||||||
_check_cancellation(cancel_flag, "Bypass cancelled by user")
|
_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:
|
if challenge_type == last_challenge_type:
|
||||||
consecutive_same_challenge += 1
|
consecutive_same_challenge += 1
|
||||||
if consecutive_same_challenge >= MAX_CONSECUTIVE_SAME_CHALLENGE:
|
if consecutive_same_challenge >= min_same_challenge_before_abort:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
f"Same challenge ({challenge_type}) detected {consecutive_same_challenge} times - aborting"
|
f"Same challenge ({challenge_type}) detected {consecutive_same_challenge} times - aborting"
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -3,7 +3,7 @@
|
|||||||
import random
|
import random
|
||||||
import time
|
import time
|
||||||
from io import BytesIO
|
from io import BytesIO
|
||||||
from threading import Event
|
from threading import Event, Thread
|
||||||
from typing import Callable, Optional
|
from typing import Callable, Optional
|
||||||
from urllib.parse import urlparse
|
from urllib.parse import urlparse
|
||||||
|
|
||||||
@@ -200,13 +200,32 @@ def html_get_page(
|
|||||||
if use_bypasser_now and _is_cf_bypass_enabled():
|
if use_bypasser_now and _is_cf_bypass_enabled():
|
||||||
logger.debug(f"GET (bypasser): {current_url}")
|
logger.debug(f"GET (bypasser): {current_url}")
|
||||||
if status_callback:
|
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:
|
try:
|
||||||
result = get_bypassed_page(current_url, selector, cancel_flag)
|
result = get_bypassed_page(current_url, selector, cancel_flag)
|
||||||
return result or ""
|
return result or ""
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning(f"Bypasser error: {type(e).__name__}: {e}")
|
logger.warning(f"Bypasser error: {type(e).__name__}: {e}")
|
||||||
return ""
|
return ""
|
||||||
|
finally:
|
||||||
|
heartbeat_stop.set()
|
||||||
|
if heartbeat_thread:
|
||||||
|
heartbeat_thread.join(timeout=1)
|
||||||
|
|
||||||
logger.debug(f"GET: {current_url}")
|
logger.debug(f"GET: {current_url}")
|
||||||
# Try with CF cookies/UA if available (from previous bypass)
|
# Try with CF cookies/UA if available (from previous bypass)
|
||||||
|
|||||||
@@ -50,6 +50,8 @@ _progress_lock = Lock()
|
|||||||
|
|
||||||
# Stall detection - track last activity time per download
|
# Stall detection - track last activity time per download
|
||||||
_last_activity: Dict[str, float] = {}
|
_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
|
STALL_TIMEOUT = 300 # 5 minutes without progress/status update = stalled
|
||||||
|
|
||||||
def search_books(query: str, filters: SearchFilters) -> List[Dict[str, Any]]:
|
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,
|
'cancelled': QueueStatus.CANCELLED,
|
||||||
}
|
}
|
||||||
|
|
||||||
queue_status_enum = status_map.get(status.lower())
|
status_key = status.lower()
|
||||||
if queue_status_enum:
|
queue_status_enum = status_map.get(status_key)
|
||||||
book_queue.update_status(book_id, queue_status_enum)
|
if not queue_status_enum:
|
||||||
|
return
|
||||||
|
|
||||||
# Track activity for stall detection
|
# Always update activity timestamp (used by stall detection) even if the status
|
||||||
with _progress_lock:
|
# event is a duplicate keep-alive update.
|
||||||
_last_activity[book_id] = time.time()
|
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)
|
book_queue.update_status(book_id, queue_status_enum)
|
||||||
if message is not None:
|
|
||||||
book_queue.update_status_message(book_id, message)
|
|
||||||
|
|
||||||
# Broadcast status update via WebSocket
|
# Update status message if provided (empty string clears the message)
|
||||||
if ws_manager:
|
if message is not None:
|
||||||
ws_manager.broadcast_status_update(queue_status())
|
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:
|
def cancel_download(book_id: str) -> bool:
|
||||||
"""Cancel a download."""
|
"""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(task_id, None)
|
||||||
_progress_last_broadcast.pop(f"{task_id}_progress", None)
|
_progress_last_broadcast.pop(f"{task_id}_progress", None)
|
||||||
_last_activity.pop(task_id, 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:
|
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:
|
with ThreadPoolExecutor(max_workers=max_workers, thread_name_prefix="Download") as executor:
|
||||||
active_futures: Dict[Future, str] = {} # Track active download futures
|
active_futures: Dict[Future, str] = {} # Track active download futures
|
||||||
|
stalled_tasks: set[str] = set() # Track tasks already cancelled due to stall
|
||||||
|
|
||||||
while True:
|
while True:
|
||||||
# Clean up completed futures
|
# Clean up completed futures
|
||||||
completed_futures = [f for f in active_futures if f.done()]
|
completed_futures = [f for f in active_futures if f.done()]
|
||||||
for future in completed_futures:
|
for future in completed_futures:
|
||||||
task_id = active_futures.pop(future)
|
task_id = active_futures.pop(future)
|
||||||
|
stalled_tasks.discard(task_id)
|
||||||
try:
|
try:
|
||||||
future.result() # This will raise any exceptions from the worker
|
future.result() # This will raise any exceptions from the worker
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
@@ -523,11 +536,14 @@ def concurrent_download_loop() -> None:
|
|||||||
current_time = time.time()
|
current_time = time.time()
|
||||||
with _progress_lock:
|
with _progress_lock:
|
||||||
for future, task_id in list(active_futures.items()):
|
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)
|
last_active = _last_activity.get(task_id, current_time)
|
||||||
if current_time - last_active > STALL_TIMEOUT:
|
if current_time - last_active > STALL_TIMEOUT:
|
||||||
logger.warning(f"Download stalled for {task_id}, cancelling")
|
logger.warning(f"Download stalled for {task_id}, cancelling")
|
||||||
book_queue.cancel_download(task_id)
|
book_queue.cancel_download(task_id)
|
||||||
book_queue.update_status_message(task_id, f"Download stalled (no activity for {STALL_TIMEOUT}s)")
|
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
|
# Start new downloads if we have capacity
|
||||||
while len(active_futures) < max_workers:
|
while len(active_futures) < max_workers:
|
||||||
|
|||||||
Reference in New Issue
Block a user