Files
2026-07-24 13:57:13 +02:00

608 lines
22 KiB
Python

"""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()