Compare commits

...
11 changed files with 589 additions and 153 deletions
+153 -46
View File
@@ -1,17 +1,19 @@
# Contrainte : la réponse de l'API Mock est une entrée hostile, pas une source de confiance. # Contrainte : la réponse de l'API Mock est une entrée hostile, pas une source de confiance.
# Voir OWASP API10 dans docs/architecture/owasp-traceabilite.md. Rien de ce qu'elle renvoie # Voir OWASP API10 dans docs/architecture/owasp-traceabilite.md. Rien de ce qu'elle renvoie
# n'atteint la base sans passer par build_site_row() ou build_reading_row() : seuls les champs # n'atteint la base sans passer par build_site_row() ou build_reading_row() : seuls les champs
# attendus sont recopiés, les grandeurs physiques sont bornées par PHYSICAL_BOUNDS et la taille # attendus sont recopiés, les grandeurs physiques sont bornées par PHYSICAL_BOUNDS, la taille des
# des tableaux est plafonnée par MAX_SITES et par --limit. Une valeur hors bornes devient NULL # tableaux est plafonnée par MAX_SITES et par limit_for_window() (dérivé de la fenêtre, jamais
# et laisse sa trace dans null_reasons plutôt que de lever : le mock émet des anomalies par # fourni par l'appelant), et les lectures dont le timestamp déborde de la fenêtre demandée sont
# construction, et raw_data conserve de toute façon la réponse d'origine intacte. # écartées (fetch_readings). Une valeur hors bornes devient NULL et laisse sa trace dans
# null_reasons plutôt que de lever : le mock émet des anomalies par construction, et raw_data
# conserve de toute façon la réponse d'origine intacte.
from __future__ import annotations from __future__ import annotations
import argparse import argparse
import asyncio import asyncio
import json import json
from datetime import datetime from datetime import UTC, datetime
from typing import Any from typing import Any
import httpx import httpx
@@ -19,6 +21,7 @@ from sqlalchemy import text
from sqlalchemy.ext.asyncio import AsyncConnection, create_async_engine from sqlalchemy.ext.asyncio import AsyncConnection, create_async_engine
from app.core.config import get_settings from app.core.config import get_settings
from app.etl.historical_import import SOURCE_NAME as SOURCE_CSV
SOURCE_HISTORY = "api_history" SOURCE_HISTORY = "api_history"
@@ -181,6 +184,19 @@ async def upsert_sites(
) )
def _timestamp_in_window(reading: dict[str, Any], start_time: datetime, end_time: datetime) -> bool:
valeur = reading.get("timestamp")
if not isinstance(valeur, str):
return False
try:
instant = parse_datetime(valeur)
except ValueError:
return False
return start_time <= instant < end_time
async def fetch_readings( async def fetch_readings(
client: httpx.AsyncClient, client: httpx.AsyncClient,
site_id: str, site_id: str,
@@ -208,7 +224,22 @@ async def fetch_readings(
if len(payload) > limit: if len(payload) > limit:
raise ValueError(f"La réponse /api/v1/readings dépasse la limite demandée de {limit}.") raise ValueError(f"La réponse /api/v1/readings dépasse la limite demandée de {limit}.")
return payload # Le garde-fou `refuse_if_overlaps_historical_dataset` ne vérifie que la fenêtre demandée :
# une réponse (bug du mock, ou hostile) dont les `timestamp` débordent de
# `[start_time, end_time)` contournerait ce contrôle et écrirait exactement le doublon
# inter-source qu'il doit empêcher. Écarter ces lectures ici rend le contrôle par fenêtre
# suffisant.
dans_la_fenetre = [
lecture
for lecture in payload
if isinstance(lecture, dict) and _timestamp_in_window(lecture, start_time, end_time)
]
if len(dans_la_fenetre) != len(payload):
ecartees = len(payload) - len(dans_la_fenetre)
print(f"{site_id}: {ecartees} lecture(s) hors fenêtre écartée(s).")
return dans_la_fenetre
def build_reading_row( def build_reading_row(
@@ -244,6 +275,36 @@ def build_reading_row(
} }
# `uq_reading_source` autorise deux lignes au même (site_id, timestamp) dès que `source` diffère :
# sans ce garde-fou, importer une fenêtre déjà couverte par le dataset historique (source='csv')
# dupliquerait silencieusement chaque point plutôt que de lever une erreur. Ce garde-fou protège
# l'ingestion ; il ne dit rien de la lecture (`GET /readings` renvoie les deux lignes en cas de
# doublon malgré tout, cf. la section réconciliation de 40-data.md).
OVERLAP_CHECK = text(
"SELECT count(*) FROM reading WHERE source = :source_csv "
"AND timestamp >= :start_time AND timestamp < :end_time"
)
async def refuse_if_overlaps_historical_dataset(
connection: AsyncConnection,
start_time: datetime,
end_time: datetime,
) -> None:
resultat = await connection.execute(
OVERLAP_CHECK,
{"source_csv": SOURCE_CSV, "start_time": start_time, "end_time": end_time},
)
nombre = resultat.scalar_one()
if nombre > 0:
raise ValueError(
f"La fenêtre [{start_time.isoformat()}, {end_time.isoformat()}) recouvre "
f"{nombre} lecture(s) déjà importée(s) du dataset historique (source='{SOURCE_CSV}') : "
"import refusé pour éviter un doublon inter-source."
)
# Le conflit vise l'index unique uq_reading_source plutôt que la table entière : sans cible # Le conflit vise l'index unique uq_reading_source plutôt que la table entière : sans cible
# nommée, DO NOTHING avalerait aussi une violation de clé primaire. # nommée, DO NOTHING avalerait aussi une violation de clé primaire.
READING_INSERT = text( READING_INSERT = text(
@@ -302,41 +363,57 @@ def build_reading_batch(
return [build_reading_row(reading) for reading in readings] return [build_reading_row(reading) for reading in readings]
def limit_for_window(start_time: datetime, end_time: datetime) -> int:
"""Nombre de lectures à demander pour que l'API Mock en rende une par heure, alignée.
L'API ne renvoie pas un flux à un rythme naturel : elle répartit exactement `limit` lectures,
espacées uniformément, sur toute la fenêtre `[start_time, end_time)` demandée, la première
au tout début de la fenêtre (vérifié empiriquement). Deux façons d'obtenir une lecture
alignée sur l'heure :
- une fenêtre d'exactement N heures (`start_time` sur l'heure) donne, avec `limit=N`, N
lectures espacées d'1h pile, la première à `start_time` : c'est le chemin du backfill
manuel (plusieurs jours d'historique en un seul appel).
- une fenêtre plus courte qu'une heure, ou qui n'est pas un multiple entier d'heure, ne peut
espacer plusieurs lectures d'1h pile (l'espacement de l'API vaut toujours
`durée / limit`) : seule `limit=1` reste alignée, la lecture unique atterrissant à
`start_time`. C'est le chemin du DAG horaire, dont la fenêtre part de l'heure pile qui
précède son déclenchement jusqu'à l'instant du déclenchement lui-même (`:45`), donc plus
courte qu'une heure.
Dans les deux cas, `start_time` doit tomber pile sur l'heure : c'est elle qui ancre
l'alignement, jamais `end_time`. Un `limit` plus grand que celui rendu ici fabriquerait des
lectures infra-horaires, incompatibles avec les lags positionnels de `build_features`.
"""
if start_time.minute or start_time.second or start_time.microsecond:
raise ValueError(
f"La fenêtre doit démarrer pile sur l'heure : {start_time.isoformat()} ne l'est pas."
)
duree = end_time - start_time
heures, reste = divmod(duree.total_seconds(), 3600)
# Fenêtre plus courte qu'une heure, ou pas un multiple entier : aucun `limit` supérieur à 1
# n'espacerait ses lectures d'1h pile (l'espacement vaut toujours durée / limit). Seule la
# lecture unique, ancrée sur `start_time`, reste alignée.
limit = int(heures) if reste == 0 and heures >= 1 else 1
if limit > MAX_LIMIT:
raise ValueError(
f"La fenêtre demandée couvre {limit}h, au-delà du plafond de {MAX_LIMIT} "
"lectures accepté par l'API Mock."
)
return limit
async def import_mock_api_history( async def import_mock_api_history(
start_time: datetime, start_time: datetime,
end_time: datetime, end_time: datetime,
limit: int,
dry_run: bool, dry_run: bool,
) -> None: ) -> None:
settings = get_settings() settings = get_settings()
limit = limit_for_window(start_time, end_time)
async with create_mock_api_client() as client:
sites = await fetch_sites(client)
print(f"Sites récupérés : {len(sites)}")
all_readings: list[dict[str, Any]] = []
for site in sites:
site_id = read_text(site, "site_id")
readings = await fetch_readings(
client=client,
site_id=site_id,
start_time=start_time,
end_time=end_time,
limit=limit,
)
print(f"{site_id}: {len(readings)} lectures")
all_readings.extend(readings)
print(f"Lectures récupérées : {len(all_readings)}")
if dry_run:
print("Dry-run terminé : aucune donnée écrite.")
return
engine = create_async_engine( engine = create_async_engine(
str(settings.database_url), str(settings.database_url),
@@ -344,6 +421,39 @@ async def import_mock_api_history(
) )
try: try:
# Garde-fou d'abord, y compris en dry-run : il est en lecture seule, et annoncer un
# succès pour une fenêtre que l'import réel refusera serait trompeur.
async with engine.connect() as connection:
await refuse_if_overlaps_historical_dataset(connection, start_time, end_time)
async with create_mock_api_client() as client:
sites = await fetch_sites(client)
print(f"Sites récupérés : {len(sites)}")
all_readings: list[dict[str, Any]] = []
for site in sites:
site_id = read_text(site, "site_id")
readings = await fetch_readings(
client=client,
site_id=site_id,
start_time=start_time,
end_time=end_time,
limit=limit,
)
print(f"{site_id}: {len(readings)} lectures")
all_readings.extend(readings)
print(f"Lectures récupérées : {len(all_readings)}")
if dry_run:
print("Dry-run terminé : aucune donnée écrite.")
return
async with engine.begin() as connection: async with engine.begin() as connection:
await upsert_sites( await upsert_sites(
connection, connection,
@@ -365,7 +475,14 @@ async def import_mock_api_history(
def parse_datetime(value: str) -> datetime: def parse_datetime(value: str) -> datetime:
return datetime.fromisoformat(value.replace("Z", "+00:00")) # Sans fuseau, l'API le traite comme reçu, telle quelle, mais l'encodeur `timestamptz`
# d'asyncpg lirait un datetime naif dans le fuseau *local du processus* (correct dans le
# conteneur Airflow en UTC, décalé de 1-2h pour un import manuel lancé depuis un poste en
# Europe/Paris). Poser `tzinfo=UTC` explicitement, même pattern que `_vers_utc()` dans
# `app/services/reading.py`, garantit que la borne envoyée à l'API et celle comparée en SQL
# (refuse_if_overlaps_historical_dataset) désignent le même instant.
instant = datetime.fromisoformat(value.replace("Z", "+00:00"))
return instant if instant.tzinfo is not None else instant.replace(tzinfo=UTC)
def parse_args() -> argparse.Namespace: def parse_args() -> argparse.Namespace:
@@ -383,12 +500,6 @@ def parse_args() -> argparse.Namespace:
type=parse_datetime, type=parse_datetime,
) )
parser.add_argument(
"--limit",
type=int,
default=MAX_LIMIT,
)
parser.add_argument( parser.add_argument(
"--dry-run", "--dry-run",
action="store_true", action="store_true",
@@ -400,9 +511,6 @@ def parse_args() -> argparse.Namespace:
def main() -> None: def main() -> None:
args = parse_args() args = parse_args()
if args.limit < 1 or args.limit > MAX_LIMIT:
raise ValueError(f"--limit doit être compris entre 1 et {MAX_LIMIT}.")
if args.start_time >= args.end_time: if args.start_time >= args.end_time:
raise ValueError("--start-time doit être antérieur à --end-time.") raise ValueError("--start-time doit être antérieur à --end-time.")
@@ -410,7 +518,6 @@ def main() -> None:
import_mock_api_history( import_mock_api_history(
start_time=args.start_time, start_time=args.start_time,
end_time=args.end_time, end_time=args.end_time,
limit=args.limit,
dry_run=args.dry_run, dry_run=args.dry_run,
) )
) )
+216 -68
View File
@@ -1,6 +1,6 @@
import json import json
import sys import sys
from datetime import datetime from datetime import UTC, datetime
from types import SimpleNamespace from types import SimpleNamespace
from typing import Any from typing import Any
from unittest.mock import AsyncMock, MagicMock from unittest.mock import AsyncMock, MagicMock
@@ -110,8 +110,8 @@ async def test_fetch_readings_sends_expected_query_parameters() -> None:
transport = MockTransport(handler) transport = MockTransport(handler)
start_time = datetime.fromisoformat("2024-06-15T12:00:00") start_time = datetime.fromisoformat("2024-06-15T12:00:00+00:00")
end_time = datetime.fromisoformat("2024-06-15T13:00:00") end_time = datetime.fromisoformat("2024-06-15T13:00:00+00:00")
async with AsyncClient( async with AsyncClient(
transport=transport, transport=transport,
@@ -127,11 +127,61 @@ async def test_fetch_readings_sends_expected_query_parameters() -> None:
assert len(readings) == 1 assert len(readings) == 1
assert captured_params["site_id"] == "SITE001" assert captured_params["site_id"] == "SITE001"
assert captured_params["start_time"] == "2024-06-15T12:00:00" assert captured_params["start_time"] == "2024-06-15T12:00:00+00:00"
assert captured_params["end_time"] == "2024-06-15T13:00:00" assert captured_params["end_time"] == "2024-06-15T13:00:00+00:00"
assert captured_params["limit"] == "60" assert captured_params["limit"] == "60"
async def test_fetch_readings_discards_a_reading_outside_the_requested_window() -> None:
# Le garde-fou `refuse_if_overlaps_historical_dataset` ne vérifie que la fenêtre demandée :
# une réponse dont un `timestamp` déborde de `[start_time, end_time)` (bug du mock, ou
# hostile) contournerait ce contrôle si elle atteignait la base telle quelle.
dans_la_fenetre = make_reading()
dans_la_fenetre["timestamp"] = "2024-06-15T12:00:00Z"
hors_fenetre = make_reading()
hors_fenetre["timestamp"] = "2023-01-01T00:00:00Z"
def handler(request: Request) -> Response:
return Response(status_code=200, json=[dans_la_fenetre, hors_fenetre])
async with AsyncClient(
transport=MockTransport(handler),
base_url="https://mock.test",
) as client:
readings = await fetch_readings(
client=client,
site_id="SITE001",
start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"),
end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"),
limit=2,
)
assert readings == [dans_la_fenetre]
async def test_fetch_readings_discards_a_reading_with_an_unparseable_timestamp() -> None:
invalide = make_reading()
invalide["timestamp"] = "pas une date"
def handler(request: Request) -> Response:
return Response(status_code=200, json=[invalide])
async with AsyncClient(
transport=MockTransport(handler),
base_url="https://mock.test",
) as client:
readings = await fetch_readings(
client=client,
site_id="SITE001",
start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"),
end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"),
limit=1,
)
assert readings == []
async def test_fetch_readings_rejects_non_list_response() -> None: async def test_fetch_readings_rejects_non_list_response() -> None:
def handler(request: Request) -> Response: def handler(request: Request) -> Response:
return Response( return Response(
@@ -152,8 +202,8 @@ async def test_fetch_readings_rejects_non_list_response() -> None:
await fetch_readings( await fetch_readings(
client=client, client=client,
site_id="SITE001", site_id="SITE001",
start_time=datetime.fromisoformat("2024-06-15T12:00:00"), start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"),
end_time=datetime.fromisoformat("2024-06-15T13:00:00"), end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"),
limit=60, limit=60,
) )
@@ -175,8 +225,8 @@ async def test_fetch_readings_raises_on_http_error() -> None:
await fetch_readings( await fetch_readings(
client=client, client=client,
site_id="SITE999", site_id="SITE999",
start_time=datetime.fromisoformat("2024-06-15T12:00:00"), start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"),
end_time=datetime.fromisoformat("2024-06-15T13:00:00"), end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"),
limit=60, limit=60,
) )
@@ -357,6 +407,100 @@ async def test_upsert_sites_with_empty_list_does_nothing() -> None:
connection.execute.assert_not_awaited() connection.execute.assert_not_awaited()
async def test_refuse_if_overlaps_historical_dataset_lets_a_clear_window_through() -> None:
connection = AsyncMock()
connection.execute.return_value.scalar_one = MagicMock(return_value=0)
await mock_api_import.refuse_if_overlaps_historical_dataset(
connection,
datetime.fromisoformat("2026-01-01T00:00:00+00:00"),
datetime.fromisoformat("2026-01-01T01:00:00+00:00"),
)
connection.execute.assert_awaited_once()
async def test_refuse_if_overlaps_historical_dataset_rejects_a_window_already_in_the_csv() -> None:
connection = AsyncMock()
connection.execute.return_value.scalar_one = MagicMock(return_value=5)
with pytest.raises(ValueError, match="doublon inter-source"):
await mock_api_import.refuse_if_overlaps_historical_dataset(
connection,
datetime.fromisoformat("2023-06-15T12:00:00+00:00"),
datetime.fromisoformat("2023-06-15T13:00:00+00:00"),
)
def test_limit_for_window_returns_one_per_hour() -> None:
limite = mock_api_import.limit_for_window(
datetime.fromisoformat("2026-09-02T12:00:00+00:00"),
datetime.fromisoformat("2026-09-23T12:00:00+00:00"),
)
assert limite == 21 * 24
def test_limit_for_window_falls_back_to_one_reading_under_an_hour() -> None:
# Le DAG horaire (`:45`) demande desormais [heure pile precedente, instant du declenchement) :
# une fenetre plus courte qu'une heure, dont l'espacement `duree/limit` ne peut jamais valoir
# 1h pile pour plus d'une lecture. Seule `limit=1`, ancree sur `start_time`, reste alignee.
limite = mock_api_import.limit_for_window(
datetime.fromisoformat("2026-09-02T12:00:00+00:00"),
datetime.fromisoformat("2026-09-02T12:45:00+00:00"),
)
assert limite == 1
def test_limit_for_window_falls_back_to_one_reading_for_a_non_whole_hour_span() -> None:
# Meme raisonnement pour une fenetre de plus d'une heure mais qui n'en est pas un multiple
# entier : aucun `limit > 1` ne donnerait un espacement d'1h pile.
limite = mock_api_import.limit_for_window(
datetime.fromisoformat("2026-09-02T12:00:00+00:00"),
datetime.fromisoformat("2026-09-02T13:30:00+00:00"),
)
assert limite == 1
def test_limit_for_window_rejects_a_start_time_not_on_the_hour() -> None:
with pytest.raises(ValueError, match="pile sur l'heure"):
mock_api_import.limit_for_window(
datetime.fromisoformat("2026-09-02T12:05:00+00:00"),
datetime.fromisoformat("2026-09-02T13:05:00+00:00"),
)
def test_limit_for_window_rejects_a_window_above_the_api_cap() -> None:
with pytest.raises(ValueError, match="au-delà du plafond"):
mock_api_import.limit_for_window(
datetime.fromisoformat("2020-01-01T00:00:00+00:00"),
datetime.fromisoformat("2020-03-01T00:00:00+00:00"),
)
def _mock_engine(*, overlap_count: int = 0) -> tuple[MagicMock, AsyncMock]:
"""Engine dont `.connect()` (garde-fou) et `.begin()` (écriture) rendent tous deux la même
connexion, dont `scalar_one()` renvoie `overlap_count` : `import_mock_api_history` ouvre
désormais le garde-fou via `.connect()`, y compris en dry-run."""
connection = AsyncMock()
connection.execute.return_value.scalar_one = MagicMock(return_value=overlap_count)
def _context() -> MagicMock:
context = MagicMock()
context.__aenter__ = AsyncMock(return_value=connection)
context.__aexit__ = AsyncMock(return_value=None)
return context
engine = MagicMock()
engine.connect.return_value = _context()
engine.begin.return_value = _context()
engine.dispose = AsyncMock()
return engine, connection
async def test_import_mock_api_history_dry_run_does_not_write( async def test_import_mock_api_history_dry_run_does_not_write(
monkeypatch: pytest.MonkeyPatch, monkeypatch: pytest.MonkeyPatch,
) -> None: ) -> None:
@@ -396,22 +540,24 @@ async def test_import_mock_api_history_dry_run_does_not_write(
), ),
) )
create_engine_mock = MagicMock() engine, connection = _mock_engine(overlap_count=0)
monkeypatch.setattr( monkeypatch.setattr(
mock_api_import, mock_api_import,
"create_async_engine", "create_async_engine",
create_engine_mock, MagicMock(return_value=engine),
) )
await mock_api_import.import_mock_api_history( await mock_api_import.import_mock_api_history(
start_time=datetime.fromisoformat("2024-06-15T12:00:00"), start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"),
end_time=datetime.fromisoformat("2024-06-15T13:00:00"), end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"),
limit=60,
dry_run=True, dry_run=True,
) )
create_engine_mock.assert_not_called() # Le garde-fou tourne quand même (lecture seule), mais aucune écriture n'a lieu.
connection.execute.assert_awaited_once()
engine.begin.assert_not_called()
engine.dispose.assert_awaited_once()
async def test_import_mock_api_history_loads_data( async def test_import_mock_api_history_loads_data(
@@ -453,23 +599,8 @@ async def test_import_mock_api_history_loads_data(
), ),
) )
connection = AsyncMock() engine, connection = _mock_engine(overlap_count=0)
create_engine_mock = MagicMock(return_value=engine)
transaction_context = MagicMock()
transaction_context.__aenter__ = AsyncMock(
return_value=connection,
)
transaction_context.__aexit__ = AsyncMock(
return_value=None,
)
engine = MagicMock()
engine.begin.return_value = transaction_context
engine.dispose = AsyncMock()
create_engine_mock = MagicMock(
return_value=engine,
)
upsert_sites_mock = AsyncMock() upsert_sites_mock = AsyncMock()
@@ -486,9 +617,8 @@ async def test_import_mock_api_history_loads_data(
) )
await mock_api_import.import_mock_api_history( await mock_api_import.import_mock_api_history(
start_time=datetime.fromisoformat("2024-06-15T12:00:00"), start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"),
end_time=datetime.fromisoformat("2024-06-15T13:00:00"), end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"),
limit=60,
dry_run=False, dry_run=False,
) )
@@ -502,7 +632,39 @@ async def test_import_mock_api_history_loads_data(
[make_site()], [make_site()],
) )
# Un appel pour le garde-fou (via .connect()), un pour READING_INSERT (via .begin()).
assert connection.execute.await_count == 2
dernier_appel = connection.execute.await_args_list[-1]
assert dernier_appel.args[0] is READING_INSERT
engine.dispose.assert_awaited_once()
async def test_import_mock_api_history_refuses_when_it_overlaps_the_historical_dataset(
monkeypatch: pytest.MonkeyPatch,
) -> None:
monkeypatch.setattr(
mock_api_import,
"get_settings",
lambda: SimpleNamespace(database_url="postgresql+asyncpg://test:test@localhost/test"),
)
engine, connection = _mock_engine(overlap_count=3)
monkeypatch.setattr(mock_api_import, "create_async_engine", MagicMock(return_value=engine))
# Le garde-fou tourne avant tout appel à l'API Mock : create_mock_api_client() ne doit
# jamais être invoqué pour une fenêtre refusée.
create_client_mock = MagicMock()
monkeypatch.setattr(mock_api_import, "create_mock_api_client", create_client_mock)
with pytest.raises(ValueError, match="doublon inter-source"):
await mock_api_import.import_mock_api_history(
start_time=datetime.fromisoformat("2023-06-15T12:00:00+00:00"),
end_time=datetime.fromisoformat("2023-06-15T13:00:00+00:00"),
dry_run=False,
)
connection.execute.assert_awaited_once() connection.execute.assert_awaited_once()
create_client_mock.assert_not_called()
engine.dispose.assert_awaited_once() engine.dispose.assert_awaited_once()
@@ -516,6 +678,23 @@ def test_parse_datetime_accepts_z_suffix() -> None:
) )
def test_parse_datetime_attaches_utc_to_a_naive_string() -> None:
# `--start-time`/`--end-time` du DAG sont formatés sans fuseau (Jinja `strftime`) : sans ce
# comportement, l'encodeur `timestamptz` d'asyncpg lirait le datetime naïf dans le fuseau
# *local du processus*, pas UTC, et le garde-fou comparerait une autre fenêtre que celle
# envoyée à l'API.
result = mock_api_import.parse_datetime("2024-06-15T12:00:00")
assert result == datetime.fromisoformat("2024-06-15T12:00:00+00:00")
assert result.tzinfo is UTC
def test_parse_datetime_keeps_a_non_utc_offset_as_is() -> None:
result = mock_api_import.parse_datetime("2024-06-15T12:00:00+02:00")
assert result == datetime.fromisoformat("2024-06-15T12:00:00+02:00")
def test_parse_args_reads_cli_parameters( def test_parse_args_reads_cli_parameters(
monkeypatch: pytest.MonkeyPatch, monkeypatch: pytest.MonkeyPatch,
) -> None: ) -> None:
@@ -528,8 +707,6 @@ def test_parse_args_reads_cli_parameters(
"2024-06-15T12:00:00Z", "2024-06-15T12:00:00Z",
"--end-time", "--end-time",
"2024-06-15T13:00:00Z", "2024-06-15T13:00:00Z",
"--limit",
"60",
"--dry-run", "--dry-run",
], ],
) )
@@ -542,34 +719,9 @@ def test_parse_args_reads_cli_parameters(
assert args.end_time == datetime.fromisoformat( assert args.end_time == datetime.fromisoformat(
"2024-06-15T13:00:00+00:00", "2024-06-15T13:00:00+00:00",
) )
assert args.limit == 60
assert args.dry_run is True assert args.dry_run is True
def test_main_rejects_limit_out_of_bounds(
monkeypatch: pytest.MonkeyPatch,
) -> None:
monkeypatch.setattr(
sys,
"argv",
[
"mock_api_import",
"--start-time",
"2024-06-15T12:00:00Z",
"--end-time",
"2024-06-15T13:00:00Z",
"--limit",
"0",
],
)
with pytest.raises(
ValueError,
match="--limit doit être compris entre 1 et 1000",
):
mock_api_import.main()
def test_main_rejects_invalid_period( def test_main_rejects_invalid_period(
monkeypatch: pytest.MonkeyPatch, monkeypatch: pytest.MonkeyPatch,
) -> None: ) -> None:
@@ -582,8 +734,6 @@ def test_main_rejects_invalid_period(
"2024-06-15T14:00:00Z", "2024-06-15T14:00:00Z",
"--end-time", "--end-time",
"2024-06-15T13:00:00Z", "2024-06-15T13:00:00Z",
"--limit",
"60",
], ],
) )
@@ -612,7 +762,6 @@ def test_main_runs_import(
lambda: SimpleNamespace( lambda: SimpleNamespace(
start_time=start_time, start_time=start_time,
end_time=end_time, end_time=end_time,
limit=60,
dry_run=True, dry_run=True,
), ),
) )
@@ -628,7 +777,6 @@ def test_main_runs_import(
import_mock.assert_awaited_once_with( import_mock.assert_awaited_once_with(
start_time=start_time, start_time=start_time,
end_time=end_time, end_time=end_time,
limit=60,
dry_run=True, dry_run=True,
) )
@@ -673,8 +821,8 @@ async def test_fetch_readings_rejects_a_response_above_the_requested_limit() ->
await fetch_readings( await fetch_readings(
client=client, client=client,
site_id="SITE001", site_id="SITE001",
start_time=datetime.fromisoformat("2024-06-15T12:00:00"), start_time=datetime.fromisoformat("2024-06-15T12:00:00+00:00"),
end_time=datetime.fromisoformat("2024-06-15T13:00:00"), end_time=datetime.fromisoformat("2024-06-15T13:00:00+00:00"),
limit=2, limit=2,
) )
+13 -6
View File
@@ -70,11 +70,17 @@ Le lien `front -.-> api` reste en pointillé : le frontend appelle bien une API,
intercepteur répond à sa place tant que les endpoints n'existent pas. Voir intercepteur répond à sa place tant que les endpoints n'existent pas. Voir
[30-frontend.md](30-frontend.md). [30-frontend.md](30-frontend.md).
Le lien `airflow --> db` est maintenant en trait plein : cinq DAGs tournent, deux pour Le lien `airflow --> db` est maintenant en trait plein : six DAGs tournent, deux pour
l'entraînement et le scoring du modèle ML (issue #115), un pour la détection d'alertes et la 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 génération des recommandations (issue #116), un pour la surveillance de dérive (issue #45),
(issue #119) et `mock_api_import` pour l'ingestion horaire de l'API Mock (issue #15). `historical_import` pour le dataset historique (issue #119) et `mock_api_import` pour l'ingestion
La réconciliation globale des données provenant des deux sources reste à compléter dans l'issue #15. horaire de l'API Mock (issue #15). La réconciliation entre les deux sources de lectures (issue
#15) est tranchée : le trou entre la fin de l'historique (31/12/2024) et le début de l'ingestion
API Mock est accepté comme définitivement perdu, aucune mesure réelle n'existant pour cette
période. `mock_api_import` refuse toute fenêtre qui recouvrirait des lectures déjà importées du
CSV plutôt que de laisser les deux sources dupliquer silencieusement un même instant, et le
pipeline ML déduplique par construction (`DISTINCT ON`, source `csv` préférée) au cas où un
recouvrement se produirait malgré tout, voir [40-data.md](40-data.md).
Les liens de la supervision sont en trait plein depuis le 23/09 (issue #26) : Prometheus scrute Les liens de la supervision sont en trait plein depuis le 23/09 (issue #26) : Prometheus scrute
`/metrics` avec un jeton, Grafana lit Prometheus et, par un rôle en lecture seule, les tables `/metrics` avec un jeton, Grafana lit Prometheus et, par un rôle en lecture seule, les tables
@@ -92,7 +98,7 @@ ailleurs ([ADR 0016](../adr/0016-supervision-en-profil-compose.md),
| 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 livrée côté backend (`app.monitoring.drift`, table `drift_report`, `GET /monitoring/drift`, DAG `derive`), voir [ADR 0013](../adr/0013-surveillance-de-derive-dans-le-backend.md) | | 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 livrée côté backend (`app.monitoring.drift`, table `drift_report`, `GET /monitoring/drift`, DAG `derive`), voir [ADR 0013](../adr/0013-surveillance-de-derive-dans-le-backend.md) |
| 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 | | 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` | `Fait` | Profil Compose `monitoring`, actif en prod : Prometheus et trois exporteurs (PostgreSQL, hôte, conteneurs), neuf règles d'alerte testées par `promtool`, Alertmanager vers Mailpit, trois tableaux de bord Grafana provisionnés. Voir [60-observabilite.md](60-observabilite.md) | | Monitoring | Prometheus, Grafana, Alertmanager | `monitoring` | `Fait` | Profil Compose `monitoring`, actif en prod : Prometheus et trois exporteurs (PostgreSQL, hôte, conteneurs), neuf règles d'alerte testées par `promtool`, Alertmanager vers Mailpit, trois tableaux de bord Grafana provisionnés. Voir [60-observabilite.md](60-observabilite.md) |
| ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver et scheduler avec LocalExecutor via Docker Compose, sur une base PostgreSQL dédiée. Six DAGs en sous-processus `uv run` : `ml_train`, `ml_score`, `alertes`, `historical_import`, `mock_api_import` et `derive` (quotidien, surveillance de dérive). L'import historique reste manuel et l'import API Mock s'exécute chaque heure. La réconciliation globale des deux sources reste à compléter dans l'issue #15. | | ETL | Apache Airflow | `etl/airflow` | `En cours` | Webserver et scheduler avec LocalExecutor via Docker Compose, sur une base PostgreSQL dédiée. Six DAGs en sous-processus `uv run` : `ml_train`, `ml_score`, `alertes`, `historical_import`, `mock_api_import` et `derive` (quotidien, surveillance de dérive). L'import historique reste manuel et l'import API Mock s'exécute chaque heure. Réconciliation entre les deux sources (issue #15) : trou temporel accepté, recouvrement refusé à l'ingestion et dédupliqué en défense côté ML, voir [40-data.md](40-data.md). |
| CI/CD | GitHub Actions | `.github/workflows` | `En cours` | Un orchestrateur `ci.yml` qui n'appelle que les composants modifiés ([ADR 0014](../adr/0014-pipeline-ci-unique-et-deploiement-conditionne.md)) : 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, Terraform, Compose et supervision, parcours Playwright et tirs k6 contre la stack de prod ([ADR 0015](../adr/0015-tests-e2e-et-de-charge-contre-la-stack-compose.md)). Déploiement vers la VM ENI par `deploy.yml`, appelé une fois « CI ok » vert, `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é : le runner n'est pas enregistré sur la machine. Détail dans [50-cicd.md](50-cicd.md) | | CI/CD | GitHub Actions | `.github/workflows` | `En cours` | Un orchestrateur `ci.yml` qui n'appelle que les composants modifiés ([ADR 0014](../adr/0014-pipeline-ci-unique-et-deploiement-conditionne.md)) : 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, Terraform, Compose et supervision, parcours Playwright et tirs k6 contre la stack de prod ([ADR 0015](../adr/0015-tests-e2e-et-de-charge-contre-la-stack-compose.md)). Déploiement vers la VM ENI par `deploy.yml`, appelé une fois « CI ok » vert, `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é : le runner n'est pas enregistré sur la machine. Détail dans [50-cicd.md](50-cicd.md) |
## Flux bout en bout ## Flux bout en bout
@@ -102,7 +108,8 @@ Statut : `En cours`. **Le chemin de lecture tourne** entre la base, l'API et le
dataset CSV/JSON sur déclenchement manuel et `mock_api_import` collecte chaque heure les mesures dataset CSV/JSON sur déclenchement manuel et `mock_api_import` collecte chaque heure les mesures
de l'API Mock. Les DAGs `ml_train` et `ml_score` (issue #115), `alertes` (issue #116) et `derive` de l'API Mock. Les DAGs `ml_train` et `ml_score` (issue #115), `alertes` (issue #116) et `derive`
(issue #45) portent le pipeline ML, la détection d'alertes et la surveillance de dérive. La (issue #45) portent le pipeline ML, la détection d'alertes et la surveillance de dérive. La
réconciliation globale des données provenant des deux sources reste à compléter dans l'issue #15. réconciliation entre les deux sources de lectures (issue #15) est close : voir
[40-data.md](40-data.md) pour le détail du garde-fou d'ingestion et de la déduplication ML.
```mermaid ```mermaid
sequenceDiagram sequenceDiagram
+4 -2
View File
@@ -99,8 +99,10 @@ modifier.
Le DAG `mock_api_import` exécute le pipeline API Mock toutes les heures, à la minute `:45`. 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 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 déclenchement manuel, mais la fenêtre transmise au script backend part de l'heure pile qui
les tables communes `site` et `reading`. Le décalage à `:45` laisse quinze minutes avant le précède le déclenchement (pas de l'intervalle Airflow tel quel), pour que la mesure importée
tombe à :00 et non à :45, voir [40-data.md](40-data.md). Le pipeline charge la mesure 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`. 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. `max_active_runs=1` empêche deux exécutions du DAG de se chevaucher.
+67 -2
View File
@@ -437,10 +437,29 @@ Les paramètres de ligne de commande disponibles pour l'import sont :
```text ```text
--start-time --start-time
--end-time --end-time
--limit
--dry-run --dry-run
``` ```
**Piège sur `limit`, corrigé dans le code plutôt que documenté** : l'API ne renvoie pas un flux à
un rythme naturel, elle répartit exactement `limit` lectures, espacées uniformément, sur toute la
fenêtre `[start_time, end_time)` demandée, la première au tout début de la fenêtre (vérifié
empiriquement en interrogeant directement l'API). Une fenêtre d'une heure avec `limit=1000`, le
réglage d'origine, renvoyait donc 1000 lectures espacées de 3,6 secondes à l'intérieur de cette
heure, pas une lecture horaire, incompatible avec les lags positionnels de `build_features`.
Plutôt que documenter la règle « `limit` = nombre d'heures de la fenêtre » et compter sur chaque
appelant pour la respecter, `limit_for_window()` la porte : `import_mock_api_history()` calcule
`limit` depuis la fenêtre reçue, refuse une fenêtre dont `start_time` ne tombe pas pile sur
l'heure (c'est elle qui ancre l'alignement), et refuse un intervalle de plus de 1000 heures (le
plafond `limit` de l'API). `--limit` n'existe donc plus côté CLI. Deux formes de fenêtre sont
gérées : un multiple entier d'heures (`limit` = ce nombre d'heures, une lecture par heure
espacée d'1h pile, chemin du backfill manuel) ou une fenêtre plus courte qu'une heure, ou qui
n'en est pas un multiple entier (`limit=1`, seule valeur qui reste alignée quand l'espacement
`durée / limit` ne peut valoir 1h pile). Le DAG `mock_api_import` est dans ce second cas : il
demande la fenêtre `[heure pile précédant le déclenchement, instant du déclenchement)`, plus
courte qu'une heure, plutôt que l'intervalle Airflow `[data_interval_start, data_interval_end)`
tel quel (`[:45, :45)`) qui aurait placé l'unique lecture à :45, hors de la grille horaire du
reste du schéma.
### Flux d'ingestion API Mock ### Flux d'ingestion API Mock
```text ```text
@@ -492,7 +511,7 @@ réponse est donc traitée comme une entrée hostile, conformément à API10 dan
[la traçabilité OWASP](owasp-traceabilite.md). Le risque premier n'est pas la fausse alerte, [la traçabilité OWASP](owasp-traceabilite.md). Le risque premier n'est pas la fausse alerte,
c'est l'empoisonnement du jeu d'entraînement du modèle de prédiction. c'est l'empoisonnement du jeu d'entraînement du modèle de prédiction.
Quatre garde-fous, tous dans `mock_api_import.py` : Cinq garde-fous, tous dans `mock_api_import.py` :
| Garde-fou | Mise en œuvre | | Garde-fou | Mise en œuvre |
|---|---| |---|---|
@@ -500,6 +519,7 @@ Quatre garde-fous, tous dans `mock_api_import.py` :
| Taille de tableau plafonnée | `MAX_SITES` sites, et au plus `--limit` mesures par site | | Taille de tableau plafonnée | `MAX_SITES` sites, et au plus `--limit` mesures par site |
| Bornes physiques | `PHYSICAL_BOUNDS`, une plage par grandeur | | Bornes physiques | `PHYSICAL_BOUNDS`, une plage par grandeur |
| Frontière d'anti-corruption | `build_site_row()` et `build_reading_row()`, qui ne recopient que les champs attendus | | Frontière d'anti-corruption | `build_site_row()` et `build_reading_row()`, qui ne recopient que les champs attendus |
| Refus de recouvrir l'historique | `refuse_if_overlaps_historical_dataset()`, voir ci-dessous |
Une valeur hors bornes, d'un type inattendu, `NaN` ou infinie devient `NULL`. Elle laisse sa Une valeur hors bornes, d'un type inattendu, `NaN` ou infinie devient `NULL`. Elle laisse sa
trace dans `null_reasons` sous la forme `out_of_physical_bounds:<colonne>`, et `data_quality` trace dans `null_reasons` sous la forme `out_of_physical_bounds:<colonne>`, et `data_quality`
@@ -510,6 +530,51 @@ d'origine intacte : rien n'est perdu, seule son exploitation est bornée.
Le plafond de taille s'applique après désérialisation de la réponse. Borner le corps HTTP Le plafond de taille s'applique après désérialisation de la réponse. Borner le corps HTTP
lui-même demanderait une lecture en flux, et reste à faire. lui-même demanderait une lecture en flux, et reste à faire.
### Réconciliation entre les deux sources (issue #15)
`historical_import` (source `csv`) et `mock_api_import` (source `api_history`) écrivent toutes
deux dans `reading`. Trois décisions ferment cette réconciliation :
- **Le trou temporel est accepté.** Le dataset historique s'arrête au 31/12/2024, et
`mock_api_import` n'importe que l'heure précédant chaque déclenchement : rien ne comble
automatiquement la période intermédiaire, et rien ne le pourra jamais, aucune mesure réelle
n'existe pour ces instants. Conséquence pour le ML, pas nouvelle mais que ce trou rend
définitive : `build_features()` calcule ses lags par `shift(n)` positionnel, et `train.py`
n'écarte que les lignes où `lag_168h` est `NaN`. Pour un site présent dans les deux sources, les
168 premières lectures `api_history` qui suivent le trou héritent donc de lags et de moyennes
glissantes calculés sur décembre 2024 (et tant que l'ingestion a moins de 7 jours, c'est le cas
de toutes les lectures). Même effet, plus ponctuel, pour chaque heure que le DAG manque
(`mock_api_import` en échec, Airflow arrêté). Aucun garde-fou ne détecte aujourd'hui un lag
calculé sur un écart réel différent de celui attendu ; issue de suivi à ouvrir.
- **Le recouvrement est refusé à l'ingestion.** `uq_reading_source` autorise deux lignes au même
`(site_id, timestamp)` dès que `source` diffère : rien dans le schéma n'empêche donc un import
Mock API manuel avec une fenêtre passée (le script accepte `--start-time`/`--end-time`
arbitraires) de dupliquer un point déjà couvert par le CSV. `import_mock_api_history()` appelle
`refuse_if_overlaps_historical_dataset()` avant toute écriture, y compris en `--dry-run` (le
contrôle est en lecture seule) et avant le moindre appel à l'API Mock : si la fenêtre demandée
recouvre au moins une lecture `source='csv'`, l'import est refusé (`ValueError`) plutôt que
d'écrire un doublon inter-source silencieux. Le contrôle ne porte que sur la fenêtre demandée,
pas sur les lectures reçues : `fetch_readings()` écarte donc toute lecture dont le `timestamp`
déborde de `[start_time, end_time)`, pour qu'une réponse hors fenêtre (bug du mock, ou hostile)
ne puisse pas le contourner. Ce contrôle compare des instants, pas des chaînes : `parse_datetime()`
pose `tzinfo=UTC` sur une entrée sans fuseau (même pattern que `_vers_utc()` dans
`app/services/reading.py`), sans quoi l'encodeur `timestamptz` d'asyncpg lirait un datetime naïf
dans le fuseau local du **processus**, correct dans le conteneur Airflow (UTC) mais décalé pour
un import manuel lancé depuis un poste en Europe/Paris.
- **Le pipeline ML déduplique en défense.** Le garde-fou ci-dessus protège l'ingestion, pas
la lecture : si un recouvrement se produisait malgré tout (import direct en base, contournement
du script), `ml/enervision_ml/data.py` ne doit pas casser silencieusement l'hypothèse de
`build_features` (« une ligne par `(site_id, timestamp)` »). `load_from_database()` et
`load_recent_from_database()` utilisent donc `SELECT DISTINCT ON (site_id, timestamp)`, `source
= 'csv'` gagnant sur `'api_history'` en cas d'égalité, l'historique étant une source vérifiée,
l'API Mock une entrée hostile (cf. ci-dessus). **Cette préférence est spécifique au chargeur
ML.** `GET /readings` renvoie les deux lignes sans les fusionner, et `DriftRepository` /
`ReadingRepository.latest_by_site()` / `.latest_for_site()` départagent par `reading_id` le plus
grand (en pratique la ligne insérée en dernier, pas forcément `csv`) : en cas de recouvrement, la
dérive comparerait alors une prévision à une valeur différente de celle sur laquelle le modèle a
appris. Pas d'incohérence aujourd'hui tant que le recouvrement reste refusé à l'ingestion ; à
aligner si ce garde-fou devait un jour être contourné.
### Qualité des données de l'API Mock ### Qualité des données de l'API Mock
Les valeurs `NULL` ne sont pas remplacées pendant l'ingestion. Les valeurs `NULL` ne sont pas remplacées pendant l'ingestion.
+12 -11
View File
@@ -668,17 +668,18 @@ le pipeline ML (`ml_train` et `ml_score`, issue #115), la détection d'alertes e
des recommandations (`alertes`, issue #116), l'import historique (`historical_import`, 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). issue #119) et l'import périodique de l'API Mock (`mock_api_import`, issue #15).
Le DAG `mock_api_import` s'exécute chaque heure, à la minute `:45`. Il appelle Le DAG `mock_api_import` s'exécute chaque heure, à la minute `:45`, sur une fenêtre qui part de
`app.etl.mock_api_import` sur l'intervalle qui va de l'heure pile à son déclenchement, avec une l'heure pile précédant son déclenchement jusqu'à l'instant du déclenchement lui-même (pas
limite d'une lecture par site. L'API Mock génère autant de points que la limite demandée, l'intervalle Airflow `data_interval_start`/`end` tel quel). L'API Mock génère autant de points que
répartis sur l'intervalle et le premier à son début : un seul donne la mesure de :00, au pas la limite demandée, répartis sur la fenêtre et le premier à son début :
horaire du dataset historique, que les features ML supposent en décalant par ligne. Les deux pipelines normalisent leurs données vers les tables communes `site` et `app.etl.mock_api_import.limit_for_window()` dérive donc `limit` de la fenêtre reçue (une seule
`reading`, tout en conservant leur source (`csv` ou `api_history`). La réconciliation globale lecture ici, ancrée sur l'heure pile) plutôt que de dépendre d'une valeur fixée à la main côté
des deux sources reste à compléter dans l'issue #15. DAG, et refuse une fenêtre qui ne démarre pas pile sur l'heure. Une fenêtre calée sur l'intervalle
Airflow tel quel (`[:45, :45)`) placerait cette lecture à :45, hors de la grille horaire du reste
Le DAG `mock_api_import` exécute `app.etl.mock_api_import` toutes les heures. Chaque exécution du schéma (vérifié empiriquement contre l'API Mock) ; partir de l'heure pile évite ce décalage.
importe la mesure de l'heure pile qui précède son déclenchement. Les deux pipelines normalisent leurs données vers les Les deux pipelines normalisent leurs données vers les tables communes `site` et `reading`, tout en
tables communes `site` et `reading`, tout en conservant leur source (`csv` ou `api_history`). conservant leur source (`csv` ou `api_history`). La réconciliation entre les deux sources
(issue #15) est close : voir `docs/architecture/40-data.md`.
Airflow permet de planifier les traitements, gérer leur ordre d'exécution, suivre leur état et remonter les erreurs. Il ne remplace pas la logique ETL Python existante : les scripts actuels restent responsables de l'extraction, de la validation, de la transformation et du chargement. `etl/airflow/dags/ml_train.py`, `ml_score.py`, `alertes.py`, `historical_import.py` et Airflow permet de planifier les traitements, gérer leur ordre d'exécution, suivre leur état et remonter les erreurs. Il ne remplace pas la logique ETL Python existante : les scripts actuels restent responsables de l'extraction, de la validation, de la transformation et du chargement. `etl/airflow/dags/ml_train.py`, `ml_score.py`, `alertes.py`, `historical_import.py` et
`mock_api_import.py` montrent le patron retenu (des `BashOperator` qui invoquent le script tel quel, dans l'environnement `uv` que l'image embarque pour lui). `mock_api_import.py` montrent le patron retenu (des `BashOperator` qui invoquent le script tel quel, dans l'environnement `uv` que l'image embarque pour lui).
+14 -7
View File
@@ -18,10 +18,6 @@ from airflow.timetables.trigger import CronTriggerTimetable
# Le backend possède son propre environnement uv dans l'image Airflow (ADR 0008). # 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" COMMANDE_BACKEND = "cd /opt/backend && env -u VIRTUAL_ENV uv run --no-sync python -m"
# Contrainte : l'API Mock génère `limit` points répartis sur l'intervalle, le premier à son début.
# Un seul, depuis l'heure pile, donne la mesure de :00 au pas du CSV que suppose `shift(168)`.
LIMITE_LECTURES = 1
# Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste # Deux reprises donnent trois tentatives au total. Même dans le pire cas, l'exécution reste
# inférieure au pas horaire du DAG. # inférieure au pas horaire du DAG.
NOMBRE_REPRISES = 2 NOMBRE_REPRISES = 2
@@ -30,7 +26,11 @@ PLAFOND_PAR_TENTATIVE = timedelta(minutes=10)
# L'intervalle est déclaré explicitement pour ne pas dépendre de la valeur du paramètre Airflow # 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`, # `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. # exécuté à l'heure pile, puis avant `alertes`, exécuté à :15. La fenêtre demandée à l'API Mock
# (voir `bash_command` ci-dessous) ne suit pas cet intervalle Airflow tel quel : elle part de
# l'heure pile qui précède le déclenchement, pas de `data_interval_start`, pour que l'unique
# lecture demandée (`app.etl.mock_api_import.limit_for_window()`) atterrisse à :00 et non à :45
# (vérifié empiriquement sur l'API Mock), au pas horaire du reste du schéma, cf. 40-data.md.
PLANIFICATION = CronTriggerTimetable( PLANIFICATION = CronTriggerTimetable(
"45 * * * *", "45 * * * *",
timezone="UTC", timezone="UTC",
@@ -51,9 +51,16 @@ with DAG(
task_id="import_mock_api", task_id="import_mock_api",
bash_command=( bash_command=(
f"{COMMANDE_BACKEND} app.etl.mock_api_import " f"{COMMANDE_BACKEND} app.etl.mock_api_import "
# `--start-time` part de l'heure pile qui précède le déclenchement, pas de
# `data_interval_start` : sur `[:45, :45)`, l'API aurait placé son unique lecture
# à :45, hors de la grille horaire du reste du schéma (vérifié empiriquement).
"--start-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:00:00') }}\" " "--start-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:00:00') }}\" "
"--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\" " "--end-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:%M:%S') }}\""
f"--limit {LIMITE_LECTURES}" # Pas de --limit : app.etl.mock_api_import.limit_for_window() le dérive de la
# fenêtre (ici plus courte qu'une heure, donc une seule lecture, ancrée sur
# --start-time) et refuse une fenêtre qui ne démarre pas pile sur l'heure. Porter
# la règle dans le code, pas dans ce DAG, évite qu'un appel manuel oublie de la
# respecter.
), ),
retries=NOMBRE_REPRISES, retries=NOMBRE_REPRISES,
retry_delay=DELAI_ENTRE_REPRISES, retry_delay=DELAI_ENTRE_REPRISES,
+2 -1
View File
@@ -118,7 +118,8 @@ def test_mock_api_import_asks_for_the_on_the_hour_reading(dagbag: DagBag) -> Non
assert "--start-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:00:00') }}\"" in commande assert "--start-time \"{{ data_interval_end.strftime('%Y-%m-%dT%H:00:00') }}\"" in commande
assert "--end-time \"{{ data_interval_end.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 commande.endswith("--limit 1") # Pas de --limit : app.etl.mock_api_import.limit_for_window() le dérive de la fenêtre.
assert "--limit" not in commande
@pytest.mark.parametrize("task_id", ["detection", "recommandations"]) @pytest.mark.parametrize("task_id", ["detection", "recommandations"])
+10 -4
View File
@@ -47,9 +47,15 @@ NUMERIC_COLUMNS = [
# `bool` : `astype(bool)` ferait un `True` d'une absence, et les deux chargeurs divergeraient. # `bool` : `astype(bool)` ferait un `True` d'une absence, et les deux chargeurs divergeraient.
FLAG_COLUMNS = ["is_working_hours"] FLAG_COLUMNS = ["is_working_hours"]
# `uq_reading_source` autorise deux lignes au meme (site_id, timestamp) des que `source` differe
# (cf. `app/etl/mock_api_import.py`, qui refuse desormais d'importer une fenetre deja couverte par
# le CSV, mais ne protege pas le sens inverse). `build_features` suppose une ligne par
# (site_id, timestamp) sans doublon : le `DISTINCT ON` l'impose plutot que de la supposer.
# 'csv' gagne sur 'api_history' en cas de recouvrement, l'historique etant une source verifiee
# alors que l'API Mock est traitee comme une entree hostile (cf. OWASP API10).
_READING_QUERY = text( _READING_QUERY = text(
""" """
SELECT SELECT DISTINCT ON (r.site_id, r.timestamp)
r.site_id, r.site_id,
r.timestamp, r.timestamp,
r.consumption_kwh, r.consumption_kwh,
@@ -61,14 +67,14 @@ _READING_QUERY = text(
s.capacity_kw s.capacity_kw
FROM reading r FROM reading r
JOIN site s ON s.site_id = r.site_id JOIN site s ON s.site_id = r.site_id
ORDER BY r.site_id, r.timestamp ORDER BY r.site_id, r.timestamp, (r.source = 'csv') DESC, r.reading_id DESC
""" """
) )
_RECENT_READING_QUERY = text( _RECENT_READING_QUERY = text(
""" """
SELECT SELECT DISTINCT ON (r.site_id, r.timestamp)
r.site_id, r.site_id,
r.timestamp, r.timestamp,
r.consumption_kwh, r.consumption_kwh,
@@ -81,7 +87,7 @@ _RECENT_READING_QUERY = text(
FROM reading r FROM reading r
JOIN site s ON s.site_id = r.site_id JOIN site s ON s.site_id = r.site_id
WHERE r.timestamp >= :since AND r.timestamp <= :until WHERE r.timestamp >= :since AND r.timestamp <= :until
ORDER BY r.site_id, r.timestamp ORDER BY r.site_id, r.timestamp, (r.source = 'csv') DESC, r.reading_id DESC
""" """
) )
+38 -5
View File
@@ -16,7 +16,7 @@ from collections.abc import Iterator
from dataclasses import dataclass, field from dataclasses import dataclass, field
from datetime import UTC, datetime, timedelta from datetime import UTC, datetime, timedelta
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any, cast
from uuid import uuid4 from uuid import uuid4
import lightgbm as lgb import lightgbm as lgb
@@ -44,20 +44,30 @@ _INSERT_SITE = text(
""" """
) )
# `source = 'api_history'` impose `dataset_id IS NULL` (ck_reading_dataset_source), ce qui evite # `source = 'api_history'` impose `dataset_id IS NULL` (ck_reading_dataset_source) : le defaut
# de creer une ligne `dataset`. `raw_data` est NOT NULL, d'ou le litteral jsonb. # `dataset_id=None` evite de creer une ligne `dataset` pour la plupart des tests. `source='csv'`
# impose l'inverse, d'ou `insere_dataset()` quand un test a besoin de cette source precise.
# `raw_data` est NOT NULL, d'ou le litteral jsonb.
_INSERT_READING = text( _INSERT_READING = text(
""" """
INSERT INTO reading ( INSERT INTO reading (
site_id, timestamp, source, consumption_kwh, temperature_celsius, site_id, timestamp, source, dataset_id, consumption_kwh, temperature_celsius,
humidity_percent, solar_irradiance_wm2, is_working_hours, raw_data humidity_percent, solar_irradiance_wm2, is_working_hours, raw_data
) VALUES ( ) VALUES (
:site_id, :timestamp, :source, :consumption_kwh, :temperature_celsius, :site_id, :timestamp, :source, :dataset_id, :consumption_kwh, :temperature_celsius,
:humidity_percent, :solar_irradiance_wm2, :is_working_hours, '{}'::jsonb :humidity_percent, :solar_irradiance_wm2, :is_working_hours, '{}'::jsonb
) )
""" """
) )
_INSERT_DATASET = text(
"""
INSERT INTO dataset (dataset_name, archive_sha256, storage_uri, source_timezone, metadata)
VALUES (:dataset_name, :archive_sha256, :storage_uri, 'UTC', '{}'::jsonb)
RETURNING dataset_id
"""
)
_SELECT_PREDICTIONS = text( _SELECT_PREDICTIONS = text(
""" """
SELECT target_at, predicted_value, status, failure_reason, model_reference SELECT target_at, predicted_value, status, failure_reason, model_reference
@@ -111,6 +121,25 @@ def insere_site(
return site_id return site_id
def insere_dataset(connexion: Connection) -> int:
"""Ligne `dataset` minimale, requise pour inserer une lecture `source='csv'`
(`ck_reading_dataset_source` impose `dataset_id IS NOT NULL` pour cette seule source).
"""
marque = uuid4().hex
return cast(
int,
connexion.execute(
_INSERT_DATASET,
{
"dataset_name": f"jeu de test {marque}",
"archive_sha256": marque.rjust(64, "0"),
"storage_uri": f"file:///test/{marque}.csv",
},
).scalar_one(),
)
def insere_lectures( def insere_lectures(
connexion: Connection, connexion: Connection,
site_id: str, site_id: str,
@@ -119,6 +148,7 @@ def insere_lectures(
fin: datetime, fin: datetime,
valeur: float = 50.0, valeur: float = 50.0,
source: str = "api_history", source: str = "api_history",
dataset_id: int | None = None,
is_working_hours: bool | None = True, is_working_hours: bool | None = True,
) -> list[datetime]: ) -> list[datetime]:
"""Grille horaire contigue finissant a `fin`, incluse. """Grille horaire contigue finissant a `fin`, incluse.
@@ -134,6 +164,7 @@ def insere_lectures(
"site_id": site_id, "site_id": site_id,
"timestamp": instant, "timestamp": instant,
"source": source, "source": source,
"dataset_id": dataset_id,
"consumption_kwh": valeur + math.sin(rang / 12.0) * 10.0, "consumption_kwh": valeur + math.sin(rang / 12.0) * 10.0,
"temperature_celsius": 15.0, "temperature_celsius": 15.0,
"humidity_percent": 50.0, "humidity_percent": 50.0,
@@ -153,6 +184,7 @@ def insere_lecture(
instant: datetime, instant: datetime,
consumption_kwh: float | None = 50.0, consumption_kwh: float | None = 50.0,
source: str = "api_history", source: str = "api_history",
dataset_id: int | None = None,
is_working_hours: bool | None = True, is_working_hours: bool | None = True,
) -> None: ) -> None:
"""Une lecture isolee, quand le test pilote sa valeur plutot que sa forme.""" """Une lecture isolee, quand le test pilote sa valeur plutot que sa forme."""
@@ -162,6 +194,7 @@ def insere_lecture(
"site_id": site_id, "site_id": site_id,
"timestamp": instant, "timestamp": instant,
"source": source, "source": source,
"dataset_id": dataset_id,
"consumption_kwh": consumption_kwh, "consumption_kwh": consumption_kwh,
"temperature_celsius": 15.0, "temperature_celsius": 15.0,
"humidity_percent": 50.0, "humidity_percent": 50.0,
+60 -1
View File
@@ -11,7 +11,7 @@ from enervision_ml.data import (
load_from_database, load_from_database,
load_recent_from_database, load_recent_from_database,
) )
from tests.conftest import ANCRAGE, insere_lecture, insere_lectures, insere_site from tests.conftest import ANCRAGE, insere_dataset, insere_lecture, insere_lectures, insere_site
pytestmark = pytest.mark.integration pytestmark = pytest.mark.integration
@@ -39,6 +39,65 @@ def test_load_from_database_joins_the_site_attributes_to_every_reading(
assert set(mien["capacity_kw"]) == {250.0} assert set(mien["capacity_kw"]) == {250.0}
def test_load_from_database_deduplicates_two_sources_at_the_same_instant(
connexion_ml: Connection,
) -> None:
# `uq_reading_source` autorise deux lignes au meme (site_id, timestamp) des que `source`
# differe : le garde-fou vit dans `mock_api_import.py`, pas dans le schema. Le chargeur ML
# doit donc imposer lui-meme "une ligne par (site_id, timestamp)", pas la supposer.
#
# `csv` est inseree en premier (reading_id le plus bas) et `api_history` en second (le plus
# haut) : un depart par `reading_id DESC` seul choisirait `api_history` a tort. Seule la
# preference explicite pour `source='csv'` fait gagner le bon reading_id ici, et le test
# cesserait de proteger cette regle si l'ordre d'insertion etait inverse.
site_id = insere_site(connexion_ml)
dataset_id = insere_dataset(connexion_ml)
insere_lecture(
connexion_ml,
site_id,
instant=ANCRAGE,
consumption_kwh=99.0,
source="csv",
dataset_id=dataset_id,
)
insere_lecture(
connexion_ml, site_id, instant=ANCRAGE, consumption_kwh=10.0, source="api_history"
)
frame = load_from_database(connexion_ml)
mien = frame[frame["site_id"] == site_id]
assert len(mien) == 1
assert mien["consumption_kwh"].iloc[0] == 99.0
def test_load_recent_from_database_prefers_csv_when_two_sources_share_an_instant(
connexion_ml: Connection,
) -> None:
# Meme ordre d'insertion que ci-dessus, et pour la meme raison : `csv` doit gagner malgre un
# `reading_id` plus bas que celui d'`api_history`.
site_id = insere_site(connexion_ml)
dataset_id = insere_dataset(connexion_ml)
insere_lecture(
connexion_ml,
site_id,
instant=ANCRAGE,
consumption_kwh=99.0,
source="csv",
dataset_id=dataset_id,
)
insere_lecture(
connexion_ml, site_id, instant=ANCRAGE, consumption_kwh=10.0, source="api_history"
)
frame = load_recent_from_database(
connexion_ml, since=ANCRAGE, until=ANCRAGE + timedelta(hours=3)
)
assert len(frame) == 1
assert frame["consumption_kwh"].iloc[0] == 99.0
def test_load_recent_from_database_excludes_readings_before_the_since_bound( def test_load_recent_from_database_excludes_readings_before_the_since_bound(
connexion_ml: Connection, connexion_ml: Connection,
) -> None: ) -> None: