"""Persistent, explicitly estimated operational demand tracking. No inference from a configured cap. Samples are device-reception observations; last-value integration is a labelled control estimate, not settlement metering. No reset, stale sample or long communication gap is bridged silently. """ from datetime import datetime,timedelta,timezone import json from .domain import number,utc,month_key,quarter_start,ZURICH from .peak_policy import basis_record def schema(con): con.executescript(''' CREATE TABLE IF NOT EXISTS planner_peak_assumptions( plant TEXT NOT NULL,month TEXT NOT NULL,kw REAL NOT NULL,value TEXT NOT NULL, PRIMARY KEY(plant,month)); CREATE TABLE IF NOT EXISTS planner_runtime_samples( plant TEXT NOT NULL,meter_id TEXT NOT NULL,at INTEGER NOT NULL, power_w REAL NOT NULL,total_kwh REAL,policy_id TEXT NOT NULL, PRIMARY KEY(plant,meter_id,at)); CREATE TABLE IF NOT EXISTS planner_runtime_quarters( plant TEXT NOT NULL,meter_id TEXT NOT NULL,start INTEGER NOT NULL, import_kwh REAL NOT NULL,policy_id TEXT NOT NULL,value TEXT NOT NULL, PRIMARY KEY(plant,meter_id,start)); CREATE INDEX IF NOT EXISTS idx_runtime_samples_time ON planner_runtime_samples(plant,at); ''') def validate_observation(obs, observed_at): if not isinstance(obs,dict) or set(obs)!={'meterId','sampleAt','powerW','totalImportKwh','controlPolicyId'}: raise ValueError('Explicit meter observation schema required') import re if not isinstance(obs['meterId'],str) or not re.fullmatch(r'symcon-active-import:[0-9a-f]{64}',obs['meterId']): raise ValueError('Active import identity required') if not isinstance(obs['controlPolicyId'],str) or not 1<=len(obs['controlPolicyId'])<=160: raise ValueError('Control policy identity required') t=utc(obs['sampleAt']) if t.microsecond or not 0<=(utc(observed_at)-t).total_seconds()<=60: raise ValueError('Fresh whole-second acquisition timestamp required') number(obs['powerW'],'meter power',-1e9,1e9) number(obs['totalImportKwh'],'active import total',0,1e12) def integrate_power(samples,start,end,max_gap=120): """samples sorted tuples (epoch,power,total,policy). Return None on gaps/reset. A source can keep the same value while receiving fresh telemetry. The caller records all acquisitions, not only changes. A quarter is never labelled exact. """ if end<=start:return None rows=sorted(samples,key=lambda r:r[0]) if len(rows)<2:return None covered=energy=0.;policies=set();largest_gap=0 previous=None for row in rows: if previous is not None: ta,pa,ea,pola=previous;tb,pb,eb,polb=row if tb<=ta:return None left=max(start,ta);right=min(end,tb) if right>left: if tb-ta>max_gap or pola!=polb:return None if ea is not None and eb is not None and eb=validated['kw']:return old con.execute('INSERT INTO planner_peak_assumptions VALUES(?,?,?,?) ON CONFLICT(plant,month) DO UPDATE SET kw=excluded.kw,value=excluded.value', (plant,month,validated['kw'],json.dumps(validated,sort_keys=True,allow_nan=False))) return validated def assumptions(con,plant): return {r[0]:json.loads(r[1]) for r in con.execute('SELECT month,value FROM planner_peak_assumptions WHERE plant=?',(plant,))} def observe(con,plant,obs,at): validate_observation(obs,at) mid=obs['meterId'];stamp=int(utc(obs['sampleAt']).timestamp()) existing=con.execute('SELECT power_w,total_kwh,policy_id FROM planner_runtime_samples WHERE plant=? AND meter_id=? AND at=?',(plant,mid,stamp)).fetchone() values=(float(obs['powerW']),float(obs['totalImportKwh']),obs['controlPolicyId']) if existing and tuple(existing)!=values:raise ValueError('Conflicting meter acquisition at same time') con.execute('INSERT OR IGNORE INTO planner_runtime_samples VALUES(?,?,?,?,?,?)',(plant,mid,stamp,*values)) # Only a just-completed quarter is finalised; late acquisition never promotes # a historical gap to complete data without all supporting samples. current=stamp//900*900 rows=[tuple(r) for r in con.execute('SELECT at,power_w,total_kwh,policy_id FROM planner_runtime_samples WHERE plant=? AND meter_id=? AND at>=? AND at<=? ORDER BY at', (plant,mid,current-900-120,stamp))] report=integrate_power(rows,current-900,current) if report: quarter=current-900 previous=con.execute('SELECT import_kwh FROM planner_runtime_quarters WHERE plant=? AND meter_id=? AND start=?',(plant,mid,quarter)).fetchone() if previous is None: con.execute('INSERT INTO planner_runtime_quarters VALUES(?,?,?,?,?,?)',(plant,mid,quarter,report['importKwh'],report['controlPolicyId'],json.dumps(report))) m=month_key(datetime.fromtimestamp(quarter,timezone.utc)) save_assumption(con,plant,m,{'kw':report['importKwh']/.25,'quality':'estimated','source':'sampled_power_estimate', 'observedAt':utc(at).isoformat(),'meterId':mid,'notes':'Maximum of available sampled quarters; earlier month may be incomplete'},at) # Samples are retained for reproducibility; production retention job required. return report def current_quarter(con,plant,obs,decision): q=quarter_start(decision);start=int(q.timestamp());end=int(utc(decision).timestamp()) if start==end:return None rows=[tuple(r) for r in con.execute('SELECT at,power_w,total_kwh,policy_id FROM planner_runtime_samples WHERE plant=? AND meter_id=? AND at>=? AND at<=? ORDER BY at', (plant,obs['meterId'],start-120,end))] report=integrate_power(rows,start,end) if not report:return None return {'start':q.isoformat(),'measuredSeconds':end-start,'importKwh':report['importKwh'], 'quality':'estimated','source':'sampled_power_estimate','coverage':1.0,'notes':'Acquisition power estimate, not an exact billing counter boundary'} def daily_peaks(con,plant,meter_id,now): rows=con.execute('SELECT start,import_kwh,policy_id FROM planner_runtime_quarters WHERE plant=? AND meter_id=? AND start>=? ORDER BY start', (plant,meter_id,int((utc(now)-timedelta(days=91)).timestamp()))) grouped={} for start,energy,policy in rows: t=datetime.fromtimestamp(start,timezone.utc).astimezone(ZURICH) grouped.setdefault((t.strftime('%Y-%m-%d'),policy),[]).append((start,energy)) result=[] for (day,policy),items in grouped.items(): begin=datetime.strptime(day,'%Y-%m-%d').replace(tzinfo=ZURICH);end=begin+timedelta(days=1) if utc(end)>utc(now):continue expected=set(range(int(begin.timestamp()),int(end.timestamp()),900)) complete={t for t,e in items}==expected result.append({'day':day,'peakKw':max(e/.25 for t,e in items),'controlPolicyId':policy,'complete':complete, 'observedAt':utc(end).isoformat(),'quality':'estimated'}) return result