diff --git a/apps/api/migrations/005_fix_checkpoints_and_data.sql b/apps/api/migrations/005_fix_checkpoints_and_data.sql new file mode 100644 index 0000000..4c38fd5 --- /dev/null +++ b/apps/api/migrations/005_fix_checkpoints_and_data.sql @@ -0,0 +1,47 @@ +-- 005: corriger les checkpoints + backfill enriched depuis les données de l'ancien scraper + +-- 1. Renommer la clé checkpoint morte 'last_additions_check' → 'last_created_check' +-- (le code utilise last_{order_by}_check donc 'last_created_check', pas 'last_additions_check') +INSERT INTO crawl_checkpoint (key, value) + VALUES ('last_created_check', NULL) + ON CONFLICT (key) DO NOTHING; + +DELETE FROM crawl_checkpoint WHERE key = 'last_additions_check'; + +-- 2. Backfill depuis les données de l'ancien scraper : +-- 85 829 bands ont data->'band_page' rempli mais enriched=false et colonnes top-level vides. +-- On remplit les colonnes à partir du JSONB existant. +UPDATE bands SET + enriched = true, + status = COALESCE( + NULLIF(trim(status), ''), + data->'band_page'->>'status' + ), + genre = COALESCE( + NULLIF(trim(genre), ''), + data->'band_page'->>'genre', + data->'band_page'->'info_raw'->>'genre' + ), + themes = COALESCE( + NULLIF(trim(themes), ''), + NULLIF(data->'band_page'->>'lyrical_themes', 'N/A'), + NULLIF(data->'band_page'->'info_raw'->>'themes', 'N/A') + ), + formed_year = COALESCE( + formed_year, + CASE + WHEN (data->'band_page'->>'formed_in') ~ '^\d{4}$' + THEN (data->'band_page'->>'formed_in')::int + WHEN (data->'band_page'->>'formed_in') IS NOT NULL + THEN (regexp_match(data->'band_page'->>'formed_in', '\d{4}'))[1]::int + WHEN (data->'band_page'->'info_raw'->>'formed in') IS NOT NULL + THEN (regexp_match(data->'band_page'->'info_raw'->>'formed in', '\d{4}'))[1]::int + ELSE NULL + END + ), + location_text = COALESCE( + NULLIF(trim(location_text), ''), + data->'band_page'->>'location', + data->'band_page'->'info_raw'->>'location' + ) +WHERE data->'band_page' IS NOT NULL; diff --git a/apps/crawler/src/db.py b/apps/crawler/src/db.py index 2ccba3e..8036552 100644 --- a/apps/crawler/src/db.py +++ b/apps/crawler/src/db.py @@ -1,6 +1,7 @@ """ Écriture directe en DB. Toutes les opérations passent par des upserts idempotents. """ +import json import logging from contextlib import contextmanager from datetime import datetime, timezone @@ -62,13 +63,14 @@ def upsert_bands(bands: List[Dict[str, Any]]) -> Dict[str, int]: b.get("ma_modified_at"), ts, # first_seen_at (ignoré si déjà set) ts, # created_at (ignoré si déjà set) + psycopg2.extras.Json(b.get("data") or {}), )) 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) + first_seen_at, created_at, data) VALUES %s ON CONFLICT (ma_id) DO UPDATE SET name = COALESCE(EXCLUDED.name, bands.name), @@ -82,7 +84,8 @@ def upsert_bands(bands: List[Dict[str, Any]]) -> Dict[str, int]: 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) + ma_modified_at= COALESCE(EXCLUDED.ma_modified_at, bands.ma_modified_at), + data = bands.data || EXCLUDED.data RETURNING (xmax = 0) AS was_inserted """ diff --git a/apps/crawler/src/flaresolverr.py b/apps/crawler/src/flaresolverr.py index b8c84b3..e11829c 100644 --- a/apps/crawler/src/flaresolverr.py +++ b/apps/crawler/src/flaresolverr.py @@ -1,16 +1,17 @@ """ -Client FlareSolverr pour bypasser Cloudflare. +Client FlareSolverr — interface stateless. 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 + 2. session_id = fs.create_session() + 3. fs.get(warmup_url, session_id) → résout le challenge CF + 4. status, body = fs.get(ajax_url, session_id) → réutilise le même Chrome + 5. fs.destroy_session(session_id) → libère la ressource Chrome """ import logging +from typing import Optional, Tuple + import requests -from typing import Dict, List, Optional log = logging.getLogger(__name__) @@ -23,42 +24,6 @@ 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: @@ -67,19 +32,47 @@ class FlareSolverr: except Exception: return False - # ------------------------------------------------------------------ - # Internal - # ------------------------------------------------------------------ + def create_session(self) -> str: + """Ouvre une session Chrome persistante dans FlareSolverr, retourne le session_id.""" + data = self._call({"cmd": "sessions.create"}) + session_id = data["session"] + log.info(f"[fs] created session {session_id}") + return session_id - def _request(self, cmd: str, url: str) -> Dict: - payload = {"cmd": cmd, "url": url, "maxTimeout": self.timeout_ms} + def destroy_session(self, session_id: str): + """Libère la session Chrome.""" try: - r = requests.post(f"{self.base_url}/v1", json=payload, timeout=self.timeout_ms / 1000 + 10) + self._call({"cmd": "sessions.destroy", "session": session_id}) + log.info(f"[fs] destroyed session {session_id}") + except Exception as e: + log.warning(f"[fs] destroy session {session_id} failed: {e}") + + def get(self, url: str, session_id: Optional[str] = None) -> Tuple[int, str]: + """ + GET via FlareSolverr. Retourne (http_status_code, response_body). + Si session_id fourni, réutilise le Chrome existant (pas de re-warmup). + """ + payload: dict = {"cmd": "request.get", "url": url, "maxTimeout": self.timeout_ms} + if session_id: + payload["session"] = session_id + data = self._call(payload) + sol = data.get("solution", {}) + status = sol.get("status", 0) + body = sol.get("response", "") + log.debug(f"[fs] GET {url} → {status}") + return status, body + + def _call(self, payload: dict) -> dict: + 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/ma_http.py b/apps/crawler/src/ma_http.py index c8d4dc8..1bee730 100644 --- a/apps/crawler/src/ma_http.py +++ b/apps/crawler/src/ma_http.py @@ -1,79 +1,88 @@ """ -Session HTTP vers Metal Archives. +Session HTTP vers Metal Archives via FlareSolverr persistent sessions. -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. +Cloudflare bloque les requêtes `requests` directes (fingerprint TLS/JA3). +Toutes les requêtes passent par FlareSolverr qui utilise un Chrome headless réel. +Une session Chrome est créée une fois et réutilisée pour toute la durée du crawl +(refresh automatique toutes les 4h ou sur 403/429). """ +import json import logging +import re import time +from html import unescape as html_unescape from typing import Any, Dict, Optional +from urllib.parse import urlencode -import requests - -from .config import MA_BASE, FS_TIMEOUT_MS, AJAX_PAGE_SIZE, LIST_MIN_DELAY, LIST_MAX_DELAY +from .config import MA_BASE, 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" +_SESSION_TTL = 3600 * 4 # refresh session toutes les 4h class MASession: def __init__(self, fs: FlareSolverr): self.fs = fs - self._session: Optional[requests.Session] = None - self._cookie_age = 0.0 + self._session_id: Optional[str] = None + self._session_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") + age = time.time() - self._session_age + if self._session_id is None or force_refresh or age > _SESSION_TTL: + if self._session_id: + self.fs.destroy_session(self._session_id) + self._session_id = None + log.info("[ma_http] creating FlareSolverr persistent session...") + self._session_id = self.fs.create_session() + log.info(f"[ma_http] warming up CF cookies on {_WARMUP_URL}") + status, _ = self.fs.get(_WARMUP_URL, self._session_id) + log.info(f"[ma_http] warmup done (status={status}), session ready") + self._session_age = time.time() def get_json(self, url: str, params: Dict = None, retries: int = 2) -> Any: """GET JSON depuis un endpoint AJAX MA.""" self.ensure_session() + full_url = _build_url(url, params) for attempt in range(retries + 1): + status, body = self.fs.get(full_url, self._session_id) + if status in (403, 429, 503) and attempt < retries: + log.warning(f"[ma_http] {status} on {url}, refreshing session...") + self.ensure_session(force_refresh=True) + sleep_range(LIST_MIN_DELAY * 2, LIST_MAX_DELAY * 2) + continue + if status != 200: + raise RuntimeError(f"HTTP {status} on {full_url}") 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: + return _extract_json(body) + except (json.JSONDecodeError, ValueError): 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.).""" + """GET HTML rendu (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: + status, body = self.fs.get(url, self._session_id) + if status in (403, 429, 503) and attempt < retries: + log.warning(f"[ma_http] {status} on {url}, refreshing session...") + self.ensure_session(force_refresh=True) + sleep_range(LIST_MIN_DELAY * 2, LIST_MAX_DELAY * 2) + continue + if status != 200: if attempt < retries: sleep_range(3, 6) continue - raise + raise RuntimeError(f"HTTP {status} on {url}") + return body raise RuntimeError(f"Failed after {retries} retries: {url}") # ------------------------------------------------------------------ @@ -154,7 +163,6 @@ class MASession: 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 @@ -172,13 +180,29 @@ class MASession: # ------------------------------------------------------------------ -# Row parsers +# Helpers internes # ------------------------------------------------------------------ -import re -from typing import Optional as Opt +def _build_url(base: str, params: Dict = None) -> str: + if not params: + return base + return f"{base}?{urlencode(params)}" +def _extract_json(body: str) -> Any: + """ + FlareSolverr wraps les réponses JSON dans
avec HTML-encoding des <>. + On extrait le contenu, on unescape les entités HTML, puis on parse le JSON. + """ + m = re.search(r"]*>([\s\S]*?)", body) + text = html_unescape(m.group(1)) if m else body + return json.loads(text) + + +# ------------------------------------------------------------------ +# Row parsers +# ------------------------------------------------------------------ + def _extract_link(html_str: str): """Extrait href et texte d'un fragment HTML 'text'.""" m = re.search(r'href=["\']([^"\']+)["\']', html_str or "") @@ -187,7 +211,7 @@ def _extract_link(html_str: str): return href, text -def _extract_ma_id(url: str) -> Opt[int]: +def _extract_ma_id(url: str) -> Optional[int]: m = re.search(r"/bands/[^/]+/(\d+)", url or "") return int(m.group(1)) if m else None @@ -196,10 +220,8 @@ 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] - """ +def _parse_list_row(row, country_code: str) -> Optional[Dict]: + """Ligne du listing pays : [name_html, genre, location, status]""" if len(row) < 3: return None href, name = _extract_link(row[0]) @@ -217,10 +239,8 @@ def _parse_list_row(row, country_code: str) -> Opt[Dict]: } -def _parse_archive_row(row, order_by: str) -> Opt[Dict]: - """ - Ligne des archives : [name_html, country_html, genre, date_str] - """ +def _parse_archive_row(row, order_by: str) -> Optional[Dict]: + """Ligne des archives : [name_html, country_html, genre, date_str]""" if len(row) < 4: return None href, name = _extract_link(row[0]) @@ -228,8 +248,7 @@ def _parse_archive_row(row, order_by: str) -> Opt[Dict]: 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 + country_href, _ = _extract_link(row[1]) m = re.search(r"/lists/([A-Z]{1,3})", country_href or "") cc = m.group(1) if m else _clean(row[1]) @@ -240,5 +259,5 @@ def _parse_archive_row(row, order_by: str) -> Opt[Dict]: "country": cc, "genre": _clean(row[2]), "date_str": _clean(row[3]), - "date_type": order_by, # 'created' ou 'modified' + "date_type": order_by, } diff --git a/apps/crawler/src/scraper_band.py b/apps/crawler/src/scraper_band.py index b515c36..70a1b46 100644 --- a/apps/crawler/src/scraper_band.py +++ b/apps/crawler/src/scraper_band.py @@ -1,5 +1,8 @@ """ Parser d'une page individuelle de band Metal Archives. + +Produit un dict compatible avec upsert_band_enriched() ET cohérent avec +le schéma band_page historique stocké en DB (lineup structuré, discography_url, etc.) """ import hashlib import re @@ -7,85 +10,96 @@ from typing import Any, Dict, List, Optional from bs4 import BeautifulSoup +# Section header → clé lineup +_LINEUP_SECTIONS = { + "current": "current", + "past": "past", + "live": "live", + "session": "session", +} + 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") - ) + 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 + # --- band_stats : dt/dd pairs --- + stats_div = 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 dl in stats_div.find_all("dl"): for dt in dl.find_all("dt"): dd = dt.find_next_sibling("dd") k = _text(dt).rstrip(":") - v = _text(dd) + v = _text(dd) if dd else None if k and v: kv[k] = v + out["info_raw"] = {k.lower(): v for k, v in kv.items()} - 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["formed_in"] = _pick(kv, "Formed in", "Formed") out["years_active"] = _pick(kv, "Years active") + out["location"] = _pick(kv, "Location") + # MA uses either "Lyrical themes", "Lyrical Themes", or simply "Themes" + out["lyrical_themes"] = _pick(kv, "Lyrical themes", "Lyrical Themes", "Themes") + # Alias for upsert_band_enriched which reads data.get("themes") + out["themes"] = out["lyrical_themes"] + out["label"] = _pick(kv, "Current label", "Last label", "Label") - # 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")] + # --- Lineup (grouped by section, with artist URLs) --- + lineup: Dict[str, List[Dict]] = {"current": [], "past": [], "live": [], "session": []} + for table in soup.find_all("table", class_="lineupTable"): + section_key = _lineup_section(table) + bucket = lineup.setdefault(section_key, []) 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 + a = tds[0].find("a", href=re.compile(r"/artists/")) + if not a: + continue + role = _text(tds[1]) if len(tds) > 1 else None + bucket.append({ + "name": _text(a), + "url": a.get("href", ""), + "role": role, + }) + out["lineup"] = lineup - # Liens externes + # --- Discography URL --- + disc = soup.find("a", href=re.compile(r"/band/discography/")) + out["discography_url"] = disc["href"] if disc else None + + # --- 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) - ] + out["links"] = [ + {"text": _text(a), "url": a["href"]} + for a in links_div.find_all("a", href=True) + ] if links_div else [] + + # --- Dates audit trail --- + audit = soup.find(id="auditTrail") or soup.find("div", class_=re.compile(r"audit", re.I)) + if audit: + trail = _text(audit) + m_added = re.search(r"Added on:\s*(\S+\s+\S+)", trail, re.I) + m_modified = re.search(r"Last modified on:\s*(\S+\s+\S+)", trail, 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 return out def page_hash(html: str) -> str: - """MD5 du HTML pour détecter les changements.""" return hashlib.md5(html.encode()).hexdigest() +# ------------------------------------------------------------------ +# Helpers +# ------------------------------------------------------------------ + def _text(el) -> str: return (el.get_text(" ", strip=True) if el else "").strip() @@ -95,3 +109,15 @@ def _pick(kv: Dict, *keys: str) -> Optional[str]: if k in kv: return kv[k] return None + + +def _lineup_section(table) -> str: + """Déduit la section lineup depuis le h2/h3 qui précède la table.""" + el = table.find_previous(["h2", "h3"]) + if not el: + return "current" + text = _text(el).lower() + for kw, key in _LINEUP_SECTIONS.items(): + if kw in text: + return key + return "current"