"""In-memory model of discovered OMS devices. Reset on every app start.""" from __future__ import annotations import statistics import time from collections import deque from dataclasses import dataclass, field from typing import Deque, Dict, List, Optional, Tuple from .manufacturers import manufacturer_name @dataclass class Device: id: str manufacturer: str = "" # 3-letter FLAG code, e.g. "KAM" media: str = "" # e.g. "Cold water meter" version: str = "" # hex version byte driver: str = "" # wmbusmeters driver, e.g. "kamwater" c_field: str = "" # DLL C field (technical) rssi: Optional[int] = None # encryption / link-layer detail (from ELL / TPL verbose lines) enc: str = "" # "AES-CTR", "AES-CBC", or "" (none) enc_layer: str = "" # "ELL" or "TPL" sec_mode: int = 0 # TPL security mode number (5, 7, 0=none) tpl_ci: str = "" # TPL CI byte ell_ci: str = "" # ELL CI byte ell_sn: str = "" # ELL session number / SN enc_blocks: str = "" # number of encrypted TPL blocks (nb=) first_seen: float = field(default_factory=time.time) last_seen: float = field(default_factory=time.time) count: int = 0 key: Optional[str] = None # AES-128 key (32 hex chars) if configured encrypted: bool = False # has any telegram used encryption? (sticky) key_decrypted: bool = False # have we decoded an *encrypted* frame with our key? decrypted: bool = False # do we have readable values (any source)? fields_source: str = "" # "key" (real decrypt) or "plain" (unencrypted read) fields: Dict[str, object] = field(default_factory=dict) # latest decoded values last_decoded: Optional[float] = None # when `fields` was last refreshed # recent telegrams as (hh:mm:ss, event) for the live stream telegrams: Deque[Tuple[str, dict]] = field( default_factory=lambda: deque(maxlen=500) ) _last_count_at: float = 0.0 # for de-duplicating multi-meter DLL bursts # arrival times of counted telegrams, for estimating the send interval rx_times: Deque[float] = field(default_factory=lambda: deque(maxlen=64)) @property def intervals(self) -> List[float]: t = list(self.rx_times) return [b - a for a, b in zip(t, t[1:])] @property def interval_samples(self) -> int: return max(0, len(self.rx_times) - 1) @property def interval_estimate(self) -> Optional[float]: """Best estimate of the transmit period: median of observed gaps. Median is robust to missed telegrams (which show up as gaps that are multiples of the true period) and to the odd retransmission burst. """ gaps = self.intervals return statistics.median(gaps) if gaps else None @property def interval_last(self) -> Optional[float]: gaps = self.intervals return gaps[-1] if gaps else None @property def interval_min(self) -> Optional[float]: gaps = self.intervals return min(gaps) if gaps else None @property def encryption(self) -> str: """Human summary of the encryption in use, incl. the OMS security mode.""" if self.enc_layer == "ELL" and self.enc: return f"{self.enc} (ELL)" if self.sec_mode: return f"mode {self.sec_mode} ({self.enc})" if self.enc: return self.enc return "none" @property def manufacturer_name(self) -> str: return manufacturer_name(self.manufacturer) @property def status(self) -> str: # Encryption status is authoritative (from TPL/ELL), not "do we have data". if self.key: return "decrypted" if self.key_decrypted else "key set" if self.encrypted: return "locked" # encrypted frames we can't read yet if self.decrypted: return "open" # unencrypted meter, readable as-is return "listening" # seen, not yet classified @property def status_icon(self) -> str: return { "decrypted": "🔓", "open": "📖", "key set": "🔑", "locked": "🔒", "listening": "📡", }[self.status] @property def status_color(self) -> str: return { "decrypted": "green", "open": "green", "key set": "yellow", "locked": "red", "listening": "cyan", }[self.status] # A single physical telegram can be logged by several matching meters (a keyed # meter *and* the wildcard), producing back-to-back DLL lines microseconds # apart. Real re-transmissions from a meter are seconds apart, so we treat # same-device DLLs within this window as one telegram. COUNT_DEDUP_WINDOW = 0.5 class DeviceStore: """Devices keyed by id, preserving discovery order.""" def __init__(self) -> None: self.devices: Dict[str, Device] = {} self.order: List[str] = [] def __len__(self) -> int: return len(self.order) def get(self, device_id: str) -> Optional[Device]: return self.devices.get(device_id) def ordered(self) -> List[Device]: return [self.devices[i] for i in self.order] @staticmethod def _note_telegram(dev: Device, now: float) -> None: """Count one telegram and record its arrival, de-duplicating the multi-line burst a single telegram produces (DLL[+DLL]+ELL+TPL+JSON, all within microseconds) so counts and the interval stay accurate. Driven by *any* telegram-level event, so a DLL line we fail to parse can't silently zero out the count and send-interval. """ if dev._last_count_at == 0.0 or (now - dev._last_count_at) > COUNT_DEDUP_WINDOW: dev.count += 1 dev.rx_times.append(now) dev._last_count_at = now def _get_or_create(self, device_id: str, now: float) -> Tuple[Device, bool]: dev = self.devices.get(device_id) if dev is None: dev = Device(id=device_id, first_seen=now, last_seen=now) self.devices[device_id] = dev self.order.append(device_id) return dev, True return dev, False def upsert_discovery(self, ev: dict, now: Optional[float] = None) -> Tuple[Device, bool]: """Apply a DLL discovery event. Returns (device, is_new). Counts the telegram.""" now = time.time() if now is None else now dev, is_new = self._get_or_create(ev["id"], now) dev.manufacturer = ev.get("manufacturer") or dev.manufacturer dev.media = ev.get("media") or dev.media dev.version = ev.get("version") or dev.version dev.driver = ev.get("driver") or dev.driver dev.c_field = ev.get("c_field") or dev.c_field if ev.get("rssi") is not None: dev.rssi = ev["rssi"] dev.last_seen = now self._note_telegram(dev, now) return dev, is_new def apply_ell(self, ev: dict, device_id: Optional[str]) -> Optional[Device]: """Apply an ELL (extended link layer) event to the current telegram's device.""" if not device_id: return None dev = self.devices.get(device_id) if dev is None: return None dev.ell_ci = ev.get("ci") or dev.ell_ci dev.ell_sn = ev.get("session") or ev.get("sn") or dev.ell_sn if ev.get("enc"): dev.enc = ev["enc"] dev.enc_layer = "ELL" dev.encrypted = True return dev def apply_tpl(self, ev: dict, device_id: Optional[str]) -> Optional[Device]: """Apply a TPL (transport layer) event to the current telegram's device.""" if not device_id: return None dev = self.devices.get(device_id) if dev is None: return None dev.tpl_ci = ev.get("ci") or dev.tpl_ci # Only let the TPL override encryption if ELL didn't already claim it. if ev.get("sec_mode"): if dev.enc_layer != "ELL": dev.enc = ev.get("enc") or dev.enc dev.enc_layer = "TPL" dev.sec_mode = ev["sec_mode"] dev.enc_blocks = ev.get("nb") or dev.enc_blocks dev.encrypted = True return dev def apply_json(self, ev: dict, now: Optional[float] = None) -> Tuple[Device, bool]: """Apply a decrypted JSON telegram event. Does not count (the DLL line does).""" now = time.time() if now is None else now dev, is_new = self._get_or_create(ev["id"], now) dev.driver = ev.get("driver") or dev.driver dev.media = ev.get("media") or dev.media if ev.get("fields"): # The wildcard meter is named "scan"; a keyed meter is "m_". # Only the latter means we actually decrypted an encrypted frame. name = ev.get("name") or "" from_key = bool(name) and name != "scan" dev.fields = dict(ev["fields"]) dev.fields_source = "key" if from_key else "plain" dev.decrypted = True dev.last_decoded = now if from_key: dev.key_decrypted = True dev.last_seen = now # Also drive the counter from JSON, so telegrams still count even if # their DLL line failed to parse. The dedup window collapses the DLL # and its JSON (same telegram) into one. self._note_telegram(dev, now) return dev, is_new