"""OMS / wireless M-Bus monitor — a Textual TUI around wmbusmeters. Left: live list of detected OMS devices. Right: details of the selected device (or a live stream of all telegrams when nothing is selected). Keys can be set live and the list exported/imported as CSV (AES keys included) for auto-decrypt. """ from __future__ import annotations import argparse import os import time from collections import deque from typing import Deque, Optional, Tuple from rich import box from rich.console import Group from rich.table import Table from rich.text import Text from textual.app import App, ComposeResult from textual.binding import Binding from textual.containers import Horizontal, Vertical, VerticalScroll from textual.message import Message from textual.screen import ModalScreen from textual.widgets import DataTable, Footer, Header, Input, Label, RichLog, Static from .csvio import export_devices, import_rows from .models import Device, DeviceStore from .wmbus import WMBusSource HEX = set("0123456789ABCDEF") def rel_time(ts: float) -> str: d = max(0, int(time.time() - ts)) if d < 1: return "now" if d < 60: return f"{d}s ago" if d < 3600: return f"{d // 60}m {d % 60}s ago" return f"{d // 3600}h {(d % 3600) // 60}m ago" def fmt_duration(seconds: Optional[float]) -> str: if seconds is None: return "?" s = int(round(seconds)) if s < 60: return f"{s}s" if s < 3600: return f"{s // 60}m {s % 60:02d}s" if s < 86400: return f"{s // 3600}h {(s % 3600) // 60:02d}m" return f"{s // 86400}d {(s % 86400) // 3600:02d}h" def stream_text(ts_str: str, ev: dict) -> Text: if ev.get("type") == "json": fields = ev.get("fields", {}) preview = " ".join(f"{k}={v}" for k, v in list(fields.items())[:5]) line = Text() line.append(f"{ts_str} ", style="dim") line.append("🔓 ", style="green") line.append(f"{ev['id']} ", style="bold green") line.append(f"{ev.get('driver', '')} ", style="cyan") line.append(preview) return line line = Text() line.append(f"{ts_str} ", style="dim") line.append("📡 ", style="yellow") line.append(f"{ev['id']} ", style="bold") line.append(f"{ev.get('manufacturer', '')} ", style="magenta") line.append(f"{ev.get('media', '')}", style="") rssi = ev.get("rssi") if rssi is not None: line.append(f" RSSI={rssi}", style="dim") return line class TelegramEvent(Message): """A parsed telegram event handed off from the wmbusmeters reader.""" def __init__(self, ev: dict) -> None: self.ev = ev super().__init__() class PromptScreen(ModalScreen[Optional[str]]): """A single-line modal prompt. Returns the entered string, or None on cancel.""" BINDINGS = [("escape", "cancel", "Cancel")] def __init__(self, prompt: str, default: str = "") -> None: super().__init__() self._prompt = prompt self._default = default def compose(self) -> ComposeResult: with Vertical(id="dialog"): yield Label(self._prompt) yield Input(value=self._default, id="prompt-input") def on_mount(self) -> None: self.query_one(Input).focus() def on_input_submitted(self, event: Input.Submitted) -> None: self.dismiss(event.value) def action_cancel(self) -> None: self.dismiss(None) class OmsApp(App[None]): CSS_PATH = "app.tcss" TITLE = "OMS Monitor" SUB_TITLE = "wireless M-Bus smart meters" BINDINGS = [ ("k", "set_key", "Set AES key"), ("e", "export", "Export CSV"), ("i", "import", "Import CSV"), ("a", "show_all", "Show all / stream"), ("enter", "select", "Select device"), ("q", "quit", "Quit"), # priority so Ctrl+C quits even while a modal Input is focused Binding("ctrl+c", "quit", "Quit", priority=True, show=False), ] COLUMNS = [ ("id", "ID"), ("manuf", "Manuf"), ("media", "Media"), ("last", "Last seen"), ("msgs", "Msgs"), ("st", ""), ] def __init__( self, device: str = "rtlwmbus", modes: str = "c1,t1", wmbusmeters: str = "wmbusmeters", stdin_data: Optional[bytes] = None, import_csv: Optional[str] = None, autostart: bool = True, ) -> None: super().__init__() self.store = DeviceStore() self.selected_id: Optional[str] = None self._current_tid: Optional[str] = None # id of the telegram being parsed self.show_all = True self.global_stream: Deque[Tuple[str, dict]] = deque(maxlen=1000) self._autostart = autostart and device != "none" self._import_csv = import_csv self.source = WMBusSource( device=device, modes=modes, wmbusmeters=wmbusmeters, stdin_data=stdin_data ) # ---- layout ----------------------------------------------------------- def compose(self) -> ComposeResult: yield Header(show_clock=True) with Horizontal(id="body"): with Vertical(id="left"): yield Static("OMS devices (0)", id="left-title") yield DataTable(id="devices", zebra_stripes=True, cursor_type="row") with Vertical(id="right"): with VerticalScroll(id="detail-scroll"): yield Static(id="detail") yield Static("Live stream", id="stream-title") yield RichLog(id="stream", wrap=True, highlight=False, markup=False) yield Footer() def on_mount(self) -> None: table = self.query_one("#devices", DataTable) for key, label in self.COLUMNS: table.add_column(label, key=key) self._refresh_detail() self.set_interval(1.0, self._tick) if self._import_csv: self._do_import(self._import_csv, announce=False) if self._autostart: self.run_worker(self._start_source(), exclusive=False) async def _start_source(self) -> None: keys = {d.id: d.key for d in self.store.devices.values() if d.key} self.source.set_keys(keys) self.source.stdin_data = self.source.stdin_data # keep replay if any await self.source.start(lambda ev: self.post_message(TelegramEvent(ev))) async def on_unmount(self) -> None: await self.source.stop() # ---- event handling --------------------------------------------------- def on_telegram_event(self, message: TelegramEvent) -> None: ev = message.ev etype = ev.get("type") # ELL/TPL lines have no id of their own; they belong to the telegram # whose DLL line preceded them. if etype in ("ell", "tpl"): apply = self.store.apply_ell if etype == "ell" else self.store.apply_tpl dev = apply(ev, self._current_tid) if dev is not None: self._upsert_row(dev) # encryption may change the status icon if self.selected_id == dev.id: self._refresh_detail() return if etype not in ("dll", "json"): return now = time.time() ts_str = time.strftime("%H:%M:%S", time.localtime(now)) if etype == "dll": self._current_tid = ev["id"] dev, _ = self.store.upsert_discovery(ev, now) else: dev, _ = self.store.apply_json(ev, now) item = (ts_str, ev) self.global_stream.append(item) dev.telegrams.append(item) self._upsert_row(dev) self.query_one("#left-title", Static).update(f"OMS devices ({len(self.store)})") stream = self.query_one("#stream", RichLog) if self.show_all: stream.write(stream_text(ts_str, ev)) elif self.selected_id == dev.id: stream.write(stream_text(ts_str, ev)) self._refresh_detail() def on_data_table_row_highlighted(self, event: DataTable.RowHighlighted) -> None: row_key = getattr(event.row_key, "value", None) or str(event.row_key) if not self.show_all: self.selected_id = row_key self._refresh_detail() self._repaint_stream() def on_data_table_row_selected(self, event: DataTable.RowSelected) -> None: row_key = getattr(event.row_key, "value", None) or str(event.row_key) self.show_all = False self.selected_id = row_key self._update_stream_title() self._refresh_detail() self._repaint_stream() # ---- table ------------------------------------------------------------ def _row_cells(self, dev: Device): return ( dev.id, dev.manufacturer or "?", dev.media or "?", rel_time(dev.last_seen), str(dev.count), dev.status_icon, ) def _upsert_row(self, dev: Device) -> None: table = self.query_one("#devices", DataTable) cells = self._row_cells(dev) if dev.id in {getattr(rk, "value", rk) for rk in table.rows}: for (col_key, _), value in zip(self.COLUMNS, cells): table.update_cell(dev.id, col_key, value) else: table.add_row(*cells, key=dev.id) def _tick(self) -> None: table = self.query_one("#devices", DataTable) existing = {getattr(rk, "value", rk) for rk in table.rows} for dev in self.store.ordered(): if dev.id in existing: table.update_cell(dev.id, "last", rel_time(dev.last_seen)) if not self.show_all and self.selected_id: self._refresh_detail() # ---- detail panel ----------------------------------------------------- def _update_stream_title(self) -> None: title = self.query_one("#stream-title", Static) if self.show_all or not self.selected_id: title.update("Live stream — all devices") else: title.update(f"Live stream — {self.selected_id}") def _refresh_detail(self) -> None: detail = self.query_one("#detail", Static) if self.show_all or not self.selected_id: detail.update(self._summary_renderable()) return dev = self.store.get(self.selected_id) if dev is None: detail.update(self._summary_renderable()) return detail.update(self._device_renderable(dev)) def _summary_renderable(self): total = len(self.store) decrypted = sum(1 for d in self.store.devices.values() if d.decrypted) keyed = sum(1 for d in self.store.devices.values() if d.key) head = Text() head.append("No device selected\n", style="bold") head.append("Showing the live stream of all telegrams on the right.\n\n", style="dim") head.append(f"Detected: {total} ", style="") head.append(f"keyed: {keyed} ", style="cyan") head.append(f"decrypted: {decrypted}\n\n", style="green") head.append("Enter", style="bold") head.append(" select · ", style="dim") head.append("k", style="bold") head.append(" set key · ", style="dim") head.append("e", style="bold") head.append(" export · ", style="dim") head.append("i", style="bold") head.append(" import · ", style="dim") head.append("a", style="bold") head.append(" back to stream", style="dim") return head def _device_renderable(self, dev: Device): header = Text() header.append(f"{dev.status_icon} {dev.id}", style="bold") header.append(f" {dev.status}", style=dev.status_color) meta = Table.grid(padding=(0, 2)) meta.add_column(style="dim", justify="right") meta.add_column() mfct = dev.manufacturer_name if dev.manufacturer and mfct != dev.manufacturer: mfct = f"{mfct} ({dev.manufacturer})" meta.add_row("manufacturer", mfct or "?") meta.add_row("media / type", dev.media or "?") meta.add_row("driver", dev.driver or "?") meta.add_row("version", dev.version or "?") # encryption / security enc = dev.encryption enc_style = "red" if enc == "none" else "magenta" enc_text = Text(enc, style=enc_style) if dev.enc_layer: enc_text.append(f" layer={dev.enc_layer}", style="dim") if dev.enc_blocks: enc_text.append(f" blocks={dev.enc_blocks}", style="dim") meta.add_row("encryption", enc_text) if dev.sec_mode: meta.add_row("security mode", str(dev.sec_mode)) if dev.tpl_ci: meta.add_row("TPL CI", dev.tpl_ci) if dev.ell_ci: meta.add_row("ELL CI", f"{dev.ell_ci} session={dev.ell_sn or '?'}") meta.add_row("link (C)", dev.c_field or "?") meta.add_row("rssi", "?" if dev.rssi is None else f"{dev.rssi} dBm") meta.add_row("telegrams", str(dev.count)) if dev.interval_samples >= 1: est = fmt_duration(dev.interval_estimate) last = fmt_duration(dev.interval_last) detail = f"~{est}" extra = f"last {last}, {dev.interval_samples} sample" extra += "s" if dev.interval_samples != 1 else "" interval_text = Text(detail, style="bold") interval_text.append(f" ({extra})", style="dim") meta.add_row("send interval", interval_text) else: meta.add_row("send interval", Text("— (need ≥2 telegrams)", style="dim")) meta.add_row( "first seen", time.strftime("%H:%M:%S", time.localtime(dev.first_seen)) ) meta.add_row("last seen", rel_time(dev.last_seen)) if dev.encrypted or dev.key: meta.add_row("AES key", dev.key or "— (press k to set)") else: meta.add_row("AES key", "— (not needed, unencrypted)") parts = [header, Text(""), meta] if dev.fields: when = ( time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(dev.last_decoded)) if dev.last_decoded else "?" ) rel = rel_time(dev.last_decoded) if dev.last_decoded else "?" # Distinguish a real key-decrypt from a plaintext read of an # (otherwise) encrypted meter, so we never claim more than we know. partial = dev.encrypted and dev.fields_source != "key" if partial: title = f"Plaintext records (unencrypted part) — {when} ({rel})" border = "yellow" else: title = f"Latest decoded values — {when} ({rel})" border = "green" if dev.status in ("open", "decrypted") else "blue" values = Table( title=title, title_style="bold", title_justify="left", expand=True, box=box.ROUNDED, border_style=border, ) values.add_column("field", style="cyan") values.add_column("value", justify="right", style="bold") for key, value in dev.fields.items(): values.add_row(key, str(value)) parts += [Text(""), values] if partial: parts.append( Text("Some records are encrypted — press k to add the AES key.", style="yellow") ) elif dev.key: parts.append( Text("\nKey set — waiting for the next telegram to decrypt…", style="yellow") ) elif dev.encrypted: parts.append( Text("\nEncrypted. Press k to add this meter's AES key to decrypt.", style="dim") ) else: parts.append( Text("\nListening… waiting to decode a telegram from this meter.", style="dim") ) return Group(*parts) def _repaint_stream(self) -> None: stream = self.query_one("#stream", RichLog) stream.clear() if self.show_all or not self.selected_id: source = self.global_stream else: dev = self.store.get(self.selected_id) source = dev.telegrams if dev else deque() for ts_str, ev in source: stream.write(stream_text(ts_str, ev)) # ---- actions ---------------------------------------------------------- def action_show_all(self) -> None: self.show_all = True self.selected_id = None self._update_stream_title() self._refresh_detail() self._repaint_stream() def action_select(self) -> None: table = self.query_one("#devices", DataTable) if table.row_count == 0: return try: row_key = table.coordinate_to_cell_key(table.cursor_coordinate).row_key except Exception: return self.show_all = False self.selected_id = getattr(row_key, "value", None) or str(row_key) self._update_stream_title() self._refresh_detail() self._repaint_stream() def action_set_key(self) -> None: if self.show_all or not self.selected_id: self.notify("Select a device first (Enter).", severity="warning") return dev = self.store.get(self.selected_id) if dev is None: return def handle(value: Optional[str]) -> None: if value is None: return v = value.strip().upper() if v == "": dev.key = None elif len(v) == 32 and all(c in HEX for c in v): dev.key = v dev.decrypted = False else: self.notify( "Key must be 32 hex chars (or empty to clear).", severity="error" ) return self._upsert_row(dev) self._refresh_detail() self.notify(f"Key {'set' if dev.key else 'cleared'} for {dev.id}. Restarting…") self.run_worker(self._apply_keys(), exclusive=True) self.push_screen( PromptScreen( f"AES-128 key for {dev.id} (32 hex chars, empty to clear):", dev.key or "", ), handle, ) async def _apply_keys(self) -> None: keys = {d.id: d.key for d in self.store.devices.values() if d.key} self.source.set_keys(keys) if self.source.running: await self.source.restart() def action_export(self) -> None: def handle(path: Optional[str]) -> None: if not path: return try: n = export_devices(path, self.store.ordered()) except OSError as exc: self.notify(f"Export failed: {exc}", severity="error") return self.notify(f"Exported {n} devices to {path}") self.push_screen(PromptScreen("Export CSV to path:", "oms_devices.csv"), handle) def action_import(self) -> None: def handle(path: Optional[str]) -> None: if not path: return self._do_import(path) self.push_screen(PromptScreen("Import CSV from path:", "oms_devices.csv"), handle) def _do_import(self, path: str, announce: bool = True) -> None: try: rows = import_rows(path) except OSError as exc: if announce: self.notify(f"Import failed: {exc}", severity="error") return keys_loaded = 0 for row in rows: device_id = row["id"] dev = self.store.get(device_id) if dev is None: dev = Device(id=device_id, count=0) self.store.devices[device_id] = dev self.store.order.append(device_id) dev.manufacturer = dev.manufacturer or row.get("manufacturer", "") dev.media = dev.media or row.get("media", "") dev.driver = dev.driver or row.get("driver", "") dev.version = dev.version or row.get("version", "") key = row.get("aes_key", "") if key: dev.key = key.upper() keys_loaded += 1 if self.is_running: self._upsert_row(dev) if self.is_running: self.query_one("#left-title", Static).update( f"OMS devices ({len(self.store)})" ) self._refresh_detail() if announce: self.notify(f"Imported {keys_loaded} keys from {path}. Restarting…") self.run_worker(self._apply_keys(), exclusive=True) def main() -> None: parser = argparse.ArgumentParser( prog="oms-tui", description="Monitor OMS / wireless M-Bus smart meters." ) parser.add_argument( "--device", default=os.environ.get("OMS_WMBUS_DEVICE", "rtlwmbus"), help="wmbusmeters device (default: rtlwmbus; drives the RTL-SDR)", ) parser.add_argument( "--listento", default=os.environ.get("OMS_WMBUS_MODES", "c1,t1"), help="link modes to listen for (default: c1,t1)", ) parser.add_argument( "--wmbusmeters", default="wmbusmeters", help="path to the wmbusmeters binary" ) parser.add_argument( "--replay", metavar="FILE", help="replay an rtlwmbus-format capture instead of using the radio", ) parser.add_argument( "--import", dest="import_csv", metavar="FILE", help="preload AES keys from a CSV before listening", ) args = parser.parse_args() stdin_data = None device = args.device if args.replay: with open(args.replay, "rb") as f: stdin_data = f.read() device = "stdin:rtlwmbus" app = OmsApp( device=device, modes=args.listento, wmbusmeters=args.wmbusmeters, stdin_data=stdin_data, import_csv=args.import_csv, ) app.run() if __name__ == "__main__": main()