From 260200e2519fe6023fdcba2d08dba79c6c0fafa3 Mon Sep 17 00:00:00 2001 From: Nicolas Fryder Date: Sat, 27 Jun 2026 15:15:29 +0200 Subject: [PATCH] feat(crawler): nouveau service crawler server-side via FlareSolverr MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - apps/crawler/ : service Python complet, remplace les scripts locaux - FlareSolverr pour bypasser Cloudflare (cookies CF → session requests) - Crawl incrémental : /archives/band-list/by/created et /by/modified - Crawl complet Europe : pagination AJAX /browse/ajax-country/ - Enrichissement : pages individuelles de bands (themes, membres, label, hash) - Écriture directe en DB (upserts bulk, idempotents) - Scheduler intégré (schedule library) : incrémental 4h, enrich 2h, full le 1er du mois - Tracking via crawl_run et crawl_checkpoint (migration 004) - docker-compose.dev.yml : flaresolverr + crawler ajoutés Full crawl désactivé en dev (CRAWLER_SCHED_FULL_DAY=0) Co-Authored-By: Claude Sonnet 4.6 --- apps/crawler/Dockerfile | 9 ++ apps/crawler/requirements.txt | 5 + apps/crawler/src/__init__.py | 0 apps/crawler/src/config.py | 28 ++++ apps/crawler/src/db.py | 233 +++++++++++++++++++++++++++++ apps/crawler/src/europe_codes.py | 8 + apps/crawler/src/flaresolverr.py | 85 +++++++++++ apps/crawler/src/jobs.py | 169 +++++++++++++++++++++ apps/crawler/src/ma_http.py | 244 +++++++++++++++++++++++++++++++ apps/crawler/src/main.py | 106 ++++++++++++++ apps/crawler/src/polite.py | 12 ++ apps/crawler/src/scraper_band.py | 98 +++++++++++++ docker-compose.dev.yml | 25 +++- 13 files changed, 1021 insertions(+), 1 deletion(-) create mode 100644 apps/crawler/Dockerfile create mode 100644 apps/crawler/requirements.txt create mode 100644 apps/crawler/src/__init__.py create mode 100644 apps/crawler/src/config.py create mode 100644 apps/crawler/src/db.py create mode 100644 apps/crawler/src/europe_codes.py create mode 100644 apps/crawler/src/flaresolverr.py create mode 100644 apps/crawler/src/jobs.py create mode 100644 apps/crawler/src/ma_http.py create mode 100644 apps/crawler/src/main.py create mode 100644 apps/crawler/src/polite.py create mode 100644 apps/crawler/src/scraper_band.py diff --git a/apps/crawler/Dockerfile b/apps/crawler/Dockerfile new file mode 100644 index 0000000..fd26a5e --- /dev/null +++ b/apps/crawler/Dockerfile @@ -0,0 +1,9 @@ +FROM python:3.12-slim + +WORKDIR /app +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +COPY src ./src + +CMD ["python", "-m", "src.main"] diff --git a/apps/crawler/requirements.txt b/apps/crawler/requirements.txt new file mode 100644 index 0000000..26d12a9 --- /dev/null +++ b/apps/crawler/requirements.txt @@ -0,0 +1,5 @@ +requests==2.32.3 +psycopg2-binary==2.9.9 +beautifulsoup4==4.12.3 +lxml==5.3.0 +schedule==1.2.2 diff --git a/apps/crawler/src/__init__.py b/apps/crawler/src/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/apps/crawler/src/config.py b/apps/crawler/src/config.py new file mode 100644 index 0000000..b063018 --- /dev/null +++ b/apps/crawler/src/config.py @@ -0,0 +1,28 @@ +import os + +DATABASE_URL = os.environ["DATABASE_URL"] +FLARESOLVERR_URL = os.getenv("FLARESOLVERR_URL", "http://flaresolverr:8191") + +# Delays between requests (seconds) +LIST_MIN_DELAY = float(os.getenv("CRAWLER_LIST_MIN_DELAY", "2.0")) +LIST_MAX_DELAY = float(os.getenv("CRAWLER_LIST_MAX_DELAY", "4.0")) +BAND_MIN_DELAY = float(os.getenv("CRAWLER_BAND_MIN_DELAY", "2.5")) +BAND_MAX_DELAY = float(os.getenv("CRAWLER_BAND_MAX_DELAY", "5.0")) +# Cooldown pause every N band pages +COOLDOWN_EVERY = int(os.getenv("CRAWLER_COOLDOWN_EVERY", "60")) +COOLDOWN_MIN = float(os.getenv("CRAWLER_COOLDOWN_MIN", "15")) +COOLDOWN_MAX = float(os.getenv("CRAWLER_COOLDOWN_MAX", "30")) + +# FlareSolverr request timeout (ms) +FS_TIMEOUT_MS = int(os.getenv("CRAWLER_FS_TIMEOUT_MS", "90000")) + +# How many bands to request per AJAX page (max MA allows) +AJAX_PAGE_SIZE = int(os.getenv("CRAWLER_AJAX_PAGE_SIZE", "500")) + +# Schedules (hours between runs, 0 = disabled) +SCHED_INCREMENTAL_H = int(os.getenv("CRAWLER_SCHED_INCREMENTAL_H", "4")) +SCHED_ENRICH_H = int(os.getenv("CRAWLER_SCHED_ENRICH_H", "1")) +# Full crawl: day of month (1-28), 0 = disabled +SCHED_FULL_DAY = int(os.getenv("CRAWLER_SCHED_FULL_DAY", "1")) + +MA_BASE = "https://www.metal-archives.com" diff --git a/apps/crawler/src/db.py b/apps/crawler/src/db.py new file mode 100644 index 0000000..2ccba3e --- /dev/null +++ b/apps/crawler/src/db.py @@ -0,0 +1,233 @@ +""" +Écriture directe en DB. Toutes les opérations passent par des upserts idempotents. +""" +import logging +from contextlib import contextmanager +from datetime import datetime, timezone +from typing import Any, Dict, List, Optional + +import psycopg2 +import psycopg2.extras + +from .config import DATABASE_URL + +log = logging.getLogger(__name__) + + +@contextmanager +def get_conn(): + conn = psycopg2.connect(DATABASE_URL) + try: + yield conn + conn.commit() + except Exception: + conn.rollback() + raise + finally: + conn.close() + + +def now_utc() -> datetime: + return datetime.now(timezone.utc) + + +# ------------------------------------------------------------------ +# Bands +# ------------------------------------------------------------------ + +def upsert_bands(bands: List[Dict[str, Any]]) -> Dict[str, int]: + """ + Upsert en bulk. Préserve les champs existants si la nouvelle valeur est None. + Retourne {"inserted": N, "updated": N}. + """ + if not bands: + return {"inserted": 0, "updated": 0} + + ts = now_utc() + rows = [] + for b in bands: + rows.append(( + b["ma_id"], + b.get("name") or "", + b.get("country"), + b.get("location_text"), + b.get("status"), + b.get("genre"), + b.get("formed_year"), + b.get("themes"), + b.get("enriched", False), + b.get("crawled_at"), + b.get("crawled_hash"), + b.get("ma_created_at"), + b.get("ma_modified_at"), + ts, # first_seen_at (ignoré si déjà set) + ts, # created_at (ignoré si déjà set) + )) + + sql = """ + INSERT INTO bands + (ma_id, name, country, location_text, status, genre, formed_year, + themes, enriched, crawled_at, crawled_hash, ma_created_at, ma_modified_at, + first_seen_at, created_at) + VALUES %s + ON CONFLICT (ma_id) DO UPDATE SET + name = COALESCE(EXCLUDED.name, bands.name), + country = COALESCE(EXCLUDED.country, bands.country), + location_text = COALESCE(EXCLUDED.location_text, bands.location_text), + status = COALESCE(EXCLUDED.status, bands.status), + genre = COALESCE(EXCLUDED.genre, bands.genre), + formed_year = COALESCE(EXCLUDED.formed_year, bands.formed_year), + themes = COALESCE(EXCLUDED.themes, bands.themes), + enriched = GREATEST(EXCLUDED.enriched, bands.enriched), + crawled_at = COALESCE(EXCLUDED.crawled_at, bands.crawled_at), + crawled_hash = COALESCE(EXCLUDED.crawled_hash, bands.crawled_hash), + ma_created_at = COALESCE(EXCLUDED.ma_created_at, bands.ma_created_at), + ma_modified_at= COALESCE(EXCLUDED.ma_modified_at, bands.ma_modified_at) + RETURNING (xmax = 0) AS was_inserted + """ + + with get_conn() as conn: + with conn.cursor() as cur: + results = psycopg2.extras.execute_values(cur, sql, rows, fetch=True) + inserted = sum(1 for r in results if r[0]) + updated = len(results) - inserted + + log.info(f"[db] upsert_bands: {inserted} inserted, {updated} updated") + return {"inserted": inserted, "updated": updated} + + +def upsert_band_enriched(ma_id: int, data: Dict[str, Any], html_hash: str) -> bool: + """Met à jour les champs d'enrichissement d'un band.""" + sql = """ + UPDATE bands SET + status = COALESCE($2, status), + genre = COALESCE($3, genre), + themes = COALESCE($4, themes), + formed_year = COALESCE($5, formed_year), + data = data || $6::jsonb, + enriched = true, + crawled_at = $7, + crawled_hash = $8, + ma_created_at = COALESCE($9, ma_created_at), + ma_modified_at= COALESCE($10, ma_modified_at) + WHERE ma_id = $1 + """ + ts = now_utc() + with get_conn() as conn: + with conn.cursor() as cur: + import json + cur.execute(sql, ( + ma_id, + data.get("status"), + data.get("genre"), + data.get("themes"), + _parse_year(data.get("formed_in")), + json.dumps({"band_page": data, "enriched_at": ts.isoformat()}), + ts, + html_hash, + data.get("ma_created_at"), + data.get("ma_modified_at"), + )) + return cur.rowcount > 0 + + +def get_bands_to_enrich(country: Optional[str] = None, limit: int = 50) -> List[Dict]: + """Bands qui n'ont jamais été enrichies (data->'band_page' IS NULL).""" + sql = """ + SELECT ma_id, data->>'url' AS url, name, country + FROM bands + WHERE data->'band_page' IS NULL + AND data->>'url' IS NOT NULL + """ + params = [] + if country: + sql += " AND country = %s" + params.append(country) + sql += " ORDER BY first_seen_at ASC NULLS LAST LIMIT %s" + params.append(limit) + + with get_conn() as conn: + with conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur: + cur.execute(sql, params) + return [dict(r) for r in cur.fetchall()] + + +def get_bands_stale(hours: int = 24 * 30, limit: int = 100) -> List[Dict]: + """Bands enrichies il y a plus de N heures (pour re-crawl périodique).""" + sql = """ + SELECT ma_id, data->>'url' AS url, name, country + FROM bands + WHERE enriched = true + AND (crawled_at IS NULL OR crawled_at < now() - make_interval(hours => %s)) + ORDER BY crawled_at ASC NULLS FIRST + LIMIT %s + """ + with get_conn() as conn: + with conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur: + cur.execute(sql, (hours, limit)) + return [dict(r) for r in cur.fetchall()] + + +# ------------------------------------------------------------------ +# Crawl run tracking +# ------------------------------------------------------------------ + +def start_crawl_run(run_type: str, countries: Optional[List[str]] = None) -> int: + sql = "INSERT INTO crawl_run (run_type, countries) VALUES (%s, %s) RETURNING id" + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute(sql, (run_type, countries)) + return cur.fetchone()[0] + + +def finish_crawl_run(run_id: int, stats: Dict[str, int], error: Optional[str] = None): + status = "error" if error else "done" + sql = """ + UPDATE crawl_run SET + status = %s, finished_at = now(), + bands_seen = %s, bands_new = %s, bands_updated = %s, bands_enriched = %s, + error = %s + WHERE id = %s + """ + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute(sql, ( + status, + stats.get("seen", 0), stats.get("new", 0), + stats.get("updated", 0), stats.get("enriched", 0), + error, run_id, + )) + + +# ------------------------------------------------------------------ +# Checkpoints +# ------------------------------------------------------------------ + +def get_checkpoint(key: str) -> Optional[str]: + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute("SELECT value FROM crawl_checkpoint WHERE key = %s", (key,)) + row = cur.fetchone() + return row[0] if row else None + + +def set_checkpoint(key: str, value: str): + sql = """ + INSERT INTO crawl_checkpoint (key, value, updated_at) + VALUES (%s, %s, now()) + ON CONFLICT (key) DO UPDATE SET value = EXCLUDED.value, updated_at = now() + """ + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute(sql, (key, value)) + + +# ------------------------------------------------------------------ +# Helpers +# ------------------------------------------------------------------ + +def _parse_year(s: Optional[str]) -> Optional[int]: + if not s: + return None + m = __import__("re").search(r"\b(1[89]\d\d|20\d\d)\b", str(s)) + return int(m.group(1)) if m else None diff --git a/apps/crawler/src/europe_codes.py b/apps/crawler/src/europe_codes.py new file mode 100644 index 0000000..73c3610 --- /dev/null +++ b/apps/crawler/src/europe_codes.py @@ -0,0 +1,8 @@ +# Codes pays Europe (Metal Archives country codes) +EUROPE_COUNTRY_CODES = [ + "AL", "AD", "AT", "BY", "BE", "BA", "BG", "HR", "CY", "CZ", + "DK", "EE", "FI", "FR", "GE", "DE", "GR", "HU", "IS", "IE", + "IT", "XK", "LV", "LI", "LT", "LU", "MT", "MD", "MC", "ME", + "NL", "MK", "NO", "PL", "PT", "RO", "RU", "SM", "RS", "SK", + "SI", "ES", "SE", "CH", "TR", "UA", "GB", +] diff --git a/apps/crawler/src/flaresolverr.py b/apps/crawler/src/flaresolverr.py new file mode 100644 index 0000000..b8c84b3 --- /dev/null +++ b/apps/crawler/src/flaresolverr.py @@ -0,0 +1,85 @@ +""" +Client FlareSolverr pour bypasser Cloudflare. + +Workflow : + 1. fs = FlareSolverr(url) + 2. fs.warmup(ma_url) → solve le challenge CF, stocke les cookies + 3. session = fs.build_session() → requests.Session avec les cookies CF + 4. Toutes les requêtes suivantes passent par session (requests standard, rapide) + 5. Si 403 : appeler fs.warmup() à nouveau pour rafraîchir les cookies +""" +import logging +import requests +from typing import Dict, List, Optional + +log = logging.getLogger(__name__) + + +class FlareSolverrError(Exception): + pass + + +class FlareSolverr: + def __init__(self, base_url: str, timeout_ms: int = 90_000): + self.base_url = base_url.rstrip("/") + self.timeout_ms = timeout_ms + self._cookies: List[Dict] = [] + self._user_agent: Optional[str] = None + + # ------------------------------------------------------------------ + # Public + # ------------------------------------------------------------------ + + def warmup(self, url: str): + """Charge une page MA pour résoudre le challenge CF et stocker les cookies.""" + log.info(f"[fs] warmup on {url}") + resp = self._request("request.get", url) + solution = resp.get("solution", {}) + self._cookies = solution.get("cookies", []) + self._user_agent = solution.get("userAgent") + log.info(f"[fs] warmup ok, got {len(self._cookies)} cookies, ua={self._user_agent[:40] if self._user_agent else None}") + + def build_session(self) -> requests.Session: + """Retourne un requests.Session avec les cookies et user-agent CF.""" + if not self._cookies: + raise FlareSolverrError("No cookies — call warmup() first") + + s = requests.Session() + s.headers.update({ + "User-Agent": self._user_agent or "Mozilla/5.0", + "Accept-Language": "en-US,en;q=0.9", + "Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8", + "Referer": "https://www.metal-archives.com/", + }) + for c in self._cookies: + s.cookies.set(c["name"], c["value"], domain=c.get("domain", ""), path=c.get("path", "/")) + return s + + def get_raw(self, url: str) -> str: + """Fait une requête via FlareSolverr et retourne le HTML rendu.""" + resp = self._request("request.get", url) + return resp.get("solution", {}).get("response", "") + + def healthy(self) -> bool: + try: + r = requests.get(f"{self.base_url}/health", timeout=5) + return r.ok + except Exception: + return False + + # ------------------------------------------------------------------ + # Internal + # ------------------------------------------------------------------ + + def _request(self, cmd: str, url: str) -> Dict: + payload = {"cmd": cmd, "url": url, "maxTimeout": self.timeout_ms} + try: + r = requests.post(f"{self.base_url}/v1", json=payload, timeout=self.timeout_ms / 1000 + 10) + r.raise_for_status() + data = r.json() + except Exception as e: + raise FlareSolverrError(f"FlareSolverr request failed: {e}") from e + + if data.get("status") != "ok": + raise FlareSolverrError(f"FlareSolverr error: {data.get('message', data)}") + return data diff --git a/apps/crawler/src/jobs.py b/apps/crawler/src/jobs.py new file mode 100644 index 0000000..b573353 --- /dev/null +++ b/apps/crawler/src/jobs.py @@ -0,0 +1,169 @@ +""" +Jobs de crawl. Chaque fonction est un job indépendant appelé par le scheduler. +""" +import logging +from datetime import datetime, timezone +from typing import List, Optional + +from .config import ( + BAND_MIN_DELAY, BAND_MAX_DELAY, + COOLDOWN_EVERY, COOLDOWN_MIN, COOLDOWN_MAX, + MA_BASE, +) +from .db import ( + upsert_bands, upsert_band_enriched, + get_bands_to_enrich, get_bands_stale, + start_crawl_run, finish_crawl_run, + get_checkpoint, set_checkpoint, +) +from .europe_codes import EUROPE_COUNTRY_CODES +from .ma_http import MASession +from .polite import sleep_range, cooldown +from .scraper_band import parse_band_page, page_hash + +log = logging.getLogger(__name__) + +CHUNK = 300 # taille des batches d'upsert + + +# ------------------------------------------------------------------ +# Crawl complet (tous les pays d'un coup) +# ------------------------------------------------------------------ + +def run_full_crawl(session: MASession, countries: List[str] = None): + countries = countries or EUROPE_COUNTRY_CODES + run_id = start_crawl_run("full_europe", countries) + stats = {"seen": 0, "new": 0, "updated": 0, "enriched": 0} + error = None + + try: + for cc in countries: + log.info(f"[full] country={cc}") + buf = [] + for band in session.fetch_country_bands(cc): + band["data"] = {"url": band.pop("url", None)} + buf.append(band) + if len(buf) >= CHUNK: + r = upsert_bands(buf) + stats["seen"] += len(buf) + stats["new"] += r["inserted"] + stats["updated"] += r["updated"] + buf = [] + if buf: + r = upsert_bands(buf) + stats["seen"] += len(buf) + stats["new"] += r["inserted"] + stats["updated"] += r["updated"] + + set_checkpoint("last_full_crawl_at", _now_iso()) + log.info(f"[full] done: {stats}") + except Exception as e: + error = str(e) + log.error(f"[full] error: {e}", exc_info=True) + finally: + finish_crawl_run(run_id, stats, error) + + +# ------------------------------------------------------------------ +# Crawl incrémental (dernières additions / modifications) +# ------------------------------------------------------------------ + +def run_incremental(session: MASession, order_by: str): + """ + order_by : 'created' | 'modified' + S'arrête dès qu'on voit des entrées déjà connues (via le checkpoint). + """ + checkpoint_key = f"last_{order_by}_check" + since = get_checkpoint(checkpoint_key) + run_id = start_crawl_run(f"incremental_{order_by}") + stats = {"seen": 0, "new": 0, "updated": 0, "enriched": 0} + error = None + latest_date = since + + log.info(f"[incr/{order_by}] since={since}") + try: + buf = [] + for band in session.fetch_archive_bands(order_by, since_date=since): + date_str = band.pop("date_str", None) + if date_str and (latest_date is None or date_str > latest_date): + latest_date = date_str + + band["data"] = {"url": band.pop("url", None)} + buf.append(band) + if len(buf) >= CHUNK: + r = upsert_bands(buf) + stats["seen"] += len(buf) + stats["new"] += r["inserted"] + stats["updated"] += r["updated"] + buf = [] + + if buf: + r = upsert_bands(buf) + stats["seen"] += len(buf) + stats["new"] += r["inserted"] + stats["updated"] += r["updated"] + + if latest_date: + set_checkpoint(checkpoint_key, latest_date) + log.info(f"[incr/{order_by}] done: {stats}") + except Exception as e: + error = str(e) + log.error(f"[incr/{order_by}] error: {e}", exc_info=True) + finally: + finish_crawl_run(run_id, stats, error) + + +# ------------------------------------------------------------------ +# Enrichissement des pages individuelles de bands +# ------------------------------------------------------------------ + +def run_enrich(session: MASession, limit: int = 100, country: Optional[str] = None): + """ + Visite les pages individuelles des bands non encore enrichies. + Stocke themes, membres, label, dates MA, hash HTML. + """ + run_id = start_crawl_run("enrich", [country] if country else None) + stats = {"seen": 0, "new": 0, "updated": 0, "enriched": 0} + error = None + + try: + bands = get_bands_to_enrich(country=country, limit=limit) + log.info(f"[enrich] {len(bands)} bands to enrich") + + for i, band in enumerate(bands): + url = band.get("url") + if not url: + continue + if not url.startswith("http"): + url = MA_BASE + url + + try: + html = session.get_html(url) + data = parse_band_page(html) + h = page_hash(html) + ok = upsert_band_enriched(band["ma_id"], data, h) + if ok: + stats["enriched"] += 1 + stats["seen"] += 1 + except Exception as e: + log.warning(f"[enrich] failed ma_id={band['ma_id']}: {e}") + continue + + sleep_range(BAND_MIN_DELAY, BAND_MAX_DELAY) + if COOLDOWN_EVERY and (i + 1) % COOLDOWN_EVERY == 0: + cooldown(COOLDOWN_MIN, COOLDOWN_MAX) + + log.info(f"[enrich] done: {stats}") + except Exception as e: + error = str(e) + log.error(f"[enrich] error: {e}", exc_info=True) + finally: + finish_crawl_run(run_id, stats, error) + + +# ------------------------------------------------------------------ +# Helpers +# ------------------------------------------------------------------ + +def _now_iso() -> str: + return datetime.now(timezone.utc).strftime("%Y-%m-%d") diff --git a/apps/crawler/src/ma_http.py b/apps/crawler/src/ma_http.py new file mode 100644 index 0000000..c8d4dc8 --- /dev/null +++ b/apps/crawler/src/ma_http.py @@ -0,0 +1,244 @@ +""" +Session HTTP vers Metal Archives. + +Combine FlareSolverr (pour les cookies CF) et requests (pour les appels AJAX rapides). +Si un appel retourne 403/429, on rafraîchit automatiquement les cookies. +""" +import logging +import time +from typing import Any, Dict, Optional + +import requests + +from .config import MA_BASE, FS_TIMEOUT_MS, AJAX_PAGE_SIZE, LIST_MIN_DELAY, LIST_MAX_DELAY +from .flaresolverr import FlareSolverr, FlareSolverrError +from .polite import sleep_range + +log = logging.getLogger(__name__) + +_WARMUP_URL = f"{MA_BASE}/lists/FR" + + +class MASession: + def __init__(self, fs: FlareSolverr): + self.fs = fs + self._session: Optional[requests.Session] = None + self._cookie_age = 0.0 + + # ------------------------------------------------------------------ + # Session management + # ------------------------------------------------------------------ + + def ensure_session(self, force_refresh: bool = False): + age = time.time() - self._cookie_age + if self._session is None or force_refresh or age > 3600 * 4: # refresh every 4h + log.info("[ma_http] acquiring Cloudflare cookies via FlareSolverr...") + self.fs.warmup(_WARMUP_URL) + self._session = self.fs.build_session() + self._cookie_age = time.time() + log.info("[ma_http] session ready") + + def get_json(self, url: str, params: Dict = None, retries: int = 2) -> Any: + """GET JSON depuis un endpoint AJAX MA.""" + self.ensure_session() + for attempt in range(retries + 1): + try: + r = self._session.get(url, params=params, timeout=30) + if r.status_code in (403, 429, 503) and attempt < retries: + log.warning(f"[ma_http] {r.status_code} on {url}, refreshing cookies...") + self.ensure_session(force_refresh=True) + sleep_range(LIST_MIN_DELAY * 2, LIST_MAX_DELAY * 2) + continue + r.raise_for_status() + return r.json() + except requests.exceptions.JSONDecodeError: + log.warning(f"[ma_http] non-JSON response from {url}") + raise + raise RuntimeError(f"Failed after {retries} retries: {url}") + + def get_html(self, url: str, retries: int = 2) -> str: + """GET HTML brut (pages band, etc.).""" + self.ensure_session() + for attempt in range(retries + 1): + try: + r = self._session.get(url, timeout=30) + if r.status_code in (403, 429, 503) and attempt < retries: + log.warning(f"[ma_http] {r.status_code} on {url}, refreshing...") + self.ensure_session(force_refresh=True) + sleep_range(LIST_MIN_DELAY * 2, LIST_MAX_DELAY * 2) + continue + r.raise_for_status() + return r.text + except Exception: + if attempt < retries: + sleep_range(3, 6) + continue + raise + raise RuntimeError(f"Failed after {retries} retries: {url}") + + # ------------------------------------------------------------------ + # Metal Archives AJAX helpers + # ------------------------------------------------------------------ + + def fetch_country_bands(self, country_code: str): + """ + Pagine l'endpoint AJAX country list MA et yield chaque band. + Endpoint : /browse/ajax-country/c/{cc}/json/1 + Réponse : {"iTotalRecords": N, "aaData": [["name", genre, location, status], ...]} + """ + url = f"{MA_BASE}/browse/ajax-country/c/{country_code}/json/1" + start = 0 + total = None + echo = 1 + + while True: + params = { + "sEcho": echo, + "iColumns": 4, + "iDisplayStart": start, + "iDisplayLength": AJAX_PAGE_SIZE, + } + log.info(f"[ma_http] country={country_code} start={start}") + data = self.get_json(url, params=params) + echo += 1 + + if total is None: + total = int(data.get("iTotalRecords") or data.get("iTotalDisplayRecords") or 0) + log.info(f"[ma_http] {country_code}: {total} total bands") + + rows = data.get("aaData") or [] + if not rows: + break + + for row in rows: + parsed = _parse_list_row(row, country_code) + if parsed: + yield parsed + + start += len(rows) + if start >= total: + break + + sleep_range(LIST_MIN_DELAY, LIST_MAX_DELAY) + + def fetch_archive_bands(self, order_by: str, since_date: Optional[str] = None): + """ + Pagine l'endpoint AJAX des archives (latest additions ou latest modified). + order_by : 'created' | 'modified' + Endpoint : /archives/ajax-band-list/by/{order_by}/json/1 + Réponse : {"iTotalRecords": N, "aaData": [["name", country, genre, date], ...]} + + Si since_date (ISO date string), s'arrête dès qu'on voit des entrées plus anciennes. + """ + url = f"{MA_BASE}/archives/ajax-band-list/by/{order_by}/json/1" + start = 0 + echo = 1 + + while True: + params = { + "sEcho": echo, + "iColumns": 4, + "iDisplayStart": start, + "iDisplayLength": AJAX_PAGE_SIZE, + } + log.info(f"[ma_http] archives/{order_by} start={start}") + data = self.get_json(url, params=params) + echo += 1 + + rows = data.get("aaData") or [] + if not rows: + break + + stop = False + for row in rows: + parsed = _parse_archive_row(row, order_by) + if not parsed: + continue + # Si on a déjà vu jusqu'à cette date, on s'arrête + if since_date and parsed.get("date_str") and parsed["date_str"] <= since_date: + stop = True + break + yield parsed + + if stop: + break + + total = int(data.get("iTotalRecords") or 0) + start += len(rows) + if start >= total: + break + + sleep_range(LIST_MIN_DELAY, LIST_MAX_DELAY) + + +# ------------------------------------------------------------------ +# Row parsers +# ------------------------------------------------------------------ + +import re +from typing import Optional as Opt + + +def _extract_link(html_str: str): + """Extrait href et texte d'un fragment HTML 'text'.""" + m = re.search(r'href=["\']([^"\']+)["\']', html_str or "") + href = m.group(1) if m else None + text = re.sub(r"<[^>]+>", "", html_str or "").strip() + return href, text + + +def _extract_ma_id(url: str) -> Opt[int]: + m = re.search(r"/bands/[^/]+/(\d+)", url or "") + return int(m.group(1)) if m else None + + +def _clean(s) -> str: + return re.sub(r"<[^>]+>", "", str(s or "")).strip() + + +def _parse_list_row(row, country_code: str) -> Opt[Dict]: + """ + Ligne du listing pays : [name_html, genre, location, status] + """ + if len(row) < 3: + return None + href, name = _extract_link(row[0]) + ma_id = _extract_ma_id(href) + if not ma_id or not name: + return None + return { + "ma_id": ma_id, + "name": name, + "url": href, + "country": country_code, + "genre": _clean(row[1]), + "location_text": _clean(row[2]), + "status": _clean(row[3]) if len(row) > 3 else None, + } + + +def _parse_archive_row(row, order_by: str) -> Opt[Dict]: + """ + Ligne des archives : [name_html, country_html, genre, date_str] + """ + if len(row) < 4: + return None + href, name = _extract_link(row[0]) + ma_id = _extract_ma_id(href) + if not ma_id or not name: + return None + + country_href, country_code = _extract_link(row[1]) + # country_href est genre /lists/FR → extraire le code + m = re.search(r"/lists/([A-Z]{1,3})", country_href or "") + cc = m.group(1) if m else _clean(row[1]) + + return { + "ma_id": ma_id, + "name": name, + "url": href, + "country": cc, + "genre": _clean(row[2]), + "date_str": _clean(row[3]), + "date_type": order_by, # 'created' ou 'modified' + } diff --git a/apps/crawler/src/main.py b/apps/crawler/src/main.py new file mode 100644 index 0000000..ba0a962 --- /dev/null +++ b/apps/crawler/src/main.py @@ -0,0 +1,106 @@ +""" +Point d'entrée du crawler. Lance le scheduler et exécute les jobs périodiques. + +Jobs : + - Toutes les SCHED_INCREMENTAL_H heures : crawl incrémental (additions + modifications) + - Toutes les SCHED_ENRICH_H heures : enrichissement des bands non enrichies + - Le SCHED_FULL_DAY du mois : crawl complet Europe + +Variables d'env : + DATABASE_URL (obligatoire) + FLARESOLVERR_URL (défaut: http://flaresolverr:8191) + CRAWLER_SCHED_INCREMENTAL_H (défaut: 4) + CRAWLER_SCHED_ENRICH_H (défaut: 1) + CRAWLER_SCHED_FULL_DAY (défaut: 1, mettre 0 pour désactiver) +""" +import logging +import os +import sys +import time +from datetime import datetime, timezone + +import schedule + +from .config import ( + FLARESOLVERR_URL, FS_TIMEOUT_MS, + SCHED_INCREMENTAL_H, SCHED_ENRICH_H, SCHED_FULL_DAY, +) +from .flaresolverr import FlareSolverr, FlareSolverrError +from .jobs import run_full_crawl, run_incremental, run_enrich +from .ma_http import MASession + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(levelname)s %(message)s", + datefmt="%Y-%m-%dT%H:%M:%S", + stream=sys.stdout, +) +log = logging.getLogger(__name__) + + +def _wait_flaresolverr(fs: FlareSolverr, max_wait: int = 120): + log.info("[main] waiting for FlareSolverr...") + for _ in range(max_wait // 5): + if fs.healthy(): + log.info("[main] FlareSolverr is up") + return + time.sleep(5) + raise RuntimeError("FlareSolverr not reachable after timeout") + + +def main(): + log.info("[main] crawler starting") + fs = FlareSolverr(FLARESOLVERR_URL, timeout_ms=FS_TIMEOUT_MS) + _wait_flaresolverr(fs) + + ma = MASession(fs) + + # ------------------------------------------------------------------ + # Helpers pour lancer les jobs avec log + # ------------------------------------------------------------------ + + def job_incremental(): + log.info("[scheduler] → incremental additions") + run_incremental(ma, "created") + log.info("[scheduler] → incremental modifications") + run_incremental(ma, "modified") + + def job_enrich(): + log.info("[scheduler] → enrich") + run_enrich(ma, limit=200) + + def job_full(): + if SCHED_FULL_DAY and datetime.now(timezone.utc).day == SCHED_FULL_DAY: + log.info("[scheduler] → full Europe crawl") + run_full_crawl(ma) + + # ------------------------------------------------------------------ + # Planification + # ------------------------------------------------------------------ + + if SCHED_INCREMENTAL_H > 0: + schedule.every(SCHED_INCREMENTAL_H).hours.do(job_incremental) + log.info(f"[main] incremental scheduled every {SCHED_INCREMENTAL_H}h") + + if SCHED_ENRICH_H > 0: + schedule.every(SCHED_ENRICH_H).hours.do(job_enrich) + log.info(f"[main] enrich scheduled every {SCHED_ENRICH_H}h") + + if SCHED_FULL_DAY > 0: + # Vérifier chaque jour à 03:00 UTC si c'est le bon jour + schedule.every().day.at("03:00").do(job_full) + log.info(f"[main] full crawl scheduled day={SCHED_FULL_DAY} of each month at 03:00 UTC") + + # Premier run au démarrage + log.info("[main] running initial jobs at startup") + job_incremental() + job_enrich() + + log.info("[main] entering scheduler loop") + while True: + schedule.run_pending() + time.sleep(60) + + +if __name__ == "__main__": + main() diff --git a/apps/crawler/src/polite.py b/apps/crawler/src/polite.py new file mode 100644 index 0000000..e1e9f74 --- /dev/null +++ b/apps/crawler/src/polite.py @@ -0,0 +1,12 @@ +import random +import time + + +def sleep_range(min_s: float, max_s: float): + time.sleep(random.uniform(min_s, max_s)) + + +def cooldown(min_s: float, max_s: float): + secs = random.uniform(min_s, max_s) + print(f"[polite] cooldown {secs:.1f}s") + time.sleep(secs) diff --git a/apps/crawler/src/scraper_band.py b/apps/crawler/src/scraper_band.py new file mode 100644 index 0000000..ee08454 --- /dev/null +++ b/apps/crawler/src/scraper_band.py @@ -0,0 +1,98 @@ +""" +Parser d'une page individuelle de band Metal Archives. +Adapté de C:\Users\nicol\Documents\crawlerbm\src\scrape_band_page.py +""" +import hashlib +import re +from typing import Any, Dict, List, Optional + +from bs4 import BeautifulSoup + + +def parse_band_page(html: str) -> Dict[str, Any]: + soup = BeautifulSoup(html, "lxml") + out: Dict[str, Any] = {} + + title = soup.find("title") + out["title"] = _text(title) + + h1 = ( + soup.find("h1", class_=re.compile(r"band_name", re.I)) + or soup.find("h1") + ) + out["name"] = _text(h1) + + # Bloc #band_stats : dt/dd pairs + stats = soup.find("div", id="band_stats") or soup.find("div", id="band_info") or soup + kv: Dict[str, str] = {} + for dl in stats.find_all("dl"): + for dt in dl.find_all("dt"): + dd = dt.find_next_sibling("dd") + k = _text(dt).rstrip(":") + v = _text(dd) + if k and v: + kv[k] = v + + out["info"] = kv + out["status"] = _pick(kv, "Status") + out["formed_in"] = _pick(kv, "Formed in", "Formed") + out["genre"] = _pick(kv, "Genre") + out["themes"] = _pick(kv, "Lyrical themes", "Lyrical Themes") + out["label"] = _pick(kv, "Current label", "Label") + out["years_active"] = _pick(kv, "Years active") + + # Dates "Added on" / "Last modified" (souvent dans #auditTrail ou en bas de page) + audit = soup.find(id="auditTrail") or soup.find("div", class_=re.compile(r"audit", re.I)) + if audit: + trail_text = _text(audit) + m_added = re.search(r"Added on:\s*([^\n,]+)", trail_text, re.I) + m_modified = re.search(r"Last modified on:\s*([^\n,]+)", trail_text, re.I) + out["ma_created_at"] = m_added.group(1).strip() if m_added else None + out["ma_modified_at"] = m_modified.group(1).strip() if m_modified else None + + # Membres + members: List[Dict] = [] + for table in soup.find_all("table"): + cls = " ".join(table.get("class") or []) + if "lineupTable" not in cls: + continue + section_el = table.find_previous(["h2", "h3"]) + section = _text(section_el) if section_el else None + headers = [_text(th) for th in table.find_all("th")] + for tr in table.find_all("tr")[1:]: + tds = tr.find_all(["td", "th"]) + if not tds: + continue + cells = [_text(td) for td in tds] + row = dict(zip(headers, cells)) if headers and len(headers) == len(cells) else {"cells": cells} + if section: + row["section"] = section + members.append(row) + if members: + out["members"] = members + + # Liens externes + links_div = soup.find("div", id="band_links") + if links_div: + out["links"] = [ + {"text": _text(a), "url": a["href"]} + for a in links_div.find_all("a", href=True) + ] + + return out + + +def page_hash(html: str) -> str: + """MD5 du HTML pour détecter les changements.""" + return hashlib.md5(html.encode()).hexdigest() + + +def _text(el) -> str: + return (el.get_text(" ", strip=True) if el else "").strip() + + +def _pick(kv: Dict, *keys: str) -> Optional[str]: + for k in keys: + if k in kv: + return kv[k] + return None diff --git a/docker-compose.dev.yml b/docker-compose.dev.yml index 31f4d27..6f4f370 100644 --- a/docker-compose.dev.yml +++ b/docker-compose.dev.yml @@ -7,6 +7,30 @@ services: networks: - bm_internal + flaresolverr: + image: ghcr.io/flaresolverr/flaresolverr:latest + environment: + LOG_LEVEL: info + CAPTCHA_SOLVER: none + networks: + - bm_internal + restart: unless-stopped + + crawler: + build: + context: apps/crawler + environment: + DATABASE_URL: ${DATABASE_URL} + FLARESOLVERR_URL: http://flaresolverr:8191 + CRAWLER_SCHED_INCREMENTAL_H: "4" + CRAWLER_SCHED_ENRICH_H: "2" + CRAWLER_SCHED_FULL_DAY: "0" + depends_on: + - flaresolverr + networks: + - bm_internal + restart: unless-stopped + geocoder-worker: build: context: apps/geocoder @@ -22,7 +46,6 @@ services: command: ["python", "src/worker.py"] networks: - bm_internal - - coolify restart: "no" pgadmin: