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: