Files
ENI-projet-piscine/etl/airflow/tests/test_dags.py
Johan LEROY f8d08c8686
Airflow / Lint et intégrité des DAGs (push) Successful in 56s
Airflow / Construction de l'image (push) Successful in 3m20s
fix(etl): borne le DAG alertes sur son pire cas et couvre sa seconde commande en CI
`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.
2026-09-21 13:28:23 +02:00

142 lines
5.8 KiB
Python

"""Tests d'integrite des DAGs : s'importent sans erreur, structure attendue. Pas d'execution
reelle des taches (ca reclamerait le conteneur avec `uv`/`enervision_ml`), juste la definition."""
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:
return DagBag(dag_folder=str(DAGS_FOLDER), include_examples=False)
def test_dags_folder_has_no_import_error(dagbag: DagBag) -> None:
assert dagbag.import_errors == {}
def test_every_expected_dag_is_discovered(dagbag: DagBag) -> None:
assert set(dagbag.dag_ids) == set(DAG_IDS)
def test_ml_train_has_no_schedule(dagbag: DagBag) -> None:
assert dagbag.dags["ml_train"].timetable.summary == "None"
def test_ml_score_runs_every_hour(dagbag: DagBag) -> None:
# `@hourly` est un alias Airflow pour ce cron, c'est sous cette forme que `.summary` le rend.
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
def test_ml_score_task_calls_the_scoring_module(dagbag: DagBag) -> None:
tache = dagbag.dags["ml_score"].get_task("score")
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
chemin_modele = "/opt/ml/state/models/lightgbm-consumption.txt"
assert chemin_modele in entrainement
assert chemin_modele in scoring
@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`, deux detections analyseraient la meme fenetre.
assert dagbag.dags[dag_id].max_active_runs == 1
@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
def test_ml_score_execution_timeout_stays_below_its_hourly_step(dagbag: DagBag) -> None:
timeout = dagbag.dags["ml_score"].get_task("score").execution_timeout
assert timeout is not None
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("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 le projet a chaque execution.
assert "--no-sync" in dagbag.dags[dag_id].get_task(task_id).bash_command