Tests / test (push) Successful in 1m1s
Approved by Daniel Haefliger for develop and beta. Author dh_Agent, authenticated account dh. Preserve published battery, charging and overall Energy Pie changes. No deployment or plant control authorization.
487 lines
35 KiB
Python
487 lines
35 KiB
Python
"""Isolated V4 service: immutable inputs, coalesced replan queue, SHADOW publication.
|
|
No external actuator endpoint, no reads of users.db, no legacy schedule changes.
|
|
"""
|
|
from contextlib import asynccontextmanager
|
|
from dataclasses import asdict,replace
|
|
from datetime import datetime,timedelta,timezone
|
|
from pathlib import Path
|
|
from threading import Event,Thread
|
|
from uuid import UUID
|
|
import json
|
|
import os
|
|
import secrets
|
|
import logging
|
|
from fastapi import FastAPI,Header,HTTPException
|
|
from .domain import Battery,Limits,Price,QuarterPast,Step,month_key,quarter_start,utc,number,priced_prefix
|
|
from .store import PlannerStore,canonical
|
|
from .selection import choose_family
|
|
from .optimizer import optimize
|
|
from . import meter_runtime, controlled_trial, economic_replay, retention, measurement_pipeline, workflow
|
|
from .forecast_quality import assess_family
|
|
from .receiver_contract import provenance
|
|
from .peak_policy import basis_record, RestMonthOutlook, empirical_rest_month
|
|
|
|
KINDS={'forecast','operation','tariffs','prices','planning_basis','peak_outlook'}
|
|
|
|
def latest(store,plant,kind):
|
|
row=store.con.execute('''SELECT value FROM planner_inputs i JOIN planner_input_current c
|
|
ON i.plant=c.plant AND i.kind=c.kind AND i.event_id=c.event_id WHERE i.plant=? AND i.kind=?''',(plant,kind)).fetchone()
|
|
return json.loads(row[0]) if row else None
|
|
|
|
def batteries(value):
|
|
result=[]
|
|
for b in value:
|
|
result.append(Battery(asset_id=b['id'],capacity_kwh=b['capacityKwh'],soc_percent=b['socPercent'],
|
|
min_soc_percent=b['minSocPercent'],max_soc_percent=b['maxSocPercent'],
|
|
max_charge_w=b['maxChargeW'],max_discharge_w=b['maxDischargeW'],measured_at=utc(b['measuredAt']),
|
|
grid_charging=b['gridCharging'],throughput_chf_kwh=b.get('throughputChfKwh',0.),
|
|
terminal_soc_min_percent=b.get('terminalMinSocPercent'),terminal_value_chf_kwh=b.get('terminalValueChfKwh',0.),
|
|
recovery_allowed=b.get('recoveryAllowed',False),physical_min_soc_percent=b.get('physicalMinSocPercent',0.),
|
|
discharge_blocked=b.get('dischargeBlocked',False),rearm_soc_percent=b.get('rearmSocPercent')))
|
|
if type(b['gridCharging']) is not bool:raise ValueError('Explicit boolean grid-charging permission required')
|
|
if len(result)>20 or len({b.asset_id for b in result})!=len(result):raise ValueError('Duplicate/too many batteries')
|
|
return result
|
|
|
|
def validate(kind,value,now,registry):
|
|
if kind not in KINDS or type(value.get('version')) is not int or value['version']!=1:raise ValueError('Unsupported event version/kind')
|
|
identifier=value['eventId']
|
|
if not isinstance(identifier,str) or not 1<=len(identifier)<=160:raise ValueError('Event ID required')
|
|
observed=utc(value['observedAt'])
|
|
if observed>now+timedelta(seconds=30):raise ValueError('Observation from future')
|
|
fields={'version','eventId','observedAt'}
|
|
if kind=='forecast':
|
|
fields|={'families','modelVersions'}
|
|
if not value['families']:raise ValueError('No forecast families')
|
|
for key,family in value['families'].items():
|
|
registry.get(key)
|
|
if family['loadBasis'] not in ('base_load','house_total'):raise ValueError('Explicit metering basis required')
|
|
evidence=family.get('accountingEvidenceId')
|
|
if evidence is not None and (family['loadBasis']!='base_load' or not isinstance(evidence,str) or not 8<=len(evidence)<=160):raise ValueError('Explicit base-load accounting evidence required')
|
|
if family.get('trainedUntil') and utc(family['trainedUntil'])>observed:raise ValueError('Training leakage')
|
|
if not 1<=len(family['points'])<=576:raise ValueError('Need 1..576 forecast intervals')
|
|
previous=None
|
|
for p in family['points']:
|
|
t=utc(p['time'])
|
|
if t.second or t.microsecond or t.minute%5:raise ValueError('Forecast interval alignment')
|
|
if previous and t-previous!=timedelta(minutes=5):raise ValueError('Forecast gap/overlap')
|
|
previous=t;number(p['pvW'],'PV',0,1e9);number(p['loadW'],'load',0,1e9);number(p.get('externalW',0.),'external',-1e9,1e9)
|
|
elif kind=='operation':
|
|
fields|={'gridW','meteringBoundary','batteries','limits','measuredPeaks','quarterPast','planningPeaks','quarterEstimate','meterObservation'}
|
|
if value['meteringBoundary']!='common_pcc':raise ValueError('Common metering boundary required')
|
|
number(value['gridW'],'grid W',-1e9,1e9)
|
|
for b in batteries(value['batteries']):b.validate(observed)
|
|
limits=value['limits']
|
|
if set(limits)!={'importW','exportW','managerMonthLimitsW'}:raise ValueError('Explicit limits required; null unlimited, zero zero')
|
|
for k in ('importW','exportW'):
|
|
if limits[k] is not None:number(limits[k],k,0,1e9)
|
|
for k,v in limits['managerMonthLimitsW'].items():
|
|
if str(int(k))!=k or not 1<=int(k)<=12:raise ValueError('Invalid manager month')
|
|
number(v,'manager limit',0,1e9)
|
|
for m,p in value.get('measuredPeaks',{}).items():
|
|
datetime.strptime(m,'%Y-%m');number(p['kw'],'peak kW',0)
|
|
if m>month_key(observed) or p['source'] not in ('meter_month_register','verified_month_history','verified_new_month'):raise ValueError('Measured peak source invalid; cap is not paid peak')
|
|
past=value.get('quarterPast')
|
|
if past:
|
|
q=quarter_start(observed)
|
|
if utc(past['start'])!=q or type(past['measuredSeconds']) is not int or past['measuredSeconds']!=int((observed-q).total_seconds()):raise ValueError('Quarter measurement timestamp mismatch')
|
|
number(past['importKwh'],'quarter energy',0)
|
|
for m,p in value.get('planningPeaks',{}).items():
|
|
record=basis_record(p,m,observed,allow_estimates=True)
|
|
if record['quality']!='estimated':raise ValueError('planningPeaks contains estimates only')
|
|
estimate=value.get('quarterEstimate')
|
|
if estimate:
|
|
if set(estimate)-{'start','measuredSeconds','importKwh','quality','source','coverage','notes'}:raise ValueError('Unknown quarter estimate field')
|
|
if estimate.get('quality')!='estimated' or estimate.get('source') not in ('power_history_estimate','sampled_power_estimate','counter_interval_estimate'):raise ValueError('Explicit quarter estimate provenance required')
|
|
q=quarter_start(observed)
|
|
if utc(estimate['start'])!=q or type(estimate['measuredSeconds']) is not int or estimate['measuredSeconds']!=int((observed-q).total_seconds()):raise ValueError('Estimated quarter timing mismatch')
|
|
number(estimate['importKwh'],'estimated quarter energy',0)
|
|
if number(estimate.get('coverage'),'quarter estimate coverage',0,1)<1.:raise ValueError('Missing current-quarter coverage')
|
|
if value.get('meterObservation') is not None:meter_runtime.validate_observation(value['meterObservation'],observed)
|
|
elif kind=='planning_basis':
|
|
fields|={'peaks'}
|
|
if not isinstance(value.get('peaks'),dict) or not value['peaks']:raise ValueError('Explicit peak estimates required')
|
|
for m,p in value['peaks'].items():
|
|
if basis_record(p,m,observed,allow_estimates=True)['quality']!='estimated':raise ValueError('Planning basis is not a metering import')
|
|
elif kind=='peak_outlook':
|
|
fields|={'outlooks'}
|
|
for m,v in value['outlooks'].items():
|
|
outlook=RestMonthOutlook.from_dict(v)
|
|
if m!=outlook.month:raise ValueError('Outlook month mismatch')
|
|
outlook.validate(observed,observed)
|
|
elif kind=='tariffs':
|
|
fields|={'import','export','peakChfKwMonth'}
|
|
for side in ('import','export'):
|
|
p=value[side]
|
|
if p['mode'] not in ('static','dynamic') or not p['tariffId']:raise ValueError('Explicit price mode/id required')
|
|
if p['mode']=='static':number(p['staticChfKwh'],'static price')
|
|
peaks=value['peakChfKwMonth']
|
|
if isinstance(peaks,dict):
|
|
for m,v in peaks.items():datetime.strptime(m,'%Y-%m');number(v,'peak tariff',0)
|
|
else:number(peaks,'peak tariff',0)
|
|
else:
|
|
fields|={'periods'}
|
|
if len(value['periods'])>3000:raise ValueError('Too many price intervals')
|
|
for p in value['periods']:
|
|
if p['unit'] not in ('CHF/kWh','CHF_kWh','Rp/kWh','CHF/MWh') or p['side'] not in ('import','export'):raise ValueError('Explicit price unit/direction required')
|
|
if p['sourceKind'] not in ('published_interval','estimate','carried_forward'):raise ValueError('Explicit price provenance required')
|
|
number(p['value'],'price')
|
|
if utc(p['end'])<=utc(p['start']) or utc(p['observedAt'])>observed:raise ValueError('Invalid price interval or observation')
|
|
if p.get('publishedAt') and utc(p['publishedAt'])>utc(p['observedAt']):raise ValueError('Price not published when observed')
|
|
if set(value)-fields:raise ValueError('Unknown fields: extra device data/credentials must not be submitted')
|
|
canonical(value)
|
|
|
|
def ingest(store,plant,kind,value,now):
|
|
validate(kind,value,now,store.registry);data=canonical(value);con=store.con;con.execute('BEGIN IMMEDIATE')
|
|
try:
|
|
old=con.execute('SELECT value FROM planner_inputs WHERE plant=? AND kind=? AND event_id=?',(plant,kind,value['eventId'])).fetchone()
|
|
if old:
|
|
if old[0]!=data:raise ValueError('Immutable event conflict')
|
|
con.commit();return {'status':'duplicate','queued':False}
|
|
current=latest(store,plant,kind)
|
|
if current and utc(current['observedAt'])==utc(value['observedAt']):
|
|
left,right=dict(current),dict(value);left.pop('eventId');right.pop('eventId')
|
|
if canonical(left)!=canonical(right):raise ValueError('Conflicting simultaneous observations')
|
|
con.execute('INSERT INTO planner_inputs VALUES(?,?,?,?,?)',(plant,kind,value['eventId'],utc(value['observedAt']).isoformat(timespec='microseconds'),data))
|
|
newer=not current or utc(current['observedAt'])<=utc(value['observedAt'])
|
|
if newer:
|
|
if kind=='operation':
|
|
known=store.peaks(plant)
|
|
for m,p in value.get('measuredPeaks',{}).items():
|
|
if m in known and p['kw']<known[m]-1e-9:raise ValueError('Measured peak decreased')
|
|
con.execute('''INSERT INTO planner_month_peaks VALUES(?,?,?,?,?) ON CONFLICT(plant,month) DO UPDATE SET peak_kw=excluded.peak_kw,source=excluded.source,updated_at=excluded.updated_at''',(plant,m,p['kw'],p['source'],utc(now).isoformat()))
|
|
if kind=='operation':
|
|
for m,p in value.get('planningPeaks',{}).items():meter_runtime.save_assumption(con,plant,m,p,now)
|
|
if value.get('meterObservation'):meter_runtime.observe(con,plant,value['meterObservation'],utc(value['observedAt']))
|
|
elif kind=='planning_basis':
|
|
for m,p in value['peaks'].items():meter_runtime.save_assumption(con,plant,m,p,now)
|
|
con.execute('INSERT INTO planner_input_current VALUES(?,?,?) ON CONFLICT(plant,kind) DO UPDATE SET event_id=excluded.event_id',(plant,kind,value['eventId']))
|
|
store._request(plant,store.settings(plant)['revision'],kind+'_changed',now)
|
|
con.commit();return {'status':'stored' if newer else 'archived_older','queued':newer}
|
|
except Exception:con.rollback();raise
|
|
|
|
def price_at(store,plant,side,config,start,end,now):
|
|
if config['mode']=='static':return Price(config['staticChfKwh'])
|
|
rows=store.con.execute('''SELECT value FROM planner_inputs WHERE plant=? AND kind='prices' AND observed_at<=? ORDER BY observed_at DESC''',(plant,utc(now).isoformat(timespec='microseconds')))
|
|
candidates=[]
|
|
for row in rows:
|
|
for p in json.loads(row[0])['periods']:
|
|
if (p['tariffId']==config['tariffId'] and p['side']==side and p['sourceKind']=='published_interval'
|
|
and utc(p['start'])<=start and utc(p['end'])>=end and utc(p['observedAt'])<=now):
|
|
factor={'CHF/kWh':1.,'CHF_kWh':1.,'Rp/kWh':.01,'CHF/MWh':.001}[p['unit']]
|
|
candidates.append((utc(p['observedAt']),p['value']*factor))
|
|
if not candidates:return None
|
|
last=max(t for t,v in candidates);values={v for t,v in candidates if t==last}
|
|
if len(values)!=1:raise ValueError('Conflicting published price intervals')
|
|
return Price(values.pop(),last,'dynamic')
|
|
|
|
class AwaitingInput(ValueError):pass
|
|
|
|
def assemble(store,plant,family,now,*,include_forecast_points=False):
|
|
values={k:latest(store,plant,k) for k in ('operation','forecast','tariffs')}
|
|
for k,v in values.items():
|
|
if not v:raise AwaitingInput('Missing '+k+' input')
|
|
op,forecast,tariffs=(values[k] for k in ('operation','forecast','tariffs'))
|
|
decision=utc(op['observedAt']).replace(microsecond=0)
|
|
settings=store.settings(plant);allow_estimates=settings['measurementPolicy']=='allow_estimates'
|
|
input_quality={'quarter':'verified','warnings':[]}
|
|
if not 0<=(now-decision).total_seconds()<=120:raise AwaitingInput('Fresh manager observation required (120s maximum)')
|
|
if (now-utc(forecast['observedAt'])).total_seconds()>5400:raise AwaitingInput('Forecast older than 90 minutes')
|
|
if utc(forecast['observedAt'])>decision or utc(tariffs['observedAt'])>decision:raise AwaitingInput('New data awaiting fresh manager observation')
|
|
if settings['forecastSource']=='corrected_profile':
|
|
try:forecast=measurement_pipeline.apply_load_forecast(store.con,plant,settings['measurementDataset'],forecast,int(decision.timestamp()))
|
|
except ValueError as exc:raise AwaitingInput(str(exc)) from exc
|
|
if family not in forecast['families']:raise AwaitingInput('Chosen family is unavailable; no silent switch')
|
|
source=forecast['families'][family];steps=[]
|
|
future_source={**source,'points':[p for p in source['points'] if utc(p['time'])+timedelta(minutes=5)>decision]}
|
|
assessment=assess_family(future_source)
|
|
if not assessment['valid']:raise AwaitingInput(assessment['reason'])
|
|
input_quality['forecastAssessment']=assessment
|
|
if source.get('dataPipeline'):
|
|
input_quality['dataPipeline']=source['dataPipeline']
|
|
input_quality['warnings'].append('Corrected physical-load profile uses configured measurement mapping')
|
|
if source['dataPipeline'].get('externalPolicy')=='last_sdl_request_persistence_estimate':
|
|
input_quality['warnings'].append('External SDL is an explicitly labelled last-request persistence scenario, not a published future SDL schedule')
|
|
for p in source['points']:
|
|
start=utc(p['time']);end=start+timedelta(minutes=5)
|
|
if end<=decision:continue
|
|
start=max(start,decision)
|
|
external=p.get('externalW',0.) if source['loadBasis']=='base_load' else 0.
|
|
steps.append(Step(start,p['loadW'],p['pvW'],price_at(store,plant,'import',tariffs['import'],start,end,decision),price_at(store,plant,'export',tariffs['export'],start,end,decision),external,int((end-start).total_seconds())))
|
|
if not steps or steps[0].start!=decision:raise AwaitingInput('Forecast has no current interval')
|
|
if include_forecast_points:
|
|
# Preserve the exact selected input horizon separately from price-limited dispatch.
|
|
input_quality['forecastPoints']=[{'time':s.start.isoformat(),'validUntil':s.end.isoformat(),
|
|
'pvW':float(s.pv_w),'loadW':float(s.base_load_w),
|
|
'externalW':float(s.external_w) if source['loadBasis']=='base_load' else None} for s in steps]
|
|
full_end=steps[-1].end;steps=priced_prefix(steps,decision)
|
|
if not steps:raise AwaitingInput('No complete published-price billing quarter')
|
|
q=quarter_start(decision);elapsed=int((decision-q).total_seconds());past={}
|
|
if elapsed:
|
|
p=op.get('quarterPast')
|
|
if not p and allow_estimates:
|
|
p=op.get('quarterEstimate')
|
|
if not p and op.get('meterObservation'):p=meter_runtime.current_quarter(store.con,plant,op['meterObservation'],decision)
|
|
if p:
|
|
input_quality['quarter']='estimated'
|
|
input_quality['warnings'].append('Current quarter uses an explicitly estimated energy value, not a billing measurement')
|
|
if not p or utc(p['start'])!=q or p['measuredSeconds']!=elapsed:raise AwaitingInput('Current-quarter energy missing (measured or explicitly permitted estimate)')
|
|
past[q]=QuarterPast(p['importKwh'],elapsed)
|
|
peaks=store.peaks(plant);months={month_key(s.start) for s in steps};contexts={}
|
|
estimates=meter_runtime.assumptions(store.con,plant)
|
|
for m in months:
|
|
if m>month_key(decision):
|
|
peaks[m]=0.;contexts[m]={'kw':0.,'quality':'new_month','source':'new_month','observedAt':decision.isoformat()}
|
|
elif m in peaks:
|
|
r=store.con.execute('SELECT source,updated_at FROM planner_month_peaks WHERE plant=? AND month=?',(plant,m)).fetchone()
|
|
known_at=op['observedAt'] if m in op.get('measuredPeaks',{}) else r['updated_at']
|
|
contexts[m]={'kw':peaks[m],'quality':'verified','source':r['source'],'observedAt':known_at}
|
|
# A later acquired larger quarter may raise an older verified baseline,
|
|
# but the resulting combined planning basis must then say estimated.
|
|
if allow_estimates and m in estimates and estimates[m]['kw']>peaks[m] and utc(estimates[m]['observedAt'])>utc(known_at):
|
|
contexts[m]=basis_record(estimates[m],m,decision,allow_estimates=True)
|
|
peaks[m]=contexts[m]['kw']
|
|
input_quality['warnings'].append('A newer sampled quarter increased the earlier verified peak baseline; current planning maximum is estimated')
|
|
elif allow_estimates and m in estimates:
|
|
contexts[m]=basis_record(estimates[m],m,decision,allow_estimates=True);peaks[m]=contexts[m]['kw']
|
|
input_quality['warnings'].append('Monthly peak '+m+' is a planning estimate, not an authoritative billing maximum')
|
|
else:raise AwaitingInput('Peak basis missing; supply a verified maximum or explicitly permit a labelled estimate')
|
|
peak_prices=tariffs['peakChfKwMonth']
|
|
if not isinstance(peak_prices,dict):peak_prices={m:peak_prices for m in months}
|
|
lim=op['limits'];limits=Limits(lim['exportW'],lim['importW'],{int(k):v for k,v in lim['managerMonthLimitsW'].items()})
|
|
outlooks={};horizon_end=steps[-1].end
|
|
if settings['peakOutlookPolicy']=='empirical_if_available':
|
|
explicit=latest(store,plant,'peak_outlook')
|
|
for m in months:
|
|
if explicit and m in explicit['outlooks']:
|
|
candidate=RestMonthOutlook.from_dict(explicit['outlooks'][m])
|
|
try:
|
|
candidate.validate(decision,horizon_end)
|
|
policy=op.get('meterObservation',{}).get('controlPolicyId')
|
|
if not policy or candidate.control_policy_id!=policy:raise ValueError('Different or unknown control policy')
|
|
except ValueError:input_quality['warnings'].append('Stale, overlapping or incomparable rest-month outlook ignored; full incremental tariff used')
|
|
else:outlooks[m]=replace(candidate,reliance=min(candidate.reliance,settings['peakOutlookReliance']))
|
|
elif op.get('meterObservation'):
|
|
observation=op['meterObservation']
|
|
candidate=empirical_rest_month(meter_runtime.daily_peaks(store.con,plant,observation['meterId'],decision),
|
|
month=m,at=decision,horizon_end=horizon_end,control_policy_id=observation['controlPolicyId'],
|
|
reliance=settings['peakOutlookReliance'])
|
|
if candidate is not None:outlooks[m]=candidate
|
|
if not outlooks:input_quality['warnings'].append('Insufficient comparable rest-month history; full incremental peak tariff used, no arbitrary free peak allowance')
|
|
data={'steps':steps,'batteries':batteries(op['batteries']),'limits':limits,'observed_peaks':peaks,'peak_prices':peak_prices,
|
|
'quarter_history':past,'at':decision,'peak_context':contexts,'peak_outlooks':outlooks}
|
|
if source['loadBasis']=='base_load' and source.get('accountingEvidenceId'):
|
|
input_quality['accountingEvidenceId']=source['accountingEvidenceId']
|
|
return data,full_end,{'loadBasis':source['loadBasis'],**input_quality,**provenance(values)}
|
|
|
|
def run_once(store,now):
|
|
retention.maintain(store,now)
|
|
stamp=int(now.timestamp())//300
|
|
plants=[r[0] for r in store.con.execute('SELECT plant FROM planner_settings UNION SELECT DISTINCT plant FROM planner_input_current UNION SELECT plant FROM planner_data_sets')]
|
|
for plant in plants:
|
|
if not workflow.learning_enabled(store.con,plant):continue
|
|
config=store.settings(plant)
|
|
economic_replay.advance(store,plant,now)
|
|
datasets=[r[0] for r in store.con.execute('SELECT dataset FROM planner_data_sets WHERE plant=?',(plant,))]
|
|
for dataset in datasets:
|
|
try:measurement_pipeline.advance(store.con,plant,dataset,config,int(now.timestamp()))
|
|
except ValueError as exc:logging.getLogger(__name__).warning('Data pipeline unavailable for configured dataset: %s',type(exc).__name__)
|
|
with store.con:
|
|
old=store.con.execute('SELECT tick FROM planner_ticks WHERE plant=?',(plant,)).fetchone()
|
|
if not old or old[0]!=stamp:
|
|
store._request(plant,store.settings(plant)['revision'],'five_minute_tick',now)
|
|
store.con.execute('INSERT INTO planner_ticks VALUES(?,?) ON CONFLICT(plant) DO UPDATE SET tick=excluded.tick',(plant,stamp))
|
|
claim=store.claim(now)
|
|
if not claim:return {'status':'idle'}
|
|
plant=claim['plant'];result={'status':'internal_error','executable':False,'points':[]}
|
|
try:
|
|
if not workflow.learning_enabled(store.con,plant):
|
|
result={'status':'learning_disabled','executable':False,'points':[]}
|
|
return result
|
|
settings=store.settings(plant);previous=store.current(plant)
|
|
current=previous['sourceFamily'] if previous else store.registry.entries()[0].key
|
|
selection=choose_family(settings['family'],current,economic_replay.scores(store,plant,now),registry=store.registry,now=now,
|
|
lookback_days=settings['autoLookbackDays'],minimum_days=settings['autoMinimumDays'],
|
|
minimum_coverage=settings['autoMinimumCoverage'],margin_chf=settings['autoSwitchMarginChf'])
|
|
data,full_end,quality=assemble(store,plant,selection['family'],now,include_forecast_points=True)
|
|
data['batteries']=[replace(b,roundtrip_efficiency=settings['roundtripEfficiency']) for b in data['batteries']]
|
|
result=optimize(**data,config_revision=settings['revision'],family=selection['family'])
|
|
if result['executable']:
|
|
result['forecastPoints']=quality.pop('forecastPoints')
|
|
result.update({'installationId':plant,'inputRefs':quality.pop('inputRefs'),'controlContext':quality.pop('controlContext'),'runMode':'shadow','liveEnabled':False,'sourceSelection':selection,'forecastUntil':full_end.isoformat(),'pricesKnownUntil':result['validUntil'],'inputQuality':quality,'warnings':quality['warnings']+([] if quality['loadBasis']=='base_load' else ['Aggregate house forecast: base-load/SDL separation not verified; shadow only'])})
|
|
store.publish_shadow(plant,result,settings['revision'],now,claim['sequence'],claim['lease_token'])
|
|
economic_replay.capture(store,plant,now,assemble,optimize)
|
|
return result
|
|
except (ValueError,TypeError,KeyError) as exc:
|
|
result={'status':'awaiting_inputs' if isinstance(exc,AwaitingInput) else 'invalid_inputs','reason':str(exc)[:300],'executable':False,'points':[]}
|
|
return result
|
|
finally:
|
|
detail={k:v for k,v in result.items() if k in ('status','reason','planId','configRevision','sourceFamily')}
|
|
with store.con:store.con.execute('INSERT INTO planner_run_status VALUES(?,?,?,?) ON CONFLICT(plant) DO UPDATE SET updated_at=excluded.updated_at,status=excluded.status,detail=excluded.detail',(plant,now.isoformat(),result['status'],canonical(detail)))
|
|
store.finish(claim)
|
|
|
|
def status(store,plant,now):
|
|
plan=store.current(plant);settings=store.settings(plant)
|
|
row=store.con.execute('SELECT * FROM planner_run_status WHERE plant=?',(plant,)).fetchone()
|
|
pending=store.con.execute('SELECT revision,reasons,requested_at FROM planner_work WHERE plant=?',(plant,)).fetchone()
|
|
ack=store.con.execute('SELECT * FROM planner_ack WHERE plant=?',(plant,)).fetchone()
|
|
fresh=bool(plan and plan['configRevision']==settings['revision'] and utc(plan['validUntil'])>now and 0<=(now-utc(plan['generatedAt'])).total_seconds()<=900 and not pending and row and row['status'] in ('optimal','feasible_time_limit'))
|
|
state={'receiverProtocolVersion':1,'installationId':plant,'checkedAt':utc(now).isoformat(),'settings':settings,'peakPlanningBases':meter_runtime.assumptions(store.con,plant),'families':[asdict(f) for f in store.registry.entries()],'plan':plan,'fresh':fresh,'pending':dict(pending) if pending else None,'lastRun':{**dict(row),'detail':json.loads(row['detail'])} if row else None,'acknowledgement':dict(ack) if ack else None,'liveEnabled':False,'dataPipeline':measurement_pipeline.pipeline_status(store.con,plant),'economicComparison':economic_replay.status(store,plant),'maintenance':retention.status(store.con)}
|
|
state['workflow']=workflow.view(store,plant,state,now)
|
|
return state
|
|
|
|
def create_app(db_path,service_token,plants,*,start_worker=True,controlled_trial_plants=()):
|
|
allowed={str(UUID(p)) for p in plants}
|
|
trial_allowed={str(UUID(p)) for p in controlled_trial_plants}
|
|
if not trial_allowed <= allowed:raise ValueError('Trial allowlist must be a subset of plant allowlist')
|
|
if not service_token or len(service_token)<24:raise ValueError('Private service token required')
|
|
path=Path(db_path).resolve()
|
|
if path.name in ('users.db','portal.sqlite','settings.json'):raise ValueError('Dedicated planner database required')
|
|
path.parent.mkdir(parents=True,exist_ok=True);stop=Event()
|
|
def factory():return PlannerStore(str(path))
|
|
def loop():
|
|
while not stop.is_set():
|
|
s=factory()
|
|
try:run_once(s,datetime.now(timezone.utc).replace(microsecond=0))
|
|
except Exception as exc:logging.getLogger(__name__).error('V4 worker error: %s',type(exc).__name__)
|
|
finally:s.close()
|
|
stop.wait(1.)
|
|
@asynccontextmanager
|
|
async def lifespan(app):
|
|
thread=Thread(target=loop,name='v4-shadow',daemon=True)
|
|
if start_worker:thread.start()
|
|
yield
|
|
stop.set()
|
|
if start_worker:thread.join(35)
|
|
app=FastAPI(title='ENELIX V4 - Schattenbetrieb',lifespan=lifespan)
|
|
app.state.store_factory=factory
|
|
def authorize(plant,token,*,onboarding=False):
|
|
if not secrets.compare_digest(token or '',service_token):raise HTTPException(401,'Unauthorized')
|
|
try:
|
|
if str(UUID(plant))!=plant:raise ValueError('Canonical UUID required')
|
|
except ValueError:raise HTTPException(400,'Invalid installation ID')
|
|
s=factory()
|
|
if not onboarding and plant not in allowed and not workflow.enrolled(s.con,plant):
|
|
s.close();raise HTTPException(403,'Installation not enabled for planning')
|
|
return s
|
|
@app.post('/internal/v2/prognosis/{plant}/planner/setup')
|
|
def setup(plant:str,payload:dict,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token,onboarding=True)
|
|
try:return workflow.setup(s,plant,payload,datetime.now(timezone.utc))
|
|
except (ValueError,KeyError,TypeError,AttributeError) as exc:raise HTTPException(409 if 'immutable' in str(exc) else 400,str(exc)[:200])
|
|
finally:s.close()
|
|
@app.post('/internal/v2/prognosis/{plant}/planner/workflow')
|
|
def workflow_report(plant:str,payload:dict,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token)
|
|
try:return workflow.report(s,plant,payload,datetime.now(timezone.utc))
|
|
except (ValueError,KeyError,TypeError,AttributeError) as exc:raise HTTPException(400,str(exc)[:200])
|
|
finally:s.close()
|
|
@app.get('/health')
|
|
def health():return {'status':'ok','mode':'shadow','liveEnabled':False,'receiverProtocolVersion':1,'applicationRelease':'unified-rc1','mappingCompatibilityVersion':1,'economicReplayVersion':1,'archiveVersion':1}
|
|
@app.get('/internal/v2/planner/installations')
|
|
def planning_installations(token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
if not secrets.compare_digest(token or '',service_token):raise HTTPException(401,'Unauthorized')
|
|
s=factory()
|
|
try:
|
|
registered={r[0] for r in s.con.execute('SELECT plant FROM planner_enrollment WHERE learning_enabled=1')}
|
|
paused={r[0] for r in s.con.execute('SELECT plant FROM planner_enrollment WHERE learning_enabled=0')}
|
|
return {'installationIds':sorted((allowed|registered)-paused),'controlEnabled':False}
|
|
finally:s.close()
|
|
@app.get('/internal/v2/prognosis/{plant}/planner')
|
|
def read(plant:str,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token)
|
|
try:
|
|
now=datetime.now(timezone.utc);view=status(s,plant,now)
|
|
view['controlledTrial']=controlled_trial.authority(s,plant,view,now,trial_allowed)
|
|
view['controlledTrialAuthorized']=view['controlledTrial'] is not None
|
|
return view
|
|
finally:s.close()
|
|
@app.post('/internal/v2/prognosis/{plant}/planner/trial/arm')
|
|
def arm_trial(plant:str,payload:dict,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token)
|
|
try:
|
|
now=datetime.now(timezone.utc)
|
|
return controlled_trial.arm(s,plant,payload,status(s,plant,now),now,trial_allowed)
|
|
except (ValueError,KeyError,TypeError) as exc:raise HTTPException(409,str(exc)[:300])
|
|
finally:s.close()
|
|
@app.post('/internal/v2/prognosis/{plant}/planner/trial/revoke')
|
|
def revoke_trial(plant:str,payload:dict,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token)
|
|
try:
|
|
if set(payload)!={'sessionId'}:raise ValueError('Session ID only')
|
|
return controlled_trial.revoke(s,plant,payload['sessionId'],datetime.now(timezone.utc))
|
|
except (ValueError,KeyError,TypeError) as exc:raise HTTPException(400,str(exc)[:300])
|
|
finally:s.close()
|
|
@app.put('/internal/v2/prognosis/{plant}/planner/datasets/{dataset}/mapping-compatibility')
|
|
def mapping_compatibility(plant:str,dataset:str,payload:dict,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token)
|
|
try:
|
|
from .mapping_identity import register
|
|
cfg=measurement_pipeline.configuration(s.con,plant,dataset)
|
|
return register(s.con,plant,cfg,payload,int(datetime.now(timezone.utc).timestamp()))
|
|
except (ValueError,KeyError,TypeError,AttributeError):
|
|
raise HTTPException(400, 'Mapping compatibility proof rejected')
|
|
finally:s.close()
|
|
|
|
@app.put('/internal/v2/prognosis/{plant}/planner/settings')
|
|
def save(plant:str,payload:dict,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token)
|
|
try:return s.save_settings(plant,payload['changes'],payload['expectedRevision'],datetime.now(timezone.utc))
|
|
except (ValueError,KeyError,TypeError) as exc:raise HTTPException(409 if 'Revision conflict' in str(exc) else 400,str(exc))
|
|
finally:s.close()
|
|
@app.post('/internal/v2/prognosis/{plant}/planner/replan')
|
|
def replan(plant:str,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token)
|
|
try:s.request(plant,'manual',datetime.now(timezone.utc));return {'status':'queued','liveEnabled':False}
|
|
finally:s.close()
|
|
@app.post('/internal/v2/prognosis/{plant}/planner/inputs/{kind}')
|
|
def input_event(plant:str,kind:str,payload:dict,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token)
|
|
try:return ingest(s,plant,kind,payload,datetime.now(timezone.utc))
|
|
except (ValueError,KeyError,TypeError,AttributeError) as exc:raise HTTPException(400,str(exc)[:300])
|
|
finally:s.close()
|
|
@app.put('/internal/v2/prognosis/{plant}/planner/datasets/{dataset}')
|
|
def configure_dataset(plant:str,dataset:str,payload:dict,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token)
|
|
try:
|
|
if payload.get('datasetId')!=dataset:raise ValueError('Dataset path/payload mismatch')
|
|
return measurement_pipeline.register_dataset(s.con,plant,payload,int(datetime.now(timezone.utc).timestamp()))
|
|
except (ValueError,KeyError,TypeError) as exc:raise HTTPException(400,str(exc)[:200])
|
|
finally:s.close()
|
|
@app.post('/internal/v2/prognosis/{plant}/planner/measurements')
|
|
def measurement_batch(plant:str,payload:dict,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token)
|
|
try:
|
|
if not workflow.learning_enabled(s.con,plant):raise ValueError('Learning is disabled for this installation')
|
|
now=datetime.now(timezone.utc)
|
|
result=measurement_pipeline.ingest_batch(s.con,plant,payload,int(now.timestamp()))
|
|
# Existing five-minute worker handles rollup/training; no per-record optimizer flood.
|
|
return result
|
|
except (ValueError,KeyError,TypeError) as exc:raise HTTPException(400,str(exc)[:200])
|
|
finally:s.close()
|
|
@app.post('/internal/v2/prognosis/{plant}/planner/ack')
|
|
def ack(plant:str,payload:dict,token:str=Header(default='',alias='X-Enelix-Service-Token')):
|
|
s=authorize(plant,token)
|
|
try:s.acknowledge(plant,payload['planId'],payload['revision'],datetime.now(timezone.utc),payload.get('step'),payload.get('status','shadow_seen'));return {'status':'recorded'}
|
|
except (ValueError,KeyError,TypeError) as exc:raise HTTPException(400,str(exc)[:300])
|
|
finally:s.close()
|
|
from starlette.responses import JSONResponse
|
|
class BodyLimit:
|
|
def __init__(self,app):self.app=app
|
|
async def __call__(self,scope,receive,send):
|
|
if scope['type']!='http' or scope['method'] not in ('POST','PUT'):return await self.app(scope,receive,send)
|
|
chunks=[];total=0
|
|
while True:
|
|
msg=await receive()
|
|
if msg['type']=='http.disconnect':return
|
|
total+=len(msg.get('body',b''))
|
|
if total>2000000:return await JSONResponse({'detail':'Request too large'},status_code=413)(scope,receive,send)
|
|
chunks.append(msg)
|
|
if not msg.get('more_body',False):break
|
|
async def replay():return chunks.pop(0) if chunks else await receive()
|
|
return await self.app(scope,replay,send)
|
|
app.add_middleware(BodyLimit)
|
|
return app
|
|
|
|
def from_environment():
|
|
return create_app(os.environ.get('NETPLAN_V4_DB','/data/netplan-v4.sqlite'),os.environ.get('PROGNOSIS_SERVICE_TOKEN',''),[p.strip() for p in os.environ.get('NETPLAN_V4_PLANTS','').split(',') if p.strip()], controlled_trial_plants=[p.strip() for p in os.environ.get('NETPLAN_V4_CONTROL_TRIAL_PLANTS','').split(',') if p.strip()])
|