diff --git a/apps/api/migrations/016_updated_at_only_on_change.sql b/apps/api/migrations/016_updated_at_only_on_change.sql new file mode 100644 index 0000000..bd440eb --- /dev/null +++ b/apps/api/migrations/016_updated_at_only_on_change.sql @@ -0,0 +1,48 @@ +-- 016_updated_at_only_on_change.sql +-- Corrige la boucle de ré-enrichissement infinie. +-- +-- Le trigger bands_set_geom (migration 002) faisait `NEW.updated_at := now()` +-- de façon INCONDITIONNELLE sur toute UPDATE. Or upsert_bands (crawler) fait un +-- `ON CONFLICT (ma_id) DO UPDATE SET name = COALESCE(EXCLUDED.name, bands.name), ...` +-- SANS clause WHERE : Postgres exécute donc l'UPDATE (et déclenche le trigger) +-- pour chaque ligne existante, même quand AUCUNE valeur ne change réellement. +-- +-- Conséquence : chaque crawl incrémental/complet remettait updated_at=now() sur +-- des dizaines de milliers de bands inchangées. get_bands_to_enrich considérait +-- alors `updated_at > crawled_at + 1 min` comme vrai pour quasi toute la table, +-- ré-enrichissant en boucle des pages inchangées → charge inutile et permanente +-- vers metal-archives.com. +-- +-- Fix : ne bumper updated_at que si la ligne change VRAIMENT (NEW IS DISTINCT +-- FROM OLD). Un upsert qui réécrit des valeurs identiques ne déclenche plus rien. + +CREATE OR REPLACE FUNCTION bands_set_geom() RETURNS trigger AS $$ +BEGIN + -- Ne recalculer geom que si lat ou lon a vraiment changé + IF TG_OP = 'INSERT' + OR OLD.lat IS DISTINCT FROM NEW.lat + OR OLD.lon IS DISTINCT FROM NEW.lon + THEN + IF NEW.lat IS NOT NULL AND NEW.lon IS NOT NULL THEN + NEW.geom := ST_SetSRID(ST_MakePoint(NEW.lon, NEW.lat), 4326)::geography; + ELSE + NEW.geom := NULL; + END IF; + END IF; + + -- Ne bumper updated_at que sur changement réel de la ligne. + -- À ce stade NEW.updated_at == OLD.updated_at (pas encore modifié), donc la + -- comparaison ne se compare pas elle-même ; si toutes les autres colonnes sont + -- identiques, NEW IS DISTINCT FROM OLD est faux et updated_at reste inchangé. + IF TG_OP = 'INSERT' OR NEW IS DISTINCT FROM OLD THEN + NEW.updated_at := now(); + END IF; + + RETURN NEW; +END; +$$ LANGUAGE plpgsql; + +DROP TRIGGER IF EXISTS trg_bands_set_geom ON bands; +CREATE TRIGGER trg_bands_set_geom +BEFORE INSERT OR UPDATE ON bands +FOR EACH ROW EXECUTE FUNCTION bands_set_geom(); diff --git a/apps/api/migrations/017_enrich_pending_signal.sql b/apps/api/migrations/017_enrich_pending_signal.sql new file mode 100644 index 0000000..4ffa394 --- /dev/null +++ b/apps/api/migrations/017_enrich_pending_signal.sql @@ -0,0 +1,22 @@ +-- 017_enrich_pending_signal.sql +-- Découple « a besoin d'être ré-enrichi » de updated_at. +-- +-- Après la migration 016, updated_at ne bouge plus que sur changement réel d'un +-- champ de listing (nom, pays, genre, statut, lieu). Problème : quand Metal +-- Archives modifie un groupe sur un champ visible UNIQUEMENT sur la page du +-- groupe (line-up, albums, thèmes), la liste archives/modified ne montre aucun +-- diff au niveau listing → sans signal dédié, ce groupe ne serait ré-enrichi +-- qu'au bout de 30 jours (filet stale). +-- +-- Solution : le crawler incrémental « modified » capture l'horodatage "modified" +-- que MA affiche lui-même dans sa liste (ma_modified_seen). Quand ce texte change +-- pour un groupe, on lève enrich_pending — une seule fois par modification MA, +-- sans boucle. L'enrichissement remet enrich_pending à false. + +ALTER TABLE bands + ADD COLUMN IF NOT EXISTS ma_modified_seen TEXT, + ADD COLUMN IF NOT EXISTS enrich_pending BOOLEAN NOT NULL DEFAULT false; + +-- Index partiel pour la file d'enrichissement (peu de lignes à true à la fois) +CREATE INDEX IF NOT EXISTS idx_bands_enrich_pending + ON bands (enrich_pending) WHERE enrich_pending; diff --git a/apps/crawler/src/db.py b/apps/crawler/src/db.py index 5b5940b..2f5c3ac 100644 --- a/apps/crawler/src/db.py +++ b/apps/crawler/src/db.py @@ -3,6 +3,7 @@ """ import json import logging +import re import time from contextlib import contextmanager from datetime import UTC, datetime @@ -15,6 +16,13 @@ from .config import DATABASE_URL log = logging.getLogger(__name__) +# Types de job que CE process sait exécuter (voir main._check_job_triggers). +# 'geocoder_enqueue' en est volontairement absent : il est consommé par +# apps/geocoder/src/enqueue.py. Sans ce filtre, le crawler raflait ces triggers +# et les refermait en « job_type inconnu », selon lequel des deux daemons +# interrogeait la table en premier. +CRAWLER_JOB_TYPES = ["enrich", "incremental", "full_crawl"] + @contextmanager def get_conn(): @@ -62,16 +70,22 @@ def upsert_bands(bands: list[dict[str, Any]]) -> dict[str, int]: b.get("crawled_hash"), b.get("ma_created_at"), b.get("ma_modified_at"), + b.get("ma_modified_seen"), ts, # first_seen_at (ignoré si déjà set) ts, # created_at (ignoré si déjà set) psycopg2.extras.Json(b.get("data") or {}), )) + # enrich_pending : levé quand l'horodatage "modified" vu dans la liste MA + # change réellement (nouvelle modif côté MA). COALESCE ne peut pas servir + # ici, il faut comparer l'ancienne et la nouvelle valeur → CASE explicite. + # Toutes les expressions du SET voient la ligne AVANT update, donc l'ordre + # entre ma_modified_seen et enrich_pending n'a pas d'importance. 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, data) + ma_modified_seen, first_seen_at, created_at, data) VALUES %s ON CONFLICT (ma_id) DO UPDATE SET name = COALESCE(EXCLUDED.name, bands.name), @@ -86,6 +100,13 @@ def upsert_bands(bands: list[dict[str, Any]]) -> dict[str, int]: 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), + enrich_pending = CASE + WHEN EXCLUDED.ma_modified_seen IS NOT NULL + AND EXCLUDED.ma_modified_seen IS DISTINCT FROM bands.ma_modified_seen + THEN true + ELSE bands.enrich_pending + END, + ma_modified_seen = COALESCE(EXCLUDED.ma_modified_seen, bands.ma_modified_seen), data = bands.data || EXCLUDED.data RETURNING (xmax = 0) AS was_inserted """ @@ -147,6 +168,8 @@ def upsert_band_enriched(ma_id: int, data: dict[str, Any], html_hash: str) -> bo "ma_created_at = COALESCE(%s, ma_created_at)", "ma_modified_at = COALESCE(%s, ma_modified_at)", "crawler_pending = %s::jsonb", + # On vient d'enrichir : le signal de ré-enrichissement est consommé. + "enrich_pending = false", ] params += [ band_data, ts, html_hash, @@ -162,9 +185,11 @@ def get_bands_to_enrich(country: str | None = None, limit: int = 50) -> list[dic """ File de priorité pour l'enrichissement : 1. Nouveaux bands (band_page absent) — jamais enrichis - 2. Modifiés depuis le dernier enrichissement (updated_at > crawled_at + 1 min) - 3. Héritage ancien scraper (crawled_at IS NULL, band_page présent) - 4. Stale (enrichis il y a > 30 jours par le système actuel) + 2. enrich_pending : MA a modifié le groupe (horodatage "modified" changé), + y compris sur des champs invisibles au niveau listing + 3. Modifiés depuis le dernier enrichissement (updated_at > crawled_at + 1 min) + 4. Héritage ancien scraper (crawled_at IS NULL, band_page présent) + 5. Stale (enrichis il y a > 30 jours par le système actuel) """ sql = """ SELECT ma_id, data->>'url' AS url, name, country @@ -173,6 +198,7 @@ def get_bands_to_enrich(country: str | None = None, limit: int = 50) -> list[dic AND ( data->'band_page' IS NULL OR crawled_at IS NULL + OR enrich_pending OR (crawled_at IS NOT NULL AND updated_at > crawled_at + interval '1 minute') OR crawled_at < now() - interval '30 days' ) @@ -185,9 +211,10 @@ def get_bands_to_enrich(country: str | None = None, limit: int = 50) -> list[dic ORDER BY CASE WHEN data->'band_page' IS NULL THEN 1 - WHEN crawled_at IS NOT NULL AND updated_at > crawled_at + interval '1 minute' THEN 2 - WHEN crawled_at IS NULL THEN 3 - ELSE 4 + WHEN enrich_pending THEN 2 + WHEN crawled_at IS NOT NULL AND updated_at > crawled_at + interval '1 minute' THEN 3 + WHEN crawled_at IS NULL THEN 4 + ELSE 5 END ASC, updated_at DESC NULLS LAST LIMIT %s @@ -293,6 +320,44 @@ def raise_if_cancelled(run_id: int): raise RunCancelled(f"run #{run_id} annulé par un administrateur") +def recover_stuck_runs(): + """Au démarrage : referme tout crawl_run / job_trigger resté en 'running'. + + Le crawler tourne en instance unique : une ligne encore 'running' ne peut + provenir que d'une instance précédente tuée brutalement (OOM, SIGKILL) avant + d'atteindre son bloc finally. On les marque en erreur pour ne pas polluer le + suivi admin et libérer la file de jobs. + + Ne touche QUE les job_triggers consommés par le crawler : ceux du geocoder + (geocoder_enqueue) appartiennent à un autre process, qui fait sa propre + récupération au démarrage. + """ + try: + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute( + """UPDATE crawl_run + SET status='error', finished_at=now(), + error='interrompu (redémarrage du crawler)' + WHERE status='running' + AND run_type <> 'geocoder_enqueue'""" + ) + n_runs = cur.rowcount + cur.execute( + """UPDATE job_triggers + SET status='error', finished_at=now(), + error='interrompu (redémarrage du crawler)' + WHERE status='running' + AND job_type = ANY(%s)""", + (CRAWLER_JOB_TYPES,), + ) + n_jobs = cur.rowcount + if n_runs or n_jobs: + log.warning(f"[db] recover_stuck_runs: {n_runs} run(s) + {n_jobs} job(s) refermés") + except Exception as e: + log.warning(f"[db] recover_stuck_runs failed: {e}") + + def finish_crawl_run( run_id: int, stats: dict[str, int], @@ -346,12 +411,19 @@ def set_checkpoint(key: str, value: str): # ------------------------------------------------------------------ def claim_job_trigger(job_type: str | None = None) -> int | None: - """Récupère et verrouille un job trigger en attente. Retourne son id ou None.""" + """Récupère et verrouille un job trigger en attente. Retourne son id ou None. + + Sans `job_type` explicite, la recherche est bornée aux CRAWLER_JOB_TYPES : + réclamer un trigger qu'on ne sait pas exécuter revient à le détruire. + """ where = "status = 'pending'" params: list = [] if job_type: where += " AND job_type = %s" params.append(job_type) + else: + where += " AND job_type = ANY(%s)" + params.append(CRAWLER_JOB_TYPES) params.append(1) try: with get_conn() as conn: @@ -391,5 +463,5 @@ def finish_job_trigger(trigger_id: int, error: str | None = None): def _parse_year(s: str | None) -> int | None: if not s: return None - m = __import__("re").search(r"\b(1[89]\d\d|20\d\d)\b", str(s)) + m = 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/jobs.py b/apps/crawler/src/jobs.py index e8da363..44f69cd 100644 --- a/apps/crawler/src/jobs.py +++ b/apps/crawler/src/jobs.py @@ -10,6 +10,8 @@ from .config import ( COOLDOWN_EVERY, COOLDOWN_MAX, COOLDOWN_MIN, + LIST_MAX_DELAY, + LIST_MIN_DELAY, MA_BASE, ) from .db import ( @@ -72,6 +74,12 @@ def run_full_crawl(session: MASession, countries: list[str] = None): stats["new"] += r["inserted"] stats["updated"] += r["updated"] + update_crawl_run_progress(run_id, stats) + # Politesse entre pays : la dernière page d'un pays ne dort pas (le + # générateur sort sur un break), d'où cette pause explicite avant + # d'attaquer le suivant. + sleep_range(LIST_MIN_DELAY, LIST_MAX_DELAY) + set_checkpoint("last_full_crawl_at", _now_iso()) log.info(f"[full] done: {stats}") log_event("info", f"full crawl done: {stats}", run_id=run_id) @@ -178,6 +186,7 @@ def run_enrich(session: MASession, limit: int = 100, country: str | None = None) if not url.startswith("http"): url = MA_BASE + url + failed = False try: html = session.get_html(url) data = parse_band_page(html) @@ -187,11 +196,18 @@ def run_enrich(session: MASession, limit: int = 100, country: str | None = None) stats["enriched"] += 1 stats["seen"] += 1 except Exception as e: + failed = True log.warning(f"[enrich] failed ma_id={band['ma_id']}: {e}") log_event("warning", f"enrich failed: {e}", run_id=run_id, ma_id=band["ma_id"]) - continue - sleep_range(BAND_MIN_DELAY, BAND_MAX_DELAY) + # La politesse s'applique TOUJOURS, y compris après un échec. Le + # `continue` d'avant la sautait : or un échec de get_html est le plus + # souvent un 403/429/503 de Metal Archives, soit le pire moment pour + # enchaîner sans délai. On ralentit donc davantage dans ce cas. + if failed: + sleep_range(BAND_MAX_DELAY, BAND_MAX_DELAY * 3) + else: + sleep_range(BAND_MIN_DELAY, BAND_MAX_DELAY) if COOLDOWN_EVERY and (i + 1) % COOLDOWN_EVERY == 0: cooldown(COOLDOWN_MIN, COOLDOWN_MAX) diff --git a/apps/crawler/src/ma_http.py b/apps/crawler/src/ma_http.py index d603fdc..560ac8c 100644 --- a/apps/crawler/src/ma_http.py +++ b/apps/crawler/src/ma_http.py @@ -15,7 +15,7 @@ from typing import Any from urllib.parse import urlencode from .config import AJAX_PAGE_SIZE, LIST_MAX_DELAY, LIST_MIN_DELAY, MA_BASE -from .flaresolverr import FlareSolverr +from .flaresolverr import FlareSolverr, FlareSolverrError from .polite import sleep_range log = logging.getLogger(__name__) @@ -52,13 +52,29 @@ class MASession: 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) + try: + status, body = self.fs.get(full_url, self._session_id) + except FlareSolverrError as e: + # Blip réseau / crash du Chrome headless : on retente avec une + # session neuve plutôt que de laisser l'erreur remonter et perdre + # tout le run en cours. + if attempt < retries: + log.warning(f"[ma_http] FlareSolverr KO sur {url} ({e}), refresh session…") + self.ensure_session(force_refresh=True) + sleep_range(LIST_MIN_DELAY * 2, LIST_MAX_DELAY * 2) + continue + raise RuntimeError( + f"FlareSolverr échec après {retries} tentatives sur {url}: {e}" + ) from e 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: + # Volontairement PAS de retry ici : un 500/502 n'est pas un + # blocage anti-bot, et réessayer ne ferait que marteler un + # serveur déjà en difficulté. Verrouillé par un test. raise RuntimeError(f"HTTP {status} on {full_url}") try: return _extract_json(body) @@ -77,7 +93,20 @@ class MASession: """GET HTML rendu (pages band, etc.).""" self.ensure_session() for attempt in range(retries + 1): - status, body = self.fs.get(url, self._session_id) + try: + status, body = self.fs.get(url, self._session_id) + except FlareSolverrError as e: + # Blip réseau / crash du Chrome headless : on retente avec une + # session neuve plutôt que de laisser l'erreur remonter et perdre + # tout le run en cours. + if attempt < retries: + log.warning(f"[ma_http] FlareSolverr KO sur {url} ({e}), refresh session…") + self.ensure_session(force_refresh=True) + sleep_range(LIST_MIN_DELAY * 2, LIST_MAX_DELAY * 2) + continue + raise RuntimeError( + f"FlareSolverr échec après {retries} tentatives sur {url}: {e}" + ) from e 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) @@ -265,7 +294,25 @@ def _parse_archive_row(row, order_by: str) -> dict | None: country_href, _ = _extract_link(row[2]) m = re.search(r"/lists/([A-Z]{1,3})", country_href or "") - cc = m.group(1) if m else _clean(row[2]) + if m: + cc = m.group(1) + else: + # Lien pays absent ou malformé : on retombe sur le texte brut (ex. + # "Germany"), qui ne matchera pas _EU_SET côté jobs.py et sera écarté. + # Le log rend visible un éventuel changement de structure HTML chez MA. + cc = _clean(row[2]) + log.warning( + f"[ma_http] code pays introuvable dans '{country_href}', " + f"fallback texte='{cc}' (ma_id={ma_id})" + ) + + # Horodatage "modified" que MA affiche lui-même dans sa liste (ex. "Jun 1st, + # 02:04"). Sert de signal de ré-enrichissement : quand ce texte change pour + # un groupe, on sait que MA l'a modifié même si aucun champ du listing ne + # bouge. Sans objet pour la liste des créations, où il est immuable. + ma_modified_seen = None + if order_by == "modified": + ma_modified_seen = _clean(row[4]) if len(row) > 4 else _clean(row[0]) return { "ma_id": ma_id, @@ -273,6 +320,7 @@ def _parse_archive_row(row, order_by: str) -> dict | None: "url": href, "country": cc, "genre": _clean(row[3]), + "ma_modified_seen": ma_modified_seen, "date_str": None, # pas d'année dans le format MA → pas de filtre date "date_type": order_by, } diff --git a/apps/crawler/src/main.py b/apps/crawler/src/main.py index a3ee113..0617e9c 100644 --- a/apps/crawler/src/main.py +++ b/apps/crawler/src/main.py @@ -30,7 +30,7 @@ from .config import ( SCHED_ENRICH_H, SCHED_INCREMENTAL_H, ) -from .db import claim_job_trigger, finish_job_trigger, get_checkpoint +from .db import claim_job_trigger, finish_job_trigger, get_checkpoint, recover_stuck_runs from .flaresolverr import FlareSolverr from .health import database_probe, freshness_probe, http_probe, maybe_run from .jobs import run_enrich, run_full_crawl, run_incremental @@ -55,8 +55,26 @@ def _wait_flaresolverr(fs: FlareSolverr, max_wait: int = 120): raise RuntimeError("FlareSolverr not reachable after timeout") +def _safe(label: str, fn): + """Exécute un job en isolant toute exception. + + Un job qui plante ne doit jamais tuer le process : `restart: unless-stopped` + le relancerait en boucle, et une panne DB transitoire suffirait à mettre le + crawler en crash-loop au lieu de le faire attendre le cycle suivant. + """ + try: + fn() + except Exception as e: + log.error(f"[main] job '{label}' a échoué (ignoré, on continue): {e}", exc_info=True) + + def main(): log.info("[main] crawler starting") + + # Runs/jobs laissés en 'running' par une instance précédente tuée brutalement + # (OOM, SIGKILL) : ils ne peuvent pas appartenir à ce process, on les referme. + _safe("recover_stuck_runs", recover_stuck_runs) + fs = FlareSolverr(FLARESOLVERR_URL, timeout_ms=FS_TIMEOUT_MS) _wait_flaresolverr(fs) @@ -110,15 +128,15 @@ def main(): # Premier run au démarrage (job_full rattrape un redémarrage qui aurait loupé la fenêtre) log.info("[main] running initial jobs at startup") - job_incremental() - job_enrich() - job_full() + _safe("incremental", job_incremental) + _safe("enrich", job_enrich) + _safe("full", job_full) log.info("[main] entering scheduler loop") while True: - schedule.run_pending() - _check_job_triggers(ma) - _check_health() + _safe("run_pending", schedule.run_pending) + _safe("check_job_triggers", lambda: _check_job_triggers(ma)) + _safe("check_health", _check_health) time.sleep(60) diff --git a/apps/crawler/src/polite.py b/apps/crawler/src/polite.py index e1e9f74..6c9ad3b 100644 --- a/apps/crawler/src/polite.py +++ b/apps/crawler/src/polite.py @@ -1,6 +1,9 @@ +import logging import random import time +log = logging.getLogger(__name__) + def sleep_range(min_s: float, max_s: float): time.sleep(random.uniform(min_s, max_s)) @@ -8,5 +11,5 @@ def sleep_range(min_s: float, max_s: float): def cooldown(min_s: float, max_s: float): secs = random.uniform(min_s, max_s) - print(f"[polite] cooldown {secs:.1f}s") + log.info(f"[polite] cooldown {secs:.1f}s") time.sleep(secs)