feat(etl): ordonnance la détection d'alertes et les recommandations par un DAG Airflow

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 `<projet>/.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.
This commit is contained in:
Johan LEROY
2026-09-21 12:13:55 +02:00
parent 3cd9a6b272
commit ae58a896d9
7 changed files with 175 additions and 17 deletions
+58 -6
View File
@@ -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