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>
196 lines
7 KiB
Python
196 lines
7 KiB
Python
"""
|
|
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
|
|
- Filet de sécurité (calcul glissant) : crawl complet Europe si plus de
|
|
CRAWLER_FULL_CRAWL_INTERVAL_DAYS jours depuis le dernier (checkpoint
|
|
last_full_crawl_at), vérifié chaque jour à 03:00 UTC + au démarrage
|
|
|
|
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_FULL_CRAWL_INTERVAL_DAYS (défaut: 60, mettre 0 pour désactiver)
|
|
"""
|
|
import logging
|
|
import sys
|
|
import time
|
|
from datetime import UTC, datetime
|
|
|
|
import schedule
|
|
|
|
from .config import (
|
|
ENRICH_LIMIT,
|
|
FLARESOLVERR_URL,
|
|
FS_TIMEOUT_MS,
|
|
FULL_CRAWL_INTERVAL_DAYS,
|
|
SCHED_ENRICH_H,
|
|
SCHED_INCREMENTAL_H,
|
|
)
|
|
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
|
|
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 _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)
|
|
|
|
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=ENRICH_LIMIT)
|
|
|
|
def job_full():
|
|
if not FULL_CRAWL_INTERVAL_DAYS:
|
|
return
|
|
last = get_checkpoint("last_full_crawl_at")
|
|
if last:
|
|
try:
|
|
last_dt = datetime.strptime(last, "%Y-%m-%d").replace(tzinfo=UTC)
|
|
elapsed_days = (datetime.now(UTC) - last_dt).days
|
|
if elapsed_days < FULL_CRAWL_INTERVAL_DAYS:
|
|
return
|
|
except ValueError:
|
|
pass # checkpoint mal formé, on relance par sécurité
|
|
log.info(f"[scheduler] → full Europe crawl (filet {FULL_CRAWL_INTERVAL_DAYS}j écoulé)")
|
|
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 FULL_CRAWL_INTERVAL_DAYS > 0:
|
|
# Vérifié chaque jour à 03:00 UTC ; ne se déclenche que si l'intervalle est écoulé
|
|
schedule.every().day.at("03:00").do(job_full)
|
|
log.info(f"[main] full crawl: filet de {FULL_CRAWL_INTERVAL_DAYS}j, vérifié chaque jour à 03:00 UTC")
|
|
|
|
# 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")
|
|
_safe("incremental", job_incremental)
|
|
_safe("enrich", job_enrich)
|
|
_safe("full", job_full)
|
|
|
|
log.info("[main] entering scheduler loop")
|
|
while True:
|
|
_safe("run_pending", schedule.run_pending)
|
|
_safe("check_job_triggers", lambda: _check_job_triggers(ma))
|
|
_safe("check_health", _check_health)
|
|
time.sleep(60)
|
|
|
|
|
|
def _check_health():
|
|
"""Contrôle horaire : base joignable, FlareSolverr joignable, données fraîches."""
|
|
import requests
|
|
|
|
from .config import FLARESOLVERR_URL
|
|
from .db import get_conn, log_event
|
|
|
|
maybe_run(
|
|
get_conn,
|
|
"crawler",
|
|
[
|
|
("database", database_probe(get_conn)),
|
|
# FlareSolverr est indispensable au crawl : injoignable, plus aucune
|
|
# page de Metal Archives n'est récupérable.
|
|
("flaresolverr", http_probe(requests.get, f"{FLARESOLVERR_URL}/health")),
|
|
# Le crawl incrémental tourne toutes les 4 h : au-delà de 12 h sans
|
|
# nouveau groupe vu, quelque chose est bloqué en amont.
|
|
("fraicheur_donnees", freshness_probe(
|
|
get_conn,
|
|
"SELECT max(started_at) FROM crawl_run WHERE status = 'done'",
|
|
12, "dernier run terminé")),
|
|
],
|
|
log_event=lambda lvl, msg: log_event(lvl, msg),
|
|
)
|
|
|
|
|
|
def _check_job_triggers(ma: "MASession"):
|
|
"""Consomme les job_triggers en attente créés depuis l'admin."""
|
|
result = claim_job_trigger()
|
|
if not result:
|
|
return
|
|
trigger_id, job_type = result
|
|
log.info(f"[main] job trigger #{trigger_id}: {job_type}")
|
|
error = None
|
|
try:
|
|
if job_type == "enrich":
|
|
run_enrich(ma, limit=ENRICH_LIMIT)
|
|
elif job_type == "incremental":
|
|
run_incremental(ma, "created")
|
|
run_incremental(ma, "modified")
|
|
elif job_type == "full_crawl":
|
|
run_full_crawl(ma)
|
|
else:
|
|
error = f"job_type inconnu: {job_type}"
|
|
log.warning(f"[main] {error}")
|
|
except Exception as e:
|
|
error = str(e)
|
|
log.error(f"[main] job trigger #{trigger_id} error: {e}", exc_info=True)
|
|
finally:
|
|
finish_job_trigger(trigger_id, error)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|