Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
62b57a7e56 | ||
|
|
d80c40065e |
No files matched your search
+2
-2
@@ -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
|
||||
|
||||
@@ -16,13 +16,18 @@ 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 |
|
||||
|
||||
## Running
|
||||
|
||||
The compose file pulls the published multi-arch image
|
||||
`hiimmilan/health-check` (amd64, arm64, arm/v7):
|
||||
|
||||
```sh
|
||||
docker compose up -d --build
|
||||
docker compose pull && docker compose up -d
|
||||
```
|
||||
|
||||
Then open `http://<device>:9999/`.
|
||||
@@ -86,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
|
||||
|
||||
@@ -97,6 +103,23 @@ availability to `healthcheck/<device_id>/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)
|
||||
@@ -111,3 +134,10 @@ Entities per device:
|
||||
python3 -m unittest discover -s tests
|
||||
DB_PATH=./data/hc.db python3 app.py
|
||||
```
|
||||
|
||||
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.3.0 --push .
|
||||
```
|
||||
@@ -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()
|
||||
|
||||
|
||||
+1
-2
@@ -1,7 +1,6 @@
|
||||
services:
|
||||
healthcheck:
|
||||
build: .
|
||||
image: healthcheck-container
|
||||
image: hiimmilan/health-check:latest
|
||||
container_name: healthcheck
|
||||
restart: unless-stopped
|
||||
# Host networking exposes the Wi-Fi interface to the container. The app
|
||||
|
||||
@@ -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())
|
||||
+190
@@ -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
|
||||
+114
-1
@@ -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; } }
|
||||
</style>
|
||||
</head>
|
||||
@@ -69,6 +77,7 @@
|
||||
<div class="card" id="card-dns"><div class="label">DNS</div><div class="value">…</div><div class="sub"></div></div>
|
||||
<div class="card" id="card-ip"><div class="label">Public IP</div><div class="value">…</div><div class="sub"></div></div>
|
||||
<div class="card" id="card-wifi"><div class="label">Wi-Fi</div><div class="value">…</div><div class="sub"></div></div>
|
||||
<div class="card" id="card-mqtt"><div class="label">MQTT</div><div class="value">…</div><div class="sub"></div></div>
|
||||
</section>
|
||||
|
||||
<h2>History</h2>
|
||||
@@ -106,6 +115,30 @@
|
||||
<div class="panel chart"><h3>Wi-Fi retries / failed / beacon loss per interval</h3><canvas id="ch-errors"></canvas><div class="legend" id="lg-errors"></div></div>
|
||||
</div>
|
||||
|
||||
<h2>Home Assistant (MQTT)</h2>
|
||||
<div class="panel">
|
||||
<div id="mqtt-notice"></div>
|
||||
<div id="mqtt-details"></div>
|
||||
<details id="mqtt-payload-wrap" style="margin-top:12px">
|
||||
<summary style="cursor:pointer;color:var(--muted)">Last state payload</summary>
|
||||
<pre id="mqtt-payload"></pre>
|
||||
</details>
|
||||
</div>
|
||||
<div class="panel" id="probe-panel" style="margin-top:12px">
|
||||
<div style="display:flex;justify-content:space-between;align-items:center;flex-wrap:wrap;gap:8px;margin-bottom:8px">
|
||||
<strong>Broker connection probe</strong>
|
||||
<span style="display:flex;align-items:center;gap:10px">
|
||||
<span class="meta" id="probe-meta" style="color:var(--muted);font-size:12px"></span>
|
||||
<button id="probe-run" class="btn">Run probe now</button>
|
||||
</span>
|
||||
</div>
|
||||
<ul id="probe-summary" class="probe"></ul>
|
||||
<details id="probe-json-wrap">
|
||||
<summary style="cursor:pointer;color:var(--muted)">Raw JSON</summary>
|
||||
<pre id="probe-json"></pre>
|
||||
</details>
|
||||
</div>
|
||||
|
||||
</main>
|
||||
|
||||
<script>
|
||||
@@ -198,6 +231,86 @@ async function loadWifi() {
|
||||
'</tbody></table>';
|
||||
}
|
||||
|
||||
const MQTT_STATUS = {
|
||||
connected: ['ok', 'connected'], connecting: ['warn', 'connecting'], failed: ['bad', 'unreachable'],
|
||||
refused: ['bad', 'refused'], disconnected: ['bad', 'disconnected'], not_started: ['', 'not started'],
|
||||
};
|
||||
|
||||
async function loadMqtt() {
|
||||
const m = await (await fetch('/api/mqtt')).json();
|
||||
if (!m.enabled) {
|
||||
card('card-mqtt', 'warn', 'not configured', 'set MQTT_HOST to publish to Home Assistant');
|
||||
$('mqtt-notice').innerHTML = '';
|
||||
$('mqtt-details').innerHTML = '<div class="empty">MQTT is disabled. Set <code>MQTT_HOST</code> to enable Home Assistant publishing.</div>';
|
||||
$('mqtt-payload-wrap').hidden = true;
|
||||
$('probe-panel').hidden = true;
|
||||
return;
|
||||
}
|
||||
const [cls, label] = MQTT_STATUS[m.status] || ['', m.status];
|
||||
const since = m.connected ? m.connected_at : (m.disconnected_at || m.last_attempt_at);
|
||||
card('card-mqtt', cls, label,
|
||||
`${esc(m.host)}:${m.port}` + (since ? ` · ${m.connected ? 'since' : 'for'} ${ago(since, m.now).replace(' ago', '')}` : '') +
|
||||
(m.last_publish_at ? ` · published ${ago(m.last_publish_at, m.now)}` : ''),
|
||||
m.connected ? m.last_publish_error : m.last_error);
|
||||
|
||||
let notice = '';
|
||||
if (!m.connected) notice = `Not connected to the broker (${esc(m.last_error || m.status)}). Home Assistant shows the device as unavailable because the broker sent the retained "offline" will message and nothing has replaced it.`;
|
||||
else if (m.publish_dropped && m.last_publish_error) notice = `Connected, but the last publish failed: ${esc(m.last_publish_error)}.`;
|
||||
else if (!m.availability_at) notice = 'Connected, but availability has not been published yet.';
|
||||
else if (!m.last_payload) notice = 'Connected and discovery sent, waiting for the first monitor result to publish state.';
|
||||
$('mqtt-notice').innerHTML = notice ? `<div class="notice">${notice}</div>` : '';
|
||||
|
||||
const when = (ts) => ts ? `${fmtTime(ts)} (${ago(ts, m.now)})` : 'never';
|
||||
const rows = [
|
||||
['Broker', `${m.host}:${m.port}` + (m.username ? ` as ${m.username}` : ' (no auth)')],
|
||||
['Client id', m.client_id],
|
||||
['Status', label + (m.last_error ? ` · ${m.last_error}` : '')],
|
||||
['Publisher started', when(m.started_at)],
|
||||
['Connected', when(m.connected_at)],
|
||||
['Disconnected', when(m.disconnected_at) + (m.disconnects ? ` · ${m.disconnects} total` : '')],
|
||||
['Connect failures', String(m.connect_failures) + (m.connect_failures ? ` · last attempt ${when(m.last_attempt_at)}` : '')],
|
||||
['Discovery sent', when(m.discovery_at) + (m.discovery_count ? ` · ${m.discovery_count} times` : '')],
|
||||
['Availability "online" sent', when(m.availability_at)],
|
||||
['Home Assistant status', m.ha_status ? `${m.ha_status} at ${when(m.ha_status_at)}` : 'nothing received on ' + m.discovery_prefix + '/status'],
|
||||
['Messages published', `${m.publish_count} ok · ${m.publish_dropped} dropped` + (m.last_publish_error ? ` · last error: ${m.last_publish_error}` : '')],
|
||||
['Last publish', when(m.last_publish_at)],
|
||||
['State topic', m.state_topic],
|
||||
['Availability topic', m.availability_topic],
|
||||
['Discovery prefix', m.discovery_prefix],
|
||||
];
|
||||
$('mqtt-details').innerHTML = '<dl class="kv">' + rows.map(([k, v]) => `<dt>${k}</dt><dd>${esc(v)}</dd>`).join('') + '</dl>';
|
||||
$('mqtt-payload-wrap').hidden = !m.last_payload;
|
||||
$('mqtt-payload').textContent = m.last_payload ? JSON.stringify(m.last_payload, null, 2) : '';
|
||||
$('probe-panel').hidden = false;
|
||||
renderProbe(m.probe, m.now, m.probe_running);
|
||||
}
|
||||
|
||||
function renderProbe(p, now, running) {
|
||||
$('probe-run').disabled = !!running;
|
||||
$('probe-run').textContent = running ? 'Running…' : 'Run probe now';
|
||||
if (!p) {
|
||||
$('probe-meta').textContent = 'not run yet; runs automatically when a connection attempt fails';
|
||||
$('probe-summary').innerHTML = ''; $('probe-json-wrap').hidden = true; return;
|
||||
}
|
||||
$('probe-meta').textContent = `ran ${fmtTime(p.ts)} (${ago(p.ts, now)})` + (p.ms != null ? ` in ${p.ms} ms` : '');
|
||||
$('probe-summary').innerHTML = p.summary.map(line => {
|
||||
const cls = /fail|no reply|NXDOMAIN|SERVFAIL|REFUSED|disagree|No address|No nameservers|crashed/.test(line) ? 'bad' : /connected in/.test(line) ? 'ok' : '';
|
||||
return `<li class="${cls}">${esc(line)}</li>`;
|
||||
}).join('');
|
||||
$('probe-json-wrap').hidden = false;
|
||||
$('probe-json').textContent = JSON.stringify(p, null, 2);
|
||||
}
|
||||
|
||||
$('probe-run').addEventListener('click', async () => {
|
||||
renderProbe(null, null, true);
|
||||
try {
|
||||
const r = await fetch('/api/mqtt/probe', { method: 'POST' });
|
||||
const p = await r.json();
|
||||
if (!r.ok) { $('probe-meta').textContent = p.error || 'probe failed'; $('probe-run').disabled = false; $('probe-run').textContent = 'Run probe now'; return; }
|
||||
renderProbe(p, Date.now() / 1000, false);
|
||||
} catch (e) { console.error(e); $('probe-run').disabled = false; $('probe-run').textContent = 'Run probe now'; }
|
||||
});
|
||||
|
||||
async function loadEvents() {
|
||||
const events = await (await fetch('/api/events?limit=300')).json();
|
||||
$('timeline').innerHTML = events.length ? events.map(e =>
|
||||
@@ -270,7 +383,7 @@ $('range').addEventListener('click', (e) => {
|
||||
});
|
||||
window.addEventListener('resize', () => loadCharts());
|
||||
|
||||
async function refreshFast() { try { await Promise.all([loadStatus(), loadWifi(), loadEvents()]); } catch (e) { console.error(e); } }
|
||||
async function refreshFast() { try { await Promise.all([loadStatus(), loadWifi(), loadMqtt(), loadEvents()]); } catch (e) { console.error(e); } }
|
||||
refreshFast(); loadCharts();
|
||||
setInterval(refreshFast, 10000);
|
||||
setInterval(loadCharts, 60000);
|
||||
|
||||
@@ -0,0 +1,140 @@
|
||||
import time
|
||||
import unittest
|
||||
from types import SimpleNamespace
|
||||
from unittest import mock
|
||||
|
||||
import paho.mqtt.client as mqtt
|
||||
|
||||
import hamqtt
|
||||
|
||||
|
||||
def make_cfg(**overrides):
|
||||
values = dict(device_id="pi", device_name="Pi", mqtt_host="broker", mqtt_port=1883,
|
||||
mqtt_username="", mqtt_password="", mqtt_discovery_prefix="homeassistant",
|
||||
wifi_interface="")
|
||||
values.update(overrides)
|
||||
return SimpleNamespace(**values)
|
||||
|
||||
|
||||
class FakeStorage:
|
||||
def __init__(self):
|
||||
self.events = []
|
||||
|
||||
def add_event(self, kind, message):
|
||||
self.events.append((kind, message))
|
||||
|
||||
|
||||
def ok_publish(*args, **kwargs):
|
||||
return SimpleNamespace(rc=mqtt.MQTT_ERR_SUCCESS)
|
||||
|
||||
|
||||
def no_conn_publish(*args, **kwargs):
|
||||
return SimpleNamespace(rc=mqtt.MQTT_ERR_NO_CONN)
|
||||
|
||||
|
||||
class PublisherDiagnosticsTest(unittest.TestCase):
|
||||
def setUp(self):
|
||||
patcher = mock.patch("hamqtt.mqtt.Client")
|
||||
self.client_cls = patcher.start()
|
||||
self.addCleanup(patcher.stop)
|
||||
self.client = self.client_cls.return_value
|
||||
self.client.publish.side_effect = ok_publish
|
||||
self.storage = FakeStorage()
|
||||
self.publisher = hamqtt.HomeAssistantPublisher(make_cfg(), self.storage)
|
||||
|
||||
def test_initial_state(self):
|
||||
diag = self.publisher.diagnostics()
|
||||
self.assertTrue(diag["enabled"])
|
||||
self.assertEqual(diag["status"], "not_started")
|
||||
self.assertFalse(diag["connected"])
|
||||
self.assertEqual(diag["host"], "broker")
|
||||
self.assertEqual(diag["client_id"], "healthcheck_pi")
|
||||
self.assertEqual(diag["state_topic"], "healthcheck/pi/state")
|
||||
self.assertIsNone(diag["username"])
|
||||
|
||||
def test_start_marks_connecting(self):
|
||||
self.publisher.start()
|
||||
self.assertEqual(self.publisher.diagnostics()["status"], "connecting")
|
||||
self.assertEqual(self.client.on_connect_fail, self.publisher._on_connect_fail)
|
||||
|
||||
def test_connect_announces_and_records(self):
|
||||
self.publisher.publish_state({"internet_up": True})
|
||||
self.publisher._on_connect(self.client, None, {}, SimpleNamespace(is_failure=False), None)
|
||||
diag = self.publisher.diagnostics()
|
||||
self.assertEqual(diag["status"], "connected")
|
||||
self.assertTrue(diag["connected"])
|
||||
self.assertIsNotNone(diag["connected_at"])
|
||||
self.assertEqual(diag["discovery_count"], 1)
|
||||
self.assertIsNotNone(diag["availability_at"])
|
||||
self.assertEqual(diag["publish_dropped"], 0)
|
||||
self.assertGreater(diag["publish_count"], len(hamqtt.SENSORS))
|
||||
self.assertEqual(diag["last_payload"]["internet_up"], True)
|
||||
self.assertEqual(self.storage.events[0][0], "mqtt_connected")
|
||||
|
||||
def test_connect_fail_records_first_event_only(self):
|
||||
self.publisher._on_connect_fail(self.client, None)
|
||||
self.publisher._on_connect_fail(self.client, None)
|
||||
diag = self.publisher.diagnostics()
|
||||
self.assertEqual(diag["status"], "failed")
|
||||
self.assertEqual(diag["connect_failures"], 2)
|
||||
self.assertIn("failed", diag["last_error"])
|
||||
self.assertEqual([k for k, _ in self.storage.events], ["mqtt_connect_failed"])
|
||||
|
||||
def test_refused(self):
|
||||
self.publisher._on_connect(self.client, None, {}, SimpleNamespace(is_failure=True), None)
|
||||
diag = self.publisher.diagnostics()
|
||||
self.assertEqual(diag["status"], "refused")
|
||||
self.assertEqual(diag["connect_failures"], 1)
|
||||
self.assertFalse(diag["connected"])
|
||||
|
||||
def test_disconnect(self):
|
||||
self.publisher._on_connect(self.client, None, {}, SimpleNamespace(is_failure=False), None)
|
||||
self.publisher._on_disconnect(self.client, None, {}, "Keep alive timeout", None)
|
||||
diag = self.publisher.diagnostics()
|
||||
self.assertEqual(diag["status"], "disconnected")
|
||||
self.assertEqual(diag["disconnects"], 1)
|
||||
self.assertIn("Keep alive timeout", diag["last_error"])
|
||||
self.assertEqual(self.storage.events[-1][0], "mqtt_disconnected")
|
||||
|
||||
def test_dropped_publish_is_counted(self):
|
||||
self.client.publish.side_effect = no_conn_publish
|
||||
self.publisher.publish_state({"internet_up": False})
|
||||
diag = self.publisher.diagnostics()
|
||||
self.assertEqual(diag["publish_count"], 0)
|
||||
self.assertEqual(diag["publish_dropped"], 1)
|
||||
self.assertIsNotNone(diag["last_publish_error"])
|
||||
self.assertIsNone(diag["last_publish_at"])
|
||||
|
||||
def test_ha_status_message(self):
|
||||
message = SimpleNamespace(topic="homeassistant/status", payload=b"online")
|
||||
self.publisher._on_message(self.client, None, message)
|
||||
diag = self.publisher.diagnostics()
|
||||
self.assertEqual(diag["ha_status"], "online")
|
||||
self.assertEqual(diag["discovery_count"], 1)
|
||||
|
||||
def test_connect_fail_runs_probe_once_per_interval(self):
|
||||
probe = {"ts": time.time(), "summary": ["System resolver failed"]}
|
||||
with mock.patch("hamqtt.netprobe.probe", return_value=probe) as probe_fn, \
|
||||
mock.patch("hamqtt.threading.Thread") as thread_cls:
|
||||
thread_cls.return_value.start.side_effect = lambda: self.publisher.run_probe()
|
||||
self.publisher._on_connect_fail(self.client, None)
|
||||
self.publisher._on_connect_fail(self.client, None)
|
||||
probe_fn.assert_called_once_with("broker", 1883)
|
||||
diag = self.publisher.diagnostics()
|
||||
self.assertEqual(diag["probe"], probe)
|
||||
self.assertFalse(diag["probe_running"])
|
||||
|
||||
def test_run_probe_survives_crash(self):
|
||||
with mock.patch("hamqtt.netprobe.probe", side_effect=RuntimeError("boom")):
|
||||
result = self.publisher.run_probe()
|
||||
self.assertIn("boom", result["error"])
|
||||
self.assertFalse(self.publisher.diagnostics()["probe_running"])
|
||||
|
||||
def test_publish_state_dedupes(self):
|
||||
self.publisher.publish_state({"internet_up": True})
|
||||
self.publisher.publish_state({"internet_up": True})
|
||||
self.assertEqual(self.publisher.diagnostics()["publish_count"], 1)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,125 @@
|
||||
import os
|
||||
import socket
|
||||
import struct
|
||||
import tempfile
|
||||
import threading
|
||||
import unittest
|
||||
|
||||
import netprobe
|
||||
|
||||
|
||||
def answer(ident, name, rcode=0, addresses=(), cname=None):
|
||||
header = struct.pack(">HHHHHH", ident, 0x8180 | rcode, 1, len(addresses) + (1 if cname else 0), 0, 0)
|
||||
question = netprobe.build_query(name, ident)[12:]
|
||||
records = b""
|
||||
if cname:
|
||||
target = b"".join(bytes([len(l)]) + l.encode() for l in cname.split(".")) + b"\x00"
|
||||
records += b"\xc0\x0c" + struct.pack(">HHIH", 5, 1, 60, len(target)) + target
|
||||
for address in addresses:
|
||||
records += b"\xc0\x0c" + struct.pack(">HHIH", 1, 1, 60, 4) + socket.inet_aton(address)
|
||||
return header + question + records
|
||||
|
||||
|
||||
class ParseResponseTest(unittest.TestCase):
|
||||
def test_addresses(self):
|
||||
data = answer(7, "broker.lan", addresses=["192.168.0.10", "192.168.0.11"])
|
||||
parsed = netprobe.parse_response(data, 7)
|
||||
self.assertEqual(parsed["rcode"], "NOERROR")
|
||||
self.assertEqual(parsed["addresses"], ["192.168.0.10", "192.168.0.11"])
|
||||
|
||||
def test_nxdomain(self):
|
||||
parsed = netprobe.parse_response(answer(1, "nope.invalid", rcode=3), 1)
|
||||
self.assertEqual(parsed["rcode"], "NXDOMAIN")
|
||||
self.assertEqual(parsed["addresses"], [])
|
||||
|
||||
def test_cname(self):
|
||||
parsed = netprobe.parse_response(answer(2, "ha.example", cname="real.example", addresses=["10.0.0.1"]), 2)
|
||||
self.assertEqual(parsed["cnames"], ["real.example"])
|
||||
self.assertEqual(parsed["addresses"], ["10.0.0.1"])
|
||||
|
||||
def test_id_mismatch(self):
|
||||
with self.assertRaises(ValueError):
|
||||
netprobe.parse_response(answer(3, "x.y"), 4)
|
||||
|
||||
|
||||
class DnsQueryTest(unittest.TestCase):
|
||||
def test_against_fake_server(self):
|
||||
server = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
||||
server.bind(("127.0.0.1", 0))
|
||||
port = server.getsockname()[1]
|
||||
|
||||
def serve():
|
||||
data, addr = server.recvfrom(512)
|
||||
ident = struct.unpack(">H", data[:2])[0]
|
||||
server.sendto(answer(ident, "broker.lan", addresses=["192.168.0.10"]), addr)
|
||||
server.close()
|
||||
|
||||
threading.Thread(target=serve, daemon=True).start()
|
||||
result = netprobe.dns_query("127.0.0.1", "broker.lan", timeout=2, port=port)
|
||||
self.assertEqual(result["addresses"], ["192.168.0.10"])
|
||||
self.assertEqual(result["rcode"], "NOERROR")
|
||||
|
||||
|
||||
class ResolvConfTest(unittest.TestCase):
|
||||
def test_parse(self):
|
||||
with tempfile.NamedTemporaryFile("w", delete=False) as fh:
|
||||
fh.write("# comment\nnameserver 192.168.0.1\nnameserver 8.8.8.8 # public\nsearch lan home\noptions ndots:1\n")
|
||||
try:
|
||||
parsed = netprobe.read_resolv_conf(fh.name)
|
||||
finally:
|
||||
os.unlink(fh.name)
|
||||
self.assertEqual(parsed["nameservers"], ["192.168.0.1", "8.8.8.8"])
|
||||
self.assertEqual(parsed["search"], ["lan", "home"])
|
||||
self.assertEqual(parsed["options"], ["ndots:1"])
|
||||
|
||||
def test_missing(self):
|
||||
parsed = netprobe.read_resolv_conf("/nonexistent/resolv.conf")
|
||||
self.assertEqual(parsed["nameservers"], [])
|
||||
self.assertIn("error", parsed)
|
||||
|
||||
def test_hosts(self):
|
||||
with tempfile.NamedTemporaryFile("w", delete=False) as fh:
|
||||
fh.write("127.0.0.1 localhost\n192.168.0.5 ha.local homeassistant\n")
|
||||
try:
|
||||
self.assertEqual(netprobe.hosts_entries("homeassistant", fh.name), ["192.168.0.5"])
|
||||
self.assertEqual(netprobe.hosts_entries("other", fh.name), [])
|
||||
finally:
|
||||
os.unlink(fh.name)
|
||||
|
||||
|
||||
class TcpTest(unittest.TestCase):
|
||||
def test_connect_ok_and_refused(self):
|
||||
listener = socket.socket()
|
||||
listener.bind(("127.0.0.1", 0))
|
||||
listener.listen(1)
|
||||
port = listener.getsockname()[1]
|
||||
try:
|
||||
ok = netprobe.tcp_connect("127.0.0.1", port, 2)
|
||||
finally:
|
||||
listener.close()
|
||||
self.assertTrue(ok["ok"])
|
||||
refused = netprobe.tcp_connect("127.0.0.1", port, 2)
|
||||
self.assertFalse(refused["ok"])
|
||||
self.assertIn("error", refused)
|
||||
|
||||
|
||||
class SummaryTest(unittest.TestCase):
|
||||
def test_disagreement_flagged(self):
|
||||
result = {
|
||||
"host": "ha.milans.cloud", "port": 1883, "hosts_file": [],
|
||||
"resolv_conf": {"nameservers": ["192.168.0.1", "8.8.8.8"]},
|
||||
"getaddrinfo": {"addresses": [], "error": "Name does not resolve"},
|
||||
"nameservers": [
|
||||
{"server": "192.168.0.1", "rcode": "NOERROR", "addresses": ["192.168.0.20"]},
|
||||
{"server": "8.8.8.8", "rcode": "NXDOMAIN", "addresses": []},
|
||||
],
|
||||
"tcp": [{"address": "192.168.0.20", "ok": True, "ms": 2.0}],
|
||||
}
|
||||
lines = netprobe.summarize(result)
|
||||
self.assertTrue(any("disagree" in line for line in lines))
|
||||
self.assertTrue(any("System resolver failed" in line for line in lines))
|
||||
self.assertTrue(any("connected in" in line for line in lines))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
Reference in new issue
Block a user