feat: validation syntaxique du SQL + supervision des services de fond
Some checks are pending
CI / javascript (push) Waiting to run
CI / python (push) Waiting to run
CI / mutation (push) Waiting to run

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 <noreply@anthropic.com>
This commit is contained in:
Nicolas Fryder 2026-08-18 18:20:15 +02:00
parent c30656544c
commit 5237f3666d
16 changed files with 977 additions and 8 deletions

View file

@ -28,11 +28,11 @@ pip install -r requirements-dev.txt
| Commande | Ce que ça fait | Froid | Incrémental | | 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 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` | 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: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 test:e2e:install` | Récupère Chromium (clone neuf) | — | — |
| `npm run check:full` | `check` + e2e + audits + mutation complète | — | ~1 min | | `npm run check:full` | `check` + e2e + audits + mutation complète | — | ~1 min |
| `npm run check:clean` | Purge les caches ESLint / tsc | — | — | | `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 | | **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 | | **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) | | **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 ## 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 ne calcule pas les styles. Le dashboard admin, lui, est couvert par axe dans
un vrai navigateur. un vrai navigateur.
- `migrate.js` : s'exécute au démarrage du conteneur et appelle `process.exit`. - `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 - Crawler et workers de géocodage : le parseur, l'annulation et les contrôles
couverts, le reste est de l'I/O réseau et base. 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 - `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 chargé via `fs` + `vm` (comme le fait le navigateur), donc l'instrumentation
Stryker ne l'atteindrait pas et afficherait un score faussement parfait. Stryker ne l'atteindrait pas et afficherait un score faussement parfait.

View file

@ -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 = "") { function healthCard(label, value, sub, tone = "") {
return `<div class="health-card ${tone}"> return `<div class="health-card ${tone}">
<div class="h-label">${esc(label)}</div> <div class="h-label">${esc(label)}</div>
@ -382,6 +402,7 @@ async function loadPilotage() {
`${fmtNum(errors)} erreurs · ${fmtNum(llmNeeded)} LLM`, `${fmtNum(errors)} erreurs · ${fmtNum(llmNeeded)} LLM`,
errors + llmNeeded === 0 ? "ok" : "err"), errors + llmNeeded === 0 ? "ok" : "err"),
healthCard("Coût LLM", `$${cost.toFixed(2)}`, `${fmtNum(geo.llm_cache?.n || 0)} appels`), healthCard("Coût LLM", `$${cost.toFixed(2)}`, `${fmtNum(geo.llm_cache?.n || 0)} appels`),
servicesCard(live.services || []),
].join(""); ].join("");
// ---- Alertes actionnables ---- // ---- Alertes actionnables ----
@ -419,6 +440,16 @@ async function loadPilotage() {
actions: [{ label: "Crawl incrémental", job: "incremental" }], 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( const stuck = (live.active_runs || []).filter(
(r) => Date.now() - new Date(r.started_at) > 30 * 60000 && !r.cancel_requested); (r) => Date.now() - new Date(r.started_at) > 30 * 60000 && !r.cancel_requested);
if (stuck.length) { if (stuck.length) {

View file

@ -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 }) => { test("le bandeau de santé résume l'état du pipeline", async ({ page }) => {
const strip = page.locator("#p-health"); const strip = page.locator("#p-health");
await expect(strip.locator(".health-card")).toHaveCount(5); await expect(strip.locator(".health-card")).toHaveCount(6);
await expect(strip).toContainText("Géocodage"); for (const label of ["Groupes", "Géocodage", "En file", "Bloqués", "Coût LLM", "Services"]) {
await expect(strip).toContainText("Coût LLM"); await expect(strip).toContainText(label);
}
}); });
// Le cœur de la refonte : plus de « je constate ici, j'agis ailleurs ». // Le cœur de la refonte : plus de « je constate ici, j'agis ailleurs ».

View file

@ -105,6 +105,11 @@ export function seedHandlers(over = {}) {
{ match: "WHERE status = 'pending'", result: rows(...pendingJobs) }, { match: "WHERE status = 'pending'", result: rows(...pendingJobs) },
{ match: "geocode_status = 'processing'", result: rows() }, { match: "geocode_status = 'processing'", result: rows() },
{ match: "GROUP BY geocode_status", result: rows(...geoStatuses) }, { 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" }) }, { 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 ---- // ---- Géocodage / LLM ----

View file

@ -293,3 +293,39 @@ test.describe("Localisations — pagination et recherche", () => {
await expect(page.locator("#l-table .err-box")).toContainText("Erreur liste localisations"); 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");
});
});

View file

@ -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.';

View file

@ -894,7 +894,7 @@ export default async function adminRoutes(fastify, opts) {
// ------------------------------------------------------------------ // ------------------------------------------------------------------
fastify.get("/admin/api/live", async (req, reply) => { fastify.get("/admin/api/live", async (req, reply) => {
try { 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(` pool.query(`
SELECT id, run_type, countries, status, started_at, SELECT id, run_type, countries, status, started_at,
bands_seen, bands_new, bands_updated, bands_enriched, error, 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 WHERE bl.geocode_status = 'processing' LIMIT 3
`), `),
pool.query(`SELECT key, value, updated_at FROM crawl_checkpoint ORDER BY key`), 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 { return {
ok: true, ok: true,
@ -921,6 +930,7 @@ export default async function adminRoutes(fastify, opts) {
geo_queue: geo_queue.rows, geo_queue: geo_queue.rows,
geo_processing: geo_processing.rows, geo_processing: geo_processing.rows,
checkpoints: checkpoints.rows, checkpoints: checkpoints.rows,
services: services.rows,
}; };
} catch (err) { } catch (err) {
fastify.log.error(err); fastify.log.error(err);

View file

@ -424,6 +424,41 @@ describe("GET /admin/api/live", () => {
expect(body.checkpoints).toHaveLength(1); 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 () => { it("limite l'aperçu des localisations en cours de traitement", async () => {
const handlers = anyPool(); const handlers = anyPool();
const app = buildApp(handlers); const app = buildApp(handlers);

View file

@ -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 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();
});
});

133
apps/crawler/src/health.py Normal file
View file

@ -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)

View file

@ -32,6 +32,7 @@ from .config import (
) )
from .db import claim_job_trigger, finish_job_trigger, get_checkpoint from .db import claim_job_trigger, finish_job_trigger, get_checkpoint
from .flaresolverr import FlareSolverr 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 .jobs import run_enrich, run_full_crawl, run_incremental
from .ma_http import MASession from .ma_http import MASession
@ -117,9 +118,36 @@ def main():
while True: while True:
schedule.run_pending() schedule.run_pending()
_check_job_triggers(ma) _check_job_triggers(ma)
_check_health()
time.sleep(60) 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"): def _check_job_triggers(ma: "MASession"):
"""Consomme les job_triggers en attente créés depuis l'admin.""" """Consomme les job_triggers en attente créés depuis l'admin."""
result = claim_job_trigger() result = claim_job_trigger()

View file

@ -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

133
apps/geocoder/src/health.py Normal file
View file

@ -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)

View file

@ -21,6 +21,7 @@ import time
import psycopg2 import psycopg2
import requests import requests
from health import database_probe, freshness_probe, http_probe, maybe_run
from parser import build_fallback_queries from parser import build_fallback_queries
GEOAPIFY_API_KEY = os.environ.get("GEOAPIFY_API_KEY", "").strip() 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}") 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 processed = 0
with conn.cursor() as cur: with conn.cursor() as cur:
while processed < MAX_PER_RUN: while processed < MAX_PER_RUN:
@ -196,6 +226,7 @@ def main():
row = cur.fetchone() row = cur.fetchone()
if not row: if not row:
print("[worker] rien à traiter, attente 60s") print("[worker] rien à traiter, attente 60s")
health_check()
time.sleep(60) time.sleep(60)
continue continue

32
package-lock.json generated
View file

@ -21,6 +21,7 @@
"eslint": "^9.17.0", "eslint": "^9.17.0",
"globals": "^15.14.0", "globals": "^15.14.0",
"jsdom": "^25.0.1", "jsdom": "^25.0.1",
"node-sql-parser": "^5.4.0",
"typescript": "^5.7.2", "typescript": "^5.7.2",
"vitest": "^2.1.8" "vitest": "^2.1.8"
} }
@ -2475,6 +2476,13 @@
"undici-types": "~6.21.0" "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": { "node_modules/@types/wrap-ansi": {
"version": "3.0.0", "version": "3.0.0",
"resolved": "https://registry.npmjs.org/@types/wrap-ansi/-/wrap-ansi-3.0.0.tgz", "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==", "integrity": "sha512-V/Hy/X9Vt7f3BbPJEi8BdVFMByHi+jNXrYkW3huaybV/kQ0KJg0Y6PkEMbn+zeT+i+SiKZ/HMqJGIIt4LZDqNQ==",
"license": "MIT" "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": { "node_modules/bm-api": {
"resolved": "apps/api", "resolved": "apps/api",
"link": true "link": true
@ -4923,6 +4941,20 @@
"node": ">=18" "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": { "node_modules/npm-run-path": {
"version": "6.0.0", "version": "6.0.0",
"resolved": "https://registry.npmjs.org/npm-run-path/-/npm-run-path-6.0.0.tgz", "resolved": "https://registry.npmjs.org/npm-run-path/-/npm-run-path-6.0.0.tgz",

View file

@ -46,6 +46,7 @@
"eslint": "^9.17.0", "eslint": "^9.17.0",
"globals": "^15.14.0", "globals": "^15.14.0",
"jsdom": "^25.0.1", "jsdom": "^25.0.1",
"node-sql-parser": "^5.4.0",
"typescript": "^5.7.2", "typescript": "^5.7.2",
"vitest": "^2.1.8" "vitest": "^2.1.8"
} }