From 72a52a8d3b20dcaa96397f981c173332ec833af3 Mon Sep 17 00:00:00 2001 From: Nicolas Fryder Date: Thu, 2 Jul 2026 20:32:16 +0200 Subject: [PATCH] feat(geocoder): seuil de confiance strict, dedup, purge ancien pipeline MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Durcissement fiabilité (suite audit). Confiance & granularité - worker: un résultat n'est accepté ('done') que si confiance >= GEOCODE_MIN_CONFIDENCE (0.7) ET granularité non-grossière (rejette country/state/county/region). Sinon on continue les fallbacks, puis -> llm_needed. - fix majeur: sur cache hit la confiance était écrite 0.5 en dur (97% des lignes faussées). Elle est désormais lue depuis geocode_cache (recalculée du raw). - migration 012: colonnes confidence+granularity sur geocode_cache (recalcul des entrées Geoapify depuis le raw), geocode_granularity sur band_locations. Dedup (7x moins de travail) - fast-path: un band_location dont le (lieu,pays) est déjà 'done' copie le résultat sans appel API. Index fonctionnel lower(location_raw). Ancien pipeline retiré - geocode_queue n'est plus lu nulle part (endpoints /geocoding/reset-errors et /requeue-all supprimés, /live et /geocoding et Monitor basculés sur band_locations, panneau admin "ancien pipeline" retiré). Boutons reset (onglet Géocodage, zone dangereuse) - POST /admin/api/locations/reset-all: remet tout en queue + efface coords - POST /admin/api/geocode-cache/purge-nominatim: purge le cache Nominatim Divers - groq_worker: coût calculé par modèle (70B vs 8B) au lieu du tarif 70B fixe - GEOCODE_MIN_CONFIDENCE ajouté aux deux compose Co-Authored-By: Claude Opus 4.8 --- apps/admin/site/app.js | 128 ++--------- .../api/migrations/012_geocode_confidence.sql | 25 +++ apps/api/src/adminRoutes.js | 52 ++--- apps/geocoder/src/groq_worker.py | 14 +- apps/geocoder/src/worker.py | 202 ++++++++++++------ docker-compose.dev.yml | 1 + docker-compose.yml | 1 + 7 files changed, 227 insertions(+), 196 deletions(-) create mode 100644 apps/api/migrations/012_geocode_confidence.sql diff --git a/apps/admin/site/app.js b/apps/admin/site/app.js index 4c3e98c..0f2be69 100644 --- a/apps/admin/site/app.js +++ b/apps/admin/site/app.js @@ -995,7 +995,7 @@ function conflictCard(b) { // Géocodage // ------------------------------------------------------------------ async function renderGeocoding() { - if (!document.getElementById("geo-stats")) { + if (!document.getElementById("loc-stats")) { content().innerHTML = `
Pipeline géocodage
@@ -1004,7 +1004,7 @@ async function renderGeocoding() {
-

Nouveau pipeline — band_locations

+

Pipeline band_locations

@@ -1021,20 +1021,19 @@ async function renderGeocoding() {
-
-
-

Ancien pipeline — geocode_queue

-
- - -
+
+

⚠️ Zone dangereuse — tout recommencer

+

+ À utiliser si le géocodage est corrompu. « Reset complet » remet toutes les + localisations en queue et efface leurs coordonnées (les appels seront refaits, cache réutilisé + quand possible). « Purger cache Nominatim » supprime les vieux géocodages de l'ancien pipeline + pour forcer un re-géocodage Geoapify frais. +

+
+ +
-
-
-
-
-
-
+
@@ -1082,37 +1081,11 @@ async function renderGeocoding() { makeLocAction("loc-reset-llm", "loc-feedback", "/admin/api/locations/reset-llm", {}, "Remettre en queue les entrées llm_needed et manual ?"); makeLocAction("loc-requeue-all", "loc-feedback", "/admin/api/locations/requeue-all", { include_done: true }, "Re-queue TOUTES les band_locations (erreurs + LLM + manual + done) ?"); - // Ancien pipeline — actions - document.getElementById("geo-reset-errors").addEventListener("click", async () => { - const btn = document.getElementById("geo-reset-errors"); - btn.disabled = true; - const fb = document.getElementById("geo-feedback"); - fb.textContent = "Réinitialisation…"; fb.style.color = "var(--muted)"; - try { - const r = await api("/admin/api/geocoding/reset-errors", { method: "POST", body: JSON.stringify({}) }); - fb.textContent = `✓ ${r.count} entrée(s) remise(s) en queue.`; - fb.style.color = "var(--ok)"; - await loadGeocoding(); - } catch (e) { - fb.textContent = `Erreur: ${e.message}`; fb.style.color = "var(--err)"; - } finally { btn.disabled = false; } - }); - - document.getElementById("geo-requeue-done").addEventListener("click", async () => { - if (!confirm("Réinitialiser TOUTES les entrées geocode_queue (erreurs + fait) ?")) return; - const btn = document.getElementById("geo-requeue-done"); - btn.disabled = true; - const fb = document.getElementById("geo-feedback"); - fb.textContent = "Réinitialisation…"; fb.style.color = "var(--muted)"; - try { - const r = await api("/admin/api/geocoding/requeue-all", { method: "POST", body: JSON.stringify({ include_done: true }) }); - fb.textContent = `✓ ${r.count} entrée(s) re-queueée(s).`; - fb.style.color = "var(--ok)"; - await loadGeocoding(); - } catch (e) { - fb.textContent = `Erreur: ${e.message}`; fb.style.color = "var(--err)"; - } finally { btn.disabled = false; } - }); + // Zone dangereuse — reset complet & purge cache + makeLocAction("loc-reset-all", "danger-feedback", "/admin/api/locations/reset-all", {}, + "⚠️ RESET COMPLET : remettre TOUTES les localisations en queue et effacer leurs coordonnées ? Le worker va tout re-géocoder (cache réutilisé quand possible)."); + makeLocAction("cache-purge-nominatim", "danger-feedback", "/admin/api/geocode-cache/purge-nominatim", {}, + "⚠️ Supprimer toutes les entrées de cache Nominatim (ancien pipeline) ? Ces lieux seront re-géocodés via Geoapify (coûte des appels API)."); } await loadGeocoding(); @@ -1169,31 +1142,10 @@ async function loadGeocoding() { const provTxt = (r.providers || []).map((p) => `${esc(p.provider)}: ${p.n.toLocaleString("fr-FR")}${p.avg_conf != null ? ` (conf~${p.avg_conf})` : ""}`).join(" · "); const modelTxt = (r.llm_models || []).map((m) => `${esc(m.model)}: ${m.n} (null ${m.n_null})`).join(" · "); provEl.innerHTML = `
Sources géocodage : ${provTxt || "—"}
` + - (modelTxt ? `
Modèles LLM : ${modelTxt}
` : ""); + (modelTxt ? `
Modèles LLM : ${modelTxt}
` : "") + + `
Cache géocodage : ${(r.cache_size || 0).toLocaleString("fr-FR")} entrées
`; } - // Ancien pipeline — geocode_queue - const total = r.queue.reduce((s, x) => s + x.n, 0); - const done = r.queue.find((x) => x.status === "done")?.n || 0; - const queued = r.queue.find((x) => x.status === "queued")?.n || 0; - const errored = r.queue.find((x) => x.status === "error")?.n || 0; - const processing = r.queue.find((x) => x.status === "processing")?.n || 0; - const pct = total ? Math.round((done / total) * 100) : 0; - - const statsEl = document.getElementById("geo-stats"); - if (statsEl) statsEl.innerHTML = ` - ${statCard("Total queue", total)} - ${statCard("Géocodés", done, "ok")} - ${statCard("En attente", queued, queued ? "warn" : "")} - ${statCard("En cours", processing, processing ? "ok" : "")} - ${statCard("Erreurs", errored, errored ? "err" : "")} - ${statCard("Cache Geoapify", r.cache_size)} - `; - const bar = document.getElementById("geo-bar"); - if (bar) bar.style.width = pct + "%"; - const pctEl = document.getElementById("geo-pct"); - if (pctEl) pctEl.textContent = `${pct}% — ${done.toLocaleString("fr-FR")} / ${total.toLocaleString("fr-FR")}${processing ? ` (${processing} en cours)` : ""}`; - const recentWrap = document.getElementById("geo-recent-wrap"); if (recentWrap) recentWrap.innerHTML = ` @@ -1210,7 +1162,7 @@ async function loadGeocoding() {
MA IDNomPaysProviderQueryGéocodé leErreur
`; } catch (e) { - const pctEl = document.getElementById("geo-pct"); + const pctEl = document.getElementById("loc-pct"); if (pctEl) pctEl.textContent = `Erreur : ${esc(e.message)}`; } } @@ -1330,8 +1282,7 @@ async function renderJobs() {

Le géocodeur et le worker Groq tournent en continu. L'enqueue parse les bands en band_locations.

- - + Les resets/purges du géocodage sont dans l'onglet Géocodage.
@@ -1394,39 +1345,6 @@ async function renderJobs() { } finally { btn.disabled = false; } }); - // Geocoder reset errors - document.getElementById("cmd-reset-errors").addEventListener("click", async () => { - const btn = document.getElementById("cmd-reset-errors"); - btn.disabled = true; - const fb = document.getElementById("geo-cmd-feedback"); - fb.textContent = "Réinitialisation…"; fb.style.color = "var(--muted)"; - try { - const r = await api("/admin/api/geocoding/reset-errors", { method: "POST", body: JSON.stringify({}) }); - fb.textContent = `✓ ${r.count} entrée(s) remise(s) en queue.`; - fb.style.color = "var(--ok)"; - await loadCmdLogTail(); - } catch (e) { - fb.textContent = `Erreur: ${e.message}`; fb.style.color = "var(--err)"; - } finally { btn.disabled = false; } - }); - - // Geocoder requeue all - document.getElementById("cmd-requeue-all").addEventListener("click", async () => { - if (!confirm("Réinitialiser TOUTES les entrées (erreurs + déjà géocodées) ? Ceci relance le géocodage depuis zéro.")) return; - const btn = document.getElementById("cmd-requeue-all"); - btn.disabled = true; - const fb = document.getElementById("geo-cmd-feedback"); - fb.textContent = "Réinitialisation…"; fb.style.color = "var(--muted)"; - try { - const r = await api("/admin/api/geocoding/requeue-all", { method: "POST", body: JSON.stringify({ include_done: true }) }); - fb.textContent = `✓ ${r.count} entrée(s) re-queueée(s).`; - fb.style.color = "var(--ok)"; - await loadCmdLogTail(); - } catch (e) { - fb.textContent = `Erreur: ${e.message}`; fb.style.color = "var(--err)"; - } finally { btn.disabled = false; } - }); - // Cleanup document.getElementById("cmd-cleanup").addEventListener("click", async () => { const btn = document.getElementById("cmd-cleanup"); diff --git a/apps/api/migrations/012_geocode_confidence.sql b/apps/api/migrations/012_geocode_confidence.sql new file mode 100644 index 0000000..05b33ea --- /dev/null +++ b/apps/api/migrations/012_geocode_confidence.sql @@ -0,0 +1,25 @@ +-- 012_geocode_confidence.sql +-- Durcissement du géocodage : on rend la confiance et la granularité exploitables. +-- +-- Bug corrigé côté worker : sur un cache hit, la confiance était écrite à 0.5 en +-- dur, jetant la vraie valeur présente dans le raw. On ajoute des colonnes +-- confidence + granularity à geocode_cache et on les recalcule depuis le raw +-- Geoapify (FeatureCollection). Aucune purge ici (la purge Nominatim + le reset +-- complet sont des boutons admin, pour maîtriser le budget API). + +ALTER TABLE band_locations + ADD COLUMN IF NOT EXISTS geocode_granularity TEXT; + +ALTER TABLE geocode_cache + ADD COLUMN IF NOT EXISTS confidence REAL, + ADD COLUMN IF NOT EXISTS granularity TEXT; + +-- Recalcul pour les entrées Geoapify (raw = FeatureCollection avec rank.confidence) +UPDATE geocode_cache +SET confidence = NULLIF(raw->'features'->0->'properties'->'rank'->>'confidence', '')::real, + granularity = raw->'features'->0->'properties'->>'result_type' +WHERE raw ? 'features' + AND confidence IS NULL; + +-- Fast-path de dedup au niveau step : lookup par lieu normalisé +CREATE INDEX IF NOT EXISTS idx_bl_locraw_lower ON band_locations (lower(location_raw)); diff --git a/apps/api/src/adminRoutes.js b/apps/api/src/adminRoutes.js index 3dd9427..9f299f4 100644 --- a/apps/api/src/adminRoutes.js +++ b/apps/api/src/adminRoutes.js @@ -454,10 +454,7 @@ export default async function adminRoutes(fastify, opts) { // ------------------------------------------------------------------ fastify.get("/admin/api/geocoding", async (req, reply) => { try { - const [queue, cache, recent, locations, llm, providers, llmModels] = await Promise.all([ - pool.query(` - SELECT status, count(*)::int AS n FROM geocode_queue GROUP BY status ORDER BY n DESC - `), + const [cache, recent, locations, llm, providers, llmModels] = await Promise.all([ pool.query(`SELECT count(*)::int AS n FROM geocode_cache`), pool.query(` SELECT b.ma_id, b.name, b.country, b.geocode_provider, b.geocoded_at, @@ -494,7 +491,6 @@ export default async function adminRoutes(fastify, opts) { ]); return { ok: true, - queue: queue.rows, cache_size: cache.rows[0].n, recent: recent.rows, locations: locations.rows, @@ -587,45 +583,48 @@ export default async function adminRoutes(fastify, opts) { }); // ------------------------------------------------------------------ - // Géocodage — actions directes sur la queue + // Géocodage — RESET complet & purge cache (pour tout relancer proprement) // ------------------------------------------------------------------ - fastify.post("/admin/api/geocoding/reset-errors", async (req, reply) => { + // Remet TOUTES les band_locations (hors pays-seul) en 'queued' et efface les + // résultats : lat/lon, provider, confiance, granularité, query, erreurs, essais. + fastify.post("/admin/api/locations/reset-all", async (req, reply) => { try { const r = await pool.query(` - UPDATE geocode_queue - SET status='queued', next_run_at=now(), updated_at=now() - WHERE status='error' + UPDATE band_locations + SET geocode_status='queued', + lat=NULL, lon=NULL, + geocode_provider=NULL, geocode_confidence=NULL, geocode_granularity=NULL, + geocode_query=NULL, geocode_error=NULL, + geocode_tries_geo=0, geocode_tries_llm=0, + geocode_next_at=now(), updated_at=now() + WHERE is_country_only = FALSE `); const count = r.rowCount; await pool.query( `INSERT INTO crawl_log (level, message) VALUES ($1, $2)`, - ["info", `[admin:${req.adminUsername}] geocoding reset-errors: ${count} entrée(s) remise(s) en queue`] + ["warning", `[admin:${req.adminUsername}] RESET géocodage complet: ${count} band_locations remises à zéro`] ).catch(() => {}); return { ok: true, count }; } catch (err) { fastify.log.error(err); - return reply.code(500).send({ ok: false, error: "Erreur reset-errors" }); + return reply.code(500).send({ ok: false, error: "Erreur reset-all" }); } }); - fastify.post("/admin/api/geocoding/requeue-all", async (req, reply) => { + // Purge les entrées de cache issues de l'ancien pipeline Nominatim (raw avec + // place_rank). Force le worker à re-géocoder ces lieux via Geoapify. + fastify.post("/admin/api/geocode-cache/purge-nominatim", async (req, reply) => { try { - const { include_done = false } = req.body || {}; - const statuses = include_done ? ["error", "done"] : ["error"]; - const r = await pool.query(` - UPDATE geocode_queue - SET status='queued', next_run_at=now(), tries=0, last_error=NULL, updated_at=now() - WHERE status = ANY($1::text[]) - `, [statuses]); + const r = await pool.query(`DELETE FROM geocode_cache WHERE raw ? 'place_rank'`); const count = r.rowCount; await pool.query( `INSERT INTO crawl_log (level, message) VALUES ($1, $2)`, - ["info", `[admin:${req.adminUsername}] geocoding requeue-all (include_done=${include_done}): ${count} re-queueée(s)`] + ["warning", `[admin:${req.adminUsername}] purge cache Nominatim: ${count} entrée(s) supprimée(s)`] ).catch(() => {}); return { ok: true, count }; } catch (err) { fastify.log.error(err); - return reply.code(500).send({ ok: false, error: "Erreur requeue-all" }); + return reply.code(500).send({ ok: false, error: "Erreur purge-nominatim" }); } }); @@ -712,11 +711,12 @@ export default async function adminRoutes(fastify, opts) { SELECT id, job_type, status, requested_by, created_at FROM job_triggers WHERE status = 'pending' ORDER BY created_at ASC `), - pool.query(`SELECT status, count(*)::int AS n FROM geocode_queue GROUP BY status ORDER BY n DESC`), + pool.query(`SELECT geocode_status AS status, count(*)::int AS n FROM band_locations GROUP BY geocode_status ORDER BY n DESC`), pool.query(` - SELECT gq.ma_id, b.name, b.country, gq.tries, gq.last_error, gq.updated_at - FROM geocode_queue gq JOIN bands b ON b.ma_id = gq.ma_id - WHERE gq.status = 'processing' LIMIT 3 + SELECT bl.ma_id, b.name, b.country, bl.location_raw, bl.geocode_tries_geo AS tries, + bl.geocode_error AS last_error, bl.updated_at + FROM band_locations bl JOIN bands b ON b.ma_id = bl.ma_id + WHERE bl.geocode_status = 'processing' LIMIT 3 `), pool.query(`SELECT key, value, updated_at FROM crawl_checkpoint ORDER BY key`), ]); diff --git a/apps/geocoder/src/groq_worker.py b/apps/geocoder/src/groq_worker.py index 53730a6..81463c1 100644 --- a/apps/geocoder/src/groq_worker.py +++ b/apps/geocoder/src/groq_worker.py @@ -39,6 +39,17 @@ MODELS = [ ("llama-3.1-8b-instant", 30, 14400), ] +# Tarifs Groq approximatifs par modèle : (usd/token_in, usd/token_out) +PRICING = { + "llama-3.3-70b-versatile": (0.59e-6, 0.79e-6), + "llama-3.1-8b-instant": (0.05e-6, 0.08e-6), +} + + +def _cost(model: str, tok_in: int, tok_out: int) -> float: + pin, pout = PRICING.get(model, (0.59e-6, 0.79e-6)) + return tok_in * pin + tok_out * pout + _SYSTEM = ( "You are a precise location data extraction assistant. " "Your only job is to extract a city name and ISO 3166-1 alpha-2 country code " @@ -224,8 +235,7 @@ def main(): city = (parsed or {}).get("city") iso2 = (parsed or {}).get("country") or country is_null = not city - # Coût approximatif llama 70B - cost_usd = (tok_in * 0.00000059 + tok_out * 0.00000079) + cost_usd = _cost(chosen_model, tok_in, tok_out) prompt_text = _USER_TMPL.format( country=country or "unknown (European metal band)", diff --git a/apps/geocoder/src/worker.py b/apps/geocoder/src/worker.py index 202042a..640dc6c 100644 --- a/apps/geocoder/src/worker.py +++ b/apps/geocoder/src/worker.py @@ -2,11 +2,16 @@ worker.py — Géocodeur Geoapify pour band_locations. Boucle infinie, traite geocode_status='queued' un par un avec SKIP LOCKED. -Stratégie progressive : - 1. Vérifier geocode_cache (requête déjà faite → gratuit) - 2. Essayer les requêtes fallback du parser (du plus précis au plus vague) - 3. Après MAX_GEO_TRIES échecs → status='llm_needed' pour groq_worker - 4. Succès → mise à jour band_locations + sync bands.lat/lon +Stratégie : + 0. Fast-path dedup : si un autre band_location du MÊME (lieu, pays) est déjà + 'done', on copie son résultat sans appel API. + 1. Sinon, essayer les requêtes fallback du parser (précis → vague). + Pour chaque requête : cache d'abord (geocode_cache), puis Geoapify. + 2. Un résultat n'est accepté (done) que s'il est FIABLE : + confiance >= MIN_CONFIDENCE ET granularité non-grossière. + Sinon on continue les fallbacks. + 3. Après MAX_GEO_TRIES sans résultat fiable → status='llm_needed'. + 4. Succès → band_locations + sync bands.lat/lon (origine). """ import os @@ -21,10 +26,14 @@ from parser import build_fallback_queries GEOAPIFY_API_KEY = os.environ.get("GEOAPIFY_API_KEY", "").strip() GEOAPIFY_BASE = "https://api.geoapify.com/v1/geocode/search" -MIN_DELAY = float(os.environ.get("GEOCODER_MIN_DELAY", "0.22")) -JITTER = float(os.environ.get("GEOCODER_JITTER", "0.10")) -MAX_PER_RUN = int(os.environ.get("GEOCODE_MAX_PER_RUN", "100000")) -MAX_GEO_TRIES = int(os.environ.get("GEOCODE_MAX_GEO_TRIES", "3")) +MIN_DELAY = float(os.environ.get("GEOCODER_MIN_DELAY", "0.22")) +JITTER = float(os.environ.get("GEOCODER_JITTER", "0.10")) +MAX_PER_RUN = int(os.environ.get("GEOCODE_MAX_PER_RUN", "100000")) +MAX_GEO_TRIES = int(os.environ.get("GEOCODE_MAX_GEO_TRIES", "3")) +MIN_CONFIDENCE = float(os.environ.get("GEOCODE_MIN_CONFIDENCE", "0.7")) + +# Granularités jugées trop grossières pour une carte de villes : on refuse (→ LLM). +COARSE_TYPES = {"country", "state", "county", "region", "province", "political"} _last_call = 0.0 @@ -38,6 +47,17 @@ def _polite_sleep(): _last_call = time.monotonic() +def is_reliable(confidence, granularity, is_country_only) -> bool: + """Un géocodage est fiable s'il est assez confiant ET assez précis.""" + if is_country_only: + return True # résolu par centroïde à l'enqueue, granularité pays assumée + if confidence is None or confidence < MIN_CONFIDENCE: + return False + if granularity and granularity.lower() in COARSE_TYPES: + return False + return True + + def geoapify_search(query: str) -> dict | None: if not GEOAPIFY_API_KEY: raise RuntimeError("GEOAPIFY_API_KEY not set") @@ -55,13 +75,15 @@ def geoapify_search(query: str) -> dict | None: coords = features[0].get("geometry", {}).get("coordinates", []) if len(coords) < 2: return None - rank = props.get("rank") or {} - confidence = rank.get("confidence", 0.5) if isinstance(rank, dict) else 0.5 + rank = props.get("rank") or {} + confidence = rank.get("confidence") if isinstance(rank, dict) else None + granularity = props.get("result_type") return { - "lat": props.get("lat") or coords[1], - "lon": props.get("lon") or coords[0], - "confidence": confidence, - "raw": data, + "lat": props.get("lat") or coords[1], + "lon": props.get("lon") or coords[0], + "confidence": confidence, + "granularity": granularity, + "raw": data, } @@ -69,8 +91,7 @@ def sync_bands_primary(cur, ma_id: int): """Met à jour bands.lat/lon avec le lieu d'ORIGINE géocodé (step_order le plus bas). Carte « metal from Europe » : on veut le point de formation du groupe, pas - sa dernière localisation (ex. un groupe de Thessaloniki parti à Boston doit - rester à Thessaloniki, en Europe). + sa dernière localisation. """ cur.execute( """ @@ -102,17 +123,63 @@ def sync_bands_primary(cur, ma_id: int): ) +def mark_done(cur, loc_id, ma_id, lat, lon, query, provider, confidence, granularity): + cur.execute( + """ + UPDATE band_locations + SET lat=%s, lon=%s, + geocode_status='done', + geocode_query=%s, + geocode_provider=%s, + geocode_confidence=%s, + geocode_granularity=%s, + geocode_error=NULL, + updated_at=now() + WHERE id=%s + """, + (lat, lon, query, provider, confidence, granularity, loc_id), + ) + sync_bands_primary(cur, ma_id) + + +def try_dedup(cur, loc_id, ma_id, location_raw, country) -> bool: + """Fast-path : copier le résultat d'un band_location identique déjà 'done'.""" + cur.execute( + """ + SELECT bl2.lat, bl2.lon, bl2.geocode_confidence, bl2.geocode_granularity, + bl2.geocode_query + FROM band_locations bl2 + JOIN bands b2 ON b2.ma_id = bl2.ma_id + WHERE lower(bl2.location_raw) = lower(%s) + AND COALESCE(b2.country,'') = COALESCE(%s,'') + AND bl2.geocode_status = 'done' + AND bl2.lat IS NOT NULL + LIMIT 1 + """, + (location_raw, country), + ) + hit = cur.fetchone() + if not hit: + return False + lat, lon, conf, gran, query = hit + mark_done(cur, loc_id, ma_id, lat, lon, query, 'dedup', conf, gran) + print(f"[worker] id={loc_id} DEDUP '{location_raw}' ({country})") + return True + + def main(): dsn = os.environ["DATABASE_URL"] conn = psycopg2.connect(dsn) conn.autocommit = True + print(f"[worker] démarré — seuil confiance={MIN_CONFIDENCE}, max_tries={MAX_GEO_TRIES}") + processed = 0 with conn.cursor() as cur: while processed < MAX_PER_RUN: cur.execute( """ - SELECT bl.id, bl.ma_id, bl.location_raw, + SELECT bl.id, bl.ma_id, bl.location_raw, bl.is_country_only, bl.geocode_tries_geo, bl.geocode_query, b.country FROM band_locations bl JOIN bands b ON b.ma_id = bl.ma_id @@ -129,21 +196,25 @@ def main(): time.sleep(60) continue - loc_id, ma_id, location_raw, tries, llm_query, country = row + loc_id, ma_id, location_raw, is_country_only, tries, llm_query, country = row cur.execute( "UPDATE band_locations SET geocode_status='processing', updated_at=now() WHERE id=%s", (loc_id,), ) - # Liste de requêtes à essayer : requête LLM en premier si disponible + # 0. Fast-path dedup (aucun appel API) + if try_dedup(cur, loc_id, ma_id, location_raw, country): + processed += 1 + continue + + # Requêtes à essayer : requête LLM en premier si disponible fallbacks = build_fallback_queries(location_raw, country) if llm_query and llm_query not in fallbacks: queries = [llm_query] + fallbacks else: queries = fallbacks - # Déduplique en gardant l'ordre seen: set[str] = set() unique: list[str] = [] for q in queries: @@ -151,53 +222,71 @@ def main(): seen.add(q) unique.append(q) - success = False - rate_limit = False - used_query = None - lat = lon = None - confidence = 0.5 - provider = 'geoapify' - last_err = None + success = False + rate_limit = False + used_query = None + lat = lon = None + confidence = None + granularity = None + provider = 'geoapify' + last_err = None for query in unique: - # Cache hit ? + # Cache hit ? (avec confiance/granularité recalculées) cur.execute( - "SELECT lat, lon FROM geocode_cache WHERE query=%s", + "SELECT lat, lon, confidence, granularity FROM geocode_cache WHERE query=%s", (query,), ) cached = cur.fetchone() if cached and cached[0] is not None: - lat, lon = cached - used_query = query - provider = 'geoapify-cache' - success = True - print(f"[worker] id={loc_id} cache HIT '{query}'") - break + c_lat, c_lon, c_conf, c_gran = cached + if is_reliable(c_conf, c_gran, is_country_only): + lat, lon = c_lat, c_lon + confidence = c_conf + granularity = c_gran + used_query = query + provider = 'geoapify-cache' + success = True + print(f"[worker] id={loc_id} cache HIT '{query}' conf={c_conf}") + break + # cache présent mais peu fiable → inutile de rappeler l'API pour cette requête + last_err = f"low_conf_cache:{c_conf}:'{query}'" + continue try: res = geoapify_search(query) if res: - lat = float(res["lat"]) - lon = float(res["lon"]) - confidence = float(res.get("confidence") or 0.5) - used_query = query + r_lat = float(res["lat"]) + r_lon = float(res["lon"]) + r_conf = res.get("confidence") + r_gran = res.get("granularity") + # Toujours mettre en cache (même peu fiable) pour éviter de rappeler cur.execute( """ INSERT INTO geocode_cache - (query, provider, lat, lon, geom, raw, updated_at) + (query, provider, lat, lon, geom, raw, confidence, granularity, updated_at) VALUES (%s,'geoapify',%s,%s, ST_SetSRID(ST_MakePoint(%s,%s),4326)::geography, - %s, now()) + %s, %s, %s, now()) ON CONFLICT (query) DO UPDATE SET lat=EXCLUDED.lat, lon=EXCLUDED.lon, geom=EXCLUDED.geom, raw=EXCLUDED.raw, + confidence=EXCLUDED.confidence, + granularity=EXCLUDED.granularity, provider=EXCLUDED.provider, updated_at=now() """, - (query, lat, lon, lon, lat, json.dumps(res["raw"])), + (query, r_lat, r_lon, r_lon, r_lat, + json.dumps(res["raw"]), r_conf, r_gran), ) - success = True - print(f"[worker] id={loc_id} OK '{query}' lat={lat:.4f} lon={lon:.4f}") - break + if is_reliable(r_conf, r_gran, is_country_only): + lat, lon = r_lat, r_lon + confidence = r_conf + granularity = r_gran + used_query = query + success = True + print(f"[worker] id={loc_id} OK '{query}' conf={r_conf} gran={r_gran}") + break + last_err = f"low_conf:{r_conf}:{r_gran}:'{query}'" else: last_err = f"no_result:'{query}'" except RuntimeError as exc: @@ -224,21 +313,8 @@ def main(): continue if success: - cur.execute( - """ - UPDATE band_locations - SET lat=%s, lon=%s, - geocode_status='done', - geocode_query=%s, - geocode_provider=%s, - geocode_confidence=%s, - geocode_error=NULL, - updated_at=now() - WHERE id=%s - """, - (lat, lon, used_query, provider, confidence, loc_id), - ) - sync_bands_primary(cur, ma_id) + mark_done(cur, loc_id, ma_id, lat, lon, used_query, + provider, confidence, granularity) else: new_tries = tries + 1 if new_tries >= MAX_GEO_TRIES: @@ -251,9 +327,9 @@ def main(): updated_at=now() WHERE id=%s """, - (new_tries, last_err or "all_fallbacks_failed", loc_id), + (new_tries, last_err or "no_reliable_match", loc_id), ) - print(f"[worker] id={loc_id} → llm_needed après {new_tries} essais") + print(f"[worker] id={loc_id} → llm_needed après {new_tries} essais ({last_err})") else: backoff = min(1440, 30 * (2 ** new_tries)) cur.execute( diff --git a/docker-compose.dev.yml b/docker-compose.dev.yml index 024e9f0..b97cf34 100644 --- a/docker-compose.dev.yml +++ b/docker-compose.dev.yml @@ -61,6 +61,7 @@ services: GEOCODER_JITTER: "0.10" GEOCODE_MAX_PER_RUN: "100000" GEOCODE_MAX_GEO_TRIES: "3" + GEOCODE_MIN_CONFIDENCE: "0.7" command: ["python", "src/worker.py"] networks: - coolify diff --git a/docker-compose.yml b/docker-compose.yml index 5767230..76237c8 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -33,6 +33,7 @@ services: GEOCODER_JITTER: "0.10" GEOCODE_MAX_PER_RUN: "100000" GEOCODE_MAX_GEO_TRIES: "3" + GEOCODE_MIN_CONFIDENCE: "0.7" command: ["python", "src/worker.py"] networks: - coolify