Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions src/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
# ============================
KB = 1024 # Number of bytes in a kilobyte.
CHUNK_SIZE = 16 * KB # The size of each chunk when downloading files.
MAX_WORKERS = 5 # The number of concurrent image downloads.

# ============================
# HTTP / Network
Expand Down
65 changes: 49 additions & 16 deletions src/crawler/crawler.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,9 @@ def __init__(
live_manager: LiveManager,
) -> None:
"""Initialize the Crawler with album URL, initial soup, and live manager."""
self.url = url
parsed = urlparse(url)
path = parsed.path if parsed.path.endswith("/") else f"{parsed.path}/"
self.url = f"{parsed.scheme}://{parsed.netloc}{path}"
self.initial_soup = initial_soup
self.live_manager = live_manager
self.album_pages = self._generate_album_pages()
Expand All @@ -53,22 +55,42 @@ def collect_album_pages_soups(self) -> list[BeautifulSoup]:
)
return album_pages_soups

def get_reloaded_pages(self, picture_pages: list[str]) -> list[str]:
"""Generate reloaded image page URLs."""
reloaded_pages = []
for picture_page in picture_pages:
def get_reloaded_page(self, picture_page: str) -> str | None:
"""Generate reloaded image page URL for a single picture page."""
try:
soup = fetch_page(picture_page)
nl_container = soup.find("a", {"id": "loadfail", "onclick": True})
nl_value = re.search(r"nl\('([^']+)'\)", nl_container["onclick"]).group(1)
if not nl_container or "onclick" not in nl_container.attrs:
self.live_manager.update_log(
"Missing 'nl' container",
f"No 'loadfail' link with 'onclick' found for {picture_page}.",
)
return None

if nl_value:
reloaded_page = generate_reloaded_page(picture_page, nl_value)
reloaded_pages.append(reloaded_page)
else:
match = re.search(r"nl\('([^']+)'\)", nl_container["onclick"])
if not match:
self.live_manager.update_log(
"Missing 'nl' value", f"No 'nl' value found for {picture_page}.",
"Missing 'nl' value",
f"No 'nl' value found in onclick for {picture_page}.",
)
return None

nl_value = match.group(1)
return generate_reloaded_page(picture_page, nl_value)
except Exception as err:
self.live_manager.update_log(
"Crawler error",
f"Error getting reloaded page for {picture_page}: {err}",
)
return None

def get_reloaded_pages(self, picture_pages: list[str]) -> list[str]:
"""Generate reloaded image page URLs."""
reloaded_pages = []
for picture_page in picture_pages:
reloaded_page = self.get_reloaded_page(picture_page)
if reloaded_page:
reloaded_pages.append(reloaded_page)
return reloaded_pages

def _generate_album_pages(self) -> list[str]:
Expand All @@ -79,9 +101,20 @@ def _generate_album_pages(self) -> list[str]:
{"href": pattern, "onclick": "return false"},
)

last_page_url = next_pages[-2].get("href")
match = re.search(r"\?p=(\d+)", last_page_url)
last_page = int(match.group(1))
album_pages = [f"{self.url}?p={page}" for page in range(1, last_page)]
album_pages.append(last_page_url)
if not next_pages:
return []

page_numbers = []
for a in next_pages:
href = a.get("href")
if href:
match = re.search(r"\?p=(\d+)", href)
if match:
page_numbers.append(int(match.group(1)))

if not page_numbers:
return []

last_page = max(page_numbers)
album_pages = [f"{self.url}?p={page}" for page in range(1, last_page + 1)]
return album_pages
99 changes: 77 additions & 22 deletions src/downloader/album_downloader.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,13 +8,17 @@
import random
import time
from pathlib import Path
from urllib.parse import urlparse

import requests
from requests import Response, Session

from concurrent.futures import ThreadPoolExecutor, as_completed

from src.config import (
CHUNK_SIZE,
HTTP_RATE_LIMIT,
MAX_WORKERS,
RATE_LIMIT_SLEEPING_TIME,
)
from src.crawler.crawler import Crawler
Expand All @@ -35,7 +39,9 @@ class AlbumDownloader:

def __init__(self, url: str, live_manager: LiveManager) -> None:
"""Initialize the AlbumDownloader with album URL and live manager."""
self.url = url
parsed = urlparse(url)
path = parsed.path if parsed.path.endswith("/") else f"{parsed.path}/"
self.url = f"{parsed.scheme}://{parsed.netloc}{path}"
self.live_manager = live_manager
self.initial_soup = fetch_page(self.url)
self.crawler = Crawler(
Expand All @@ -59,10 +65,9 @@ def download_album(self) -> None:
for current_task, soup in enumerate(album_pages_soups):
containers = soup.find_all("a", {"href": True})
picture_pages = get_picture_pages(containers)
reloaded_pages = self.crawler.get_reloaded_pages(picture_pages)

failed_downloads = self._extract_and_download(
session, reloaded_pages, current_task,
session, picture_pages, current_task,
)

if failed_downloads:
Expand Down Expand Up @@ -97,41 +102,91 @@ def download_picture(
def _extract_and_download(
self,
session: Session,
reloaded_pages: list[str],
picture_pages: list[str],
current_task: int,
) -> list[str]:
"""Extract image links and download them."""
"""Extract image links and download them concurrently."""
failed_downloads = []
num_pictures = len(reloaded_pages)
num_pictures = len(picture_pages)
task = self.live_manager.add_task(current_task=current_task, total=num_pictures)

for reloaded_page in reloaded_pages:
soup = fetch_page(reloaded_page)
download_link_container = soup.find("img", {"id": "img", "src": True})
download_link = download_link_container["src"]
# Thread lock for safely appending to failed_downloads list
import threading
failed_lock = threading.Lock()

def download_worker(picture_page: str) -> None:
nonlocal failed_downloads

# 1. Fetch the picture page to extract the 'nl' value
reloaded_page = self.crawler.get_reloaded_page(picture_page)
if not reloaded_page:
with failed_lock:
failed_downloads.append(picture_page)
self.live_manager.update_task(task, advance=1)
return

# 2. Fetch the reloaded page to extract the direct image source link
try:
soup = fetch_page(reloaded_page)
download_link_container = soup.find("img", {"id": "img", "src": True})
if not download_link_container:
self.live_manager.update_log(
"Image not found",
f"Could not find img with id='img' on {reloaded_page}.",
)
with failed_lock:
failed_downloads.append(picture_page)
self.live_manager.update_task(task, advance=1)
return
download_link = download_link_container["src"]
except Exception as err:
self.live_manager.update_log(
"Page fetch error",
f"Error reading reloaded page {reloaded_page}: {err}",
)
with failed_lock:
failed_downloads.append(picture_page)
self.live_manager.update_task(task, advance=1)
return

# 3. Download the actual image file with retries
headers = prepare_headers(download_link)
session.headers.update(headers)
response = fetch_with_retries(session, download_link, self.live_manager)
response = fetch_with_retries(
session=session,
url=download_link,
live_manager=self.live_manager,
headers=headers,
)

if response is None:
self.live_manager.update_log(
"Failed download",
f"None response from {download_link}, check the log file",
)
failed_downloads.append(download_link)
with failed_lock:
failed_downloads.append(download_link)
write_on_session_log(download_link)
continue

if response.status_code == HTTP_RATE_LIMIT:
self.live_manager.update_log(
"Rate limit",
"Rate limit hit. Sleeping for a while...",
)
time.sleep(RATE_LIMIT_SLEEPING_TIME)
self.live_manager.update_task(task, advance=1)
return

filename = download_link.split("/")[-1]
self.download_picture(response, filename, task)
time.sleep(random.uniform(1.5, 4.0)) # noqa: S311

# Polite sleep between tasks in each thread
time.sleep(random.uniform(1.0, 3.0)) # noqa: S311

with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor:
futures = [
executor.submit(download_worker, picture_page)
for picture_page in picture_pages
]
for future in as_completed(futures):
try:
future.result()
except Exception as err:
self.live_manager.update_log(
"Worker error",
f"Unhandled exception in download worker: {err}",
)

return failed_downloads
36 changes: 29 additions & 7 deletions src/downloader/download_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,9 @@
from typing import TYPE_CHECKING
from urllib.parse import urlparse

from requests.exceptions import ConnectTimeout, RequestException, Timeout
from requests.exceptions import ConnectTimeout, HTTPError, RequestException, Timeout

from src.config import CONNECTION_TIMEOUT, prepare_user_agent
from src.config import CONNECTION_TIMEOUT, RATE_LIMIT_SLEEPING_TIME, prepare_user_agent

if TYPE_CHECKING:
from requests import Response, Session
Expand Down Expand Up @@ -46,25 +46,47 @@ def fetch_with_retries(
url: str,
live_manager: LiveManager,
retries: int = 5,
headers: dict | None = None,
) -> Response | None:
"""Fetch a URL with retry logic and exponential backoff."""
"""Fetch a URL with retry logic and exponential backoff/rate-limit handling."""
for attempt in range(retries):
try:
response = session.get(url, timeout=CONNECTION_TIMEOUT)
response = session.get(url, timeout=CONNECTION_TIMEOUT, headers=headers)
response.raise_for_status()

except HTTPError as http_err:
status_code = (
http_err.response.status_code
if http_err.response is not None
else 0
)
if status_code in (429, 509):
live_manager.update_log(
"Rate limit",
f"Rate limit ({status_code}) hit for {url}. "
f"Sleeping for {RATE_LIMIT_SLEEPING_TIME} seconds...",
)
time.sleep(RATE_LIMIT_SLEEPING_TIME)
continue
else:
live_manager.update_log(
"Request failed",
f"HTTP error {status_code} for {url}: {http_err}",
)
break

except ConnectTimeout as conn_timeout:
live_manager.update_log(
"ConnectTimeout",
f"Connect timeout for {url}: {conn_timeout}. "
"Retrying ({attempt + 1}/{retries})...",
f"Retrying ({attempt + 1}/{retries})...",
)

except Timeout as timeout_err:
live_manager.update_log(
"Timeout error",
f"Timeout for {url}: {timeout_err}. "
"Retrying ({attempt + 1}/{retries})...",
f"Retrying ({attempt + 1}/{retries})...",
)

except RequestException as req_err:
Expand All @@ -80,7 +102,7 @@ def fetch_with_retries(
if attempt < retries - 1:
live_manager.update_log(
"Fetch attempt failed",
f"Fetch attemp failed for {url}. Retrying ({attempt + 1}/{retries})...",
f"Fetch attempt failed for {url}. Retrying ({attempt + 1}/{retries})...",
)
delay = 2 ** (attempt + 1) + random.uniform(1, 2) # noqa: S311
time.sleep(delay)
Expand Down
7 changes: 5 additions & 2 deletions src/managers/live_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
from __future__ import annotations

import datetime
import threading
import time
from typing import TYPE_CHECKING

Expand Down Expand Up @@ -36,6 +37,7 @@ def __init__(
self.progress_manager = progress_manager
self.progress_table = self.progress_manager.create_progress_table()
self.logger = logger
self._lock = threading.Lock()
self.live = Live(
self._render_live_view(),
refresh_per_second=refresh_per_second,
Expand Down Expand Up @@ -63,8 +65,9 @@ def update_task(

def update_log(self, event: str, details: str) -> None:
"""Log an event and refreshes the live display."""
self.logger.log(event, details)
self.live.update(self._render_live_view())
with self._lock:
self.logger.log(event, details)
self.live.update(self._render_live_view())

def start(self) -> None:
"""Start the live display."""
Expand Down
17 changes: 10 additions & 7 deletions src/managers/progress_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
from __future__ import annotations

from collections import deque
import threading

from rich.panel import Panel
from rich.progress import (
Expand Down Expand Up @@ -38,6 +39,7 @@ def __init__(
self.task_progress = self._create_progress_bar()
self.num_tasks = 0
self.overall_buffer = deque(maxlen=overall_buffer_size)
self._lock = threading.Lock()

def add_overall_task(self, description: str, num_tasks: int) -> None:
"""Add an overall progress task with a given description and total tasks."""
Expand Down Expand Up @@ -66,13 +68,14 @@ def update_task(
visible: bool = True,
) -> None:
"""Update the progress of an individual task and the overall progress."""
self.task_progress.update(
task_id,
completed=completed if completed is not None else None,
advance=advance if completed is None else None,
visible=visible,
)
self._update_overall_task(task_id)
with self._lock:
self.task_progress.update(
task_id,
completed=completed if completed is not None else None,
advance=advance if completed is None else None,
visible=visible,
)
self._update_overall_task(task_id)

def create_progress_table(self) -> Table:
"""Create a formatted progress table for tracking the download."""
Expand Down