fix(crawler): boucle de ré-enrichissement infinie, politesse et résilience

Le trigger bands_set_geom (migration 002) faisait `NEW.updated_at := now()`
sans condition. Or upsert_bands exécute un `ON CONFLICT DO UPDATE` sans clause
WHERE : Postgres déclenche donc le trigger pour chaque ligne vue, même quand
aucune valeur ne change. get_bands_to_enrich filtrant sur
`updated_at > crawled_at + 1 min`, la table entière redevenait « à enrichir »
après chaque crawl — le crawler repartait indéfiniment chercher des pages
inchangées sur Metal Archives.

La migration 016 ne bumpe plus updated_at que sur changement réel de la ligne.
La 017 introduit alors le signal manquant : quand MA modifie un groupe sur un
champ visible seulement sur sa page (line-up, albums), le listing ne montre
aucun diff. Le crawler capture donc l'horodatage "modified" affiché par MA
lui-même et lève enrich_pending quand ce texte change — une fois par
modification, sans boucle.

Autres correctifs :

- claim_job_trigger() réclamait n'importe quel trigger en attente. Le crawler
  raflait donc les 'geocoder_enqueue', qu'il ne sait pas exécuter, et les
  refermait en « job_type inconnu » — selon lequel des deux daemons
  interrogeait la table en premier. Borné aux CRAWLER_JOB_TYPES.
- run_enrich sautait la politesse après un échec, via un `continue` placé
  avant le sleep. Un échec de get_html étant le plus souvent un 403/429,
  c'était le pire moment pour enchaîner sans délai.
- FlareSolverrError n'était rattrapée nulle part : un blip réseau ou un crash
  du Chrome headless faisait perdre le run entier.
- La boucle du scheduler n'isolait aucune exception : une panne DB transitoire
  tuait le process, que restart: unless-stopped relançait en crash-loop.
- recover_stuck_runs() referme au démarrage les runs laissés en 'running' par
  une instance tuée brutalement.
- full_crawl publie sa progression et respire entre deux pays.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
Nicolas FRYDER 2026-08-20 16:43:48 +02:00
parent 228c678789
commit aa4ef68dd8
7 changed files with 250 additions and 23 deletions

View file

@ -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();

View file

@ -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;

View file

@ -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

View file

@ -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,10 +196,17 @@ 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
# 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)

View file

@ -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):
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):
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,
}

View file

@ -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)

View file

@ -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)