From 62b57a7e569920e62f11fec5fea8a74cf45f570f Mon Sep 17 00:00:00 2001 From: Milan Pandurov Date: Tue, 29 Sep 2026 12:08:48 +0200 Subject: [PATCH] Add MQTT diagnostics, broker connection probe, and switch to Debian base The MQTT publisher now tracks connection state, publish counters, discovery and availability timestamps, and the last error, exposed at /api/mqtt and on the dashboard as an MQTT card and a Home Assistant section. Connection failures, disconnects, and dropped publishes are logged; paho's on_connect_fail callback was not registered before, so failed attempts were silent. MQTT connect, disconnect, and unreachable events go to the timeline. When a connection attempt fails the service runs a probe from inside the container: system resolver, each nameserver from resolv.conf queried directly, /etc/hosts, and a TCP connect to every address found. The result is shown as a summary and raw JSON on the dashboard and can be rerun via POST /api/mqtt/probe. The image base moves from Alpine to Debian slim. musl queries all nameservers in parallel and accepts the first reply, so a fast public NXDOMAIN beats a slower local server that knows the name. glibc asks the nameservers in order. Co-Authored-By: Claude Fable 5.1 --- Dockerfile | 4 +- README.md | 22 ++++- app.py | 18 +++- hamqtt.py | 137 +++++++++++++++++++++++++++-- netprobe.py | 190 +++++++++++++++++++++++++++++++++++++++++ templates/index.html | 115 ++++++++++++++++++++++++- tests/test_hamqtt.py | 140 ++++++++++++++++++++++++++++++ tests/test_netprobe.py | 125 +++++++++++++++++++++++++++ 8 files changed, 737 insertions(+), 14 deletions(-) create mode 100644 netprobe.py create mode 100644 tests/test_hamqtt.py create mode 100644 tests/test_netprobe.py diff --git a/Dockerfile b/Dockerfile index dd81f83..7aa11c9 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,5 +1,5 @@ -FROM python:3.12-alpine -RUN apk add --no-cache iw +FROM python:3.12-slim +RUN apt-get update && apt-get install -y --no-install-recommends iw && rm -rf /var/lib/apt/lists/* WORKDIR /app COPY requirements.txt /app/ RUN pip install --no-cache-dir -r requirements.txt diff --git a/README.md b/README.md index f41acde..47d9111 100644 --- a/README.md +++ b/README.md @@ -16,6 +16,8 @@ Assistant as a device with sensors and no manual configuration. | `/` | Dashboard: status, Wi-Fi, graphs, timeline | | `/api/status` | Latest state of every monitor as JSON | | `/api/wifi` | Interface list, link details, last scan results | +| `/api/mqtt` | MQTT connection state, counters, last error, last published payload, last probe | +| `POST /api/mqtt/probe` | Run the broker connection probe now and return its result | | `/api/events?limit=200` | Timeline events, newest first | | `/api/samples?prefix=reach.&range=3600` | Bucketed samples for graphs | @@ -89,6 +91,7 @@ All settings are environment variables. Intervals are in seconds. - Wi-Fi connected, disconnected, roamed to another BSSID - Host rebooted - Health check service started +- MQTT connected, disconnected, cannot reach broker ## Home Assistant @@ -100,6 +103,23 @@ availability to `healthcheck//availability`. When Home Assistant restarts it announces itself on `homeassistant/status` and the service re-publishes discovery. +The dashboard has an MQTT card and a Home Assistant section showing the +connection state, when it last connected or failed, how many messages were +published or dropped, when discovery and availability were last sent, and the +last state payload. The same data is at `/api/mqtt`. Use it to tell an app-side +problem (not connected, publishes dropped) from a Home Assistant-side one +(connected and publishing, but the device is still unavailable). + +When a connection attempt fails the service runs a probe from inside the +container, at most every two minutes: it resolves the broker name with the +system resolver, queries each nameserver from `/etc/resolv.conf` directly, +checks `/etc/hosts`, and tries a TCP connect to every address it found. The +result is shown as a summary and as raw JSON in the Home Assistant section, and +can be rerun from the dashboard. If the nameservers disagree about the broker +name, for example a local server that knows it and a public one that does not, +the probe says so. The image is Debian based, so glibc asks the nameservers in +resolv.conf order and only moves on when one does not answer. + Entities per device: - Internet, DNS, Wi-Fi (binary sensors, `connectivity` class) @@ -119,5 +139,5 @@ Publish a new image for all Raspberry Pi architectures: ```sh docker buildx build --platform linux/amd64,linux/arm64,linux/arm/v7 \ - -t hiimmilan/health-check:latest -t hiimmilan/health-check:2.0.0 --push . + -t hiimmilan/health-check:latest -t hiimmilan/health-check:2.3.0 --push . ``` diff --git a/app.py b/app.py index d2e887e..236f9d5 100644 --- a/app.py +++ b/app.py @@ -13,6 +13,7 @@ log = logging.getLogger("app") storage = Storage(cfg.db_path) state = State() +publisher = None app = Flask(__name__) WIFI_KEYS = ("wifi_interface", "wifi_interfaces", "wifi_status", "wifi_error", @@ -51,6 +52,20 @@ def api_wifi(): return jsonify(result) +@app.route('/api/mqtt') +def api_mqtt(): + result = publisher.diagnostics() if publisher else {"enabled": False} + result["now"] = time.time() + return jsonify(result) + + +@app.route('/api/mqtt/probe', methods=['POST']) +def api_mqtt_probe(): + if not publisher: + return jsonify({"error": "MQTT is not configured"}), 404 + return jsonify(publisher.run_probe()) + + @app.route('/api/events') def api_events(): limit = min(max(request.args.get("limit", 200, type=int), 1), 1000) @@ -67,12 +82,13 @@ def api_samples(): def start_background(): + global publisher storage.add_event("service_started", "Health check service started") storage.prune_samples(cfg.sample_retention_days) if cfg.mqtt_host: from hamqtt import HomeAssistantPublisher - publisher = HomeAssistantPublisher(cfg) + publisher = HomeAssistantPublisher(cfg, storage) state.on_change(publisher.publish_state) publisher.start() diff --git a/hamqtt.py b/hamqtt.py index 846b5f4..ba24b16 100644 --- a/hamqtt.py +++ b/hamqtt.py @@ -1,11 +1,17 @@ import json import logging +import threading +import time import paho.mqtt.client as mqtt +import netprobe + log = logging.getLogger(__name__) -VERSION = "2.0.0" +PROBE_INTERVAL = 120 + +VERSION = "2.3.0" CONNECTIVITY = {"device_class": "connectivity"} @@ -56,53 +62,163 @@ def flatten(snapshot): class HomeAssistantPublisher: - def __init__(self, cfg): + def __init__(self, cfg, storage=None): self.cfg = cfg + self.storage = storage self.node = f"healthcheck_{cfg.device_id}" base = f"healthcheck/{cfg.device_id}" self.state_topic = f"{base}/state" self.availability_topic = f"{base}/availability" self.status_topic = f"{cfg.mqtt_discovery_prefix}/status" self._last_payload = None + self._lock = threading.Lock() + self._diag = { + "status": "not_started", + "connected": False, + "last_error": None, + "started_at": None, + "connected_at": None, + "disconnected_at": None, + "last_attempt_at": None, + "connect_failures": 0, + "disconnects": 0, + "discovery_count": 0, + "discovery_at": None, + "availability_at": None, + "ha_status": None, + "ha_status_at": None, + "publish_count": 0, + "publish_dropped": 0, + "last_publish_at": None, + "last_publish_error": None, + "probe": None, + "probe_running": False, + } self.client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, client_id=self.node) if cfg.mqtt_username: self.client.username_pw_set(cfg.mqtt_username, cfg.mqtt_password) self.client.will_set(self.availability_topic, "offline", retain=True) self.client.on_connect = self._on_connect + self.client.on_connect_fail = self._on_connect_fail self.client.on_message = self._on_message self.client.on_disconnect = self._on_disconnect def start(self): self.client.reconnect_delay_set(min_delay=2, max_delay=60) + self._set(status="connecting", started_at=time.time(), last_attempt_at=time.time()) self.client.connect_async(self.cfg.mqtt_host, self.cfg.mqtt_port, keepalive=60) self.client.loop_start() log.info("MQTT: connecting to %s:%s as %s", self.cfg.mqtt_host, self.cfg.mqtt_port, self.node) + def diagnostics(self): + with self._lock: + diag = dict(self._diag) + diag.update({ + "enabled": True, + "host": self.cfg.mqtt_host, + "port": self.cfg.mqtt_port, + "username": self.cfg.mqtt_username or None, + "client_id": self.node, + "state_topic": self.state_topic, + "availability_topic": self.availability_topic, + "discovery_prefix": self.cfg.mqtt_discovery_prefix, + "last_payload": self._last_payload, + }) + return diag + def publish_state(self, snapshot): payload = flatten(snapshot) if payload == self._last_payload: return self._last_payload = payload - self.client.publish(self.state_topic, json.dumps(payload), retain=True) + self._publish(self.state_topic, json.dumps(payload)) + + def run_probe(self): + with self._lock: + if self._diag["probe_running"]: + return self._diag["probe"] + self._diag["probe_running"] = True + try: + result = netprobe.probe(self.cfg.mqtt_host, self.cfg.mqtt_port) + except Exception as exc: + log.exception("MQTT: probe failed") + result = {"ts": time.time(), "error": str(exc), "summary": [f"probe crashed: {exc}"]} + finally: + self._set(probe_running=False) + self._set(probe=result) + for line in result["summary"]: + log.warning("MQTT probe: %s", line) + return result + + def _probe_in_background(self): + with self._lock: + probe = self._diag["probe"] + if self._diag["probe_running"] or (probe and time.time() - probe["ts"] < PROBE_INTERVAL): + return + threading.Thread(target=self.run_probe, name="mqtt-probe", daemon=True).start() + + def _set(self, *counters, **fields): + with self._lock: + for key in counters: + self._diag[key] += 1 + self._diag.update(fields) + + def _event(self, kind, message): + if self.storage is not None: + self.storage.add_event(kind, message) + + def _publish(self, topic, payload): + info = self.client.publish(topic, payload, retain=True) + if info.rc == mqtt.MQTT_ERR_SUCCESS: + self._set("publish_count", last_publish_at=time.time(), last_publish_error=None) + else: + error = mqtt.error_string(info.rc) + self._set("publish_dropped", last_publish_error=error) + log.warning("MQTT: publish to %s failed: %s", topic, error) + return info def _on_connect(self, client, userdata, flags, reason_code, properties): if reason_code.is_failure: log.warning("MQTT: connection refused: %s", reason_code) + self._set("connect_failures", status="refused", connected=False, + last_error=f"connection refused: {reason_code}") return log.info("MQTT: connected") + with self._lock: + was_connected = self._diag["connected"] + self._diag.update(status="connected", connected=True, connected_at=time.time(), last_error=None) + if not was_connected: + self._event("mqtt_connected", f"MQTT connected to {self.cfg.mqtt_host}:{self.cfg.mqtt_port}") client.subscribe(self.status_topic) self._announce() if self._last_payload is not None: - client.publish(self.state_topic, json.dumps(self._last_payload), retain=True) + self._publish(self.state_topic, json.dumps(self._last_payload)) + + def _on_connect_fail(self, client, userdata): + with self._lock: + first = self._diag["connect_failures"] == 0 and not self._diag["connected"] + self._diag.update(status="failed", connected=False, last_attempt_at=time.time(), + last_error="connection attempt failed (unreachable, refused or unresolvable host)") + self._diag["connect_failures"] += 1 + log.warning("MQTT: connection to %s:%s failed, retrying", self.cfg.mqtt_host, self.cfg.mqtt_port) + if first: + self._event("mqtt_connect_failed", f"MQTT cannot reach {self.cfg.mqtt_host}:{self.cfg.mqtt_port}") + self._probe_in_background() def _on_disconnect(self, client, userdata, flags, reason_code, properties): log.warning("MQTT: disconnected: %s", reason_code) + self._set("disconnects", status="disconnected", connected=False, disconnected_at=time.time(), + last_error=f"disconnected: {reason_code}") + self._event("mqtt_disconnected", f"MQTT disconnected: {reason_code}") def _on_message(self, client, userdata, message): - if message.topic == self.status_topic and message.payload.decode(errors="ignore") == "online": - log.info("MQTT: Home Assistant came online, re-announcing") - self._announce() + if message.topic == self.status_topic: + status = message.payload.decode(errors="ignore") + self._set(ha_status=status, ha_status_at=time.time()) + if status == "online": + log.info("MQTT: Home Assistant came online, re-announcing") + self._announce() def _announce(self): device = { @@ -131,5 +247,8 @@ class HomeAssistantPublisher: **extra, } topic = f"{self.cfg.mqtt_discovery_prefix}/{component}/{self.node}/{key}/config" - self.client.publish(topic, json.dumps(payload), retain=True) - self.client.publish(self.availability_topic, "online", retain=True) + self._publish(topic, json.dumps(payload)) + self._set("discovery_count", discovery_at=time.time()) + info = self._publish(self.availability_topic, "online") + if info.rc == mqtt.MQTT_ERR_SUCCESS: + self._set(availability_at=time.time()) diff --git a/netprobe.py b/netprobe.py new file mode 100644 index 0000000..46ab84b --- /dev/null +++ b/netprobe.py @@ -0,0 +1,190 @@ +import random +import socket +import struct +import time + +RCODES = {0: "NOERROR", 1: "FORMERR", 2: "SERVFAIL", 3: "NXDOMAIN", 4: "NOTIMP", 5: "REFUSED"} +RESOLV_CONF = "/etc/resolv.conf" +HOSTS_FILE = "/etc/hosts" + + +def _ms(started): + return round((time.monotonic() - started) * 1000, 1) + + +def read_resolv_conf(path=RESOLV_CONF): + nameservers, search, options = [], [], [] + try: + with open(path) as fh: + for line in fh: + parts = line.split("#", 1)[0].split() + if not parts: + continue + if parts[0] == "nameserver" and len(parts) > 1: + nameservers.append(parts[1]) + elif parts[0] in ("search", "domain"): + search.extend(parts[1:]) + elif parts[0] == "options": + options.extend(parts[1:]) + except OSError as exc: + return {"error": str(exc), "nameservers": [], "search": [], "options": []} + return {"nameservers": nameservers, "search": search, "options": options} + + +def hosts_entries(name, path=HOSTS_FILE): + matches = [] + try: + with open(path) as fh: + for line in fh: + parts = line.split("#", 1)[0].split() + if len(parts) > 1 and name in parts[1:]: + matches.append(parts[0]) + except OSError: + pass + return matches + + +def _skip_name(data, offset): + while True: + length = data[offset] + if length == 0: + return offset + 1 + if length & 0xC0 == 0xC0: + return offset + 2 + offset += 1 + length + + +def _read_name(data, offset): + labels = [] + while True: + length = data[offset] + if length == 0: + return ".".join(labels) + if length & 0xC0 == 0xC0: + pointer = struct.unpack(">H", data[offset:offset + 2])[0] & 0x3FFF + return ".".join(labels + [_read_name(data, pointer)]) + offset += 1 + labels.append(data[offset:offset + length].decode(errors="replace")) + offset += length + + +def build_query(name, ident, qtype=1): + header = struct.pack(">HHHHHH", ident, 0x0100, 1, 0, 0, 0) + question = b"".join(bytes([len(label)]) + label.encode() for label in name.rstrip(".").split(".")) + return header + question + b"\x00" + struct.pack(">HH", qtype, 1) + + +def parse_response(data, ident): + if len(data) < 12: + raise ValueError("response too short") + rid, flags, qdcount, ancount, _, _ = struct.unpack(">HHHHHH", data[:12]) + if rid != ident: + raise ValueError("response id mismatch") + offset = 12 + for _ in range(qdcount): + offset = _skip_name(data, offset) + 4 + addresses, cnames = [], [] + for _ in range(ancount): + offset = _skip_name(data, offset) + rtype, _, _, rdlength = struct.unpack(">HHIH", data[offset:offset + 10]) + offset += 10 + if rtype == 1 and rdlength == 4: + addresses.append(socket.inet_ntoa(data[offset:offset + 4])) + elif rtype == 28 and rdlength == 16: + addresses.append(socket.inet_ntop(socket.AF_INET6, data[offset:offset + 16])) + elif rtype == 5: + cnames.append(_read_name(data, offset)) + offset += rdlength + rcode = flags & 0xF + return {"rcode": RCODES.get(rcode, str(rcode)), "addresses": addresses, "cnames": cnames} + + +def dns_query(server, name, timeout=2.0, qtype=1, port=53): + ident = random.randint(0, 0xFFFF) + family = socket.AF_INET6 if ":" in server else socket.AF_INET + started = time.monotonic() + result = {"server": server, "type": "AAAA" if qtype == 28 else "A"} + try: + with socket.socket(family, socket.SOCK_DGRAM) as sock: + sock.settimeout(timeout) + sock.sendto(build_query(name, ident, qtype), (server, port)) + data, _ = sock.recvfrom(4096) + result.update(parse_response(data, ident)) + except (OSError, ValueError) as exc: + result["error"] = str(exc) or type(exc).__name__ + result["ms"] = _ms(started) + return result + + +def resolve(host, port): + started = time.monotonic() + try: + infos = socket.getaddrinfo(host, port, proto=socket.IPPROTO_TCP) + addresses = sorted({info[4][0] for info in infos}) + return {"addresses": addresses, "ms": _ms(started)} + except OSError as exc: + return {"addresses": [], "error": str(exc), "ms": _ms(started)} + + +def tcp_connect(address, port, timeout): + started = time.monotonic() + try: + with socket.create_connection((address, port), timeout=timeout): + return {"address": address, "ok": True, "ms": _ms(started)} + except OSError as exc: + return {"address": address, "ok": False, "error": str(exc) or type(exc).__name__, "ms": _ms(started)} + + +def probe(host, port, timeout=3.0): + started = time.monotonic() + resolv = read_resolv_conf() + result = { + "ts": time.time(), + "host": host, + "port": port, + "resolv_conf": resolv, + "hosts_file": hosts_entries(host), + "getaddrinfo": resolve(host, port), + "nameservers": [dns_query(server, host, timeout=min(timeout, 2.0)) for server in resolv["nameservers"]], + "tcp": [], + } + candidates = list(result["getaddrinfo"]["addresses"]) + for answer in result["nameservers"]: + for address in answer.get("addresses", []): + if address not in candidates: + candidates.append(address) + result["tcp"] = [tcp_connect(address, port, timeout) for address in candidates] + result["ms"] = _ms(started) + result["summary"] = summarize(result) + return result + + +def summarize(result): + lines = [] + gai = result["getaddrinfo"] + if gai.get("error"): + lines.append(f"System resolver failed for {result['host']}: {gai['error']}") + else: + lines.append(f"System resolver: {result['host']} -> {', '.join(gai['addresses'])}") + if result["hosts_file"]: + lines.append(f"/etc/hosts maps it to {', '.join(result['hosts_file'])}") + if not result["resolv_conf"]["nameservers"]: + lines.append("No nameservers in /etc/resolv.conf") + for answer in result["nameservers"]: + if answer.get("error"): + lines.append(f"Nameserver {answer['server']}: no reply ({answer['error']})") + elif answer["addresses"]: + lines.append(f"Nameserver {answer['server']}: {answer['rcode']} -> {', '.join(answer['addresses'])}") + else: + lines.append(f"Nameserver {answer['server']}: {answer['rcode']}, no address") + answers = [set(a.get("addresses", [])) for a in result["nameservers"] if not a.get("error")] + if len(answers) > 1 and any(a != answers[0] for a in answers[1:]): + lines.append("Nameservers disagree; the result depends on which one the resolver asks first") + if not result["tcp"]: + lines.append("No address to try a TCP connection to") + for attempt in result["tcp"]: + if attempt["ok"]: + lines.append(f"TCP {attempt['address']}:{result['port']} connected in {attempt['ms']} ms") + else: + lines.append(f"TCP {attempt['address']}:{result['port']} failed: {attempt['error']}") + return lines diff --git a/templates/index.html b/templates/index.html index 6791f39..3ea9777 100644 --- a/templates/index.html +++ b/templates/index.html @@ -53,6 +53,14 @@ .kind.up, .kind.wifi_connected, .kind.public_ip { background: var(--ok); } .kind.down, .kind.wifi_disconnected, .kind.reboot { background: var(--bad); } .kind.public_ip_changed, .kind.wifi_roamed, .kind.service_started { background: var(--warn); } + .kind.mqtt_connected { background: var(--ok); } + .kind.mqtt_disconnected, .kind.mqtt_connect_failed { background: var(--bad); } + .btn { background: var(--panel); border: 1px solid var(--border); color: var(--text); padding: 4px 12px; border-radius: 6px; cursor: pointer; } + .btn:disabled { opacity: .5; cursor: default; } + ul.probe { list-style: none; margin: 0 0 8px; padding: 0; font-size: 13px; } + ul.probe li { padding: 4px 0; border-bottom: 1px solid var(--border); } + ul.probe li.bad { color: var(--bad); } ul.probe li.ok { color: var(--ok); } + pre { background: var(--bg); border: 1px solid var(--border); border-radius: 6px; padding: 10px 12px; font-size: 12px; overflow-x: auto; margin: 8px 0 0; } @media (max-width: 600px) { ul.timeline li { grid-template-columns: 1fr; gap: 2px; } } @@ -69,6 +77,7 @@
DNS
…
Public IP
…
Wi-Fi
…
+
MQTT
…

History

@@ -106,6 +115,30 @@

Wi-Fi retries / failed / beacon loss per interval

+

Home Assistant (MQTT)

+
+
+
+
+ Last state payload +

+    
+
+
+
+ Broker connection probe + + + + +
+
    +
    + Raw JSON +
    
    +    
    +
    +