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 <noreply@anthropic.com>
This commit is contained in:
1 parent
d80c40065e
commit
62b57a7e56
8 files changed
+737
-14
No files matched your search
@@ -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())
|
||||
Reference in new issue
Block a user