fix(airflow): fiabilise l'import horaire de la Mock API
Airflow / Construction de l'image (push) Successful in 1m4s
Backend / Analyse statique de sécurité (push) Successful in 7s
Backend / Tests exigeant une base (push) Failing after 4m55s
Airflow / Lint et intégrité des DAGs (push) Successful in 9m49s
Backend / Lint, typage et tests (push) Successful in 10m12s
Backend / Audit des dépendances (push) Successful in 9m36s
SonarQube / build-front (push) Successful in 9m35s
SonarQube / test-ml (push) Failing after 5m37s
SonarQube / build-back (push) Successful in 9m46s
SonarQube / test-front (push) Failing after 5m4s
SonarQube / test-back (push) Failing after 5m9s
SonarQube / SonarQube (push) Skipped
Airflow / Construction de l'image (push) Successful in 1m4s
Backend / Analyse statique de sécurité (push) Successful in 7s
Backend / Tests exigeant une base (push) Failing after 4m55s
Airflow / Lint et intégrité des DAGs (push) Successful in 9m49s
Backend / Lint, typage et tests (push) Successful in 10m12s
Backend / Audit des dépendances (push) Successful in 9m36s
SonarQube / build-front (push) Successful in 9m35s
SonarQube / test-ml (push) Failing after 5m37s
SonarQube / build-back (push) Successful in 9m46s
SonarQube / test-front (push) Failing after 5m4s
SonarQube / test-back (push) Failing after 5m9s
SonarQube / SonarQube (push) Skipped
This commit is contained in:
@@ -23,7 +23,7 @@ Ce que la documentation apporte à chacun : [docs/architecture/00-vue-ensemble.m
|
||||
| Backend | FastAPI, Python 3.14 | `apps/backend` | En place |
|
||||
| Frontend | Angular 22, Node 26 | `apps/frontend` | En place |
|
||||
| Base | PostgreSQL 17 + TimescaleDB | `db` | En place |
|
||||
| ETL | Apache Airflow | `etl/airflow` | Quatre DAGs |
|
||||
| ETL | Apache Airflow | `etl/airflow` | Cinq DAGs |
|
||||
| Infra | Terraform (k3s single-node) | `infra/terraform` | Initialise |
|
||||
| Reverse proxy | Nginx, TLS | `infra/proxy` | En place |
|
||||
| CI/CD | GitHub Actions | `.github/workflows` | En place |
|
||||
|
||||
@@ -45,15 +45,19 @@ CAPACITY_BOUNDS = (0.0, 100_000.0)
|
||||
def create_mock_api_client() -> httpx.AsyncClient:
|
||||
settings = get_settings()
|
||||
|
||||
if settings.mock_api_username is None or settings.mock_api_password is None:
|
||||
username = settings.mock_api_username
|
||||
password = (
|
||||
settings.mock_api_password.get_secret_value()
|
||||
if settings.mock_api_password is not None
|
||||
else None
|
||||
)
|
||||
|
||||
if not username or not username.strip() or not password or not password.strip():
|
||||
raise ValueError("Les identifiants de l'API Mock ne sont pas configurés.")
|
||||
|
||||
return httpx.AsyncClient(
|
||||
base_url=settings.mock_api_base_url.rstrip("/"),
|
||||
auth=(
|
||||
settings.mock_api_username,
|
||||
settings.mock_api_password.get_secret_value(),
|
||||
),
|
||||
auth=(username, password),
|
||||
timeout=settings.mock_api_timeout_seconds,
|
||||
)
|
||||
|
||||
|
||||
@@ -283,6 +283,41 @@ def test_create_mock_api_client_requires_credentials(
|
||||
mock_api_import.create_mock_api_client()
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("username", "password_value"),
|
||||
[
|
||||
("", "test-password"),
|
||||
("test-user", ""),
|
||||
(" ", "test-password"),
|
||||
("test-user", " "),
|
||||
],
|
||||
)
|
||||
def test_create_mock_api_client_rejects_empty_credentials(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
username: str,
|
||||
password_value: str,
|
||||
) -> None:
|
||||
password = MagicMock()
|
||||
password.get_secret_value.return_value = password_value
|
||||
|
||||
settings = SimpleNamespace(
|
||||
mock_api_username=username,
|
||||
mock_api_password=password,
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"get_settings",
|
||||
lambda: settings,
|
||||
)
|
||||
|
||||
with pytest.raises(
|
||||
ValueError,
|
||||
match="Les identifiants de l'API Mock ne sont pas configurés",
|
||||
):
|
||||
mock_api_import.create_mock_api_client()
|
||||
|
||||
|
||||
async def test_create_mock_api_client_uses_configuration(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
|
||||
+8
-5
@@ -42,11 +42,6 @@ x-airflow-common: &airflow-common
|
||||
# vérifient aucun jeton. Airflow permet d'exécuter du code depuis son interface (ADR 0008).
|
||||
APP_SECRET_KEY: ${AIRFLOW_APP_SECRET_KEY:-}
|
||||
|
||||
# Configuration utilisée par `app.etl.mock_api_import` dans le scheduler Airflow.
|
||||
APP_MOCK_API_BASE_URL: ${APP_MOCK_API_BASE_URL:-https://api-mock.charlieandre.fr}
|
||||
APP_MOCK_API_USERNAME: ${APP_MOCK_API_USERNAME:-}
|
||||
APP_MOCK_API_PASSWORD: ${APP_MOCK_API_PASSWORD:-}
|
||||
APP_MOCK_API_TIMEOUT_SECONDS: ${APP_MOCK_API_TIMEOUT_SECONDS:-10}
|
||||
volumes:
|
||||
- ./etl/airflow/dags:/opt/airflow/dags
|
||||
- ./etl/airflow/plugins:/opt/airflow/plugins
|
||||
@@ -174,6 +169,14 @@ services:
|
||||
airflow-scheduler:
|
||||
<<: *airflow-common
|
||||
command: scheduler
|
||||
environment:
|
||||
<<: *airflow-common-env
|
||||
# LocalExecutor exécute les tâches dans le scheduler : lui seul a besoin des
|
||||
# identifiants de l'API Mock.
|
||||
APP_MOCK_API_BASE_URL: ${APP_MOCK_API_BASE_URL:-https://api-mock.charlieandre.fr}
|
||||
APP_MOCK_API_USERNAME: ${APP_MOCK_API_USERNAME:-}
|
||||
APP_MOCK_API_PASSWORD: ${APP_MOCK_API_PASSWORD:-}
|
||||
APP_MOCK_API_TIMEOUT_SECONDS: ${APP_MOCK_API_TIMEOUT_SECONDS:-10}
|
||||
depends_on:
|
||||
db:
|
||||
condition: service_healthy
|
||||
|
||||
@@ -74,6 +74,7 @@ Le lien `airflow --> db` est maintenant en trait plein : cinq DAGs tournent, deu
|
||||
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), `historical_import` pour le dataset historique
|
||||
(issue #119) et `mock_api_import` pour l'ingestion horaire de l'API Mock (issue #15).
|
||||
La réconciliation globale des données provenant des deux sources reste à compléter dans l'issue #15.
|
||||
|
||||
Le lien `prom -.-> api` de même : l'API expose bien `/metrics` au format Prometheus, mais aucun
|
||||
collecteur ne vient le lire.
|
||||
@@ -88,15 +89,16 @@ 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 | Docker Compose, Nginx, Terraform, k3s single-node | `infra`, `docker-compose.prod.yml` | `En cours` | Reverse proxy et overlay de déploiement écrits et validés, jamais lancés sur le serveur ([ADR 0007](../adr/0007-terminaison-tls-et-reverse-proxy-nginx.md)). Provisionnement de la VM par Terraform, qui installe Docker, prépare les deux environnements et enregistre le runner, jamais appliqué ([ADR 0010](../adr/0010-terraform-provisionne-github-actions-deploie.md)). Module d'installation k3s 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. Quatre DAGs en sous-processus `uv run` : `ml_train`, `ml_score`, `alertes` et `historical_import`. Le DAG historique orchestre `app.etl.historical_import` et charge `dataset`, `site` et `reading`. L'orchestration API Mock reste à compléter dans #15 |
|
||||
| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver et scheduler avec LocalExecutor via Docker Compose, sur une base PostgreSQL dédiée. Cinq DAGs sont présents : `ml_train`, `ml_score`, `alertes`, `historical_import` et `mock_api_import`. L'import historique reste manuel et l'import API Mock est exécuté chaque heure. La réconciliation globale des deux sources reste à compléter dans l'issue #15. |
|
||||
| CI/CD | GitHub Actions | `.github/workflows` | `En cours` | 7 workflows, 19 jobs : lint, typage, tests avec seuil de couverture bloquant, tests d'intégration sur TimescaleDB réel, audit de dépendances, SAST Bandit, quality gate SonarCloud, intégrité des DAGs Airflow, formatage et validation du Terraform. Déploiement continu vers la VM ENI écrit par `deploy.yml`, `dev` en recette et `main` en production après approbation ([ADR 0009](../adr/0009-deux-environnements-compose-sur-la-vm-eni.md)), mais jamais exécuté : la machine n'est pas provisionnée et le runner n'y est pas enregistré. Détail dans [50-cicd.md](50-cicd.md) |
|
||||
|
||||
## Flux bout en bout
|
||||
|
||||
Statut : `En cours`. **Le chemin de lecture tourne** : base, API et frontend. **Le chemin
|
||||
d'ingestion dessiné ci-dessous n'existe pas** : les trois DAGs livrés (`ml_train`, `ml_score`,
|
||||
issue #115 ; `alertes`, issue #116) orchestrent le pipeline ML et la détection d'alertes, pas
|
||||
l'ingestion, qui reste lancée à la main par les scripts d'import (issues #15 et #16).
|
||||
Statut : `En cours`. **Le chemin de lecture tourne** entre la base, l'API et le frontend.
|
||||
**Le chemin d'ingestion est maintenant orchestré par Airflow** : `historical_import` charge
|
||||
manuellement le dataset CSV/JSON et `mock_api_import` collecte périodiquement les mesures de
|
||||
l'API Mock. La réconciliation globale des données provenant des deux sources reste à compléter
|
||||
dans l'issue #15.
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
|
||||
@@ -87,17 +87,23 @@ l'[ADR 0008](../adr/0008-airflow-execute-le-code-du-backend.md).
|
||||
| `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` |
|
||||
| `historical_import` | manuelle | `app.etl.historical_import`, dans `/opt/backend/.venv` ; les fichiers de `data/raw` sont montés en lecture seule dans `/opt/data/raw` |
|
||||
| `mock_api_import` | `0 * * * *` | `app.etl.mock_api_import`, dans `/opt/backend/.venv` ; importe l'intervalle horaire Airflow précédent depuis l'API Mock |
|
||||
| `mock_api_import` | `45 * * * *` | `app.etl.mock_api_import`, dans `/opt/backend/.venv` ; importe l'heure précédant son déclenchement depuis l'API Mock |
|
||||
|
||||
Le DAG `historical_import` réutilise le pipeline historique existant sans dupliquer sa logique.
|
||||
Il reste manuel, car le dataset sert à initialiser l'environnement. Le montage
|
||||
`./data/raw:/opt/data/raw:ro` permet au scheduler de lire les fichiers CSV/JSON sans pouvoir les
|
||||
modifier.
|
||||
|
||||
Le DAG `mock_api_import` exécute le pipeline API Mock toutes les heures. Il transmet
|
||||
`data_interval_start` et `data_interval_end` au script backend et charge les mesures dans les
|
||||
mêmes tables `site` et `reading` que le pipeline historique. `max_active_runs=1` empêche deux
|
||||
intervalles de s'exécuter simultanément.
|
||||
Le DAG `mock_api_import` exécute le pipeline API Mock toutes les heures, à la minute `:45`.
|
||||
Un `CronTriggerTimetable` explicite lui attribue un intervalle d'une heure, y compris lors d'un
|
||||
déclenchement manuel. Il transmet cet intervalle au script backend et charge les mesures dans
|
||||
les tables communes `site` et `reading`. Le décalage à `:45` laisse quinze minutes avant le
|
||||
scoring exécuté à l'heure pile, puis quinze minutes supplémentaires avant les alertes à `:15`.
|
||||
`max_active_runs=1` empêche deux exécutions du DAG de se chevaucher.
|
||||
|
||||
Le DAG conserve `catchup=False` pour éviter un rattrapage massif depuis sa date de démarrage.
|
||||
Une interruption du scheduler peut donc créer un intervalle manquant, qui devra être rejoué
|
||||
explicitement par une opération de backfill.
|
||||
|
||||
**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
|
||||
|
||||
+9
-7
@@ -663,14 +663,16 @@ 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`) et orchestre le pipeline
|
||||
ML (`ml_train`/`ml_score`, issue #115), la détection d'alertes et la génération des
|
||||
recommandations (`alertes`, issue #116), ainsi que l'import historique
|
||||
(`historical_import`, issue #119).
|
||||
Airflow tourne désormais réellement (`etl/airflow/`, `make airflow-up`) et orchestre cinq DAGs :
|
||||
le pipeline ML (`ml_train` et `ml_score`, issue #115), la détection d'alertes et la génération
|
||||
des recommandations (`alertes`, issue #116), l'import historique (`historical_import`,
|
||||
issue #119) et l'import périodique de l'API Mock (`mock_api_import`, issue #15).
|
||||
|
||||
Le DAG `historical_import` est déclenché manuellement. Il exécute
|
||||
`app.etl.historical_import` avec les fichiers montés en lecture seule depuis `data/raw` vers
|
||||
`/opt/data/raw`.
|
||||
Le DAG `mock_api_import` s'exécute chaque heure, à la minute `:45`. Il appelle
|
||||
`app.etl.mock_api_import` avec un intervalle explicite d'une heure et une limite de 1 000 lectures
|
||||
par site. Les deux pipelines normalisent leurs données vers les tables communes `site` et
|
||||
`reading`, tout en conservant leur source (`csv` ou `api_history`). La réconciliation globale
|
||||
des deux sources reste à compléter dans l'issue #15.
|
||||
|
||||
Le DAG `mock_api_import` exécute `app.etl.mock_api_import` toutes les heures. Chaque exécution
|
||||
traite l'intervalle Airflow précédent. Les deux pipelines normalisent leurs données vers les
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
"""DAG d'import périodique des données de l'API Mock EnerVision (issue #15).
|
||||
|
||||
Orchestre le pipeline existant `app.etl.mock_api_import` sans dupliquer sa logique ETL.
|
||||
Chaque exécution traite l'intervalle horaire Airflow précédent.
|
||||
Chaque exécution traite l'heure précédant son déclenchement.
|
||||
|
||||
Le pipeline backend reste responsable de la validation, de la normalisation, du suivi de la
|
||||
qualité, de l'idempotence et du chargement dans PostgreSQL/TimescaleDB.
|
||||
@@ -13,12 +13,14 @@ from datetime import datetime, timedelta
|
||||
|
||||
from airflow.providers.standard.operators.bash import BashOperator
|
||||
from airflow.sdk import DAG
|
||||
from airflow.timetables.trigger import CronTriggerTimetable
|
||||
|
||||
# Le backend possède son propre environnement uv dans l'image Airflow (ADR 0008).
|
||||
COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m"
|
||||
|
||||
# L'API et le pipeline backend plafonnent une réponse à 1 000 lectures par site.
|
||||
LIMITE_LECTURES = 60
|
||||
# Le pipeline backend et l'API acceptent au maximum 1 000 lectures par site.
|
||||
# Cette marge évite de perdre silencieusement une lecture si une heure en contient plus de 60.
|
||||
LIMITE_LECTURES = 1000
|
||||
|
||||
# Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste
|
||||
# inférieure au pas horaire du DAG.
|
||||
@@ -26,10 +28,19 @@ NOMBRE_REPRISES = 2
|
||||
DELAI_ENTRE_REPRISES = timedelta(minutes=2)
|
||||
PLAFOND_PAR_TENTATIVE = timedelta(minutes=10)
|
||||
|
||||
# L'intervalle est déclaré explicitement pour ne pas dépendre de la valeur du paramètre Airflow
|
||||
# `create_cron_data_intervals`. Le déclenchement à :45 laisse quinze minutes avant `ml_score`,
|
||||
# exécuté à l'heure pile, puis avant `alertes`, exécuté à :15.
|
||||
PLANIFICATION = CronTriggerTimetable(
|
||||
"45 * * * *",
|
||||
timezone="UTC",
|
||||
interval=timedelta(hours=1),
|
||||
)
|
||||
|
||||
with DAG(
|
||||
dag_id="mock_api_import",
|
||||
description="Importe chaque heure les données de l'API Mock dans site et reading.",
|
||||
schedule="@hourly",
|
||||
schedule=PLANIFICATION,
|
||||
start_date=datetime(2026, 1, 1),
|
||||
catchup=False,
|
||||
# Deux exécutions simultanées pourraient demander et traiter le même intervalle.
|
||||
@@ -40,8 +51,8 @@ with DAG(
|
||||
task_id="import_mock_api",
|
||||
bash_command=(
|
||||
f"{COMMANDE_BACKEND} app.etl.mock_api_import "
|
||||
'--start-time "{{ data_interval_start.isoformat() }}" '
|
||||
'--end-time "{{ data_interval_end.isoformat() }}" '
|
||||
"--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\" "
|
||||
"--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\" "
|
||||
f"--limit {LIMITE_LECTURES}"
|
||||
),
|
||||
retries=NOMBRE_REPRISES,
|
||||
|
||||
@@ -1,12 +1,13 @@
|
||||
"""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 datetime import datetime, timedelta
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
from airflow.dag_processing.dagbag import DagBag
|
||||
from airflow.sdk import BaseOperator
|
||||
from airflow.timetables.trigger import CronTriggerTimetable
|
||||
|
||||
DAGS_FOLDER = Path(__file__).resolve().parent.parent / "dags"
|
||||
|
||||
@@ -59,8 +60,17 @@ def test_historical_import_has_no_schedule(dagbag: DagBag) -> None:
|
||||
assert dagbag.dags["historical_import"].schedule is None
|
||||
|
||||
|
||||
def test_mock_api_import_runs_every_hour(dagbag: DagBag) -> None:
|
||||
assert dagbag.dags["mock_api_import"].timetable.expression == "0 * * * *"
|
||||
def test_mock_api_import_uses_an_explicit_hourly_interval(dagbag: DagBag) -> None:
|
||||
timetable = dagbag.dags["mock_api_import"].timetable
|
||||
|
||||
assert isinstance(timetable, CronTriggerTimetable)
|
||||
assert timetable.serialize()["expression"] == "45 * * * *"
|
||||
|
||||
manual_interval = timetable.infer_manual_data_interval(
|
||||
run_after=datetime.fromisoformat("2026-09-22T12:30:00+00:00"),
|
||||
)
|
||||
|
||||
assert manual_interval.end - manual_interval.start == timedelta(hours=1)
|
||||
|
||||
|
||||
def test_ml_train_task_calls_the_training_module(dagbag: DagBag) -> None:
|
||||
@@ -104,9 +114,9 @@ def test_mock_api_import_calls_the_existing_backend_module(dagbag: DagBag) -> No
|
||||
def test_mock_api_import_uses_the_airflow_data_interval(dagbag: DagBag) -> None:
|
||||
commande = dagbag.dags["mock_api_import"].get_task("import_mock_api").bash_command
|
||||
|
||||
assert '--start-time "{{ data_interval_start.isoformat() }}"' in commande
|
||||
assert '--end-time "{{ data_interval_end.isoformat() }}"' in commande
|
||||
assert "--limit 60" in commande
|
||||
assert "--start-time \"{{ data_interval_start.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande
|
||||
assert "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\"" in commande
|
||||
assert "--limit 1000" in commande
|
||||
|
||||
|
||||
@pytest.mark.parametrize("task_id", ["detection", "recommandations"])
|
||||
|
||||
Reference in New Issue
Block a user