diff --git a/flake.lock b/flake.lock new file mode 100644 index 0000000..0b9301e --- /dev/null +++ b/flake.lock @@ -0,0 +1,27 @@ +{ + "nodes": { + "nixpkgs": { + "locked": { + "lastModified": 1784796856, + "narHash": "sha256-wWFrV5/Qbm+lyt5x20E/bSbfJiGKMo4RCxZV8cl/WZI=", + "owner": "NixOS", + "repo": "nixpkgs", + "rev": "e2587caef70cea85dd97d7daab492899902dbf5d", + "type": "github" + }, + "original": { + "owner": "NixOS", + "ref": "nixos-unstable", + "repo": "nixpkgs", + "type": "github" + } + }, + "root": { + "inputs": { + "nixpkgs": "nixpkgs" + } + } + }, + "root": "root", + "version": 7 +} diff --git a/flake.nix b/flake.nix new file mode 100644 index 0000000..445d6df --- /dev/null +++ b/flake.nix @@ -0,0 +1,153 @@ +{ + description = "SDR toolbox — programs for exploring an RTL-SDR stick"; + + inputs = { + nixpkgs.url = "github:NixOS/nixpkgs/nixos-unstable"; + }; + + outputs = { self, nixpkgs }: + let + system = "x86_64-linux"; + pkgs = nixpkgs.legacyPackages.${system}; + lib = pkgs.lib; + + # rtl_wmbus: demodulator that turns raw rtl_sdr I/Q into wM-Bus frames. + # Not in nixpkgs, so build from source (tiny pure-C project). + rtl_wmbus = pkgs.stdenv.mkDerivation { + pname = "rtl-wmbus"; + version = "1.1.0"; + src = pkgs.fetchFromGitHub { + owner = "xaelsouth"; + repo = "rtl-wmbus"; + rev = "1.1.0"; + sha256 = "1f9psbakvi0ymqmlbbnn96cbvd3bfphc90hw6xwid6r184pnzac6"; + }; + enableParallelBuilding = true; + installPhase = '' + runHook preInstall + install -Dm755 build/rtl_wmbus $out/bin/rtl_wmbus + runHook postInstall + ''; + meta.description = "Software defined receiver for wireless M-Bus with RTL-SDR"; + }; + + # wmbusmeters: reads/decodes wM-Bus (OMS) smart-meter telegrams. Dropped + # from nixpkgs, so build from source. At runtime it shells out to + # rtl_sdr | rtl_wmbus, so wrap those onto its PATH. + wmbusmeters = pkgs.stdenv.mkDerivation { + pname = "wmbusmeters"; + version = "3.0.0"; + src = pkgs.fetchFromGitHub { + owner = "wmbusmeters"; + repo = "wmbusmeters"; + rev = "3.0.0"; + sha256 = "1nnirlnb12zdadsvsybj8h2g6fy7hc7ih488i2qaxzinm1faxagp"; + }; + nativeBuildInputs = [ pkgs.pkg-config pkgs.makeWrapper ]; + buildInputs = [ pkgs.rtl-sdr pkgs.libusb1 pkgs.libxml2 ]; + enableParallelBuilding = true; + postPatch = '' + patchShebangs scripts + ''; + installPhase = '' + runHook preInstall + install -Dm755 build/wmbusmeters $out/bin/wmbusmeters + ln -s wmbusmeters $out/bin/wmbusmetersd + wrapProgram $out/bin/wmbusmeters \ + --prefix PATH : ${lib.makeBinPath [ pkgs.rtl-sdr rtl_wmbus ]} + runHook postInstall + ''; + meta.description = "Read wM-Bus / OMS smart meters (water, heat, gas, electricity)"; + }; + + # oms-tui: full-screen TUI for monitoring OMS / wM-Bus meters. Wraps + # wmbusmeters (put on PATH) and adds a live device list, detail view, + # live AES-key entry, and CSV export/import. + oms-tui = pkgs.python3Packages.buildPythonApplication { + pname = "oms-tui"; + version = "0.1.0"; + pyproject = true; + src = ./oms-tui; + build-system = [ pkgs.python3Packages.setuptools ]; + dependencies = [ pkgs.python3Packages.textual ]; + nativeBuildInputs = [ pkgs.makeWrapper ]; + nativeCheckInputs = [ + pkgs.python3Packages.pytestCheckHook + pkgs.python3Packages.pytest-asyncio + wmbusmeters # exercises the real decode pipeline in the test suite + ]; + postFixup = '' + wrapProgram $out/bin/oms-tui \ + --prefix PATH : ${lib.makeBinPath [ wmbusmeters ]} + ''; + meta.description = "TUI monitor for OMS / wireless M-Bus smart meters"; + }; + + # wmbus-listen: convenience wrapper around wmbusmeters for OMS discovery. + wmbus-listen = pkgs.writeShellApplication { + name = "wmbus-listen"; + runtimeInputs = [ wmbusmeters ]; + text = '' + # Listen for wM-Bus / OMS smart-meter telegrams on 868.95 MHz via RTL-SDR. + # + # Usage: + # wmbus-listen # discovery: show every telegram heard + # wmbus-listen NAME DRIVER ID KEY # decode one meter's readings + # e.g. wmbus-listen Water multical21 12345678 00112233445566778899AABBCCDDEEFF + # + # Env overrides: + # WMBUS_MODES link modes to listen for (default: c1,t1 — typical OMS) + # WMBUS_FORMAT output format: json | fields | hr (default: json) + modes="''${WMBUS_MODES:-c1,t1}" + format="''${WMBUS_FORMAT:-json}" + if [ "$#" -eq 0 ]; then + echo "Listening for OMS/wM-Bus telegrams (modes: $modes) on 868.95 MHz. Ctrl-C to stop." >&2 + echo "Unknown meters are printed as they are heard — note their id/driver to decode them." >&2 + exec wmbusmeters --format="$format" --listento="$modes" --verbose rtlwmbus + else + exec wmbusmeters --format="$format" --listento="$modes" rtlwmbus "$@" + fi + ''; + }; + + sdrTools = (with pkgs; [ + rtl-sdr # rtl_test, rtl_fm, rtl_power — verify the stick works first + sdrpp # modern spectrum browser with waterfall (run: sdrpp) + gqrx # classic beginner-friendly waterfall receiver + rtl_433 # decode 433/868 MHz sensors: weather stations, doorbells, meters + dump1090-fa # ADS-B aircraft tracking on 1090 MHz with a live web map + multimon-ng # decode POCSAG pagers, AFSK, DTMF (pipe rtl_fm into it) + welle-io # DAB+ digital radio + satdump # NOAA/Meteor weather satellite images + noaa-apt # simple NOAA APT satellite image decoder + # Heavier options — uncomment if wanted: + # sdrangel # kitchen-sink SDR app, many digital modes + # gnuradio # build your own DSP flowgraphs + ]) ++ [ + rtl_wmbus # wM-Bus demodulator (rtl_sdr | rtl_wmbus) + wmbusmeters # decode OMS smart meters + wmbus-listen # convenience: listen for OMS telegrams (see below) + oms-tui # full-screen OMS monitor TUI + ]; + in + { + devShells.${system}.default = pkgs.mkShell { + packages = sdrTools; + shellHook = '' + echo "SDR shell ready. Start with: rtl_test (OMS meters: wmbus-listen)" + ''; + }; + + packages.${system} = { + inherit rtl_wmbus wmbusmeters wmbus-listen oms-tui; + default = oms-tui; + + # Also expose the whole toolbox as one installable package set: + # nix profile install .#sdr-tools + sdr-tools = pkgs.buildEnv { + name = "sdr-tools"; + paths = sdrTools; + }; + }; + }; +} diff --git a/oms-tui/oms_tui/__init__.py b/oms-tui/oms_tui/__init__.py new file mode 100644 index 0000000..5a7f444 --- /dev/null +++ b/oms-tui/oms_tui/__init__.py @@ -0,0 +1,3 @@ +"""OMS / wireless M-Bus smart-meter monitor (TUI).""" + +__version__ = "0.1.0" diff --git a/oms-tui/oms_tui/__pycache__/__init__.cpython-313.pyc b/oms-tui/oms_tui/__pycache__/__init__.cpython-313.pyc new file mode 100644 index 0000000..0a54698 Binary files /dev/null and b/oms-tui/oms_tui/__pycache__/__init__.cpython-313.pyc differ diff --git a/oms-tui/oms_tui/__pycache__/app.cpython-313.pyc b/oms-tui/oms_tui/__pycache__/app.cpython-313.pyc new file mode 100644 index 0000000..abfcdcc Binary files /dev/null and b/oms-tui/oms_tui/__pycache__/app.cpython-313.pyc differ diff --git a/oms-tui/oms_tui/__pycache__/csvio.cpython-313.pyc b/oms-tui/oms_tui/__pycache__/csvio.cpython-313.pyc new file mode 100644 index 0000000..e2f787a Binary files /dev/null and b/oms-tui/oms_tui/__pycache__/csvio.cpython-313.pyc differ diff --git a/oms-tui/oms_tui/__pycache__/manufacturers.cpython-313.pyc b/oms-tui/oms_tui/__pycache__/manufacturers.cpython-313.pyc new file mode 100644 index 0000000..6c5e1fe Binary files /dev/null and b/oms-tui/oms_tui/__pycache__/manufacturers.cpython-313.pyc differ diff --git a/oms-tui/oms_tui/__pycache__/models.cpython-313.pyc b/oms-tui/oms_tui/__pycache__/models.cpython-313.pyc new file mode 100644 index 0000000..db4053a Binary files /dev/null and b/oms-tui/oms_tui/__pycache__/models.cpython-313.pyc differ diff --git a/oms-tui/oms_tui/__pycache__/parser.cpython-313.pyc b/oms-tui/oms_tui/__pycache__/parser.cpython-313.pyc new file mode 100644 index 0000000..c0c5517 Binary files /dev/null and b/oms-tui/oms_tui/__pycache__/parser.cpython-313.pyc differ diff --git a/oms-tui/oms_tui/__pycache__/wmbus.cpython-313.pyc b/oms-tui/oms_tui/__pycache__/wmbus.cpython-313.pyc new file mode 100644 index 0000000..9ffdb0e Binary files /dev/null and b/oms-tui/oms_tui/__pycache__/wmbus.cpython-313.pyc differ diff --git a/oms-tui/oms_tui/app.py b/oms-tui/oms_tui/app.py new file mode 100644 index 0000000..8160006 --- /dev/null +++ b/oms-tui/oms_tui/app.py @@ -0,0 +1,580 @@ +"""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.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 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"), + ] + + 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)) + 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() diff --git a/oms-tui/oms_tui/app.tcss b/oms-tui/oms_tui/app.tcss new file mode 100644 index 0000000..2f4d5bd --- /dev/null +++ b/oms-tui/oms_tui/app.tcss @@ -0,0 +1,70 @@ +Screen { + layout: vertical; +} + +#body { + height: 1fr; +} + +#left { + width: 44%; + min-width: 34; + border-right: solid $primary; +} + +#left-title { + padding: 0 1; + background: $boost; + color: $text; + text-style: bold; +} + +#devices { + height: 1fr; +} + +#right { + width: 1fr; +} + +#detail-scroll { + height: auto; + max-height: 62%; + border-bottom: solid $primary; + scrollbar-size-vertical: 1; +} + +#detail { + height: auto; + padding: 1 2; +} + +#stream-title { + padding: 0 1; + background: $boost; + color: $text; + text-style: bold; +} + +#stream { + height: 1fr; + padding: 0 1; + background: $surface; +} + +/* Modal prompt for keys / file paths */ +PromptScreen { + align: center middle; +} + +#dialog { + width: 72; + height: auto; + border: thick $primary; + background: $surface; + padding: 1 2; +} + +#dialog Label { + margin-bottom: 1; +} diff --git a/oms-tui/oms_tui/csvio.py b/oms-tui/oms_tui/csvio.py new file mode 100644 index 0000000..bfa0c0c --- /dev/null +++ b/oms-tui/oms_tui/csvio.py @@ -0,0 +1,77 @@ +"""CSV export/import of the device list, including AES keys.""" + +from __future__ import annotations + +import csv +import time +from typing import Dict, Iterable, List + +from .models import Device + +FIELDNAMES = [ + "id", + "manufacturer", + "manufacturer_name", + "media", + "driver", + "version", + "first_seen", + "last_seen", + "count", + "rssi", + "encryption", + "security_mode", + "status", + "aes_key", +] + + +def _fmt_time(ts: float) -> str: + if not ts: + return "" + return time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(ts)) + + +def export_devices(path: str, devices: Iterable[Device]) -> int: + n = 0 + with open(path, "w", newline="") as f: + writer = csv.DictWriter(f, fieldnames=FIELDNAMES) + writer.writeheader() + for d in devices: + writer.writerow( + { + "id": d.id, + "manufacturer": d.manufacturer, + "manufacturer_name": d.manufacturer_name, + "media": d.media, + "driver": d.driver, + "version": d.version, + "first_seen": _fmt_time(d.first_seen), + "last_seen": _fmt_time(d.last_seen), + "count": d.count, + "rssi": "" if d.rssi is None else d.rssi, + "encryption": d.encryption, + "security_mode": d.sec_mode or "", + "status": d.status, + "aes_key": d.key or "", + } + ) + n += 1 + return n + + +def import_rows(path: str) -> List[Dict[str, str]]: + """Read a previously exported CSV. Returns a list of row dicts. + + Only ``id`` is required; ``aes_key`` and the metadata columns are optional, + so hand-written key files with just ``id,aes_key`` import fine too. + """ + rows: List[Dict[str, str]] = [] + with open(path, newline="") as f: + reader = csv.DictReader(f) + for row in reader: + device_id = (row.get("id") or "").strip() + if not device_id: + continue + rows.append({k: (v or "").strip() for k, v in row.items()}) + return rows diff --git a/oms-tui/oms_tui/manufacturers.py b/oms-tui/oms_tui/manufacturers.py new file mode 100644 index 0000000..3295448 --- /dev/null +++ b/oms-tui/oms_tui/manufacturers.py @@ -0,0 +1,51 @@ +"""Map wM-Bus 3-letter FLAG manufacturer codes to full names. + +The DLL header only carries the 3-letter FLAG code (e.g. "KAM"). This is a +best-effort lookup of common metering manufacturers; unknown codes fall back +to the code itself. +""" + +_FLAG = { + "ABB": "ABB", + "AMT": "Aquametro", + "APA": "Apator", + "APT": "Apator", + "BME": "BMeters", + "BME": "BMeters", + "DME": "Diehl Metering", + "DWZ": "Lorenz", + "EFE": "Engelmann", + "ELS": "Elster", + "ELV": "Elvaco", + "EMH": "EMH Metering", + "ESY": "EasyMeter", + "GAV": "Carlo Gavazzi", + "GWF": "GWF", + "HYD": "Diehl (Hydrometer)", + "INE": "Innotas", + "ITW": "Itron", + "ITR": "Itron", + "KAM": "Kamstrup", + "KAW": "Kamstrup", + "LUG": "Landis+Gyr", + "LGB": "Landis+Gyr", + "MAD": "Maddalena", + "MTR": "Metrona", + "NZR": "NZR", + "QDS": "Qundis", + "REL": "Relay", + "RKE": "Viterra / Ista", + "SAP": "Sappel", + "SEN": "Sensus", + "SON": "Sontex", + "SPX": "Spanner-Pollux", + "TCH": "Techem", + "WEP": "Weptech", + "ZRI": "Zenner", +} + + +def manufacturer_name(code: str) -> str: + if not code: + return "" + return _FLAG.get(code.upper(), code.upper()) diff --git a/oms-tui/oms_tui/models.py b/oms-tui/oms_tui/models.py new file mode 100644 index 0000000..7bf01ae --- /dev/null +++ b/oms-tui/oms_tui/models.py @@ -0,0 +1,195 @@ +"""In-memory model of discovered OMS devices. Reset on every app start.""" + +from __future__ import annotations + +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 + + @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] + + 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 + # Count once per telegram, ignoring the duplicate DLL burst from a + # second matching meter. + if is_new or (now - dev._last_count_at) > COUNT_DEDUP_WINDOW: + dev.count += 1 + dev._last_count_at = 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 + return dev, is_new diff --git a/oms-tui/oms_tui/parser.py b/oms-tui/oms_tui/parser.py new file mode 100644 index 0000000..ea6fb8f --- /dev/null +++ b/oms-tui/oms_tui/parser.py @@ -0,0 +1,138 @@ +"""Parse wmbusmeters output lines into structured events. + +wmbusmeters is run with ``--format=json --verbose`` plus a wildcard +``scan auto '*' NOKEY`` meter. That produces, for *every* received telegram: + +* a verbose DLL line on stderr (universal discovery), e.g.:: + + (telegram) DLL L=2a C=44 (from meter SND_NR) M=2c2d (KAM) A=76348799 \ + VER=1b TYPE=16 (Cold water meter) (driver kamwater) DEV=rtlwmbus[] RSSI=97 + +* and, when a matching keyed meter is configured, a decrypted JSON line on + stdout, e.g. ``{"_":"telegram","id":"76348799","total_m3":6.408,...}``. + +stdout and stderr are merged, so a single parser handles both. +""" + +from __future__ import annotations + +import json +import re +from typing import Optional + +_DLL_RE = re.compile( + r"\(telegram\) DLL " + r"L=(?P[0-9a-fA-F]+) " + r"C=(?P[0-9a-fA-F]+).*?" + r"M=(?P[0-9a-fA-F]+) \((?P[A-Z?]+)\) " + r"A=(?P\w+) " + r"VER=(?P[0-9a-fA-F]+) " + r"TYPE=(?P[0-9a-fA-F]+) \((?P[^)]*)\) " + r"\(driver (?P\w+)\)" + r"(?:.*?RSSI=(?P-?\d+))?" +) + +# JSON keys that are metadata, not measurement values. +_META = {"_", "media", "meter", "name", "id", "timestamp", "device", "rssi_dbm"} + +# OMS/EN13757-3 TPL security modes we can name. +_SEC_MODE = {"AES_CBC_IV": 5, "AES_CBC_NO_IV": 4, "AES_CTR": 7, "AES_DES": 2} + + +def _find(pattern: str, line: str) -> str: + m = re.search(pattern, line) + return m.group(1) if m else "" + + +def parse_dll(line: str) -> Optional[dict]: + m = _DLL_RE.search(line) + if not m: + return None + g = m.groupdict() + return { + "type": "dll", + "id": g["id"], + "manufacturer": g["mfct"], + "version": g["ver"], + "media": g["media"], + "driver": g["driver"], + "c_field": g["C"], + "rssi": int(g["rssi"]) if g["rssi"] is not None else None, + } + + +def parse_json_telegram(line: str) -> Optional[dict]: + line = line.strip() + if not line.startswith("{"): + return None + try: + obj = json.loads(line) + except json.JSONDecodeError: + return None + if obj.get("_") != "telegram": + return None + fields = {k: v for k, v in obj.items() if k not in _META} + return { + "type": "json", + "id": str(obj.get("id", "")), + "driver": obj.get("meter", ""), + "media": obj.get("media", ""), + "name": obj.get("name", ""), + "timestamp": obj.get("timestamp", ""), + "rssi": obj.get("rssi_dbm"), + "fields": fields, + } + + +def parse_ell(line: str) -> Optional[dict]: + """Extended Link Layer line — carries ELL-layer encryption (e.g. AES-CTR).""" + if "(telegram) ELL" not in line: + return None + enc = "AES-CTR" if "AES_CTR" in line else ("AES-CBC" if "AES_CBC" in line else "") + return { + "type": "ell", + "ci": _find(r"CI=(\w+)", line), + "cc": _find(r"CC=(\w+)", line), + "sn": _find(r"SN=(\w+)", line), + "session": _find(r"session=(\d+)", line), + "enc": enc, + } + + +def parse_tpl(line: str) -> Optional[dict]: + """Transport Layer line — carries the TPL security mode (5=CBC, 7=CTR).""" + if "(telegram) TPL" not in line: + return None + sec_mode = 0 + enc = "" + for name, mode in _SEC_MODE.items(): + if name in line: + sec_mode = mode + enc = "AES-CBC" if "CBC" in name else ("AES-CTR" if "CTR" in name else name) + break + return { + "type": "tpl", + "ci": _find(r"CI=(\w+)", line), + "acc": _find(r"ACC=(\w+)", line), + "sts": _find(r"STS=(\w+)", line), + "cfg": _find(r"CFG=(\w+)", line), + "nb": _find(r"nb=(\d+)", line), + "sec_mode": sec_mode, + "enc": enc, + } + + +def parse_line(line: str) -> Optional[dict]: + """Return a structured event for a wmbusmeters output line, or None.""" + stripped = line.lstrip() + if stripped.startswith("{"): + ev = parse_json_telegram(line) + if ev is not None: + return ev + if "(telegram) DLL" in line: + return parse_dll(line) + if "(telegram) ELL" in line: + return parse_ell(line) + if "(telegram) TPL" in line: + return parse_tpl(line) + return None diff --git a/oms-tui/oms_tui/wmbus.py b/oms-tui/oms_tui/wmbus.py new file mode 100644 index 0000000..84bc76f --- /dev/null +++ b/oms-tui/oms_tui/wmbus.py @@ -0,0 +1,121 @@ +"""Drive wmbusmeters as a subprocess and stream parsed telegram events.""" + +from __future__ import annotations + +import asyncio +from typing import Callable, Dict, Optional + +from .parser import parse_line + +EventCallback = Callable[[dict], None] + + +class WMBusSource: + """Runs one wmbusmeters process and feeds parsed events to a callback. + + The command is:: + + wmbusmeters --format=json --verbose [--listento=MODES] DEVICE \ + [m_ auto ...] \ + scan auto '*' NOKEY + + The trailing wildcard meter makes wmbusmeters emit a verbose DLL line for + *every* telegram (universal discovery); keyed meters additionally emit + decrypted JSON. Changing keys requires relaunching the process (:meth:`restart`). + """ + + def __init__( + self, + device: str = "rtlwmbus", + modes: str = "c1,t1", + wmbusmeters: str = "wmbusmeters", + stdin_data: Optional[bytes] = None, + ) -> None: + self.device = device + self.modes = modes + self.wmbusmeters = wmbusmeters + self.stdin_data = stdin_data + self._keys: Dict[str, str] = {} + self._proc: Optional[asyncio.subprocess.Process] = None + self._task: Optional[asyncio.Task] = None + self._on_event: Optional[EventCallback] = None + self.running = False + + def set_keys(self, keys: Dict[str, str]) -> None: + self._keys = dict(keys) + + def build_cmd(self) -> list[str]: + cmd = [self.wmbusmeters, "--format=json", "--verbose"] + if self.device.startswith("rtl"): + cmd.append(f"--listento={self.modes}") + cmd.append(self.device) + for device_id, key in self._keys.items(): + cmd += [f"m_{device_id}", "auto", device_id, key] + cmd += ["scan", "auto", "*", "NOKEY"] + return cmd + + async def start(self, on_event: EventCallback) -> None: + self._on_event = on_event + await self._spawn() + + async def _spawn(self) -> None: + cmd = self.build_cmd() + use_stdin = self.stdin_data is not None + self._proc = await asyncio.create_subprocess_exec( + *cmd, + stdin=asyncio.subprocess.PIPE if use_stdin else asyncio.subprocess.DEVNULL, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.STDOUT, + ) + self.running = True + if use_stdin and self._proc.stdin is not None: + self._proc.stdin.write(self.stdin_data) # type: ignore[arg-type] + await self._proc.stdin.drain() + self._proc.stdin.close() + self._task = asyncio.create_task(self._read_loop()) + + async def _read_loop(self) -> None: + assert self._proc is not None and self._proc.stdout is not None + try: + while True: + raw = await self._proc.stdout.readline() + if not raw: + break + line = raw.decode("utf-8", "replace").rstrip("\n") + if not self._on_event: + continue + ev = parse_line(line) + if ev is not None: + self._on_event(ev) + except asyncio.CancelledError: # pragma: no cover - shutdown path + pass + + async def _kill(self) -> None: + if self._task is not None: + self._task.cancel() + try: + await self._task + except (asyncio.CancelledError, Exception): + pass + self._task = None + if self._proc is not None and self._proc.returncode is None: + try: + self._proc.terminate() + await asyncio.wait_for(self._proc.wait(), timeout=3) + except (asyncio.TimeoutError, ProcessLookupError): + try: + self._proc.kill() + except ProcessLookupError: + pass + self._proc = None + self.running = False + + async def restart(self) -> None: + await self._kill() + # a one-shot stdin replay has been consumed; don't re-feed on restart + if self.stdin_data is not None: + self.stdin_data = None + await self._spawn() + + async def stop(self) -> None: + await self._kill() diff --git a/oms-tui/pyproject.toml b/oms-tui/pyproject.toml new file mode 100644 index 0000000..0700894 --- /dev/null +++ b/oms-tui/pyproject.toml @@ -0,0 +1,23 @@ +[build-system] +requires = ["setuptools>=61"] +build-backend = "setuptools.build_meta" + +[project] +name = "oms-tui" +version = "0.1.0" +description = "TUI for monitoring OMS / wireless M-Bus smart meters via wmbusmeters" +requires-python = ">=3.10" +dependencies = ["textual>=0.60"] + +[project.scripts] +oms-tui = "oms_tui.app:main" + +[tool.setuptools.packages.find] +where = ["."] +include = ["oms_tui*"] + +[tool.setuptools.package-data] +oms_tui = ["*.tcss"] + +[tool.pytest.ini_options] +asyncio_mode = "auto" diff --git a/oms-tui/sample.msg b/oms-tui/sample.msg new file mode 100644 index 0000000..03cca93 --- /dev/null +++ b/oms-tui/sample.msg @@ -0,0 +1,2 @@ +T1;1;1;2019-04-03 19:00:42.000;97;148;88888888;0x6e4401068888888805077a85006085bc2630713819512eb4cd87fba554fb43f67cf9654a68ee8e194088160df752e716238292e8af1ac20986202ee561d743602466915e42f1105d9c6782a54504e4f099e65a7656b930c73a30775122d2fdf074b5035cfaa7e0050bf32faae03a77 +C1;1;1;2020-01-23 10:25:13.000;97;148;76348799;0x2A442D2C998734761B168D2091D37CAC21E1D68CDAFFCD3DC452BD802913FF7B1706CA9E355D6C2701CC24 diff --git a/oms-tui/tests/__pycache__/test_app.cpython-313-pytest-9.0.3.pyc b/oms-tui/tests/__pycache__/test_app.cpython-313-pytest-9.0.3.pyc new file mode 100644 index 0000000..946074f Binary files /dev/null and b/oms-tui/tests/__pycache__/test_app.cpython-313-pytest-9.0.3.pyc differ diff --git a/oms-tui/tests/__pycache__/test_integration.cpython-313-pytest-9.0.3.pyc b/oms-tui/tests/__pycache__/test_integration.cpython-313-pytest-9.0.3.pyc new file mode 100644 index 0000000..c828ec0 Binary files /dev/null and b/oms-tui/tests/__pycache__/test_integration.cpython-313-pytest-9.0.3.pyc differ diff --git a/oms-tui/tests/__pycache__/test_models_csv.cpython-313-pytest-9.0.3.pyc b/oms-tui/tests/__pycache__/test_models_csv.cpython-313-pytest-9.0.3.pyc new file mode 100644 index 0000000..f822e31 Binary files /dev/null and b/oms-tui/tests/__pycache__/test_models_csv.cpython-313-pytest-9.0.3.pyc differ diff --git a/oms-tui/tests/__pycache__/test_parser.cpython-313-pytest-9.0.3.pyc b/oms-tui/tests/__pycache__/test_parser.cpython-313-pytest-9.0.3.pyc new file mode 100644 index 0000000..4f7fd46 Binary files /dev/null and b/oms-tui/tests/__pycache__/test_parser.cpython-313-pytest-9.0.3.pyc differ diff --git a/oms-tui/tests/test_app.py b/oms-tui/tests/test_app.py new file mode 100644 index 0000000..c204fd8 --- /dev/null +++ b/oms-tui/tests/test_app.py @@ -0,0 +1,84 @@ +import pytest + +from oms_tui.app import OmsApp, TelegramEvent + +DLL = { + "type": "dll", + "id": "76348799", + "manufacturer": "KAM", + "version": "1b", + "media": "Cold water meter", + "driver": "kamwater", + "c_field": "44", + "rssi": 97, +} +JSON = { + "type": "json", + "id": "76348799", + "driver": "kamwater", + "media": "cold water", + "fields": {"total_m3": 6.408, "status": "DRY"}, +} + + +@pytest.mark.asyncio +async def test_device_appears_and_selects(): + app = OmsApp(device="none", autostart=False) + async with app.run_test() as pilot: + app.post_message(TelegramEvent(DLL)) + await pilot.pause() + table = app.query_one("#devices") + assert table.row_count == 1 + + # nothing selected yet -> global stream mode + assert app.show_all is True + + # decrypted JSON updates the model + app.post_message(TelegramEvent(JSON)) + await pilot.pause() + dev = app.store.get("76348799") + assert dev.decrypted is True + assert dev.fields["total_m3"] == 6.408 + assert dev.count == 1 # one DLL telegram + + # select the device + app.action_select() + await pilot.pause() + assert app.show_all is False + assert app.selected_id == "76348799" + + +@pytest.mark.asyncio +async def test_set_key_via_modal(): + app = OmsApp(device="none", autostart=False) + async with app.run_test() as pilot: + app.post_message(TelegramEvent(DLL)) + await pilot.pause() + app.action_select() + await pilot.pause() + + app.action_set_key() # opens the modal prompt + await pilot.pause() + await pilot.press(*"28F64A24988064A079AA2C807D6102AE") + await pilot.press("enter") # submit + await pilot.pause() + + dev = app.store.get("76348799") + assert dev.key == "28F64A24988064A079AA2C807D6102AE" + assert dev.status == "key set" + + +@pytest.mark.asyncio +async def test_import_preloads_key(tmp_path): + path = str(tmp_path / "keys.csv") + with open(path, "w") as f: + f.write("id,aes_key\n76348799,28F64A24988064A079AA2C807D6102AE\n") + + app = OmsApp(device="none", autostart=False) + async with app.run_test() as pilot: + app._do_import(path, announce=False) + await pilot.pause() + dev = app.store.get("76348799") + assert dev is not None + assert dev.key == "28F64A24988064A079AA2C807D6102AE" + assert dev.status == "key set" diff --git a/oms-tui/tests/test_integration.py b/oms-tui/tests/test_integration.py new file mode 100644 index 0000000..9ccc165 --- /dev/null +++ b/oms-tui/tests/test_integration.py @@ -0,0 +1,43 @@ +"""End-to-end test through the real wmbusmeters binary (skipped if absent).""" + +import asyncio +import shutil + +import pytest + +from oms_tui.wmbus import WMBusSource + +TELEGRAM = ( + "C1;1;1;2020-01-23 10:25:13.000;97;148;76348799;" + "0x2A442D2C998734761B168D2091D37CAC21E1D68CDAFFCD3DC452BD802913FF7B1706CA9E355D6C2701CC24\n" +) +KEY = "28F64A24988064A079AA2C807D6102AE" + + +@pytest.mark.skipif(shutil.which("wmbusmeters") is None, reason="wmbusmeters not on PATH") +def test_pipeline_discovers_and_decrypts(): + events = [] + + async def run(): + src = WMBusSource( + device="stdin:rtlwmbus", stdin_data=(TELEGRAM * 3).encode() + ) + src.set_keys({"76348799": KEY}) + await src.start(events.append) + for _ in range(50): # up to ~5s for the process to finish replaying + await asyncio.sleep(0.1) + if any(e["type"] == "json" for e in events): + break + await src.stop() + + asyncio.run(run()) + + ids = {e.get("id") for e in events if e["type"] in ("dll", "json")} + assert "76348799" in ids + + dlls = [e for e in events if e["type"] == "dll"] + assert any(e["manufacturer"] == "KAM" for e in dlls) + + jsons = [e for e in events if e["type"] == "json"] + assert jsons, "no decrypted telegram produced" + assert any(e["fields"].get("total_m3") == 6.408 for e in jsons) diff --git a/oms-tui/tests/test_models_csv.py b/oms-tui/tests/test_models_csv.py new file mode 100644 index 0000000..9aed0fe --- /dev/null +++ b/oms-tui/tests/test_models_csv.py @@ -0,0 +1,149 @@ +import os + +from oms_tui.csvio import export_devices, import_rows +from oms_tui.models import DeviceStore + + +def _dll(device_id, mfct="KAM", media="Cold water meter"): + return { + "type": "dll", + "id": device_id, + "manufacturer": mfct, + "version": "1b", + "media": media, + "driver": "kamwater", + "c_field": "44", + "rssi": 97, + } + + +def test_duplicate_dll_burst_counts_once(): + # Two DLL lines for the SAME telegram (keyed meter + wildcard) arrive + # microseconds apart -> counted once. + store = DeviceStore() + store.upsert_discovery(_dll("111"), now=1000.0) + store.upsert_discovery(_dll("111"), now=1000.001) + assert store.get("111").count == 1 + + +def test_real_retransmissions_count_separately(): + store = DeviceStore() + store.upsert_discovery(_dll("111"), now=1000.0) + store.upsert_discovery(_dll("111"), now=1005.0) # seconds later + dev = store.get("111") + assert dev.count == 2 + assert dev.manufacturer_name == "Kamstrup" + assert dev.status == "listening" # DLL only, encryption not yet known + + +def test_encrypted_meter_with_plaintext_read_is_locked_not_open(): + # The exact contradiction the user hit: an encrypted (mode 5) meter that + # also exposes some plaintext records read by the wildcard ("scan"). + store = DeviceStore() + store.upsert_discovery(_dll("61471000")) + store.apply_tpl( + {"type": "tpl", "ci": "7a", "cfg": "8560", "nb": "6", "sec_mode": 5, + "enc": "AES-CBC"}, + "61471000", + ) + store.apply_json( + {"type": "json", "id": "61471000", "name": "scan", "driver": "x", + "media": "water", "fields": {"total_m3": 1.0}} + ) + dev = store.get("61471000") + assert dev.encrypted is True + assert dev.decrypted is True # we do have some values + assert dev.key_decrypted is False + assert dev.fields_source == "plain" + assert dev.status == "locked" # NOT "open" + assert dev.encryption == "mode 5 (AES-CBC)" + + +def test_keyed_decrypt_marks_decrypted(): + store = DeviceStore() + store.upsert_discovery(_dll("61471000")) + store.apply_tpl({"type": "tpl", "sec_mode": 5, "enc": "AES-CBC"}, "61471000") + dev = store.get("61471000") + dev.key = "00112233445566778899AABBCCDDEEFF" + store.apply_json( + {"type": "json", "id": "61471000", "name": "m_61471000", "driver": "x", + "media": "water", "fields": {"total_m3": 2.0}} + ) + assert dev.key_decrypted is True + assert dev.fields_source == "key" + assert dev.status == "decrypted" + + +def test_unencrypted_device_is_open_not_locked(): + store = DeviceStore() + store.upsert_discovery(_dll("333")) + store.apply_json( + {"type": "json", "id": "333", "driver": "apator162", "media": "water", + "fields": {"total_m3": 1.2}} + ) + dev = store.get("333") + assert dev.decrypted is True + assert dev.key is None + assert dev.status == "open" # not "locked" + assert dev.status_icon == "📖" + assert dev.status_color == "green" + + +def test_encryption_from_tpl_and_ell(): + store = DeviceStore() + store.upsert_discovery(_dll("444")) + store.apply_tpl( + {"type": "tpl", "ci": "7a", "cfg": "8560", "nb": "6", "sec_mode": 5, + "enc": "AES-CBC"}, + "444", + ) + dev = store.get("444") + assert dev.sec_mode == 5 + assert dev.encryption == "mode 5 (AES-CBC)" + assert dev.enc_blocks == "6" + + store2 = DeviceStore() + store2.upsert_discovery(_dll("555")) + store2.apply_ell( + {"type": "ell", "ci": "8d", "sn": "d37cac21", "session": "3", "enc": "AES-CTR"}, + "555", + ) + assert store2.get("555").encryption == "AES-CTR (ELL)" + + +def test_json_sets_decrypted_without_counting(): + store = DeviceStore() + store.upsert_discovery(_dll("222")) + store.apply_json( + {"type": "json", "id": "222", "driver": "kamwater", "media": "cold water", + "fields": {"total_m3": 6.408}} + ) + dev = store.get("222") + assert dev.count == 1 # JSON does not double-count + assert dev.decrypted is True + assert dev.fields["total_m3"] == 6.408 + + +def test_csv_roundtrip_with_keys(tmp_path): + store = DeviceStore() + store.upsert_discovery(_dll("76348799")) + store.get("76348799").key = "28F64A24988064A079AA2C807D6102AE" + store.upsert_discovery(_dll("88888888", mfct="TCH", media="Heat")) + + path = os.path.join(tmp_path, "out.csv") + assert export_devices(path, store.ordered()) == 2 + + rows = import_rows(path) + by_id = {r["id"]: r for r in rows} + assert by_id["76348799"]["aes_key"] == "28F64A24988064A079AA2C807D6102AE" + assert by_id["76348799"]["manufacturer"] == "KAM" + assert by_id["88888888"]["aes_key"] == "" + + +def test_import_minimal_key_file(tmp_path): + path = os.path.join(tmp_path, "keys.csv") + with open(path, "w") as f: + f.write("id,aes_key\n76348799,28F64A24988064A079AA2C807D6102AE\n") + rows = import_rows(path) + assert rows[0]["id"] == "76348799" + assert rows[0]["aes_key"].startswith("28F64A") diff --git a/oms-tui/tests/test_parser.py b/oms-tui/tests/test_parser.py new file mode 100644 index 0000000..6f13684 --- /dev/null +++ b/oms-tui/tests/test_parser.py @@ -0,0 +1,64 @@ +from oms_tui.parser import parse_line + +DLL = ( + "(telegram) DLL L=2a C=44 (from meter SND_NR) M=2c2d (KAM) A=76348799 " + "VER=1b TYPE=16 (Cold water meter) (driver kamwater) DEV=rtlwmbus[] RSSI=97" +) +JSON = ( + '{"_":"telegram","media":"cold water","meter":"kamwater","name":"Vatten",' + '"id":"76348799","total_m3":6.408,"target_m3":6.408,"status":"DRY",' + '"timestamp":"2020-01-23T10:25:13Z","device":"rtlwmbus[]","rssi_dbm":97}' +) + + +def test_parse_dll(): + ev = parse_line(DLL) + assert ev["type"] == "dll" + assert ev["id"] == "76348799" + assert ev["manufacturer"] == "KAM" + assert ev["version"] == "1b" + assert ev["media"] == "Cold water meter" + assert ev["driver"] == "kamwater" + assert ev["rssi"] == 97 + + +def test_parse_json(): + ev = parse_line(JSON) + assert ev["type"] == "json" + assert ev["id"] == "76348799" + assert ev["driver"] == "kamwater" + assert ev["fields"]["total_m3"] == 6.408 + # metadata keys are excluded from measurement fields + assert "timestamp" not in ev["fields"] + assert "rssi_dbm" not in ev["fields"] + + +def test_parse_ell_encryption(): + line = ( + "(telegram) ELL CI=8d CC=20 (slow_resp sync) ACC=91 SN=d37cac21 " + "(AES_CTR session=3 time=1755085) CRC=576c" + ) + ev = parse_line(line) + assert ev["type"] == "ell" + assert ev["enc"] == "AES-CTR" + assert ev["session"] == "3" + + +def test_parse_tpl_security_mode(): + line = "(telegram) TPL CI=7a ACC=85 STS=00 CFG=8560 (bidirectional AES_CBC_IV nb=6 cntn=0 ra=0 hc=0)" + ev = parse_line(line) + assert ev["type"] == "tpl" + assert ev["sec_mode"] == 5 + assert ev["enc"] == "AES-CBC" + assert ev["nb"] == "6" + + +def test_parse_tpl_no_security(): + ev = parse_line("(telegram) TPL CI=78") + assert ev["type"] == "tpl" + assert ev["sec_mode"] == 0 + + +def test_parse_noise_returns_none(): + assert parse_line("(config) number of meters: 0") is None + assert parse_line("") is None diff --git a/oms_devices.csv b/oms_devices.csv new file mode 100644 index 0000000..cd3f747 --- /dev/null +++ b/oms_devices.csv @@ -0,0 +1,25 @@ +id,manufacturer,manufacturer_name,media,driver,version,first_seen,last_seen,count,rssi,status,aes_key +23252256,ITW,Itron,Water meter,itron,00,2026-07-24 12:13:32,2026-07-24 12:18:31,2,122,locked, +61379612,,,warm water,qwaterv2,,2026-07-24 12:13:38,2026-07-24 12:19:30,0,,open, +61471660,QDS,Qundis,water,qwaterv2,1d,2026-07-24 12:13:39,2026-07-24 12:19:22,5,124,open, +10179914,ITW,Itron,water,itron,00,2026-07-24 12:14:11,2026-07-24 12:19:11,2,46,open, +23252389,ITW,Itron,Water meter,itron,00,2026-07-24 12:14:38,2026-07-24 12:19:38,2,130,locked, +37287397,,,radio converter (meter side),qheatv2,,2026-07-24 12:14:39,2026-07-24 12:20:15,0,,open, +09237338,ITW,Itron,Cold water meter,itron,00,2026-07-24 12:14:52,2026-07-24 12:19:52,2,61,locked, +68544330,,,heat,qheatv2,,2026-07-24 12:14:56,2026-07-24 12:14:56,0,,open, +23252290,ITW,Itron,Water meter,itron,00,2026-07-24 12:15:16,2026-07-24 12:15:16,1,34,locked, +00258334,ITW,Itron,water,itron,00,2026-07-24 12:16:02,2026-07-24 12:16:02,1,120,open, +08936799,ITW,Itron,Cold water meter,itron,00,2026-07-24 12:16:09,2026-07-24 12:16:09,1,73,locked, +23252289,ITW,Itron,Water meter,itron,00,2026-07-24 12:16:17,2026-07-24 12:16:17,1,74,locked, +09155592,ITW,Itron,Cold water meter,itron,00,2026-07-24 12:16:31,2026-07-24 12:16:31,1,57,locked, +02713603,ITW,Itron,Cold water meter,itron,00,2026-07-24 12:17:05,2026-07-24 12:17:05,1,60,locked, +00340954,ITW,Itron,water,itron,00,2026-07-24 12:17:08,2026-07-24 12:17:08,1,85,open, +23252254,ITW,Itron,Water meter,itron,00,2026-07-24 12:17:08,2026-07-24 12:17:08,1,46,locked, +23252292,ITW,Itron,Water meter,itron,00,2026-07-24 12:17:09,2026-07-24 12:17:09,1,67,locked, +00227616,ITW,Itron,Cold water meter,itron,00,2026-07-24 12:17:21,2026-07-24 12:17:21,1,37,locked, +23252281,ITW,Itron,Water meter,itron,00,2026-07-24 12:17:33,2026-07-24 12:17:33,1,123,locked, +23252286,ITW,Itron,Water meter,itron,00,2026-07-24 12:17:54,2026-07-24 12:17:54,1,100,locked, +00093643,ITW,Itron,Cold water meter,itron,00,2026-07-24 12:18:07,2026-07-24 12:18:07,1,51,locked, +00124993,ITW,Itron,water,itron,00,2026-07-24 12:18:16,2026-07-24 12:18:16,1,34,open, +00065762,ITW,Itron,Cold water meter,itron,00,2026-07-24 12:18:31,2026-07-24 12:18:31,1,51,locked, +00257608,ITW,Itron,water,itron,00,2026-07-24 12:19:08,2026-07-24 12:19:08,1,39,open,