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=", "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 = """""" def nav(active): p = 'active' if active == 'pruefstand' else '' h = 'active' if active == 'heizraum' else '' return f'' @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'{rom[-8:]}{s.get("temp","--")} C{crc}{s.get("err","0")}x' if not rows: rows = 'Warte auf Sensordaten...' ins = ''.join(f'IN{i} {x.get(f"pruefstand/in/{i}","?")}' for i in range(4)) rels = ''.join(f'Relais {i}: {x.get(f"pruefstand/relay/{i}","?")}' for i in range(2)) st = x.get('pruefstand/status','--'); sc = '#00cc44' if st=='online' else '#cc2200' return Response(f''' Pruefstand{CSS}{nav("pruefstand")}

Pruefstand

Status: {st} · {len(sns)} Sensor(en) · Refresh 4s

{rows}
ROMTempCRCErr

{ins}

{rels}

''', 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'{rom[-8:]}{temp} C' if not rows: rows = 'Warte auf Sensordaten...' ios = ''.join(f'{n}: {"AN" if str(x.get(f"heizraum/io/{n}"))=="1" else "AUS"}' 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''' Messknoten{CSS}{nav("heizraum")}

Messknoten (Heizraum)

Status: {st_txt} · {len(sens)} Sensor(en) · FW: {x.get("heizraum/state/fw","--")} · Refresh 4s

{rows}
ROMTemp

{ios}

Kessel: {"GESPERRT" if kg else "FREIGEGEBEN"}

''', content_type='text/html') if __name__ == '__main__': app.run(host='0.0.0.0', port=8765, debug=False)