gmao/app_new/core/services/meter_analytics.py

392 lines
19 KiB
Python

"""Calculs C3 sur les relevés bruts.
Les intervalles sont volontairement calculés à la demande : MeterReading reste
la source de vérité et un reset ou un changement de compteur ne peut pas être
absorbé silencieusement dans une soustraction.
"""
from dataclasses import dataclass, asdict
from datetime import date, datetime, timedelta, timezone
from statistics import median, mean
from typing import Optional
from sqlalchemy.orm import joinedload
from ...extensions import db
from ..models.planning import (
Meter, MeterReading, MeterAlert, MeterAlertEvidence, MeterAlertRule,
MeterHeatingRegime, GasConversion, MeterTariff, CollegeClosure, ClosureWorkDay,
MeterReadingSchedule, MeterReadingOccurrence,
)
from .planning_service import is_public_holiday
@dataclass
class ConsumptionInterval:
meter_id: int
reading_start: MeterReading
reading_end: MeterReading
start_date: datetime
end_date: datetime
raw_delta: float
calendar_days: int
consumption_per_day: float
quality: str = "EXACT"
context: dict = None
def as_dict(self):
data = asdict(self)
data["reading_start"] = self.reading_start.id
data["reading_end"] = self.reading_end.id
data["start_date"] = self.start_date.isoformat()
data["end_date"] = self.end_date.isoformat()
return data
def _reading_datetime(reading):
value = reading.reading_date or datetime.min.replace(tzinfo=timezone.utc)
return value if value.tzinfo else value.replace(tzinfo=timezone.utc)
def _context_for(meter, start, end):
counts = {"scolaire": 0, "vacances": 0, "fermeture": 0, "permanence": 0, "autre": 0}
d = start.date()
closures = CollegeClosure.query.filter(CollegeClosure.start_date <= end.date(), CollegeClosure.end_date >= start.date()).all()
closure_ids = [row.id for row in closures]
work_days = {row.work_date for row in ClosureWorkDay.query.filter(ClosureWorkDay.work_date >= start.date(), ClosureWorkDay.work_date < end.date()).all()} if closure_ids else set()
while d < end.date():
if meter.housing_unit_id:
counts["scolaire"] += 1
else:
closure = next((row for row in closures if row.start_date <= d <= row.end_date), None)
if closure:
kind = (closure.closure_type or "vacances").lower()
key = "permanence" if d in work_days or "perman" in kind else ("fermeture" if "fermet" in kind else "vacances")
counts[key] += 1
elif d.weekday() >= 5 or is_public_holiday(d)[0]:
counts["fermeture"] += 1
else:
counts["scolaire"] += 1
d += timedelta(days=1)
regimes = {"NORMAL": 0, "REDUCED": 0, "STOP": 0}
rows = MeterHeatingRegime.query.filter(
MeterHeatingRegime.meter_id == meter.id,
MeterHeatingRegime.valid_from <= end.date(),
db.or_(MeterHeatingRegime.valid_to.is_(None), MeterHeatingRegime.valid_to >= start.date()),
).all()
d = start.date()
while d < end.date():
row = next((r for r in rows if r.valid_from <= d and (r.valid_to is None or d <= r.valid_to)), None)
if row:
regimes[row.regime] += 1
d += timedelta(days=1)
return {"calendar": counts, "heating": regimes}
def consumption_intervals(meter, start=None, end=None):
query = MeterReading.query.filter_by(meter_id=meter.id).order_by(MeterReading.reading_date, MeterReading.id).options(joinedload(MeterReading.read_by))
if start:
query = query.filter(MeterReading.reading_date >= start)
if end:
query = query.filter(MeterReading.reading_date <= end)
readings = query.all()
result = []
for previous, current in zip(readings, readings[1:]):
# A reset belongs to a new physical segment; never derive a delta over it.
if previous.is_reset or current.is_reset or current.value < previous.value:
continue
days = max((_reading_datetime(current) - _reading_datetime(previous)).days, 1)
interval = ConsumptionInterval(
meter_id=meter.id, reading_start=previous, reading_end=current,
start_date=_reading_datetime(previous), end_date=_reading_datetime(current),
raw_delta=current.value - previous.value, calendar_days=days,
consumption_per_day=(current.value - previous.value) / days,
)
interval.context = _context_for(meter, interval.start_date, interval.end_date)
result.append(interval)
return result
def aggregate_intervals(intervals, period="DAY"):
"""Agrégation simple par jour/semaine/mois sans modifier les relevés."""
groups = {}
for interval in intervals:
d = interval.end_date.date()
key = d if period == "DAY" else (d - timedelta(days=d.weekday()) if period == "WEEK" else (d.year, d.month))
groups.setdefault(key, []).append(interval)
return [{"period": key, "consumption": sum(i.raw_delta for i in rows), "calendar_days": sum(i.calendar_days for i in rows),
"consumption_per_day": sum(i.raw_delta for i in rows) / max(sum(i.calendar_days for i in rows), 1), "intervals": rows}
for key, rows in sorted(groups.items(), key=lambda x: x[0])]
def _metric_value(intervals, metric):
"""Retourne une métrique contextualisée par répartition temporelle."""
if not intervals:
return 0.0
if metric == "consumption_interval":
return sum(item.raw_delta for item in intervals)
if metric == "consumption_per_day":
return sum(item.raw_delta for item in intervals) / max(sum(item.calendar_days for item in intervals), 1)
calendar_key = {"school_consumption": "scolaire", "vacation_consumption": "vacances", "closure_consumption": "fermeture"}
heating_key = {"heating_normal": "NORMAL", "heating_reduced": "REDUCED", "heating_stop": "STOP"}
key = calendar_key.get(metric)
if key:
days = sum((item.context or {}).get("calendar", {}).get(key, 0) for item in intervals)
else:
days = sum((item.context or {}).get("heating", {}).get(heating_key.get(metric, ""), 0) for item in intervals)
total_days = sum(item.calendar_days for item in intervals)
return sum(item.raw_delta for item in intervals) * days / total_days if total_days else 0.0
def remainder_for_interval(meter, interval):
children = meter.children.all() if hasattr(meter.children, "all") else list(meter.children)
if not children:
return None
total = 0.0
quality = "EXACT"
for child in children:
rows = [i for i in consumption_intervals(child) if i.start_date <= interval.start_date and i.end_date >= interval.end_date]
if rows:
total += rows[-1].consumption_per_day * interval.calendar_days
continue
history = consumption_intervals(child)
if len(history) >= 2:
total += mean(i.consumption_per_day for i in history[-3:]) * interval.calendar_days
quality = "ESTIMATED"
else:
quality = "PARTIAL"
return {"value": interval.raw_delta - total, "quality": quality}
def _matches(value, operator, threshold):
return {">": value > threshold, "<": value < threshold, ">=": value >= threshold, "<=": value <= threshold}[operator]
def evaluate_alert_rule(rule, interval, intervals=None, commit=True):
if not rule.is_active:
return None
intervals = intervals or [interval]
if rule.period in ("WEEK", "MONTH"):
grouped = aggregate_intervals(intervals, rule.period)
bucket = next((row["intervals"] for row in grouped if interval in row["intervals"]), [interval])
else:
bucket = [interval]
value = _metric_value(bucket, rule.metric)
if rule.context_filter:
context = {
"calendar": {key: sum((item.context or {}).get("calendar", {}).get(key, 0) for item in bucket)
for key in ("scolaire", "vacances", "fermeture", "permanence", "autre")},
"heating": {key: sum((item.context or {}).get("heating", {}).get(key, 0) for item in bucket)
for key in ("NORMAL", "REDUCED", "STOP")},
}
if rule.context_filter in ("NORMAL", "REDUCED", "STOP"):
if not (context.get("heating", {}).get(rule.context_filter) or 0): return None
elif not (context.get("calendar", {}).get(rule.context_filter) or 0): return None
zero_low = rule.detect_zero_or_low and (
value == 0 or
(rule.low_consumption_threshold is not None and value <= rule.low_consumption_threshold)
)
if not zero_low and not _matches(value, rule.operator, rule.threshold): return None
alert_type = "ZERO_OR_LOW_CONSUMPTION" if zero_low else "MANUAL_THRESHOLD"
open_alert = MeterAlert.query.filter_by(meter_id=rule.meter_id, alert_type=alert_type, metric=rule.metric, period=rule.period, context_key=rule.context_filter, rule_id=rule.id, status="OPEN").first()
if not open_alert:
open_alert = MeterAlert(meter_id=rule.meter_id, alert_type=alert_type, level=rule.level, metric=rule.metric, period=rule.period, context_key=rule.context_filter, rule_id=rule.id, observed_value=value, period_start=bucket[0].start_date.date(), period_end=bucket[-1].end_date.date())
db.session.add(open_alert); db.session.flush()
else:
open_alert.observed_value = value
evidence = MeterAlertEvidence.query.filter_by(alert_id=open_alert.id, reading_start_id=bucket[0].reading_start.id, reading_end_id=bucket[-1].reading_end.id).first()
if not evidence:
db.session.add(MeterAlertEvidence(alert_id=open_alert.id, reading_start_id=bucket[0].reading_start.id, reading_end_id=bucket[-1].reading_end.id, observed_value=value))
if commit: db.session.commit()
return open_alert
def statistical_anomaly(meter, interval, min_intervals=3, minimum_deviation_percent=30.0, rule=None, commit=True):
if rule is not None:
min_intervals = rule.min_comparable_intervals
minimum_deviation_percent = rule.min_deviation_percent
history = [i for i in consumption_intervals(meter) if i.end_date <= interval.start_date and i.context == interval.context]
history = [i for i in history if not MeterAlertEvidence.query.filter_by(reading_end_id=i.reading_end.id, excluded_from_baseline=True).first()]
if len(history) < min_intervals: return None
reference = median(i.consumption_per_day for i in history)
if reference == 0 or interval.consumption_per_day <= reference: return None
deviation = (interval.consumption_per_day - reference) / reference * 100
if deviation < minimum_deviation_percent:
return None
rule_id = rule.id if rule else None
alert = MeterAlert.query.filter_by(meter_id=meter.id, alert_type="STATISTICAL_ANOMALY", metric="consumption_per_day", period="DAY", context_key="comparable", rule_id=rule_id, status="OPEN").first()
if not alert:
alert = MeterAlert(meter_id=meter.id, alert_type="STATISTICAL_ANOMALY", level="WARNING", metric="consumption_per_day", period="DAY", context_key="comparable", rule_id=rule_id)
db.session.add(alert); db.session.flush()
alert.observed_value, alert.reference_value, alert.deviation_percent, alert.comparable_count = interval.consumption_per_day, reference, deviation, len(history)
if not MeterAlertEvidence.query.filter_by(alert_id=alert.id, reading_start_id=interval.reading_start.id, reading_end_id=interval.reading_end.id).first():
db.session.add(MeterAlertEvidence(alert_id=alert.id, reading_start_id=interval.reading_start.id, reading_end_id=interval.reading_end.id, observed_value=interval.consumption_per_day))
if commit: db.session.commit()
return alert
def gas_coefficient(meter, on_date):
rows = GasConversion.query.filter_by(meter_id=meter.id).filter(GasConversion.valid_from <= on_date, db.or_(GasConversion.valid_to.is_(None), GasConversion.valid_to >= on_date)).order_by(GasConversion.valid_from.desc()).all()
manual = next((row for row in rows if row.origin == "MANUAL"), None)
return (manual or (rows[0] if rows else None)).coefficient_kwh_per_m3 if (manual or rows) else None
def create_bill_gas_conversion(*, meter, volume_m3, billed_kwh, valid_from, valid_to=None, commit=True):
if volume_m3 is None or volume_m3 <= 0:
raise ValueError("Le volume gaz facturé doit être strictement positif.")
if billed_kwh is None or billed_kwh < 0:
raise ValueError("Les kWh facturés ne peuvent pas être négatifs.")
conversion = GasConversion(meter=meter, coefficient_kwh_per_m3=billed_kwh / volume_m3,
valid_from=valid_from, valid_to=valid_to, origin="BILL",
volume_m3=volume_m3, billed_kwh=billed_kwh)
db.session.add(conversion)
if commit: db.session.commit()
return conversion
def close_meter_alert(*, alert, conclusion, comment=None, user_id=None, commit=True):
allowed = {"CONFIRMED_LEAK", "NORMAL_EXPLAINED", "READING_ERROR", "BAD_THRESHOLD", "OTHER"}
if conclusion not in allowed:
raise ValueError("Conclusion d'alerte invalide.")
alert.status = "CLOSED"
alert.closed_at = datetime.now(timezone.utc)
alert.conclusion = conclusion
alert.comment = comment
if conclusion in {"CONFIRMED_LEAK", "READING_ERROR"}:
for evidence in alert.evidence:
evidence.excluded_from_baseline = True
if commit: db.session.commit()
return alert
def investigate_meter_alert(*, alert, comment=None, user_id=None, commit=True):
"""Passe une alerte ouverte en investigation via le service métier."""
if alert.status not in {"OPEN", "INVESTIGATING"}:
raise ValueError("Une alerte clôturée ne peut plus passer en investigation.")
alert.status = "INVESTIGATING"
if comment:
alert.comment = comment
if commit:
db.session.commit()
return alert
def visible_open_meter_alerts(user):
"""Retourne les alertes visibles dans Ma journée selon le RBAC métier."""
query = MeterAlert.query.filter(MeterAlert.status.in_(("OPEN", "INVESTIGATING")))
if user.is_admin():
return query.order_by(MeterAlert.level.desc(), MeterAlert.opened_at.desc()).all()
assigned_schedule = db.session.query(MeterReadingSchedule.meter_id).filter(
MeterReadingSchedule.assigned_to_id == user.id
)
assigned_occurrence = db.session.query(MeterReadingOccurrence.schedule_id).filter(
MeterReadingOccurrence.assigned_to_id == user.id
)
occurrence_meter = db.session.query(MeterReadingSchedule.meter_id).filter(
MeterReadingSchedule.id.in_(assigned_occurrence)
)
return query.filter(MeterAlert.meter_id.in_(assigned_schedule.union(occurrence_meter))).order_by(
MeterAlert.level.desc(), MeterAlert.opened_at.desc()
).all()
def _intervals_touching_reading(meter, reading):
"""Retourne uniquement les segments voisins d'un relevé modifié."""
return [interval for interval in consumption_intervals(meter)
if interval.reading_start.id == reading.id or interval.reading_end.id == reading.id]
def _invalidate_open_evidence_for(intervals, meter, changed_reading_ids=None):
"""Neutralise les preuves ouvertes qui ne correspondent plus à un segment valide.
Les alertes déjà clôturées restent historiques. Les preuves ouvertes
obsolètes sont retirées afin qu'une réévaluation puisse recréer une preuve
valide sans collision ; la correction C1 conserve l'audit de la mesure.
"""
changed_ids = set(changed_reading_ids or ())
changed_ids.update({reading.id for item in intervals for reading in (item.reading_start, item.reading_end)})
if not changed_ids:
return []
affected = []
alerts = MeterAlert.query.filter_by(meter_id=meter.id).filter(MeterAlert.status == "OPEN").all()
for alert in alerts:
invalidated = False
removed = []
for evidence in alert.evidence:
if evidence.reading_start_id in changed_ids or evidence.reading_end_id in changed_ids:
# Le couple de relevés peut rester valide tout en ayant une
# valeur différente : sa preuve précédente est alors obsolète.
# On la retire uniquement d'une alerte encore ouverte ; les
# alertes clôturées conservent leur audit historique.
db.session.delete(evidence)
removed.append(evidence)
invalidated = True
db.session.flush()
remaining = [evidence for evidence in alert.evidence if evidence not in removed]
if invalidated and not any(not evidence.excluded_from_baseline for evidence in remaining):
alert.status = "CLOSED"
alert.closed_at = datetime.now(timezone.utc)
alert.conclusion = "READING_ERROR"
alert.comment = "Alerte clôturée automatiquement après correction d'un relevé."
affected.append(alert)
return affected
def recalculate_meter_after_reading(*, meter, reading, commit=True):
"""Analyse le ou les nouveaux segments, sans rescanner l'historique des règles.
Le relevé est déjà ajouté à la transaction appelante et doit avoir été
flushé pour disposer de son identifiant.
"""
impacted = _intervals_touching_reading(meter, reading)
if not impacted:
if commit:
db.session.commit()
return []
rules = MeterAlertRule.query.filter_by(meter_id=meter.id, is_active=True).all()
all_intervals = consumption_intervals(meter)
for rule in rules:
for interval in impacted:
if rule.rule_type == "STATISTICAL_ANOMALY":
statistical_anomaly(meter, interval, rule=rule, commit=False)
else:
evaluate_alert_rule(rule, interval, intervals=all_intervals, commit=False)
if commit:
db.session.commit()
return impacted
def recalculate_meter_after_correction(*, meter, reading, commit=True):
"""Réévalue seulement les deux segments pouvant dépendre du relevé corrigé."""
ordered = (MeterReading.query.filter_by(meter_id=meter.id)
.order_by(MeterReading.reading_date, MeterReading.id).all())
reading_index = next((index for index, item in enumerate(ordered) if item.id == reading.id), None)
neighbor_ids = {reading.id}
if reading_index is not None:
if reading_index:
neighbor_ids.add(ordered[reading_index - 1].id)
if reading_index + 1 < len(ordered):
neighbor_ids.add(ordered[reading_index + 1].id)
impacted = _intervals_touching_reading(meter, reading)
_invalidate_open_evidence_for(impacted, meter, neighbor_ids)
rules = MeterAlertRule.query.filter_by(meter_id=meter.id, is_active=True).all()
all_intervals = consumption_intervals(meter)
for rule in rules:
for interval in impacted:
if rule.rule_type == "STATISTICAL_ANOMALY":
statistical_anomaly(meter, interval, rule=rule, commit=False)
else:
evaluate_alert_rule(rule, interval, intervals=all_intervals, commit=False)
if commit:
db.session.commit()
return impacted
def estimated_cost(meter, interval):
tariff = MeterTariff.query.filter_by(meter_id=meter.id).filter(MeterTariff.valid_from <= interval.end_date.date(), db.or_(MeterTariff.valid_to.is_(None), MeterTariff.valid_to >= interval.end_date.date())).order_by(MeterTariff.valid_from.desc()).first()
if not tariff: return None
amount = interval.raw_delta
if meter.meter_type == "gaz":
coefficient = gas_coefficient(meter, interval.end_date.date())
if coefficient is None and tariff.unit == "€/kWh": return None
amount *= coefficient or 1
return amount * tariff.unit_price