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"
}