feat(crawler): nouveau service crawler server-side via FlareSolverr
- 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 <noreply@anthropic.com>
This commit is contained in:
parent
d2fcafc6be
commit
260200e251
13 changed files with 1021 additions and 1 deletions
9
apps/crawler/Dockerfile
Normal file
9
apps/crawler/Dockerfile
Normal file
|
|
@ -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"]
|
||||
5
apps/crawler/requirements.txt
Normal file
5
apps/crawler/requirements.txt
Normal file
|
|
@ -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
|
||||
0
apps/crawler/src/__init__.py
Normal file
0
apps/crawler/src/__init__.py
Normal file
28
apps/crawler/src/config.py
Normal file
28
apps/crawler/src/config.py
Normal file
|
|
@ -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"
|
||||
233
apps/crawler/src/db.py
Normal file
233
apps/crawler/src/db.py
Normal file
|
|
@ -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
|
||||
8
apps/crawler/src/europe_codes.py
Normal file
8
apps/crawler/src/europe_codes.py
Normal file
|
|
@ -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",
|
||||
]
|
||||
85
apps/crawler/src/flaresolverr.py
Normal file
85
apps/crawler/src/flaresolverr.py
Normal file
|
|
@ -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
|
||||
169
apps/crawler/src/jobs.py
Normal file
169
apps/crawler/src/jobs.py
Normal file
|
|
@ -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")
|
||||
244
apps/crawler/src/ma_http.py
Normal file
244
apps/crawler/src/ma_http.py
Normal file
|
|
@ -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": [["<a href='url'>name</a>", 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": [["<a>name</a>", 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 '<a href="...">text</a>'."""
|
||||
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'
|
||||
}
|
||||
106
apps/crawler/src/main.py
Normal file
106
apps/crawler/src/main.py
Normal file
|
|
@ -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()
|
||||
12
apps/crawler/src/polite.py
Normal file
12
apps/crawler/src/polite.py
Normal file
|
|
@ -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)
|
||||
98
apps/crawler/src/scraper_band.py
Normal file
98
apps/crawler/src/scraper_band.py
Normal file
|
|
@ -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
|
||||
|
|
@ -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:
|
||||
|
|
|
|||
Loading…
Reference in a new issue