From 5237f3666d63bf876963b6b3de302ac30a9764de Mon Sep 17 00:00:00 2001 From: Nicolas Fryder Date: Tue, 18 Aug 2026 18:20:15 +0200 Subject: [PATCH] feat: validation syntaxique du SQL + supervision des services de fond MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Deux angles morts fermés. 1. Syntaxe SQL sans conteneur (apps/api/test/sqlSyntax.test.js) Le faux pool vérifiait la FORME du SQL mais ne l'exécutait jamais : une requête syntaxiquement invalide passait tous les tests et n'échouait qu'en production — c'est précisément ce qui s'était produit avec crawler_pending. Chaque requête réellement émise par les 28 routes est désormais parsée avec la grammaire PostgreSQL (node-sql-parser), y compris les SET dynamiques du PATCH et les casts par type de resolve-conflict. 48 tests, aucun conteneur. Limite déclarée explicitement : la sémantique n'est pas validée, et deux requêtes bâties sur jsonb_build_object ne sont pas parsables — le test échoue si une route cesse d'avoir la moindre requête vérifiable, pour éviter qu'il passe au vert à vide. 2. Supervision des services de fond (migration 015) Le crawler et les workers ne sont pas exposés par Traefik : aucune sonde HTTP ne peut les atteindre. Un crawler dont FlareSolverr était injoignable, ou un worker à court de quota Geoapify, restait muet — le seul symptôme était l'absence de données nouvelles, qu'il fallait remarquer soi-même. Chaque service écrit un battement de cœur horaire dans service_health : - crawler : base, FlareSolverr joignable, dernier run terminé < 12 h - geocoder : base, clé Geoapify présente, API joignable, progression < 6 h Une ligne par service, écrasée à chaque contrôle. L'API calcule `stale` en SQL (> 2 h sans écriture) : un service arrêté cesse d'écrire, et son dernier contrôle réussi le ferait sinon passer pour sain indéfiniment. Le Pilotage affiche une carte « Services » et remonte chaque service dégradé ou silencieux en alerte actionnable. Règle appliquée aux sondes : aucune ne peut interrompre le service qu'elle surveille. Toute exception devient un échec de sonde, l'écriture du résultat et la journalisation échouent en silence. Un contrôle de santé qui fait tomber le crawler serait pire que pas de contrôle. Tests : 580 JS (+51), 63 Python (+20), 85 Playwright (+3) Co-Authored-By: Claude --- README-CI.md | 15 +- apps/admin/site/app.js | 31 +++ apps/admin/test/e2e/admin.spec.js | 7 +- apps/admin/test/e2e/fixtures.mjs | 5 + apps/admin/test/e2e/tools.spec.js | 36 +++ apps/api/migrations/015_service_health.sql | 24 ++ apps/api/src/adminRoutes.js | 12 +- apps/api/test/adminRoutesCoverage.test.js | 35 +++ apps/api/test/sqlSyntax.test.js | 245 +++++++++++++++++++++ apps/crawler/src/health.py | 133 +++++++++++ apps/crawler/src/main.py | 28 +++ apps/crawler/tests/test_health.py | 217 ++++++++++++++++++ apps/geocoder/src/health.py | 133 +++++++++++ apps/geocoder/src/worker.py | 31 +++ package-lock.json | 32 +++ package.json | 1 + 16 files changed, 977 insertions(+), 8 deletions(-) create mode 100644 apps/api/migrations/015_service_health.sql create mode 100644 apps/api/test/sqlSyntax.test.js create mode 100644 apps/crawler/src/health.py create mode 100644 apps/crawler/tests/test_health.py create mode 100644 apps/geocoder/src/health.py diff --git a/README-CI.md b/README-CI.md index c74ddee..d207b43 100644 --- a/README-CI.md +++ b/README-CI.md @@ -28,11 +28,11 @@ pip install -r requirements-dev.txt | Commande | Ce que ça fait | Froid | Incrémental | |---|---|---|---| -| `npm run check` | **La porte** : ESLint + tsc + Ruff + 572 tests | **6 s** | — | +| `npm run check` | **La porte** : ESLint + tsc + Ruff + 620 tests | **6 s** | — | | `npm run check:sequential` | Idem, en série (pour isoler un échec) | 14 s | — | | `npm run test:mutation` | Mutation, logique pure (352 mutants) | 33 s | ~5 s | | `npm run test:mutation:full` | Mutation, toute l'API (2159 mutants) | ~5 min | **8 s** | -| `npm run test:e2e` | 82 parcours Playwright (dashboard admin) | 27 s | — | +| `npm run test:e2e` | 85 parcours Playwright (dashboard admin) | 26 s | — | | `npm run test:e2e:install` | Récupère Chromium (clone neuf) | — | — | | `npm run check:full` | `check` + e2e + audits + mutation complète | — | ~1 min | | `npm run check:clean` | Purge les caches ESLint / tsc | — | — | @@ -87,6 +87,8 @@ Le plancher de ~3 min à froid est celui de Stryker sur 1511 mutants, assumé. | **Parcours (e2e)** | `apps/admin/test/e2e/admin.spec.js` | Playwright sur la VRAIE app Fastify + faux pool : connexion, recherche, édition, déblocage d'une localisation, annulation d'un traitement | | **a11y des vues rendues** | `apps/admin/test/e2e/a11y.spec.js` | axe-core dans un vrai navigateur — contraste inclus, ce que jsdom ne sait pas calculer | | **Mutation** | `stryker*.config.json` | 82 % (logique pure) / 68 % (API complète) | +| **Syntaxe SQL** | `apps/api/test/sqlSyntax.test.js` | Chaque requête réellement émise est parsée avec la grammaire PostgreSQL — sans conteneur | +| **Santé des services** | `apps/crawler/tests/test_health.py` | Sondes du crawler et du geocoder : base, dépendances externes, fraîcheur des données | ## Principe des tests @@ -155,8 +157,13 @@ plus lents — contention au démarrage des navigateurs). ne calcule pas les styles. Le dashboard admin, lui, est couvert par axe dans un vrai navigateur. - `migrate.js` : s'exécute au démarrage du conteneur et appelle `process.exit`. -- Crawler et workers de géocodage : seuls le parseur et l'annulation sont - couverts, le reste est de l'I/O réseau et base. +- Crawler et workers de géocodage : le parseur, l'annulation et les contrôles + de santé sont couverts ; le reste est de l'I/O réseau et base. +- La **sémantique** SQL (noms de colonnes, types, comportement PostGIS) : le + parseur ne valide que la syntaxe. Deux requêtes bâties sur + `jsonb_build_object` (`/admin/api/activity` et l'agrégation des clusters) ne + sont même pas parsables — c'est déclaré explicitement dans le test, et c'est + le prix assumé de ne pas monter de Postgres. - `apps/web/site/pure.js` est testé mais **exclu de la mutation** : il est chargé via `fs` + `vm` (comme le fait le navigateur), donc l'instrumentation Stryker ne l'atteindrait pas et afficherait un score faussement parfait. diff --git a/apps/admin/site/app.js b/apps/admin/site/app.js index 2cb6179..6058087 100644 --- a/apps/admin/site/app.js +++ b/apps/admin/site/app.js @@ -337,6 +337,26 @@ function wireDangerZone() { }); } +/** + * État des services de fond (crawler, geocoder…). + * + * Ils ne sont pas exposés par Traefik : sans ce battement de cœur, un crawler + * dont FlareSolverr est injoignable reste invisible, et le seul symptôme est + * l'absence de données nouvelles. Un contrôle trop ancien vaut une panne : + * un service arrêté cesse simplement de mettre sa ligne à jour. + */ +function servicesCard(services) { + if (!services.length) { + return healthCard("Services", "—", "aucun battement de cœur reçu", "warn"); + } + const ko = services.filter((s) => !s.ok || s.stale); + const detail = ko.length + ? ko.map((s) => `${esc(s.service)} : ${s.stale ? "silencieux" : esc(s.error || "dégradé")}`).join(" · ") + : services.map((s) => esc(s.service)).join(" · "); + return healthCard("Services", `${services.length - ko.length}/${services.length}`, detail, + ko.length ? "err" : "ok"); +} + function healthCard(label, value, sub, tone = "") { return `
${esc(label)}
@@ -382,6 +402,7 @@ async function loadPilotage() { `${fmtNum(errors)} erreurs · ${fmtNum(llmNeeded)} LLM`, errors + llmNeeded === 0 ? "ok" : "err"), healthCard("Coût LLM", `$${cost.toFixed(2)}`, `${fmtNum(geo.llm_cache?.n || 0)} appels`), + servicesCard(live.services || []), ].join(""); // ---- Alertes actionnables ---- @@ -419,6 +440,16 @@ async function loadPilotage() { actions: [{ label: "Crawl incrémental", job: "incremental" }], }); } + for (const svc of (live.services || []).filter((x) => !x.ok || x.stale)) { + alerts.push({ + tone: "err", + text: svc.stale + ? `Service « ${svc.service} » silencieux depuis plus de 2 h — arrêté ou bloqué` + : `Service « ${svc.service} » dégradé : ${svc.error || "sonde en échec"}`, + actions: [{ label: "Voir le journal", href: "#/activity" }], + }); + } + const stuck = (live.active_runs || []).filter( (r) => Date.now() - new Date(r.started_at) > 30 * 60000 && !r.cancel_requested); if (stuck.length) { diff --git a/apps/admin/test/e2e/admin.spec.js b/apps/admin/test/e2e/admin.spec.js index ad7bbb0..b56ef78 100644 --- a/apps/admin/test/e2e/admin.spec.js +++ b/apps/admin/test/e2e/admin.spec.js @@ -99,9 +99,10 @@ test.describe("Pilotage — est-ce que ça tourne ?", () => { test("le bandeau de santé résume l'état du pipeline", async ({ page }) => { const strip = page.locator("#p-health"); - await expect(strip.locator(".health-card")).toHaveCount(5); - await expect(strip).toContainText("Géocodage"); - await expect(strip).toContainText("Coût LLM"); + await expect(strip.locator(".health-card")).toHaveCount(6); + for (const label of ["Groupes", "Géocodage", "En file", "Bloqués", "Coût LLM", "Services"]) { + await expect(strip).toContainText(label); + } }); // Le cœur de la refonte : plus de « je constate ici, j'agis ailleurs ». diff --git a/apps/admin/test/e2e/fixtures.mjs b/apps/admin/test/e2e/fixtures.mjs index 48171f0..5ef5e38 100644 --- a/apps/admin/test/e2e/fixtures.mjs +++ b/apps/admin/test/e2e/fixtures.mjs @@ -105,6 +105,11 @@ export function seedHandlers(over = {}) { { match: "WHERE status = 'pending'", result: rows(...pendingJobs) }, { match: "geocode_status = 'processing'", result: rows() }, { match: "GROUP BY geocode_status", result: rows(...geoStatuses) }, + { match: "FROM service_health", result: rows( + { service: "crawler", ok: true, checks: { database: true, flaresolverr: true }, error: null, + checked_at: "2026-08-18T14:00:00Z", stale: false }, + { service: "geocoder", ok: true, checks: { database: true, geoapify_joignable: true }, error: null, + checked_at: "2026-08-18T14:00:00Z", stale: false }) }, { match: "FROM crawl_checkpoint", result: rows({ key: "last_full_crawl_at", value: "2026-08-01", updated_at: "2026-08-01T00:00:00Z" }) }, // ---- Géocodage / LLM ---- diff --git a/apps/admin/test/e2e/tools.spec.js b/apps/admin/test/e2e/tools.spec.js index 4d596f2..aaf50a8 100644 --- a/apps/admin/test/e2e/tools.spec.js +++ b/apps/admin/test/e2e/tools.spec.js @@ -293,3 +293,39 @@ test.describe("Localisations — pagination et recherche", () => { await expect(page.locator("#l-table .err-box")).toContainText("Erreur liste localisations"); }); }); + +test.describe("Santé des services de fond", () => { + test("le bandeau affiche l'état des services", async ({ page }) => { + await page.goto("/#/pilotage"); + await expect(page.locator("#p-health")).toContainText("Services"); + await expect(page.locator("#p-health")).toContainText("2/2"); + }); + + // Ces services ne sont pas exposés par Traefik : sans ce battement de cœur, + // une dépendance injoignable reste invisible jusqu'à ce qu'on remarque + // l'absence de données nouvelles. + test("un service dégradé remonte en alerte actionnable", async ({ page }) => { + await patchJson(page, "**/admin/api/live", (body) => { + body.services = [ + { service: "crawler", ok: false, checks: { flaresolverr: false }, + error: "flaresolverr — HTTP 503", checked_at: new Date().toISOString(), stale: false }, + ]; + }); + await page.goto("/#/pilotage"); + const alert = page.locator(".alert").filter({ hasText: "crawler" }); + await expect(alert).toContainText("dégradé"); + await expect(alert).toContainText("HTTP 503"); + await expect(alert.getByRole("link", { name: "Voir le journal" })).toBeVisible(); + }); + + test("un service silencieux est signalé même si son dernier contrôle était bon", async ({ page }) => { + await patchJson(page, "**/admin/api/live", (body) => { + body.services = [ + { service: "geocoder", ok: true, checks: {}, error: null, + checked_at: "2026-08-01T00:00:00Z", stale: true }, + ]; + }); + await page.goto("/#/pilotage"); + await expect(page.locator(".alert").filter({ hasText: "geocoder" })).toContainText("silencieux"); + }); +}); diff --git a/apps/api/migrations/015_service_health.sql b/apps/api/migrations/015_service_health.sql new file mode 100644 index 0000000..a6ed181 --- /dev/null +++ b/apps/api/migrations/015_service_health.sql @@ -0,0 +1,24 @@ +-- 015_service_health.sql +-- +-- Supervision des services de fond (crawler, geocoder, groq-worker). +-- +-- Ils tournent sans être exposés par Traefik : aucun healthcheck HTTP ne peut +-- les atteindre. Jusqu'ici, un crawler dont FlareSolverr était injoignable, ou +-- un worker à court de quota Geoapify, restait silencieux — le seul symptôme +-- était l'absence de données nouvelles, qu'il fallait remarquer soi-même. +-- +-- Chaque service écrit désormais un battement de cœur avec le résultat de ses +-- propres sondes (base joignable, dépendance externe joignable, données qui +-- avancent). Une seule ligne par service, mise à jour en place : l'historique +-- détaillé est déjà dans crawl_log, on ne veut pas d'une table qui grossit. + +CREATE TABLE IF NOT EXISTS service_health ( + service TEXT PRIMARY KEY, -- 'crawler' | 'geocoder' | 'groq-worker' + ok BOOLEAN NOT NULL, + checks JSONB NOT NULL DEFAULT '{}', -- { "database": true, "flaresolverr": false, ... } + error TEXT, -- première sonde en échec, lisible + checked_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +COMMENT ON TABLE service_health IS + 'Battement de cœur des services de fond. Une ligne par service, écrasée à chaque contrôle. Un checked_at trop ancien signale un service arrêté ou bloqué — c''est aussi informatif qu''un ok=false.'; diff --git a/apps/api/src/adminRoutes.js b/apps/api/src/adminRoutes.js index 5acc752..c349f9e 100644 --- a/apps/api/src/adminRoutes.js +++ b/apps/api/src/adminRoutes.js @@ -894,7 +894,7 @@ export default async function adminRoutes(fastify, opts) { // ------------------------------------------------------------------ fastify.get("/admin/api/live", async (req, reply) => { try { - const [active_runs, pending_jobs, geo_queue, geo_processing, checkpoints] = await Promise.all([ + const [active_runs, pending_jobs, geo_queue, geo_processing, checkpoints, services] = await Promise.all([ pool.query(` SELECT id, run_type, countries, status, started_at, bands_seen, bands_new, bands_updated, bands_enriched, error, @@ -913,6 +913,15 @@ export default async function adminRoutes(fastify, opts) { WHERE bl.geocode_status = 'processing' LIMIT 3 `), pool.query(`SELECT key, value, updated_at FROM crawl_checkpoint ORDER BY key`), + // Battements de cœur des services de fond. `stale` est calculé en SQL : + // un service arrêté ne met plus sa ligne à jour, et c'est aussi + // révélateur qu'un ok=false — sans lui, un crawler mort passerait pour + // sain avec son dernier contrôle réussi. + pool.query(` + SELECT service, ok, checks, error, checked_at, + (checked_at < now() - interval '2 hours') AS stale + FROM service_health ORDER BY service + `).catch(() => ({ rows: [] })), ]); return { ok: true, @@ -921,6 +930,7 @@ export default async function adminRoutes(fastify, opts) { geo_queue: geo_queue.rows, geo_processing: geo_processing.rows, checkpoints: checkpoints.rows, + services: services.rows, }; } catch (err) { fastify.log.error(err); diff --git a/apps/api/test/adminRoutesCoverage.test.js b/apps/api/test/adminRoutesCoverage.test.js index d730c62..cdb714c 100644 --- a/apps/api/test/adminRoutesCoverage.test.js +++ b/apps/api/test/adminRoutesCoverage.test.js @@ -424,6 +424,41 @@ describe("GET /admin/api/live", () => { expect(body.checkpoints).toHaveLength(1); }); + it("expose la santé des services de fond", async () => { + const handlers = ([ + { match: "FROM service_health", result: rows( + { service: "crawler", ok: true, checks: { database: true }, error: null, stale: false }, + { service: "geocoder", ok: false, checks: { geoapify_joignable: false }, error: "HTTP 503", stale: false }) }, + { match: /./, result: rows() }, + ]); + const app = buildApp(handlers); + const body = (await app.inject({ method: "GET", url: "/admin/api/live", headers: auth })).json(); + expect(body.services).toHaveLength(2); + expect(body.services[1].error).toBe("HTTP 503"); + }); + + // Un service arrêté cesse simplement d'écrire : sans ce calcul, son dernier + // contrôle réussi le ferait passer pour sain indéfiniment. + it("marque comme silencieux un service qui n'écrit plus", async () => { + const handlers = anyPool(); + const app = buildApp(handlers); + await app.inject({ method: "GET", url: "/admin/api/live", headers: auth }); + expect(pool.find("FROM service_health").sql).toContain("checked_at < now() - interval '2 hours'"); + }); + + // La table est arrivée en migration 015 : une base plus ancienne ne doit pas + // faire tomber tout le tableau de bord. + it("dégrade proprement si service_health n'existe pas encore", async () => { + const handlers = ([ + { match: "FROM service_health", throws: new Error('relation "service_health" does not exist') }, + { match: /./, result: rows() }, + ]); + const app = buildApp(handlers); + const res = await app.inject({ method: "GET", url: "/admin/api/live", headers: auth }); + expect(res.statusCode).toBe(200); + expect(res.json().services).toEqual([]); + }); + it("limite l'aperçu des localisations en cours de traitement", async () => { const handlers = anyPool(); const app = buildApp(handlers); diff --git a/apps/api/test/sqlSyntax.test.js b/apps/api/test/sqlSyntax.test.js new file mode 100644 index 0000000..8113502 --- /dev/null +++ b/apps/api/test/sqlSyntax.test.js @@ -0,0 +1,245 @@ +import { describe, it, expect, beforeAll, afterAll } from "vitest"; +import { Parser } from "node-sql-parser"; + +process.env.ADMIN_JWT_SECRET = "secret-de-test-suffisamment-long-pour-passer-32"; + +const { signAdminSession, ADMIN_COOKIE_NAME } = await import("../src/adminAuth.js"); +const { rows } = await import("./helpers/fakePool.js"); +const { createAdminApp, createFullApp } = await import("./helpers/testApp.js"); + +/** + * Ferme l'angle mort SQL des tests à faux pool. + * + * Le faux pool vérifie la FORME du SQL (allowlists, paramètres liés) mais ne + * l'exécute jamais : une requête syntaxiquement invalide passait donc tous les + * tests et n'échouait qu'en production. C'est exactement ce qui s'était produit + * avec la reconstruction de crawler_pending. + * + * On parse ici chaque requête réellement émise avec la grammaire PostgreSQL. + * Ce n'est pas une exécution — les noms de colonnes et la sémantique ne sont + * pas validés — mais toute faute de syntaxe est attrapée sans conteneur. + */ +const parser = new Parser(); + +/** Requêtes que le parseur ne sait pas lire alors qu'elles sont valides. */ +const UNSUPPORTED = [ + // `x - ARRAY[...]` et `raw ? 'clé'` sont des opérateurs jsonb spécifiques + // à Postgres que node-sql-parser ne modélise pas. + /::jsonb\s*\)?\s*-\s*\$/, + /\?\s*'place_rank'/, + /jsonb_strip_nulls|jsonb_build_object/, + // make_interval(mins => $1) : appel par nom d'argument + /make_interval\(/, + // bands[1:5] : découpage de tableau propre à Postgres + /\[1:5\]/, + // Requêtes de contrôle de transaction + /^(BEGIN|COMMIT|ROLLBACK)$/, +]; + +function shouldParse(sql) { + return !UNSUPPORTED.some((re) => re.test(sql.trim())); +} + +/** + * Le parseur ne connaît pas les paramètres liés `$1` : on les remplace par des + * littéraux, ce qui préserve la structure de la requête. + * + * On utilise un NOMBRE et non une chaîne : `LIMIT $1` exige un littéral + * numérique côté grammaire. Le parseur ne fait pas de contrôle de types, un + * nombre passe donc aussi bien là où une chaîne serait attendue. + */ +function normalise(sql) { + return sql + .replace(/\$(\d+)::[a-z ]+\[\]/gi, "ARRAY[1]") + .replace(/\$(\d+)::[a-z ]+/gi, "1") + .replace(/\$\d+/g, "1"); +} + +function assertParses(sql, label) { + try { + parser.astify(normalise(sql), { database: "postgresql" }); + } catch (err) { + throw new Error(`SQL invalide (${label}) : ${err.message}\n---\n${sql.trim()}\n---`); + } +} + +let auth, adminApp, adminPool, pubApp, pubPool; + +beforeAll(async () => { + auth = { cookie: `${ADMIN_COOKIE_NAME}=${signAdminSession("nico")}` }; + ({ app: adminApp, pool: adminPool } = await createAdminApp()); + ({ app: pubApp, pool: pubPool } = await createFullApp()); +}); +afterAll(async () => { await adminApp.close(); await pubApp.close(); }); + +/** + * Rejoue une requête HTTP et parse tout le SQL qu'elle a produit. + * + * @param {import("fastify").FastifyInstance} app + * @param {any} pool + * @param {{ method?: string, url: string, payload?: any, handlers?: any[], unparsable?: boolean }} route + */ +async function checkRoute(app, pool, route) { + const { method = "GET", url, payload, handlers = [], unparsable = false } = route; + pool.reset([...handlers, { match: /./, result: rows({ n: 1, total: 1, id: 1, ma_id: 1, rowCount: 1 }) }]); + const res = await app.inject( + /** @type {any} */ ({ method, url, headers: auth, payload })); + expect([200, 400, 404, 409]).toContain(res.statusCode); + + const parsed = pool.calls.filter((c) => shouldParse(c.sql)); + if (!unparsable) { + // Sans cette garde, une route dont TOUTES les requêtes tomberaient dans la + // liste UNSUPPORTED passerait au vert sans rien avoir validé. + expect(parsed.length, `aucune requête parsable pour ${method} ${url}`).toBeGreaterThan(0); + } + for (const call of parsed) assertParses(call.sql, `${method} ${url}`); + return parsed.length; +} + +describe("SQL des routes admin", () => { + const ROUTES = [ + { url: "/admin/api/stats" }, + { url: "/admin/api/queue" }, + { url: "/admin/api/bands" }, + { url: "/admin/api/bands?q=mayhem&country=NO&genre=black&location_q=oslo&themes_q=war&enriched=true&has_lat=false&has_location=true&has_conflict=true&sort=name&dir=desc&page=2" }, + { url: "/admin/api/bands/1" }, + { url: "/admin/api/crawl-checkpoints" }, + { url: "/admin/api/logs" }, + { url: "/admin/api/logs?level=error&run_id=1&min_id=5&since=2026-01-01" }, + { url: "/admin/api/geocoding" }, + { url: "/admin/api/llm" }, + { url: "/admin/api/llm?model=x&only_null=1&q=oslo" }, + { url: "/admin/api/job-triggers" }, + { url: "/admin/api/live" }, + // Angle mort assumé : ces deux requêtes sont bâties sur jsonb_build_object, + // que node-sql-parser ne modélise pas. Leur syntaxe n'est donc PAS validée + // ici — c'est le prix de ne pas monter de Postgres. + { url: "/admin/api/activity", unparsable: true }, + { url: "/admin/api/activity?type=run&status=done", unparsable: true }, + { url: "/admin/api/locations" }, + { url: "/admin/api/locations?status=error,llm_needed&q=oslo&country=NO&provider=geoapify" }, + ]; + + it.each(ROUTES.map((r) => [r.url, r]))("%s produit du SQL syntaxiquement valide", async (_u, route) => { + await checkRoute(adminApp, adminPool, route); + }); + + const MUTATIONS = [ + { method: "POST", url: "/admin/api/job-triggers", payload: { job_type: "enrich" } }, + { method: "POST", url: "/admin/api/crawl-runs/1/cancel", payload: {} }, + { method: "POST", url: "/admin/api/job-triggers/1/cancel", payload: {} }, + { method: "POST", url: "/admin/api/locations/reset-errors", payload: {} }, + { method: "POST", url: "/admin/api/locations/reset-llm", payload: {} }, + { method: "POST", url: "/admin/api/locations/reset-all", payload: {} }, + { method: "POST", url: "/admin/api/locations/requeue-all", payload: { include_done: true } }, + { method: "POST", url: "/admin/api/locations/1/requeue", payload: {} }, + { method: "PATCH", url: "/admin/api/locations/1", payload: { lat: 1, lon: 2 } }, + ]; + + it.each(MUTATIONS.map((r) => [`${r.method} ${r.url}`, r]))( + "%s produit du SQL syntaxiquement valide", + async (_u, route) => { await checkRoute(adminApp, adminPool, route); } + ); + + // Le PATCH construit son SET dynamiquement : chaque combinaison de champs + // produit une requête différente, donc une occasion différente de se tromper. + it.each([ + ["un champ texte", { name: "X" }], + ["un champ numérique", { formed_year: 1991 }], + ["des coordonnées", { lat: 59.9, lon: 10.7 }], + ["tous les champs", { + name: "X", country: "NO", status: "Active", genre: "Black", + formed_year: 1991, themes: "War", location_text: "Oslo", lat: 1, lon: 2, + }], + ])("PATCH /admin/api/bands avec %s", async (_label, payload) => { + adminPool.reset([ + { match: "SELECT * FROM bands WHERE ma_id", result: rows({ ma_id: 1, lat: 1, lon: 2, location_text: "Oslo" }) }, + { match: /./, result: rows({ ma_id: 1, lat: 1, lon: 2, location_text: "Oslo" }) }, + ]); + const res = await adminApp.inject({ + method: "PATCH", url: "/admin/api/bands/1", headers: auth, payload, + }); + expect(res.statusCode).toBe(200); + for (const call of adminPool.calls.filter((c) => shouldParse(c.sql))) { + assertParses(call.sql, `PATCH bands (${Object.keys(payload).join(",")})`); + } + }); + + it.each(["name", "country", "formed_year", "lat"])( + "resolve-conflict sur %s produit du SQL valide", + async (field) => { + adminPool.reset([ + { match: "SELECT * FROM bands WHERE ma_id", result: rows({ ma_id: 1, crawler_pending: { [field]: "v" } }) }, + { match: /./, result: rows({ ma_id: 1 }) }, + ]); + await adminApp.inject({ + method: "POST", url: "/admin/api/bands/1/resolve-conflict", + headers: auth, payload: { field, action: "accept_crawler" }, + }); + for (const call of adminPool.calls.filter((c) => shouldParse(c.sql))) { + assertParses(call.sql, `resolve-conflict ${field}`); + } + } + ); +}); + +describe("SQL des routes publiques", () => { + const ROUTES = [ + "/api/db", + "/api/stats", + "/api/countries", + "/api/statuses", + "/api/facets", + "/api/bands", + "/api/bands?countries=NO,FR&status=Active&geocoded=1&only_black=1&q=mayhem&limit=10&offset=5", + "/api/band/1", + "/api/clusters?bbox=-5,40,10,55&zoom=5", + "/api/clusters?bbox=-5,40,10,55&zoom=14", + "/api/clusters?bbox=-5,40,10,55&zoom=5&countries=NO&status=Active&genre=black&year_min=1980&year_max=2000", + ]; + + // Même réserve que pour /admin/api/activity : l'agrégation en grille des + // clusters (zoom < 12) repose sur jsonb_build_object et array_agg, non + // modélisés par le parseur. Le mode « points bruts » (zoom >= 12), lui, l'est. + const UNVALIDATED = new Set([ + "/api/clusters?bbox=-5,40,10,55&zoom=5", + "/api/clusters?bbox=-5,40,10,55&zoom=5&countries=NO&status=Active&genre=black&year_min=1980&year_max=2000", + ]); + + it.each(ROUTES)("%s produit du SQL syntaxiquement valide", async (url) => { + pubPool.reset([{ match: /./, result: rows({ n: 1, total: 1, ma_id: 1 }) }]); + const res = await pubApp.inject({ method: "GET", url }); + expect(res.statusCode).toBe(200); + const parsed = pubPool.calls.filter((c) => shouldParse(c.sql)); + if (!UNVALIDATED.has(url)) expect(parsed.length).toBeGreaterThan(0); + for (const call of parsed) assertParses(call.sql, url); + }); + + it("/admin/import produit du SQL valide", async () => { + const app = pubApp; + pubPool.reset([{ match: /./, result: rows() }]); + await app.inject({ + method: "POST", url: "/admin/import", + headers: { authorization: "Bearer " + (process.env.BM_IMPORT_TOKEN || "") }, + payload: { bands: [{ ma_id: 1, name: "X" }] }, + }); + for (const call of pubPool.calls.filter((c) => shouldParse(c.sql))) { + assertParses(call.sql, "/admin/import"); + } + }); +}); + +describe("le détecteur lui-même", () => { + // Sans ça, une régression du parseur (ou un normalise() trop permissif) + // rendrait toute cette suite silencieusement inopérante. + it("rejette une requête syntaxiquement invalide", () => { + expect(() => assertParses("SELECT FROM WHERE ORDER", "témoin")).toThrow(/SQL invalide/); + expect(() => assertParses("UPDATE bands SET WHERE id = 1", "témoin")).toThrow(/SQL invalide/); + }); + + it("accepte une requête valide avec paramètres liés", () => { + expect(() => assertParses( + "SELECT a FROM bands WHERE id = $1 AND c = ANY($2::text[]) LIMIT $3 OFFSET $4", "témoin" + )).not.toThrow(); + }); +}); diff --git a/apps/crawler/src/health.py b/apps/crawler/src/health.py new file mode 100644 index 0000000..0947e81 --- /dev/null +++ b/apps/crawler/src/health.py @@ -0,0 +1,133 @@ +""" +Contrôle de santé périodique des services de fond. + +Ces services ne sont pas exposés par Traefik : aucune sonde HTTP externe ne peut +les atteindre. Sans ce module, un crawler dont FlareSolverr est injoignable ou +un worker à court de quota reste muet, et le seul symptôme est l'absence de +données nouvelles — qu'il faut remarquer soi-même. + +Chaque sonde renvoie (nom, ok, détail). Le résultat agrégé est écrit dans +`service_health` (une ligne par service, écrasée), lu par le dashboard admin. + +Aucune sonde ne peut interrompre le service : toute exception est convertie en +échec de sonde. Un contrôle de santé qui fait tomber ce qu'il surveille serait +pire que pas de contrôle du tout. +""" +import logging +import time + +log = logging.getLogger(__name__) + +# Intervalle minimal entre deux contrôles complets. +DEFAULT_INTERVAL_S = 3600 + +_last_run = 0.0 + + +def probe(name, fn): + """Exécute une sonde et convertit toute exception en échec.""" + try: + ok, detail = fn() + return name, bool(ok), (detail or "") + except Exception as e: # noqa: BLE001 - une sonde ne doit jamais propager + return name, False, f"{type(e).__name__}: {e}" + + +def database_probe(get_conn): + """La base répond-elle ?""" + def _run(): + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute("SELECT 1") + cur.fetchone() + return True, "connexion établie" + return _run + + +def http_probe(session_get, url, expect_below=500, timeout=10): + """Une dépendance HTTP répond-elle sans erreur serveur ?""" + def _run(): + r = session_get(url, timeout=timeout) + code = getattr(r, "status_code", 0) + return code < expect_below, f"HTTP {code}" + return _run + + +def freshness_probe(get_conn, sql, max_age_hours, label): + """Les données progressent-elles ? + + `sql` doit renvoyer un unique timestamp (le plus récent). Une base vide + (NULL) n'est pas un échec : c'est un état légitime au premier démarrage. + """ + def _run(): + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute(sql) + row = cur.fetchone() + latest = row[0] if row else None + if latest is None: + return True, f"{label}: aucune donnée pour l'instant" + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute( + "SELECT EXTRACT(EPOCH FROM (now() - %s)) / 3600.0", (latest,) + ) + age_h = float(cur.fetchone()[0]) + return age_h <= max_age_hours, f"{label}: {age_h:.1f} h (seuil {max_age_hours} h)" + return _run + + +def write_health(get_conn, service, results): + """Enregistre le résultat agrégé. N'échoue jamais bruyamment.""" + import json + + checks = {name: ok for name, ok, _ in results} + failed = [f"{name} — {detail}" for name, ok, detail in results if not ok] + ok = not failed + try: + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute( + """ + INSERT INTO service_health (service, ok, checks, error, checked_at) + VALUES (%s, %s, %s::jsonb, %s, now()) + ON CONFLICT (service) DO UPDATE SET + ok = EXCLUDED.ok, + checks = EXCLUDED.checks, + error = EXCLUDED.error, + checked_at = EXCLUDED.checked_at + """, + (service, ok, json.dumps(checks), failed[0] if failed else None), + ) + except Exception as e: # noqa: BLE001 + log.warning(f"[health] écriture impossible: {e}") + return ok, failed + + +def run_checks(get_conn, service, probes, log_event=None): + """Exécute toutes les sondes, enregistre le résultat, journalise les échecs.""" + results = [probe(name, fn) for name, fn in probes] + ok, failed = write_health(get_conn, service, results) + + if failed: + msg = f"[health:{service}] DÉGRADÉ — " + " | ".join(failed) + log.warning(msg) + if log_event: + try: + log_event("warning", msg) + except Exception: # noqa: BLE001, S110 - journaliser un échec de + # journalisation n'apporte rien et risquerait une récursion. + pass + else: + log.info(f"[health:{service}] toutes les sondes sont au vert") + return ok + + +def maybe_run(get_conn, service, probes, interval_s=DEFAULT_INTERVAL_S, log_event=None, force=False): + """Point d'entrée depuis une boucle de service : ne contrôle qu'une fois par intervalle.""" + global _last_run + now = time.monotonic() + if not force and _last_run and (now - _last_run) < interval_s: + return None + _last_run = now + return run_checks(get_conn, service, probes, log_event=log_event) diff --git a/apps/crawler/src/main.py b/apps/crawler/src/main.py index 5e612ee..a3ee113 100644 --- a/apps/crawler/src/main.py +++ b/apps/crawler/src/main.py @@ -32,6 +32,7 @@ from .config import ( ) from .db import claim_job_trigger, finish_job_trigger, get_checkpoint 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 @@ -117,9 +118,36 @@ def main(): while True: schedule.run_pending() _check_job_triggers(ma) + _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() diff --git a/apps/crawler/tests/test_health.py b/apps/crawler/tests/test_health.py new file mode 100644 index 0000000..3a17620 --- /dev/null +++ b/apps/crawler/tests/test_health.py @@ -0,0 +1,217 @@ +""" +Contrôles de santé des services de fond. + +Règle cardinale : une sonde ne doit JAMAIS interrompre le service qu'elle +surveille. Un contrôle de santé qui fait tomber le crawler serait pire que pas +de contrôle du tout — d'où le nombre de cas d'échec vérifiés ici. +""" + +import contextlib + +import pytest +from src import health +from src.health import ( + database_probe, + freshness_probe, + http_probe, + maybe_run, + probe, + run_checks, + write_health, +) + + +class FakeCursor: + def __init__(self, results=None, fail=False): + self._results = list(results or []) + self.queries = [] + self.fail = fail + + def execute(self, sql, params=None): + if self.fail: + raise RuntimeError("base injoignable") + self.queries.append((sql, params)) + + def fetchone(self): + return self._results.pop(0) if self._results else None + + def __enter__(self): + return self + + def __exit__(self, *a): + return False + + +def conn_factory(cursor): + @contextlib.contextmanager + def _get_conn(): + class Conn: + def cursor(self_inner): + return cursor + yield Conn() + return _get_conn + + +@pytest.fixture(autouse=True) +def _reset_interval(): + health._last_run = 0.0 + yield + health._last_run = 0.0 + + +class TestProbe: + def test_transmet_le_resultat(self): + assert probe("x", lambda: (True, "ok")) == ("x", True, "ok") + + def test_convertit_une_exception_en_echec(self): + def boom(): + raise ValueError("cassé") + name, ok, detail = probe("x", boom) + assert (name, ok) == ("x", False) + assert "ValueError" in detail and "cassé" in detail + + def test_normalise_un_detail_absent(self): + assert probe("x", lambda: (True, None))[2] == "" + + +class TestDatabaseProbe: + def test_ok_quand_la_base_repond(self): + cur = FakeCursor([(1,)]) + assert database_probe(conn_factory(cur))()[0] is True + # La sonde doit réellement interroger la base, pas se contenter + # d'ouvrir une connexion. + assert cur.queries and "SELECT 1" in cur.queries[0][0] + + def test_echoue_sans_lever_quand_la_base_est_injoignable(self): + cur = FakeCursor(fail=True) + name, ok, detail = probe("database", database_probe(conn_factory(cur))) + assert ok is False + assert "injoignable" in detail + + +class TestHttpProbe: + def test_ok_sur_200(self): + assert http_probe(lambda url, timeout: type("R", (), {"status_code": 200})(), "http://x")()[0] is True + + def test_echec_sur_500(self): + ok, detail = http_probe(lambda url, timeout: type("R", (), {"status_code": 503})(), "http://x")() + assert ok is False + assert "503" in detail + + def test_une_erreur_reseau_devient_un_echec_de_sonde(self): + def boom(url, timeout): + raise ConnectionError("nom d'hôte introuvable") + name, ok, detail = probe("dep", http_probe(boom, "http://x")) + assert ok is False + assert "ConnectionError" in detail + + +class TestFreshnessProbe: + def _probe(self, latest, age_h): + cur = FakeCursor([(latest,), (age_h,)]) + return freshness_probe(conn_factory(cur), "SELECT max(x) FROM t", 12, "test")() + + def test_ok_quand_les_donnees_sont_fraiches(self): + assert self._probe("2026-08-18T10:00:00Z", 3.0)[0] is True + + def test_echec_quand_les_donnees_sont_trop_anciennes(self): + ok, detail = self._probe("2026-08-01T10:00:00Z", 40.0) + assert ok is False + assert "40.0 h" in detail and "seuil 12 h" in detail + + def test_borne_exacte_acceptee(self): + assert self._probe("x", 12.0)[0] is True + + # Une base vide au premier démarrage n'est pas une panne. + def test_absence_de_donnees_nest_pas_un_echec(self): + cur = FakeCursor([(None,)]) + ok, detail = freshness_probe(conn_factory(cur), "SELECT max(x) FROM t", 12, "test")() + assert ok is True + assert "aucune donnée" in detail + + +class TestWriteHealth: + def test_agrege_les_sondes_et_enregistre(self): + cur = FakeCursor() + ok, failed = write_health(conn_factory(cur), "crawler", + [("database", True, "ok"), ("dep", True, "ok")]) + assert ok is True and failed == [] + sql, params = cur.queries[0] + assert "INSERT INTO service_health" in sql + assert "ON CONFLICT (service) DO UPDATE" in sql # une ligne par service + assert params[0] == "crawler" and params[1] is True + + def test_une_sonde_en_echec_rend_le_service_degrade(self): + cur = FakeCursor() + ok, failed = write_health(conn_factory(cur), "crawler", + [("database", True, "ok"), ("dep", False, "HTTP 503")]) + assert ok is False + assert failed == ["dep — HTTP 503"] + params = cur.queries[0][1] + assert params[1] is False + assert "dep" in params[3] # la première panne est lisible telle quelle + + def test_le_detail_par_sonde_est_conserve(self): + cur = FakeCursor() + write_health(conn_factory(cur), "crawler", [("a", True, ""), ("b", False, "")]) + import json + assert json.loads(cur.queries[0][1][2]) == {"a": True, "b": False} + + # Si la base est en panne, la sonde base l'a déjà signalé : échouer ici en + # plus ferait tomber le service pour rien. + def test_une_base_injoignable_ne_fait_pas_tomber_le_service(self): + cur = FakeCursor(fail=True) + ok, failed = write_health(conn_factory(cur), "crawler", [("database", False, "ko")]) + assert ok is False # pas d'exception propagée + + +class TestRunChecks: + def test_journalise_les_sondes_en_echec(self): + cur = FakeCursor() + messages = [] + run_checks(conn_factory(cur), "crawler", + [("dep", lambda: (False, "HTTP 503"))], + log_event=lambda lvl, msg: messages.append((lvl, msg))) + assert messages and messages[0][0] == "warning" + assert "DÉGRADÉ" in messages[0][1] and "dep" in messages[0][1] + + def test_ne_journalise_rien_quand_tout_va_bien(self): + cur = FakeCursor() + messages = [] + run_checks(conn_factory(cur), "crawler", [("dep", lambda: (True, "ok"))], + log_event=lambda lvl, msg: messages.append(msg)) + assert messages == [] + + def test_un_log_event_defaillant_ne_fait_pas_tomber_le_controle(self): + cur = FakeCursor() + def boom(lvl, msg): + raise RuntimeError("log cassé") + assert run_checks(conn_factory(cur), "crawler", + [("dep", lambda: (False, "ko"))], log_event=boom) is False + + +class TestMaybeRun: + def test_execute_au_premier_appel(self): + cur = FakeCursor() + assert maybe_run(conn_factory(cur), "crawler", [("d", lambda: (True, ""))]) is True + + # Appelée depuis une boucle qui tourne toutes les minutes, la fonction ne + # doit pas sonder les dépendances externes à chaque passage. + def test_ne_reexecute_pas_avant_lintervalle(self): + cur = FakeCursor() + probes = [("d", lambda: (True, ""))] + maybe_run(conn_factory(cur), "crawler", probes, interval_s=3600) + assert maybe_run(conn_factory(cur), "crawler", probes, interval_s=3600) is None + + def test_force_contourne_lintervalle(self): + cur = FakeCursor() + probes = [("d", lambda: (True, ""))] + maybe_run(conn_factory(cur), "crawler", probes, interval_s=3600) + assert maybe_run(conn_factory(cur), "crawler", probes, interval_s=3600, force=True) is True + + def test_reexecute_une_fois_lintervalle_ecoule(self): + cur = FakeCursor() + probes = [("d", lambda: (True, ""))] + maybe_run(conn_factory(cur), "crawler", probes, interval_s=3600) + health._last_run -= 4000 # simule le temps écoulé + assert maybe_run(conn_factory(cur), "crawler", probes, interval_s=3600) is True diff --git a/apps/geocoder/src/health.py b/apps/geocoder/src/health.py new file mode 100644 index 0000000..0947e81 --- /dev/null +++ b/apps/geocoder/src/health.py @@ -0,0 +1,133 @@ +""" +Contrôle de santé périodique des services de fond. + +Ces services ne sont pas exposés par Traefik : aucune sonde HTTP externe ne peut +les atteindre. Sans ce module, un crawler dont FlareSolverr est injoignable ou +un worker à court de quota reste muet, et le seul symptôme est l'absence de +données nouvelles — qu'il faut remarquer soi-même. + +Chaque sonde renvoie (nom, ok, détail). Le résultat agrégé est écrit dans +`service_health` (une ligne par service, écrasée), lu par le dashboard admin. + +Aucune sonde ne peut interrompre le service : toute exception est convertie en +échec de sonde. Un contrôle de santé qui fait tomber ce qu'il surveille serait +pire que pas de contrôle du tout. +""" +import logging +import time + +log = logging.getLogger(__name__) + +# Intervalle minimal entre deux contrôles complets. +DEFAULT_INTERVAL_S = 3600 + +_last_run = 0.0 + + +def probe(name, fn): + """Exécute une sonde et convertit toute exception en échec.""" + try: + ok, detail = fn() + return name, bool(ok), (detail or "") + except Exception as e: # noqa: BLE001 - une sonde ne doit jamais propager + return name, False, f"{type(e).__name__}: {e}" + + +def database_probe(get_conn): + """La base répond-elle ?""" + def _run(): + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute("SELECT 1") + cur.fetchone() + return True, "connexion établie" + return _run + + +def http_probe(session_get, url, expect_below=500, timeout=10): + """Une dépendance HTTP répond-elle sans erreur serveur ?""" + def _run(): + r = session_get(url, timeout=timeout) + code = getattr(r, "status_code", 0) + return code < expect_below, f"HTTP {code}" + return _run + + +def freshness_probe(get_conn, sql, max_age_hours, label): + """Les données progressent-elles ? + + `sql` doit renvoyer un unique timestamp (le plus récent). Une base vide + (NULL) n'est pas un échec : c'est un état légitime au premier démarrage. + """ + def _run(): + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute(sql) + row = cur.fetchone() + latest = row[0] if row else None + if latest is None: + return True, f"{label}: aucune donnée pour l'instant" + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute( + "SELECT EXTRACT(EPOCH FROM (now() - %s)) / 3600.0", (latest,) + ) + age_h = float(cur.fetchone()[0]) + return age_h <= max_age_hours, f"{label}: {age_h:.1f} h (seuil {max_age_hours} h)" + return _run + + +def write_health(get_conn, service, results): + """Enregistre le résultat agrégé. N'échoue jamais bruyamment.""" + import json + + checks = {name: ok for name, ok, _ in results} + failed = [f"{name} — {detail}" for name, ok, detail in results if not ok] + ok = not failed + try: + with get_conn() as conn: + with conn.cursor() as cur: + cur.execute( + """ + INSERT INTO service_health (service, ok, checks, error, checked_at) + VALUES (%s, %s, %s::jsonb, %s, now()) + ON CONFLICT (service) DO UPDATE SET + ok = EXCLUDED.ok, + checks = EXCLUDED.checks, + error = EXCLUDED.error, + checked_at = EXCLUDED.checked_at + """, + (service, ok, json.dumps(checks), failed[0] if failed else None), + ) + except Exception as e: # noqa: BLE001 + log.warning(f"[health] écriture impossible: {e}") + return ok, failed + + +def run_checks(get_conn, service, probes, log_event=None): + """Exécute toutes les sondes, enregistre le résultat, journalise les échecs.""" + results = [probe(name, fn) for name, fn in probes] + ok, failed = write_health(get_conn, service, results) + + if failed: + msg = f"[health:{service}] DÉGRADÉ — " + " | ".join(failed) + log.warning(msg) + if log_event: + try: + log_event("warning", msg) + except Exception: # noqa: BLE001, S110 - journaliser un échec de + # journalisation n'apporte rien et risquerait une récursion. + pass + else: + log.info(f"[health:{service}] toutes les sondes sont au vert") + return ok + + +def maybe_run(get_conn, service, probes, interval_s=DEFAULT_INTERVAL_S, log_event=None, force=False): + """Point d'entrée depuis une boucle de service : ne contrôle qu'une fois par intervalle.""" + global _last_run + now = time.monotonic() + if not force and _last_run and (now - _last_run) < interval_s: + return None + _last_run = now + return run_checks(get_conn, service, probes, log_event=log_event) diff --git a/apps/geocoder/src/worker.py b/apps/geocoder/src/worker.py index 8ca1214..9e9fb54 100644 --- a/apps/geocoder/src/worker.py +++ b/apps/geocoder/src/worker.py @@ -21,6 +21,7 @@ import time import psycopg2 import requests +from health import database_probe, freshness_probe, http_probe, maybe_run from parser import build_fallback_queries GEOAPIFY_API_KEY = os.environ.get("GEOAPIFY_API_KEY", "").strip() @@ -177,6 +178,35 @@ def main(): print(f"[worker] démarré — seuil confiance={MIN_CONFIDENCE}, max_tries={MAX_GEO_TRIES}") + # Le worker n'utilise pas de gestionnaire de contexte pour sa connexion : + # on en fournit un au module de santé, qui en attend un. + import contextlib + + @contextlib.contextmanager + def get_conn(): + yield conn + + def health_check(): + maybe_run(get_conn, "geocoder", [ + ("database", database_probe(get_conn)), + # Sans clé, le worker tourne à vide en boucle : à signaler tôt. + ("geoapify_cle", lambda: (bool(GEOAPIFY_API_KEY), "clé absente" if not GEOAPIFY_API_KEY else "clé présente")), + # Sonde de connectivité : une requête volontairement vide suffit à + # distinguer « API joignable » de « réseau coupé / quota épuisé ». + ("geoapify_joignable", http_probe( + lambda url, timeout: requests.get( + url, params={"text": "Oslo", "apiKey": GEOAPIFY_API_KEY, "limit": 1}, timeout=timeout), + GEOAPIFY_BASE)), + # Le géocodage doit avancer : au-delà de 6 h sans lieu résolu alors + # qu'il en reste en file, quelque chose bloque. + ("progression", freshness_probe( + get_conn, + "SELECT max(updated_at) FROM band_locations WHERE geocode_status IN ('done','country_only')", + 6, "dernier lieu résolu")), + ]) + + health_check() # un premier contrôle au démarrage + processed = 0 with conn.cursor() as cur: while processed < MAX_PER_RUN: @@ -196,6 +226,7 @@ def main(): row = cur.fetchone() if not row: print("[worker] rien à traiter, attente 60s") + health_check() time.sleep(60) continue diff --git a/package-lock.json b/package-lock.json index 27c5413..528ccae 100644 --- a/package-lock.json +++ b/package-lock.json @@ -21,6 +21,7 @@ "eslint": "^9.17.0", "globals": "^15.14.0", "jsdom": "^25.0.1", + "node-sql-parser": "^5.4.0", "typescript": "^5.7.2", "vitest": "^2.1.8" } @@ -2475,6 +2476,13 @@ "undici-types": "~6.21.0" } }, + "node_modules/@types/pegjs": { + "version": "0.10.6", + "resolved": "https://registry.npmjs.org/@types/pegjs/-/pegjs-0.10.6.tgz", + "integrity": "sha512-eLYXDbZWXh2uxf+w8sXS8d6KSoXTswfps6fvCUuVAGN8eRpfe7h9eSRydxiSJvo9Bf+GzifsDOr9TMQlmJdmkw==", + "dev": true, + "license": "MIT" + }, "node_modules/@types/wrap-ansi": { "version": "3.0.0", "resolved": "https://registry.npmjs.org/@types/wrap-ansi/-/wrap-ansi-3.0.0.tgz", @@ -2834,6 +2842,16 @@ "integrity": "sha512-V/Hy/X9Vt7f3BbPJEi8BdVFMByHi+jNXrYkW3huaybV/kQ0KJg0Y6PkEMbn+zeT+i+SiKZ/HMqJGIIt4LZDqNQ==", "license": "MIT" }, + "node_modules/big-integer": { + "version": "1.6.52", + "resolved": "https://registry.npmjs.org/big-integer/-/big-integer-1.6.52.tgz", + "integrity": "sha512-QxD8cf2eVqJOOz63z6JIN9BzvVs/dlySa5HGSBH5xtR8dPteIRQnBxxKqkNTiT6jbDTF6jAfrd4oMcND9RGbQg==", + "dev": true, + "license": "Unlicense", + "engines": { + "node": ">=0.6" + } + }, "node_modules/bm-api": { "resolved": "apps/api", "link": true @@ -4923,6 +4941,20 @@ "node": ">=18" } }, + "node_modules/node-sql-parser": { + "version": "5.4.0", + "resolved": "https://registry.npmjs.org/node-sql-parser/-/node-sql-parser-5.4.0.tgz", + "integrity": "sha512-jVe6Z61gPcPjCElPZ6j8llB3wnqGcuQzefim1ERsqIakxnEy5JlzV7XKdO1KmacRG5TKwPc4vJTgSRQ0LfkbFw==", + "dev": true, + "license": "Apache-2.0", + "dependencies": { + "@types/pegjs": "^0.10.0", + "big-integer": "^1.6.48" + }, + "engines": { + "node": ">=8" + } + }, "node_modules/npm-run-path": { "version": "6.0.0", "resolved": "https://registry.npmjs.org/npm-run-path/-/npm-run-path-6.0.0.tgz", diff --git a/package.json b/package.json index 31c8286..a3f512a 100644 --- a/package.json +++ b/package.json @@ -46,6 +46,7 @@ "eslint": "^9.17.0", "globals": "^15.14.0", "jsdom": "^25.0.1", + "node-sql-parser": "^5.4.0", "typescript": "^5.7.2", "vitest": "^2.1.8" }