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