import json import logging import threading import time import paho.mqtt.client as mqtt import netprobe log = logging.getLogger(__name__) PROBE_INTERVAL = 120 VERSION = "2.3.0" CONNECTIVITY = {"device_class": "connectivity"} SENSORS = [ ("internet_up", "binary_sensor", "Internet", CONNECTIVITY), ("internet_latency_ms", "sensor", "Internet latency", {"unit_of_measurement": "ms", "state_class": "measurement", "icon": "mdi:timer-outline"}), ("dns_up", "binary_sensor", "DNS", CONNECTIVITY), ("dns_latency_ms", "sensor", "DNS latency", {"unit_of_measurement": "ms", "state_class": "measurement", "icon": "mdi:dns"}), ("public_ip", "sensor", "Public IP", {"icon": "mdi:ip-network-outline"}), ("uptime_s", "sensor", "Host uptime", {"device_class": "duration", "unit_of_measurement": "s", "entity_category": "diagnostic"}), ] WIFI_SENSORS = [ ("wifi_connected", "binary_sensor", "Wi-Fi", CONNECTIVITY), ("wifi_ssid", "sensor", "Wi-Fi SSID", {"icon": "mdi:wifi"}), ("wifi_signal_dbm", "sensor", "Wi-Fi signal", {"device_class": "signal_strength", "unit_of_measurement": "dBm", "state_class": "measurement"}), ("wifi_tx_bitrate_mbps", "sensor", "Wi-Fi TX bitrate", {"device_class": "data_rate", "unit_of_measurement": "Mbit/s", "state_class": "measurement"}), ("wifi_bssid", "sensor", "Wi-Fi BSSID", {"icon": "mdi:access-point", "entity_category": "diagnostic"}), ("wifi_channel", "sensor", "Wi-Fi channel", {"icon": "mdi:radio-tower", "entity_category": "diagnostic"}), ("wifi_ap_count", "sensor", "Wi-Fi networks in range", {"icon": "mdi:access-point-network", "state_class": "measurement"}), ] def flatten(snapshot): link = snapshot.get("wifi_link") or {} scan = snapshot.get("wifi_scan") or {} return { "internet_up": snapshot.get("internet_up"), "internet_latency_ms": snapshot.get("internet_latency_ms"), "dns_up": snapshot.get("dns_up"), "dns_latency_ms": snapshot.get("dns_latency_ms"), "public_ip": snapshot.get("public_ip"), "uptime_s": snapshot.get("uptime_s"), "wifi_connected": link.get("connected"), "wifi_ssid": link.get("ssid"), "wifi_bssid": link.get("bssid"), "wifi_signal_dbm": link.get("signal_dbm"), "wifi_tx_bitrate_mbps": link.get("tx_bitrate_mbps"), "wifi_channel": link.get("channel"), "wifi_ap_count": scan.get("count"), } class HomeAssistantPublisher: 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._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: 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: 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 = { "identifiers": [self.node], "name": self.cfg.device_name, "model": "healthcheck-container", "manufacturer": "healthcheck", "sw_version": VERSION, } sensors = SENSORS + (WIFI_SENSORS if self.cfg.wifi_interface else []) for key, component, name, extra in sensors: if component == "binary_sensor": template = ( f"{{{{ 'None' if value_json.{key} is none " f"else ('ON' if value_json.{key} else 'OFF') }}}}" ) else: template = f"{{{{ value_json.{key} }}}}" payload = { "name": name, "unique_id": f"{self.node}_{key}", "state_topic": self.state_topic, "value_template": template, "availability_topic": self.availability_topic, "device": device, **extra, } topic = f"{self.cfg.mqtt_discovery_prefix}/{component}/{self.node}/{key}/config" 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())