Merge pull request #124 from ineszang/feat/dag-alertes

feat(etl): ordonnance la détection d'alertes et les recommandations par un DAG Airflow
This commit is contained in:
Johan LEROY
2026-09-21 13:31:16 +02:00
committed by GitHub
17 changed files with 343 additions and 49 deletions
+5
View File
@@ -43,6 +43,11 @@ AIRFLOW_ADMIN_USERNAME=admin
# comptes `app_user` d'EnerVision.
AIRFLOW_ADMIN_PASSWORD=change_me
AIRFLOW_ADMIN_EMAIL=admin@enervision.fr
# `APP_SECRET_KEY` du backend, que le DAG `alertes` lance en sous-processus. Distincte de
# celle de l'API : la détection ne signe aucun jeton, et Airflow exécute du code depuis son
# interface (cf. ADR 0008). Générer la vôtre :
# python -c "import secrets; print(secrets.token_urlsafe(48))"
AIRFLOW_APP_SECRET_KEY=change_me
# Stack complète derrière le reverse proxy (docker-compose.prod.yml).
# PUBLIC_HOST alimente l'origine CORS, le lien de réinitialisation et le certificat.
+20 -3
View File
@@ -4,8 +4,9 @@ name: Airflow
# (contrairement à backend.yml et ml.yml) : apache-airflow 2.10 ne supporte pas 3.14. Le 3.14 de
# ml/ ne vit que dans l'image Docker, dans son propre environnement (cf. etl/airflow/Dockerfile).
#
# Piège : l'image COPY les fichiers de dépendances et le code de ml/. Une modification de ml/
# peut donc casser sa construction, d'où ces chemins dans les déclencheurs.
# Piège : l'image COPY les fichiers de dépendances et le code de ml/ et de apps/backend/. Une
# modification de l'un ou de l'autre peut donc casser sa construction, d'où ces chemins dans
# les déclencheurs, alors même que ce workflow ne teste ni le modèle ni l'API.
on:
push:
@@ -14,6 +15,9 @@ on:
- "ml/pyproject.toml"
- "ml/uv.lock"
- "ml/enervision_ml/**"
- "apps/backend/pyproject.toml"
- "apps/backend/uv.lock"
- "apps/backend/app/**"
- ".github/workflows/airflow.yml"
pull_request:
paths:
@@ -21,6 +25,9 @@ on:
- "ml/pyproject.toml"
- "ml/uv.lock"
- "ml/enervision_ml/**"
- "apps/backend/pyproject.toml"
- "apps/backend/uv.lock"
- "apps/backend/app/**"
- ".github/workflows/airflow.yml"
permissions:
@@ -73,7 +80,7 @@ jobs:
- name: Récupère le dépôt
uses: actions/checkout@v4
- name: Construit l'image (contexte à la racine, elle COPY ml/)
- name: Construit l'image (contexte à la racine, elle COPY ml/ et apps/backend/)
run: docker build -f etl/airflow/Dockerfile -t enervision-airflow:ci .
# Vérifie ce qui ne casse qu'à l'exécution, pas à la construction : libgomp1 absent
@@ -82,3 +89,13 @@ jobs:
run: >
docker run --rm --network none enervision-airflow:ci
bash -c "cd /opt/ml && env -u VIRTUAL_ENV uv run --no-sync python -m enervision_ml.train --help"
# `--help` sort par argparse avant `get_settings()` : ni base ni secret requis, et
# l'import du module prouve que l'environnement /opt/backend est complet. Les deux
# commandes du DAG `alertes` sont couvertes, `app.cli` tirant tout FastAPI derrière lui.
- name: Vérifie que les deux commandes du DAG alertes s'importent sans réseau
run: >
docker run --rm --network none enervision-airflow:ci
bash -c "cd /opt/backend
&& env -u VIRTUAL_ENV uv run --no-sync python -m app.detection.internal_alerts --help
&& env -u VIRTUAL_ENV uv run --no-sync python -m app.cli generate-recommendations --help"
+4 -1
View File
@@ -18,7 +18,7 @@ endif
dev dev-backend dev-frontend \
lint format typecheck test test-cov test-integration check \
openapi docker-build db-up db-down db-reset db-logs db-psql migrate bootstrap-admin \
ml-lint ml-typecheck ml-test ml-check ml-train ml-score recommendations \
ml-lint ml-typecheck ml-test ml-check ml-train ml-score detect-alerts recommendations \
airflow-lint airflow-test airflow-check airflow-up airflow-down airflow-logs \
tls-selfsigned tls-acme tls-renew stack-up stack-down stack-logs
@@ -94,6 +94,9 @@ ml-train: ## Entraine le modele LightGBM. CSV=chemin optionnel, sinon lit ML_DAT
ml-score: ## Score le prochain pas horaire et l'ecrit dans `prediction`. CSV=chemin optionnel
cd $(ML) && uv run python -m enervision_ml.score $(if $(CSV),--csv $(CSV),)
detect-alerts: ## Détecte les alertes internes depuis les lectures en base. SITE= et NOW= optionnels
cd $(BACKEND) && uv run python -m app.detection.internal_alerts $(if $(SITE),--site-id $(SITE),) $(if $(NOW),--now $(NOW),)
recommendations: ## Genere les recommandations depuis les alertes en base. SITE=identifiant optionnel
cd $(BACKEND) && uv run python -m app.cli generate-recommendations $(if $(SITE),--site-id $(SITE),)
+2 -2
View File
@@ -21,7 +21,7 @@ Ce que la documentation apporte à chacun : [docs/architecture/00-vue-ensemble.m
| Backend | FastAPI, Python 3.14 | `apps/backend` | Initialise |
| Frontend | Angular 22, Node 24 LTS | `apps/frontend` | Tableau de bord |
| Base | PostgreSQL 17 + TimescaleDB | `db` | Initialise |
| ETL | Apache Airflow | `etl/airflow` | A initialiser |
| ETL | Apache Airflow | `etl/airflow` | Trois DAGs |
| Infra | Terraform (k3s single-node) | `infra/terraform` | Initialise |
| Reverse proxy | Nginx, TLS | `infra/proxy` | En place |
| CI/CD | GitHub Actions | `.github/workflows` | Backend en place |
@@ -48,7 +48,7 @@ L'etat detaille de chaque brique et les vues d'architecture sont dans
│ ├── migrations/ Migrations SQL versionnees
│ └── seeds/ Jeux de donnees de reference
├── etl/airflow/
│ ├── dags/ DAGs d'ingestion et d'agregation
│ ├── dags/ DAGs d'orchestration (pipeline ML, alertes)
│ ├── plugins/ Operateurs et hooks maison
│ ├── include/ Requetes SQL et ressources des DAGs
│ └── tests/ Tests d'integrite des DAGs
+7
View File
@@ -27,6 +27,12 @@ x-airflow-common: &airflow-common
# memes identifiants que le backend en attendant.
ML_DATABASE_URL: postgresql+psycopg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
MLFLOW_TRACKING_URI: sqlite:////opt/ml/state/mlflow.db
# Le DAG `alertes` lance le backend en sous-processus : il lit `DATABASE_URL`, en
# dialecte asyncpg, là où le pipeline ML lit `ML_DATABASE_URL`.
DATABASE_URL: postgresql+asyncpg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
# Clé distincte de celle de l'API : la détection ne signe ni ne vérifie aucun jeton, et
# Airflow permet d'exécuter du code depuis son interface (cf. ADR 0008).
APP_SECRET_KEY: ${AIRFLOW_APP_SECRET_KEY:-}
volumes:
- ./etl/airflow/dags:/opt/airflow/dags
- ./etl/airflow/plugins:/opt/airflow/plugins
@@ -129,6 +135,7 @@ services:
set -euo pipefail
: "$${AIRFLOW__CORE__FERNET_KEY:?AIRFLOW_FERNET_KEY manquant dans .env}"
: "$${AIRFLOW__WEBSERVER__SECRET_KEY:?AIRFLOW_WEBSERVER_SECRET_KEY manquant dans .env}"
: "$${APP_SECRET_KEY:?AIRFLOW_APP_SECRET_KEY manquant dans .env}"
exec airflow version
airflow-webserver:
+1
View File
@@ -14,3 +14,4 @@
| [0005](adr/0005-modele-prediction-lightgbm.md) | LightGBM pour la prédiction de consommation, un modèle global |
| [0006](adr/0006-moteur-de-regles-dans-le-backend.md) | Le moteur de règles de recommandation vit dans le backend, pas dans `ml/` |
| [0007](adr/0007-terminaison-tls-et-reverse-proxy-nginx.md) | Terminaison TLS par un reverse proxy Nginx, en Docker Compose |
| [0008](adr/0008-airflow-execute-le-code-du-backend.md) | Airflow exécute le code du backend en sous-processus, dans son propre environnement |
@@ -0,0 +1,79 @@
# 0008 - Airflow exécute le code du backend en sous-processus
- Statut : accepté
- Date : 2026-09-21
## Contexte
L'issue #116 demande un DAG d'alertes. Ce qu'il a à ordonnancer existe déjà et n'est pas à
réécrire : `AlertService.detect()` et ses cinq règles (#104), puis le moteur de recommandations
(#38). Les deux vivent dans `apps/backend/app/`, et
l'[ADR 0006](0006-moteur-de-regles-dans-le-backend.md) a précisément décidé qu'ils y restent parce
qu'ils s'appuient sur les repositories ORM de l'API plutôt que sur du SQL brut. Les deux
sont décrits par la documentation comme « lancés à la main ».
L'image Airflow livrée par #115 ne porte que `ml/`, dans un environnement `uv` distinct
(`/opt/ml/.venv`, Python 3.14) de celui d'Airflow lui-même (Python 3.12, contraint par
apache-airflow 2.10). Les DAGs `ml_train` et `ml_score` shellent vers cet environnement. Rien
d'équivalent n'existe pour `apps/backend` : un `BashOperator` sur
`python -m app.detection.internal_alerts` échouerait en `ModuleNotFoundError`.
## Décision
**L'image Airflow porte un troisième environnement, `/opt/backend/.venv`**, construit depuis le
`pyproject.toml`, le `uv.lock` et le paquet `app/` du backend. Le DAG `alertes` shelle vers lui
exactement comme `ml_score` shelle vers `/opt/ml/.venv`.
Trois raisons :
- **Le patron existe et vient d'être revu.** #115 a posé `BashOperator` + `uv run --no-sync` +
`env -u VIRTUAL_ENV`, avec les tests d'intégrité qui le verrouillent. Introduire une seconde
forme d'appel dans le même dossier `dags/` coûterait plus cher à lire qu'un second environnement
dans le même `Dockerfile`.
- **Aucune surface réseau n'est ajoutée.** La détection n'a pas de route HTTP, contrairement à la
génération de recommandations (`POST /recommendations/generate`, rôle `admin`). En créer une pour
qu'Airflow l'appelle donnerait à l'ordonnanceur un compte administrateur de l'API, en plus des
identifiants PostgreSQL complets qu'il détient déjà, et ferait dépendre la production d'alertes
de la disponibilité du conteneur `backend`.
- **La logique reste où l'ADR 0006 l'a mise.** Le DAG n'apprend rien du domaine : ni les seuils, ni
les cinq règles, ni les clés d'idempotence. Il ne sait que l'heure à laquelle appeler.
## Conséquences
- **Airflow reçoit une `APP_SECRET_KEY` délibérément distincte de celle de l'API.** La
configuration du backend refuse de se construire sans elle (`app/core/config.py`), et
`internal_alerts.main()` appelle `get_settings()` avant toute requête pour échouer tôt. Mais la
détection ne signe ni ne vérifie aucun jeton, et Airflow permet d'exécuter du code arbitraire
depuis son interface : un Airflow compromis ne doit pas livrer la clé de signature des JWT. D'où
`AIRFLOW_APP_SECRET_KEY`, avec sa propre garde dans `airflow-init`.
- **`DATABASE_URL`, en dialecte asyncpg, rejoint `ML_DATABASE_URL`** dans l'environnement du
conteneur. Le cantonnement des rôles PostgreSQL reste la dette de
l'[ADR 0003](0003-autorisation-rbac-a-trois-roles.md), et cette décision l'alourdit d'un
consommateur de plus.
- **La CI Airflow se déclenche sur les changements du backend.** L'image le `COPY` : sans
`apps/backend/app/**`, `pyproject.toml` et `uv.lock` dans les déclencheurs du workflow, une
dépendance modifiée casserait la construction sans que rien ne le signale avant le déploiement.
En contrepartie, l'image grossit de ce que pèsent SQLAlchemy, asyncpg et pandas.
- **Aucune variable ne départage les deux environnements, et c'est voulu.** `uv` place par défaut
le venv d'un projet dans `<projet>/.venv` : `cd /opt/ml` ou `cd /opt/backend` suffit à choisir le
bon. L'image ne pose donc plus de `UV_PROJECT_ENVIRONMENT` global, hérité de #115 : il vaudrait
pour les deux projets, et `uv run` dans l'un résoudrait le venv de l'autre. Le symptôme n'est pas
une construction ratée mais un `ModuleNotFoundError` à la première tâche, d'où la vérification
d'import sans réseau que la CI fait maintenant sur chacun des deux.
- Airflow lui-même reste étranger au domaine : ni LightGBM, ni SQLAlchemy, ni FastAPI n'entrent
dans son interpréteur. C'est la propriété que #115 avait établie, et elle tient toujours.
## Alternatives écartées
- **Route HTTP `POST /alerts/detect` réservée `admin`, appelée par le DAG.** L'image ne bougeait
pas, mais Airflow détenait alors un compte administrateur de l'API, la détection devenait
tributaire du conteneur `backend`, et l'API gagnait une route d'écriture dont aucun client
humain n'a l'usage. À rouvrir si un jour un tiers doit déclencher la détection.
- **`DockerOperator` lançant l'image du backend.** Demande la socket Docker de l'hôte dans le
conteneur Airflow, c'est-à-dire un équivalent root sur la machine, pour un service qui permet
déjà d'exécuter du code depuis son interface. Le provider n'est d'ailleurs pas installé.
- **Réécrire les cinq règles en SQL dans le DAG.** Contredit frontalement l'ADR 0006, duplique le
domaine, et fait diverger les deux copies au premier changement de seuil.
- **Monter `apps/backend` en volume plutôt que le copier.** L'environnement ne serait plus figé à
la construction, `uv` resynchroniserait au premier lancement, et la CI ne prouverait plus rien
de ce qui tourne réellement.
+5 -4
View File
@@ -70,9 +70,10 @@ Le lien `front -.-> api` reste en pointillé : le frontend appelle bien une API,
intercepteur répond à sa place tant que les endpoints n'existent pas. Voir
[30-frontend.md](30-frontend.md).
Le lien `airflow --> db` est maintenant en trait plein : deux DAGs orchestrent l'entraînement et
le scoring du modèle ML (issue #115), cf. plus bas et [20-backend.md](20-backend.md). Le reste du
périmètre Airflow envisagé (ingestion, issues #15/#16) reste en pointillé, non construit.
Le lien `airflow --> db` est maintenant en trait plein : trois DAGs tournent, deux pour
l'entraînement et le scoring du modèle ML (issue #115), un pour la détection d'alertes et la
génération des recommandations (issue #116), cf. plus bas et [20-backend.md](20-backend.md). Le
reste du périmètre Airflow envisagé (ingestion, issues #15/#16) reste en pointillé, non construit.
Le lien `prom -.-> api` de même : l'API expose bien `/metrics` au format Prometheus, mais aucun
collecteur ne vient le lire.
@@ -87,7 +88,7 @@ collecteur ne vient le lire.
| ML | LightGBM, MLflow | `ml` | `En cours` | Pipeline d'entraînement et de scoring (`enervision_ml.train`/`.score`, features par lags/moyennes glissantes partagées entre les deux, baseline de persistance saisonnière, suivi MLflow local), exposé en lecture via `GET /predictions`, orchestré par Airflow (`ml_train`/`ml_score`). Voir [ADR 0005](../adr/0005-modele-prediction-lightgbm.md) et [ML-START.md](../../ML-START.md). Surveillance de dérive (EC06, #44/#45) pas encore construite |
| Infra | Docker Compose, Nginx, Terraform, k3s single-node | `infra`, `docker-compose.prod.yml` | `En cours` | Reverse proxy et overlay de déploiement écrits et validés, jamais lancés sur le serveur ([ADR 0007](../adr/0007-terminaison-tls-et-reverse-proxy-nginx.md)). Module d'installation k3s jamais appliqué, aucune ressource Kubernetes déclarée |
| Monitoring | Prometheus, Grafana, Alertmanager | `monitoring` | `Cible` | Rien, hors le `/metrics` exposé par l'API |
| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver + scheduler (LocalExecutor) tournent via docker-compose, base de métadonnées Postgres dédiée. Deux DAGs (`ml_train` manuel, `ml_score` `@hourly`) orchestrent le pipeline ML existant en sous-processus `uv run` (issue #115). L'ingestion (issues #15/#16) n'a pas encore de DAG |
| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver + scheduler (LocalExecutor) tournent via docker-compose, base de métadonnées Postgres dédiée. Trois DAGs en sous-processus `uv run` : `ml_train` manuel et `ml_score` `@hourly` pour le pipeline ML (issue #115), `alertes` à `15 * * * *` pour la détection et les recommandations (issue #116, [ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md)). L'ingestion (issues #15/#16) n'a pas encore de DAG |
| CI/CD | GitHub Actions | `.github/workflows` | `Cible` | Rien |
## Flux bout en bout
+42 -9
View File
@@ -51,7 +51,7 @@ Trois pièges sont documentés en tête du `docker-compose.yml`, ils ne se devin
- `LocalExecutor` exécute les tâches comme sous-processus du **scheduler**, jamais du webserver :
c'est le scheduler qui a besoin du volume `airflow_ml_state` (modèle, magasin MLflow).
### Airflow (`ml_train`/`ml_score`, issue #115)
### Airflow (issues #115 et #116)
Trois services, `docker compose profiles` non utilisés (démarrage explicite via `make
airflow-up`, pas dans `make dev`) :
@@ -60,13 +60,37 @@ airflow-up`, pas dans `make dev`) :
|---|---|---|
| `airflow-init` | Migre la base de métadonnées, crée le compte admin | Conteneur jetable (`restart: "no"`), ne redémarre jamais. `webserver`/`scheduler` attendent qu'il se termine avec succès |
| `airflow-webserver` | UI, port `8080` | `LocalExecutor` : n'exécute aucune tâche lui-même |
| `airflow-scheduler` | Planifie et **exécute** les tâches (`LocalExecutor`) | Les DAGs y tournent en sous-processus (`uv run --frozen --no-dev python -m enervision_ml...`), c'est lui qui a besoin du volume `airflow_ml_state` |
| `airflow-scheduler` | Planifie et **exécute** les tâches (`LocalExecutor`) | Les DAGs y tournent en sous-processus (`uv run --no-sync python -m ...`), c'est lui qui a besoin du volume `airflow_ml_state` |
Construits depuis `etl/airflow/Dockerfile`, contexte `.` (racine du repo, pas `etl/airflow/`) :
l'image doit pouvoir `COPY` `ml/pyproject.toml`/`ml/uv.lock`/`ml/enervision_ml` pour se
synchroniser un second environnement Python **3.14** (`/opt/ml/.venv`, `uv sync --locked` à la
construction), distinct du Python 3.12 qui fait tourner Airflow lui-même. Les DAGs shellent vers
ce venv plutôt que d'importer LightGBM/MLflow dans le process Airflow.
l'image doit pouvoir `COPY` les sources de `ml/` **et** de `apps/backend/` pour se synchroniser
deux environnements Python **3.14** (`/opt/ml/.venv` et `/opt/backend/.venv`, `uv sync --locked` à
la construction), distincts du Python 3.12 qui fait tourner Airflow lui-même. Les DAGs shellent
vers ces venvs plutôt que d'importer LightGBM, MLflow ou SQLAlchemy dans le process Airflow.
Le choix et ses contreparties sont dans
l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md).
| DAG | Planification | Ce qu'il lance, et où |
|---|---|---|
| `ml_train` | manuelle | `enervision_ml.train`, dans `/opt/ml/.venv` |
| `ml_score` | `0 * * * *` | `enervision_ml.score`, dans `/opt/ml/.venv` |
| `alertes` | `15 * * * *` | `app.detection.internal_alerts` puis `app.cli generate-recommendations`, dans `/opt/backend/.venv` |
**Pourquoi `alertes` tourne à la quinzième minute.** Sa règle `anomaly` compare une lecture à la
`prediction` du même instant, que `ml_score` écrit à l'heure pile. Le décalage laisse le scoring
finir. Aucune dépendance n'est déclarée entre les deux DAGs pour autant, ni `ExternalTaskSensor` ni
tâche greffée : quatre règles de détection sur cinq ne touchent pas au modèle, et un modèle jamais
entraîné ne doit pas priver le parc de ses alertes. Le décalage est donc une convention et non une
garantie : le plafond de `ml_score` est de 30 minutes, et un scoring qui déborde de `:15` prive
`anomaly` de la `prediction` de l'heure, qu'elle ne retrouvera au passage suivant que si sa fenêtre
la couvre encore. Les quatre autres règles ne s'en aperçoivent pas.
Ses deux tâches s'enchaînent en revanche (`recommendation.alert_id` est une clé étrangère `NOT
NULL`), et toutes deux sont rejouables sans risque : l'idempotence est portée par la base,
`uq_alert_source_reference` et `uq_recommendation_alert_rule`. Chacune a 2 tentatives, 2 minutes
d'attente entre elles et un plafond de 5 minutes **par tentative** : au pire, reprises comprises,
l'enchaînement occupe 38 minutes, ce qui le garde sous le pas horaire qu'un `max_active_runs=1`
rend contraignant.
`airflow-init` s'appuie sur l'entrypoint de l'image (`_AIRFLOW_DB_MIGRATE`,
`_AIRFLOW_WWW_USER_*`) plutôt que sur un script maison : l'entrypoint porte le code de sortie, une
@@ -77,7 +101,15 @@ passe par l'environnement, jamais par `argv` (ni `ps`, ni `docker compose config
Les variables `AIRFLOW_*` ne sont volontairement pas en `${VAR:?}` : Compose interpole le fichier
entier avant de filtrer les services, une variable requise manquante casserait `make db-up`,
`make dev`... pour tout poste dont le `.env` est antérieur. Elles valent `${VAR:-}` et c'est
`airflow-init` qui refuse de démarrer (clé Fernet, clé Flask ou mot de passe vides).
`airflow-init` qui refuse de démarrer (clé Fernet, clé Flask, mot de passe ou
`AIRFLOW_APP_SECRET_KEY` vides).
Le conteneur reçoit deux variables du backend en plus de `ML_DATABASE_URL` : `DATABASE_URL`, en
dialecte asyncpg, et `APP_SECRET_KEY`, alimentée par `AIRFLOW_APP_SECRET_KEY`. Cette dernière est
**délibérément différente** de celle de l'API. La configuration du backend refuse de se construire
sans clé, mais la détection ne signe ni ne vérifie aucun jeton : un Airflow compromis, qui permet
déjà d'exécuter du code depuis son interface, ne doit pas livrer par-dessus la clé de signature
des JWT.
**Pourquoi `ml_train` est manuel.** Réentraîner est coûteux et sa cadence n'est pas une décision
prise. Surtout, `train.py` écrase le modèle sans comparer ses métriques à celles de l'ancien : un
@@ -86,8 +118,9 @@ déclenchement reste humain. `ml_score`, lui, est planifié à l'heure, avec `ma
(pas deux scorings simultanés dans `prediction`), 2 tentatives et un plafond de 30 minutes.
CI : `.github/workflows/airflow.yml` (Python 3.12 via `etl/airflow/.python-version`) lance lint et
tests d'intégrité des DAGs, et construit l'image (elle `COPY` `ml/`, une modification de `ml/`
peut donc la casser) avant de vérifier que le pipeline s'y importe sans réseau.
tests d'intégrité des DAGs, et construit l'image (elle `COPY` `ml/` et `apps/backend/`, une
modification de l'un ou de l'autre peut donc la casser, d'où leurs chemins dans les déclencheurs)
avant de vérifier que les deux environnements s'y importent sans réseau.
Piège à connaître : sur un volume `pgdata` déjà peuplé (poste de dev existant plutôt que premier
`make db-up`), `db/init/120-airflow-database.sql` ne se rejoue pas (PostgreSQL n'exécute
+13 -6
View File
@@ -256,12 +256,19 @@ auraient pu comparer des lectures/choisir une prévision au hasard. `_detect_spi
explicitement les paires de lectures qui partagent le même horodatage (deux `source` pour un seul
instant réel, pas une variation).
La détection est un script lancé à la main, pas encore ordonnancé par Airflow (contrairement à
`enervision_ml.score`, orchestré par le DAG `ml_score` depuis l'issue #115) : `uv run python -m app.detection.internal_alerts [--site-id ...] [--now ...]`, dans
`apps/backend` puisque les règles s'appuient sur les repositories ORM de l'API plutôt que sur une
connexion SQL directe (contrairement à `app/etl/historical_import.py`). Cette issue (#104)
débloquait #38 (moteur de règles pour recommandations), dont la FK `alert_id` `NOT NULL` n'avait
jusqu'ici rien à référencer côté `source="enervision"`.
La détection s'exécute dans `apps/backend`, puisque les règles s'appuient sur les repositories ORM
de l'API plutôt que sur une connexion SQL directe (contrairement à
`app/etl/historical_import.py`) : `uv run python -m app.detection.internal_alerts [--site-id ...]
[--now ...]`, ou `make detect-alerts`. Cette issue (#104) débloquait #38 (moteur de règles pour
recommandations), dont la FK `alert_id` `NOT NULL` n'avait jusqu'ici rien à référencer côté
`source="enervision"`.
Depuis l'issue #116, le lancement n'est plus manuel : le DAG Airflow `alertes` enchaîne cette
détection et la génération des recommandations, toutes les heures à la quinzième minute. Airflow
exécute le code du backend en sous-processus, dans son propre environnement, ce que décide
l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md) ; le détail de l'ordonnancement est
dans [10-infra.md](10-infra.md). La ligne de commande reste le moyen de rejouer une fenêtre
passée, ce que `--now` permet et que le DAG ne fait pas.
### `/health/ready`
+8 -6
View File
@@ -14,8 +14,10 @@ décrivent les éléments prévus mais pas encore réalisés.
L'ingestion des **mesures** est implémentée pour les deux sources du MVP, le dataset CSV/JSON et
l'API Mock. Celle des **alertes** de l'API Mock, `/alerts`, reste à faire : voir
l'[ADR 0006](../adr/0006-moteur-de-regles-dans-le-backend.md). L'orchestration Airflow, les
agrégats continus, la compression et la rétention restent des cibles.
l'[ADR 0006](../adr/0006-moteur-de-regles-dans-le-backend.md). Les alertes `source='enervision'`,
elles, sont produites par la détection interne, désormais ordonnancée par le DAG Airflow `alertes`
(issue #116). L'orchestration de l'ingestion, les agrégats continus, la compression et la
rétention restent des cibles.
## Trois emplacements, trois rôles
@@ -290,10 +292,10 @@ Les anomalies historiques décrites dans les JSON sont conservées dans `dataset
Elles servent à l'analyse des données et ne sont pas considérées comme des alertes actuelles.
Les lignes de `recommendation` sont écrites par le moteur de règles du backend
(`app/services/recommendation_rules.py`), déclenché par `POST /api/v1/recommendations/generate`
ou par `make recommendations`, à partir des alertes déjà en base. Le couple
`(alert_id, rule_reference)` est unique : rejouer le moteur sur les mêmes alertes n'ajoute aucune
ligne.
(`app/services/recommendation_rules.py`), déclenché par `POST /api/v1/recommendations/generate`,
par `make recommendations`, ou par la seconde tâche du DAG `alertes`, à partir des alertes déjà en
base. Le couple `(alert_id, rule_reference)` est unique : rejouer le moteur sur les mêmes alertes
n'ajoute aucune ligne.
### Relations entre les tables
+5 -2
View File
@@ -17,8 +17,11 @@ contredisent, c'est l'ADR qui fait foi et la vue qui est en retard.
| [40-data.md](40-data.md) | Frontières `db/` et `alembic/`, cycle de vie d'une mesure, modèle |
L'observabilité et la CI/CD n'ont pas de document propre : ce sont des sections des documents
ci-dessus, tant que `monitoring/` et `etl/airflow/` ne contiennent que des `.gitkeep`. Elles en
sortiront le jour où elles auront de la matière. Un fichier vide de plus n'aide personne.
ci-dessus, tant que `monitoring/` ne contient que des `.gitkeep`. Elles en sortiront le jour où
elles auront de la matière. Un fichier vide de plus n'aide personne.
L'orchestration Airflow, elle, en a depuis les issues #115 et #116 : trois DAGs, leur image et
leurs contraintes sont décrits dans [10-infra.md](10-infra.md).
La sécurité applicative, elle, a désormais de la matière : la vue consolidée reste dans
[00-vue-ensemble.md](00-vue-ensemble.md), le détail dans [20-backend.md](20-backend.md), la
+1 -1
View File
@@ -61,7 +61,7 @@ le certificat servi reste auto-signé. Le chemin ACME est livré et documenté,
| **API10 Unsafe Consumption of APIs** | **partiel, et spécifique à ce projet** | L'API Mock de l'école n'a aucune authentification, tourne en HTTP clair sur le réseau de l'école, et expose un endpoint mutatif à quiconque. Sa réponse est traitée comme une entrée hostile par `app/etl/mock_api_import.py`, son seul consommateur à ce jour : les quatre garde-fous attendus sont en place, voir la ligne correspondante plus haut. Reste ouvert : le plafond de taille s'applique après désérialisation de la réponse, borner le corps HTTP lui-même demanderait une lecture en flux ; et `APP_MOCK_API_BASE_URL` n'impose pas `https`, donc les identifiants Basic partiraient en clair sur une URL en `http`. La conséquence la plus sérieuse n'est pas la fausse alerte, c'est l'empoisonnement du jeu d'entraînement du modèle de prédiction. |
| **A08 Software and Data Integrity Failures** | **partiel** | La CI vérifie le code mais n'analyse ni les dépendances ni les images. `.terraform.lock.hcl` reste ignoré par git, ce qui contredit une chaîne d'approvisionnement maîtrisée. |
| **A10 Server-Side Request Forgery** | **sans objet aujourd'hui** | Aucune URL sortante n'est pilotée par une donnée utilisateur. Le jour où l'adresse d'une source devient un champ de configuration, il faudra une liste blanche de schémas et d'hôtes, sans suivi de redirection. |
| **Cantonnement des accès ETL et ML** | **dette assumée** | Le compte applicatif porte l'identité, le rôle PostgreSQL porterait le cantonnement. Voir ADR 0003. Plus coûteuse depuis Airflow (#115) : ce service publie le port 8080, détient les identifiants Postgres complets (`ML_DATABASE_URL`, mêmes que le backend) et permet de déclencher l'exécution de code depuis son interface. Un compte Airflow compromis atteint donc toute la base, pas seulement `reading`/`site`. Le compte admin Airflow est distinct des `app_user` et son mot de passe passe par l'environnement, jamais par `argv`. |
| **Cantonnement des accès ETL et ML** | **dette assumée** | Le compte applicatif porte l'identité, le rôle PostgreSQL porterait le cantonnement. Voir ADR 0003. Plus coûteuse depuis Airflow (#115) : ce service publie le port 8080, détient les identifiants Postgres complets (`ML_DATABASE_URL`, mêmes que le backend) et permet de déclencher l'exécution de code depuis son interface. Un compte Airflow compromis atteint donc toute la base, pas seulement `reading`/`site`. Aggravée par #116 : le conteneur reçoit aussi `DATABASE_URL` et exécute le code du backend en sous-processus (ADR 0008). Atténuations en place : le compte admin Airflow est distinct des `app_user` et son mot de passe passe par l'environnement, jamais par `argv` ; et l'`APP_SECRET_KEY` donnée à Airflow est distincte de celle de l'API, pour qu'une compromission ne livre pas la clé de signature des JWT. |
| **Non-répudiation de l'audit** | **dette assumée** | Les déclencheurs arrêtent les accidents, pas un compte détenant `ALTER TABLE`. Voir ADR 0004. |
## Ce qu'il faut répondre, et ne pas répondre
+2 -2
View File
@@ -663,8 +663,8 @@ mock_api_import.py
La logique d'extraction, de transformation et de chargement est donc disponible pour les deux sources de données du MVP.
Airflow tourne désormais réellement (`etl/airflow/`, `make airflow-up`), mais il orchestre pour l'instant le pipeline ML (`ml_train`/`ml_score`, issue #115), pas encore ces deux imports : orchestrer `historical_import.py` et `mock_api_import.py` (normalisation et chargement micro-batch, issues #15/#16) reste à faire.
Airflow tourne désormais réellement (`etl/airflow/`, `make airflow-up`) et orchestre le pipeline ML (`ml_train`/`ml_score`, issue #115) ainsi que la détection d'alertes et la génération des recommandations (`alertes`, issue #116). Il n'orchestre pas encore ces deux imports : `historical_import.py` et `mock_api_import.py` (normalisation et chargement micro-batch, issues #15/#16) restent à faire.
Airflow permet de planifier les traitements, gérer leur ordre d'exécution, suivre leur état et remonter les erreurs. Il ne remplace pas la logique ETL Python existante : les scripts actuels restent responsables de l'extraction, de la validation, de la transformation et du chargement. `etl/airflow/dags/ml_train.py` et `ml_score.py` montrent le patron retenu (des `BashOperator` qui invoquent le script tel quel).
Airflow permet de planifier les traitements, gérer leur ordre d'exécution, suivre leur état et remonter les erreurs. Il ne remplace pas la logique ETL Python existante : les scripts actuels restent responsables de l'extraction, de la validation, de la transformation et du chargement. `etl/airflow/dags/ml_train.py`, `ml_score.py` et `alertes.py` montrent le patron retenu (des `BashOperator` qui invoquent le script tel quel, dans l'environnement `uv` que l'image embarque pour lui).
Le pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning.
+19 -7
View File
@@ -1,7 +1,8 @@
# Image Airflow EnerVision : ajoute le projet ml/ dans son propre environnement Python 3.14,
# distinct du Python 3.12 qui fait tourner Airflow lui-meme (apache-airflow 2.10 ne supporte pas
# 3.14), pour que les DAGs puissent lancer `uv run python -m enervision_ml.train`/`.score` en
# sous-processus. Airflow ne devient jamais un consommateur direct de LightGBM/MLflow.
# Image Airflow EnerVision : ajoute ml/ et apps/backend/ dans leurs propres environnements Python
# 3.14, distincts du Python 3.12 qui fait tourner Airflow lui-meme (apache-airflow 2.10 ne supporte
# pas 3.14), pour que les DAGs puissent lancer `uv run python -m enervision_ml.train`/`.score`,
# `app.detection.internal_alerts` et `app.cli` en sous-processus. Airflow ne devient jamais un
# consommateur direct de LightGBM, de MLflow ou du SQLAlchemy du backend. Cf. ADR 0008.
FROM apache/airflow:2.10.4-python3.12
# LightGBM est compile contre libgomp (OpenMP), absent de l'image de base (minimale, sans
@@ -16,7 +17,8 @@ RUN apt-get update \
# `ml_train` et `ml_score` (le modele ecrit par l'un, lu par l'autre). Un volume nomme herite des
# permissions du repertoire qu'il recouvre a son premier montage ; sans ce chown prealable, il
# serait cree root:root et illisible par le conteneur, qui tourne en `airflow` (uid 50000).
RUN mkdir -p /opt/ml/state && chown -R airflow:root /opt/ml
# `/opt/backend` ne porte aucun volume, mais `WORKDIR` le creerait root meme sous `USER airflow`.
RUN mkdir -p /opt/ml/state /opt/backend && chown -R airflow:root /opt/ml /opt/backend
USER airflow
# L'image de base embarque deja un `uv`, mais trop ancien (0.4.29) pour le format de verrou de
@@ -24,9 +26,10 @@ USER airflow
# (apps/backend/Dockerfile).
COPY --from=ghcr.io/astral-sh/uv:0.11.26 /uv /home/airflow/.local/bin/uv
# Piege : pas de `UV_PROJECT_ENVIRONMENT` global. Il vaudrait pour les deux projets, et `uv run`
# dans l'un resoudrait le venv de l'autre. Par defaut, uv prend `<projet>/.venv`, donc le bon.
ENV UV_COMPILE_BYTECODE=1 \
UV_LINK_MODE=copy \
UV_PROJECT_ENVIRONMENT=/opt/ml/.venv
UV_LINK_MODE=copy
WORKDIR /opt/ml
@@ -36,4 +39,13 @@ RUN uv sync --locked --no-install-project --no-dev
COPY --chown=airflow:root ml/enervision_ml ./enervision_ml
RUN uv sync --locked --no-dev
WORKDIR /opt/backend
# `packages = ["app"]` : le reste de apps/backend (alembic, tests) n'a rien a faire dans l'image.
COPY --chown=airflow:root apps/backend/pyproject.toml apps/backend/uv.lock ./
RUN uv sync --locked --no-install-project --no-dev
COPY --chown=airflow:root apps/backend/app ./app
RUN uv sync --locked --no-dev
WORKDIR /opt/airflow
+65
View File
@@ -0,0 +1,65 @@
"""DAG de détection des alertes et de génération des recommandations (issue #116).
Ordonnance ce que `docs/architecture/20-backend.md` et l'ADR 0006 décrivent encore comme lancé à
la main. Toute la logique reste dans `apps/backend`, ce DAG ne fait que l'appeler, sur le patron
de `ml_score` (cf. `docs/architecture/10-infra.md`, section Airflow, et l'ADR 0008 pour
l'environnement `/opt/backend` que l'image embarque désormais).
Planifié à la quinzième minute plutôt qu'à l'heure pile : la règle `anomaly` compare une lecture
à la `prediction` du même instant, que `ml_score` (`@hourly`) vient d'écrire. Aucune dépendance
déclarée entre les deux DAGs pour autant, quatre règles sur cinq ne touchent pas au modèle et un
modèle jamais entraîné ne doit pas priver le parc de ses alertes.
"""
from __future__ import annotations
from datetime import datetime, timedelta
from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
# Le backend a son propre environnement uv dans l'image (ADR 0008). `--no-sync` et
# `env -u VIRTUAL_ENV` : cf. `ml_train.py`, même raisonnement.
COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m"
# Les deux tâches sont idempotentes en base (`ON CONFLICT DO NOTHING` sur
# `uq_alert_source_reference` et `uq_recommendation_alert_rule`) : reprendre ne duplique rien.
TENTATIVES = 2
DELAI_ENTRE_TENTATIVES = timedelta(minutes=2)
# `execution_timeout` vaut par tentative : c'est le pire cas des deux tâches enchaînées, reprises
# et délais compris, qui doit tenir sous le pas horaire. Les tests d'intégrité en font le calcul.
PLAFOND_PAR_TACHE = timedelta(minutes=5)
with DAG(
dag_id="alertes",
description=(
"Détecte les alertes internes puis génère les recommandations "
"(app.detection.internal_alerts, app.cli)."
),
schedule="15 * * * *",
start_date=datetime(2026, 1, 1),
catchup=False,
# Deux exécutions simultanées analyseraient la même fenêtre de 48h, et la génération relit
# l'intégralité de la table `alert` à chaque passage.
max_active_runs=1,
tags=["alertes"],
) as dag:
detection = BashOperator(
task_id="detection",
bash_command=f"{COMMANDE_BACKEND} app.detection.internal_alerts",
retries=TENTATIVES,
retry_delay=DELAI_ENTRE_TENTATIVES,
execution_timeout=PLAFOND_PAR_TACHE,
)
recommandations = BashOperator(
task_id="recommandations",
bash_command=f"{COMMANDE_BACKEND} app.cli generate-recommendations",
retries=TENTATIVES,
retry_delay=DELAI_ENTRE_TENTATIVES,
execution_timeout=PLAFOND_PAR_TACHE,
)
# `recommendation.alert_id` est une clé étrangère `NOT NULL` : la génération n'a rien à lire
# tant que la détection n'a pas écrit.
detection >> recommandations
+65 -6
View File
@@ -5,10 +5,19 @@ from datetime import timedelta
from pathlib import Path
import pytest
from airflow.models.baseoperator import BaseOperator
from airflow.models.dagbag import DagBag
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
DAG_IDS = ["ml_train", "ml_score", "alertes"]
TACHES = [
("ml_train", "train"),
("ml_score", "score"),
("alertes", "detection"),
("alertes", "recommandations"),
]
@pytest.fixture(scope="module")
def dagbag() -> DagBag:
@@ -20,7 +29,7 @@ def test_dags_folder_has_no_import_error(dagbag: DagBag) -> None:
def test_every_expected_dag_is_discovered(dagbag: DagBag) -> None:
assert set(dagbag.dag_ids) == {"ml_train", "ml_score"}
assert set(dagbag.dag_ids) == set(DAG_IDS)
def test_ml_train_has_no_schedule(dagbag: DagBag) -> None:
@@ -32,6 +41,12 @@ def test_ml_score_runs_every_hour(dagbag: DagBag) -> None:
assert dagbag.dags["ml_score"].timetable.summary == "0 * * * *"
def test_alertes_runs_after_the_hourly_scoring(dagbag: DagBag) -> None:
# Le decalage n'est pas cosmetique : la regle `anomaly` compare une lecture a la `prediction`
# du meme instant, que `ml_score` ecrit a l'heure pile.
assert dagbag.dags["alertes"].timetable.summary == "15 * * * *"
def test_ml_train_task_calls_the_training_module(dagbag: DagBag) -> None:
tache = dagbag.dags["ml_train"].get_task("train")
assert "enervision_ml.train" in tache.bash_command
@@ -42,6 +57,28 @@ def test_ml_score_task_calls_the_scoring_module(dagbag: DagBag) -> None:
assert "enervision_ml.score" in tache.bash_command
def test_alertes_detection_task_calls_the_backend_detection(dagbag: DagBag) -> None:
tache = dagbag.dags["alertes"].get_task("detection")
assert "app.detection.internal_alerts" in tache.bash_command
def test_alertes_recommendation_task_calls_the_backend_cli(dagbag: DagBag) -> None:
tache = dagbag.dags["alertes"].get_task("recommandations")
assert "app.cli generate-recommendations" in tache.bash_command
@pytest.mark.parametrize("task_id", ["detection", "recommandations"])
def test_alertes_tasks_run_in_the_backend_environment(dagbag: DagBag, task_id: str) -> None:
# Le backend a son propre venv dans l'image, distinct de celui de ml/ (ADR 0008).
assert "/opt/backend" in dagbag.dags["alertes"].get_task(task_id).bash_command
def test_alertes_generates_recommendations_after_detecting(dagbag: DagBag) -> None:
# `recommendation.alert_id` est une cle etrangere `NOT NULL` : la generation n'a rien a lire
# tant que la detection n'a pas ecrit.
assert dagbag.dags["alertes"].get_task("detection").downstream_task_ids == {"recommandations"}
def test_ml_score_reuses_the_model_path_written_by_ml_train(dagbag: DagBag) -> None:
entrainement = dagbag.dags["ml_train"].get_task("train").bash_command
scoring = dagbag.dags["ml_score"].get_task("score").bash_command
@@ -51,14 +88,14 @@ def test_ml_score_reuses_the_model_path_written_by_ml_train(dagbag: DagBag) -> N
assert chemin_modele in scoring
@pytest.mark.parametrize("dag_id", ["ml_train", "ml_score"])
@pytest.mark.parametrize("dag_id", DAG_IDS)
def test_no_two_runs_of_a_dag_overlap(dagbag: DagBag, dag_id: str) -> None:
# Deux entrainements ecriraient le meme fichier modele, deux scorings inseriraient en meme
# temps dans `prediction`.
# temps dans `prediction`, deux detections analyseraient la meme fenetre.
assert dagbag.dags[dag_id].max_active_runs == 1
@pytest.mark.parametrize(("dag_id", "task_id"), [("ml_train", "train"), ("ml_score", "score")])
@pytest.mark.parametrize(("dag_id", "task_id"), TACHES)
def test_every_task_has_an_execution_timeout(dagbag: DagBag, dag_id: str, task_id: str) -> None:
# Sans plafond, une connexion pendue immobilise un slot du scheduler indefiniment.
assert dagbag.dags[dag_id].get_task(task_id).execution_timeout is not None
@@ -70,13 +107,35 @@ def test_ml_score_execution_timeout_stays_below_its_hourly_step(dagbag: DagBag)
assert timeout < timedelta(hours=1)
def duree_au_pire(tache: BaseOperator) -> timedelta:
# `execution_timeout` plafonne une tentative, pas la tache : deux reprises occupent trois
# plafonds et deux delais d'attente.
assert tache.execution_timeout is not None
return (tache.retries + 1) * tache.execution_timeout + tache.retries * tache.retry_delay
def test_alertes_worst_case_stays_below_its_hourly_step(dagbag: DagBag) -> None:
# Les deux taches s'enchainent : c'est leur somme, reprises comprises, qui doit tenir dans le
# pas horaire, sinon `max_active_runs=1` fait attendre l'execution suivante.
taches = [
dagbag.dags["alertes"].get_task(task_id) for task_id in ("detection", "recommandations")
]
assert sum((duree_au_pire(tache) for tache in taches), timedelta()) < timedelta(hours=1)
def test_ml_score_retries_after_a_transient_failure(dagbag: DagBag) -> None:
assert dagbag.dags["ml_score"].get_task("score").retries >= 1
@pytest.mark.parametrize(("dag_id", "task_id"), [("ml_train", "train"), ("ml_score", "score")])
@pytest.mark.parametrize("task_id", ["detection", "recommandations"])
def test_alertes_retries_after_a_transient_failure(dagbag: DagBag, task_id: str) -> None:
# Les deux commandes sont idempotentes en base, une reprise ne duplique rien.
assert dagbag.dags["alertes"].get_task(task_id).retries >= 1
@pytest.mark.parametrize(("dag_id", "task_id"), TACHES)
def test_tasks_never_resync_the_baked_environment(
dagbag: DagBag, dag_id: str, task_id: str
) -> None:
# Sans `--no-sync`, `uv run` reconstruit `enervision-ml` a chaque execution.
# Sans `--no-sync`, `uv run` reconstruit le projet a chaque execution.
assert "--no-sync" in dagbag.dags[dag_id].get_task(task_id).bash_command