"""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 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']=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): 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; 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') 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: 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: 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) 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.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 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')) return {'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)} 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 allowed or not service_token or len(service_token)<24:raise ValueError('Private service token and explicit plant allowlist 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): if not secrets.compare_digest(token or '',service_token):raise HTTPException(401,'Unauthorized') try:plant=str(UUID(plant)) except ValueError:raise HTTPException(400,'Invalid installation ID') if plant not in allowed:raise HTTPException(403,'Installation not enabled for shadow trial') return factory() @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/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: 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()])