"""Behaviour tests for the WarmteLink per-source admission state machine.""" from __future__ import annotations from dataclasses import replace from datetime import UTC, datetime, timedelta from decimal import Decimal from pathlib import Path from types import SimpleNamespace import pytest from alembic import command from alembic.config import Config from sqlalchemy import create_engine, func, select from sqlalchemy.orm import Session, sessionmaker from app.integrations.p1 import IntegrityStatus, ObisField, P1Channel, P1Telegram, dsmr_crc16, parse_telegram from app.models.meter_source import MeterSource, MeterSourceBinding, MeterSourceChannel, WarmteLinkReading from app.services.warmtelink_ingest import WarmteLinkIngestor def _config(url: str) -> Config: config = Config("alembic_app.ini") config.set_main_option("sqlalchemy.url", url) return config @pytest.fixture() def database(tmp_path: Path): url = f"sqlite:///{tmp_path / 'warmtelink_ingest.db'}" command.upgrade(_config(url), "head") engine = create_engine(url, connect_args={"check_same_thread": False}) factory = sessionmaker(bind=engine, autoflush=False, class_=Session) with factory() as session: for name in ("one", "two"): now = datetime(2026, 8, 22, tzinfo=UTC) session.add( MeterSource( name=name, kind="warmtelink_serial", enabled=True, config={"path": f"/dev/{name}", "baudrate": 115200, "data_bits": 7, "parity": "N", "stop_bits": 1}, created_at=now, updated_at=now, ) ) session.commit() ids = list(session.scalars(select(MeterSource.id).order_by(MeterSource.id))) yield engine, factory, ids[-2:] engine.dispose() def _telegram( *, second: int, minute: int = 0, integrity: IntegrityStatus = IntegrityStatus.UNVERIFIABLE, water: str = "5.900", heat: str = "0.017", channels: tuple[int, ...] = (1, 2), device_type: str = "006", fingerprint: str | None = "a" * 64, channel_fingerprint: str | None = None, channel_fingerprints: dict[int, str | None] | None = None, timestamp: str | None = None, ) -> P1Telegram: fields: list[ObisField] = [] parsed_channels: list[P1Channel] = [] for number in channels: value, unit = (water, "m3") if number == 1 else (heat, "GJ") code = f"0-{number}:24.2.1" field = ObisField(code, (value + "*" + unit,), Decimal(value), unit) fields.append(field) parsed_channels.append( P1Channel( number, device_type if number == 1 else "012", ( channel_fingerprints[number] if channel_fingerprints is not None else fingerprint if channel_fingerprint is None else channel_fingerprint ), (field,), ) ) return P1Telegram( 100, integrity, "fixture", timestamp or f"26082212{minute:02d}{second:02d}S", fingerprint, tuple(fields), tuple(parsed_channels), ) def _counts(factory) -> tuple[int, int]: with factory() as session: return ( # Revision 16 supplies one DSMR channel; this task must not touch it. int(session.scalar(select(func.count(MeterSourceChannel.id))) or 0) - 1, int(session.scalar(select(func.count(WarmteLinkReading.id))) or 0), ) def _all_counts(factory) -> tuple[int, int, int]: channels, readings = _counts(factory) with factory() as session: bindings = int(session.scalar(select(func.count(MeterSourceBinding.id))) or 0) return channels, readings, bindings def _fixed_clock() -> datetime: return datetime(2026, 8, 22, 10, tzinfo=UTC) def _clock_at(value: datetime): return lambda: value class _AdvancingClock: def __init__(self, initial: datetime) -> None: self._next = initial self.calls = 0 def __call__(self) -> datetime: value = self._next self._next += timedelta(minutes=1) self.calls += 1 return value def _as_utc(value: datetime) -> datetime: return value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC) def test_summer_wall_clock_ignores_incorrect_w_marker(database): _engine, factory, (source_id, _) = database received_at = datetime(2026, 8, 23, 21, 40, tzinfo=UTC) ingest = WarmteLinkIngestor(clock=_clock_at(received_at)) with factory() as session: assert ingest.ingest( session, source_id=source_id, telegram=_telegram(second=0, integrity=IntegrityStatus.VALID, timestamp="260823234000W"), ) session.commit() with factory() as session: reading = session.scalar(select(WarmteLinkReading).order_by(WarmteLinkReading.id)) assert reading is not None assert _as_utc(reading.recorded_at) == datetime(2026, 8, 23, 21, 40, tzinfo=UTC) def test_winter_wall_clock_uses_amsterdam_standard_time_even_with_s_marker(database): _engine, factory, (source_id, _) = database received_at = datetime(2026, 1, 23, 22, 40, tzinfo=UTC) ingest = WarmteLinkIngestor(clock=_clock_at(received_at)) with factory() as session: assert ingest.ingest( session, source_id=source_id, telegram=_telegram(second=0, integrity=IntegrityStatus.VALID, timestamp="260123234000S"), ) session.commit() with factory() as session: reading = session.scalar(select(WarmteLinkReading).order_by(WarmteLinkReading.id)) assert reading is not None assert _as_utc(reading.recorded_at) == datetime(2026, 1, 23, 22, 40, tzinfo=UTC) @pytest.mark.parametrize( ("received_at", "expected"), [ (datetime(2026, 10, 25, 0, 30, tzinfo=UTC), datetime(2026, 10, 25, 0, 30, tzinfo=UTC)), (datetime(2026, 10, 25, 1, 30, tzinfo=UTC), datetime(2026, 10, 25, 1, 30, tzinfo=UTC)), ], ) def test_fall_back_fold_is_selected_from_received_at(database, received_at, expected): _engine, factory, (source_id, _) = database ingest = WarmteLinkIngestor(clock=_clock_at(received_at)) with factory() as session: assert ingest.ingest( session, source_id=source_id, telegram=_telegram(second=0, integrity=IntegrityStatus.VALID, timestamp="261025023000W"), ) session.commit() with factory() as session: reading = session.scalar(select(WarmteLinkReading).order_by(WarmteLinkReading.id)) assert reading is not None and _as_utc(reading.recorded_at) == expected @pytest.mark.parametrize("timestamp", ["260329023000S", "260823234000X"]) def test_nonexistent_or_invalid_marker_timestamp_fails_closed_with_generic_diagnostic(database, timestamp): _engine, factory, (source_id, _) = database ingest = WarmteLinkIngestor(clock=_clock_at(datetime(2026, 3, 29, 1, 30, tzinfo=UTC))) with factory() as session: assert not ingest.ingest( session, source_id=source_id, telegram=_telegram(second=0, integrity=IntegrityStatus.VALID, timestamp=timestamp), ) session.commit() source = session.get(MeterSource, source_id) assert source is not None and source.last_error == "WarmteLink frame could not be normalized" assert _counts(factory) == (0, 0) def test_timestamp_beyond_clock_skew_fails_closed_with_generic_diagnostic(database): _engine, factory, (source_id, _) = database ingest = WarmteLinkIngestor(clock=_clock_at(datetime(2026, 8, 23, 20, tzinfo=UTC))) with factory() as session: assert not ingest.ingest( session, source_id=source_id, telegram=_telegram(second=0, integrity=IntegrityStatus.VALID, timestamp="260823234000W"), ) session.commit() source = session.get(MeterSource, source_id) assert source is not None and source.last_error == "WarmteLink frame could not be normalized" assert _counts(factory) == (0, 0) @pytest.mark.parametrize("timestamp", ["260823234000X", "260823234600W"]) def test_post_clock_rejections_reuse_received_at_for_diagnostic(database, timestamp): _engine, factory, (source_id, _) = database received_at = datetime(2026, 8, 23, 21, 40, tzinfo=UTC) clock = _AdvancingClock(received_at) ingest = WarmteLinkIngestor(clock=clock) with factory() as session: assert not ingest.ingest( session, source_id=source_id, telegram=_telegram(second=0, integrity=IntegrityStatus.VALID, timestamp=timestamp), ) session.commit() source = session.get(MeterSource, source_id) assert source is not None assert _as_utc(source.updated_at) == received_at assert clock.calls == 1 def test_first_unverifiable_candidate_reuses_received_at_for_diagnostic(database): _engine, factory, (source_id, _) = database received_at = datetime(2026, 8, 22, 10, tzinfo=UTC) clock = _AdvancingClock(received_at) ingest = WarmteLinkIngestor(clock=clock) with factory() as session: assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=0)) session.commit() source = session.get(MeterSource, source_id) assert source is not None assert _as_utc(source.updated_at) == received_at assert clock.calls == 1 def test_accepted_frame_uses_received_at_for_all_source_timestamps(database): _engine, factory, (source_id, _) = database received_at = datetime(2026, 8, 22, 10, tzinfo=UTC) clock = _AdvancingClock(received_at) ingest = WarmteLinkIngestor(clock=clock) with factory() as session: assert ingest.ingest( session, source_id=source_id, telegram=_telegram(second=0, integrity=IntegrityStatus.VALID) ) session.commit() with factory() as session: source = session.get(MeterSource, source_id) reading = session.scalar(select(WarmteLinkReading).order_by(WarmteLinkReading.id)) assert source is not None and reading is not None assert _as_utc(reading.received_at) == received_at assert _as_utc(source.last_seen_at) == received_at assert _as_utc(source.updated_at) == received_at assert clock.calls == 1 @pytest.mark.parametrize( "path", ["parser", "integrity", "normalization", "fingerprint", "candidate", "continuity", "accepted"], ) def test_handle_frame_samples_the_receive_clock_once_for_each_admission_path(database, path): _engine, factory, (source_id, _) = database received_at = datetime(2026, 8, 22, 10, tzinfo=UTC) clock = _AdvancingClock(received_at) ingest = WarmteLinkIngestor(clock=clock) if path == "parser": assert not ingest.handle_frame( source_id, b"private", session_factory=factory, parser=lambda _raw: (_ for _ in ()).throw(ValueError()), ) expected_at = received_at elif path == "integrity": assert not ingest.handle_frame( source_id, b"private", session_factory=factory, parser=lambda _raw: _telegram(second=0, integrity=IntegrityStatus.INVALID), ) expected_at = received_at elif path == "normalization": assert not ingest.handle_frame( source_id, b"private", session_factory=factory, parser=lambda _raw: _telegram(second=0, timestamp="invalid-time"), ) expected_at = received_at elif path == "fingerprint": assert not ingest.handle_frame( source_id, b"private", session_factory=factory, parser=lambda _raw: _telegram(second=0, integrity=IntegrityStatus.VALID, fingerprint="invalid"), ) expected_at = received_at elif path == "candidate": assert not ingest.handle_frame( source_id, b"private", session_factory=factory, parser=lambda _raw: _telegram(second=0) ) expected_at = received_at elif path == "continuity": assert not ingest.handle_frame( source_id, b"private", session_factory=factory, parser=lambda _raw: _telegram(second=0) ) assert not ingest.handle_frame( source_id, b"private", session_factory=factory, parser=lambda _raw: _telegram(second=20, water="5.901"), ) expected_at = received_at + timedelta(minutes=1) else: assert ingest.handle_frame( source_id, b"private", session_factory=factory, parser=lambda _raw: _telegram(second=0, integrity=IntegrityStatus.VALID), ) expected_at = received_at with factory() as session: source = session.get(MeterSource, source_id) assert source is not None assert _as_utc(source.updated_at) == expected_at assert clock.calls == (2 if path == "continuity" else 1) def test_unverifiable_summer_w_frames_keep_utc_continuity(database): _engine, factory, (source_id, _) = database received_at = datetime(2026, 8, 23, 21, 40, 10, tzinfo=UTC) ingest = WarmteLinkIngestor(clock=_clock_at(received_at)) with factory() as session: assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=0, timestamp="260823234000W")) assert ingest.ingest( session, source_id=source_id, telegram=_telegram(second=10, water="5.901", timestamp="260823234010W"), ) session.commit() with factory() as session: readings = list(session.scalars(select(WarmteLinkReading).order_by(WarmteLinkReading.id))) assert [_as_utc(reading.recorded_at) for reading in readings] == [ datetime(2026, 8, 23, 21, 40, 10, tzinfo=UTC), datetime(2026, 8, 23, 21, 40, 10, tzinfo=UTC), ] def test_unverifiable_requires_two_continuous_frames_and_preserves_true_metadata(database): _engine, factory, (source_id, _) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=0)) session.commit() assert _counts(factory) == (0, 0) with factory() as session: assert ingest.ingest(session, source_id=source_id, telegram=_telegram(second=10, water="5.901")) assert ingest.ingest(session, source_id=source_id, telegram=_telegram(second=20, water="5.902")) session.commit() assert _counts(factory) == (2, 2) with factory() as session: channels = list( session.scalars( select(MeterSourceChannel) .where(MeterSourceChannel.source_id == source_id) .order_by(MeterSourceChannel.channel_key) ) ) assert [(item.channel_key, item.unit, item.device_type, item.latest_quality) for item in channels] == [ ("channel-1", "m³", "006", "unverifiable"), ("channel-2", "GJ", "012", "unverifiable"), ] assert all(item.fingerprint == "a" * 64 for item in channels) assert channels[0].latest_value == Decimal("5.902") def test_unverifiable_cadence_requires_exactly_ten_seconds_and_restarts_after_gaps(database): _engine, factory, (source_id, _) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=0)) assert ingest.ingest(session, source_id=source_id, telegram=_telegram(second=10, water="5.901")) session.commit() assert _counts(factory) == (2, 2) restarted = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not restarted.ingest(session, source_id=source_id, telegram=_telegram(second=0, water="5.902")) # 20 seconds misses one expected frame; it is only a new candidate. assert not restarted.ingest(session, source_id=source_id, telegram=_telegram(second=20, water="5.903")) assert restarted.ingest(session, source_id=source_id, telegram=_telegram(second=30, water="5.904")) # A larger gap likewise needs its own new matching successor. assert not restarted.ingest(session, source_id=source_id, telegram=_telegram(second=0, minute=2, water="5.905")) assert restarted.ingest(session, source_id=source_id, telegram=_telegram(second=10, minute=2, water="5.906")) session.commit() assert _counts(factory) == (2, 4) @pytest.mark.parametrize( "first, second", [ (_telegram(second=10), _telegram(second=9)), (_telegram(second=0), _telegram(second=30)), (_telegram(second=10), _telegram(second=20, water="5.899")), (_telegram(second=10), _telegram(second=20, device_type="007")), (_telegram(second=10), _telegram(second=20, channels=(1,))), (_telegram(second=10, channels=(1,)), _telegram(second=20)), ], ) def test_unverifiable_discontinuities_restart_two_frame_confirmation(database, first, second): _engine, factory, (source_id, _) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not ingest.ingest(session, source_id=source_id, telegram=first) assert not ingest.ingest(session, source_id=source_id, telegram=second) session.commit() assert _counts(factory) == (0, 0) with factory() as session: source = session.get(MeterSource, source_id) assert source is not None and source.status == "error" and source.last_error is not None def test_discontinuity_candidate_needs_a_new_matching_successor(database): _engine, factory, (source_id, _) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=0)) assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=10, channels=(1,))) # The restored channel set is a fresh candidate, not an admission. assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=20)) assert ingest.ingest(session, source_id=source_id, telegram=_telegram(second=30, water="5.901")) session.commit() assert _counts(factory) == (2, 2) @pytest.mark.parametrize( "first, discontinuous, successor", [ (_telegram(second=0), _telegram(second=10, water="5.899"), _telegram(second=20, water="5.900")), (_telegram(second=0), _telegram(second=10, fingerprint="b" * 64), _telegram(second=20, fingerprint="b" * 64)), (_telegram(second=0), _telegram(second=10, channels=(1,)), _telegram(second=20, channels=(1,))), (_telegram(second=0, channels=(1,)), _telegram(second=10), _telegram(second=20)), ], ) def test_each_metadata_or_channel_discontinuity_requires_its_own_two_frames( database, first, discontinuous, successor ): _engine, factory, (source_id, _) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not ingest.ingest(session, source_id=source_id, telegram=first) # The discontinuous frame becomes a new first candidate, never an admission. assert not ingest.ingest(session, source_id=source_id, telegram=discontinuous) assert ingest.ingest(session, source_id=source_id, telegram=successor) session.commit() def test_parser_failure_clears_only_its_source_candidate_and_requires_two_new_frames(database): _engine, factory, (one, two) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=0)) assert not ingest.ingest(session, source_id=two, telegram=_telegram(second=0)) session.commit() def parse_failure(_frame: bytes) -> P1Telegram: raise ValueError("raw telegram must never be reported") assert not ingest.handle_frame(one, b"private", session_factory=factory, parser=parse_failure) with factory() as session: # One's candidate was cleared; two's independent candidate survives. assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=10, water="5.901")) assert ingest.ingest(session, source_id=two, telegram=_telegram(second=10, water="5.901")) assert ingest.ingest(session, source_id=one, telegram=_telegram(second=20, water="5.902")) session.commit() assert _counts(factory) == (4, 4) def _normalization_failure(kind: str) -> P1Telegram: telegram = _telegram(second=10) if kind == "no-channels": return replace(telegram, channels=()) first, *rest = telegram.channels if kind == "empty-readings": return replace(telegram, channels=(replace(first, readings=()), *rest)) if kind == "ambiguous-readings": return replace(telegram, channels=(replace(first, readings=(first.readings[0],) * 2), *rest)) if kind == "incomplete-reading": return replace( telegram, channels=(replace(first, readings=(replace(first.readings[0], value=None),)), *rest), ) invalid_values: dict[str, object] = { "decimal-nan": Decimal("NaN"), "decimal-positive-infinity": Decimal("Infinity"), "decimal-negative-infinity": Decimal("-Infinity"), "float": 5.9, "string": "private-invalid-cumulative-value", } if kind in invalid_values: return replace( telegram, channels=( replace(first, readings=(replace(first.readings[0], value=invalid_values[kind]),)), *rest, ), ) if kind == "invalid-timestamp": return replace(telegram, timestamp="invalid-time") raise AssertionError(f"unexpected normalization failure kind: {kind}") @pytest.mark.parametrize( "kind", [ "no-channels", "empty-readings", "ambiguous-readings", "incomplete-reading", "decimal-nan", "decimal-positive-infinity", "decimal-negative-infinity", "float", "string", "invalid-timestamp", ], ) def test_dto_normalization_failure_clears_only_its_source_candidate_and_requires_two_new_frames( database, kind ): _engine, factory, (one, two) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=0)) assert not ingest.ingest(session, source_id=two, telegram=_telegram(second=0)) session.commit() failed = _normalization_failure(kind) assert not ingest.handle_frame(one, b"private", session_factory=factory, parser=lambda _raw: failed) # The rejected DTO creates no discovery/latest/history rows, and source two # retains its independent first candidate. assert _counts(factory) == (0, 0) with factory() as session: source = session.get(MeterSource, one) assert source is not None and source.last_error == "WarmteLink frame could not be normalized" assert "private" not in (source.last_error or "") assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=10, water="5.901")) assert ingest.ingest(session, source_id=two, telegram=_telegram(second=10, water="5.901")) assert ingest.ingest(session, source_id=one, telegram=_telegram(second=20, water="5.902")) session.commit() assert _counts(factory) == (4, 4) def test_valid_invalid_sampling_replay_restart_and_source_isolation(database): _engine, factory, (one, two) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert ingest.ingest(session, source_id=one, telegram=_telegram(second=0, integrity=IntegrityStatus.VALID)) assert ingest.ingest(session, source_id=one, telegram=_telegram(second=10, integrity=IntegrityStatus.VALID, water="5.901")) assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=20, integrity=IntegrityStatus.INVALID)) # Source two cannot consume source one's candidate. assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=30)) assert not ingest.ingest(session, source_id=two, telegram=_telegram(second=40)) session.commit() assert _counts(factory) == (2, 2) # candidates are never discovered; first minute only once # A new process sees the minute bucket in DB and does not create duplicates. restarted = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert restarted.ingest(session, source_id=one, telegram=_telegram(second=50, integrity=IntegrityStatus.VALID)) session.commit() assert _counts(factory) == (2, 2) def test_history_keeps_first_utc_minute_sample_and_creates_next_minute(database): _engine, factory, (source_id, _) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert ingest.ingest(session, source_id=source_id, telegram=_telegram(second=0, integrity=IntegrityStatus.VALID)) assert ingest.ingest(session, source_id=source_id, telegram=_telegram(second=10, integrity=IntegrityStatus.VALID, water="5.901")) assert ingest.ingest( session, source_id=source_id, telegram=_telegram(second=0, minute=1, integrity=IntegrityStatus.VALID, water="5.902"), ) session.commit() with factory() as session: readings = list(session.scalars(select(WarmteLinkReading).order_by(WarmteLinkReading.recorded_at))) assert len(readings) == 4 assert readings[0].recorded_at.second == 0 assert readings[2].recorded_at.minute == 1 and readings[2].recorded_at.second == 0 channel = session.scalar( select(MeterSourceChannel).where( MeterSourceChannel.source_id == source_id, MeterSourceChannel.channel_key == "channel-1", ) ) assert channel is not None and channel.latest_value == Decimal("5.902") def test_handle_frame_rolls_back_database_error_marks_source_and_recovers(database, monkeypatch): _engine, factory, (source_id, _) = database received_at = datetime(2026, 8, 22, 10, tzinfo=UTC) clock = _AdvancingClock(received_at) ingest = WarmteLinkIngestor(clock=clock) telegram = _telegram(second=0, integrity=IntegrityStatus.VALID) import app.services.warmtelink_ingest as module original = module.upsert_discovered_channel def broken(*_args, **_kwargs): raise RuntimeError("database unavailable") monkeypatch.setattr(module, "upsert_discovered_channel", broken) assert not ingest.handle_frame(source_id, b"private telegram", session_factory=factory, parser=lambda _raw: telegram) assert _counts(factory) == (0, 0) with factory() as session: source = session.get(MeterSource, source_id) assert source is not None and source.status == "error" assert "private" not in (source.last_error or "") assert _as_utc(source.updated_at) == received_at assert clock.calls == 1 monkeypatch.setattr(module, "upsert_discovered_channel", original) assert ingest.handle_frame(source_id, b"ignored", session_factory=factory, parser=lambda _raw: telegram) assert _counts(factory) == (2, 2) def test_handle_frame_commit_failure_reuses_its_receive_clock_for_error_diagnostic(database, monkeypatch): _engine, factory, (source_id, _) = database received_at = datetime(2026, 8, 22, 10, tzinfo=UTC) clock = _AdvancingClock(received_at) ingest = WarmteLinkIngestor(clock=clock) original_commit = Session.commit sessions_created = 0 def session_factory(): nonlocal sessions_created session = factory() session.info["fail_commit"] = sessions_created == 0 sessions_created += 1 return session def fail_initial_commit(self): if self.info.get("fail_commit"): raise RuntimeError("database commit failed") return original_commit(self) monkeypatch.setattr(Session, "commit", fail_initial_commit) assert not ingest.handle_frame( source_id, b"private", session_factory=session_factory, parser=lambda _raw: _telegram(second=0, integrity=IntegrityStatus.VALID), ) with factory() as session: source = session.get(MeterSource, source_id) assert source is not None assert source.status == "error" assert _as_utc(source.updated_at) == received_at assert clock.calls == 1 def test_second_channel_failure_rolls_back_first_channel_latest_and_history(database, monkeypatch): _engine, factory, (source_id, _) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) import app.services.warmtelink_ingest as module original = module.upsert_discovered_channel def broken_on_second(*args, **kwargs): if kwargs["channel_key"] == "channel-2": raise RuntimeError("second channel failed") return original(*args, **kwargs) monkeypatch.setattr(module, "upsert_discovered_channel", broken_on_second) assert not ingest.handle_frame( source_id, b"private", session_factory=factory, parser=lambda _raw: _telegram(second=0, integrity=IntegrityStatus.VALID), ) assert _counts(factory) == (0, 0) monkeypatch.setattr(module, "upsert_discovered_channel", original) assert ingest.handle_frame( source_id, b"private", session_factory=factory, parser=lambda _raw: _telegram(second=10, integrity=IntegrityStatus.VALID), ) assert _counts(factory) == (2, 2) def test_unverifiable_database_failure_clears_only_its_source_sliding_candidate(database, monkeypatch): _engine, factory, (one, two) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) import app.services.warmtelink_ingest as module with factory() as session: assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=0)) assert not ingest.ingest(session, source_id=two, telegram=_telegram(second=0)) session.commit() original = module.upsert_discovered_channel monkeypatch.setattr(module, "upsert_discovered_channel", lambda *_args, **_kwargs: (_ for _ in ()).throw(RuntimeError())) assert not ingest.handle_frame( one, b"private", session_factory=factory, parser=lambda _raw: _telegram(second=10, water="5.901") ) monkeypatch.setattr(module, "upsert_discovered_channel", original) assert _all_counts(factory) == (0, 0, 0) with factory() as session: # One needs a new pair after its failed write; two retains its own candidate. assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=20, water="5.902")) assert ingest.ingest(session, source_id=two, telegram=_telegram(second=10, water="5.901")) assert ingest.ingest(session, source_id=one, telegram=_telegram(second=30, water="5.903")) session.commit() assert _all_counts(factory) == (4, 4, 0) def _valid_crc_frame(*, device_type: str = "006") -> bytes: body = ( b"/WARMTE\r\n" b"0-0:1.0.0(260822120000S)\r\n" b"0-0:96.1.1(WARMTE-REDACTED)\r\n" + f"0-1:24.1.0({device_type})\r\n".encode() + b"0-1:96.1.0(KAM-REDACTED)\r\n" + b"0-1:24.2.1(260822120000S)(5.900*m3)\r\n" + b"0-2:24.1.0(012)\r\n" + b"0-2:96.1.0(KAM-REDACTED)\r\n" + b"0-2:24.2.1(260822120000S)(0.017*GJ)\r\n" ) payload = body + b"!" return payload + f"{dsmr_crc16(payload):04X}".encode() + b"\r\n" def test_non_enum_integrity_and_fake_metadata_reject_whole_frame_per_source(database): _engine, factory, (one, two) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=0)) assert not ingest.ingest(session, source_id=two, telegram=_telegram(second=0)) session.commit() forged_quality = replace(_telegram(second=10, integrity=IntegrityStatus.VALID), integrity=SimpleNamespace(value="valid")) assert not ingest.handle_frame(one, b"private", session_factory=factory, parser=lambda _raw: forged_quality) mixed = replace(_telegram(second=10), channels=(replace(_telegram(second=10).channels[0], device_type="raw-device"),)) assert not ingest.handle_frame(one, b"private", session_factory=factory, parser=lambda _raw: mixed) assert _all_counts(factory) == (0, 0, 0) with factory() as session: source = session.get(MeterSource, one) assert source is not None and source.last_error == "WarmteLink frame could not be normalized" # Both malformed frames clear one, while two's candidate is untouched. assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=20, water="5.901")) assert ingest.ingest(session, source_id=two, telegram=_telegram(second=10, water="5.901")) assert ingest.ingest(session, source_id=one, telegram=_telegram(second=30, water="5.902")) session.commit() assert _all_counts(factory) == (4, 4, 0) def test_valid_crc_parser_frame_with_raw_device_type_is_rejected_without_leakage(database): _engine, factory, (one, two) = database ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=0)) assert not ingest.ingest(session, source_id=two, telegram=_telegram(second=0)) session.commit() raw_identifier = "serial-number-which-must-not-persist" assert parse_telegram(_valid_crc_frame(device_type=raw_identifier)).integrity is IntegrityStatus.VALID assert not ingest.handle_frame(one, _valid_crc_frame(device_type=raw_identifier), session_factory=factory) assert _all_counts(factory) == (0, 0, 0) with factory() as session: source = session.get(MeterSource, one) assert source is not None and source.last_error == "WarmteLink frame could not be normalized" assert raw_identifier not in (source.last_error or "") assert not ingest.ingest(session, source_id=one, telegram=_telegram(second=10, water="5.901")) assert ingest.ingest(session, source_id=two, telegram=_telegram(second=10, water="5.901")) assert ingest.ingest(session, source_id=one, telegram=_telegram(second=20, water="5.902")) session.commit() assert _all_counts(factory) == (4, 4, 0) @pytest.mark.parametrize( "fingerprint, channel_fingerprint", [ (None, None), ("", ""), ("a" * 63, "a" * 63), ("g" * 64, "g" * 64), ("serial-number-which-must-not-persist", "serial-number-which-must-not-persist"), ], ) def test_invalid_fingerprint_fails_closed_without_persisting_identifiers( database, fingerprint, channel_fingerprint ): _engine, factory, (source_id, _) = database raw_identifier = "serial-number-which-must-not-persist" ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=0, integrity=IntegrityStatus.INVALID)) assert not ingest.ingest( session, source_id=source_id, telegram=_telegram( second=10, integrity=IntegrityStatus.VALID, fingerprint=fingerprint, channel_fingerprint=channel_fingerprint, ), ) session.commit() assert _counts(factory) == (0, 0) with factory() as session: source = session.get(MeterSource, source_id) assert source is not None and source.last_error == "WarmteLink frame fingerprint is invalid" assert raw_identifier not in (source.last_error or "") assert ingest.ingest(session, source_id=source_id, telegram=_telegram(second=20, integrity=IntegrityStatus.VALID)) session.commit() rendered = " ".join(str(row) for row in session.execute(select(WarmteLinkReading)).all()) assert raw_identifier not in rendered assert all("raw" not in column.name for column in WarmteLinkReading.__table__.columns) @pytest.mark.parametrize( "invalid_fingerprint", [None, "", "a" * 63, "serial-number-which-must-not-persist"], ids=["none", "empty", "wrong-length", "raw-like"], ) @pytest.mark.parametrize("invalid_channel", [1, 2]) def test_invalid_channel_fingerprint_rejects_whole_frame_and_clears_candidate( database, invalid_fingerprint, invalid_channel ): _engine, factory, (source_id, _) = database raw_identifier = "serial-number-which-must-not-persist" valid = "a" * 64 ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: # Establish a candidate, then prove an invalid channel clears only it. assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=0)) assert not ingest.ingest( session, source_id=source_id, telegram=_telegram( second=10, fingerprint=valid, channel_fingerprints={1: valid, 2: valid} | {invalid_channel: invalid_fingerprint}, ), ) session.commit() assert _counts(factory) == (0, 0) with factory() as session: source = session.get(MeterSource, source_id) assert source is not None and source.last_error == "WarmteLink frame fingerprint is invalid" assert raw_identifier not in (source.last_error or "") # The first legal frame is a fresh candidate; only its successor admits. assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=20, water="5.901")) assert ingest.ingest(session, source_id=source_id, telegram=_telegram(second=30, water="5.902")) session.commit() assert _counts(factory) == (2, 2) @pytest.mark.parametrize( "invalid_fingerprint", [None, "", "a" * 63, "serial-number-which-must-not-persist"], ids=["none", "empty", "wrong-length", "raw-like"], ) def test_invalid_top_fingerprint_rejects_whole_frame_and_clears_candidate(database, invalid_fingerprint): _engine, factory, (source_id, _) = database raw_identifier = "serial-number-which-must-not-persist" valid = "a" * 64 ingest = WarmteLinkIngestor(clock=_fixed_clock) with factory() as session: assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=0)) assert not ingest.ingest( session, source_id=source_id, telegram=_telegram( second=10, fingerprint=invalid_fingerprint, channel_fingerprints={1: valid, 2: valid}, ), ) session.commit() assert _counts(factory) == (0, 0) with factory() as session: source = session.get(MeterSource, source_id) assert source is not None and source.last_error == "WarmteLink frame fingerprint is invalid" assert raw_identifier not in (source.last_error or "") assert not ingest.ingest(session, source_id=source_id, telegram=_telegram(second=20, water="5.901")) assert ingest.ingest(session, source_id=source_id, telegram=_telegram(second=30, water="5.902")) session.commit() assert _counts(factory) == (2, 2)