From 17057ed41e2d2b77d6d99107d38d9e7cbefc5abe Mon Sep 17 00:00:00 2001 From: Tianyu Liu Date: Thu, 25 Jun 2026 15:38:21 +0200 Subject: [PATCH] M7-T03: make billing engine meter-aware (no cross-meter delta, delta sanity guard, period.meter_id) --- app/services/energy_cost.py | 232 +++++++++++++-- tests/test_energy_cost.py | 565 ++++++++++++++++++++++++++++++++++-- 2 files changed, 752 insertions(+), 45 deletions(-) diff --git a/app/services/energy_cost.py b/app/services/energy_cost.py index 4b2c535..dc06973 100644 --- a/app/services/energy_cost.py +++ b/app/services/energy_cost.py @@ -1,7 +1,7 @@ """Billing engine for DSMR 15-minute energy metering periods. This module implements the two-layer billing model described in §3.4 of the -M6 design document: +M6 design document, extended in M7-T03 to be meter-aware: **Layer 1 — per-period metering cost (immutable, price-snapshot)** ``compute_period(session, t0)`` computes the import cost, export revenue, and @@ -31,17 +31,49 @@ Design notes - **Register keys**: DSMR payload uses JSON strings like ``"20915.154"`` for cumulative kWh registers. ``register_at`` converts them to Decimal. - **Degraded vs skip semantics**: + - *No meter coverage* (``meter_at`` returns None for t0): write a + ``degraded=True`` row with ``meter_id=None``. + - *Cross-meter boundary* (m0.id != m1.id for t0/t1): write a ``degraded=True`` + row with ``meter_id=m0.id``; losing this one period at the swap boundary is + acceptable (D5 decision). - *Missing readings* (``register_at`` returns None for start or end - boundary): write a ``degraded=True`` row with costs at 0 so the period is - tracked and can be retried by ``compute_closed_periods``. + boundary within the meter window): write a ``degraded=True`` row with + ``meter_id=m0.id`` so the period is tracked and can be retried by + ``compute_closed_periods``. + - *Negative or excessively large delta* (delta sanity guard D6): write a + ``degraded=True`` row with ``meter_id=m0.id``; prevents negative costs and + grossly inflated costs from meter resets, DSMR rollover, or data spikes. - *Missing Tibber price* (``TibberPriceNotFoundError``): skip entirely (do not write a row); the period will be retried once prices arrive. - *Missing active contract version*: skip (no contract to compute against). +- **Meter-aware register lookup**: ``register_at`` now accepts a ``meter`` + parameter and restricts the DSMR reading query to readings within + ``[meter.started_at, meter.ended_at)`` (half-open), preventing old-meter + readings from leaking into a new-meter epoch. - **Lookback window in ``compute_closed_periods``**: to avoid scanning all historical DSMR data on every tick, the function looks back at most 7 days from the current time. This covers typical short outages (no data / no contract) while staying bounded. Periods older than 7 days must be recovered via an explicit ``recompute_range`` call. + +Meter-aware compute_period ordering rationale (M7-T03) +------------------------------------------------------- +The order of checks inside ``compute_period`` is: + +1. **Immutability guard** (existing non-degraded row, overwrite=False) → return False. +2. **Meter determination** (m0 = meter_at(t0), m1 = meter_at(t1)): + - No meter (m0 is None) → write degraded, meter_id=None. + - Cross-meter boundary (m0.id != m1.id) → write degraded, meter_id=m0.id. +3. **Active contract version check** → skip (no write) if absent. +4. **Boundary register readings** within m0's window → write degraded if missing. +5. **Delta sanity guard** → write degraded if any delta < 0 or > _MAX_DELTA_KWH. +6. **Price strategy** → skip (no write) if Tibber price missing. +7. **Upsert billing record** with meter_id=m0.id. + +Why meter before contract? The meter is a *structural* prerequisite: without a +known meter epoch we cannot trust the delta at all, so we commit a degraded row +immediately. The contract skip, by contrast, is transient (the period can be +re-computed once a contract is configured), so it produces no row. """ from __future__ import annotations @@ -59,8 +91,9 @@ from app.integrations.pricing.strategies import ( TibberPriceNotFoundError, get_strategy, ) -from app.models.energy import DsmrReading, EnergyCostPeriod +from app.models.energy import DsmrReading, EnergyCostPeriod, Meter from app.services.contracts import active_contract_version_at, active_contract_versions +from app.services.meters import meter_at from app.services.timezone import local_date, local_now logger = logging.getLogger(__name__) @@ -81,6 +114,15 @@ _LOOKBACK_DAYS = 7 # maximum lookback window for compute_closed_periods # producing a spurious zero-delta "successful" row. _READING_MAX_STALENESS = timedelta(minutes=_PERIOD_MINUTES) +# Maximum plausible kWh delta for a single 15-minute period (D6 sanity guard). +# A typical Dutch household uses well under 5 kWh per quarter hour even under +# heavy load. 100 kWh per 15 minutes corresponds to ~400 kW — far beyond any +# residential consumption — but is lenient enough to never fire on legitimate +# data. Any delta at or above this threshold indicates a meter reset, DSMR +# rollover, sign error, or other data anomaly, and the period is marked +# degraded to prevent negative costs or grossly inflated charges. +_MAX_DELTA_KWH = Decimal("100") + # DSMR payload register keys (cumulative kWh, JSON string values). _KEY_D1 = "electricity_delivered_1" # delivered low-tariff (dal / _1) _KEY_D2 = "electricity_delivered_2" # delivered high-tariff (normal / _2) @@ -123,15 +165,29 @@ def _existing_period(session: Session, t0: datetime) -> EnergyCostPeriod | None: # --------------------------------------------------------------------------- -# register_at — boundary reading lookup +# register_at — boundary reading lookup (meter-aware) # --------------------------------------------------------------------------- -def register_at(session: Session, boundary: datetime) -> dict[str, Decimal] | None: - """Return the four cumulative kWh register values at *boundary*. +def register_at( + session: Session, + boundary: datetime, + meter: Meter, +) -> dict[str, Decimal] | None: + """Return the four cumulative kWh register values at *boundary*, within *meter*'s window. - Queries the most recent ``DsmrReading`` with ``recorded_at ≤ boundary`` - and extracts the four energy registers from ``payload``: + Queries the most recent ``DsmrReading`` with: + ``recorded_at ≤ boundary`` + AND ``recorded_at ≥ meter.started_at`` + AND (``meter.ended_at IS NULL`` OR ``recorded_at < meter.ended_at``) + + The meter window constraint (half-open ``[started_at, ended_at)``) ensures + that readings from a previous meter epoch are never used to anchor a new + meter's computation. Without this guard, the final reading of the old meter + would be visible at the start of the new meter's epoch and produce a + cross-meter delta, defeating the isolation guarantee. + + Extracts the four energy registers from ``payload``: d1 — electricity_delivered_1 (delivered low-tariff / dal) d2 — electricity_delivered_2 (delivered high-tariff / normal) @@ -142,18 +198,40 @@ def register_at(session: Session, boundary: datetime) -> dict[str, Decimal] | No ------- dict[str, Decimal] with keys ``d1``, ``d2``, ``r1``, ``r2``, or ``None`` when: - - No ``DsmrReading`` row exists with ``recorded_at ≤ boundary``. + - No ``DsmrReading`` row exists with ``recorded_at ≤ boundary`` within + *meter*'s epoch window. + - The most recent such reading is older than ``_READING_MAX_STALENESS`` + relative to *boundary* (freshness guard). - Any of the four register keys is absent from the payload. - Any of the four register values is ``None`` (null in JSON). + + SQLite naive datetime note + -------------------------- + ``recorded_at`` is stored as a naive UTC datetime in SQLite. Comparisons + against *boundary* (always tz-aware UTC) use ``_as_utc()`` for the + freshness check. The SQL ``WHERE`` clause comparisons work correctly + because SQLAlchemy's SQLite dialect strips tzinfo when binding parameters + (leaving the wall-clock UTC value unchanged), consistent with the storage + format. """ - row: DsmrReading | None = ( - session.execute( - select(DsmrReading) - .where(DsmrReading.recorded_at <= boundary) - .order_by(DsmrReading.recorded_at.desc()) - .limit(1) - ).scalar_one_or_none() + # Build the meter-window constraints: [started_at, ended_at). + meter_lower = meter.started_at # DsmrReading.recorded_at >= meter.started_at + meter_upper = meter.ended_at # DsmrReading.recorded_at < meter.ended_at (if set) + + stmt = ( + select(DsmrReading) + .where( + DsmrReading.recorded_at <= boundary, + DsmrReading.recorded_at >= meter_lower, + ) + .order_by(DsmrReading.recorded_at.desc()) + .limit(1) ) + # Apply the upper bound only when the meter is closed (ended_at is not None). + if meter_upper is not None: + stmt = stmt.where(DsmrReading.recorded_at < meter_upper) + + row: DsmrReading | None = session.execute(stmt).scalar_one_or_none() if row is None: return None @@ -214,8 +292,14 @@ def compute_period(session: Session, t0: datetime, *, overwrite: bool = False) - Side-effects ------------ - Inserts or updates an ``EnergyCostPeriod`` row keyed on ``period_start=t0``. - - If readings are missing at either boundary: inserts/updates a degraded - row (costs=0, degraded=True). + - If no meter covers t0 (``meter_at`` returns None for t0): inserts/updates + a degraded row with ``meter_id=None``. + - If the period spans a meter boundary (``meter_at(t0).id != meter_at(t1).id``): + inserts/updates a degraded row with ``meter_id=m0.id`` (D5 decision). + - If readings are missing at either boundary within the meter window: + inserts/updates a degraded row with ``meter_id=m0.id``. + - If any delta is negative or exceeds ``_MAX_DELTA_KWH`` (D6 sanity guard): + inserts/updates a degraded row with ``meter_id=m0.id``. - If the active contract version is missing: **skips** (returns False, no write). - If the Tibber price is missing (TibberPriceNotFoundError): **skips** (returns False, no write). @@ -229,7 +313,50 @@ def compute_period(session: Session, t0: datetime, *, overwrite: bool = False) - if existing is not None and not existing.degraded and not overwrite: return False - # --- Active contract version at t0 (checked before readings) --- + # --- Meter determination (structural prerequisite, checked before contract) --- + # + # A missing or cross-boundary meter is a structural problem: we cannot trust + # the delta at all, so we write a degraded row immediately. This is different + # from the contract skip (transient, no write): the degraded row ensures the + # period appears in the history and can be revisited once the meter timeline + # is corrected and a recompute_range is triggered. + # + # Ordering rationale: + # 1. No meter (m0 is None) → degraded(meter_id=None): no epoch for t0. + # 2. Cross-meter boundary (m0.id != m1.id) → degraded(meter_id=m0.id): D5. + # 3. (Single meter, proceed) → contract check → readings → delta guard → price. + # + # We place meter before contract so that "cross-table period" is always + # marked degraded regardless of contract state. If we checked contract + # first, a missing-contract skip would silently discard the cross-table + # evidence; once a contract is added and recompute runs, the engine would + # incorrectly use cross-table reads. + m0 = meter_at(session, t0) + m1 = meter_at(session, t1) + + if m0 is None: + # No meter epoch covers t0 — degraded with no meter attribution. + logger.debug( + "compute_period(%s): no active meter at t0 — writing degraded (meter_id=None).", + t0.isoformat(), + ) + _upsert_degraded(session, t0, now, existing, meter_id=None) + return True + + if m1 is None or m0.id != m1.id: + # Period spans a meter boundary or t1 has no meter. Degrade with m0's id + # (t0's meter attribution): the period's start belongs to m0's epoch. + logger.debug( + "compute_period(%s): period crosses meter boundary " + "(m0.id=%s, m1.id=%s) — writing degraded.", + t0.isoformat(), + m0.id, + m1.id if m1 is not None else None, + ) + _upsert_degraded(session, t0, now, existing, meter_id=m0.id) + return True + + # --- Active contract version at t0 --- # If there is no active contract covering t0, skip the period entirely. # We do not write a degraded row — there is no meaningful state to recover # without a contract (we would not know which strategy to apply once data @@ -240,14 +367,13 @@ def compute_period(session: Session, t0: datetime, *, overwrite: bool = False) - logger.debug("compute_period(%s): no active contract version — skipping.", t0.isoformat()) return False - # --- Boundary readings --- - start_regs = register_at(session, t0) - end_regs = register_at(session, t1) + # --- Boundary readings within m0's meter window --- + start_regs = register_at(session, t0, m0) + end_regs = register_at(session, t1, m0) if start_regs is None or end_regs is None: - # Missing readings → write/update a degraded placeholder so the period - # is visible and can be retried by compute_closed_periods once data arrives. - _upsert_degraded(session, t0, now, existing) + # Missing readings within the meter window → degraded with m0 attribution. + _upsert_degraded(session, t0, now, existing, meter_id=m0.id) return True # a record was written (degraded) # --- Compute deltas (end − start) --- @@ -258,6 +384,26 @@ def compute_period(session: Session, t0: datetime, *, overwrite: bool = False) - r2=end_regs["r2"] - start_regs["r2"], ) + # --- Delta sanity guard (D6) --- + # Any negative delta indicates a meter reset, DSMR rollover, or data error. + # Any delta exceeding _MAX_DELTA_KWH (100 kWh per 15 min = 400 kW average) + # is implausible for residential use and indicates an anomaly. + # Both cases produce a degraded row so no negative or grossly inflated cost + # is ever written to the billing record. + all_deltas = (deltas.d1, deltas.d2, deltas.r1, deltas.r2) + if any(d < Decimal("0") for d in all_deltas) or any(d > _MAX_DELTA_KWH for d in all_deltas): + logger.debug( + "compute_period(%s): delta sanity guard triggered " + "(d1=%s, d2=%s, r1=%s, r2=%s) — writing degraded.", + t0.isoformat(), + deltas.d1, + deltas.d2, + deltas.r1, + deltas.r2, + ) + _upsert_degraded(session, t0, now, existing, meter_id=m0.id) + return True + # --- Price strategy --- strategy = get_strategy(version.contract.kind) try: @@ -288,6 +434,7 @@ def compute_period(session: Session, t0: datetime, *, overwrite: bool = False) - existing.currency = version.contract.currency existing.pricing = pricing existing.contract_version_id = version.id + existing.meter_id = m0.id existing.degraded = False existing.computed_at = now else: @@ -303,6 +450,7 @@ def compute_period(session: Session, t0: datetime, *, overwrite: bool = False) - currency=version.contract.currency, pricing=pricing, contract_version_id=version.id, + meter_id=m0.id, degraded=False, computed_at=now, ) @@ -316,9 +464,25 @@ def _upsert_degraded( t0: datetime, now: datetime, existing: EnergyCostPeriod | None, + *, + meter_id: int | None, ) -> None: """Insert or update a degraded placeholder for period *t0*. + Parameters + ---------- + session: + Active SQLAlchemy session. + t0: + UTC start of the 15-minute period. + now: + Current UTC timestamp for the ``computed_at`` field. + existing: + The existing ``EnergyCostPeriod`` row for this period, or ``None``. + meter_id: + The meter ID to attribute this degraded period to, or ``None`` when + no meter epoch covers the period (no-meter degraded case). + When *existing* is not None (row was previously written — either degraded or successful), the row is explicitly reset to the standard degraded state. This is required for the ``recompute_range`` (overwrite=True) path: if the @@ -326,6 +490,11 @@ def _upsert_degraded( since disappeared, the stale non-zero costs must be cleared so the row accurately reflects the current "missing readings" state rather than masquerading as a valid result. + + The ``meter_id`` is always updated to reflect the current meter attribution + judgment (the result of ``meter_at`` at the time of recompute). This + ensures that a retroactive ``started_at`` change + ``recompute_range`` will + re-attribute historical degraded periods to the correct meter epoch. """ if existing is not None: # Explicitly reset to degraded state — identical field values to the @@ -340,6 +509,7 @@ def _upsert_degraded( existing.net_cost = 0.0 existing.pricing = {} existing.contract_version_id = None + existing.meter_id = meter_id existing.degraded = True existing.computed_at = now else: @@ -355,6 +525,7 @@ def _upsert_degraded( currency="EUR", # placeholder; real currency known after contract lookup pricing={}, contract_version_id=None, + meter_id=meter_id, degraded=True, computed_at=now, ) @@ -381,6 +552,7 @@ def compute_closed_periods(session: Session) -> int: 3. For each boundary, calls ``compute_period(overwrite=False)``, which: - Skips periods that already have a *successful* (non-degraded) record. - Retries periods that have a *degraded* record. + - Writes degraded rows for periods with no meter or cross-meter boundaries. - Skips periods for which no active contract version exists or the Tibber price is unavailable (without writing a degraded row). @@ -429,11 +601,17 @@ def recompute_range(session: Session, start: datetime, end: datetime) -> int: This is the *explicit opt-in* path for recovering from: - Periods where readings or prices arrived late. - Price corrections (new contract version retroactively applied). + - Retroactive meter changes (``update_meter`` with new ``started_at``) — + re-running this function will re-judge meter attribution and re-compute + costs using the corrected epoch boundaries. - Any other reason to override the immutability guard. The function iterates over every UTC quarter-hour boundary in ``[floor(start), end)`` and calls ``compute_period(overwrite=True)``. - Existing rows (including successful ones) are overwritten. + Existing rows (including successful ones) are overwritten; their + ``meter_id`` fields will reflect the *current* ``meter_at`` judgment for + each period's start timestamp, naturally re-attributing periods when meter + ``started_at`` values have been retroactively corrected. Parameters ---------- diff --git a/tests/test_energy_cost.py b/tests/test_energy_cost.py index 8aa5ade..0686cc1 100644 --- a/tests/test_energy_cost.py +++ b/tests/test_energy_cost.py @@ -41,9 +41,11 @@ from app.models.energy import ( EnergyContract, EnergyContractVersion, EnergyCostPeriod, + Meter, TibberPrice, ) from app.services.energy_cost import ( + _MAX_DELTA_KWH, compute_closed_periods, compute_period, floor_to_quarter, @@ -179,6 +181,45 @@ def _make_reading( return r +def _make_meter( + session: Session, + *, + started_at: datetime, + ended_at: datetime | None = None, + label: str = "Test Meter", + commodity: str = "electricity", + reason: str = "initial", + note: str | None = None, +) -> Meter: + """Insert and flush a Meter row; return the ORM object.""" + now = datetime.now(_UTC) + m = Meter( + label=label, + commodity=commodity, + started_at=started_at, + ended_at=ended_at, + reason=reason, + note=note, + created_at=now, + ) + session.add(m) + session.flush() + return m + + +def _make_active_meter(session: Session, *, started_at: datetime | None = None) -> Meter: + """Insert and flush an active electricity meter covering the full test day. + + By default the meter starts at 2026-06-23 00:00 UTC (the beginning of the + test day), covering all boundaries used by the standard test helpers + (_T0 = 10:00, _T1 = 10:15, etc.). + """ + if started_at is None: + # Start well before any test boundary to cover the whole test day. + started_at = datetime(2026, 6, 23, 0, 0, 0, tzinfo=_UTC) + return _make_meter(session, started_at=started_at, ended_at=None) + + def _make_tibber_price( session: Session, *, @@ -252,30 +293,37 @@ class TestFloorToQuarter: class TestRegisterAt: + """register_at now requires a Meter argument; all tests use an active meter.""" + def test_returns_none_when_no_readings(self, energy_db: Session) -> None: - result = register_at(energy_db, _ts(10, 0)) + meter = _make_active_meter(energy_db) + energy_db.commit() + result = register_at(energy_db, _ts(10, 0), meter) assert result is None def test_returns_most_recent_at_or_before_boundary(self, energy_db: Session) -> None: + meter = _make_active_meter(energy_db) # Insert two readings: one before boundary, one after. _make_reading(energy_db, recorded_at=_ts(9, 55), d1="100.0", d2="200.0", r1="10.0", r2="20.0", source_id=1) _make_reading(energy_db, recorded_at=_ts(10, 5), d1="999.0", d2="999.0", r1="999.0", r2="999.0", source_id=2) energy_db.commit() - result = register_at(energy_db, _ts(10, 0)) + result = register_at(energy_db, _ts(10, 0), meter) assert result is not None assert result["d1"] == Decimal("100.0") assert result["d2"] == Decimal("200.0") def test_exact_boundary_included(self, energy_db: Session) -> None: + meter = _make_active_meter(energy_db) _make_reading(energy_db, recorded_at=_ts(10, 0), d1="500.0", d2="600.0", r1="50.0", r2="60.0", source_id=1) energy_db.commit() - result = register_at(energy_db, _ts(10, 0)) + result = register_at(energy_db, _ts(10, 0), meter) assert result is not None assert result["d1"] == Decimal("500.0") def test_missing_register_key_returns_none(self, energy_db: Session) -> None: + meter = _make_active_meter(energy_db) r = DsmrReading( recorded_at=_ts(10, 0), source_id=99, @@ -284,10 +332,11 @@ class TestRegisterAt: energy_db.add(r) energy_db.commit() - result = register_at(energy_db, _ts(10, 0)) + result = register_at(energy_db, _ts(10, 0), meter) assert result is None def test_null_register_value_returns_none(self, energy_db: Session) -> None: + meter = _make_active_meter(energy_db) r = DsmrReading( recorded_at=_ts(10, 0), source_id=88, @@ -301,20 +350,61 @@ class TestRegisterAt: energy_db.add(r) energy_db.commit() - result = register_at(energy_db, _ts(10, 0)) + result = register_at(energy_db, _ts(10, 0), meter) assert result is None def test_values_are_decimal(self, energy_db: Session) -> None: + meter = _make_active_meter(energy_db) _make_reading(energy_db, recorded_at=_ts(10, 0), d1="20915.154", d2="18372.099", r1="1234.567", r2="890.123", source_id=1) energy_db.commit() - result = register_at(energy_db, _ts(10, 0)) + result = register_at(energy_db, _ts(10, 0), meter) assert result is not None assert isinstance(result["d1"], Decimal) assert result["d1"] == Decimal("20915.154") assert result["r2"] == Decimal("890.123") + def test_reading_outside_meter_window_excluded(self, energy_db: Session) -> None: + """A reading before meter.started_at must not be returned (cross-meter isolation).""" + # Meter starts at 10:00 — a reading at 09:55 is from the old epoch. + meter = _make_meter(energy_db, started_at=_ts(10, 0), ended_at=None) + _make_reading(energy_db, recorded_at=_ts(9, 55), d1="100.0", d2="200.0", + r1="10.0", r2="20.0", source_id=1) + energy_db.commit() + + # Boundary is 10:00; the reading at 09:55 is before meter.started_at. + result = register_at(energy_db, _ts(10, 0), meter) + assert result is None, ( + "register_at must not return a reading from before meter.started_at" + ) + + def test_reading_at_meter_started_at_included(self, energy_db: Session) -> None: + """A reading exactly at meter.started_at must be included (half-open lower bound).""" + meter = _make_meter(energy_db, started_at=_ts(10, 0), ended_at=None) + _make_reading(energy_db, recorded_at=_ts(10, 0), d1="500.0", d2="600.0", + r1="50.0", r2="60.0", source_id=1) + energy_db.commit() + + result = register_at(energy_db, _ts(10, 0), meter) + assert result is not None + assert result["d1"] == Decimal("500.0") + + def test_reading_at_meter_ended_at_excluded(self, energy_db: Session) -> None: + """A reading exactly at meter.ended_at must be excluded (half-open upper bound).""" + # Meter covers [10:00, 10:15) — a reading at 10:15 belongs to the next epoch. + meter = _make_meter(energy_db, started_at=_ts(10, 0), ended_at=_ts(10, 15)) + _make_reading(energy_db, recorded_at=_ts(10, 15), d1="500.0", d2="600.0", + r1="50.0", r2="60.0", source_id=1) + energy_db.commit() + + # Boundary is 10:15, reading is at 10:15 = ended_at → excluded. + result = register_at(energy_db, _ts(10, 15), meter) + assert result is None, ( + "register_at must exclude a reading exactly at meter.ended_at " + "(half-open upper bound)" + ) + # --------------------------------------------------------------------------- # 1-2. compute_period — manual dual-tariff @@ -351,12 +441,14 @@ _END_R2 = "3000.100" def _setup_manual_scenario(session: Session) -> EnergyContractVersion: - """Create active manual contract + two boundary readings; return the version.""" + """Create active manual contract + active meter + two boundary readings; return the version.""" contract = _make_contract(session, kind="manual", active=True) version = _make_version( session, contract, _MANUAL_VALUES, effective_from=_ts(0, 0), # covers t0=10:00 ) + # Active meter covering the full test day (started before T0). + _make_active_meter(session) # Start reading (at t0) _make_reading(session, recorded_at=_T0, d1=_START_D1, d2=_START_D2, r1=_START_R1, r2=_START_R2, source_id=1) @@ -535,6 +627,7 @@ class TestComputePeriodTibber: def _setup(self, session: Session, total: float = 0.25) -> tuple[EnergyContractVersion, TibberPrice]: contract = _make_contract(session, kind="tibber", active=True) version = _make_version(session, contract, _TIBBER_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(session) _make_reading(session, recorded_at=_T0, d1=_START_D1, d2=_START_D2, r1=_START_R1, r2=_START_R2, source_id=1) _make_reading(session, recorded_at=_T1, d1=_END_D1, d2=_END_D2, @@ -600,6 +693,7 @@ class TestComputePeriodMissingTibberPrice: """When Tibber price is absent for the period, no EnergyCostPeriod is written.""" contract = _make_contract(energy_db, kind="tibber", active=True) _make_version(energy_db, contract, _TIBBER_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) _make_reading(energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, r1=_START_R1, r2=_START_R2, source_id=1) _make_reading(energy_db, recorded_at=_T1, d1=_END_D1, d2=_END_D2, @@ -619,6 +713,7 @@ class TestComputePeriodMissingTibberPrice: """A TibberPrice with starts_at > t0 must NOT be used; period is skipped.""" contract = _make_contract(energy_db, kind="tibber", active=True) _make_version(energy_db, contract, _TIBBER_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) _make_reading(energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, r1=_START_R1, r2=_START_R2, source_id=1) _make_reading(energy_db, recorded_at=_T1, d1=_END_D1, d2=_END_D2, @@ -643,9 +738,10 @@ class TestComputePeriodMissingTibberPrice: class TestComputePeriodMissingReadings: def test_degraded_when_start_reading_missing(self, energy_db: Session) -> None: - """No DsmrReading at or before t0 → degraded row written.""" + """No DsmrReading at or before t0 within the meter window → degraded row written.""" contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) # Only an end reading; no start reading. _make_reading(energy_db, recorded_at=_T1, d1=_END_D1, d2=_END_D2, r1=_END_R1, r2=_END_R2, source_id=2) @@ -663,7 +759,7 @@ class TestComputePeriodMissingReadings: assert row.net_cost == 0.0 def test_degraded_when_end_reading_missing(self, energy_db: Session) -> None: - """No DsmrReading at or before t1 → degraded row written. + """No DsmrReading at or before t1 within the meter window → degraded row written. We place a reading BEFORE t0 (so t0 boundary has data) but the first reading AT OR AFTER t1 is only after t1+5min, leaving the t1 boundary @@ -677,6 +773,7 @@ class TestComputePeriodMissingReadings: """ contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) # Only a reading AFTER t1 — no reading at or before t1. _make_reading(energy_db, recorded_at=_ts(10, 20), d1=_END_D1, d2=_END_D2, r1=_END_R1, r2=_END_R2, source_id=2) @@ -694,6 +791,7 @@ class TestComputePeriodMissingReadings: """Degraded rows written due to missing readings have contract_version_id=None.""" contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) # No readings at all. energy_db.commit() @@ -711,6 +809,7 @@ class TestComputePeriodMissingReadings: """ contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) # First compute: no readings at all → degraded. energy_db.commit() @@ -786,6 +885,8 @@ class TestCrossVersionSelection: effective_from=_ts(8, 0), effective_to=None, ) + # Active meter covering the full test day. + _make_active_meter(session) session.commit() return v1, v2 @@ -857,9 +958,10 @@ class TestCrossVersionSelection: class TestSummarize: def _setup_two_periods(self, session: Session) -> None: - """Insert contract + readings for two consecutive 15-min periods and compute them.""" + """Insert contract + meter + readings for two consecutive 15-min periods and compute them.""" contract = _make_contract(session, kind="manual", active=True) _make_version(session, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(session) # Period 1: [10:00, 10:15) _make_reading(session, recorded_at=_T0, d1=_START_D1, d2=_START_D2, @@ -954,6 +1056,7 @@ class TestSummarize: """ contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) # No readings at all → degraded period. energy_db.commit() compute_period(energy_db, _T0) @@ -977,6 +1080,7 @@ class TestSummarize: contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) energy_db.commit() # Summarize over exactly 1 UTC day (no periods in DB — only standing/credits). @@ -1294,12 +1398,28 @@ class TestSummarizePrincipleC: class TestComputeClosedPeriods: def test_skips_future_periods(self, energy_db: Session) -> None: - """Periods whose t1 > now must not be computed.""" + """Periods whose t1 > now must not be computed. + + The meter covers from before the lookback window so that all past periods + have a meter; the contract effective_from=now ensures that every past + period's contract-version lookup returns None, causing them all to be + skipped (no row written). Only the current open period [now, now+15min) + is a "future" period — and the normal tick never writes it. + + Written = 0 confirms that no row was written (neither the future period + nor any past period without a contract). + """ + now = datetime.now(_UTC) contract = _make_contract(energy_db, kind="manual", active=True) - _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=datetime.now(_UTC)) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=now) + # Meter covering from well before the 7-day lookback window so every past + # period has a meter. Without a matching contract version those past + # periods are all skipped (not written as degraded). + _make_meter(energy_db, started_at=now - timedelta(days=10), ended_at=None) energy_db.commit() # The period [now, now+15min) is still open — should not be computed. + # All closed past periods: meter OK, but no contract version → skip, no write. written = compute_closed_periods(energy_db) assert written == 0 # nothing written (no closed periods with data) @@ -1314,6 +1434,8 @@ class TestComputeClosedPeriods: contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=past_t0 - timedelta(hours=1)) + # Meter covering from before past_t0. + _make_meter(energy_db, started_at=past_t0 - timedelta(hours=1), ended_at=None) _make_reading(energy_db, recorded_at=past_t0, d1=_START_D1, d2=_START_D2, r1=_START_R1, r2=_START_R2, source_id=1) _make_reading(energy_db, recorded_at=past_t1, d1=_END_D1, d2=_END_D2, @@ -1349,6 +1471,8 @@ class TestComputeClosedPeriods: contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=past_t0 - timedelta(hours=1)) + # Meter covering from before past_t0. + _make_meter(energy_db, started_at=past_t0 - timedelta(hours=1), ended_at=None) # First compute: no readings → degraded. energy_db.commit() compute_period(energy_db, past_t0) @@ -1386,6 +1510,7 @@ class TestRecomputeRange: """recompute_range must overwrite non-degraded rows (explicit opt-in).""" contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) _make_reading(energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, r1=_START_R1, r2=_START_R2, source_id=1) _make_reading(energy_db, recorded_at=_T1, d1=_END_D1, d2=_END_D2, @@ -1415,6 +1540,7 @@ class TestRecomputeRange: """recompute_range returns the number of periods actually written.""" contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) # Add readings for two consecutive periods: [10:00,10:15), [10:15,10:30). _make_reading(energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, @@ -1436,6 +1562,7 @@ class TestRecomputeRange: """Calling recompute_range twice must not create duplicate rows.""" contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) _make_reading(energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, r1=_START_R1, r2=_START_R2, source_id=1) _make_reading(energy_db, recorded_at=_T1, d1=_END_D1, d2=_END_D2, @@ -1469,6 +1596,7 @@ class TestRecomputeRange: # Step 1: successful compute. contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) start_reading = _make_reading( energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, r1=_START_R1, r2=_START_R2, source_id=1, @@ -1528,16 +1656,19 @@ class TestRegisterAtFreshness: ``_READING_MAX_STALENESS`` (15 minutes) old at a given boundary is treated as absent so the period is marked degraded rather than producing a zero-delta fake-success row. + + All tests now pass a meter that covers the reading and boundary times. """ def test_stale_reading_returns_none(self, energy_db: Session) -> None: """A reading older than 15 min before the boundary must be rejected.""" reading_time = _ts(10, 0) # 10:00 boundary = _ts(10, 30) # 10:30 — 30 min later (> 15 min staleness) + meter = _make_active_meter(energy_db) _make_reading(energy_db, recorded_at=reading_time, source_id=1) energy_db.commit() - result = register_at(energy_db, boundary) + result = register_at(energy_db, boundary, meter) assert result is None, ( "register_at must return None when the closest reading is more than " "15 minutes before the boundary" @@ -1547,12 +1678,13 @@ class TestRegisterAtFreshness: """A reading within 15 min of the boundary must be returned normally.""" reading_time = _ts(10, 5) # 10:05 boundary = _ts(10, 15) # 10:15 — only 10 min gap (within window) + meter = _make_active_meter(energy_db) _make_reading(energy_db, recorded_at=reading_time, d1="30000.0", d2="15000.0", r1="1000.0", r2="500.0", source_id=1) energy_db.commit() - result = register_at(energy_db, boundary) + result = register_at(energy_db, boundary, meter) assert result is not None, ( "register_at must return the reading when it is within the 15-min staleness window" ) @@ -1562,12 +1694,13 @@ class TestRegisterAtFreshness: """A reading exactly 15 min before the boundary sits at the edge — accepted.""" reading_time = _ts(10, 0) # 10:00 boundary = _ts(10, 15) # 10:15 — exactly 15 min gap + meter = _make_active_meter(energy_db) _make_reading(energy_db, recorded_at=reading_time, d1="40000.0", d2="20000.0", r1="2000.0", r2="1000.0", source_id=1) energy_db.commit() - result = register_at(energy_db, boundary) + result = register_at(energy_db, boundary, meter) # boundary - reading == 15 min == staleness limit → NOT stale (strict <) assert result is not None @@ -1581,12 +1714,13 @@ class TestRegisterAtFreshness: """ # Insert a reading at a realistic "now" time. reading_time = _ts(10, 0) # 10:00 on 2026-06-23 + meter = _make_active_meter(energy_db) _make_reading(energy_db, recorded_at=reading_time, source_id=1) energy_db.commit() # Query with a far-future boundary (e.g. end of month, same day +8 hours) far_future_boundary = _ts(18, 0) # 18:00 — 8 hours later - result = register_at(energy_db, far_future_boundary) + result = register_at(energy_db, far_future_boundary, meter) assert result is None, ( "register_at must return None for a far-future boundary rather than " "the latest historical reading" @@ -1615,6 +1749,7 @@ class TestFuturePeriodDegraded: """ contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) # Only a reading at 10:00 — far from the 14:00/14:15 boundaries. _make_reading(energy_db, recorded_at=_ts(10, 0), source_id=1) energy_db.commit() @@ -1643,12 +1778,14 @@ class TestRecomputeRangeNoFuturePeriods: future end datetime; the only periods that may be written are those where t1 <= now. No row with period_start > now should appear in the DB. """ + now = datetime.now(UTC) contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + # Meter covering from before now. + _make_meter(energy_db, started_at=now - timedelta(hours=1), ended_at=None) # No readings needed — we only care that future rows are NOT written. energy_db.commit() - now = datetime.now(UTC) far_future_end = now + timedelta(days=7) # Use a past start so there are some candidate periods. past_start = floor_to_quarter(now - timedelta(minutes=30)) @@ -1678,15 +1815,21 @@ class TestRecomputeRangeNoFuturePeriods: We only verify that no future rows exist after step 1 (Fix B), because that is the core slot-squatting prevention. """ + now = datetime.now(UTC) contract = _make_contract(energy_db, kind="manual", active=True) # Contract effective from far past so all periods in scope have a version. _make_version( energy_db, contract, _MANUAL_VALUES, effective_from=datetime(2020, 1, 1, tzinfo=UTC), ) + # Meter covering from far past. + _make_meter( + energy_db, + started_at=datetime(2020, 1, 1, tzinfo=UTC), + ended_at=None, + ) energy_db.commit() - now = datetime.now(UTC) far_future = now + timedelta(days=30) past_start = floor_to_quarter(now - timedelta(hours=1)) @@ -1714,6 +1857,7 @@ class TestNormalPeriodsUnaffected: """A normal period with readings seconds before each boundary → non-degraded.""" contract = _make_contract(energy_db, kind="manual", active=True) _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) # Readings placed very close to boundaries (as DSMR normally delivers them). # t0=10:00, reading at 09:59:50 (10 s before); t1=10:15, reading at 10:14:55 (5 s before) @@ -1734,3 +1878,388 @@ class TestNormalPeriodsUnaffected: ) # import_cost matches hand-calc: same deltas as standard scenario assert abs(row.import_cost - 0.4101) < 1e-6 + + +# --------------------------------------------------------------------------- +# M7-T03 New tests: meter-aware billing engine +# --------------------------------------------------------------------------- + + +class TestMeterAwareComputePeriod: + """M7-T03 Acceptance criteria: meter-aware compute_period behaviour. + + Covers: + ① Same-meter period: delta computed correctly, meter_id attributed. + ② Cross-meter boundary period: degraded. + ③ No active meter coverage: degraded with meter_id=None. + ④ Negative delta: degraded (D6 guard — no negative costs). + ⑤ Super-large delta exceeding _MAX_DELTA_KWH: degraded (D6 guard). + ⑥ Recompute re-judges meter attribution. + ⑦ No-contract skip preserved inside single-meter path. + """ + + # ① Same-meter period: delta correct + meter_id attributed + def test_same_meter_period_computes_correctly_with_meter_id( + self, energy_db: Session + ) -> None: + """A normal period within a single meter epoch: costs correct, meter_id set.""" + contract = _make_contract(energy_db, kind="manual", active=True) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + meter = _make_active_meter(energy_db) + _make_reading(energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, + r1=_START_R1, r2=_START_R2, source_id=1) + _make_reading(energy_db, recorded_at=_T1, d1=_END_D1, d2=_END_D2, + r1=_END_R1, r2=_END_R2, source_id=2) + energy_db.commit() + + result = compute_period(energy_db, _T0) + assert result is True + + row = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + + # Cost correctness (same as manual scenario hand-calc). + assert row.degraded is False + assert abs(row.import_cost - 0.4101) < 1e-9 + assert abs(row.export_revenue - 0.005) < 1e-9 + assert abs(row.net_cost - 0.4051) < 1e-9 + + # meter_id must be set to the active meter's id. + assert row.meter_id == meter.id, ( + f"Expected meter_id={meter.id}, got {row.meter_id}" + ) + + # ② Cross-meter boundary: degraded + def test_cross_meter_boundary_period_is_degraded(self, energy_db: Session) -> None: + """A period spanning two meter epochs must be written as degraded. + + Scenario: + - Old meter: [09:00, 10:15) — covers t0=10:00 but NOT t1=10:15. + - New meter: [10:15, ∞) — covers t1=10:15. + - Period [10:00, 10:15): m0 != m1 → degraded with meter_id=m0.id. + """ + old_meter = _make_meter(energy_db, started_at=_ts(9, 0), ended_at=_T1) + new_meter = _make_meter(energy_db, started_at=_T1, ended_at=None) + contract = _make_contract(energy_db, kind="manual", active=True) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + # Readings exist, but the period still spans two meters. + _make_reading(energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, + r1=_START_R1, r2=_START_R2, source_id=1) + _make_reading(energy_db, recorded_at=_T1, d1=_END_D1, d2=_END_D2, + r1=_END_R1, r2=_END_R2, source_id=2) + energy_db.commit() + + result = compute_period(energy_db, _T0) + assert result is True # a row was written + + row = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + assert row.degraded is True, ( + "A period spanning two meter epochs must be written as degraded" + ) + # meter_id should be m0 (t0's meter), not m1. + assert row.meter_id == old_meter.id, ( + f"Cross-meter degraded row should attribute to old_meter (id={old_meter.id}), " + f"got meter_id={row.meter_id}" + ) + # No real cost should be produced. + assert row.import_cost == 0.0 + assert row.net_cost == 0.0 + # Suppress unused variable warning. + _ = new_meter + + # ③ No active meter coverage: degraded with meter_id=None + def test_no_meter_coverage_is_degraded_with_null_meter_id( + self, energy_db: Session + ) -> None: + """When no meter epoch covers t0, the period is written as degraded(meter_id=None).""" + # No meter inserted at all. + contract = _make_contract(energy_db, kind="manual", active=True) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_reading(energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, + r1=_START_R1, r2=_START_R2, source_id=1) + _make_reading(energy_db, recorded_at=_T1, d1=_END_D1, d2=_END_D2, + r1=_END_R1, r2=_END_R2, source_id=2) + energy_db.commit() + + result = compute_period(energy_db, _T0) + assert result is True + + row = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + assert row.degraded is True, "No-meter period must be written as degraded" + assert row.meter_id is None, ( + "No-meter degraded row must have meter_id=None, " + f"got meter_id={row.meter_id}" + ) + assert row.import_cost == 0.0 + + # ④ Negative delta: degraded (D6 guard) + def test_negative_delta_is_degraded(self, energy_db: Session) -> None: + """A period with any negative register delta → degraded (D6 guard). + + Negative deltas occur after a meter reset or DSMR rollover and must + never produce a negative cost row. + """ + contract = _make_contract(energy_db, kind="manual", active=True) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) + # End reading has LOWER d1 than start → negative delta for d1. + _make_reading(energy_db, recorded_at=_T0, d1="20000.500", d2="10000.000", + r1="5000.000", r2="3000.000", source_id=1) + _make_reading(energy_db, recorded_at=_T1, d1="20000.000", d2="10001.200", + r1="5000.000", r2="3000.100", source_id=2) + energy_db.commit() + + result = compute_period(energy_db, _T0) + assert result is True + + row = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + assert row.degraded is True, ( + "A period with a negative delta must be written as degraded (no negative costs)" + ) + assert row.import_cost == 0.0 + assert row.net_cost == 0.0 + + # ⑤ Super-large delta: degraded (D6 guard) + def test_super_large_delta_is_degraded(self, energy_db: Session) -> None: + """A period with a delta > _MAX_DELTA_KWH → degraded (D6 guard). + + Implausibly large deltas indicate an anomaly (wrong scale, DSMR + reporting bug) and must never produce a grossly inflated cost row. + """ + contract = _make_contract(energy_db, kind="manual", active=True) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) + # d1 delta = 200 kWh >> _MAX_DELTA_KWH (100 kWh). + _make_reading(energy_db, recorded_at=_T0, d1="10000.000", d2="10000.000", + r1="5000.000", r2="3000.000", source_id=1) + _make_reading(energy_db, recorded_at=_T1, d1="10200.000", d2="10000.000", + r1="5000.000", r2="3000.000", source_id=2) + energy_db.commit() + + result = compute_period(energy_db, _T0) + assert result is True + + row = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + assert row.degraded is True, ( + f"A delta > {_MAX_DELTA_KWH} kWh must be written as degraded (D6 guard)" + ) + assert row.import_cost == 0.0 + assert row.net_cost == 0.0 + + # ⑤b Delta exactly at limit is NOT degraded (guard fires at strictly >) + def test_delta_exactly_at_max_limit_not_degraded(self, energy_db: Session) -> None: + """A delta exactly equal to _MAX_DELTA_KWH is NOT degraded. + + The guard condition is strictly > _MAX_DELTA_KWH, so a delta of exactly + 100 kWh passes through and is computed normally. This is intentional: + the threshold is set far above any plausible residential consumption + (400 kW average over 15 min) to avoid false positives. + """ + contract = _make_contract(energy_db, kind="manual", active=True) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) + # d1 delta = _MAX_DELTA_KWH exactly (100 kWh) — must NOT degrade. + max_delta = float(_MAX_DELTA_KWH) + _make_reading(energy_db, recorded_at=_T0, d1="10000.000", d2="10000.000", + r1="5000.000", r2="3000.000", source_id=1) + _make_reading(energy_db, recorded_at=_T1, d1=str(10000.0 + max_delta), + d2="10000.000", r1="5000.000", r2="3000.000", source_id=2) + energy_db.commit() + + result = compute_period(energy_db, _T0) + assert result is True + + row = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + assert row.degraded is False, ( + f"A delta exactly at _MAX_DELTA_KWH ({max_delta} kWh) must NOT trigger " + "the D6 guard (guard condition is strictly >)" + ) + assert row.import_cost > 0 + + # ⑤c Delta strictly above limit IS degraded + def test_delta_strictly_above_max_limit_is_degraded(self, energy_db: Session) -> None: + """A delta strictly greater than _MAX_DELTA_KWH → degraded (D6 guard fires).""" + contract = _make_contract(energy_db, kind="manual", active=True) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) + # d1 delta = _MAX_DELTA_KWH + 0.001 (strictly over threshold) → degraded. + max_delta_plus = float(_MAX_DELTA_KWH) + 0.001 + _make_reading(energy_db, recorded_at=_T0, d1="10000.000", d2="10000.000", + r1="5000.000", r2="3000.000", source_id=1) + _make_reading(energy_db, recorded_at=_T1, d1=str(10000.0 + max_delta_plus), + d2="10000.000", r1="5000.000", r2="3000.000", source_id=2) + energy_db.commit() + + result = compute_period(energy_db, _T0) + assert result is True + + row = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + assert row.degraded is True, ( + f"A delta of {max_delta_plus} kWh (> _MAX_DELTA_KWH) must trigger D6 guard" + ) + + # ⑤c Delta just below limit is NOT degraded + def test_delta_just_below_max_limit_not_degraded(self, energy_db: Session) -> None: + """A delta just below _MAX_DELTA_KWH must NOT trigger the D6 guard.""" + contract = _make_contract(energy_db, kind="manual", active=True) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) + # d1 delta = 99.999 kWh — just under 100, should compute normally. + _make_reading(energy_db, recorded_at=_T0, d1="10000.000", d2="10000.000", + r1="5000.000", r2="3000.000", source_id=1) + _make_reading(energy_db, recorded_at=_T1, d1="10099.999", d2="10000.000", + r1="5000.000", r2="3000.000", source_id=2) + energy_db.commit() + + result = compute_period(energy_db, _T0) + assert result is True + + row = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + assert row.degraded is False, ( + "A delta of 99.999 kWh must NOT trigger the D6 guard (guard fires at > 100 kWh)" + ) + assert row.import_cost > 0 + + # ⑥ recompute re-judges meter attribution + def test_recompute_re_judges_meter_attribution(self, energy_db: Session) -> None: + """recompute_range must re-attribute meter_id to the current meter_at judgment. + + Scenario: + 1. Compute period [10:00, 10:15) with meter M1 active → row has meter_id=M1.id. + 2. Retroactively close M1 at 10:00 and open M2 from 09:00 (so M2 covers both + boundaries, by updating the meter row directly — simulating update_meter). + 3. Call recompute_range → row's meter_id must now be M2.id. + """ + # Step 1: initial compute with M1. + contract = _make_contract(energy_db, kind="manual", active=True) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + m1 = _make_active_meter(energy_db) + _make_reading(energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, + r1=_START_R1, r2=_START_R2, source_id=1) + _make_reading(energy_db, recorded_at=_T1, d1=_END_D1, d2=_END_D2, + r1=_END_R1, r2=_END_R2, source_id=2) + energy_db.commit() + + compute_period(energy_db, _T0) + energy_db.commit() + + row_before = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + assert row_before.meter_id == m1.id, "Pre-condition: meter_id should be m1.id" + + # Step 2: retroactively replace M1 with M2. + # Close M1 (make it a zero-width closed epoch before any test reading). + # Open M2 covering the whole day. + m1.ended_at = datetime(2026, 6, 22, 0, 0, 0, tzinfo=_UTC) # before test day + m2 = _make_meter( + energy_db, + started_at=datetime(2026, 6, 22, 0, 0, 0, tzinfo=_UTC), + ended_at=None, + label="Replacement Meter", + ) + energy_db.commit() + + # Step 3: recompute the range — should re-attribute to m2. + count = recompute_range(energy_db, _T0, _T1) + assert count == 1 + + energy_db.expire_all() + row_after = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + assert row_after.meter_id == m2.id, ( + f"After recompute, meter_id should be m2.id={m2.id}, " + f"got {row_after.meter_id}" + ) + assert row_after.degraded is False + + # ⑦ No-contract skip preserved inside single-meter path + def test_no_contract_is_skip_not_degraded(self, energy_db: Session) -> None: + """Single-meter path with no active contract → skip (no row written). + + This ensures that 'meter OK, contract missing' still produces a skip + (not a degraded row), preserving the existing skip semantics. + """ + # Active meter but NO contract. + _make_active_meter(energy_db) + _make_reading(energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, + r1=_START_R1, r2=_START_R2, source_id=1) + _make_reading(energy_db, recorded_at=_T1, d1=_END_D1, d2=_END_D2, + r1=_END_R1, r2=_END_R2, source_id=2) + energy_db.commit() + + result = compute_period(energy_db, _T0) + assert result is False, ( + "Single-meter, no-contract period must be skipped (return False, no row)" + ) + + rows = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalars().all() + assert len(rows) == 0, ( + "No EnergyCostPeriod row must be written when skipping due to missing contract" + ) + + # Degraded within-meter rows get the correct meter_id + def test_degraded_within_meter_has_meter_id(self, energy_db: Session) -> None: + """Degraded rows due to missing readings within a known meter get meter_id set.""" + contract = _make_contract(energy_db, kind="manual", active=True) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + meter = _make_active_meter(energy_db) + # No readings at all → will degrade due to missing readings. + energy_db.commit() + + result = compute_period(energy_db, _T0) + assert result is True + + row = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + assert row.degraded is True + assert row.meter_id == meter.id, ( + f"Within-meter degraded row must carry meter_id={meter.id}, " + f"got {row.meter_id}" + ) + + # Delta sanity guard with all-zero deltas (zero is fine, not negative) + def test_zero_delta_not_degraded(self, energy_db: Session) -> None: + """A period with all-zero deltas (no energy used) must NOT be degraded. + + Zero is valid — perfectly matching start/end readings means no energy + was consumed or produced in that window. + """ + contract = _make_contract(energy_db, kind="manual", active=True) + _make_version(energy_db, contract, _MANUAL_VALUES, effective_from=_ts(0, 0)) + _make_active_meter(energy_db) + # Identical start and end readings → all deltas = 0. + _make_reading(energy_db, recorded_at=_T0, d1=_START_D1, d2=_START_D2, + r1=_START_R1, r2=_START_R2, source_id=1) + _make_reading(energy_db, recorded_at=_T1, d1=_START_D1, d2=_START_D2, + r1=_START_R1, r2=_START_R2, source_id=2) + energy_db.commit() + + result = compute_period(energy_db, _T0) + assert result is True + + row = energy_db.execute( + select(EnergyCostPeriod).where(EnergyCostPeriod.period_start == _T0) + ).scalar_one() + assert row.degraded is False, "All-zero deltas must not trigger the D6 guard" + assert row.import_cost == 0.0 + assert row.net_cost == 0.0