from netplan_v4_publisher import publish_forecasts as _v4_publish_forecasts from model_isolation import collect_predictions import os import sqlite3 import json import time import datetime import traceback import warnings import subprocess import sys import threading import numpy as np import pandas as pd import pytz from influxdb_client import InfluxDBClient, Point, WritePrecision from influxdb_client.client.write_api import SYNCHRONOUS from soc_diagnostics import battery_soc_points from telemetry_quality import require_recent_telemetry, sanitize_measured_frame try: from influxdb_client.client.warnings import MissingPivotFunction warnings.simplefilter("ignore", MissingPivotFunction) except Exception: pass import methods.var_1 as v1 import methods.var_2 as v2 import methods.var_3 as v3 import methods.var_10 as v10 import methods.var_11 as v11 import methods.var_13 as v13 import methods.var_21 as v21 import methods.var_22 as v22 import methods.var_23 as v23 INFLUX_URL = os.getenv("INFLUX_URL", "http://influxdb:8086") INFLUX_TOKEN = os.environ["INFLUX_TOKEN"] INFLUX_ORG = os.getenv("INFLUX_ORG", "belevo") INFLUX_BUCKET = os.getenv("INFLUX_BUCKET", "energy_data") SQLITE_DB_PATH = os.getenv("SQLITE_DB_PATH", "/app/data/users.db") HISTORY_START = os.getenv("FORECAST_HISTORY_START", "1970-01-01T00:00:00Z") QUALITY_LOOKBACK_DAYS = int(os.getenv("FORECAST_QUALITY_LOOKBACK_DAYS", "14")) LOCAL_TZ = pytz.timezone(os.getenv("TZ", "Europe/Zurich")) INFLUX_TIMEOUT_MS = int(os.getenv("FORECAST_INFLUX_TIMEOUT_MS", "120000")) FORECAST_TIMEOUT_SECONDS = int(os.getenv("FORECAST_RUN_TIMEOUT_SECONDS", "900")) FORECAST_HORIZON_HOURS = max(24, min(72, int(os.getenv("FORECAST_HORIZON_HOURS", "48")))) TRAINING_TIMEOUT_SECONDS = int(os.getenv("FORECAST_TRAINING_TIMEOUT_SECONDS", "3600")) FORECAST_STALE_SECONDS = int(os.getenv("FORECAST_STALE_SECONDS", "5400")) WATCHDOG_INTERVAL_SECONDS = int(os.getenv("FORECAST_WATCHDOG_INTERVAL_SECONDS", "300")) TRAINING_HOUR = int(os.getenv("FORECAST_TRAINING_HOUR", "2")) LAST_SUCCESS_PATH = os.getenv("FORECAST_LAST_SUCCESS_PATH", "/tmp/forecast_engine_last_success.json") LAST_TRAINING_PATH = os.getenv("FORECAST_LAST_TRAINING_PATH", "/app/data/forecast_training_status.json") last_trained_day = None _FORECAST_PROCESS_LOCK = threading.Lock() _TRAINING_THREAD = None MODEL_MODULES = {1: v1, 2: v2, 3: v3, 10: v10, 11: v11, 13: v13, 21: v21, 22: v22, 23: v23} QUALITY_TARGETS = {1: "PV", 10: "PV", 21: "PV", 2: "Hausverbrauch", 11: "Hausverbrauch", 22: "Hausverbrauch"} def active(config, n): return bool(int(config.get(f"prog_var_{n}", config.get(f"var_{n}", 0)) or 0)) def get_configs(): conn = sqlite3.connect(SQLITE_DB_PATH) conn.row_factory = sqlite3.Row columns = {row["name"] for row in conn.execute("PRAGMA table_info(anlagen_meta)")} if "batt_grid_charging_enabled" not in columns: conn.execute( "ALTER TABLE anlagen_meta ADD COLUMN batt_grid_charging_enabled INTEGER NOT NULL DEFAULT 0" ) conn.commit() rows = conn.execute("SELECT * FROM anlagen_meta").fetchall() conn.close() configs = [] for r in rows: d = dict(r) try: d["daecher"] = json.loads(d.get("daecher") or "[]") except Exception: d["daecher"] = [] for key, default in [ ("ac_leistung", 10.0), ("batt_capacity_kwh", 0.0), ("batt_power_kw", 0.0), ("tarif_bezug_fest", 0.30), ("tarif_einspeisung_fest", 0.10), ("tarif_peak_fest", 5.0), ]: raw_value = d.get(key) d[key] = float(default if raw_value is None or raw_value == "" else raw_value) configs.append(d) return configs def _time_literal(dt): return dt.strftime("%Y-%m-%dT%H:%M:%SZ") def _query_df(query): client = InfluxDBClient( url=INFLUX_URL, token=INFLUX_TOKEN, org=INFLUX_ORG, timeout=INFLUX_TIMEOUT_MS, ) try: df = client.query_api().query_data_frame(org=INFLUX_ORG, query=query) finally: client.close() if isinstance(df, list): df = pd.concat(df, ignore_index=True) if df else pd.DataFrame() if df is None or df.empty or "_time" not in df.columns: return pd.DataFrame() df["_time"] = pd.to_datetime(df["_time"], utc=True).dt.tz_localize(None) return df def _pivot_frame(df, fields): if df.empty: return pd.DataFrame() out = df.set_index("_time") keep = [c for c in fields if c in out.columns] out = out[keep] if keep else pd.DataFrame(index=out.index) out = out.apply(pd.to_numeric, errors="coerce") out = out[~out.index.duplicated(keep="last")] return out.sort_index().resample("5min").mean(numeric_only=True) def _tariff_frame(df, config=None): if df is None or df.empty: return pd.DataFrame() df = df.copy() if "_time" in df.columns: df["_time"] = pd.to_datetime(df["_time"], errors="coerce") df = df.dropna(subset=["_time"]).set_index("_time") if df.empty: return pd.DataFrame() price = df["price_chf_kwh"] if "price_chf_kwh" in df.columns else df.get("_value") if price is None: return pd.DataFrame() price = pd.to_numeric(price, errors="coerce") model = df.get("tariff_model", pd.Series("", index=df.index)).astype(str).str.lower() typ = df.get("type", pd.Series("", index=df.index)).astype(str).str.lower() tariff_name = df.get("tariff_name", pd.Series("", index=df.index)).astype(str).str.lower() provider = df.get("provider", pd.Series("", index=df.index)).astype(str).str.lower() key = (model + " " + tariff_name + " " + provider).str.lower() rows = pd.DataFrame({"price": price, "key": key, "type": typ}, index=df.index) rows = rows[pd.notna(rows["price"])].sort_index() if rows.empty: return pd.DataFrame() cfg = config or {} import_choice = str(cfg.get("tarif_bezug", "") or "").lower() export_choice = str(cfg.get("tarif_einspeisung", "") or "").lower() out = pd.DataFrame(index=rows.index.unique().sort_values()) import_base = ( rows["type"].str.contains("consumption|import|bezug", regex=True, na=False) | rows["key"].str.contains("dynamic|dynamisch|home|business", regex=True, na=False) ) if "business" in import_choice: import_mask = import_base & rows["key"].str.contains("business|gewerbe|commercial", regex=True, na=False) elif "home" in import_choice or "privat" in import_choice: import_mask = import_base & rows["key"].str.contains("home|privat|private", regex=True, na=False) elif "dynam" in import_choice: import_mask = import_base else: import_mask = pd.Series(False, index=rows.index) if not import_mask.any() and "dynam" in import_choice: import_mask = import_base export_base = ( rows["type"].str.contains("feed|einspeis|export", regex=True, na=False) | rows["key"].str.contains("referenzmarktpreis|marktpreis|reference|feed", regex=True, na=False) ) if "referenz" in export_choice or "marktpreis" in export_choice or "market" in export_choice: export_mask = export_base & rows["key"].str.contains("referenzmarktpreis|marktpreis|reference|belevo", regex=True, na=False) else: export_mask = rows["key"].str.contains("standard_feedin|ckw statisch", regex=True, na=False) & export_base if not export_mask.any() and ("referenz" in export_choice or "marktpreis" in export_choice or "market" in export_choice): export_mask = export_base if import_mask.any(): out["import_price"] = rows.loc[import_mask, "price"].groupby(level=0).last() if export_mask.any(): out["export_price"] = rows.loc[export_mask, "price"].groupby(level=0).last() if out.empty: return out return out.sort_index().resample("5min").mean().ffill().bfill() def fetch_influx_frames(config, training): aid = config["anlagen_id"] start = HISTORY_START if training else "-14d" future_stop = _time_literal(datetime.datetime.utcnow() + datetime.timedelta(hours=FORECAST_HORIZON_HOURS)) q_tel = f''' from(bucket: "{INFLUX_BUCKET}") |> range(start: {start}) |> filter(fn: (r) => r["_measurement"] == "api_telemetry") |> filter(fn: (r) => r["anlagen_id"] == "{aid}") |> filter(fn: (r) => r["_field"] == "PV" or r["_field"] == "Hausverbrauch" or r["_field"] == "Netzleistung" or r["_field"] == "SOC") |> filter(fn: (r) => not exists r["data_type"] or (r["data_type"] != "forecast" and r["data_type"] != "forecast_snapshot")) |> aggregateWindow(every: 5m, fn: mean, createEmpty: false, timeSrc: "_start") |> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value") ''' df_tel = _pivot_frame(_query_df(q_tel), ["PV", "Hausverbrauch", "Netzleistung", "SOC"]) q_wea = f''' from(bucket: "{INFLUX_BUCKET}") |> range(start: {start}, stop: {future_stop}) |> filter(fn: (r) => r["_measurement"] == "weather_forecast") |> aggregateWindow(every: 5m, fn: mean, createEmpty: false) |> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value") ''' df_wea = _pivot_frame(_query_df(q_wea), ["temp_c", "temperature", "cloud", "cloud_cover", "precip_mm", "wind_kph", "chance_of_snow"]) if "temperature" in df_wea.columns and "temp_c" not in df_wea.columns: df_wea["temp_c"] = df_wea["temperature"] if "cloud_cover" in df_wea.columns and "cloud" not in df_wea.columns: df_wea["cloud"] = df_wea["cloud_cover"] if "cloud" in df_wea.columns and "cloud_cover" not in df_wea.columns: df_wea["cloud_cover"] = df_wea["cloud"] q_tar = f''' from(bucket: "{INFLUX_BUCKET}") |> range(start: -7d, stop: {future_stop}) |> filter(fn: (r) => r["_measurement"] == "tariffs") |> filter(fn: (r) => r["_field"] == "price_chf_kwh") ''' df_tar = _tariff_frame(_query_df(q_tar), config) return df_tel, df_wea, df_tar def _consistent_tail(df): if df.empty or "PV" not in df.columns or "Hausverbrauch" not in df.columns: return df probe = df[["PV", "Hausverbrauch"]].copy() filled = probe.interpolate(limit=3, limit_direction="both") valid = filled.notna().all(axis=1) if not valid.any(): return df.iloc[0:0] run = 0 last_break = -1 for i, ok in enumerate(valid.to_numpy()): if ok: run = 0 else: run += 1 if run >= 4: last_break = i if last_break >= 0: after = np.where(valid.iloc[last_break + 1:].to_numpy())[0] if len(after): return df.iloc[last_break + 1 + after[0]:] return df.iloc[0:0] return df.loc[valid[valid].index[0]:] def _longest_consistent_segment(df, column): if df.empty or column not in df.columns: return pd.DataFrame() series = pd.to_numeric(df[column], errors="coerce") valid = series.notna().to_numpy() if not valid.any(): return pd.DataFrame() best_start = best_end = None start = 0 gap = 0 for i, ok in enumerate(valid): if ok: gap = 0 else: gap += 1 if gap >= 4: end = i - gap if end >= start and (best_start is None or end - start > best_end - best_start): best_start, best_end = start, end start = i + 1 gap = 0 end = len(valid) - 1 if end >= start and (best_start is None or end - start > best_end - best_start): best_start, best_end = start, end if best_start is None: return pd.DataFrame() segment = df.iloc[best_start:best_end + 1].copy() first = pd.to_numeric(segment[column], errors="coerce").first_valid_index() last = pd.to_numeric(segment[column], errors="coerce").last_valid_index() if first is None or last is None: return pd.DataFrame() return segment.loc[first:last] def _add_time_features(df): hour = df.index.hour + df.index.minute / 60.0 doy = df.index.dayofyear df["hour_float"] = hour df["hour_sin"] = np.sin(2 * np.pi * hour / 24.0) df["hour_cos"] = np.cos(2 * np.pi * hour / 24.0) df["sin_year"] = np.sin(2 * np.pi * doy / 365.25) df["cos_year"] = np.cos(2 * np.pi * doy / 365.25) df["weekday"] = df.index.weekday df["is_weekday"] = (df.index.weekday < 5).astype(int) return df def _fill_defaults(df, history): defaults = { "PV": np.nan, "Hausverbrauch": np.nan, "SOC": np.nan, "Netzleistung": np.nan, "temp_c": 15.0, "cloud": 20.0, "cloud_cover": 20.0, "precip_mm": 0.0, "wind_kph": 0.0, "chance_of_snow": 0.0, "import_price": np.nan, "export_price": np.nan, } for col, default in defaults.items(): if col not in df.columns: df[col] = default df[col] = pd.to_numeric(df[col], errors="coerce") # Measurements are never interpolated/filled here. Outages are not zero load, # zero PV, a fresh SOC, or a measured grid peak. Weather defaults are separate. df = sanitize_measured_frame(df) for col, default in defaults.items(): if col not in ["PV", "Hausverbrauch", "SOC", "Netzleistung"]: df[col] = df[col].ffill().bfill() if np.isfinite(default): df[col] = df[col].fillna(default) if "cloud_cover" in df.columns: df["cloud"] = df["cloud"].fillna(df["cloud_cover"]) df["import_price"] = df["import_price"].fillna(0.30) df["export_price"] = df["export_price"].fillna(0.10) return _add_time_features(df) def build_data_object(config, training=False): now = datetime.datetime.utcnow().replace(second=0, microsecond=0) now = now - datetime.timedelta(minutes=now.minute % 5) df_tel, df_wea, df_tar = fetch_influx_frames(config, training) frames = [f for f in [df_tel, df_wea, df_tar] if not f.empty] combined = frames[0] if frames else pd.DataFrame() for f in frames[1:]: combined = combined.join(f, how="outer") combined = combined.sort_index() full_hist_raw = combined.loc[:now - datetime.timedelta(minutes=5)] if not combined.empty else pd.DataFrame() # A trailing telemetry gap must not discard all earlier valid observations. hist_raw = full_hist_raw.copy() # Freshness is determined from measurements, not the outer-joined weather grid. recent_raw = sanitize_measured_frame(df_tel.loc[:now - datetime.timedelta(minutes=5)].copy()) if not df_tel.empty else pd.DataFrame() start_hist = hist_raw.index.min().floor("5min") if training and not hist_raw.empty else now - datetime.timedelta(days=14) idx_hist = pd.date_range(start=start_hist, end=now - datetime.timedelta(minutes=5), freq="5min") df_hist = pd.DataFrame(index=idx_hist).join(hist_raw, how="left") df_hist = _fill_defaults(df_hist, history=True) df_load_training = pd.DataFrame() if training and "Hausverbrauch" in full_hist_raw.columns: load_segment = _longest_consistent_segment(full_hist_raw, "Hausverbrauch") if not load_segment.empty: idx_load = pd.date_range( start=load_segment.index.min().floor("5min"), end=load_segment.index.max().floor("5min"), freq="5min", ) df_load_training = pd.DataFrame(index=idx_load).join(load_segment, how="left") df_load_training = _fill_defaults(df_load_training, history=True) df_pv_training = pd.DataFrame() if training and "PV" in full_hist_raw.columns: pv_first = full_hist_raw["PV"].first_valid_index() if pv_first is not None: idx_pv = pd.date_range( start=pv_first.floor("5min"), end=now - datetime.timedelta(minutes=5), freq="5min", ) df_pv_training = pd.DataFrame(index=idx_pv).join(full_hist_raw, how="left") df_pv_training = _fill_defaults(df_pv_training, history=True) idx_fut = pd.date_range(start=now, periods=FORECAST_HORIZON_HOURS * 12, freq="5min") fut_raw = combined.reindex(combined.index.union(idx_fut)).sort_index() if not combined.empty else pd.DataFrame(index=idx_fut) df_fut = pd.DataFrame(index=idx_fut).join(fut_raw, how="left") df_fut = _fill_defaults(df_fut, history=False) month_hist = df_hist[(df_hist.index.year == now.year) & (df_hist.index.month == now.month)] current_peak_kw = 0.0 if not month_hist.empty and "Netzleistung" in month_hist.columns: measured_grid = pd.to_numeric(month_hist["Netzleistung"], errors="coerce").dropna() measured_peak = measured_grid.resample( "15min", origin="start_day", label="left", closed="left" ).mean() if not measured_peak.empty: current_peak_kw = max(0.0, float(measured_peak.max()) / 1000.0) if current_peak_kw <= 0.0 and not month_hist.empty and {"Hausverbrauch", "PV"}.issubset(month_hist.columns): fallback_residual = (month_hist["Hausverbrauch"] - month_hist["PV"]).resample( "15min", origin="start_day", label="left", closed="left" ).mean() if not fallback_residual.empty: current_peak_kw = max(0.0, float(fallback_residual.max()) / 1000.0) min_soc = float( config.get("batt_min_soc", config.get("batt_min_soc_percent", 0.0)) or 0.0 ) current_soc = max(0.0, min(100.0, min_soc)) current_soc_source = "safe_minimum" current_soc_age_minutes = None if "SOC" in full_hist_raw.columns: soc_values = pd.to_numeric(full_hist_raw["SOC"], errors="coerce").dropna() if not soc_values.empty: latest_soc_time = pd.Timestamp(soc_values.index[-1]) current_soc_age_minutes = max( 0.0, (pd.Timestamp(now) - latest_soc_time).total_seconds() / 60.0, ) try: max_soc_age_minutes = max( 5.0, float(config.get("batt_soc_max_age_minutes", 30.0) or 30.0), ) except (TypeError, ValueError): max_soc_age_minutes = 30.0 if current_soc_age_minutes <= max_soc_age_minutes: current_soc = float(soc_values.iloc[-1]) current_soc_source = "telemetry" current_soc = max(0.0, min(100.0, current_soc)) return { "config": config, "metadata": config, "now": now, "df_hist": df_hist, "df_recent_raw": recent_raw, "df_load_training": df_load_training, "df_pv_training": df_pv_training, "df_fut": df_fut, "current_soc_perc": current_soc, "current_soc_source": current_soc_source, "current_soc_age_minutes": current_soc_age_minutes, "current_month_peak_kw": current_peak_kw, } def run_training(): for config in get_configs(): if not any(active(config, n) for n in MODEL_MODULES): continue print(f"Training Anlage {config['anlagen_id']}...") data_obj = build_data_object(config, training=True) for n, module in MODEL_MODULES.items(): if active(config, n) and hasattr(module, "train"): try: target = QUALITY_TARGETS.get(n) if target: require_recent_telemetry(data_obj, [target]) print(f" var_{n}: {module.train(data_obj)}") except Exception: print(f" var_{n}: Training fehlgeschlagen") traceback.print_exc() def _forecast_point(aid, field, t, value): return ( Point("api_telemetry") .tag("anlagen_id", aid) .tag("data_type", "forecast") .field(field, float(value)) .time(t.to_pydatetime(), WritePrecision.S) ) def _snapshot_point(aid, field, run_hour, t, value): return ( Point("api_telemetry") .tag("anlagen_id", aid) .tag("data_type", "forecast_snapshot") .tag("run_hour", run_hour) .field(f"{field}_run_{run_hour}", float(value)) .time(t.to_pydatetime(), WritePrecision.S) ) def _quality_query(aid, forecast_field, target_field): stop = _time_literal(datetime.datetime.utcnow() - datetime.timedelta(minutes=10)) return f''' from(bucket: "{INFLUX_BUCKET}") |> range(start: -{QUALITY_LOOKBACK_DAYS}d, stop: {stop}) |> filter(fn: (r) => r["_measurement"] == "api_telemetry") |> filter(fn: (r) => r["anlagen_id"] == "{aid}") |> filter(fn: (r) => r["_field"] == "{target_field}" or r["_field"] == "{forecast_field}") |> aggregateWindow(every: 5m, fn: mean, createEmpty: false) |> group() |> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value") ''' def _forecast_quality(aid, variant, target_field): forecast_field = f"prog_var_{variant}" df = _query_df(_quality_query(aid, forecast_field, target_field)) if df.empty or forecast_field not in df.columns or target_field not in df.columns: return None pair = df[[target_field, forecast_field]].apply(pd.to_numeric, errors="coerce").dropna() if len(pair) < 3: return None actual = pair[target_field].to_numpy(dtype=float) pred = pair[forecast_field].to_numpy(dtype=float) ss_res = float(np.sum((actual - pred) ** 2)) ss_tot = float(np.sum((actual - np.mean(actual)) ** 2)) r2 = None if ss_tot <= 0 else 1.0 - (ss_res / ss_tot) mae = float(np.mean(np.abs(actual - pred))) rmse = float(np.sqrt(np.mean((actual - pred) ** 2))) return {"r2": r2, "mae": mae, "rmse": rmse, "samples": int(len(pair))} def _quality_point(aid, variant, target, metrics): p = ( Point("forecast_metrics") .tag("anlagen_id", aid) .tag("forecast", f"prog_var_{variant}") .tag("target", target) .field("samples", int(metrics["samples"])) .field("mae", float(metrics["mae"])) .field("rmse", float(metrics["rmse"])) .time(datetime.datetime.utcnow(), WritePrecision.S) ) if metrics["r2"] is not None: p.field("r2", float(metrics["r2"])) return p def write_quality_metrics(write_api, aid, active_variants): points = [] for variant, target in QUALITY_TARGETS.items(): if variant not in active_variants: continue try: metrics = _forecast_quality(aid, variant, target) if metrics: points.append(_quality_point(aid, variant, target, metrics)) r2_text = "nan" if metrics["r2"] is None else f"{metrics['r2']:.3f}" print(f"R2 Anlage {aid} prog_var_{variant}: {r2_text} / samples={metrics['samples']}") except Exception: print(f"R2 Anlage {aid} prog_var_{variant}: Berechnung fehlgeschlagen") traceback.print_exc() if points: try: write_api.write(bucket=INFLUX_BUCKET, org=INFLUX_ORG, record=points) except Exception: print(f"R2 Anlage {aid}: Schreiben der Qualitaetswerte fehlgeschlagen") traceback.print_exc() def run_forecast(only_anlagen_id=None, manual=False): configs = get_configs() if only_anlagen_id: configs = [c for c in configs if c.get("anlagen_id") == only_anlagen_id] client = InfluxDBClient( url=INFLUX_URL, token=INFLUX_TOKEN, org=INFLUX_ORG, timeout=INFLUX_TIMEOUT_MS, ) write_api = client.write_api(write_options=SYNCHRONOUS) completed = [] failures = [] try: for config in configs: aid = config["anlagen_id"] try: data_obj = build_data_object(config, training=False) targets = [] if any(active(config, n) for n in (1, 3, 10, 13, 21, 23)): targets.append("PV") if any(active(config, n) for n in (2, 3, 11, 13, 22, 23)): targets.append("Hausverbrauch") if float(config.get("batt_capacity_kwh", 0.0) or 0.0) > 0 and any(active(config, n) for n in (3, 13, 23)): targets.append("SOC") # Fail BEFORE model prediction, snapshot publication or V1/V4 plan writes. data_obj["telemetry_quality"] = require_recent_telemetry(data_obj, targets) soc_age = data_obj.get("current_soc_age_minutes") soc_age_text = "keine Messung" if soc_age is None else f"{soc_age:.1f} min" print( f"Forecast Anlage {aid}: Batterie-SOC {data_obj['current_soc_perc']:.1f}% " f"({data_obj['current_soc_source']}, Alter {soc_age_text})." ) forecasts, model_errors = collect_predictions( data_obj, config, {1: v1, 2: v2, 10: v10, 11: v11, 21: v21, 22: v22}, active, ) data_obj['forecast_model_status'] = model_errors for variant, detail in model_errors.items(): print(f"Forecast Anlage {aid} prog_var_{variant}: unavailable ({detail['errorType']}); keine Nullwerte eingesetzt.") p_1, p_2, p_10, p_11, p_21, p_22 = ( forecasts[n] for n in (1, 2, 10, 11, 21, 22) ) # ENELIX_V4_SHADOW_BRIDGE _v4_publish_forecasts(config, [(3,p_1,p_2),(13,p_10,p_11),(23,p_21,p_22)]) p_3 = v3.predict(data_obj, p_1, p_2) if active(config, 3) and p_1 and p_2 else {} p_13 = v13.predict(data_obj, p_10, p_11) if active(config, 13) and p_10 and p_11 else {} p_23 = v23.predict(data_obj, p_21, p_22) if active(config, 23) and p_21 and p_22 else {} forecast_sets = [ (1, p_1), (2, p_2), (3, p_3), (10, p_10), (11, p_11), (13, p_13), (21, p_21), (22, p_22), (23, p_23), ] run_hour = str(int(datetime.datetime.now(LOCAL_TZ).strftime("%H"))) points = [] for variant_id, pv_src, load_src, grid_src in [ (3, p_1, p_2, p_3), (13, p_10, p_11, p_13), (23, p_21, p_22, p_23), ]: if grid_src and pv_src and load_src: points.extend(battery_soc_points(data_obj, variant_id, pv_src, load_src, grid_src)) active_variants = set() battery_plans = data_obj.get("battery_plans", {}) for t in data_obj["df_fut"].index: for n, values in forecast_sets: if t in values: field = f"prog_var_{n}" value = float(values[t]) points.append(_forecast_point(aid, field, t, value)) points.append(_snapshot_point(aid, field, run_hour, t, value)) if n in battery_plans and t in battery_plans[n].get("battery", {}): points.append(_forecast_point( aid, f"prog_var_{n}_battery", t, float(battery_plans[n]["battery"][t]), )) active_variants.add(n) if points: write_api.write(bucket=INFLUX_BUCKET, org=INFLUX_ORG, record=points) print(f"Forecast Anlage {aid}: {len(points)} Punkte geschrieben inkl. run_{run_hour}.") else: print(f"Forecast Anlage {aid}: keine aktiven Prognosen oder keine Daten.") write_quality_metrics(write_api, aid, active_variants) completed.append(aid) except Exception as exc: failures.append({"anlagen_id": aid, "error": f"{type(exc).__name__}: {exc}"}) print(f"Forecast Anlage {aid}: Lauf fehlgeschlagen.") traceback.print_exc() finally: client.close() if failures: raise RuntimeError(f"Forecast-Fehler: {failures}") return {"completed": completed, "manual": bool(manual)} def _write_json_atomic(path, payload): directory = os.path.dirname(path) or "." os.makedirs(directory, exist_ok=True) tmp_path = f"{path}.tmp.{os.getpid()}" try: with open(tmp_path, "w", encoding="utf-8") as handle: json.dump(payload, handle) os.replace(tmp_path, path) finally: if os.path.exists(tmp_path): os.remove(tmp_path) def _read_json(path): try: with open(path, "r", encoding="utf-8") as handle: return json.load(handle) except Exception: return {} def _last_success_age_seconds(): stamp = _read_json(LAST_SUCCESS_PATH).get("timestamp") if not stamp: return None try: then = datetime.datetime.fromisoformat(str(stamp).replace("Z", "+00:00")) now = datetime.datetime.now(datetime.timezone.utc) return max(0.0, (now - then).total_seconds()) except Exception: return None def _child_environment(training=False): env = os.environ.copy() env["PYTHONUNBUFFERED"] = "1" if training: max_threads = str(max(1, int(env.get("FORECAST_TRAINING_MAX_THREADS", "1")))) for name in ( "OMP_NUM_THREADS", "OPENBLAS_NUM_THREADS", "MKL_NUM_THREADS", "NUMEXPR_NUM_THREADS", "LOKY_MAX_CPU_COUNT", ): env[name] = max_threads return env def _run_child(mode, timeout_seconds, only_anlagen_id=None): command = [sys.executable, "-u", os.path.abspath(__file__), mode] if only_anlagen_id: command.append(str(only_anlagen_id)) label = "Training" if mode == "--train-once" else "Forecast" print(f"{label}-Kindprozess startet (Timeout {timeout_seconds}s).") process = subprocess.Popen(command, env=_child_environment(training=mode == "--train-once")) try: return_code = process.wait(timeout=timeout_seconds) except subprocess.TimeoutExpired: print(f"{label}-Kindprozess hat das Zeitlimit erreicht und wird beendet.") process.terminate() try: process.wait(timeout=15) except subprocess.TimeoutExpired: process.kill() process.wait(timeout=15) return 124 if return_code != 0: print(f"{label}-Kindprozess beendet mit Status {return_code}.") return return_code def run_forecast_isolated(only_anlagen_id=None, source="scheduler"): if not _FORECAST_PROCESS_LOCK.acquire(blocking=False): print(f"Forecast-Aufruf ({source}) uebersprungen: bereits ein Lauf aktiv.") return {"status": "busy", "source": source} try: return_code = _run_child("--forecast-once", FORECAST_TIMEOUT_SECONDS, only_anlagen_id) return { "status": "ok" if return_code == 0 else "error", "source": source, "return_code": return_code, } finally: _FORECAST_PROCESS_LOCK.release() def _training_worker(day_text): global _TRAINING_THREAD started = datetime.datetime.now(datetime.timezone.utc).isoformat() _write_json_atomic(LAST_TRAINING_PATH, { "day": day_text, "status": "running", "started_at": started, }) return_code = 1 error = None try: return_code = _run_child("--train-once", TRAINING_TIMEOUT_SECONDS) except Exception as exc: error = f"{type(exc).__name__}: {exc}" traceback.print_exc() finally: status = "ok" if return_code == 0 else ("timeout" if return_code == 124 else "error") payload = { "day": day_text, "status": status, "started_at": started, "finished_at": datetime.datetime.now(datetime.timezone.utc).isoformat(), "return_code": return_code, } if error: payload["error"] = error _write_json_atomic(LAST_TRAINING_PATH, payload) print(f"Nachttraining beendet: {status}.") _TRAINING_THREAD = None def start_training_async(day): global _TRAINING_THREAD if _TRAINING_THREAD is not None and _TRAINING_THREAD.is_alive(): print("Nachttraining laeuft bereits.") return False day_text = day.isoformat() _TRAINING_THREAD = threading.Thread(target=_training_worker, args=(day_text,), daemon=True) _TRAINING_THREAD.start() return True def _watchdog_loop(): while True: time.sleep(max(60, WATCHDOG_INTERVAL_SECONDS)) try: age = _last_success_age_seconds() if age is None or age > FORECAST_STALE_SECONDS: age_text = "unbekannt" if age is None else f"{age / 60.0:.1f} Minuten" print(f"Forecast-Watchdog: letzter erfolgreicher Lauf {age_text}; neuer Lauf wird gestartet.") run_forecast_isolated(source="watchdog") except Exception: print("Forecast-Watchdog: Pruefung fehlgeschlagen.") traceback.print_exc() # BEGIN EMS RESIMULATE HTTP SERVER _RESIMULATE_SERVER_STARTED = False def _start_resimulate_server(): global _RESIMULATE_SERVER_STARTED if _RESIMULATE_SERVER_STARTED: return _RESIMULATE_SERVER_STARTED = True import json as _json import os as _os import threading as _threading import traceback as _traceback import urllib.parse as _urlparse from http.server import BaseHTTPRequestHandler as _BaseHTTPRequestHandler, ThreadingHTTPServer as _ThreadingHTTPServer class _Handler(_BaseHTTPRequestHandler): def log_message(self, fmt, *args): return def _send(self, code, body): raw = _json.dumps(body).encode("utf-8") self.send_response(code) self.send_header("Content-Type", "application/json") self.send_header("Content-Length", str(len(raw))) self.end_headers() self.wfile.write(raw) def do_GET(self): self._handle() def do_POST(self): self._handle() def _handle(self): try: parsed = _urlparse.urlparse(self.path) if parsed.path == "/health": age = _last_success_age_seconds() healthy = age is not None and age <= FORECAST_STALE_SECONDS self._send(200 if healthy else 503, { "status": "ok" if healthy else "stale", "last_success_age_seconds": age, "training": _read_json(LAST_TRAINING_PATH), }) return if parsed.path not in ("/run_now", "/resimulate"): self._send(404, {"detail": "not found"}) return params = _urlparse.parse_qs(parsed.query) anlagen_id = (params.get("anlagen_id") or [None])[0] result = run_forecast_isolated(only_anlagen_id=anlagen_id, source="http") result["anlagen_id"] = anlagen_id code = 200 if result["status"] == "ok" else (409 if result["status"] == "busy" else 500) self._send(code, result) except Exception as exc: _traceback.print_exc() self._send(500, {"detail": str(exc)}) port = int(_os.getenv("FORECAST_ENGINE_HTTP_PORT", "9000")) server = _ThreadingHTTPServer(("0.0.0.0", port), _Handler) thread = _threading.Thread(target=server.serve_forever, daemon=True) thread.start() print(f"Forecast Resimulate HTTP Server startet auf Port {port}.") # END EMS RESIMULATE HTTP SERVER def main(): _start_resimulate_server() global last_trained_day print("Forecast Engine startet.") threading.Thread(target=_watchdog_loop, daemon=True).start() run_forecast_isolated(source="startup") while True: try: now_local = datetime.datetime.now(LOCAL_TZ) if now_local.hour == TRAINING_HOUR and last_trained_day != now_local.date(): if start_training_async(now_local.date()): last_trained_day = now_local.date() next_run = (now_local + datetime.timedelta(hours=1)).replace(minute=0, second=0, microsecond=0) time.sleep(max(60.0, (next_run - now_local).total_seconds())) run_forecast_isolated(source="scheduler") except Exception: traceback.print_exc() time.sleep(60) if __name__ == "__main__": if len(sys.argv) >= 2 and sys.argv[1] == "--forecast-once": selected_anlage = sys.argv[2] if len(sys.argv) >= 3 else None run_forecast(only_anlagen_id=selected_anlage, manual=True) _write_json_atomic(LAST_SUCCESS_PATH, { "timestamp": datetime.datetime.now(datetime.timezone.utc).isoformat(), "anlagen_id": selected_anlage, "pid": os.getpid(), }) elif len(sys.argv) >= 2 and sys.argv[1] == "--train-once": run_training() else: main()