fix(robustesse): session Chrome abandonnée, requêtes vides re-payées, nettoyage trop zélé
Le crawler n'avait aucun arrêt propre. Sa session FlareSolverr n'est détruite qu'au moment d'en ouvrir une neuve : chaque redeploy abandonnait donc une instance Chrome persistante côté FlareSolverr. Il gère désormais SIGTERM/SIGINT, dort par tranches d'une seconde pour ne pas faire attendre une minute à l'arrêt, et rend sa session dans un finally. Le worker de géocodage ne mémorisait pas les résultats VIDES. Le commentaire annonçait « toujours mettre en cache pour éviter de rappeler », mais l'insertion était à l'intérieur du `if res:` : une requête sans résultat était re-payée à chaque tentative — jusqu'à MAX_GEO_TRIES passages, multipliés par les requêtes de repli, puis de nouveau après chaque aller-retour LLM qui remet les compteurs à zéro. Elles sont désormais mises en cache avec lat/lon NULL, et la lecture distingue « absent du cache » de « connu sans résultat ». crawl-runs/cleanup utilisait 30 minutes par défaut, plus court qu'un crawl complet Europe qui dure des heures. Le déclencher pendant un crawl légitime le marquait en erreur alors qu'il tournait toujours, et update_crawl_run_progress (filtré sur status='running') cessait silencieusement de publier : l'affichage restait figé jusqu'à la fin. Le défaut passe à 24 h, et les runs dont l'annulation est déjà demandée sont laissés au chemin coopératif. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
parent
1081ea4965
commit
2043705fc7
4 changed files with 84 additions and 8 deletions
|
|
@ -409,7 +409,15 @@ export default async function adminRoutes(fastify, opts) {
|
||||||
// ------------------------------------------------------------------
|
// ------------------------------------------------------------------
|
||||||
fastify.post("/admin/api/crawl-runs/cleanup", async (req, reply) => {
|
fastify.post("/admin/api/crawl-runs/cleanup", async (req, reply) => {
|
||||||
try {
|
try {
|
||||||
const olderThanMinutes = Math.max(1, Number(req.body?.older_than_minutes) || 30);
|
// 30 minutes par défaut était plus court qu'un crawl complet Europe, qui
|
||||||
|
// dure des heures : déclencher ce nettoyage pendant un crawl légitime le
|
||||||
|
// marquait en erreur alors qu'il tournait toujours, et
|
||||||
|
// update_crawl_run_progress (filtré sur status='running') cessait
|
||||||
|
// silencieusement de publier — l'affichage restait figé jusqu'à la fin.
|
||||||
|
// C'est exactement le double mensonge que la migration 014 avait corrigé
|
||||||
|
// sur le bouton Annuler. Le défaut couvre désormais le plus long run
|
||||||
|
// attendu ; un seuil plus court reste possible explicitement.
|
||||||
|
const olderThanMinutes = Math.max(1, Number(req.body?.older_than_minutes) || 24 * 60);
|
||||||
const r = await pool.query(`
|
const r = await pool.query(`
|
||||||
UPDATE crawl_run
|
UPDATE crawl_run
|
||||||
SET status = 'error',
|
SET status = 'error',
|
||||||
|
|
@ -417,6 +425,10 @@ export default async function adminRoutes(fastify, opts) {
|
||||||
error = 'annulé manuellement (run bloqué)'
|
error = 'annulé manuellement (run bloqué)'
|
||||||
WHERE status = 'running'
|
WHERE status = 'running'
|
||||||
AND started_at < now() - make_interval(mins => $1)
|
AND started_at < now() - make_interval(mins => $1)
|
||||||
|
-- Une annulation déjà demandée suit le chemin coopératif : le
|
||||||
|
-- crawler la lit et écrira lui-même 'cancelled'. Forcer 'error'
|
||||||
|
-- par-dessus reviendrait à mentir sur l'issue du run.
|
||||||
|
AND cancel_requested = FALSE
|
||||||
RETURNING id, run_type, started_at
|
RETURNING id, run_type, started_at
|
||||||
`, [olderThanMinutes]);
|
`, [olderThanMinutes]);
|
||||||
await writeAuditLog(pool, req.adminUsername, "cleanup_stuck_runs", "crawl_run", null, null, { cleaned: r.rows });
|
await writeAuditLog(pool, req.adminUsername, "cleanup_stuck_runs", "crawl_run", null, null, { cleaned: r.rows });
|
||||||
|
|
|
||||||
|
|
@ -368,11 +368,23 @@ describe("POST /admin/api/crawl-runs/cleanup", () => {
|
||||||
{ match: "INSERT INTO admin_audit_log", result: rows() },
|
{ match: "INSERT INTO admin_audit_log", result: rows() },
|
||||||
]);
|
]);
|
||||||
|
|
||||||
it("utilise 30 minutes par défaut", async () => {
|
it("le seuil par défaut dépasse la durée d'un crawl complet", async () => {
|
||||||
|
// 30 minutes était plus court qu'un crawl complet Europe, qui dure des
|
||||||
|
// heures : le nettoyage marquait alors en erreur un run parfaitement vivant,
|
||||||
|
// et update_crawl_run_progress (filtré sur status='running') cessait de
|
||||||
|
// publier — l'affichage restait figé jusqu'à la fin du run.
|
||||||
const handlers = cleanupPool();
|
const handlers = cleanupPool();
|
||||||
const app = buildApp(handlers);
|
const app = buildApp(handlers);
|
||||||
await app.inject({ method: "POST", url: "/admin/api/crawl-runs/cleanup", headers: auth, payload: {} });
|
await app.inject({ method: "POST", url: "/admin/api/crawl-runs/cleanup", headers: auth, payload: {} });
|
||||||
expect(pool.find("UPDATE crawl_run").values).toEqual([30]);
|
expect(pool.find("UPDATE crawl_run").values).toEqual([24 * 60]);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("laisse tranquille un run dont l'annulation est déjà demandée", async () => {
|
||||||
|
// L'annulation est coopérative : le crawler écrira lui-même 'cancelled'.
|
||||||
|
const handlers = cleanupPool();
|
||||||
|
const app = buildApp(handlers);
|
||||||
|
await app.inject({ method: "POST", url: "/admin/api/crawl-runs/cleanup", headers: auth, payload: {} });
|
||||||
|
expect(pool.find("UPDATE crawl_run").sql).toMatch(/cancel_requested = FALSE/);
|
||||||
});
|
});
|
||||||
|
|
||||||
it("ne touche que les runs marqués running", async () => {
|
it("ne touche que les runs marqués running", async () => {
|
||||||
|
|
|
||||||
|
|
@ -16,6 +16,7 @@ Variables d'env :
|
||||||
CRAWLER_FULL_CRAWL_INTERVAL_DAYS (défaut: 60, mettre 0 pour désactiver)
|
CRAWLER_FULL_CRAWL_INTERVAL_DAYS (défaut: 60, mettre 0 pour désactiver)
|
||||||
"""
|
"""
|
||||||
import logging
|
import logging
|
||||||
|
import signal
|
||||||
import sys
|
import sys
|
||||||
import time
|
import time
|
||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
|
|
@ -55,6 +56,18 @@ def _wait_flaresolverr(fs: FlareSolverr, max_wait: int = 120):
|
||||||
raise RuntimeError("FlareSolverr not reachable after timeout")
|
raise RuntimeError("FlareSolverr not reachable after timeout")
|
||||||
|
|
||||||
|
|
||||||
|
# Arrêt propre sur SIGTERM/SIGINT. Sans lui, chaque redeploy abandonnait une
|
||||||
|
# session Chrome persistante côté FlareSolverr : elle n'est détruite qu'au
|
||||||
|
# moment d'en ouvrir une neuve, jamais à l'extinction du crawler.
|
||||||
|
_shutdown = False
|
||||||
|
|
||||||
|
|
||||||
|
def _handle_shutdown(signum, frame):
|
||||||
|
global _shutdown
|
||||||
|
_shutdown = True
|
||||||
|
log.info(f"[main] signal {signum} reçu — arrêt propre demandé")
|
||||||
|
|
||||||
|
|
||||||
def _safe(label: str, fn):
|
def _safe(label: str, fn):
|
||||||
"""Exécute un job en isolant toute exception.
|
"""Exécute un job en isolant toute exception.
|
||||||
|
|
||||||
|
|
@ -69,6 +82,9 @@ def _safe(label: str, fn):
|
||||||
|
|
||||||
|
|
||||||
def main():
|
def main():
|
||||||
|
signal.signal(signal.SIGTERM, _handle_shutdown)
|
||||||
|
signal.signal(signal.SIGINT, _handle_shutdown)
|
||||||
|
|
||||||
log.info("[main] crawler starting")
|
log.info("[main] crawler starting")
|
||||||
|
|
||||||
# Runs/jobs laissés en 'running' par une instance précédente tuée brutalement
|
# Runs/jobs laissés en 'running' par une instance précédente tuée brutalement
|
||||||
|
|
@ -80,6 +96,13 @@ def main():
|
||||||
|
|
||||||
ma = MASession(fs)
|
ma = MASession(fs)
|
||||||
|
|
||||||
|
def _liberer_session():
|
||||||
|
"""Rend sa session Chrome à FlareSolverr avant de sortir."""
|
||||||
|
sid = getattr(ma, "_session_id", None)
|
||||||
|
if sid:
|
||||||
|
fs.destroy_session(sid)
|
||||||
|
ma._session_id = None
|
||||||
|
|
||||||
# ------------------------------------------------------------------
|
# ------------------------------------------------------------------
|
||||||
# Helpers pour lancer les jobs avec log
|
# Helpers pour lancer les jobs avec log
|
||||||
# ------------------------------------------------------------------
|
# ------------------------------------------------------------------
|
||||||
|
|
@ -133,11 +156,19 @@ def main():
|
||||||
_safe("full", job_full)
|
_safe("full", job_full)
|
||||||
|
|
||||||
log.info("[main] entering scheduler loop")
|
log.info("[main] entering scheduler loop")
|
||||||
while True:
|
try:
|
||||||
|
while not _shutdown:
|
||||||
_safe("run_pending", schedule.run_pending)
|
_safe("run_pending", schedule.run_pending)
|
||||||
_safe("check_job_triggers", lambda: _check_job_triggers(ma))
|
_safe("check_job_triggers", lambda: _check_job_triggers(ma))
|
||||||
_safe("check_health", _check_health)
|
_safe("check_health", _check_health)
|
||||||
time.sleep(60)
|
# Sommeil fractionné : un SIGTERM ne doit pas attendre une minute.
|
||||||
|
for _ in range(60):
|
||||||
|
if _shutdown:
|
||||||
|
break
|
||||||
|
time.sleep(1)
|
||||||
|
finally:
|
||||||
|
_safe("liberer_session", _liberer_session)
|
||||||
|
log.info("[main] crawler arrêté")
|
||||||
|
|
||||||
|
|
||||||
def _check_health():
|
def _check_health():
|
||||||
|
|
|
||||||
|
|
@ -368,6 +368,10 @@ def main():
|
||||||
(query,),
|
(query,),
|
||||||
)
|
)
|
||||||
cached = cur.fetchone()
|
cached = cur.fetchone()
|
||||||
|
if cached and cached[0] is None:
|
||||||
|
# Connu, et connu SANS résultat : inutile de rappeler l'API.
|
||||||
|
last_err = f"no_result_cache:'{query}'"
|
||||||
|
continue
|
||||||
if cached and cached[0] is not None:
|
if cached and cached[0] is not None:
|
||||||
c_lat, c_lon, c_conf, c_gran = cached
|
c_lat, c_lon, c_conf, c_gran = cached
|
||||||
if is_reliable(c_conf, c_gran, is_country_only):
|
if is_reliable(c_conf, c_gran, is_country_only):
|
||||||
|
|
@ -418,6 +422,23 @@ def main():
|
||||||
break
|
break
|
||||||
last_err = f"low_conf:{r_conf}:{r_gran}:'{query}'"
|
last_err = f"low_conf:{r_conf}:{r_gran}:'{query}'"
|
||||||
else:
|
else:
|
||||||
|
# Résultat vide mémorisé lui aussi (lat/lon NULL). Le
|
||||||
|
# commentaire ci-dessus disait « toujours mettre en cache »
|
||||||
|
# mais l'insertion était à l'intérieur du `if res:` : une
|
||||||
|
# requête sans résultat était donc re-payée à chaque
|
||||||
|
# tentative — jusqu'à MAX_GEO_TRIES passages, multipliés
|
||||||
|
# par les requêtes de repli, puis de nouveau après chaque
|
||||||
|
# aller-retour LLM qui remet les compteurs à zéro.
|
||||||
|
cur.execute(
|
||||||
|
"""
|
||||||
|
INSERT INTO geocode_cache
|
||||||
|
(query, provider, lat, lon, raw, confidence, granularity, updated_at)
|
||||||
|
VALUES (%s,'geoapify',NULL,NULL,%s,NULL,NULL,now())
|
||||||
|
ON CONFLICT (query) DO UPDATE
|
||||||
|
SET raw = EXCLUDED.raw, updated_at = now()
|
||||||
|
""",
|
||||||
|
(query, json.dumps({"features": []})),
|
||||||
|
)
|
||||||
last_err = f"no_result:'{query}'"
|
last_err = f"no_result:'{query}'"
|
||||||
except RuntimeError as exc:
|
except RuntimeError as exc:
|
||||||
if "429" in str(exc):
|
if "429" in str(exc):
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue