514 lines
16 KiB
Python
514 lines
16 KiB
Python
# 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
|
|
# 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, la taille des
|
|
# tableaux est plafonnée par MAX_SITES et par limit_for_window() (dérivé de la fenêtre, jamais
|
|
# fourni par l'appelant), et les lectures dont le timestamp déborde de la fenêtre demandée sont
|
|
# é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
|
|
|
|
import argparse
|
|
import asyncio
|
|
import json
|
|
from datetime import UTC, datetime
|
|
from typing import Any
|
|
|
|
import httpx
|
|
from sqlalchemy import text
|
|
from sqlalchemy.ext.asyncio import AsyncConnection, create_async_engine
|
|
|
|
from app.core.config import get_settings
|
|
from app.etl.historical_import import SOURCE_NAME as SOURCE_CSV
|
|
|
|
SOURCE_HISTORY = "api_history"
|
|
|
|
MAX_SITES = 100
|
|
|
|
MAX_LIMIT = 1000
|
|
|
|
# Les quatre seules valeurs que la contrainte ck_reading_quality accepte.
|
|
ACCEPTED_QUALITIES = frozenset({"good", "partial", "degraded", "critical"})
|
|
|
|
PHYSICAL_BOUNDS: dict[str, tuple[float, float]] = {
|
|
"consumption_kw": (0.0, 100_000.0),
|
|
"consumption_kwh": (0.0, 100_000.0),
|
|
"voltage_v": (0.0, 1_000.0),
|
|
"current_a": (0.0, 10_000.0),
|
|
"power_factor": (0.0, 1.0),
|
|
"temperature_celsius": (-90.0, 60.0),
|
|
"humidity_percent": (0.0, 100.0),
|
|
}
|
|
|
|
CAPACITY_BOUNDS = (0.0, 100_000.0)
|
|
|
|
|
|
def create_mock_api_client() -> httpx.AsyncClient:
|
|
settings = get_settings()
|
|
|
|
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=(username, password),
|
|
timeout=settings.mock_api_timeout_seconds,
|
|
)
|
|
|
|
|
|
def read_text(payload: dict[str, Any], key: str) -> str:
|
|
value = payload.get(key)
|
|
|
|
if not isinstance(value, str) or not value:
|
|
raise ValueError(f"Champ {key} absent ou invalide dans la réponse de l'API Mock.")
|
|
|
|
return value
|
|
|
|
|
|
def optional_text(value: Any) -> str | None:
|
|
return value if isinstance(value, str) else None
|
|
|
|
|
|
def coerce_measure(
|
|
value: Any,
|
|
bounds: tuple[float, float],
|
|
) -> float | None:
|
|
if isinstance(value, bool) or not isinstance(value, int | float):
|
|
return None
|
|
|
|
lower, upper = bounds
|
|
|
|
# Écarte aussi NaN et les infinis, qu'aucune comparaison de bornes ne retient.
|
|
return float(value) if lower <= value <= upper else None
|
|
|
|
|
|
def resolve_quality(
|
|
value: Any,
|
|
rejected: list[str],
|
|
) -> str | None:
|
|
quality = value if isinstance(value, str) and value in ACCEPTED_QUALITIES else None
|
|
|
|
if rejected:
|
|
return "critical" if quality == "critical" else "degraded"
|
|
|
|
return quality
|
|
|
|
|
|
def resolve_null_reasons(
|
|
value: Any,
|
|
rejected: list[str],
|
|
) -> list[str]:
|
|
reported = [str(reason) for reason in value] if isinstance(value, list) else []
|
|
|
|
return reported + rejected
|
|
|
|
|
|
async def fetch_sites(
|
|
client: httpx.AsyncClient,
|
|
) -> list[dict[str, Any]]:
|
|
response = await client.get("/api/v1/sites")
|
|
|
|
response.raise_for_status()
|
|
|
|
payload = response.json()
|
|
|
|
if not isinstance(payload, list):
|
|
raise ValueError("La réponse /api/v1/sites doit être une liste.")
|
|
|
|
if len(payload) > MAX_SITES:
|
|
raise ValueError(f"La réponse /api/v1/sites dépasse le plafond de {MAX_SITES} sites.")
|
|
|
|
return payload
|
|
|
|
|
|
def build_site_row(
|
|
site: dict[str, Any],
|
|
) -> dict[str, Any]:
|
|
return {
|
|
"site_id": read_text(site, "site_id"),
|
|
"site_type": read_text(site, "site_type"),
|
|
"site_name": read_text(site, "site_name"),
|
|
"location": optional_text(site.get("location")),
|
|
"capacity_kw": coerce_measure(site.get("capacity_kw"), CAPACITY_BOUNDS),
|
|
"status": optional_text(site.get("status")),
|
|
}
|
|
|
|
|
|
async def upsert_sites(
|
|
connection: AsyncConnection,
|
|
sites: list[dict[str, Any]],
|
|
) -> None:
|
|
rows = [build_site_row(site) for site in sites]
|
|
|
|
if not rows:
|
|
return
|
|
|
|
await connection.execute(
|
|
text(
|
|
"""
|
|
INSERT INTO site (
|
|
site_id,
|
|
site_type,
|
|
site_name,
|
|
location,
|
|
capacity_kw,
|
|
status
|
|
)
|
|
VALUES (
|
|
:site_id,
|
|
:site_type,
|
|
:site_name,
|
|
:location,
|
|
:capacity_kw,
|
|
:status
|
|
)
|
|
ON CONFLICT (site_id)
|
|
DO UPDATE SET
|
|
site_type = EXCLUDED.site_type,
|
|
site_name = EXCLUDED.site_name,
|
|
location = EXCLUDED.location,
|
|
capacity_kw = EXCLUDED.capacity_kw,
|
|
status = EXCLUDED.status
|
|
"""
|
|
),
|
|
rows,
|
|
)
|
|
|
|
|
|
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(
|
|
client: httpx.AsyncClient,
|
|
site_id: str,
|
|
start_time: datetime,
|
|
end_time: datetime,
|
|
limit: int = MAX_LIMIT,
|
|
) -> list[dict[str, Any]]:
|
|
response = await client.get(
|
|
"/api/v1/readings",
|
|
params={
|
|
"site_id": site_id,
|
|
"start_time": start_time.isoformat(),
|
|
"end_time": end_time.isoformat(),
|
|
"limit": limit,
|
|
},
|
|
)
|
|
|
|
response.raise_for_status()
|
|
|
|
payload = response.json()
|
|
|
|
if not isinstance(payload, list):
|
|
raise ValueError("La réponse /api/v1/readings doit être une liste.")
|
|
|
|
if len(payload) > limit:
|
|
raise ValueError(f"La réponse /api/v1/readings dépasse la limite demandée de {limit}.")
|
|
|
|
# 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(
|
|
reading: dict[str, Any],
|
|
) -> dict[str, Any]:
|
|
measures: dict[str, float | None] = {}
|
|
rejected: list[str] = []
|
|
|
|
for name, bounds in PHYSICAL_BOUNDS.items():
|
|
received = reading.get(name)
|
|
measures[name] = coerce_measure(received, bounds)
|
|
|
|
if received is not None and measures[name] is None:
|
|
rejected.append(f"out_of_physical_bounds:{name}")
|
|
|
|
return {
|
|
"site_id": read_text(reading, "site_id"),
|
|
"timestamp": parse_datetime(read_text(reading, "timestamp")),
|
|
"source": SOURCE_HISTORY,
|
|
"dataset_id": None,
|
|
**measures,
|
|
"consumption_euros": None,
|
|
"solar_irradiance_wm2": None,
|
|
"is_working_hours": None,
|
|
"data_quality": resolve_quality(reading.get("data_quality"), rejected),
|
|
"null_reasons": resolve_null_reasons(reading.get("null_reasons"), rejected),
|
|
"imputed_values": None,
|
|
"imputation_method": None,
|
|
"raw_data": json.dumps(
|
|
reading,
|
|
ensure_ascii=False,
|
|
),
|
|
}
|
|
|
|
|
|
# `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
|
|
# nommée, DO NOTHING avalerait aussi une violation de clé primaire.
|
|
READING_INSERT = text(
|
|
"""
|
|
INSERT INTO reading (
|
|
site_id,
|
|
timestamp,
|
|
source,
|
|
dataset_id,
|
|
consumption_kw,
|
|
consumption_kwh,
|
|
consumption_euros,
|
|
voltage_v,
|
|
current_a,
|
|
power_factor,
|
|
temperature_celsius,
|
|
humidity_percent,
|
|
solar_irradiance_wm2,
|
|
is_working_hours,
|
|
data_quality,
|
|
null_reasons,
|
|
imputed_values,
|
|
imputation_method,
|
|
raw_data
|
|
)
|
|
VALUES (
|
|
:site_id,
|
|
:timestamp,
|
|
:source,
|
|
:dataset_id,
|
|
:consumption_kw,
|
|
:consumption_kwh,
|
|
:consumption_euros,
|
|
:voltage_v,
|
|
:current_a,
|
|
:power_factor,
|
|
:temperature_celsius,
|
|
:humidity_percent,
|
|
:solar_irradiance_wm2,
|
|
:is_working_hours,
|
|
:data_quality,
|
|
:null_reasons,
|
|
CAST(:imputed_values AS jsonb),
|
|
:imputation_method,
|
|
CAST(:raw_data AS jsonb)
|
|
)
|
|
ON CONFLICT (site_id, timestamp, source, (coalesce(dataset_id, 0)))
|
|
DO NOTHING
|
|
"""
|
|
)
|
|
|
|
|
|
def build_reading_batch(
|
|
readings: list[dict[str, Any]],
|
|
) -> list[dict[str, Any]]:
|
|
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 (vérifié
|
|
empiriquement). `limit` = nombre d'heures de la fenêtre est donc le seul réglage cohérent avec
|
|
le grain horaire du reste du schéma (`period_minutes=60`, historique CSV à une ligne/heure) ;
|
|
un `limit` plus grand fabriquerait des lectures infra-horaires, incompatibles avec les lags
|
|
positionnels de `build_features`. La fenêtre doit donc couvrir un nombre entier d'heures.
|
|
"""
|
|
duree = end_time - start_time
|
|
heures, reste = divmod(duree.total_seconds(), 3600)
|
|
|
|
if reste != 0:
|
|
raise ValueError(
|
|
"La fenêtre doit couvrir un nombre entier d'heures pour obtenir une lecture par "
|
|
f"heure alignée : [{start_time.isoformat()}, {end_time.isoformat()}) n'en couvre pas."
|
|
)
|
|
|
|
if heures > MAX_LIMIT:
|
|
raise ValueError(
|
|
f"La fenêtre demandée couvre {int(heures)}h, au-delà du plafond de {MAX_LIMIT} "
|
|
"lectures accepté par l'API Mock."
|
|
)
|
|
|
|
return int(heures)
|
|
|
|
|
|
async def import_mock_api_history(
|
|
start_time: datetime,
|
|
end_time: datetime,
|
|
dry_run: bool,
|
|
) -> None:
|
|
settings = get_settings()
|
|
limit = limit_for_window(start_time, end_time)
|
|
|
|
engine = create_async_engine(
|
|
str(settings.database_url),
|
|
pool_pre_ping=True,
|
|
)
|
|
|
|
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:
|
|
await upsert_sites(
|
|
connection,
|
|
sites,
|
|
)
|
|
|
|
rows = build_reading_batch(all_readings)
|
|
|
|
if rows:
|
|
await connection.execute(
|
|
READING_INSERT,
|
|
rows,
|
|
)
|
|
|
|
finally:
|
|
await engine.dispose()
|
|
|
|
print("Import API Mock terminé.")
|
|
|
|
|
|
def parse_datetime(value: str) -> datetime:
|
|
# 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:
|
|
parser = argparse.ArgumentParser(description=("Import historique depuis l'API Mock EnerVision"))
|
|
|
|
parser.add_argument(
|
|
"--start-time",
|
|
required=True,
|
|
type=parse_datetime,
|
|
)
|
|
|
|
parser.add_argument(
|
|
"--end-time",
|
|
required=True,
|
|
type=parse_datetime,
|
|
)
|
|
|
|
parser.add_argument(
|
|
"--dry-run",
|
|
action="store_true",
|
|
)
|
|
|
|
return parser.parse_args()
|
|
|
|
|
|
def main() -> None:
|
|
args = parse_args()
|
|
|
|
if args.start_time >= args.end_time:
|
|
raise ValueError("--start-time doit être antérieur à --end-time.")
|
|
|
|
asyncio.run(
|
|
import_mock_api_history(
|
|
start_time=args.start_time,
|
|
end_time=args.end_time,
|
|
dry_run=args.dry_run,
|
|
)
|
|
)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|