107 lines
4.2 KiB
Python
107 lines
4.2 KiB
Python
"""Measured telemetry integrity and repeat-profile helpers.
|
|
No IO, no fabricated measurements and no fixed household-load fallback.
|
|
Forecast freshness is checked against raw telemetry, never against filled features.
|
|
"""
|
|
from __future__ import annotations
|
|
import datetime
|
|
import math
|
|
import numpy as np
|
|
import pandas as pd
|
|
|
|
MEASURED_COLUMNS = ('PV', 'Hausverbrauch', 'Netzleistung', 'SOC')
|
|
|
|
class TelemetryUnavailable(ValueError):
|
|
"""No publishable forecast can be derived from the supplied observations."""
|
|
|
|
|
|
def sanitize_measured_frame(frame):
|
|
out = frame.copy()
|
|
for column in MEASURED_COLUMNS:
|
|
if column not in out:
|
|
continue
|
|
raw = out[column]
|
|
numeric = pd.to_numeric(raw, errors='coerce').astype(float)
|
|
bad = ~np.isfinite(numeric) | raw.map(lambda x: isinstance(x, (bool, np.bool_)))
|
|
if column in ('PV', 'Hausverbrauch', 'SOC'):
|
|
bad |= numeric < 0
|
|
if column == 'SOC':
|
|
bad |= numeric > 100
|
|
out[column] = numeric.mask(bad)
|
|
return out
|
|
|
|
|
|
def _naive_utc(value):
|
|
stamp = pd.Timestamp(value)
|
|
if stamp.tzinfo is not None:
|
|
stamp = stamp.tz_convert('UTC').tz_localize(None)
|
|
return stamp
|
|
|
|
|
|
def require_recent_telemetry(data_obj, fields=('PV','Hausverbrauch'), max_age_minutes=30.0):
|
|
"""Refuse missing/stale inputs without changing real zeros or raw samples."""
|
|
if isinstance(max_age_minutes, bool) or not math.isfinite(max_age_minutes) or max_age_minutes <= 0:
|
|
raise ValueError('Positive telemetry age limit required')
|
|
now = _naive_utc(data_obj['now'])
|
|
raw = data_obj.get('df_recent_raw')
|
|
raw = sanitize_measured_frame(raw) if raw is not None else pd.DataFrame()
|
|
if not raw.empty:
|
|
raw.index = pd.DatetimeIndex([_naive_utc(t) for t in raw.index])
|
|
raw = raw.loc[raw.index < now].sort_index()
|
|
report = {}
|
|
errors = []
|
|
for field in fields:
|
|
if field not in MEASURED_COLUMNS:
|
|
raise ValueError('Unknown telemetry target')
|
|
values = raw[field].dropna() if field in raw else pd.Series(dtype=float)
|
|
if values.empty:
|
|
errors.append(field + ': keine gemessenen Werte')
|
|
continue
|
|
stamp = values.index[-1]
|
|
age = (now-stamp).total_seconds()/60.0
|
|
report[field] = {'lastObservedInterval': stamp.isoformat()+'Z', 'ageMinutes': age,
|
|
'observedIntervals': int(len(values)), 'lastValue': float(values.iloc[-1])}
|
|
if age > max_age_minutes:
|
|
errors.append(field + ': Messdaten veraltet (' + format(age,'.1f') + ' min)')
|
|
if errors:
|
|
raise TelemetryUnavailable('; '.join(errors) + '. Keine neuen Prognosen/Fahrplaene veroeffentlicht.')
|
|
return report
|
|
|
|
|
|
def _finite_nonnegative(value):
|
|
if isinstance(value, (bool, np.bool_)):
|
|
return None
|
|
try:
|
|
value = float(value)
|
|
except (ValueError, TypeError):
|
|
return None
|
|
return value if math.isfinite(value) and value >= 0.0 else None
|
|
|
|
|
|
def profile_source_value(history, at, column, predictions=None):
|
|
"""Repeat yesterday; prefer same weekday if yesterday is missing, then older days.
|
|
Explicitly generated first-day values may be repeated on the second forecast day.
|
|
Missing historical values are not zeros. No future measurement is read.
|
|
"""
|
|
if column not in ('PV','Hausverbrauch'):
|
|
raise ValueError('Unsupported repeat-profile target')
|
|
predictions = {} if predictions is None else predictions
|
|
at = pd.Timestamp(at)
|
|
for day in (1, 7, 2, 3, 4, 5, 6, 8, 9, 10, 11, 12, 13, 14):
|
|
source = at - datetime.timedelta(days=day)
|
|
if source in predictions:
|
|
value = _finite_nonnegative(predictions[source])
|
|
if value is not None:
|
|
return value
|
|
if history is not None and column in history and source in history.index:
|
|
value = _finite_nonnegative(history.at[source,column])
|
|
if value is not None:
|
|
return value
|
|
raise TelemetryUnavailable(column + ': kein gemessener Tagesprofilwert fuer ' + str(at))
|
|
|
|
|
|
def repeat_daily_profile(history, future_index, column):
|
|
result = {}
|
|
for at in future_index:
|
|
result[at] = profile_source_value(history, at, column, result)
|
|
return result
|