608 lines
22 KiB
Python
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()
|