Compare commits

..
Author SHA1 Message Date
Johan LEROY 3eb5a0e8dc style(backend): applique ruff format au modele data
Backend / Lint, typage et tests (push) Successful in 1m7s
La cible make check ne lance que ruff check ; la CI lance en plus
ruff format --check, qui refusait la contrainte unique repliee.
2026-09-16 08:43:12 +02:00
Johan LEROY c733ccfc62 refactor(backend): passe les tables data au singulier et clarifie alert_id
La convention de docs/architecture/40-data.md impose des noms de tables au
singulier, que les quatre tables d'authentification respectent déjà. Les six
tables data passent donc au singulier, avec leurs contraintes et leurs index.

La révision n'étant appliquée que sur des bases locales, elle est modifiée sur
place plutôt que doublée d'une migration de renommage.

alert_id désignait deux colonnes différentes : la clé métier text de l'API Mock
et la clé étrangère bigint de recommendation. La première devient
source_alert_id, la seconde pointe désormais vers alert.alert_id.
2026-09-16 08:41:48 +02:00
Johan LEROY e3e0e843d0 fix(backend): rebranche la revision data sur la tete d'authentification
Le merge de dev apporte trois revisions d'authentification qui partent de la
meme racine 5353c0e4f094 que la revision data. Git ne signale rien, mais
alembic upgrade head refuse de choisir entre deux tetes.

La revision data se greffe desormais sur 821f71be74c0, ce qui rend la chaine
lineaire.
2026-09-15 16:46:37 +02:00
Johan LEROY c3b7c818aa Merge remote-tracking branch 'origin/dev' into feat/data-schema
# Conflicts:
#	apps/backend/app/models/__init__.py
2026-09-15 16:46:16 +02:00
Meryemel-gham c04ce9a9ae test(apps): verifie les contraintes du schema data 2026-09-15 16:16:40 +02:00
Meryemel-gham 128133761f feat(apps): cree les six tables data et l'hypertable readings 2026-09-15 16:16:18 +02:00
Meryemel-gham b032f084fc fix(apps): gere les caracteres encodes dans l'URL Alembic 2026-09-15 16:15:21 +02:00
6 changed files with 768 additions and 46 deletions
+1 -1
View File
@@ -19,7 +19,7 @@ config = context.config
if config.config_file_name is not None: if config.config_file_name is not None:
fileConfig(config.config_file_name) fileConfig(config.config_file_name)
config.set_main_option("sqlalchemy.url", get_settings().database_url) config.set_main_option("sqlalchemy.url", get_settings().database_url.replace("%", "%%"))
target_metadata = Base.metadata target_metadata = Base.metadata
@@ -0,0 +1,218 @@
"""Création des six tables Data et de l'hypertable reading.
Revision ID: e6d2026091501
Revises: 821f71be74c0
"""
from alembic import op
import sqlalchemy as sa
from sqlalchemy.dialects import postgresql
revision = "e6d2026091501"
down_revision = "821f71be74c0"
branch_labels = None
depends_on = None
def upgrade() -> None:
# ### commands auto generated by Alembic - please adjust! ###
op.create_table(
"dataset",
sa.Column("dataset_id", sa.BigInteger(), autoincrement=True, nullable=False),
sa.Column("dataset_name", sa.Text(), nullable=False),
sa.Column("archive_sha256", sa.String(length=64), nullable=False),
sa.Column("storage_uri", sa.Text(), nullable=False),
sa.Column("source_timezone", sa.Text(), nullable=True),
sa.Column(
"metadata", postgresql.JSONB(none_as_null=True, astext_type=sa.Text()), nullable=False
),
sa.CheckConstraint("dataset_id > 0", name="ck_dataset_positive_id"),
sa.PrimaryKeyConstraint("dataset_id"),
sa.UniqueConstraint("archive_sha256", name="uq_dataset_archive_sha256"),
)
op.create_table(
"site",
sa.Column("site_id", sa.Text(), nullable=False),
sa.Column("site_name", sa.Text(), nullable=False),
sa.Column("site_type", sa.Text(), nullable=False),
sa.Column("location", sa.Text(), nullable=True),
sa.Column("capacity_kw", sa.Double(), nullable=True),
sa.Column("status", sa.Text(), nullable=True),
sa.PrimaryKeyConstraint("site_id"),
)
op.create_table(
"prediction",
sa.Column("prediction_id", sa.BigInteger(), autoincrement=True, nullable=False),
sa.Column("site_id", sa.Text(), nullable=False),
sa.Column(
"created_at",
sa.DateTime(timezone=True),
server_default=sa.text("now()"),
nullable=False,
),
sa.Column("target_at", sa.DateTime(timezone=True), nullable=False),
sa.Column("target_metric", sa.Text(), nullable=False),
sa.Column("period_minutes", sa.Integer(), nullable=True),
sa.Column("predicted_value", sa.Double(), nullable=True),
sa.Column("model_reference", sa.Text(), nullable=False),
sa.Column("status", sa.Text(), nullable=False),
sa.Column("failure_reason", sa.Text(), nullable=True),
sa.CheckConstraint(
"(status = 'available' AND predicted_value IS NOT NULL AND failure_reason IS NULL) OR (status IN ('insufficient_data', 'error') AND predicted_value IS NULL AND failure_reason IS NOT NULL)",
name="ck_prediction_status",
),
sa.CheckConstraint(
"target_metric <> 'consumption_kwh' OR period_minutes IS NOT NULL",
name="ck_prediction_energy_period",
),
sa.CheckConstraint(
"target_metric IN ('consumption_kwh', 'consumption_kw')", name="ck_prediction_metric"
),
sa.CheckConstraint(
"period_minutes IS NULL OR period_minutes > 0", name="ck_prediction_period"
),
sa.ForeignKeyConstraint(
["site_id"], ["site.site_id"], name="fk_prediction_site", ondelete="RESTRICT"
),
sa.PrimaryKeyConstraint("prediction_id"),
sa.UniqueConstraint("prediction_id", "site_id", name="uq_prediction_id_site"),
)
op.create_index(
"ix_prediction_site_target", "prediction", ["site_id", "target_at"], unique=False
)
op.create_table(
"reading",
sa.Column("reading_id", sa.BigInteger(), autoincrement=True, nullable=False),
sa.Column("site_id", sa.Text(), nullable=False),
sa.Column("timestamp", sa.DateTime(timezone=True), nullable=False),
sa.Column("source", sa.Text(), nullable=False),
sa.Column("dataset_id", sa.BigInteger(), nullable=True),
sa.Column("consumption_kw", sa.Double(), nullable=True),
sa.Column("consumption_kwh", sa.Double(), nullable=True),
sa.Column("consumption_euros", sa.Numeric(precision=14, scale=2), nullable=True),
sa.Column("voltage_v", sa.Double(), nullable=True),
sa.Column("current_a", sa.Double(), nullable=True),
sa.Column("power_factor", sa.Double(), nullable=True),
sa.Column("temperature_celsius", sa.Double(), nullable=True),
sa.Column("humidity_percent", sa.Double(), nullable=True),
sa.Column("solar_irradiance_wm2", sa.Double(), nullable=True),
sa.Column("is_working_hours", sa.Boolean(), nullable=True),
sa.Column("data_quality", sa.Text(), nullable=True),
sa.Column("null_reasons", postgresql.ARRAY(sa.Text()), nullable=True),
sa.Column(
"imputed_values", postgresql.JSONB(none_as_null=True, astext_type=sa.Text()), nullable=True
),
sa.Column("imputation_method", sa.Text(), nullable=True),
sa.Column(
"ingested_at",
sa.DateTime(timezone=True),
server_default=sa.text("now()"),
nullable=False,
),
sa.Column(
"raw_data", postgresql.JSONB(none_as_null=True, astext_type=sa.Text()), nullable=False
),
sa.CheckConstraint(
"(source = 'csv' AND dataset_id IS NOT NULL) OR (source IN ('api_current', 'api_history') AND dataset_id IS NULL)",
name="ck_reading_dataset_source",
),
sa.CheckConstraint(
"data_quality IS NULL OR data_quality IN ('good', 'partial', 'degraded', 'critical')",
name="ck_reading_quality",
),
sa.CheckConstraint(
"source IN ('csv', 'api_current', 'api_history')", name="ck_reading_source"
),
sa.CheckConstraint(
"(imputed_values IS NULL AND imputation_method IS NULL) OR (imputed_values IS NOT NULL AND imputation_method IS NOT NULL)",
name="ck_reading_imputation",
),
sa.ForeignKeyConstraint(
["dataset_id"], ["dataset.dataset_id"], name="fk_reading_dataset", ondelete="RESTRICT"
),
sa.ForeignKeyConstraint(
["site_id"], ["site.site_id"], name="fk_reading_site", ondelete="RESTRICT"
),
sa.PrimaryKeyConstraint("reading_id", "timestamp"),
)
op.create_index("ix_reading_dataset_id", "reading", ["dataset_id"], unique=False)
op.create_index(
"ix_reading_site_timestamp", "reading", ["site_id", "timestamp"], unique=False
)
op.create_index(
"uq_reading_source",
"reading",
["site_id", "timestamp", "source", sa.literal_column("coalesce(dataset_id, 0)")],
unique=True,
)
op.execute(
"SELECT create_hypertable('reading', by_range('timestamp'), create_default_indexes => FALSE)"
)
op.create_table(
"alert",
sa.Column("alert_id", sa.BigInteger(), autoincrement=True, nullable=False),
sa.Column("source_alert_id", sa.Text(), nullable=False),
sa.Column("site_id", sa.Text(), nullable=False),
sa.Column("source", sa.Text(), nullable=False),
sa.Column("timestamp", sa.DateTime(timezone=True), nullable=False),
sa.Column("type", sa.Text(), nullable=False),
sa.Column("severity", sa.Text(), nullable=False),
sa.Column("message", sa.Text(), nullable=False),
sa.Column("value", sa.Double(), nullable=True),
sa.Column("threshold", sa.Double(), nullable=True),
sa.Column("metric", sa.Text(), nullable=True),
sa.Column("prediction_id", sa.BigInteger(), nullable=True),
sa.Column(
"raw_data", postgresql.JSONB(none_as_null=True, astext_type=sa.Text()), nullable=False
),
sa.CheckConstraint(
"severity IN ('low', 'medium', 'high', 'critical')", name="ck_alert_severity"
),
sa.CheckConstraint("source IN ('api_mock', 'enervision')", name="ck_alert_source"),
sa.CheckConstraint(
"type IN ('spike', 'threshold', 'anomaly', 'outage', 'sensor')", name="ck_alert_type"
),
sa.ForeignKeyConstraint(
["prediction_id", "site_id"],
["prediction.prediction_id", "prediction.site_id"],
name="fk_alert_prediction_site",
ondelete="RESTRICT",
),
sa.ForeignKeyConstraint(
["site_id"], ["site.site_id"], name="fk_alert_site", ondelete="RESTRICT"
),
sa.PrimaryKeyConstraint("alert_id"),
sa.UniqueConstraint(
"source", "site_id", "source_alert_id", name="uq_alert_source_reference"
),
)
op.create_index("ix_alert_site_timestamp", "alert", ["site_id", "timestamp"], unique=False)
op.create_table(
"recommendation",
sa.Column("recommendation_id", sa.BigInteger(), autoincrement=True, nullable=False),
sa.Column("alert_id", sa.BigInteger(), nullable=False),
sa.Column("action", sa.Text(), nullable=False),
sa.Column("explanation", sa.Text(), nullable=False),
sa.Column("rule_reference", sa.Text(), nullable=False),
sa.Column(
"created_at",
sa.DateTime(timezone=True),
server_default=sa.text("now()"),
nullable=False,
),
sa.ForeignKeyConstraint(
["alert_id"], ["alert.alert_id"], name="fk_recommendation_alert", ondelete="RESTRICT"
),
sa.PrimaryKeyConstraint("recommendation_id"),
sa.UniqueConstraint("alert_id", "rule_reference", name="uq_recommendation_alert_rule"),
)
# ### end Alembic commands ###
def downgrade() -> None:
op.drop_table("recommendation")
op.drop_table("alert")
op.drop_table("reading")
op.drop_table("prediction")
op.drop_table("site")
op.drop_table("dataset")
+13 -1
View File
@@ -2,8 +2,20 @@
# --autogenerate`, qui générerait alors un drop de sa table. # --autogenerate`, qui générerait alors un drop de sa table.
from app.models.audit_log import AuditLog 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.login_attempt import LoginAttempt
from app.models.refresh_token import RefreshToken from app.models.refresh_token import RefreshToken
from app.models.user import AppUser from app.models.user import AppUser
__all__ = ["AppUser", "AuditLog", "LoginAttempt", "RefreshToken"] __all__ = [
"Alert",
"AppUser",
"AuditLog",
"Dataset",
"LoginAttempt",
"Prediction",
"Reading",
"Recommendation",
"RefreshToken",
"Site",
]
+210
View File
@@ -0,0 +1,210 @@
"""Tables du modèle de données EnerVision (CSV, API Mock et résultats ML)."""
from datetime import datetime
from decimal import Decimal
from typing import Any
from sqlalchemy import (
BigInteger,
Boolean,
CheckConstraint,
DateTime,
Double,
ForeignKey,
ForeignKeyConstraint,
Index,
Integer,
Numeric,
String,
Text,
UniqueConstraint,
func,
text,
)
from sqlalchemy.dialects.postgresql import ARRAY, JSONB
from sqlalchemy.orm import Mapped, mapped_column
from app.db.base import Base
class Dataset(Base):
__tablename__ = "dataset"
__table_args__ = (
CheckConstraint("dataset_id > 0", name="ck_dataset_positive_id"),
UniqueConstraint("archive_sha256", name="uq_dataset_archive_sha256"),
)
dataset_id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
dataset_name: Mapped[str] = mapped_column(Text)
archive_sha256: Mapped[str] = mapped_column(String(64))
storage_uri: Mapped[str] = mapped_column(Text)
source_timezone: Mapped[str | None] = mapped_column(Text)
# "metadata" est réservé par SQLAlchemy ; le nom SQL reste inchangé.
dataset_metadata: Mapped[dict[str, Any]] = mapped_column("metadata", JSONB(none_as_null=True))
class Site(Base):
__tablename__ = "site"
site_id: Mapped[str] = mapped_column(Text, primary_key=True)
site_name: Mapped[str] = mapped_column(Text)
site_type: Mapped[str] = mapped_column(Text)
location: Mapped[str | None] = mapped_column(Text)
capacity_kw: Mapped[float | None] = mapped_column(Double)
status: Mapped[str | None] = mapped_column(Text)
class Reading(Base):
__tablename__ = "reading"
__table_args__ = (
CheckConstraint(
"source IN ('csv', 'api_current', 'api_history')", name="ck_reading_source"
),
CheckConstraint(
"(source = 'csv' AND dataset_id IS NOT NULL) OR "
"(source IN ('api_current', 'api_history') AND dataset_id IS NULL)",
name="ck_reading_dataset_source",
),
CheckConstraint(
"data_quality IS NULL OR data_quality IN ('good', 'partial', 'degraded', 'critical')",
name="ck_reading_quality",
),
CheckConstraint(
"(imputed_values IS NULL AND imputation_method IS NULL) OR "
"(imputed_values IS NOT NULL AND imputation_method IS NOT NULL)",
name="ck_reading_imputation",
),
Index("ix_reading_site_timestamp", "site_id", "timestamp"),
Index("ix_reading_dataset_id", "dataset_id"),
)
reading_id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
site_id: Mapped[str] = mapped_column(
Text, ForeignKey("site.site_id", name="fk_reading_site", ondelete="RESTRICT")
)
timestamp: Mapped[datetime] = mapped_column(DateTime(timezone=True), primary_key=True)
source: Mapped[str] = mapped_column(Text)
dataset_id: Mapped[int | None] = mapped_column(
BigInteger,
ForeignKey("dataset.dataset_id", name="fk_reading_dataset", ondelete="RESTRICT"),
)
consumption_kw: Mapped[float | None] = mapped_column(Double)
consumption_kwh: Mapped[float | None] = mapped_column(Double)
consumption_euros: Mapped[Decimal | None] = mapped_column(Numeric(14, 2))
voltage_v: Mapped[float | None] = mapped_column(Double)
current_a: Mapped[float | None] = mapped_column(Double)
power_factor: Mapped[float | None] = mapped_column(Double)
temperature_celsius: Mapped[float | None] = mapped_column(Double)
humidity_percent: Mapped[float | None] = mapped_column(Double)
solar_irradiance_wm2: Mapped[float | None] = mapped_column(Double)
is_working_hours: Mapped[bool | None] = mapped_column(Boolean)
data_quality: Mapped[str | None] = mapped_column(Text)
null_reasons: Mapped[list[str] | None] = mapped_column(ARRAY(Text))
imputed_values: Mapped[dict[str, Any] | None] = mapped_column(JSONB(none_as_null=True))
imputation_method: Mapped[str | None] = mapped_column(Text)
ingested_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now()
)
raw_data: Mapped[dict[str, Any]] = mapped_column(JSONB(none_as_null=True))
Index(
"uq_reading_source",
Reading.site_id,
Reading.timestamp,
Reading.source,
func.coalesce(Reading.dataset_id, text("0")),
unique=True,
)
class Prediction(Base):
__tablename__ = "prediction"
__table_args__ = (
UniqueConstraint("prediction_id", "site_id", name="uq_prediction_id_site"),
Index("ix_prediction_site_target", "site_id", "target_at"),
CheckConstraint(
"target_metric IN ('consumption_kwh', 'consumption_kw')",
name="ck_prediction_metric",
),
CheckConstraint(
"period_minutes IS NULL OR period_minutes > 0", name="ck_prediction_period"
),
CheckConstraint(
"target_metric <> 'consumption_kwh' OR period_minutes IS NOT NULL",
name="ck_prediction_energy_period",
),
CheckConstraint(
"(status = 'available' AND predicted_value IS NOT NULL AND failure_reason IS NULL) OR "
"(status IN ('insufficient_data', 'error') AND predicted_value IS NULL "
"AND failure_reason IS NOT NULL)",
name="ck_prediction_status",
),
)
prediction_id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
site_id: Mapped[str] = mapped_column(
Text, ForeignKey("site.site_id", name="fk_prediction_site", ondelete="RESTRICT")
)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now())
target_at: Mapped[datetime] = mapped_column(DateTime(timezone=True))
target_metric: Mapped[str] = mapped_column(Text)
period_minutes: Mapped[int | None] = mapped_column(Integer)
predicted_value: Mapped[float | None] = mapped_column(Double)
model_reference: Mapped[str] = mapped_column(Text)
status: Mapped[str] = mapped_column(Text)
failure_reason: Mapped[str | None] = mapped_column(Text)
class Alert(Base):
__tablename__ = "alert"
__table_args__ = (
UniqueConstraint("source", "site_id", "source_alert_id", name="uq_alert_source_reference"),
Index("ix_alert_site_timestamp", "site_id", "timestamp"),
ForeignKeyConstraint(
["prediction_id", "site_id"],
["prediction.prediction_id", "prediction.site_id"],
name="fk_alert_prediction_site",
ondelete="RESTRICT",
),
CheckConstraint("source IN ('api_mock', 'enervision')", name="ck_alert_source"),
CheckConstraint(
"type IN ('spike', 'threshold', 'anomaly', 'outage', 'sensor')", name="ck_alert_type"
),
CheckConstraint(
"severity IN ('low', 'medium', 'high', 'critical')", name="ck_alert_severity"
),
)
alert_id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
source_alert_id: Mapped[str] = mapped_column(Text)
site_id: Mapped[str] = mapped_column(
Text, ForeignKey("site.site_id", name="fk_alert_site", ondelete="RESTRICT")
)
source: Mapped[str] = mapped_column(Text)
timestamp: Mapped[datetime] = mapped_column(DateTime(timezone=True))
type: Mapped[str] = mapped_column(Text)
severity: Mapped[str] = mapped_column(Text)
message: Mapped[str] = mapped_column(Text)
value: Mapped[float | None] = mapped_column(Double)
threshold: Mapped[float | None] = mapped_column(Double)
metric: Mapped[str | None] = mapped_column(Text)
prediction_id: Mapped[int | None] = mapped_column(BigInteger)
raw_data: Mapped[dict[str, Any]] = mapped_column(JSONB(none_as_null=True))
class Recommendation(Base):
__tablename__ = "recommendation"
__table_args__ = (
UniqueConstraint("alert_id", "rule_reference", name="uq_recommendation_alert_rule"),
)
recommendation_id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
alert_id: Mapped[int] = mapped_column(
BigInteger,
ForeignKey("alert.alert_id", name="fk_recommendation_alert", ondelete="RESTRICT"),
)
action: Mapped[str] = mapped_column(Text)
explanation: Mapped[str] = mapped_column(Text)
rule_reference: Mapped[str] = mapped_column(Text)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now())
+261
View File
@@ -0,0 +1,261 @@
from collections.abc import AsyncIterator
from datetime import UTC, datetime
from uuid import uuid4
import pytest
from sqlalchemy import insert, select, text
from sqlalchemy.engine import make_url
from sqlalchemy.exc import IntegrityError
from sqlalchemy.ext.asyncio import AsyncConnection, create_async_engine
from app.core.config import get_settings
from app.models.energy import Alert, Dataset, Prediction, Reading, Recommendation, Site
pytestmark = pytest.mark.integration
MOMENT = datetime(2024, 1, 1, tzinfo=UTC)
@pytest.fixture
async def data_connection() -> AsyncIterator[AsyncConnection]:
url = make_url(get_settings().database_url)
if url.database != "enervision_test":
pytest.fail("Ces tests exigent DATABASE_URL vers enervision_test.")
engine = create_async_engine(url)
try:
async with engine.connect() as connection:
transaction = await connection.begin()
try:
yield connection
finally:
await transaction.rollback()
finally:
await engine.dispose()
@pytest.fixture
async def data_site(data_connection: AsyncConnection) -> str:
site_id = f"TEST-{uuid4()}"
await data_connection.execute(
insert(Site).values(site_id=site_id, site_name="Site de test", site_type="office")
)
return site_id
async def test_reading_is_a_time_hypertable_when_migrated(
data_connection: AsyncConnection,
) -> None:
query = text(
"SELECT column_name FROM timescaledb_information.dimensions "
"WHERE hypertable_schema = 'public' AND hypertable_name = 'reading'"
)
result = await data_connection.execute(query)
assert result.scalars().all() == ["timestamp"]
async def test_reading_preserves_null_and_zero_when_inserted(
data_connection: AsyncConnection, data_site: str
) -> None:
statement = insert(Reading).values(
site_id=data_site,
timestamp=MOMENT,
source="api_current",
consumption_kw=None,
consumption_kwh=0,
data_quality="partial",
null_reasons=["sensor_failure"],
raw_data={"consumption_kw": None},
imputed_values=None,
imputation_method=None,
)
await data_connection.execute(statement)
result = (
await data_connection.execute(
select(
Reading.consumption_kw,
Reading.consumption_kwh,
Reading.raw_data,
Reading.imputed_values,
).where(Reading.site_id == data_site)
)
).one()
assert tuple(result) == (None, 0, {"consumption_kw": None}, None)
@pytest.mark.parametrize("source", ["csv", "api_current", "api_history"])
async def test_duplicate_reading_is_rejected_when_key_matches(
data_connection: AsyncConnection, data_site: str, source: str
) -> None:
dataset_id = None
if source == "csv":
dataset_id = (
await data_connection.execute(
insert(Dataset.__table__)
.values(
dataset_name="Archive de test",
archive_sha256=uuid4().hex + uuid4().hex,
storage_uri="test://archive",
metadata={},
)
.returning(Dataset.dataset_id)
)
).scalar_one()
statement = insert(Reading).values(
site_id=data_site,
timestamp=MOMENT,
source=source,
dataset_id=dataset_id,
raw_data={},
)
await data_connection.execute(statement)
with pytest.raises(IntegrityError):
async with data_connection.begin_nested():
await data_connection.execute(statement)
@pytest.mark.parametrize(
"changes",
[
{"source": "csv"},
{"source": "unknown"},
{"site_id": "UNKNOWN-SITE"},
{"data_quality": "unknown"},
{"imputed_values": {"consumption_kw": 12}},
{"imputation_method": "mean-v1"},
],
ids=[
"csv_sans_dataset",
"source_inconnue",
"site_absent",
"qualite_inconnue",
"imputation_sans_methode",
"methode_sans_imputation",
],
)
async def test_invalid_reading_is_rejected_when_constraints_fail(
data_connection: AsyncConnection, data_site: str, changes: dict[str, object]
) -> None:
values: dict[str, object] = {
"site_id": data_site,
"timestamp": MOMENT,
"source": "api_current",
"raw_data": {},
}
values.update(changes)
with pytest.raises(IntegrityError):
async with data_connection.begin_nested():
await data_connection.execute(insert(Reading).values(**values))
async def test_prediction_requires_period_when_energy_is_predicted(
data_connection: AsyncConnection, data_site: str
) -> None:
statement = insert(Prediction).values(
site_id=data_site,
target_at=MOMENT,
target_metric="consumption_kwh",
predicted_value=12,
status="available",
model_reference="test-model/1",
)
with pytest.raises(IntegrityError):
async with data_connection.begin_nested():
await data_connection.execute(statement)
async def test_unavailable_prediction_preserves_null_when_inserted(
data_connection: AsyncConnection, data_site: str
) -> None:
statement = (
insert(Prediction)
.values(
site_id=data_site,
target_at=MOMENT,
target_metric="consumption_kw",
status="insufficient_data",
failure_reason="Historique trop court",
model_reference="test-model/1",
)
.returning(Prediction.predicted_value)
)
value = (await data_connection.execute(statement)).scalar_one()
assert value is None
async def test_alert_rejects_prediction_when_site_differs(
data_connection: AsyncConnection, data_site: str
) -> None:
other_site = f"TEST-{uuid4()}"
await data_connection.execute(
insert(Site).values(site_id=other_site, site_name="Autre site", site_type="office")
)
prediction_id = (
await data_connection.execute(
insert(Prediction)
.values(
site_id=data_site,
target_at=MOMENT,
target_metric="consumption_kw",
predicted_value=12,
status="available",
model_reference="test-model/1",
)
.returning(Prediction.prediction_id)
)
).scalar_one()
with pytest.raises(IntegrityError):
async with data_connection.begin_nested():
await data_connection.execute(
insert(Alert).values(
source_alert_id=str(uuid4()),
site_id=other_site,
source="enervision",
timestamp=MOMENT,
type="spike",
severity="high",
message="Test",
prediction_id=prediction_id,
raw_data={},
)
)
async def test_recommendation_is_unique_when_alert_and_rule_match(
data_connection: AsyncConnection, data_site: str
) -> None:
alert_id = (
await data_connection.execute(
insert(Alert)
.values(
source_alert_id=str(uuid4()),
site_id=data_site,
source="api_mock",
timestamp=MOMENT,
type="spike",
severity="high",
message="Test",
raw_data={},
)
.returning(Alert.alert_id)
)
).scalar_one()
statement = insert(Recommendation).values(
alert_id=alert_id,
action="Vérifier la consommation",
explanation="Pic détecté",
rule_reference="spike-v1",
)
await data_connection.execute(statement)
with pytest.raises(IntegrityError):
async with data_connection.begin_nested():
await data_connection.execute(statement)
+65 -44
View File
@@ -4,12 +4,12 @@ PostgreSQL 17 avec l'extension TimescaleDB. Le choix, ses alternatives et ses co
dans l'[ADR 0001](../adr/0001-postgresql-timescaledb.md), qui fait foi. Ce document décrit le dans l'[ADR 0001](../adr/0001-postgresql-timescaledb.md), qui fait foi. Ce document décrit le
système qui en découle. système qui en découle.
## Ce que couvre ce document ## Avertissement
**Dix tables applicatives existent** : quatre pour l'authentification, six pour les données **Aucune table applicative n'existe à ce jour.** `Base.metadata` est vide, `app/models/` ne
d'énergie, dont l'hypertable `reading`. Les sections marquées `Fait` relèvent le code. Celles contient qu'un commentaire, l'unique révision Alembic ne crée aucune table, et aucune hypertable
marquées `Cible` décrivent ce qui n'est pas écrit, au premier rang desquelles la chaîne n'a été déclarée. Tout ce qui suit sous le statut `Cible` est une proposition de structure, pas un
d'ingestion, les agrégats continus, la compression et la rétention. relevé du code. Le modèle sera arrêté au jalon J2.
## Trois emplacements, trois rôles ## Trois emplacements, trois rôles
@@ -35,7 +35,7 @@ Statut : `Fait`.
- `db/init/100-extensions.sql` crée l'extension `timescaledb`. - `db/init/100-extensions.sql` crée l'extension `timescaledb`.
- `db/init/110-test-database.sql` crée `enervision_test`, dont le nom est attendu en dur par - `db/init/110-test-database.sql` crée `enervision_test`, dont le nom est attendu en dur par
`apps/backend/tests/conftest.py`. `apps/backend/tests/conftest.py`.
- Cinq révisions Alembic. La première, `5353c0e4f094`, **ne crée aucune table** : elle - Quatre révisions Alembic. La première, `5353c0e4f094`, **ne crée aucune table** : elle
établit `alembic_version` et refuse de s'appliquer si l'extension manque : établit `alembic_version` et refuse de s'appliquer si l'extension manque :
```sql ```sql
@@ -48,18 +48,16 @@ Cette garde forme paire avec le 503 de `/api/v1/health/ready`. Un bootstrap saut
au démarrage de l'API : ces deux gardes le rendent visible tôt, des deux côtés. au démarrage de l'API : ces deux gardes le rendent visible tôt, des deux côtés.
Les trois suivantes créent les tables de l'authentification, décrites plus bas : `app_user`, Les trois suivantes créent les tables de l'authentification, décrites plus bas : `app_user`,
puis `login_attempt` et `audit_log`, puis `refresh_token`. La cinquième, `e6d2026091501`, crée puis `login_attempt` et `audit_log`, puis `refresh_token`.
les six tables de données décrites en fin de document et déclare l'hypertable `reading`.
## Cycle de vie d'une mesure ## Cycle de vie d'une mesure
Statut : `Cible`, sauf l'hypertable `reading` qui existe. Ni l'ingestion, ni les agrégats Statut : `Cible`. Aucun de ces maillons n'existe.
continus, ni la compression, ni la rétention ne sont écrits.
```mermaid ```mermaid
flowchart LR flowchart LR
src["Source de mesures"] -.-> ing["Ingestion Airflow"] src["Source de mesures"] -.-> ing["Ingestion Airflow"]
ing -.-> hy[("Hypertable reading")] ing -.-> hy[("Hypertable mesure")]
hy -.-> agg[("Agrégat continu")] hy -.-> agg[("Agrégat continu")]
hy -.-> comp["Compression"] hy -.-> comp["Compression"]
hy -.-> ret["Rétention"] hy -.-> ret["Rétention"]
@@ -135,46 +133,67 @@ donc **pas** une hypertable : une politique de rétention émettrait des `DELETE
refuseraient. `login_attempt`, à l'inverse, est faite pour se purger, puisque son volume est refuseraient. `login_attempt`, à l'inverse, est faite pour se purger, puisque son volume est
piloté par l'attaquant. piloté par l'attaquant.
## Modèle métier
Statut : `Cible`. Les entités ci-dessous sont des **candidates**, à valider en J2. Elles
s'appuient sur les gabarits de [`apps/backend/TESTING.md`](../../apps/backend/TESTING.md), qui
évoquent déjà un modèle `Site`, un `SiteRepository` et un `ConsumptionService` exposant un
`total_kwh(site_id)`.
```mermaid
erDiagram
SITE ||--o{ POINT_DE_MESURE : porte
POINT_DE_MESURE ||--o{ MESURE : produit
SITE {
int id PK
string nom
}
POINT_DE_MESURE {
int id PK
int site_id FK
string libelle
string unite
}
MESURE {
timestamptz horodatage PK
int point_id PK
double valeur
}
```
`MESURE` est la table destinée à devenir une hypertable, partitionnée sur `horodatage`. Sa clé
primaire doit inclure la colonne de temps : TimescaleDB l'exige, une clé sur le seul identifiant
de point serait refusée.
## Gabarit de révision créant une hypertable ## Gabarit de révision créant une hypertable
Conforme à la règle de l'ADR 0001 : table et hypertable dans la même révision. La révision Conforme à la règle de l'ADR 0001 : table et hypertable dans la même révision.
`e6d2026091501` en est l'exemple réel, réduit ici à l'essentiel.
```python ```python
def upgrade() -> None: def upgrade() -> None:
op.create_table( op.create_table(
"reading", "mesure",
sa.Column("reading_id", sa.BigInteger(), autoincrement=True, nullable=False), sa.Column("horodatage", sa.DateTime(timezone=True), nullable=False),
sa.Column("site_id", sa.Text(), nullable=False), sa.Column("point_id", sa.Integer(), sa.ForeignKey("point_de_mesure.id"), nullable=False),
sa.Column("timestamp", sa.DateTime(timezone=True), nullable=False), sa.Column("valeur", sa.Float(), nullable=False),
sa.PrimaryKeyConstraint("reading_id", "timestamp"), sa.PrimaryKeyConstraint("horodatage", "point_id"),
)
op.execute(
"SELECT create_hypertable('reading', by_range('timestamp'), "
"create_default_indexes => FALSE)"
) )
op.execute("SELECT create_hypertable('mesure', by_range('horodatage'))")
def downgrade() -> None: def downgrade() -> None:
op.drop_table("reading") op.drop_table("mesure")
``` ```
La clé primaire inclut la colonne de temps parce que TimescaleDB l'exige : toute contrainte
unique d'une hypertable doit porter la colonne de partitionnement, et une clé sur le seul
`reading_id` serait refusée par `create_hypertable`.
`create_default_indexes => FALSE` écarte l'index que TimescaleDB pose d'office sur la seule
colonne de temps : les index déclarés dans la révision le couvrent déjà.
`drop_table` suffit au retour arrière : supprimer la table supprime l'hypertable et ses partitions. `drop_table` suffit au retour arrière : supprimer la table supprime l'hypertable et ses partitions.
## Conventions ## Conventions
- **Noms au singulier**, en minuscules, sans préfixe de table : `app_user`, `reading`. - **Noms au singulier**, en minuscules, sans préfixe de table.
- **Toute colonne de temps en `timestamptz`.** Jamais de `timestamp` nu : une mesure sans fuseau - **Toute colonne de temps en `timestamptz`.** Jamais de `timestamp` nu : une mesure sans fuseau
devient ininterprétable dès le premier changement d'heure. devient ininterprétable dès le premier changement d'heure.
- **La colonne de partitionnement entre dans la clé primaire.** Dans `reading` elle s'appelle - **La colonne de partitionnement s'appelle `horodatage`** et entre dans la clé primaire.
`timestamp` : c'est un nom de colonne, son type reste `timestamptz`.
- **Les politiques de rétention et de compression** vont dans `db/migrations/`, pas dans Alembic : - **Les politiques de rétention et de compression** vont dans `db/migrations/`, pas dans Alembic :
elles ne découlent pas du schéma applicatif. elles ne découlent pas du schéma applicatif.
- **Tout modèle doit être importé dans `app/models/__init__.py`**, sans quoi - **Tout modèle doit être importé dans `app/models/__init__.py`**, sans quoi
@@ -182,12 +201,13 @@ colonne de temps : les index déclarés dans la révision le couvrent déjà.
## Questions ouvertes ## Questions ouvertes
Elles relèvent du jalon J2, « valider le périmètre retenu ». Le schéma est livré : ce qui suit Elles relèvent du jalon J2, « valider le périmètre retenu », et bloquent le modèle définitif.
porte sur son exploitation, plus sur sa forme.
- **Quelles sources de mesures**, et selon quel protocole elles sont collectées.
- **Quelle granularité** à l'ingestion : la seconde, la minute, le quart d'heure. - **Quelle granularité** à l'ingestion : la seconde, la minute, le quart d'heure.
- **Quels agrégats continus**, et sur quelles fenêtres. - **Quels agrégats continus**, et sur quelles fenêtres.
- **Quelle profondeur de rétention** en données brutes, et à partir de quand on compresse. - **Quelle profondeur de rétention** en données brutes, et à partir de quand on compresse.
- **Quelles unités** sont manipulées, et si une même table les mélange.
- **Multi-tenant ou non** : un site appartient-il à un client, et faut-il cloisonner les lectures. - **Multi-tenant ou non** : un site appartient-il à un client, et faut-il cloisonner les lectures.
## Modélisation détaillée des données ## Modélisation détaillée des données
@@ -200,11 +220,12 @@ jusqu’aux recommandations proposées à l’utilisateur.
### Schéma de données ### Schéma de données
Le diagramme ci-dessous présente les tables et leurs relations. Le diagramme ci-dessous présente les tables et leurs relations.
La révision `e6d2026091501` les crée. Il décrit une structure de conception ; les migrations correspondantes
restent à implémenter.
![Schéma de données EnerVision](images/EnerVision-schema-donnees.png) ![Schéma de données EnerVision](images/EnerVision-schema-donnees.png)
*Figure : Modélisation des données EnerVision.* *Figure — Modélisation des données EnerVision.*
### Description des tables ### Description des tables
@@ -213,15 +234,15 @@ des données.
| Table | Rôle | Origine des informations | | Table | Rôle | Origine des informations |
|---|---|---| |---|---|---|
| `dataset` | Identifier les jeux historiques, retrouver leurs fichiers et conserver leurs métadonnées | Archive CSV/JSON et informations ajoutées lors de l’import | | `datasets` | Identifier les jeux historiques, retrouver leurs fichiers et conserver leurs métadonnées | Archive CSV/JSON et informations ajoutées lors de l’import |
| `site` | Regrouper les informations des sites : identifiant, nom, type et caractéristiques disponibles | CSV et API Mock `/api/v1/sites` | | `sites` | Regrouper les informations des sites : identifiant, nom, type et caractéristiques disponibles | CSV et API Mock `/api/v1/sites` |
| `reading` | Stocker les mesures, leur provenance, leur qualité et les éventuelles valeurs imputées | CSV et API Mock `/current` et `/readings` | | `readings` | Stocker les mesures, leur provenance, leur qualité et les éventuelles valeurs imputées | CSV et API Mock `/current` et `/readings` |
| `prediction` | Conserver les prévisions, leur période cible et la référence du modèle utilisé | Traitements ML d’EnerVision | | `predictions` | Conserver les prévisions, leur période cible et la référence du modèle utilisé | Traitements ML d’EnerVision |
| `alert` | Enregistrer les alertes, leur type, leur gravité et leur message | API Mock `/alerts` et détections EnerVision | | `alerts` | Enregistrer les alertes, leur type, leur gravité et leur message | API Mock `/alerts` et détections EnerVision |
| `recommendation` | Proposer des actions et expliquer la règle qui les motive | Règles métier d’EnerVision | | `recommendations` | Proposer des actions et expliquer la règle qui les motive | Règles métier d’EnerVision |
Les anomalies historiques décrites dans les JSON sont conservées Les anomalies historiques décrites dans les JSON sont conservées
dans `dataset.metadata`. Elles servent à l’analyse des données dans `datasets.metadata`. Elles servent à l’analyse des données
et ne sont pas considérées comme des alertes actuelles. et ne sont pas considérées comme des alertes actuelles.
### Relations entre les tables ### Relations entre les tables