"""Read-only probe executed INSIDE the existing forecast-engine container. No get_configs() (it may migrate SQL), no training/prediction, no forecast writing, no manual /run_now request, and no user credentials in output. Reads one plant's stored raw 5-minute forecast values and the input frames used by the engine. """ import contextlib from datetime import datetime, timezone import hashlib import io import json import math import os from pathlib import Path import sqlite3 from urllib.parse import quote from uuid import UUID class Discard(io.TextIOBase): def write(self,text):return len(text) def clean_time(value): if hasattr(value,'to_pydatetime'):value=value.to_pydatetime() if value.tzinfo is None:value=value.replace(tzinfo=timezone.utc) return value.astimezone(timezone.utc).isoformat() def run(): aid=str(UUID(os.environ['ENELIX_ACCEPTANCE_PLANT'])) os.environ['FORECAST_INFLUX_TIMEOUT_MS']='20000' # These assignments affect only this diagnostic process, not the service. with contextlib.redirect_stdout(Discard()),contextlib.redirect_stderr(Discard()): import main import numpy as np import pandas as pd path=Path(main.SQLITE_DB_PATH) if not path.is_file():raise RuntimeError('Existing configuration database missing') con=sqlite3.connect('file:'+quote(str(path))+'?mode=ro',uri=True) con.row_factory=sqlite3.Row try:row=con.execute('SELECT * FROM anlagen_meta WHERE anlagen_id=?',(aid,)).fetchone() finally:con.close() if row is None:raise RuntimeError('Installation not found') cfg=dict(row) cfg['daecher']=json.loads(cfg.get('daecher') or '[]') for key,default in [('ac_leistung',10.),('batt_capacity_kwh',0.),('batt_power_kw',0.),('tarif_bezug_fest',.3),('tarif_einspeisung_fest',.1),('tarif_peak_fest',5.)]: value=cfg.get(key);cfg[key]=float(default if value is None or value=='' else value) data=main.build_data_object(cfg,training=False) frames={} for key in ('df_hist','df_recent_raw','df_load_training','df_fut'): frame=data.get(key) if frame is None:frames[key]={'available':False};continue item={'rows':len(frame),'from':clean_time(frame.index.min()) if len(frame) else None,'until':clean_time(frame.index.max()) if len(frame) else None,'columns':{}} for col in ('Hausverbrauch','PV','Netzleistung','SOC'): if col not in frame.columns:continue values=pd.to_numeric(frame[col],errors='coerce');valid=values[np.isfinite(values)] item['columns'][col]={'finite':len(valid),'zeros':int((valid==0).sum()),'median':float(valid.median()) if len(valid) else None,'maximum':float(valid.max()) if len(valid) else None,'lastFiniteAt':clean_time(valid.index[-1]) if len(valid) else None,'recentValues':[{'time':clean_time(t),'value':float(v)} for t,v in valid.tail(12).items()]} frames[key]=item start=datetime.now(timezone.utc).replace(second=0,microsecond=0) from datetime import timedelta end=start+timedelta(hours=48) fields=('prog_var_1','prog_var_2','prog_var_10','prog_var_11','prog_var_21','prog_var_22') field_filter=' or '.join('r["_field"] == '+json.dumps(f) for f in fields) query='''from(bucket: %s) |> range(start: %s, stop: %s) |> filter(fn: (r) => r["_measurement"] == "api_telemetry") |> filter(fn: (r) => r["anlagen_id"] == %s) |> filter(fn: (r) => r["data_type"] == "forecast") |> filter(fn: (r) => %s) |> keep(columns: ["_time", "_field", "_value"]) ''' % (json.dumps(main.INFLUX_BUCKET),start.isoformat(),end.isoformat(),json.dumps(aid),field_filter) client=main.InfluxDBClient(url=main.INFLUX_URL,token=main.INFLUX_TOKEN,org=main.INFLUX_ORG,timeout=20000) series={field:[] for field in fields} try: tables=client.query_api().query(org=main.INFLUX_ORG,query=query) for table in tables: for record in table.records: value=record.get_value();field=record.get_field() if field in series and isinstance(value,(int,float)) and math.isfinite(value): series[field].append({'time':clean_time(record.get_time()),'value':float(value)}) finally:client.close() for field in series:series[field].sort(key=lambda p:p['time']) versions={} for relative in ('main.py','methods/var_1.py','methods/var_2.py','methods/var_10.py','methods/var_11.py','methods/var_21.py','methods/var_22.py'): source=Path('/app')/relative if source.is_file():versions[relative]=hashlib.sha256(source.read_bytes()).hexdigest() summary={} for field,points in series.items(): values=[p['value'] for p in points] summary[field]={'points':len(values),'allZero':bool(values) and max(abs(v) for v in values)==0,'maximumW':max(values) if values else None,'from':points[0]['time'] if points else None,'until':points[-1]['time'] if points else None} native=True for points in series.values(): stamps=[datetime.fromisoformat(p['time']).timestamp() for p in points] if not stamps or any(t%300 for t in stamps) or any(b-a!=300 for a,b in zip(stamps,stamps[1:])):native=False return {'status':'read_only_acquired','installationId':aid,'observedAt':datetime.now(timezone.utc).isoformat(),'queryResolution':'raw_no_chart_resampling','nativeFiveMinuteForecast':native,'generationTimeVerified':False,'inputFrames':frames,'forecastSummary':summary,'forecastSeries':series,'runningSourceHashes':versions,'modelWrite':False,'liveControlChanged':False} try: result=run() except Exception as error: # Never echo exception details containing SQL payloads, credentials or URLs. result={'status':'read_only_probe_failed','errorType':type(error).__name__,'liveControlChanged':False} print(json.dumps(result,allow_nan=False))