home-monitoring/collector.py
2026-07-31 09:17:22 +02:00

384 lines
15 KiB
Python

#!/usr/bin/env python3
import json
import logging
import math
import os
import queue
import re
import signal
import threading
import time
import urllib.error
import urllib.parse
import urllib.request
from collections import Counter
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
import paho.mqtt.client as mqtt
logging.basicConfig(level=os.getenv("LOG_LEVEL", "INFO"), format="%(asctime)s %(levelname)s %(message)s")
LOG = logging.getLogger("home-monitor")
MQTT_HOST = os.getenv("MQTT_HOST", "10.10.1.1")
MQTT_PORT = int(os.getenv("MQTT_PORT", "1883"))
MQTT_TOPICS = [x.strip() for x in os.getenv("MQTT_TOPICS", "#,$SYS/#").split(",") if x.strip()]
MQTT_USERNAME = os.getenv("MQTT_USERNAME", "")
MQTT_PASSWORD = os.getenv("MQTT_PASSWORD", "")
INFLUX_URL = os.getenv("INFLUX_URL", "http://influxdb:8086").rstrip("/")
INFLUX_ORG = os.getenv("INFLUX_ORG", "influx.sascha-lutz.de")
INFLUX_BUCKET = os.getenv("INFLUX_BUCKET", "telegraf")
INFLUX_TOKEN = os.getenv("INFLUX_TOKEN", "")
BUTLER_URL = os.getenv("BUTLER_URL", "http://10.4.1.116:8888").rstrip("/")
BUTLER_TOKEN = os.getenv("BUTLER_TOKEN", "")
WRITE_INTERVAL = int(os.getenv("WRITE_INTERVAL", "15"))
TOPIC_STATS_INTERVAL = int(os.getenv("TOPIC_STATS_INTERVAL", "300"))
HA_POLL_INTERVAL = int(os.getenv("HA_POLL_INTERVAL", "60"))
MAX_TOPICS = int(os.getenv("MAX_TOPICS", "500"))
HEALTH_PORT = int(os.getenv("HEALTH_PORT", "9199"))
SENSITIVE = re.compile(r"(^|[_.\-/])(token|secret|password|pass|credential|api[_-]?key|latitude|longitude|lat|lon|gps|location|ssid|bssid|mac|ip|user|person)([_.\-/]|$)", re.I)
BLOCKED_TOPIC = re.compile(r"(^|/)(camera|image|snapshot|video|audio|recording|clip)(/|$)", re.I)
BOOLS = {
"on": 1.0, "off": 0.0, "true": 1.0, "false": 0.0,
"online": 1.0, "offline": 0.0, "open": 1.0, "closed": 0.0,
"home": 1.0, "not_home": 0.0, "locked": 1.0, "unlocked": 0.0,
"available": 1.0, "unavailable": 0.0,
}
SAFE_HA_DOMAINS = {
"sensor", "binary_sensor", "switch", "light", "climate", "cover", "lock",
"fan", "humidifier", "input_number", "number", "water_heater", "sun",
}
stop_event = threading.Event()
write_queue: queue.Queue[str] = queue.Queue(maxsize=20000)
lock = threading.Lock()
state = {
"started": time.time(), "mqtt_up": 0, "messages_total": 0, "parse_errors": 0,
"numeric_samples": 0, "topics": {}, "roots": Counter(), "last_message": 0.0,
"influx_up": 0, "influx_writes": 0, "influx_errors": 0,
"ha_up": 0, "ha_http_status": 0, "ha_entities": 0, "ha_numeric": 0,
"ha_latency_ms": 0.0, "last_ha_poll": 0.0,
}
def esc_measurement(value):
return str(value).replace("\\", "\\\\").replace(",", "\\,").replace(" ", "\\ ")
def esc_tag(value):
return str(value).replace("\\", "\\\\").replace(",", "\\,").replace(" ", "\\ ").replace("=", "\\=")
def esc_field_key(value):
return esc_tag(value)
def field_value(value):
if isinstance(value, bool):
return "true" if value else "false"
if isinstance(value, int):
return f"{value}i"
if isinstance(value, float):
if not math.isfinite(value):
raise ValueError("non-finite float")
return repr(value)
text = str(value).replace("\\", "\\\\").replace('"', '\\"')
return f'"{text}"'
def line(measurement, tags, fields, timestamp=None):
tag_text = "".join(f",{esc_tag(k)}={esc_tag(v)}" for k, v in sorted(tags.items()) if str(v) != "")
field_text = ",".join(f"{esc_field_key(k)}={field_value(v)}" for k, v in sorted(fields.items()))
stamp = f" {int(timestamp)}" if timestamp is not None else ""
return f"{esc_measurement(measurement)}{tag_text} {field_text}{stamp}"
def safe_key(key):
key = re.sub(r"[^A-Za-z0-9_.-]+", "_", str(key)).strip("_.-")[:100]
return key and not SENSITIVE.search(key)
def flatten(obj, prefix="", depth=0):
out = {}
if depth > 2 or not isinstance(obj, dict):
return out
for raw_key, value in obj.items():
key = str(raw_key)
full = f"{prefix}.{key}" if prefix else key
if not safe_key(full):
continue
if isinstance(value, bool):
out[full] = 1.0 if value else 0.0
elif isinstance(value, (int, float)) and not isinstance(value, bool) and math.isfinite(float(value)):
out[full] = float(value)
elif isinstance(value, str) and value.strip().lower() in BOOLS:
out[full] = BOOLS[value.strip().lower()]
elif isinstance(value, dict):
out.update(flatten(value, full, depth + 1))
return out
def extract_values(payload):
if len(payload) > 65536:
return {}
try:
text = payload.decode("utf-8").strip()
except UnicodeDecodeError:
return {}
if not text:
return {}
low = text.lower()
if low in BOOLS:
return {"value": BOOLS[low]}
try:
value = float(text)
if math.isfinite(value):
return {"value": value}
except ValueError:
pass
if text[:1] == "{" and text[-1:] == "}":
try:
return flatten(json.loads(text))
except (json.JSONDecodeError, TypeError, ValueError):
return {}
return {}
def topic_allowed(topic):
if not topic or len(topic) > 250 or SENSITIVE.search(topic) or BLOCKED_TOPIC.search(topic):
return False
# Discovery documents contain names, identifiers and configuration, not time-series values.
if topic.startswith("homeassistant/") and topic.endswith("/config"):
return False
return True
def enqueue(metric_line):
try:
write_queue.put_nowait(metric_line)
except queue.Full:
with lock:
state["influx_errors"] += 1
def on_connect(client, userdata, flags, reason_code, properties=None):
ok = int(reason_code) == 0
with lock:
state["mqtt_up"] = 1 if ok else 0
if ok:
for topic in MQTT_TOPICS:
client.subscribe(topic, qos=0)
LOG.info("MQTT connected; %d read-only subscriptions active", len(MQTT_TOPICS))
else:
LOG.warning("MQTT connection rejected with code %s", reason_code)
def on_disconnect(client, userdata, disconnect_flags, reason_code, properties=None):
with lock:
state["mqtt_up"] = 0
if not stop_event.is_set():
LOG.warning("MQTT disconnected with code %s", reason_code)
def on_message(client, userdata, msg):
now = int(time.time())
topic = msg.topic
values = extract_values(msg.payload) if topic_allowed(topic) else {}
root = topic.split("/", 1)[0][:80]
with lock:
state["messages_total"] += 1
state["last_message"] = now
state["roots"][root] += 1
topics = state["topics"]
if topic in topics or len(topics) < MAX_TOPICS:
row = topics.setdefault(topic, {"messages": 0, "last_seen": 0, "retained": 0, "numeric": 0})
row["messages"] += 1
row["last_seen"] = now
row["retained"] = 1 if msg.retain else row["retained"]
row["numeric"] += len(values)
state["numeric_samples"] += len(values)
for field, value in values.items():
enqueue(line("mqtt_value", {"topic": topic, "root": root, "field": field}, {"value": value}, now))
if topic.startswith("$SYS/broker/") and values:
metric = topic.removeprefix("$SYS/broker/").replace("/", "_")
for value in values.values():
enqueue(line("mqtt_broker", {"metric": metric}, {"value": value}, now))
break
def post_influx(lines):
if not lines or not INFLUX_TOKEN:
return False
query = urllib.parse.urlencode({"org": INFLUX_ORG, "bucket": INFLUX_BUCKET, "precision": "s"})
req = urllib.request.Request(
f"{INFLUX_URL}/api/v2/write?{query}", data=("\n".join(lines)).encode(), method="POST",
headers={"Authorization": f"Token {INFLUX_TOKEN}", "Content-Type": "text/plain; charset=utf-8"},
)
try:
with urllib.request.urlopen(req, timeout=10) as response:
ok = response.status in (204, 200)
except (urllib.error.URLError, TimeoutError, OSError) as exc:
LOG.warning("InfluxDB write failed: %s", type(exc).__name__)
ok = False
with lock:
state["influx_up"] = 1 if ok else 0
state["influx_writes" if ok else "influx_errors"] += 1
return ok
def snapshot_metrics(include_topics=False):
now = int(time.time())
with lock:
s = dict(state)
roots = dict(state["roots"])
topics = {k: dict(v) for k, v in state["topics"].items()}
age = now - int(s["last_message"]) if s["last_message"] else -1
lines = [line("home_monitoring", {"source": "home-monitor"}, {
"mqtt_up": int(s["mqtt_up"]), "influx_up": int(s["influx_up"]),
"messages_total": int(s["messages_total"]), "topics_seen": len(topics),
"numeric_samples_total": int(s["numeric_samples"]), "parse_errors_total": int(s["parse_errors"]),
"last_message_age_seconds": age, "influx_writes_total": int(s["influx_writes"]),
"influx_errors_total": int(s["influx_errors"]), "homeassistant_up": int(s["ha_up"]),
"homeassistant_http_status": int(s["ha_http_status"]), "homeassistant_entities": int(s["ha_entities"]),
"homeassistant_numeric_entities": int(s["ha_numeric"]), "homeassistant_latency_ms": float(s["ha_latency_ms"]),
}, now)]
for root, count in roots.items():
lines.append(line("mqtt_root", {"root": root}, {"messages_total": int(count)}, now))
if include_topics:
for topic, row in topics.items():
lines.append(line("mqtt_topic", {"topic": topic, "root": topic.split('/', 1)[0]}, {
"messages_total": int(row["messages"]), "last_seen": int(row["last_seen"]),
"retained": int(row["retained"]), "numeric_samples_total": int(row["numeric"]),
}, now))
return lines
def poll_homeassistant():
started = time.monotonic()
url = f"{BUTLER_URL}/homeassistant/api/states"
req = urllib.request.Request(url, headers={"Authorization": f"Bearer {BUTLER_TOKEN}"})
status = 0
entities = []
try:
with urllib.request.urlopen(req, timeout=15) as response:
status = response.status
if status == 200:
entities = json.load(response)
except urllib.error.HTTPError as exc:
status = exc.code
except (urllib.error.URLError, TimeoutError, OSError):
status = 0
now = int(time.time())
numeric = 0
if status == 200 and isinstance(entities, list):
for entity in entities:
entity_id = str(entity.get("entity_id", ""))
if "." not in entity_id:
continue
domain = entity_id.split(".", 1)[0]
if domain not in SAFE_HA_DOMAINS or SENSITIVE.search(entity_id):
continue
raw = str(entity.get("state", "")).strip()
if raw.lower() in BOOLS:
value = BOOLS[raw.lower()]
else:
try:
value = float(raw)
if not math.isfinite(value):
continue
except ValueError:
continue
attrs = entity.get("attributes") or {}
unit = str(attrs.get("unit_of_measurement", ""))[:30]
device_class = str(attrs.get("device_class", ""))[:40]
enqueue(line("homeassistant_state", {
"entity_id": entity_id, "domain": domain, "unit": unit, "device_class": device_class,
}, {"value": value}, now))
numeric += 1
latency = (time.monotonic() - started) * 1000.0
with lock:
state["ha_up"] = 1 if status == 200 else 0
state["ha_http_status"] = status
state["ha_entities"] = len(entities)
state["ha_numeric"] = numeric
state["ha_latency_ms"] = latency
state["last_ha_poll"] = now
def worker():
last_status = 0.0
last_topics = 0.0
last_ha = 0.0
batch = []
while not stop_event.wait(1):
now = time.monotonic()
while len(batch) < 5000:
try:
batch.append(write_queue.get_nowait())
except queue.Empty:
break
if now - last_ha >= HA_POLL_INTERVAL:
poll_homeassistant()
last_ha = now
include_topics = now - last_topics >= TOPIC_STATS_INTERVAL
if include_topics:
last_topics = now
if now - last_status >= WRITE_INTERVAL:
batch.extend(snapshot_metrics(include_topics=include_topics))
post_influx(batch)
batch.clear()
last_status = now
class HealthHandler(BaseHTTPRequestHandler):
def do_GET(self):
if self.path not in ("/health", "/"):
self.send_error(404)
return
with lock:
payload = {
"status": "ok" if state["influx_up"] else "degraded",
"mqtt_up": bool(state["mqtt_up"]), "influx_up": bool(state["influx_up"]),
"homeassistant_up": bool(state["ha_up"]), "homeassistant_http_status": state["ha_http_status"],
"topics_seen": len(state["topics"]), "messages_total": state["messages_total"],
"numeric_samples_total": state["numeric_samples"],
"uptime_seconds": int(time.time() - state["started"]),
}
body = json.dumps(payload, separators=(",", ":")).encode()
self.send_response(200 if payload["influx_up"] else 503)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
def log_message(self, fmt, *args):
return
def main():
if not INFLUX_TOKEN or not BUTLER_TOKEN:
raise SystemExit("INFLUX_TOKEN and BUTLER_TOKEN are required")
threading.Thread(target=worker, name="writer", daemon=True).start()
server = ThreadingHTTPServer(("0.0.0.0", HEALTH_PORT), HealthHandler)
threading.Thread(target=server.serve_forever, name="health", daemon=True).start()
client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, client_id="pfannkuchen-home-monitor", clean_session=True)
if MQTT_USERNAME:
client.username_pw_set(MQTT_USERNAME, MQTT_PASSWORD)
client.on_connect = on_connect
client.on_disconnect = on_disconnect
client.on_message = on_message
client.reconnect_delay_set(min_delay=2, max_delay=60)
client.connect_async(MQTT_HOST, MQTT_PORT, keepalive=60)
for sig in (signal.SIGTERM, signal.SIGINT):
signal.signal(sig, lambda *_: stop_event.set())
client.loop_start()
LOG.info("Home monitor started: MQTT read-only, raw payload persistence disabled")
while not stop_event.wait(1):
pass
client.disconnect()
client.loop_stop()
server.shutdown()
if __name__ == "__main__":
main()