From ae58a896d974a03a0e72f7f68e3a0588e5dad029 Mon Sep 17 00:00:00 2001 From: Johan LEROY Date: Mon, 21 Sep 2026 12:13:55 +0200 Subject: [PATCH 1/3] =?UTF-8?q?feat(etl):=20ordonnance=20la=20d=C3=A9tecti?= =?UTF-8?q?on=20d'alertes=20et=20les=20recommandations=20par=20un=20DAG=20?= =?UTF-8?q?Airflow?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Le DAG `alertes` enchaîne `app.detection.internal_alerts` puis `app.cli generate-recommendations`, à la quinzième minute de chaque heure. Le décalage laisse finir `ml_score`, qui écrit à l'heure pile les prédictions dont la règle `anomaly` a besoin, sans créer de dépendance entre les deux DAGs : 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. L'image Airflow porte un second environnement uv, `/opt/backend/.venv`, puisque la logique vit dans le backend (ADR 0006) et qu'aucune route HTTP ne l'expose. Le `UV_PROJECT_ENVIRONMENT` global hérité de l'issue #115 disparaît : il vaut pour tous les projets, donc `uv run` depuis `/opt/ml` résolvait le venv du backend. uv prend `/.venv` par défaut, se placer dans le dossier suffit. La CI vérifie maintenant que les deux environnements s'importent sans réseau. Le conteneur reçoit `DATABASE_URL` en asyncpg et une `APP_SECRET_KEY` distincte de celle de l'API, alimentée par `AIRFLOW_APP_SECRET_KEY` : la détection ne signe aucun jeton, et Airflow permet d'exécuter du code depuis son interface. --- .env.example | 5 +++ .github/workflows/airflow.yml | 20 +++++++++-- Makefile | 5 ++- docker-compose.yml | 7 ++++ etl/airflow/Dockerfile | 26 ++++++++++---- etl/airflow/dags/alertes.py | 65 ++++++++++++++++++++++++++++++++++ etl/airflow/tests/test_dags.py | 64 +++++++++++++++++++++++++++++---- 7 files changed, 175 insertions(+), 17 deletions(-) create mode 100644 etl/airflow/dags/alertes.py diff --git a/.env.example b/.env.example index e0abef3..b3ff600 100644 --- a/.env.example +++ b/.env.example @@ -38,3 +38,8 @@ 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 diff --git a/.github/workflows/airflow.yml b/.github/workflows/airflow.yml index a79f97f..9616da2 100644 --- a/.github/workflows/airflow.yml +++ b/.github/workflows/airflow.yml @@ -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,10 @@ 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. + - name: Vérifie que la détection d'alertes s'importe 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" diff --git a/Makefile b/Makefile index 66f4298..51b6e67 100644 --- a/Makefile +++ b/Makefile @@ -8,7 +8,7 @@ AIRFLOW := etl/airflow 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 help: ## Liste les cibles disponibles @@ -83,6 +83,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),) diff --git a/docker-compose.yml b/docker-compose.yml index 5f6c9c0..73ffa64 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -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: diff --git a/etl/airflow/Dockerfile b/etl/airflow/Dockerfile index 058928f..b934588 100644 --- a/etl/airflow/Dockerfile +++ b/etl/airflow/Dockerfile @@ -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 `/.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 diff --git a/etl/airflow/dags/alertes.py b/etl/airflow/dags/alertes.py new file mode 100644 index 0000000..f221ab1 --- /dev/null +++ b/etl/airflow/dags/alertes.py @@ -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) +# La somme des deux plafonds reste sous le pas horaire : une exécution pendue ne doit pas +# empiéter sur la suivante. +PLAFOND_PAR_TACHE = timedelta(minutes=15) + +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 diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index 58f6ab9..ccfbe5d 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -9,6 +9,14 @@ 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 +28,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 +40,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 +56,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 +87,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 +106,29 @@ def test_ml_score_execution_timeout_stays_below_its_hourly_step(dagbag: DagBag) assert timeout < timedelta(hours=1) +def test_alertes_execution_timeouts_stay_below_its_hourly_step(dagbag: DagBag) -> None: + # Les deux taches s'enchainent : c'est leur somme qui doit tenir dans le pas horaire. + plafonds = [ + dagbag.dags["alertes"].get_task(task_id).execution_timeout + for task_id in ("detection", "recommandations") + ] + assert all(plafond is not None for plafond in plafonds) + assert sum(plafonds, 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 From 26868801853bd83eb66a3085480383b7929d98de Mon Sep 17 00:00:00 2001 From: Johan LEROY Date: Mon, 21 Sep 2026 12:14:06 +0200 Subject: [PATCH 2/3] =?UTF-8?q?docs:=20acte=20l'ordonnancement=20des=20ale?= =?UTF-8?q?rtes=20par=20l'ADR=200008=20et=20met=20=C3=A0=20jour=20les=20vu?= =?UTF-8?q?es?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit L'ADR 0008 décide qu'Airflow exécute le code du backend en sous-processus plutôt que d'appeler l'API, et assume ce que cela coûte : une image plus lourde, la CI Airflow déclenchée par les changements du backend, une clé applicative de plus. Les vues suivent. Trois DAGs dans 10-infra.md et dans la vue d'ensemble, avec le motif du décalage horaire. La détection n'est plus « lancée à la main » dans 20-backend.md. La génération des recommandations gagne son troisième déclencheur dans 40-data.md. La dette de cantonnement ETL et ML porte l'aggravation comme l'atténuation. L'affirmation selon laquelle `etl/airflow/` ne contient que des `.gitkeep`, fausse depuis l'issue #115, disparaît. L'index des décisions omettait les ADR 0005 et 0006, il les récupère au passage. --- README.md | 4 +- docs/README.md | 3 + ...0008-airflow-execute-le-code-du-backend.md | 79 +++++++++++++++++++ docs/architecture/00-vue-ensemble.md | 9 ++- docs/architecture/10-infra.md | 44 ++++++++--- docs/architecture/20-backend.md | 19 +++-- docs/architecture/40-data.md | 14 ++-- docs/architecture/README.md | 7 +- docs/architecture/owasp-traceabilite.md | 2 +- etl/README.md | 4 +- 10 files changed, 153 insertions(+), 32 deletions(-) create mode 100644 docs/adr/0008-airflow-execute-le-code-du-backend.md diff --git a/README.md b/README.md index 33426fb..ec5ffd9 100644 --- a/README.md +++ b/README.md @@ -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 | | CI/CD | GitHub Actions | `.github/workflows` | Backend en place | | Monitoring | Prometheus, Grafana, Alertmanager | `monitoring` | A initialiser | @@ -47,7 +47,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 diff --git a/docs/README.md b/docs/README.md index 17859f5..bbd6978 100644 --- a/docs/README.md +++ b/docs/README.md @@ -11,3 +11,6 @@ | [0002](adr/0002-authentification-jwt-et-refresh-opaque.md) | Authentification par JWT d'accès et jeton de rafraîchissement opaque | | [0003](adr/0003-autorisation-rbac-a-trois-roles.md) | Autorisation RBAC à trois rôles, relecture du compte à chaque requête | | [0004](adr/0004-journal-d-audit-en-ajout-seul.md) | Journal d'audit en ajout seul, garanti par PostgreSQL | +| [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/` | +| [0008](adr/0008-airflow-execute-le-code-du-backend.md) | Airflow exécute le code du backend en sous-processus, dans son propre environnement | diff --git a/docs/adr/0008-airflow-execute-le-code-du-backend.md b/docs/adr/0008-airflow-execute-le-code-du-backend.md new file mode 100644 index 0000000..98e50c1 --- /dev/null +++ b/docs/adr/0008-airflow-execute-le-code-du-backend.md @@ -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 `/.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. diff --git a/docs/architecture/00-vue-ensemble.md b/docs/architecture/00-vue-ensemble.md index b031a50..5eda630 100644 --- a/docs/architecture/00-vue-ensemble.md +++ b/docs/architecture/00-vue-ensemble.md @@ -67,9 +67,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. @@ -84,7 +85,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 | Terraform, k3s single-node | `infra/terraform` | `En cours` | Module d'installation du cluster. 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 diff --git a/docs/architecture/10-infra.md b/docs/architecture/10-infra.md index 1c44eb0..a74855a 100644 --- a/docs/architecture/10-infra.md +++ b/docs/architecture/10-infra.md @@ -50,7 +50,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`) : @@ -59,13 +59,30 @@ 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. 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`. `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 @@ -76,7 +93,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 @@ -85,8 +110,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 diff --git a/docs/architecture/20-backend.md b/docs/architecture/20-backend.md index b9abfa4..2c52815 100644 --- a/docs/architecture/20-backend.md +++ b/docs/architecture/20-backend.md @@ -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` diff --git a/docs/architecture/40-data.md b/docs/architecture/40-data.md index 47fdc8a..82f769b 100644 --- a/docs/architecture/40-data.md +++ b/docs/architecture/40-data.md @@ -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 diff --git a/docs/architecture/README.md b/docs/architecture/README.md index 78a82c9..1e3a260 100644 --- a/docs/architecture/README.md +++ b/docs/architecture/README.md @@ -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 diff --git a/docs/architecture/owasp-traceabilite.md b/docs/architecture/owasp-traceabilite.md index 132e107..967791d 100644 --- a/docs/architecture/owasp-traceabilite.md +++ b/docs/architecture/owasp-traceabilite.md @@ -57,7 +57,7 @@ règles Bandit. Ajouter Bandit à la CI serait redondant, contrairement à ce qu | **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 diff --git a/etl/README.md b/etl/README.md index b442ae3..698ffca 100644 --- a/etl/README.md +++ b/etl/README.md @@ -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. From f8d08c86863682abf13a757cdebc2e699eb69794 Mon Sep 17 00:00:00 2001 From: Johan LEROY Date: Mon, 21 Sep 2026 13:28:23 +0200 Subject: [PATCH 3/3] fix(etl): borne le DAG alertes sur son pire cas et couvre sa seconde commande en CI MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `execution_timeout` plafonne une tentative, pas la tâche. Avec deux reprises, quinze minutes par tentative autorisaient quarante-neuf minutes par tâche et quatre-vingt-dix-huit pour l'enchaînement, quand le commentaire annonçait une somme tenant sous le pas horaire. Le plafond passe à cinq minutes, ce qui borne le pire cas à trente-huit minutes, et le test d'intégrité calcule désormais ce pire cas plutôt que la somme des plafonds : reprises et délais d'attente compris, c'est la durée qu'un `max_active_runs=1` fait payer à l'exécution suivante. La CI vérifie aussi `app.cli generate-recommendations --help` sans réseau. C'est la seconde commande du DAG, et son import tire FastAPI, les repositories et les services, donc une part de l'environnement `/opt/backend` que la détection seule ne touche pas. `10-infra.md` nomme enfin ce que le décalage de quinze minutes ne garantit pas : le plafond de `ml_score` valant trente minutes, un scoring qui déborde prive la règle `anomaly` de la prédiction de l'heure, qu'elle ne retrouvera au passage suivant que si sa fenêtre la couvre encore. --- .github/workflows/airflow.yml | 9 ++++++--- docs/architecture/10-infra.md | 15 +++++++++++---- etl/airflow/dags/alertes.py | 6 +++--- etl/airflow/tests/test_dags.py | 21 ++++++++++++++------- 4 files changed, 34 insertions(+), 17 deletions(-) diff --git a/.github/workflows/airflow.yml b/.github/workflows/airflow.yml index 9616da2..d5a722a 100644 --- a/.github/workflows/airflow.yml +++ b/.github/workflows/airflow.yml @@ -91,8 +91,11 @@ jobs: 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. - - name: Vérifie que la détection d'alertes s'importe sans réseau + # 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" + 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" diff --git a/docs/architecture/10-infra.md b/docs/architecture/10-infra.md index b3ecfd2..c6121e7 100644 --- a/docs/architecture/10-infra.md +++ b/docs/architecture/10-infra.md @@ -80,10 +80,17 @@ l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md). `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. 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`. +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 diff --git a/etl/airflow/dags/alertes.py b/etl/airflow/dags/alertes.py index f221ab1..4c043d2 100644 --- a/etl/airflow/dags/alertes.py +++ b/etl/airflow/dags/alertes.py @@ -26,9 +26,9 @@ COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync pytho # `uq_alert_source_reference` et `uq_recommendation_alert_rule`) : reprendre ne duplique rien. TENTATIVES = 2 DELAI_ENTRE_TENTATIVES = timedelta(minutes=2) -# La somme des deux plafonds reste sous le pas horaire : une exécution pendue ne doit pas -# empiéter sur la suivante. -PLAFOND_PAR_TACHE = timedelta(minutes=15) +# `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", diff --git a/etl/airflow/tests/test_dags.py b/etl/airflow/tests/test_dags.py index ccfbe5d..555849d 100644 --- a/etl/airflow/tests/test_dags.py +++ b/etl/airflow/tests/test_dags.py @@ -5,6 +5,7 @@ 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" @@ -106,14 +107,20 @@ def test_ml_score_execution_timeout_stays_below_its_hourly_step(dagbag: DagBag) assert timeout < timedelta(hours=1) -def test_alertes_execution_timeouts_stay_below_its_hourly_step(dagbag: DagBag) -> None: - # Les deux taches s'enchainent : c'est leur somme qui doit tenir dans le pas horaire. - plafonds = [ - dagbag.dags["alertes"].get_task(task_id).execution_timeout - for task_id in ("detection", "recommandations") +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 all(plafond is not None for plafond in plafonds) - assert sum(plafonds, timedelta()) < timedelta(hours=1) + 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: