homelab-brain/infra/heizung-api/pruefstand_web.py
Homelab Cursor 5161f429df Ölknoten: Grafana-Zeile auf Heizung-Dashboard + MCP/API-Fix
Heizungs-API und Hermes-MCP liefern Ölknoten-Daten aus MQTT/sensors-Influx
statt fälschlich offline zu melden wenn nur ioBroker mqtt.0 tot ist.
2026-07-01 15:07:39 +02:00

420 lines
21 KiB
Python

import threading, time, json, re
from urllib.request import urlopen
from urllib.parse import quote
from flask import Flask, Response, request, jsonify
import paho.mqtt.client as mqtt
app = Flask(__name__)
d = {}
lock = threading.Lock()
last_seen = {}
sensor_ts = {}
SENSOR_TIMEOUT = 60
INFLUX_URL = "http://192.168.178.36:8086"
INFLUX_DB = "iobroker"
INFLUX_DB_SENSORS = "sensors"
OELKNOTEN_SENSORS = {
"raum": "Raum",
"hk_vorlauf": "HK Vorlauf",
"vl_erd_an": "Erdtrasse VL",
"rl_erd_ab": "Erdtrasse RL",
}
# Sensor-Rollen (ROM-ID -> Klartext). Erweiterbar sobald Zuordnung bekannt.
SENSOR_ROLLEN = {
"28BE941600000093": "Heizraum (historisiert)",
}
# Bekannte Heiz-Measurements fuer Discovery
HEIZ_MEASUREMENTS = {
"brennerlaufzeit": "Brenner-Laufzeit pro Lauf (Minuten)",
"brennerstarts": "Brennerstart-Ereignis (1=Start)",
"brennerstatus": "Brenner an/aus (1/0)",
"brenner_heute": "Tageswerte: liter, starts, stunden",
"mqtt.0.Oelkessel.Oelkessel_VL.Vorlauf": "Oelkessel Vorlauftemperatur (C)",
"mqtt.0.Holzvergaser_Sensoren_6.temp1.temperature1": "Holzvergaser Temp 1 (C)",
"mqtt.0.Holzvergaser_Sensoren_6.temp2.temperature2": "Holzvergaser Temp 2 (C)",
"mqtt.0.heizraum.io.brenner": "Heizraum-Knoten Brenner-Signal",
"mqtt.0.heizraum.io.pumpe1": "Pumpe 1 (1/0)",
"mqtt.0.heizraum.io.pumpe2": "Pumpe 2 (1/0)",
"mqtt.0.heizraum.io.pumpe3": "Pumpe 3 (1/0)",
"mqtt.0.heizraum.io.pumpe4": "Pumpe 4 (1/0)",
"mqtt.0.heizraum.state.kessel": "Kessel gesperrt (1/0)",
}
# ── InfluxDB Helper ─────────────────────────────────────────────────────────
def influx_query(q, epoch="ms", db=INFLUX_DB):
url = f"{INFLUX_URL}/query?db={db}&epoch={epoch}&q={quote(q)}"
try:
with urlopen(url, timeout=8) as r:
return json.loads(r.read())
except Exception as e:
return {"error": str(e)}
def influx_series(raw):
"""Extrahiert (columns, values) aus erster Serie, robust gegen Fehler."""
try:
s = raw["results"][0]["series"][0]
return s.get("columns", []), s.get("values", [])
except Exception:
return [], []
def oelknoten_from_mqtt(x, ts, now):
"""Live-Temperaturen oelknoten aus MQTT-Cache."""
temps = {}
for sensor in OELKNOTEN_SENSORS:
topic = f"oelknoten/{sensor}"
if ts.get(topic, 0) >= now - SENSOR_TIMEOUT:
try:
temps[sensor] = round(float(x.get(topic, "")), 2)
except (TypeError, ValueError):
pass
age = now - last_seen.get("oelknoten", 0) if last_seen.get("oelknoten") else None
online = age is not None and age < 90
if x.get("oelknoten/status") == "online":
online = True
return {
"online": online,
"status": x.get("oelknoten/status"),
"firmware": x.get("oelknoten/state/fw"),
"letzte_aktivitaet_vor_sek": round(age, 1) if age else None,
"temperaturen": temps,
}
def oelknoten_from_influx():
"""Fallback: frische Werte aus InfluxDB sensors (mqtt-influx-bridge)."""
raw = influx_query(
"SELECT last(value) FROM temperature WHERE node='oelknoten' AND time > now()-5m GROUP BY sensor",
db=INFLUX_DB_SENSORS,
)
temps = {}
try:
for series in raw.get("results", [{}])[0].get("series", []) or []:
sensor = series.get("tags", {}).get("sensor")
vals = series.get("values", [])
if sensor and vals and vals[0][1] is not None:
temps[sensor] = round(float(vals[0][1]), 2)
except Exception:
pass
st = influx_series(influx_query(
"SELECT last(value) FROM node_status WHERE node='oelknoten' AND time > now()-10m",
db=INFLUX_DB_SENSORS,
))
online = bool(temps) or (st[1] and st[1][0][1] == 1)
return {
"online": online,
"status": "online" if online else "offline",
"firmware": None,
"letzte_aktivitaet_vor_sek": None,
"temperaturen": temps,
"quelle": "influx_sensors",
}
# ── MQTT ──────────────────────────────────────────────────────────────────────
def on_msg(c, u, m):
topic = m.topic
val = m.payload.decode('utf-8', 'replace')
with lock:
d[topic] = val
if topic.startswith('heizraum/'):
last_seen['heizraum'] = time.time()
if topic.startswith('oelknoten/'):
last_seen['oelknoten'] = time.time()
if topic.count('/') == 1 and topic.split('/')[1] in OELKNOTEN_SENSORS:
sensor_ts[topic] = time.time()
if topic.startswith('pruefstand/'):
last_seen['pruefstand'] = time.time()
if '/sensor/' in topic:
parts = topic.split('/')
if len(parts) == 4:
sensor_ts[parts[0] + '/' + parts[2]] = time.time()
if topic.startswith('heizraum/28'):
sensor_ts[topic] = time.time()
def mq():
c = mqtt.Client()
c.on_message = on_msg
c.connect('192.168.178.36', 1883, 60)
c.subscribe('pruefstand/#')
c.subscribe('heizraum/#')
c.subscribe('oelknoten/#')
c.loop_forever()
threading.Thread(target=mq, daemon=True).start()
# ════════════════════════════════════════════════════════════════════════════
# API
# ════════════════════════════════════════════════════════════════════════════
@app.route('/api/hilfe')
def api_hilfe():
return jsonify({
"basis_url": "http://100.109.101.12:8765",
"beschreibung": "Heizungs- und PV-Daten des Homelabs. Live via MQTT, Historie via InfluxDB.",
"endpunkte": {
"GET /api/hilfe": "Diese Uebersicht",
"GET /api/status": "System-Uebersicht (online/offline aller Knoten)",
"GET /api/heizung/now": "Live: Temperaturen, Pumpen, Brenner, Kessel (MQTT)",
"GET /api/heizung/history": "Temp-Verlauf. Params: ?hours=24&sensor=<ROM|all>",
"GET /api/heizung/brenner": "Brenner Takt-Analyse. Params: ?days=7",
"GET /api/heizung/kessel": "Oelkessel Vorlauf-Verlauf. Params: ?hours=48",
"GET /api/pv/now": "Live PV: Ertrag heute, Netz, Verbrauch",
"GET /api/pv/history": "PV-Tageserträge. Params: ?days=14",
"GET /api/measurements": "Alle verfuegbaren InfluxDB-Measurements + Zeitraum",
"GET /api/query": "GENERISCH: beliebige read-only InfluxQL. Param: ?q=SELECT... (nur SELECT/SHOW erlaubt)",
},
"influxdb": {
"hinweis": "Fuer eigene Analysen /api/query nutzen. Beispiel:",
"beispiel": "/api/query?q=SELECT count(value) FROM brennerstarts WHERE time > now()-7d",
"db": INFLUX_DB,
},
"wichtig": [
"Daten erst ab Dez 2025 / Jan 2026 (DB-Migration vom alten Raspi).",
"Aeltere Logs (Jahre) noch nicht importiert - liegen auf alter Raspi-SD-Karte.",
"Brenner-Daten (brennerlaufzeit/starts/status) enden 2026-05-18.",
"Heizraum-ROM-IDs noch ohne vollstaendige Rollenzuordnung.",
],
})
@app.route('/api/status')
def api_status():
now = time.time()
with lock:
x = dict(d)
ts = dict(sensor_ts)
ls = dict(last_seen)
hz_age = now - ls.get('heizraum', 0) if ls.get('heizraum') else None
ps_age = now - ls.get('pruefstand', 0) if ls.get('pruefstand') else None
hz_sensors = [k for k in x if k.startswith('heizraum/28') and ts.get(k, 0) >= now - SENSOR_TIMEOUT]
ok_mqtt = oelknoten_from_mqtt(x, ts, now)
ok = ok_mqtt if ok_mqtt.get("temperaturen") else oelknoten_from_influx()
ok_age = now - ls.get('oelknoten', 0) if ls.get('oelknoten') else None
return jsonify({
"timestamp_utc": time.strftime('%Y-%m-%dT%H:%M:%SZ', time.gmtime()),
"heizraum_knoten": {
"online": hz_age is not None and hz_age < 90,
"letzte_aktivitaet_vor_sek": round(hz_age, 1) if hz_age else None,
"aktive_sensoren": len(hz_sensors),
},
"oelknoten": {
"online": ok.get("online", False),
"letzte_aktivitaet_vor_sek": round(ok_age, 1) if ok_age else ok.get("letzte_aktivitaet_vor_sek"),
"aktive_sensoren": len(ok.get("temperaturen", {})),
"firmware": ok.get("firmware"),
},
"pruefstand": {
"online": ps_age is not None and ps_age < 90,
"letzte_aktivitaet_vor_sek": round(ps_age, 1) if ps_age else None,
},
"influxdb_erreichbar": "error" not in influx_query("SHOW DATABASES"),
"datenquellen": {"influxdb": INFLUX_URL, "mqtt": "192.168.178.36:1883"},
})
@app.route('/api/heizung/now')
def api_heizung_now():
now = time.time()
with lock:
x = dict(d)
ts = dict(sensor_ts)
ls = dict(last_seen)
sensors = {}
for k, v in x.items():
if k.startswith('heizraum/28') and ts.get(k, 0) >= now - SENSOR_TIMEOUT:
rom = k.split('/')[1]
try:
sensors[rom] = {"temp_c": round(float(v), 2), "rolle": SENSOR_ROLLEN.get(rom, "unbekannt")}
except:
sensors[rom] = {"temp_c": v, "rolle": SENSOR_ROLLEN.get(rom, "unbekannt")}
age = now - ls.get('heizraum', 0) if ls.get('heizraum') else None
heizraum_online = age is not None and age < 90
ok = oelknoten_from_mqtt(x, ts, now)
if not ok.get("temperaturen"):
ok = oelknoten_from_influx()
return jsonify({
"timestamp_utc": time.strftime('%Y-%m-%dT%H:%M:%SZ', time.gmtime()),
"online": heizraum_online or ok.get("online", False),
"heizraum_online": heizraum_online,
"oelknoten": ok,
"temperaturen": sensors,
"pumpen": {n: x.get(f'heizraum/io/{n}') == '1' for n in ['pumpe1','pumpe2','pumpe3','pumpe4']},
"brenner_an": x.get('heizraum/io/brenner') == '1',
"kessel_gesperrt": x.get('heizraum/state/kessel') == '1',
"hinweis": "mosquitto+mqtt-influx-bridge laufen unabhaengig vom ioBroker-Adapter mqtt.0",
})
@app.route('/api/heizung/history')
def api_heizung_history():
hours = int(request.args.get('hours', 24))
sensor = request.args.get('sensor', '28BE941600000093')
if sensor == 'all':
cols, vals = influx_series(influx_query("SHOW MEASUREMENTS"))
sns = [v[0] for v in vals if 'heizraum' in v[0] or 'Oelkessel' in v[0] or 'Holzvergaser' in v[0]]
return jsonify({"verfuegbare_sensoren": sns})
meas = f"mqtt.0.heizraum.{sensor}" if len(sensor) >= 12 and not sensor.startswith('mqtt') else sensor
raw = influx_query(f'SELECT mean(value) FROM "{meas}" WHERE time > now()-{hours}h GROUP BY time(5m) fill(previous)')
cols, vals = influx_series(raw)
pts = [{"ts_ms": v[0], "temp_c": round(v[1], 2)} for v in vals if v[1] is not None]
return jsonify({"sensor": sensor, "measurement": meas, "zeitraum_h": hours,
"punkte": len(pts), "daten": pts[-300:]})
@app.route('/api/heizung/brenner')
def api_heizung_brenner():
days = int(request.args.get('days', 7))
out = {"zeitraum_tage": days, "timestamp_utc": time.strftime('%Y-%m-%dT%H:%M:%SZ', time.gmtime())}
# Starts (Takten)
raw = influx_query(f'SELECT count(value) FROM brennerstarts WHERE time > now()-{days}d')
_, v = influx_series(raw)
out["starts_gesamt"] = v[0][1] if v else 0
# Starts pro Tag
raw = influx_query(f'SELECT count(value) FROM brennerstarts WHERE time > now()-{days}d GROUP BY time(1d) tz(\'Europe/Berlin\')')
_, v = influx_series(raw)
out["starts_pro_tag"] = [{"ts_ms": r[0], "starts": r[1] or 0} for r in v]
# Laufzeit-Summe
raw = influx_query(f'SELECT sum(value) FROM brennerlaufzeit WHERE time > now()-{days}d')
_, v = influx_series(raw)
out["laufzeit_summe_min"] = round(v[0][1], 1) if v and v[0][1] else 0
# Mittlere Laufzeit pro Start (Takt-Indikator)
if out["starts_gesamt"] and out["laufzeit_summe_min"]:
out["mittlere_laufzeit_pro_start_min"] = round(out["laufzeit_summe_min"] / out["starts_gesamt"], 1)
out["hinweis"] = "Kurze mittlere Laufzeit + viele Starts = Takten (ineffizient). Daten enden 2026-05-18."
return jsonify(out)
@app.route('/api/heizung/kessel')
def api_heizung_kessel():
hours = int(request.args.get('hours', 48))
raw = influx_query(f'SELECT mean(value) FROM "mqtt.0.Oelkessel.Oelkessel_VL.Vorlauf" WHERE time > now()-{hours}h GROUP BY time(10m) fill(previous)')
_, v = influx_series(raw)
pts = [{"ts_ms": r[0], "vorlauf_c": round(r[1], 2)} for r in v if r[1] is not None]
temps = [p["vorlauf_c"] for p in pts]
return jsonify({
"zeitraum_h": hours, "punkte": len(pts),
"min_c": min(temps) if temps else None,
"max_c": max(temps) if temps else None,
"daten": pts[-300:],
})
@app.route('/api/pv/now')
def api_pv_now():
res = {}
for name, meas in [("ertrag_kwh_heute","mqtt.1.openWB.pv.DailyYieldKwh"),
("netz_watt","mqtt.1.openWB.evu.W"),
("verbrauch_wh","mqtt.1.openWB.global.WHouseConsumption")]:
_, v = influx_series(influx_query(f'SELECT last(value) FROM "{meas}"'))
res[name] = round(float(v[0][1]), 2) if v else None
nw = res.get('netz_watt') or 0
return jsonify({"timestamp_utc": time.strftime('%Y-%m-%dT%H:%M:%SZ', time.gmtime()),
**res, "einspeisung": nw < 0, "netzbezug": nw > 0})
@app.route('/api/pv/history')
def api_pv_history():
days = int(request.args.get('days', 14))
raw = influx_query(f'SELECT max(value)-min(value) AS e FROM "mqtt.1.openWB.pv.DailyYieldKwh" WHERE time > now()-{days}d GROUP BY time(1d) fill(0) tz(\'Europe/Berlin\')')
_, v = influx_series(raw)
tage = [{"datum": time.strftime('%Y-%m-%d', time.gmtime(r[0]/1000)), "ertrag_kwh": round(r[1],2)}
for r in v if r[1] and r[1] > 0]
return jsonify({"zeitraum_tage": days, "gesamt_kwh": round(sum(t["ertrag_kwh"] for t in tage),2), "tage": tage})
@app.route('/api/measurements')
def api_measurements():
cols, vals = influx_series(influx_query("SHOW MEASUREMENTS"))
return jsonify({"db": INFLUX_DB, "measurements": [v[0] for v in vals],
"bekannte_heizung": HEIZ_MEASUREMENTS})
# ── GENERISCHER read-only InfluxQL-Passthrough ──────────────────────────────
_FORBIDDEN = re.compile(r'\b(DROP|DELETE|INSERT|ALTER|CREATE|GRANT|REVOKE|UPDATE|KILL)\b', re.IGNORECASE)
@app.route('/api/query')
def api_query():
q = request.args.get('q', '').strip()
if not q:
return jsonify({"error": "Parameter q fehlt. Beispiel: ?q=SELECT count(value) FROM brennerstarts WHERE time > now()-7d"}), 400
if not re.match(r'^\s*(SELECT|SHOW)\b', q, re.IGNORECASE):
return jsonify({"error": "Nur SELECT- und SHOW-Abfragen erlaubt (read-only)."}), 403
if _FORBIDDEN.search(q):
return jsonify({"error": "Verbotenes Schluesselwort. Nur read-only."}), 403
raw = influx_query(q)
if "error" in raw:
return jsonify(raw), 502
return jsonify(raw)
# ════════════════════════════════════════════════════════════════════════════
# Web UI
# ════════════════════════════════════════════════════════════════════════════
CSS = """<style>
*{box-sizing:border-box}
body{background:#111;color:#eee;font-family:monospace;padding:10px;max-width:640px;margin:0 auto}
h2{margin:8px 0 4px}
table{width:100%;border-collapse:collapse}
th{background:#1a1a2e;padding:8px 6px;text-align:left;color:#88aacc;font-size:14px}
td{padding:9px 6px;border-bottom:1px solid #222;font-size:14px}
.ok{color:#00dd55}.err{color:#ff4422}
.tag{display:inline-block;padding:6px 14px;border-radius:5px;margin:4px 4px 0 0;color:#fff;font-size:15px}
nav{display:flex;gap:8px;margin-bottom:10px}
nav a{flex:1;text-align:center;padding:12px 6px;background:#1a1a2e;color:#88ccff;
text-decoration:none;border-radius:6px;font-size:16px;border:1px solid #334}
nav a.active{background:#003366;color:#fff;border-color:#66aaff}
.tempval{font-size:26px;font-weight:bold}
.status{font-size:13px;color:#666;margin-bottom:8px}
</style>"""
def nav(active):
p = 'active' if active == 'pruefstand' else ''
h = 'active' if active == 'heizraum' else ''
return f'<nav><a href="/" class="{p}">Pruefstand</a><a href="/heizraum" class="{h}">Messknoten</a></nav>'
@app.route('/')
def pruefstand():
now = time.time()
with lock:
x = dict(d); ts = dict(sensor_ts)
sns = {}
for k, v in x.items():
if '/sensor/' in k:
parts = k.split('/')
if len(parts) == 4 and ts.get(parts[0]+'/'+parts[2], 0) >= now - SENSOR_TIMEOUT:
sns.setdefault(parts[2], {})[parts[3]] = v
rows = ''
for rom, s in sorted(sns.items()):
crc = s.get('crc','--'); cls = 'ok' if crc=='ok' else 'err'
rows += f'<tr><td style="color:#888">{rom[-8:]}</td><td class="{cls} tempval">{s.get("temp","--")} C</td><td class="{cls}">{crc}</td><td style="color:#888">{s.get("err","0")}x</td></tr>'
if not rows: rows = '<tr><td colspan=4 style="color:#444">Warte auf Sensordaten...</td></tr>'
ins = ''.join(f'<span class="tag" style="background:{"#005500" if x.get(f"pruefstand/in/{i}")=="HIGH" else "#550000"}">IN{i} {x.get(f"pruefstand/in/{i}","?")}</span>' for i in range(4))
rels = ''.join(f'<span class="tag" style="background:{"#005500" if x.get(f"pruefstand/relay/{i}")=="ON" else "#222"}">Relais {i}: {x.get(f"pruefstand/relay/{i}","?")}</span>' for i in range(2))
st = x.get('pruefstand/status','--'); sc = '#00cc44' if st=='online' else '#cc2200'
return Response(f'''<!DOCTYPE html><html><head><meta charset=utf-8><meta name=viewport content="width=device-width,initial-scale=1">
<meta http-equiv=refresh content=4><title>Pruefstand</title>{CSS}</head><body>{nav("pruefstand")}
<h2>Pruefstand</h2><p class=status>Status: <b style="color:{sc}">{st}</b> · {len(sns)} Sensor(en) · Refresh 4s</p>
<table><tr><th>ROM</th><th>Temp</th><th>CRC</th><th>Err</th></tr>{rows}</table>
<p style="margin-top:12px">{ins}</p><p>{rels}</p></body></html>''', content_type='text/html')
@app.route('/heizraum')
def heizraum():
now = time.time()
with lock:
x = dict(d); ts = dict(sensor_ts)
sens = {k.split('/')[1]: v for k, v in x.items()
if k.startswith('heizraum/28') and len(k.split('/'))==2 and ts.get(k,0) >= now - SENSOR_TIMEOUT}
rows = ''
for rom, temp in sorted(sens.items()):
try: cls = 'ok' if -10 < float(temp) < 110 else 'err'
except: cls = 'err'
rows += f'<tr><td style="color:#888">{rom[-8:]}</td><td class="{cls} tempval">{temp} C</td></tr>'
if not rows: rows = '<tr><td colspan=2 style="color:#444">Warte auf Sensordaten...</td></tr>'
ios = ''.join(f'<span class="tag" style="background:{"#005500" if str(x.get(f"heizraum/io/{n}"))=="1" else "#222"}">{n}: {"AN" if str(x.get(f"heizraum/io/{n}"))=="1" else "AUS"}</span>' for n in ['pumpe1','pumpe2','pumpe3','pumpe4','brenner'])
kg = str(x.get('heizraum/state/kessel'))=='1'
ls = last_seen.get('heizraum', 0); age = now - ls if ls else 9999
st_txt, st_col = ('unbekannt','#888') if ls==0 else (('online','#00cc44') if age<90 else ('offline','#cc2200'))
return Response(f'''<!DOCTYPE html><html><head><meta charset=utf-8><meta name=viewport content="width=device-width,initial-scale=1">
<meta http-equiv=refresh content=4><title>Messknoten</title>{CSS}</head><body>{nav("heizraum")}
<h2>Messknoten (Heizraum)</h2><p class=status>Status: <b style="color:{st_col}">{st_txt}</b> · {len(sens)} Sensor(en) · FW: {x.get("heizraum/state/fw","--")} · Refresh 4s</p>
<table><tr><th>ROM</th><th>Temp</th></tr>{rows}</table><p style="margin-top:12px">{ios}</p>
<p><span class="tag" style="background:{"#660000" if kg else "#005500"};font-size:16px">Kessel: {"GESPERRT" if kg else "FREIGEGEBEN"}</span></p></body></html>''', content_type='text/html')
if __name__ == '__main__':
app.run(host='0.0.0.0', port=8765, debug=False)