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: