Merge remote-tracking branch 'origin/dev' into feat/design-system
# Conflicts: # apps/frontend/src/app/features/auth/change-password/change-password.html # apps/frontend/src/app/features/auth/change-password/change-password.ts # apps/frontend/src/app/features/auth/login/login.html # apps/frontend/src/app/features/auth/login/login.ts
This commit is contained in:
@@ -16,6 +16,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.core.config import Settings, get_settings
|
||||
from app.core.hashing import Argon2Hasher, build_hasher
|
||||
from app.core.mailer import Mailer, SmtpConfig
|
||||
from app.core.principal import Principal
|
||||
from app.core.roles import AccountKind, Role, has_at_least
|
||||
from app.core.security import TokenExpiredError, TokenInvalidError, TokenPolicy
|
||||
@@ -24,13 +25,16 @@ from app.db.session import get_session
|
||||
from app.repositories.alert import AlertRepository
|
||||
from app.repositories.audit_log import AuditLogRepository
|
||||
from app.repositories.login_attempt import LoginAttemptRepository
|
||||
from app.repositories.password_reset_attempt import PasswordResetAttemptRepository
|
||||
from app.repositories.password_reset_token import PasswordResetTokenRepository
|
||||
from app.repositories.reading import ReadingRepository
|
||||
from app.repositories.recommendation import RecommendationRepository
|
||||
from app.repositories.refresh_token import RefreshTokenRepository
|
||||
from app.repositories.site import SiteRepository
|
||||
from app.repositories.user import UserRepository
|
||||
from app.services.alert import AlertService
|
||||
from app.services.auth import AuthService, LoginPolicy
|
||||
from app.services.auth import AuthService, LoginPolicy, PasswordResetPolicy
|
||||
from app.services.reading import ReadingService
|
||||
from app.services.recommendation import RecommendationService
|
||||
from app.services.sensor import SensorService
|
||||
from app.services.site import SiteService
|
||||
@@ -97,11 +101,27 @@ def get_client_ip(request: Request, settings: SettingsDep) -> str | None:
|
||||
return request.client.host if request.client else None
|
||||
|
||||
|
||||
def get_mailer(settings: SettingsDep) -> Mailer:
|
||||
return Mailer(
|
||||
SmtpConfig(
|
||||
host=settings.smtp_host,
|
||||
port=settings.smtp_port,
|
||||
username=settings.smtp_username,
|
||||
password=(
|
||||
settings.smtp_password.get_secret_value() if settings.smtp_password else None
|
||||
),
|
||||
use_tls=settings.smtp_use_tls,
|
||||
from_address=settings.smtp_from_address,
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
def get_auth_service(
|
||||
session: SessionDep,
|
||||
settings: SettingsDep,
|
||||
hasher: Annotated[Argon2Hasher, Depends(get_hasher)],
|
||||
token_policy: Annotated[TokenPolicy, Depends(get_token_policy)],
|
||||
mailer: Annotated[Mailer, Depends(get_mailer)],
|
||||
) -> AuthService:
|
||||
return AuthService(
|
||||
users=UserRepository(session),
|
||||
@@ -118,6 +138,16 @@ def get_auth_service(
|
||||
max_failures_per_identifier=settings.login_max_failures_per_identifier,
|
||||
),
|
||||
refresh_ttl=timedelta(seconds=settings.refresh_token_ttl_seconds),
|
||||
reset_tokens=PasswordResetTokenRepository(session),
|
||||
reset_attempts=PasswordResetAttemptRepository(session),
|
||||
reset_policy=PasswordResetPolicy(
|
||||
window_seconds=settings.password_reset_window_seconds,
|
||||
max_requests_per_identifier=settings.password_reset_max_requests_per_identifier,
|
||||
max_requests_per_ip=settings.password_reset_max_requests_per_ip,
|
||||
token_ttl=timedelta(seconds=settings.password_reset_ttl_seconds),
|
||||
frontend_reset_url=settings.frontend_reset_password_url,
|
||||
),
|
||||
mailer=mailer,
|
||||
)
|
||||
|
||||
|
||||
@@ -168,6 +198,13 @@ def get_stats_service(session: SessionDep) -> StatsService:
|
||||
StatsServiceDep = Annotated[StatsService, Depends(get_stats_service)]
|
||||
|
||||
|
||||
def get_reading_service(session: SessionDep) -> ReadingService:
|
||||
return ReadingService(readings=ReadingRepository(session))
|
||||
|
||||
|
||||
ReadingServiceDep = Annotated[ReadingService, Depends(get_reading_service)]
|
||||
|
||||
|
||||
def get_sensor_service(session: SessionDep) -> SensorService:
|
||||
return SensorService(sites=SiteRepository(session), readings=ReadingRepository(session))
|
||||
|
||||
|
||||
@@ -71,6 +71,14 @@ TAGS: Final[list[dict[str, Any]]] = [
|
||||
"description": "Statistiques agrégées de consommation. Accessible à partir du rôle "
|
||||
"`lecteur`.",
|
||||
},
|
||||
{
|
||||
"name": "readings",
|
||||
"description": (
|
||||
"Historique des lectures de consommation. Fenêtre temporelle plafonnée à 90 jours, "
|
||||
"24 dernières heures par défaut si `start`/`end` sont omis. Accessible à partir du "
|
||||
"rôle `lecteur`."
|
||||
),
|
||||
},
|
||||
{
|
||||
"name": "sensors",
|
||||
"description": "État de santé des capteurs par site. Réservé au rôle `admin`.",
|
||||
@@ -156,3 +164,16 @@ REPONSE_ORIGINE_REFUSEE: Final[Reponses] = {
|
||||
"description": "Origine non autorisée (protection CSRF de `require_trusted_origin`).",
|
||||
},
|
||||
}
|
||||
|
||||
REPONSE_LIMITE: Final[Reponses] = {
|
||||
429: {
|
||||
"model": ErrorResponse,
|
||||
"description": "Trop de demandes sur cette fenêtre glissante.",
|
||||
"headers": {
|
||||
"Retry-After": {
|
||||
"description": "Secondes à attendre avant une nouvelle tentative.",
|
||||
"schema": {"type": "integer"},
|
||||
}
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
# d'accès ne va jamais dans un cookie. C'est ce qui réduit la surface CSRF aux trois routes de
|
||||
# ce module : partout ailleurs, le navigateur n'attache rien de lui-même.
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Request, Response, status
|
||||
from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException, Request, Response, status
|
||||
|
||||
from app.api.deps import (
|
||||
AuthServiceDep,
|
||||
@@ -12,6 +12,7 @@ from app.api.deps import (
|
||||
require_trusted_origin,
|
||||
)
|
||||
from app.api.openapi import (
|
||||
REPONSE_LIMITE,
|
||||
REPONSE_ORIGINE_REFUSEE,
|
||||
REPONSE_VALIDATION,
|
||||
REPONSES_AUTHENTIFIEES,
|
||||
@@ -21,15 +22,19 @@ from app.api.openapi import (
|
||||
from app.core.cookies import RefreshCookie, cookie_name
|
||||
from app.core.logging import get_logger
|
||||
from app.schemas.auth import (
|
||||
ForgotPasswordRequest,
|
||||
LoginRequest,
|
||||
PasswordChangeRequest,
|
||||
PrincipalResponse,
|
||||
ResetPasswordRequest,
|
||||
ResetTokenValidationResponse,
|
||||
TokenResponse,
|
||||
)
|
||||
from app.schemas.errors import ErrorResponse
|
||||
from app.services.auth import (
|
||||
AuthenticatedSession,
|
||||
InvalidCredentialsError,
|
||||
InvalidOrExpiredResetTokenError,
|
||||
RateLimitedError,
|
||||
SessionRejectedError,
|
||||
)
|
||||
@@ -39,6 +44,7 @@ logger = get_logger(__name__)
|
||||
|
||||
DETAIL_IDENTIFIANTS = "Identifiants invalides"
|
||||
DETAIL_SESSION = "Session invalide"
|
||||
DETAIL_LIEN_RESET = "Lien invalide ou expiré"
|
||||
|
||||
REPONSES_LOGIN: Reponses = {
|
||||
**REPONSE_VALIDATION,
|
||||
@@ -85,6 +91,20 @@ REPONSES_MOT_DE_PASSE: Reponses = {
|
||||
},
|
||||
}
|
||||
|
||||
REPONSES_FORGOT_PASSWORD: Reponses = {
|
||||
**REPONSE_VALIDATION,
|
||||
**REPONSE_LIMITE,
|
||||
}
|
||||
|
||||
REPONSES_RESET_PASSWORD: Reponses = {
|
||||
**REPONSE_VALIDATION,
|
||||
**REPONSE_ORIGINE_REFUSEE,
|
||||
400: {
|
||||
"model": ErrorResponse,
|
||||
"description": "Lien invalide, déjà utilisé, ou expiré (durée de vie : 15 minutes).",
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def repond(
|
||||
response: Response, settings: SettingsDep, session: AuthenticatedSession
|
||||
@@ -267,3 +287,79 @@ async def change_password(
|
||||
|
||||
logger.info("auth.password_changed user_id=%s", principal.id)
|
||||
return repond(response, settings, session)
|
||||
|
||||
|
||||
@router.post(
|
||||
"/forgot-password",
|
||||
status_code=status.HTTP_202_ACCEPTED,
|
||||
summary="Demande un lien de réinitialisation par email",
|
||||
responses=REPONSES_FORGOT_PASSWORD,
|
||||
)
|
||||
async def forgot_password(
|
||||
payload: ForgotPasswordRequest,
|
||||
request: Request,
|
||||
response: Response,
|
||||
service: AuthServiceDep,
|
||||
background_tasks: BackgroundTasks,
|
||||
client_ip: str | None = Depends(get_client_ip),
|
||||
) -> None:
|
||||
response.headers["Cache-Control"] = "no-store"
|
||||
|
||||
try:
|
||||
await service.request_password_reset(
|
||||
email=payload.email,
|
||||
client_ip=client_ip,
|
||||
user_agent=request.headers.get("user-agent"),
|
||||
background_tasks=background_tasks,
|
||||
)
|
||||
except RateLimitedError as erreur:
|
||||
logger.warning("auth.password_reset.rate_limited ip=%s", client_ip)
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
|
||||
detail="Trop de demandes, réessayez plus tard",
|
||||
headers={"Retry-After": str(erreur.retry_after)},
|
||||
) from erreur
|
||||
|
||||
|
||||
@router.get(
|
||||
"/reset-password/validate",
|
||||
response_model=ResetTokenValidationResponse,
|
||||
summary="Vérifie sans le consommer si un lien de réinitialisation est encore valide",
|
||||
responses=REPONSE_VALIDATION,
|
||||
)
|
||||
async def validate_reset_token(token: str, service: AuthServiceDep) -> ResetTokenValidationResponse:
|
||||
return ResetTokenValidationResponse(valid=await service.is_reset_token_valid(token=token))
|
||||
|
||||
|
||||
@router.post(
|
||||
"/reset-password",
|
||||
response_model=TokenResponse,
|
||||
summary="Choisit un nouveau mot de passe depuis un lien reçu par email",
|
||||
dependencies=[Depends(require_trusted_origin)],
|
||||
responses=REPONSES_RESET_PASSWORD,
|
||||
)
|
||||
async def reset_password(
|
||||
payload: ResetPasswordRequest,
|
||||
request: Request,
|
||||
response: Response,
|
||||
settings: SettingsDep,
|
||||
service: AuthServiceDep,
|
||||
client_ip: str | None = Depends(get_client_ip),
|
||||
) -> TokenResponse:
|
||||
response.headers["Cache-Control"] = "no-store"
|
||||
|
||||
try:
|
||||
session = await service.confirm_password_reset(
|
||||
token=payload.token,
|
||||
new_password=payload.new_password,
|
||||
client_ip=client_ip,
|
||||
user_agent=request.headers.get("user-agent"),
|
||||
)
|
||||
except InvalidOrExpiredResetTokenError as erreur:
|
||||
logger.warning("auth.password_reset.invalid_token ip=%s", client_ip)
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST, detail=DETAIL_LIEN_RESET
|
||||
) from erreur
|
||||
|
||||
logger.info("auth.password_reset.success user_id=%s", session.principal.id)
|
||||
return repond(response, settings, session)
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
from datetime import datetime
|
||||
|
||||
from fastapi import APIRouter, HTTPException, Query, status
|
||||
|
||||
from app.api.deps import LecteurDep, ReadingServiceDep
|
||||
from app.api.openapi import REPONSE_VALIDATION, Reponses
|
||||
from app.schemas.errors import ErrorResponse
|
||||
from app.schemas.reading import ReadingResponse
|
||||
from app.services.reading import FenetreInverseeError, FenetreTropLargeError
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
REPONSES_FENETRE: Reponses = {
|
||||
**REPONSE_VALIDATION,
|
||||
400: {
|
||||
"model": ErrorResponse,
|
||||
"description": (
|
||||
"Fenêtre temporelle invalide : `start` postérieur ou égal à `end`, ou écart entre "
|
||||
"les deux supérieur à 90 jours."
|
||||
),
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@router.get(
|
||||
"",
|
||||
response_model=list[ReadingResponse],
|
||||
summary="Liste l'historique des lectures",
|
||||
responses=REPONSES_FENETRE,
|
||||
)
|
||||
async def list_readings(
|
||||
_: LecteurDep,
|
||||
service: ReadingServiceDep,
|
||||
site_id: str | None = None,
|
||||
start: datetime | None = None,
|
||||
end: datetime | None = None,
|
||||
limit: int = Query(500, ge=1, le=2000),
|
||||
offset: int = Query(0, ge=0),
|
||||
) -> list[ReadingResponse]:
|
||||
try:
|
||||
lectures = await service.list_history(
|
||||
site_id=site_id, start=start, end=end, limit=limit, offset=offset
|
||||
)
|
||||
except FenetreInverseeError as erreur:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
detail="`start` doit être strictement antérieur à `end`",
|
||||
) from erreur
|
||||
except FenetreTropLargeError as erreur:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
detail="L'écart entre `start` et `end` ne peut pas dépasser 90 jours",
|
||||
) from erreur
|
||||
return [ReadingResponse.model_validate(lecture) for lecture in lectures]
|
||||
@@ -1,7 +1,17 @@
|
||||
from fastapi import APIRouter
|
||||
|
||||
from app.api.openapi import REPONSE_SERVEUR, REPONSES_ADMIN, REPONSES_LECTEUR
|
||||
from app.api.v1.endpoints import alerts, auth, health, recommendations, sensors, sites, stats, users
|
||||
from app.api.v1.endpoints import (
|
||||
alerts,
|
||||
auth,
|
||||
health,
|
||||
readings,
|
||||
recommendations,
|
||||
sensors,
|
||||
sites,
|
||||
stats,
|
||||
users,
|
||||
)
|
||||
|
||||
api_router = APIRouter(responses=REPONSE_SERVEUR)
|
||||
api_router.include_router(health.router, prefix="/health", tags=["health"])
|
||||
@@ -18,6 +28,9 @@ api_router.include_router(
|
||||
responses=REPONSES_LECTEUR,
|
||||
)
|
||||
api_router.include_router(stats.router, prefix="/stats", tags=["stats"], responses=REPONSES_LECTEUR)
|
||||
api_router.include_router(
|
||||
readings.router, prefix="/readings", tags=["readings"], responses=REPONSES_LECTEUR
|
||||
)
|
||||
api_router.include_router(
|
||||
sensors.router, prefix="/sensors", tags=["sensors"], responses=REPONSES_ADMIN
|
||||
)
|
||||
|
||||
+24
-4
@@ -9,6 +9,7 @@ import argparse
|
||||
import asyncio
|
||||
import json
|
||||
import secrets
|
||||
import string
|
||||
import sys
|
||||
from getpass import getpass
|
||||
from pathlib import Path
|
||||
@@ -22,9 +23,9 @@ from app.core.roles import Role
|
||||
from app.db.session import get_session_factory
|
||||
from app.main import create_app
|
||||
from app.repositories.user import UserRepository
|
||||
from app.schemas.auth import PASSWORD_MIN_LENGTH, SPECIAL_CHARACTERS, valide_complexite
|
||||
|
||||
LONGUEUR_MOT_DE_PASSE_GENERE = 24
|
||||
LONGUEUR_MINIMALE = 12
|
||||
CHEMIN_CONTRAT = Path(__file__).resolve().parent.parent / "openapi.json"
|
||||
|
||||
|
||||
@@ -111,15 +112,34 @@ def build_parser() -> argparse.ArgumentParser:
|
||||
return parser
|
||||
|
||||
|
||||
def genere_mot_de_passe() -> str:
|
||||
tirage = secrets.SystemRandom()
|
||||
classes = [
|
||||
string.ascii_uppercase,
|
||||
string.ascii_lowercase,
|
||||
string.digits,
|
||||
SPECIAL_CHARACTERS,
|
||||
]
|
||||
reste = LONGUEUR_MOT_DE_PASSE_GENERE - len(classes)
|
||||
caracteres = [tirage.choice(classe) for classe in classes]
|
||||
caracteres += [tirage.choice("".join(classes)) for _ in range(reste)]
|
||||
tirage.shuffle(caracteres)
|
||||
return "".join(caracteres)
|
||||
|
||||
|
||||
def read_password(*, generate: bool) -> str:
|
||||
if generate:
|
||||
mot_de_passe = secrets.token_urlsafe(LONGUEUR_MOT_DE_PASSE_GENERE)
|
||||
mot_de_passe = genere_mot_de_passe()
|
||||
print(f"Mot de passe généré, il ne sera plus affiché : {mot_de_passe}")
|
||||
return mot_de_passe
|
||||
|
||||
mot_de_passe = getpass("Mot de passe : ")
|
||||
if len(mot_de_passe) < LONGUEUR_MINIMALE:
|
||||
raise SystemExit(f"Le mot de passe doit faire au moins {LONGUEUR_MINIMALE} caractères")
|
||||
if len(mot_de_passe) < PASSWORD_MIN_LENGTH:
|
||||
raise SystemExit(f"Le mot de passe doit faire au moins {PASSWORD_MIN_LENGTH} caractères")
|
||||
try:
|
||||
valide_complexite(mot_de_passe)
|
||||
except ValueError as erreur:
|
||||
raise SystemExit(str(erreur)) from erreur
|
||||
if mot_de_passe != getpass("Confirmation : "):
|
||||
raise SystemExit("Les deux saisies diffèrent")
|
||||
return mot_de_passe
|
||||
|
||||
@@ -54,6 +54,19 @@ class Settings(BaseSettings):
|
||||
login_max_failures_per_ip: int = Field(default=20, ge=1)
|
||||
login_max_failures_per_identifier: int = Field(default=50, ge=1)
|
||||
|
||||
password_reset_ttl_seconds: int = Field(default=900, ge=60, le=3600)
|
||||
password_reset_window_seconds: int = Field(default=900, ge=60)
|
||||
password_reset_max_requests_per_identifier: int = Field(default=3, ge=1)
|
||||
password_reset_max_requests_per_ip: int = Field(default=10, ge=1)
|
||||
|
||||
smtp_host: str = "localhost"
|
||||
smtp_port: int = Field(default=587, ge=1, le=65535)
|
||||
smtp_username: str | None = None
|
||||
smtp_password: SecretStr | None = None
|
||||
smtp_use_tls: bool = False
|
||||
smtp_from_address: str = "no-reply@enervision.fr"
|
||||
frontend_reset_password_url: str = "http://localhost:4200/reset-password" # noqa: S105
|
||||
|
||||
trust_proxy_headers: bool = False
|
||||
expose_api_docs: bool | None = None
|
||||
metrics_token: SecretStr | None = None
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
# Piège : l'URL de réinitialisation porte le jeton en clair. Ne jamais la journaliser :
|
||||
# `send_password_reset_email()` ne logue que le destinataire, jamais `reset_url`.
|
||||
|
||||
from dataclasses import dataclass
|
||||
from email.message import EmailMessage
|
||||
|
||||
import aiosmtplib
|
||||
|
||||
from app.core.logging import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class SmtpConfig:
|
||||
host: str
|
||||
port: int
|
||||
username: str | None
|
||||
password: str | None
|
||||
use_tls: bool
|
||||
from_address: str
|
||||
|
||||
|
||||
class Mailer:
|
||||
def __init__(self, config: SmtpConfig) -> None:
|
||||
self._config = config
|
||||
|
||||
async def send_password_reset_email(self, *, to: str, reset_url: str) -> None:
|
||||
message = EmailMessage()
|
||||
message["From"] = self._config.from_address
|
||||
message["To"] = to
|
||||
message["Subject"] = "Réinitialisation de votre mot de passe EnerVision"
|
||||
message.set_content(
|
||||
"Une réinitialisation de mot de passe a été demandée pour ce compte.\n\n"
|
||||
f"Ouvrez ce lien dans les 15 minutes pour choisir un nouveau mot de passe : "
|
||||
f"{reset_url}\n\n"
|
||||
"Si vous n'êtes pas à l'origine de cette demande, ignorez cet email."
|
||||
)
|
||||
|
||||
_, message_recu = await aiosmtplib.send(
|
||||
message,
|
||||
hostname=self._config.host,
|
||||
port=self._config.port,
|
||||
username=self._config.username,
|
||||
password=self._config.password,
|
||||
use_tls=self._config.use_tls,
|
||||
)
|
||||
logger.info("mailer.password_reset_sent to=%s smtp_response=%s", to, message_recu)
|
||||
@@ -0,0 +1,621 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import asyncio
|
||||
import hashlib
|
||||
import json
|
||||
from pathlib import Path
|
||||
from typing import Any, cast
|
||||
|
||||
import pandas as pd
|
||||
from sqlalchemy import text
|
||||
from sqlalchemy.ext.asyncio import AsyncConnection, create_async_engine
|
||||
|
||||
from app.core.config import get_settings
|
||||
|
||||
REQUIRED_COLUMNS = {
|
||||
"timestamp",
|
||||
"site_id",
|
||||
"site_type",
|
||||
"site_name",
|
||||
"consumption_kwh",
|
||||
"consumption_euros",
|
||||
"temperature_celsius",
|
||||
"humidity_percent",
|
||||
"solar_irradiance_wm2",
|
||||
"hour",
|
||||
"day_of_week",
|
||||
"day_name",
|
||||
"month",
|
||||
"is_weekend",
|
||||
"is_working_hours",
|
||||
}
|
||||
|
||||
MEASURE_COLUMNS = [
|
||||
"consumption_kwh",
|
||||
"consumption_euros",
|
||||
"temperature_celsius",
|
||||
"humidity_percent",
|
||||
"solar_irradiance_wm2",
|
||||
]
|
||||
|
||||
SOURCE_NAME = "csv"
|
||||
|
||||
|
||||
def compute_sha256(path: Path) -> str:
|
||||
"""Calcule l'empreinte SHA-256 du fichier source."""
|
||||
sha256 = hashlib.sha256()
|
||||
|
||||
with path.open("rb") as source:
|
||||
for block in iter(lambda: source.read(1024 * 1024), b""):
|
||||
sha256.update(block)
|
||||
|
||||
return sha256.hexdigest()
|
||||
|
||||
|
||||
def load_metadata(path: Path) -> dict[str, Any]:
|
||||
"""Charge les métadonnées fournies avec le dataset."""
|
||||
with path.open("r", encoding="utf-8") as source:
|
||||
metadata = json.load(source)
|
||||
|
||||
if not isinstance(metadata, dict):
|
||||
raise ValueError("Le fichier de métadonnées doit contenir un objet JSON.")
|
||||
|
||||
return cast(dict[str, Any], metadata)
|
||||
|
||||
|
||||
def classify_quality(
|
||||
row: dict[str, Any],
|
||||
) -> tuple[str, list[str]]:
|
||||
"""
|
||||
Déduit une qualité technique à partir des champs manquants.
|
||||
|
||||
Les valeurs NULL sont conservées. On ne cherche pas ici à
|
||||
déterminer la cause physique exacte de leur absence.
|
||||
"""
|
||||
missing = [column for column in MEASURE_COLUMNS if pd.isna(row.get(column))]
|
||||
|
||||
if not missing:
|
||||
quality = "good"
|
||||
elif len(missing) == len(MEASURE_COLUMNS):
|
||||
quality = "critical"
|
||||
elif "consumption_kwh" in missing:
|
||||
quality = "degraded"
|
||||
else:
|
||||
quality = "partial"
|
||||
|
||||
reasons = [f"missing:{column}" for column in missing]
|
||||
|
||||
return quality, reasons
|
||||
|
||||
|
||||
def validate_source(
|
||||
frame: pd.DataFrame,
|
||||
metadata: dict[str, Any],
|
||||
) -> None:
|
||||
"""Valide le dataset avant tout chargement en base."""
|
||||
missing_columns = REQUIRED_COLUMNS.difference(frame.columns)
|
||||
|
||||
if missing_columns:
|
||||
raise ValueError(f"Colonnes obligatoires absentes : {sorted(missing_columns)}")
|
||||
|
||||
expected_records = int(metadata["total_records"])
|
||||
|
||||
if len(frame) != expected_records:
|
||||
raise ValueError(f"Nombre de lignes inattendu : {len(frame)} au lieu de {expected_records}")
|
||||
|
||||
expected_sites = set(metadata["sites"].keys())
|
||||
actual_sites = set(frame["site_id"].unique())
|
||||
|
||||
if actual_sites != expected_sites:
|
||||
raise ValueError(
|
||||
f"Sites incohérents. Attendus={sorted(expected_sites)}, trouvés={sorted(actual_sites)}"
|
||||
)
|
||||
|
||||
duplicated = frame.duplicated(subset=["site_id", "timestamp"]).sum()
|
||||
|
||||
if duplicated:
|
||||
raise ValueError(f"{duplicated} doublons (site_id, timestamp) détectés")
|
||||
|
||||
static_variants = frame.groupby("site_id")[["site_type", "site_name"]].nunique()
|
||||
|
||||
if (static_variants > 1).any().any():
|
||||
raise ValueError("Un site possède plusieurs valeurs de site_type ou site_name.")
|
||||
|
||||
# Vérifie également que tous les timestamps
|
||||
# peuvent être interprétés correctement.
|
||||
pd.to_datetime(
|
||||
frame["timestamp"],
|
||||
errors="raise",
|
||||
)
|
||||
|
||||
|
||||
def normalize_timestamps(
|
||||
frame: pd.DataFrame,
|
||||
source_timezone: str,
|
||||
) -> pd.DataFrame:
|
||||
"""
|
||||
Normalise les timestamps et leur associe une timezone.
|
||||
|
||||
Les timestamps originaux sont conservés dans une colonne
|
||||
temporaire afin de pouvoir les stocker dans raw_data.
|
||||
"""
|
||||
normalized = frame.copy()
|
||||
|
||||
normalized["_source_timestamp"] = normalized["timestamp"]
|
||||
|
||||
timestamps = pd.to_datetime(
|
||||
normalized["timestamp"],
|
||||
errors="raise",
|
||||
)
|
||||
|
||||
if timestamps.dt.tz is None:
|
||||
timestamps = timestamps.dt.tz_localize(source_timezone)
|
||||
else:
|
||||
timestamps = timestamps.dt.tz_convert(source_timezone)
|
||||
|
||||
normalized["timestamp"] = timestamps
|
||||
|
||||
return normalized
|
||||
|
||||
|
||||
def to_json_value(value: Any) -> Any:
|
||||
"""
|
||||
Convertit une valeur Pandas/Numpy en valeur
|
||||
compatible JSON.
|
||||
"""
|
||||
if value is None:
|
||||
return None
|
||||
|
||||
try:
|
||||
if pd.isna(value):
|
||||
return None
|
||||
except TypeError, ValueError:
|
||||
pass
|
||||
|
||||
if isinstance(value, pd.Timestamp):
|
||||
return value.isoformat()
|
||||
|
||||
if hasattr(value, "item"):
|
||||
return value.item()
|
||||
|
||||
return value
|
||||
|
||||
|
||||
async def ensure_dataset(
|
||||
connection: AsyncConnection,
|
||||
metadata: dict[str, Any],
|
||||
sha256: str,
|
||||
source_timezone: str,
|
||||
storage_uri: str,
|
||||
) -> int:
|
||||
"""
|
||||
Crée l'entrée dataset si elle n'existe pas.
|
||||
|
||||
Le SHA-256 permet de reconnaître un fichier déjà importé
|
||||
et participe à l'idempotence et à la traçabilité.
|
||||
"""
|
||||
result = await connection.execute(
|
||||
text(
|
||||
"""
|
||||
SELECT dataset_id
|
||||
FROM dataset
|
||||
WHERE archive_sha256 = :sha256
|
||||
LIMIT 1
|
||||
"""
|
||||
),
|
||||
{
|
||||
"sha256": sha256,
|
||||
},
|
||||
)
|
||||
|
||||
existing = result.scalar_one_or_none()
|
||||
|
||||
if existing is not None:
|
||||
return int(existing)
|
||||
|
||||
metadata_summary = {
|
||||
"generator_version": metadata.get("generator_version"),
|
||||
"total_sites": metadata.get("total_sites"),
|
||||
"total_records": metadata.get("total_records"),
|
||||
"date_range": metadata.get("date_range"),
|
||||
"frequency": metadata.get("frequency"),
|
||||
"null_injection_enabled": metadata.get("null_injection_enabled"),
|
||||
"null_strategies": metadata.get("null_strategies"),
|
||||
"importer": "historical_import_v1",
|
||||
}
|
||||
|
||||
result = await connection.execute(
|
||||
text(
|
||||
"""
|
||||
INSERT INTO dataset (
|
||||
dataset_name,
|
||||
archive_sha256,
|
||||
storage_uri,
|
||||
source_timezone,
|
||||
"metadata"
|
||||
)
|
||||
VALUES (
|
||||
:dataset_name,
|
||||
:archive_sha256,
|
||||
:storage_uri,
|
||||
:source_timezone,
|
||||
CAST(:metadata AS jsonb)
|
||||
)
|
||||
RETURNING dataset_id
|
||||
"""
|
||||
),
|
||||
{
|
||||
"dataset_name": ("EnerVision historical dataset 2023-2024"),
|
||||
"archive_sha256": sha256,
|
||||
"storage_uri": storage_uri,
|
||||
"source_timezone": source_timezone,
|
||||
"metadata": json.dumps(
|
||||
metadata_summary,
|
||||
ensure_ascii=False,
|
||||
),
|
||||
},
|
||||
)
|
||||
|
||||
return int(result.scalar_one())
|
||||
|
||||
|
||||
async def upsert_sites(
|
||||
connection: AsyncConnection,
|
||||
frame: pd.DataFrame,
|
||||
) -> None:
|
||||
"""Insère ou met à jour les sites du dataset."""
|
||||
sites = cast(
|
||||
list[dict[str, Any]],
|
||||
frame[
|
||||
[
|
||||
"site_id",
|
||||
"site_type",
|
||||
"site_name",
|
||||
]
|
||||
]
|
||||
.drop_duplicates(subset=["site_id"])
|
||||
.to_dict(orient="records"),
|
||||
)
|
||||
|
||||
await connection.execute(
|
||||
text(
|
||||
"""
|
||||
INSERT INTO site (
|
||||
site_id,
|
||||
site_type,
|
||||
site_name
|
||||
)
|
||||
VALUES (
|
||||
:site_id,
|
||||
:site_type,
|
||||
:site_name
|
||||
)
|
||||
ON CONFLICT (site_id)
|
||||
DO UPDATE SET
|
||||
site_type = EXCLUDED.site_type,
|
||||
site_name = EXCLUDED.site_name
|
||||
"""
|
||||
),
|
||||
sites,
|
||||
)
|
||||
|
||||
|
||||
def build_reading_batch(
|
||||
chunk: pd.DataFrame,
|
||||
dataset_id: int,
|
||||
) -> list[dict[str, Any]]:
|
||||
"""
|
||||
Transforme un chunk Pandas en lignes prêtes
|
||||
à être chargées dans la table reading.
|
||||
"""
|
||||
rows: list[dict[str, Any]] = []
|
||||
|
||||
records = cast(
|
||||
list[dict[str, Any]],
|
||||
chunk.to_dict(orient="records"),
|
||||
)
|
||||
|
||||
for record in records:
|
||||
quality, reasons = classify_quality(record)
|
||||
|
||||
raw_data = {
|
||||
column: to_json_value(value)
|
||||
for column, value in record.items()
|
||||
if column != "_source_timestamp"
|
||||
}
|
||||
|
||||
# Dans raw_data, on conserve le timestamp
|
||||
# exactement tel qu'il était dans le CSV.
|
||||
raw_data["timestamp"] = to_json_value(record["_source_timestamp"])
|
||||
|
||||
rows.append(
|
||||
{
|
||||
"site_id": record["site_id"],
|
||||
"timestamp": record["timestamp"],
|
||||
"source": SOURCE_NAME,
|
||||
"dataset_id": dataset_id,
|
||||
# Non fourni par le dataset historique.
|
||||
"consumption_kw": None,
|
||||
"consumption_kwh": to_json_value(record["consumption_kwh"]),
|
||||
"consumption_euros": to_json_value(record["consumption_euros"]),
|
||||
# Non fournis par le CSV historique.
|
||||
"voltage_v": None,
|
||||
"current_a": None,
|
||||
"power_factor": None,
|
||||
"temperature_celsius": (to_json_value(record["temperature_celsius"])),
|
||||
"humidity_percent": (to_json_value(record["humidity_percent"])),
|
||||
"solar_irradiance_wm2": (to_json_value(record["solar_irradiance_wm2"])),
|
||||
"is_working_hours": bool(record["is_working_hours"]),
|
||||
"data_quality": quality,
|
||||
"null_reasons": reasons,
|
||||
# Aucune imputation pendant l'ingestion RAW.
|
||||
# Les valeurs manquantes sont conservées telles quelles
|
||||
# afin de préserver la donnée source.
|
||||
"imputed_values": None,
|
||||
"imputation_method": None,
|
||||
# Conservation de la donnée source
|
||||
# pour la traçabilité.
|
||||
"raw_data": json.dumps(
|
||||
raw_data,
|
||||
ensure_ascii=False,
|
||||
),
|
||||
}
|
||||
)
|
||||
|
||||
return rows
|
||||
|
||||
|
||||
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 DO NOTHING
|
||||
"""
|
||||
)
|
||||
|
||||
|
||||
async def import_historical(
|
||||
csv_path: Path,
|
||||
metadata_path: Path,
|
||||
source_timezone: str,
|
||||
batch_size: int,
|
||||
dry_run: bool,
|
||||
storage_uri: str,
|
||||
) -> None:
|
||||
"""
|
||||
Exécute le pipeline ETL historique EnerVision.
|
||||
|
||||
Étapes :
|
||||
1. Extract
|
||||
2. Validate
|
||||
3. Transform
|
||||
4. Load
|
||||
"""
|
||||
metadata = load_metadata(metadata_path)
|
||||
|
||||
frame = pd.read_csv(csv_path)
|
||||
|
||||
validate_source(
|
||||
frame,
|
||||
metadata,
|
||||
)
|
||||
|
||||
print(f"Lignes : {len(frame)}")
|
||||
print(f"Sites : {frame['site_id'].nunique()}")
|
||||
print(f"Période : {frame['timestamp'].min()} -> {frame['timestamp'].max()}")
|
||||
print(f"Doublons : {frame.duplicated(['site_id', 'timestamp']).sum()}")
|
||||
|
||||
print("\nValeurs NULL :")
|
||||
print(frame[MEASURE_COLUMNS].isna().sum())
|
||||
|
||||
sha256 = compute_sha256(csv_path)
|
||||
|
||||
print(f"\nSHA-256 : {sha256}")
|
||||
|
||||
if dry_run:
|
||||
print("\nDry-run terminé : aucune donnée écrite.")
|
||||
return
|
||||
|
||||
normalized = normalize_timestamps(
|
||||
frame,
|
||||
source_timezone,
|
||||
)
|
||||
|
||||
settings = get_settings()
|
||||
|
||||
engine = create_async_engine(
|
||||
str(settings.database_url),
|
||||
pool_pre_ping=True,
|
||||
)
|
||||
|
||||
try:
|
||||
async with engine.begin() as connection:
|
||||
dataset_id = await ensure_dataset(
|
||||
connection=connection,
|
||||
metadata=metadata,
|
||||
sha256=sha256,
|
||||
source_timezone=source_timezone,
|
||||
storage_uri=storage_uri,
|
||||
)
|
||||
|
||||
await upsert_sites(
|
||||
connection,
|
||||
normalized,
|
||||
)
|
||||
|
||||
result = await connection.execute(
|
||||
text(
|
||||
"""
|
||||
SELECT COUNT(*)
|
||||
FROM reading
|
||||
WHERE dataset_id = :dataset_id
|
||||
AND source = :source
|
||||
"""
|
||||
),
|
||||
{
|
||||
"dataset_id": dataset_id,
|
||||
"source": SOURCE_NAME,
|
||||
},
|
||||
)
|
||||
|
||||
before = int(result.scalar_one())
|
||||
|
||||
for start in range(
|
||||
0,
|
||||
len(normalized),
|
||||
batch_size,
|
||||
):
|
||||
chunk = normalized.iloc[start : start + batch_size]
|
||||
|
||||
rows = build_reading_batch(
|
||||
chunk,
|
||||
dataset_id,
|
||||
)
|
||||
|
||||
await connection.execute(
|
||||
READING_INSERT,
|
||||
rows,
|
||||
)
|
||||
|
||||
loaded = min(
|
||||
start + batch_size,
|
||||
len(normalized),
|
||||
)
|
||||
|
||||
print(f"Chargement : {loaded}/{len(normalized)}")
|
||||
|
||||
result = await connection.execute(
|
||||
text(
|
||||
"""
|
||||
SELECT COUNT(*)
|
||||
FROM reading
|
||||
WHERE dataset_id = :dataset_id
|
||||
AND source = :source
|
||||
"""
|
||||
),
|
||||
{
|
||||
"dataset_id": dataset_id,
|
||||
"source": SOURCE_NAME,
|
||||
},
|
||||
)
|
||||
|
||||
after = int(result.scalar_one())
|
||||
|
||||
print("\nImport terminé.")
|
||||
print(f"dataset_id : {dataset_id}")
|
||||
print(f"lectures avant : {before}")
|
||||
print(f"lectures après : {after}")
|
||||
print(f"nouvelles lectures : {after - before}")
|
||||
|
||||
finally:
|
||||
await engine.dispose()
|
||||
|
||||
|
||||
def parse_args() -> argparse.Namespace:
|
||||
"""Définit les arguments CLI de l'import."""
|
||||
parser = argparse.ArgumentParser(description=("Import historique EnerVision"))
|
||||
|
||||
parser.add_argument(
|
||||
"--csv",
|
||||
type=Path,
|
||||
required=True,
|
||||
help="Chemin vers le CSV historique.",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--metadata",
|
||||
type=Path,
|
||||
required=True,
|
||||
help=("Chemin vers le fichier dataset_metadata.json."),
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--source-timezone",
|
||||
default="UTC",
|
||||
help=("Timezone associée aux timestamps du dataset. Défaut : UTC."),
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--batch-size",
|
||||
type=int,
|
||||
default=1000,
|
||||
help=("Nombre de lignes insérées par batch. Défaut : 1000."),
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--dry-run",
|
||||
action="store_true",
|
||||
help=("Valide les données sans écrire en base."),
|
||||
)
|
||||
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
def main() -> None:
|
||||
"""Point d'entrée CLI du pipeline."""
|
||||
args = parse_args()
|
||||
|
||||
if args.batch_size <= 0:
|
||||
raise ValueError("--batch-size doit être strictement supérieur à 0.")
|
||||
|
||||
# resolve() est volontairement exécuté ici,
|
||||
# dans la partie synchrone du programme.
|
||||
# Cela évite une opération filesystem bloquante
|
||||
# à l'intérieur d'une fonction async.
|
||||
storage_uri = args.csv.resolve().as_uri()
|
||||
|
||||
asyncio.run(
|
||||
import_historical(
|
||||
csv_path=args.csv,
|
||||
metadata_path=args.metadata,
|
||||
source_timezone=(args.source_timezone),
|
||||
batch_size=args.batch_size,
|
||||
dry_run=args.dry_run,
|
||||
storage_uri=storage_uri,
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -4,6 +4,8 @@
|
||||
from app.models.audit_log import AuditLog
|
||||
from app.models.energy import Alert, Dataset, Prediction, Reading, Recommendation, Site
|
||||
from app.models.login_attempt import LoginAttempt
|
||||
from app.models.password_reset_attempt import PasswordResetAttempt
|
||||
from app.models.password_reset_token import PasswordResetToken
|
||||
from app.models.refresh_token import RefreshToken
|
||||
from app.models.user import AppUser
|
||||
|
||||
@@ -13,6 +15,8 @@ __all__ = [
|
||||
"AuditLog",
|
||||
"Dataset",
|
||||
"LoginAttempt",
|
||||
"PasswordResetAttempt",
|
||||
"PasswordResetToken",
|
||||
"Prediction",
|
||||
"Reading",
|
||||
"Recommendation",
|
||||
|
||||
@@ -29,6 +29,8 @@ class AuditAction(StrEnum):
|
||||
COMPTE_ACTIVE = "user.enabled"
|
||||
COMPTE_MOT_DE_PASSE_REINITIALISE = "user.password_reset_by_admin"
|
||||
COMPTE_MOT_DE_PASSE_CHANGE = "user.password_changed"
|
||||
MOT_DE_PASSE_OUBLIE_DEMANDE = "auth.password_reset_requested"
|
||||
MOT_DE_PASSE_REINITIALISE_PAR_SOI = "auth.password_reset_self_service"
|
||||
REFRESH_REUTILISE = "auth.refresh_reuse_detected"
|
||||
SESSIONS_REVOQUEES = "auth.all_sessions_revoked"
|
||||
LIMITE_PAR_IDENTIFIANT = "auth.identifier_throttled"
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
# Pourquoi : même séparation que `login_attempt` par rapport à `audit_log` : ce compteur est
|
||||
# piloté par l'attaquant (une campagne de demandes) et se purge, l'audit log est en ajout seul.
|
||||
# Piège : la tentative est enregistrée même quand l'email est inconnu, sinon le 429 apprendrait
|
||||
# qu'un compte existe.
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import BigInteger, DateTime, Identity, Index, String, func
|
||||
from sqlalchemy.dialects.postgresql import INET
|
||||
from sqlalchemy.orm import Mapped, mapped_column
|
||||
|
||||
from app.db.base import Base
|
||||
|
||||
|
||||
class PasswordResetAttempt(Base):
|
||||
__tablename__ = "password_reset_attempt"
|
||||
__table_args__ = (
|
||||
Index("ix_password_reset_attempt_email_date", "email_tried", "occurred_at"),
|
||||
Index("ix_password_reset_attempt_ip_date", "client_ip", "occurred_at"),
|
||||
)
|
||||
|
||||
id: Mapped[int] = mapped_column(BigInteger, Identity(always=True), primary_key=True)
|
||||
occurred_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), nullable=False, server_default=func.now()
|
||||
)
|
||||
email_tried: Mapped[str] = mapped_column(String(320), nullable=False)
|
||||
client_ip: Mapped[str | None] = mapped_column(INET, nullable=True)
|
||||
@@ -0,0 +1,40 @@
|
||||
# Pourquoi : même schéma que `refresh_token` (chaîne opaque, jamais un JWT) pour la même
|
||||
# raison : un jeton de réinitialisation doit être révocable d'un coup, et un JWT ne figure
|
||||
# dans aucune ligne à invalider.
|
||||
|
||||
import uuid
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import DateTime, ForeignKey, Index, LargeBinary, Text, func
|
||||
from sqlalchemy.dialects.postgresql import INET
|
||||
from sqlalchemy.dialects.postgresql import UUID as PG_UUID
|
||||
from sqlalchemy.orm import Mapped, mapped_column
|
||||
|
||||
from app.db.base import Base
|
||||
|
||||
|
||||
class PasswordResetToken(Base):
|
||||
__tablename__ = "password_reset_token"
|
||||
__table_args__ = (
|
||||
Index("ix_password_reset_token_user", "user_id"),
|
||||
Index(
|
||||
"ix_password_reset_token_vivants",
|
||||
"user_id",
|
||||
postgresql_where="consumed_at is null",
|
||||
),
|
||||
)
|
||||
|
||||
id: Mapped[uuid.UUID] = mapped_column(
|
||||
PG_UUID(as_uuid=True), primary_key=True, server_default=func.gen_random_uuid()
|
||||
)
|
||||
user_id: Mapped[uuid.UUID] = mapped_column(
|
||||
PG_UUID(as_uuid=True), ForeignKey("app_user.id", ondelete="CASCADE"), nullable=False
|
||||
)
|
||||
token_hash: Mapped[bytes] = mapped_column(LargeBinary, nullable=False, unique=True)
|
||||
issued_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), nullable=False, server_default=func.now()
|
||||
)
|
||||
expires_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
|
||||
consumed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
|
||||
client_ip: Mapped[str | None] = mapped_column(INET, nullable=True)
|
||||
user_agent: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
@@ -0,0 +1,42 @@
|
||||
from dataclasses import dataclass
|
||||
from datetime import UTC, datetime, timedelta
|
||||
|
||||
from sqlalchemy import func, select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.models.password_reset_attempt import PasswordResetAttempt
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class ResetRequestCounts:
|
||||
per_identifier: int
|
||||
per_ip: int
|
||||
|
||||
|
||||
class PasswordResetAttemptRepository:
|
||||
def __init__(self, session: AsyncSession) -> None:
|
||||
self._session = session
|
||||
|
||||
async def record(self, *, email: str, client_ip: str | None) -> None:
|
||||
self._session.add(
|
||||
PasswordResetAttempt(email_tried=email.strip().lower(), client_ip=client_ip)
|
||||
)
|
||||
|
||||
async def count_recent(
|
||||
self, *, email: str, client_ip: str | None, window_seconds: int
|
||||
) -> ResetRequestCounts:
|
||||
identifiant = email.strip().lower()
|
||||
meme_email = PasswordResetAttempt.email_tried == identifiant
|
||||
meme_ip = PasswordResetAttempt.client_ip == client_ip
|
||||
|
||||
requete = select(
|
||||
func.count().filter(meme_email),
|
||||
func.count().filter(meme_ip),
|
||||
).where(
|
||||
PasswordResetAttempt.occurred_at
|
||||
> datetime.now(UTC) - timedelta(seconds=window_seconds),
|
||||
meme_email | meme_ip,
|
||||
)
|
||||
|
||||
par_identifiant, par_ip = (await self._session.execute(requete)).one()
|
||||
return ResetRequestCounts(per_identifier=par_identifiant, per_ip=par_ip)
|
||||
@@ -0,0 +1,78 @@
|
||||
# Piège : `consume()` est une seule instruction, sur le modèle de `claim_for_rotation()` du
|
||||
# jeton de rafraîchissement. Un SELECT puis un UPDATE laisseraient une fenêtre où deux
|
||||
# soumissions concurrentes du même lien réussiraient toutes les deux.
|
||||
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from uuid import UUID
|
||||
|
||||
from sqlalchemy import func, select, update
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.models.password_reset_token import PasswordResetToken
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class ConsumedResetToken:
|
||||
id: UUID
|
||||
user_id: UUID
|
||||
|
||||
|
||||
class PasswordResetTokenRepository:
|
||||
def __init__(self, session: AsyncSession) -> None:
|
||||
self._session = session
|
||||
|
||||
async def create(
|
||||
self,
|
||||
*,
|
||||
user_id: UUID,
|
||||
token_hash: bytes,
|
||||
expires_at: datetime,
|
||||
client_ip: str | None,
|
||||
user_agent: str | None,
|
||||
) -> PasswordResetToken:
|
||||
jeton = PasswordResetToken(
|
||||
user_id=user_id,
|
||||
token_hash=token_hash,
|
||||
expires_at=expires_at,
|
||||
client_ip=client_ip,
|
||||
user_agent=user_agent,
|
||||
)
|
||||
self._session.add(jeton)
|
||||
await self._session.flush()
|
||||
return jeton
|
||||
|
||||
async def consume(self, token_hash: bytes) -> ConsumedResetToken | None:
|
||||
requete = (
|
||||
update(PasswordResetToken)
|
||||
.where(
|
||||
PasswordResetToken.token_hash == token_hash,
|
||||
PasswordResetToken.consumed_at.is_(None),
|
||||
PasswordResetToken.expires_at > func.clock_timestamp(),
|
||||
)
|
||||
.values(consumed_at=func.clock_timestamp())
|
||||
.returning(PasswordResetToken.id, PasswordResetToken.user_id)
|
||||
)
|
||||
ligne = (await self._session.execute(requete)).one_or_none()
|
||||
if ligne is None:
|
||||
return None
|
||||
return ConsumedResetToken(id=ligne.id, user_id=ligne.user_id)
|
||||
|
||||
# Piège : simple SELECT, volontairement pas atomique avec la consommation. Sert seulement
|
||||
# au feedback UX (jeton encore valide ?) ; `consume()` reste la seule source de vérité.
|
||||
async def exists_valid(self, token_hash: bytes) -> bool:
|
||||
requete = select(PasswordResetToken.id).where(
|
||||
PasswordResetToken.token_hash == token_hash,
|
||||
PasswordResetToken.consumed_at.is_(None),
|
||||
PasswordResetToken.expires_at > func.clock_timestamp(),
|
||||
)
|
||||
return (await self._session.execute(requete)).first() is not None
|
||||
|
||||
async def invalidate_all_for_user(self, user_id: UUID) -> int:
|
||||
resultat = await self._session.execute(
|
||||
update(PasswordResetToken)
|
||||
.where(PasswordResetToken.user_id == user_id, PasswordResetToken.consumed_at.is_(None))
|
||||
.values(consumed_at=func.clock_timestamp())
|
||||
.returning(PasswordResetToken.id)
|
||||
)
|
||||
return len(resultat.all())
|
||||
@@ -1,4 +1,5 @@
|
||||
from collections.abc import Sequence
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
@@ -19,3 +20,23 @@ class ReadingRepository:
|
||||
.order_by(Reading.site_id, Reading.timestamp.desc())
|
||||
)
|
||||
return (await self._session.execute(requete)).scalars().all()
|
||||
|
||||
async def list_history(
|
||||
self,
|
||||
*,
|
||||
start: datetime,
|
||||
end: datetime,
|
||||
site_id: str | None = None,
|
||||
limit: int,
|
||||
offset: int,
|
||||
) -> Sequence[Reading]:
|
||||
requete = (
|
||||
select(Reading)
|
||||
.where(Reading.timestamp >= start, Reading.timestamp < end)
|
||||
.order_by(Reading.timestamp.desc(), Reading.reading_id.desc())
|
||||
.limit(limit)
|
||||
.offset(offset)
|
||||
)
|
||||
if site_id is not None:
|
||||
requete = requete.where(Reading.site_id == site_id)
|
||||
return (await self._session.scalars(requete)).all()
|
||||
|
||||
@@ -1,17 +1,45 @@
|
||||
# Contrainte : le mot de passe est borné à 128 caractères. Sans plafond, une chaîne de dix
|
||||
# mégaoctets ferait travailler Argon2 gratuitement, à la charge du serveur.
|
||||
# Contrainte : `SPECIAL_CHARACTERS` doit rester identique à `password.validator.ts` côté
|
||||
# frontend. `\w`/`\d` divergent entre Python (Unicode) et JavaScript (ASCII) : une classe
|
||||
# explicite, plutôt qu'une négation, évite qu'un mot de passe soit accepté d'un côté et
|
||||
# rejeté de l'autre (ex. "Sécurité1", où "é" comptait comme "spécial" pour Python seul).
|
||||
|
||||
import re
|
||||
from typing import Literal, Self
|
||||
from uuid import UUID
|
||||
|
||||
from pydantic import BaseModel, ConfigDict, EmailStr, Field
|
||||
from pydantic import BaseModel, ConfigDict, EmailStr, Field, field_validator
|
||||
|
||||
from app.core.principal import Principal
|
||||
from app.core.roles import AccountKind, Role
|
||||
|
||||
PASSWORD_MIN_LENGTH = 12
|
||||
PASSWORD_MIN_LENGTH = 8
|
||||
PASSWORD_MAX_LENGTH = 128
|
||||
|
||||
SPECIAL_CHARACTERS = "!@#$%^&*()-_=+[]{};:,.?"
|
||||
|
||||
_MAJUSCULE = re.compile(r"[A-ZÀ-ÖØ-Þ]")
|
||||
_MINUSCULE = re.compile(r"[a-zà-öø-þ]")
|
||||
_CHIFFRE = re.compile(r"[0-9]")
|
||||
_SPECIAL = re.compile(r"[" + re.escape(SPECIAL_CHARACTERS) + r"]")
|
||||
|
||||
|
||||
def valide_complexite(mot_de_passe: str) -> str:
|
||||
manquants = [
|
||||
nom
|
||||
for nom, motif in (
|
||||
("une majuscule", _MAJUSCULE),
|
||||
("une minuscule", _MINUSCULE),
|
||||
("un chiffre", _CHIFFRE),
|
||||
("un caractère spécial", _SPECIAL),
|
||||
)
|
||||
if not motif.search(mot_de_passe)
|
||||
]
|
||||
if manquants:
|
||||
raise ValueError(f"Le mot de passe doit contenir au moins {', '.join(manquants)}")
|
||||
return mot_de_passe
|
||||
|
||||
|
||||
class LoginRequest(BaseModel):
|
||||
email: EmailStr
|
||||
@@ -22,6 +50,25 @@ class PasswordChangeRequest(BaseModel):
|
||||
current_password: str = Field(min_length=1, max_length=PASSWORD_MAX_LENGTH)
|
||||
new_password: str = Field(min_length=PASSWORD_MIN_LENGTH, max_length=PASSWORD_MAX_LENGTH)
|
||||
|
||||
@field_validator("new_password")
|
||||
@classmethod
|
||||
def _new_password_est_complexe(cls, valeur: str) -> str:
|
||||
return valide_complexite(valeur)
|
||||
|
||||
|
||||
class ForgotPasswordRequest(BaseModel):
|
||||
email: EmailStr
|
||||
|
||||
|
||||
class ResetPasswordRequest(BaseModel):
|
||||
token: str = Field(min_length=1)
|
||||
new_password: str = Field(min_length=PASSWORD_MIN_LENGTH, max_length=PASSWORD_MAX_LENGTH)
|
||||
|
||||
@field_validator("new_password")
|
||||
@classmethod
|
||||
def _new_password_est_complexe(cls, valeur: str) -> str:
|
||||
return valide_complexite(valeur)
|
||||
|
||||
|
||||
class PrincipalResponse(BaseModel):
|
||||
model_config = ConfigDict(from_attributes=True)
|
||||
@@ -37,6 +84,10 @@ class PrincipalResponse(BaseModel):
|
||||
return cls.model_validate(principal)
|
||||
|
||||
|
||||
class ResetTokenValidationResponse(BaseModel):
|
||||
valid: bool
|
||||
|
||||
|
||||
class TokenResponse(BaseModel):
|
||||
access_token: str
|
||||
token_type: Literal["bearer"] = "bearer" # noqa: S105
|
||||
|
||||
@@ -0,0 +1,45 @@
|
||||
from datetime import datetime
|
||||
from decimal import Decimal
|
||||
from enum import StrEnum
|
||||
from typing import Any
|
||||
|
||||
from pydantic import BaseModel, ConfigDict
|
||||
|
||||
|
||||
class ReadingSource(StrEnum):
|
||||
CSV = "csv"
|
||||
API_CURRENT = "api_current"
|
||||
API_HISTORY = "api_history"
|
||||
|
||||
|
||||
class ReadingDataQuality(StrEnum):
|
||||
GOOD = "good"
|
||||
PARTIAL = "partial"
|
||||
DEGRADED = "degraded"
|
||||
CRITICAL = "critical"
|
||||
|
||||
|
||||
class ReadingResponse(BaseModel):
|
||||
model_config = ConfigDict(from_attributes=True)
|
||||
|
||||
reading_id: int
|
||||
site_id: str
|
||||
timestamp: datetime
|
||||
source: ReadingSource
|
||||
consumption_kw: float | None
|
||||
consumption_kwh: float | None
|
||||
# Piège : `Decimal` (miroir de `Numeric(14, 2)` en base, pour ne pas arrondir un montant)
|
||||
# sérialise en chaîne dans le JSON, pas en nombre — un consommateur qui ferait un `parseFloat`
|
||||
# naïf perdrait la précision que ce choix visait à garder.
|
||||
consumption_euros: Decimal | None
|
||||
voltage_v: float | None
|
||||
current_a: float | None
|
||||
power_factor: float | None
|
||||
temperature_celsius: float | None
|
||||
humidity_percent: float | None
|
||||
solar_irradiance_wm2: float | None
|
||||
is_working_hours: bool | None
|
||||
data_quality: ReadingDataQuality | None
|
||||
null_reasons: list[str] | None
|
||||
imputed_values: dict[str, Any] | None
|
||||
imputation_method: str | None
|
||||
@@ -14,7 +14,11 @@ from datetime import UTC, datetime, timedelta
|
||||
from typing import NoReturn, Protocol
|
||||
from uuid import UUID, uuid4
|
||||
|
||||
from fastapi import BackgroundTasks
|
||||
|
||||
from app.core.hashing import Argon2Hasher
|
||||
from app.core.logging import get_logger
|
||||
from app.core.mailer import Mailer
|
||||
from app.core.principal import Principal
|
||||
from app.core.roles import AccountKind, Role
|
||||
from app.core.security import (
|
||||
@@ -28,9 +32,13 @@ from app.models.login_attempt import LoginOutcome
|
||||
from app.models.refresh_token import RevocationReason
|
||||
from app.repositories.audit_log import AuditLogRepository
|
||||
from app.repositories.login_attempt import LoginAttemptRepository
|
||||
from app.repositories.password_reset_attempt import PasswordResetAttemptRepository
|
||||
from app.repositories.password_reset_token import PasswordResetTokenRepository
|
||||
from app.repositories.refresh_token import RefreshTokenRepository
|
||||
from app.repositories.user import UserRepository
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
class Transaction(Protocol):
|
||||
async def commit(self) -> None: ...
|
||||
@@ -54,6 +62,10 @@ class RateLimitedError(AuthError):
|
||||
self.retry_after = retry_after
|
||||
|
||||
|
||||
class InvalidOrExpiredResetTokenError(AuthError):
|
||||
pass
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class LoginPolicy:
|
||||
window_seconds: int
|
||||
@@ -62,6 +74,15 @@ class LoginPolicy:
|
||||
max_failures_per_identifier: int
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class PasswordResetPolicy:
|
||||
window_seconds: int
|
||||
max_requests_per_identifier: int
|
||||
max_requests_per_ip: int
|
||||
token_ttl: timedelta
|
||||
frontend_reset_url: str
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class AuthenticatedSession:
|
||||
principal: Principal
|
||||
@@ -83,6 +104,10 @@ class AuthService:
|
||||
token_policy: TokenPolicy,
|
||||
login_policy: LoginPolicy,
|
||||
refresh_ttl: timedelta,
|
||||
reset_tokens: PasswordResetTokenRepository,
|
||||
reset_attempts: PasswordResetAttemptRepository,
|
||||
reset_policy: PasswordResetPolicy,
|
||||
mailer: Mailer,
|
||||
) -> None:
|
||||
self._users = users
|
||||
self._attempts = attempts
|
||||
@@ -93,6 +118,10 @@ class AuthService:
|
||||
self._token_policy = token_policy
|
||||
self._login_policy = login_policy
|
||||
self._refresh_ttl = refresh_ttl
|
||||
self._reset_tokens = reset_tokens
|
||||
self._reset_attempts = reset_attempts
|
||||
self._reset_policy = reset_policy
|
||||
self._mailer = mailer
|
||||
|
||||
async def authenticate(
|
||||
self, *, email: str, password: str, client_ip: str | None, user_agent: str | None
|
||||
@@ -200,6 +229,102 @@ class AuthService:
|
||||
rafraichi = await self._users.get_by_id(principal.id)
|
||||
return self._session(self._en_principal(rafraichi or compte), secret)
|
||||
|
||||
async def request_password_reset(
|
||||
self,
|
||||
*,
|
||||
email: str,
|
||||
client_ip: str | None,
|
||||
user_agent: str | None,
|
||||
background_tasks: BackgroundTasks,
|
||||
) -> None:
|
||||
await self._refuse_si_limite_reset(email=email, client_ip=client_ip)
|
||||
|
||||
compte = await self._users.get_by_email(email)
|
||||
# Piège : le hachage factice équilibre le temps de réponse sur un compte inconnu, comme
|
||||
# `authenticate()`. La réponse et sa forme restent identiques dans tous les cas : compte
|
||||
# inconnu, compte inactif, ou email envoyé avec succès. L'envoi SMTP lui-même est différé
|
||||
# en tâche de fond : le laisser dans le chemin de réponse rouvrirait le même oracle par le
|
||||
# temps (aller-retour réseau) et par la forme (500 si le relais SMTP échoue, contre 202).
|
||||
if compte is None or not compte.is_active or compte.kind != AccountKind.HUMAIN.value:
|
||||
await self._hasher.verify_dummy()
|
||||
await self._reset_attempts.record(email=email, client_ip=client_ip)
|
||||
await self._transaction.commit()
|
||||
return
|
||||
|
||||
await self._reset_tokens.invalidate_all_for_user(compte.id)
|
||||
secret = generate_refresh_secret()
|
||||
await self._reset_tokens.create(
|
||||
user_id=compte.id,
|
||||
token_hash=fingerprint_refresh(secret),
|
||||
expires_at=datetime.now(UTC) + self._reset_policy.token_ttl,
|
||||
client_ip=client_ip,
|
||||
user_agent=user_agent,
|
||||
)
|
||||
await self._reset_attempts.record(email=email, client_ip=client_ip)
|
||||
await self._audit.record(
|
||||
action=AuditAction.MOT_DE_PASSE_OUBLIE_DEMANDE,
|
||||
actor_label=compte.email,
|
||||
target_type="app_user",
|
||||
target_id=str(compte.id),
|
||||
client_ip=client_ip,
|
||||
user_agent=user_agent,
|
||||
)
|
||||
await self._transaction.commit()
|
||||
|
||||
lien = f"{self._reset_policy.frontend_reset_url}?token={secret}"
|
||||
background_tasks.add_task(self._envoie_email_reset, compte.email, lien)
|
||||
|
||||
async def _envoie_email_reset(self, email: str, reset_url: str) -> None:
|
||||
try:
|
||||
await self._mailer.send_password_reset_email(to=email, reset_url=reset_url)
|
||||
except Exception:
|
||||
logger.exception("auth.password_reset.mail_failed")
|
||||
|
||||
# Piège : lecture seule, pas d'appel à `consume()`. Aucune limitation de débit n'est
|
||||
# nécessaire ici : le jeton est un secret de 256 bits (`generate_refresh_secret`), donc
|
||||
# non brute-forçable, et cette route n'apprend rien sur l'existence d'un compte ou d'un
|
||||
# email, seulement si le lien déjà en main du visiteur est encore valide.
|
||||
async def is_reset_token_valid(self, token: str) -> bool:
|
||||
return await self._reset_tokens.exists_valid(fingerprint_refresh(token))
|
||||
|
||||
async def confirm_password_reset(
|
||||
self, *, token: str, new_password: str, client_ip: str | None, user_agent: str | None
|
||||
) -> AuthenticatedSession:
|
||||
revendique = await self._reset_tokens.consume(fingerprint_refresh(token))
|
||||
if revendique is None:
|
||||
raise InvalidOrExpiredResetTokenError("Lien invalide ou expiré")
|
||||
|
||||
# Piège : le jeton peut avoir été émis avant une désactivation du compte. Sans cette
|
||||
# relecture, un lien encore valide (15 min) changerait quand même le mot de passe d'un
|
||||
# compte désactivé, réutilisable dès sa réactivation.
|
||||
compte = await self._users.get_by_id(revendique.user_id)
|
||||
if compte is None or not compte.is_active or compte.kind != AccountKind.HUMAIN.value:
|
||||
raise InvalidOrExpiredResetTokenError("Lien invalide ou expiré")
|
||||
|
||||
await self._users.update_password(
|
||||
revendique.user_id, await self._hasher.hash(new_password), must_change_password=False
|
||||
)
|
||||
revoquees = await self._refresh.revoke_all_for_user(
|
||||
revendique.user_id, RevocationReason.CHANGEMENT_MOT_DE_PASSE
|
||||
)
|
||||
secret = await self._ouvre_une_famille(
|
||||
user_id=revendique.user_id, client_ip=client_ip, user_agent=user_agent
|
||||
)
|
||||
await self._audit.record(
|
||||
action=AuditAction.MOT_DE_PASSE_REINITIALISE_PAR_SOI,
|
||||
target_type="app_user",
|
||||
target_id=str(revendique.user_id),
|
||||
client_ip=client_ip,
|
||||
user_agent=user_agent,
|
||||
detail={"sessions_revoquees": revoquees},
|
||||
)
|
||||
await self._transaction.commit()
|
||||
|
||||
compte = await self._users.get_by_id(revendique.user_id)
|
||||
if compte is None:
|
||||
raise SessionRejectedError("Compte introuvable")
|
||||
return self._session(self._en_principal(compte), secret)
|
||||
|
||||
async def logout_all(self, principal: Principal) -> int:
|
||||
revoquees = await self._refresh.revoke_all_for_user(
|
||||
principal.id, RevocationReason.DECONNEXION
|
||||
@@ -307,6 +432,23 @@ class AuthService:
|
||||
await self._transaction.commit()
|
||||
raise RateLimitedError(politique.window_seconds)
|
||||
|
||||
async def _refuse_si_limite_reset(self, *, email: str, client_ip: str | None) -> None:
|
||||
politique = self._reset_policy
|
||||
compteurs = await self._reset_attempts.count_recent(
|
||||
email=email, client_ip=client_ip, window_seconds=politique.window_seconds
|
||||
)
|
||||
|
||||
depasse = (
|
||||
compteurs.per_identifier >= politique.max_requests_per_identifier
|
||||
or compteurs.per_ip >= politique.max_requests_per_ip
|
||||
)
|
||||
if not depasse:
|
||||
return
|
||||
|
||||
await self._reset_attempts.record(email=email, client_ip=client_ip)
|
||||
await self._transaction.commit()
|
||||
raise RateLimitedError(politique.window_seconds)
|
||||
|
||||
async def _echoue(
|
||||
self,
|
||||
email: str,
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
from collections.abc import Sequence
|
||||
from datetime import UTC, datetime, timedelta
|
||||
|
||||
from app.models.energy import Reading
|
||||
from app.repositories.reading import ReadingRepository
|
||||
|
||||
FENETRE_PAR_DEFAUT = timedelta(hours=24)
|
||||
FENETRE_MAXIMALE = timedelta(days=90)
|
||||
|
||||
|
||||
class FenetreInverseeError(Exception):
|
||||
"""`start` est postérieur ou égal à `end`."""
|
||||
|
||||
|
||||
class FenetreTropLargeError(Exception):
|
||||
"""L'écart entre `start` et `end` dépasse `FENETRE_MAXIMALE`."""
|
||||
|
||||
|
||||
class ReadingService:
|
||||
def __init__(self, *, readings: ReadingRepository) -> None:
|
||||
self._readings = readings
|
||||
|
||||
async def list_history(
|
||||
self,
|
||||
*,
|
||||
site_id: str | None = None,
|
||||
start: datetime | None = None,
|
||||
end: datetime | None = None,
|
||||
limit: int,
|
||||
offset: int,
|
||||
) -> Sequence[Reading]:
|
||||
debut, fin = self._resoudre_fenetre(start, end)
|
||||
return await self._readings.list_history(
|
||||
site_id=site_id, start=debut, end=fin, limit=limit, offset=offset
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _resoudre_fenetre(
|
||||
start: datetime | None, end: datetime | None
|
||||
) -> tuple[datetime, datetime]:
|
||||
# Piège : un datetime naïf (sans fuseau dans la chaîne ISO reçue) fait échouer la
|
||||
# comparaison à `reading.timestamp` (`timestamptz`) au niveau du pilote, en 500 plutôt
|
||||
# qu'un refus propre. On le traite comme de l'UTC plutôt que de le rejeter.
|
||||
debut = _vers_utc(start)
|
||||
fin = _vers_utc(end) or datetime.now(UTC)
|
||||
if debut is None:
|
||||
debut = fin - FENETRE_PAR_DEFAUT
|
||||
|
||||
if debut >= fin:
|
||||
raise FenetreInverseeError
|
||||
if fin - debut > FENETRE_MAXIMALE:
|
||||
raise FenetreTropLargeError
|
||||
return debut, fin
|
||||
|
||||
|
||||
def _vers_utc(instant: datetime | None) -> datetime | None:
|
||||
if instant is None:
|
||||
return None
|
||||
return instant if instant.tzinfo is not None else instant.replace(tzinfo=UTC)
|
||||
Reference in New Issue
Block a user