"""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, ) 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 = interval.context or {} 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 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 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