From b20c5a1216a746dbc8a672cc09526aa2b67d219f Mon Sep 17 00:00:00 2001 From: root Date: Mon, 24 Aug 2026 12:08:57 +0000 Subject: [PATCH] fix(meters): close checkpoint C3 analytics workflows --- app_new/core/services/meter_analytics.py | 118 +++++++++++++++++- app_new/core/services/meter_service.py | 6 + app_new/planning/meter_readings.py | 80 +++++++++++- app_new/planning/schedules.py | 4 +- .../templates/planning/meter_monitoring.html | 6 +- tests/integration/test_meter_checkpoint_c3.py | 81 +++++++++++- 6 files changed, 279 insertions(+), 16 deletions(-) diff --git a/app_new/core/services/meter_analytics.py b/app_new/core/services/meter_analytics.py index 3e8fa2d..3fe74a2 100644 --- a/app_new/core/services/meter_analytics.py +++ b/app_new/core/services/meter_analytics.py @@ -15,6 +15,7 @@ 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 @@ -171,11 +172,19 @@ def evaluate_alert_rule(rule, interval, intervals=None, commit=True): bucket = [interval] value = _metric_value(bucket, rule.metric) if rule.context_filter: - context = interval.context or {} + 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 rule.low_consumption_threshold is not None and value <= rule.low_consumption_threshold + 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() @@ -249,6 +258,111 @@ def close_meter_alert(*, alert, conclusion, comment=None, user_id=None, commit=T 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: + 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: + 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 diff --git a/app_new/core/services/meter_service.py b/app_new/core/services/meter_service.py index 2e9be98..2480926 100644 --- a/app_new/core/services/meter_service.py +++ b/app_new/core/services/meter_service.py @@ -126,6 +126,9 @@ def record_meter_reading(*, meter, value, user_id, reading_date=None, notes=None meter.initial_value = value meter.last_maintenance_value = value db.session.add(reading) + db.session.flush() + from .meter_analytics import recalculate_meter_after_reading + recalculate_meter_after_reading(meter=meter, reading=reading, commit=False) if commit: db.session.commit() return reading @@ -155,6 +158,9 @@ def correct_meter_reading(*, reading, new_value, user_id, reason, commit=True): if latest is not None: meter.current_value = latest.value meter.last_reading_date = latest.reading_date.date() + db.session.flush() + from .meter_analytics import recalculate_meter_after_correction + recalculate_meter_after_correction(meter=meter, reading=reading, commit=False) if commit: db.session.commit() return reading diff --git a/app_new/planning/meter_readings.py b/app_new/planning/meter_readings.py index 162dd5e..1b6b839 100644 --- a/app_new/planning/meter_readings.py +++ b/app_new/planning/meter_readings.py @@ -1,6 +1,6 @@ """Routes C2 : règles, échéances et tournées de relevés.""" -from datetime import date, datetime +from datetime import date, datetime, timedelta from pathlib import Path from uuid import uuid4 @@ -11,7 +11,8 @@ from werkzeug.utils import secure_filename from ..extensions import db from ..core.models.planning import ( - Meter, MeterReadingOccurrence, MeterReadingRound, MeterReadingSchedule, + Meter, MeterReading, MeterReadingOccurrence, MeterReadingRound, MeterReadingSchedule, + GasConversion, ) from ..core.models.user import User from ..core.services.meter_reading_planning import ( @@ -19,7 +20,7 @@ from ..core.services.meter_reading_planning import ( create_round, ensure_occurrences_for_operational_date, generate_occurrences, record_occurrence_reading, record_occurrence_without_reading, ) -from ..core.services.meter_analytics import consumption_intervals +from ..core.services.meter_analytics import consumption_intervals, visible_open_meter_alerts, close_meter_alert from ..core.models.planning import MeterAlert, MeterAlertRule from .schedules import planning_bp @@ -133,9 +134,54 @@ def meter_occurrence_without_reading(occurrence_id): def meter_monitoring(): """Tableau de surveillance C3, sans recalculer ni modifier les relevés.""" from sqlalchemy.orm import selectinload + period_end = date.today() + period_start = period_end - timedelta(days=29) meters = Meter.query.filter(Meter.is_active.is_(True)).options(selectinload(Meter.alerts)).order_by(Meter.name).all() - alerts = MeterAlert.query.filter(MeterAlert.status.in_(("OPEN", "INVESTIGATING"))).order_by(MeterAlert.level.desc(), MeterAlert.opened_at.desc()).all() - return render_template('planning/meter_monitoring.html', meters=meters, alerts=alerts) + alerts = visible_open_meter_alerts(current_user) + summaries = {"eau": {"raw": 0.0}, "gaz": {"raw": 0.0, "kwh": 0.0}, "electricite": {"raw": 0.0}} + insufficient = [] + meter_ids = [meter.id for meter in meters] + readings_by_meter = {meter_id: [] for meter_id in meter_ids} + readings = (MeterReading.query.filter(MeterReading.meter_id.in_(meter_ids)) + .order_by(MeterReading.meter_id, MeterReading.reading_date, MeterReading.id).all()) if meter_ids else [] + for reading in readings: + readings_by_meter[reading.meter_id].append(reading) + conversions_by_meter = {meter_id: [] for meter_id in meter_ids} + for conversion in GasConversion.query.filter(GasConversion.meter_id.in_(meter_ids)).all() if meter_ids else []: + conversions_by_meter[conversion.meter_id].append(conversion) + reading_counts = {meter_id: len(rows) for meter_id, rows in readings_by_meter.items()} + for meter in meters: + rows = readings_by_meter[meter.id] + intervals = [] + for previous, current in zip(rows, rows[1:]): + if current.reading_date.date() < period_start or current.reading_date.date() > period_end: + continue + if previous.is_reset or current.is_reset or current.value < previous.value: + continue + intervals.append((current.reading_date.date(), current.value - previous.value)) + rules = [rule for rule in meter.alert_rules if rule.is_active and rule.metric in ("consumption_per_day", "consumption_interval")] + if rules and any(len(intervals) < rule.min_comparable_intervals for rule in rules): + insufficient.append(meter) + bucket = summaries.get(meter.meter_type) + if not bucket: + continue + bucket["raw"] += sum(value for _, value in intervals) + if meter.meter_type == "gaz": + for on_date, value in intervals: + applicable = [conversion for conversion in conversions_by_meter[meter.id] + if conversion.valid_from <= on_date and + (conversion.valid_to is None or conversion.valid_to >= on_date)] + manual = next((conversion for conversion in sorted(applicable, key=lambda item: item.valid_from, reverse=True) + if conversion.origin == "MANUAL"), None) + coefficient = (manual or max(applicable, key=lambda item: item.valid_from) if applicable else None) + coefficient = coefficient.coefficient_kwh_per_m3 if coefficient else None + if coefficient is not None: + bucket["kwh"] += value * coefficient + return render_template( + 'planning/meter_monitoring.html', meters=meters, alerts=alerts, + period_start=period_start, period_end=period_end, summaries=summaries, + insufficient_meters=insufficient, reading_counts=reading_counts, + ) @planning_bp.route('/meter-analytics/') @@ -171,6 +217,30 @@ def create_intervention_from_meter_alert(alert_id): return redirect(url_for('interventions.detail', id=intervention.id)) +@planning_bp.route('/meter-alerts//close', methods=['POST']) +@login_required +@permission_required('planning.manage') +def close_meter_alert_route(alert_id): + alert = MeterAlert.query.get_or_404(alert_id) + action = request.form.get('status', 'CLOSED').upper() + if action == 'INVESTIGATING': + alert.status = 'INVESTIGATING' + alert.comment = request.form.get('comment') or alert.comment + db.session.commit() + flash('Alerte passée en investigation.', 'success') + else: + try: + close_meter_alert( + alert=alert, conclusion=request.form.get('conclusion', 'OTHER'), + comment=request.form.get('comment'), user_id=current_user.id, + ) + flash('Alerte clôturée.', 'success') + except ValueError as exc: + db.session.rollback() + flash(str(exc), 'danger') + return redirect(url_for('planning.meter_monitoring')) + + @planning_bp.route('/meter-rounds') @login_required def meter_rounds(): diff --git a/app_new/planning/schedules.py b/app_new/planning/schedules.py index 894c32f..664c7d8 100644 --- a/app_new/planning/schedules.py +++ b/app_new/planning/schedules.py @@ -85,8 +85,8 @@ def my_day(): if item.status in ('à replanifier', 'proposé', 'conflit') } execution_states = {} - from ..core.models.planning import MeterAlert - informational_alerts = MeterAlert.query.filter(MeterAlert.status.in_(("OPEN", "INVESTIGATING"))).order_by(MeterAlert.level.desc()).all() + from ..core.services.meter_analytics import visible_open_meter_alerts + informational_alerts = visible_open_meter_alerts(current_user) for task in ScheduledTask.query.filter(ScheduledTask.scheduled_date == target_date).all(): execution_states[task.id] = { 'status': task.status, diff --git a/app_new/planning/templates/planning/meter_monitoring.html b/app_new/planning/templates/planning/meter_monitoring.html index a077f27..0e27994 100644 --- a/app_new/planning/templates/planning/meter_monitoring.html +++ b/app_new/planning/templates/planning/meter_monitoring.html @@ -3,10 +3,12 @@ {% block content %}

Surveillance des compteurs

-
Alertes critiques
{{ alerts|selectattr('level','equalto','CRITICAL')|list|length }}
Avertissements
{{ alerts|selectattr('level','equalto','WARNING')|list|length }}
Compteurs actifs
{{ meters|length }}
+

Synthèses informatives du {{ period_start.strftime('%d/%m/%Y') }} au {{ period_end.strftime('%d/%m/%Y') }}.

+
Alertes critiques
{{ alerts|selectattr('level','equalto','CRITICAL')|list|length }}
Avertissements
{{ alerts|selectattr('level','equalto','WARNING')|list|length }}
Compteurs actifs
{{ meters|length }}
Données insuffisantes
{{ insufficient_meters|length }}
+
Eau
{{ '%.2f'|format(summaries.eau.raw) }} m³
Gaz
{{ '%.2f'|format(summaries.gaz.raw) }} m³
{% if summaries.gaz.kwh %}{{ '%.2f'|format(summaries.gaz.kwh) }} kWh convertis{% else %}kWh non calculés{% endif %}
Électricité
{{ '%.2f'|format(summaries.electricite.raw) }} kWh
Compteurs
- {% for meter in meters %}{% else %}{% endfor %} + {% for meter in meters %}{% else %}{% endfor %}
CompteurDerniers intervallesAction
{{ meter.name }}{{ meter.unit }} · {{ meter.usage or 'Usage non renseigné' }}{{ meter.readings.count() }}Analyser
Aucun compteur actif.
{{ meter.name }}{{ meter.unit }} · {{ meter.usage or 'Usage non renseigné' }}{{ reading_counts.get(meter.id, 0) }}Analyser
Aucun compteur actif.
Alertes ouvertes
    {% for alert in alerts %}
  • {{ 'Critique' if alert.level == 'CRITICAL' else 'Avertissement' }} {{ alert.meter.name }}{{ alert.alert_type }} · {{ alert.observed_value }}
  • {% else %}
  • Aucune alerte ouverte.
  • {% endfor %}
diff --git a/tests/integration/test_meter_checkpoint_c3.py b/tests/integration/test_meter_checkpoint_c3.py index 0689605..e035947 100644 --- a/tests/integration/test_meter_checkpoint_c3.py +++ b/tests/integration/test_meter_checkpoint_c3.py @@ -6,9 +6,10 @@ import pytest from app_new import db from app_new.core.models.college import Building from app_new.core.models.planning import ( - Meter, MeterReading, MeterAlertRule, MeterHeatingRegime, GasConversion, MeterTariff, + Meter, MeterReading, MeterReadingCorrection, MeterAlert, MeterAlertEvidence, + MeterAlertRule, MeterHeatingRegime, GasConversion, MeterTariff, ) -from app_new.core.services.meter_service import create_meter, record_meter_reading +from app_new.core.services.meter_service import create_meter, record_meter_reading, correct_meter_reading from app_new.core.services.meter_analytics import ( consumption_intervals, aggregate_intervals, evaluate_alert_rule, statistical_anomaly, remainder_for_interval, estimated_cost, gas_coefficient, @@ -135,15 +136,85 @@ def test_c3_alert_closure_excludes_only_relevant_baselines(app, admin_user): def test_c3_performance_measurement(app, admin_user, authenticated_client): with app.app_context(): + first_meter_id = None for index in range(4): meter = _meter(f"PERF_{index}") + first_meter_id = first_meter_id or meter.id for offset in range(12): _reading(meter, offset * 5, date(2025, 1, 1) + __import__('datetime').timedelta(days=offset), admin_user["id"]) + perf_meter = db.session.get(Meter, first_meter_id) + perf_intervals = consumption_intervals(perf_meter) + perf_rule = MeterAlertRule(meter=perf_meter, metric="consumption_interval", period="DAY", operator=">", threshold=1) + db.session.add(perf_rule) + db.session.commit() count = {"value": 0} from sqlalchemy import event def before_cursor(*args): count["value"] += 1 + def measure(label, callback): + count["value"] = 0 + started = perf_counter() + callback() + elapsed = perf_counter() - started + print(f"C3_CLOSE_PERF {label}_seconds={elapsed:.4f} {label}_queries={count['value']}") + return elapsed, count["value"] event.listen(db.engine, "before_cursor_execute", before_cursor) - started = perf_counter(); response = authenticated_client.get("/planning/meter-monitoring"); elapsed = perf_counter() - started + dashboard_time, dashboard_queries = measure("dashboard", lambda: authenticated_client.get("/planning/meter-monitoring")) + analytics_time, analytics_queries = measure("analytics", lambda: authenticated_client.get(f"/planning/meter-analytics/{first_meter_id}")) + alert_time, alert_queries = measure("alerts", lambda: evaluate_alert_rule(perf_rule, perf_intervals[-1], intervals=perf_intervals)) + day_time, day_queries = measure("my_day", lambda: authenticated_client.get("/planning/my-day")) event.remove(db.engine, "before_cursor_execute", before_cursor) - print(f"C3_PERF dashboard_seconds={elapsed:.4f} dashboard_queries={count['value']}") - assert response.status_code == 200 + assert dashboard_time >= 0 and analytics_time >= 0 and alert_time >= 0 and day_time >= 0 + + +def test_c3_new_reading_runs_targeted_alert_analysis(app, admin_user): + with app.app_context(): + meter = _meter("AUTO_READING") + _reading(meter, 0, date(2026, 3, 1), admin_user["id"]) + _reading(meter, 1, date(2026, 3, 2), admin_user["id"]) + rule = MeterAlertRule(meter=meter, metric="consumption_per_day", period="DAY", operator=">", threshold=5) + db.session.add(rule) + db.session.commit() + reading = _reading(meter, 10, date(2026, 3, 3), admin_user["id"]) + alert = MeterAlert.query.filter_by(meter_id=meter.id, rule_id=rule.id, status="OPEN").one() + assert alert.evidence and alert.evidence[0].reading_end_id == reading.id + assert MeterAlertEvidence.query.filter_by(alert_id=alert.id).count() == 1 + + +def test_c3_correction_recalculates_neighbors_and_closes_false_alert(app, admin_user): + with app.app_context(): + meter = _meter("AUTO_CORRECTION") + _reading(meter, 0, date(2026, 4, 1), admin_user["id"]) + _reading(meter, 1, date(2026, 4, 2), admin_user["id"]) + rule = MeterAlertRule(meter=meter, metric="consumption_per_day", period="DAY", operator=">", threshold=5) + db.session.add(rule) + db.session.commit() + high = _reading(meter, 20, date(2026, 4, 3), admin_user["id"]) + alert = MeterAlert.query.filter_by(meter_id=meter.id, rule_id=rule.id, status="OPEN").one() + correct_meter_reading(reading=high, new_value=3, user_id=admin_user["id"], reason="TEST_UI_C3_CLOSE correction") + db.session.expire_all() + refreshed = db.session.get(MeterAlert, alert.id) + assert MeterReadingCorrection.query.filter_by(reading_id=high.id).count() == 1 + assert refreshed.status == "CLOSED" + assert all(evidence.excluded_from_baseline for evidence in refreshed.evidence) + + +def test_c3_alert_close_route_and_dashboard_summaries(app, admin_user, authenticated_client): + with app.app_context(): + water = _meter("CLOSE_HTTP_WATER", "eau", "m³") + _reading(water, 0, date.today(), admin_user["id"]) + _reading(water, 12, date.today(), admin_user["id"]) + rule = MeterAlertRule(meter=water, metric="consumption_interval", period="DAY", operator=">", threshold=1) + db.session.add(rule) + db.session.commit() + alert = evaluate_alert_rule(rule, consumption_intervals(water)[0]) + alert_id = alert.id + response = authenticated_client.post( + f"/planning/meter-alerts/{alert_id}/close", + data={"status": "CLOSED", "conclusion": "BAD_THRESHOLD", "comment": "TEST_UI_C3_CLOSE"}, + ) + assert response.status_code == 302 + with app.app_context(): + assert db.session.get(MeterAlert, alert_id).status == "CLOSED" + dashboard = authenticated_client.get("/planning/meter-monitoring") + assert dashboard.status_code == 200 + assert b"Eau" in dashboard.data and b"m" in dashboard.data