Files
Enelix-EMS/services/netplan-v4/netplan_v4/measurement_pipeline.py
T

514 lines
31 KiB
Python

"""Application data path: versioned numeric observations -> physical load -> trained profiles.
Lives in the existing planner service/database; no separate diagnostic service.
The device may append only to an operator-configured dataset. Original observations,
model revisions and prediction vintages are preserved. Output is never an actuator grant.
"""
from __future__ import annotations
from bisect import bisect_right
from collections import defaultdict
from datetime import datetime, timedelta, timezone
from hashlib import sha256
from math import isfinite
from statistics import median
from zoneinfo import ZoneInfo
import json
from .history_timing import validate_policy, endpoint_bridges
from . import mapping_identity
UTC = timezone.utc
LOCAL = ZoneInfo('Europe/Zurich')
FAMILIES = ('3', '13', '23')
def canonical(value):
return json.dumps(value, sort_keys=True, separators=(',', ':'), allow_nan=False)
def epoch(value):
if not isinstance(value, str):
raise ValueError('UTC timestamp required')
t = datetime.fromisoformat(value.replace('Z', '+00:00'))
if t.tzinfo is None or t.utcoffset().total_seconds() != 0 or t.microsecond:
raise ValueError('Explicit whole-second UTC timestamp required')
return int(t.timestamp())
def iso(t):
return datetime.fromtimestamp(t, UTC).isoformat()
def numeric(value, bound=1e12):
return type(value) in (int, float) and isfinite(value) and abs(value) <= bound
def schema(con):
mapping_identity.schema(con)
con.executescript('''
CREATE TABLE IF NOT EXISTS planner_data_sets(
plant TEXT NOT NULL, dataset TEXT NOT NULL, config TEXT NOT NULL,
created_at INTEGER NOT NULL, PRIMARY KEY(plant,dataset));
CREATE TABLE IF NOT EXISTS planner_observations(
plant TEXT NOT NULL, dataset TEXT NOT NULL, captured_at INTEGER NOT NULL,
received_at INTEGER NOT NULL, fingerprint TEXT NOT NULL, value TEXT NOT NULL,
PRIMARY KEY(plant,dataset,captured_at));
CREATE TABLE IF NOT EXISTS planner_load_windows(
plant TEXT NOT NULL, dataset TEXT NOT NULL, start INTEGER NOT NULL,
available_at INTEGER NOT NULL, coverage REAL NOT NULL, value TEXT NOT NULL,
PRIMARY KEY(plant,dataset,start));
CREATE TABLE IF NOT EXISTS planner_load_models(
plant TEXT NOT NULL, dataset TEXT NOT NULL, model_id TEXT NOT NULL,
trained_at INTEGER NOT NULL, trained_through INTEGER NOT NULL, value TEXT NOT NULL,
PRIMARY KEY(plant,dataset,model_id));
CREATE TABLE IF NOT EXISTS planner_model_current(
plant TEXT NOT NULL, dataset TEXT NOT NULL, model_id TEXT NOT NULL,
PRIMARY KEY(plant,dataset));
CREATE TABLE IF NOT EXISTS planner_pipeline_state(
plant TEXT NOT NULL, dataset TEXT NOT NULL, tick INTEGER NOT NULL,
status TEXT NOT NULL, detail TEXT NOT NULL, PRIMARY KEY(plant,dataset));
CREATE TABLE IF NOT EXISTS planner_prediction_vintages(
plant TEXT NOT NULL, dataset TEXT NOT NULL, issued_at INTEGER NOT NULL,
target INTEGER NOT NULL, family TEXT NOT NULL, model_id TEXT NOT NULL,
load_w REAL NOT NULL, pv_w REAL NOT NULL,
PRIMARY KEY(plant,dataset,issued_at,target,family));
CREATE INDEX IF NOT EXISTS planner_observation_window
ON planner_observations(plant,dataset,captured_at);
CREATE INDEX IF NOT EXISTS planner_prediction_target
ON planner_prediction_vintages(plant,dataset,target);
''')
def validate_config(c):
fields = {'datasetId', 'mappingSha256', 'inventorySha256', 'sources',
'formula', 'solarReference', 'minimumCoverage', 'maximumGapSeconds',
'minimumTrainingHours', 'historyDays'}
optional = {'sourceDatasetId', 'historyTimingPolicy'}
if not isinstance(c, dict) or not fields <= set(c) or set(c)-fields-optional:
raise ValueError('Explicit dataset configuration required')
name = c['datasetId']
if not isinstance(name, str) or not 1 <= len(name) <= 80 or any(x not in 'abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789-_' for x in name):
raise ValueError('Invalid dataset ID')
for field in ('mappingSha256', 'inventorySha256'):
h = c[field]
if not isinstance(h, str) or len(h) != 64 or any(x not in '0123456789abcdef' for x in h):
raise ValueError('Explicit mapping/inventory fingerprint required')
if c['formula'] not in ('physical_sum_v1', 'solar_terminal_v1'):
raise ValueError('Unknown physical formula')
if not numeric(c['minimumCoverage']) or not .90 <= c['minimumCoverage'] <= 1:
raise ValueError('Coverage must be .90..1; recorded gaps remain visible')
for field, lo, hi in (('maximumGapSeconds', 1, 10), ('minimumTrainingHours', 1, 168), ('historyDays', 2, 90)):
if type(c[field]) is not int or not lo <= c[field] <= hi:
raise ValueError('Invalid '+field)
sources = c['sources']
if not isinstance(sources, list) or not 3 <= len(sources) <= 80:
raise ValueError('Source list required')
seen, ids = set(), set()
roles = {'grid', 'pv', 'physical_storage', 'flexible_load', 'reference', 'sdl_request', 'solar_raw', 'solar_scale'}
for s in sources:
if set(s) != {'key', 'variableId', 'role', 'factorToW', 'maxAgeSeconds'}:
raise ValueError('Explicit source definition required')
k = s['key']
if not isinstance(k, str) or not 1 <= len(k) <= 64 or k in seen or type(s['variableId']) is not int or not 1 <= s['variableId'] <= 99999 or s['variableId'] in ids:
raise ValueError('Duplicate/invalid source')
if s['role'] not in roles or not numeric(s['factorToW'], 1e6) or s['factorToW'] == 0:
raise ValueError('Source role/factor invalid')
if type(s['maxAgeSeconds']) is not int or not 1 <= s['maxAgeSeconds'] <= 300:
raise ValueError('Source lifetime invalid')
seen.add(k); ids.add(s['variableId'])
if sum(s['role'] == 'grid' for s in sources) != 1 or not any(s['role'] == 'pv' for s in sources):
raise ValueError('Grid and PV measurement sources required')
sr = c['solarReference']
if c['formula'] == 'solar_terminal_v1':
if not isinstance(sr, dict) or set(sr) != {'pvKey', 'batteryKey', 'rawKey', 'scaleKey'}:
raise ValueError('Solar terminal sources required')
bykey = {s['key']: s['role'] for s in sources}
if any(bykey.get(sr[k]) != role for k, role in (('pvKey','pv'),('batteryKey','physical_storage'),('rawKey','solar_raw'),('scaleKey','solar_scale'))):
raise ValueError('Solar origin roles mismatch')
elif sr is not None:
raise ValueError('No unused solar mapping allowed')
validate_policy(c.get('historyTimingPolicy'), sources)
source = c.get('sourceDatasetId')
if source is not None and (not isinstance(source, str) or not 1 <= len(source) <= 80 or source == name or any(x not in 'abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789-_' for x in source)):
raise ValueError('Invalid source dataset reference')
canonical(c)
return c
def register_dataset(con, plant, c, now):
"""Operator endpoint only; device append endpoint cannot change units or limits."""
validate_config(c)
if c.get('sourceDatasetId'):
origin = configuration(con, plant, c['sourceDatasetId'])
if origin.get('sourceDatasetId'):
raise ValueError('Dataset reference chains are not allowed')
comparable = lambda v: {k:x for k,x in v.items() if k not in ('datasetId','sourceDatasetId','historyTimingPolicy')}
if canonical(comparable(origin)) != canonical(comparable(c)):
raise ValueError('Referenced observations must retain identical measurement meaning')
value = canonical(c)
con.execute('BEGIN IMMEDIATE')
try:
old = con.execute('SELECT config FROM planner_data_sets WHERE plant=? AND dataset=?', (plant,c['datasetId'])).fetchone()
if old and old[0] != value:
raise ValueError('Dataset is immutable; use a new datasetId for changed measurement meaning')
con.execute('INSERT OR IGNORE INTO planner_data_sets VALUES(?,?,?,?)', (plant,c['datasetId'],value,now))
con.commit()
except Exception:
con.rollback(); raise
return {'status':'configured', 'datasetId':c['datasetId'], 'mappingSha256':c['mappingSha256'], 'controlEnabled':False}
def configuration(con, plant, dataset):
row = con.execute('SELECT config FROM planner_data_sets WHERE plant=? AND dataset=?', (plant,dataset)).fetchone()
if row is None:
raise ValueError('Dataset not configured for this installation')
return json.loads(row[0])
def project(record, c, plant, now, approved_mappings=()):
if not isinstance(record, dict) or type(record.get('schemaVersion')) is not int or record.get('schemaVersion') != 1 or record.get('kind') != 'raw_accounting_capture' or record.get('installationId') != plant:
raise ValueError('Wrong capture identity')
incoming = record.get('mappingSha256')
if not isinstance(incoming, str) or (incoming != c['mappingSha256'] and incoming not in approved_mappings) or record.get('reportedInventorySha256') != c['inventorySha256']:
raise ValueError('Wrong capture mapping or inventory')
t = epoch(record.get('capturedAt')); start = epoch(record.get('captureStartedAt'))
if start > t or t > now+30 or t < now-90*86400:
raise ValueError('Capture timestamp outside permitted range')
if not isinstance(record.get('raw'), dict):
raise ValueError('Numeric raw observations required')
out = {}; issues = []
for s in c['sources']:
r = record['raw'].get(s['key'], {})
if not isinstance(r, dict):
r = {}
v, at = r.get('value'), r.get('sourceUpdatedAt')
good = r.get('variableId') == s['variableId'] and numeric(v) and type(at) is int and 0 < at <= t and r.get('issues') == []
if not good:
v = at = None
issues.append(s['key'])
# Unknown/free-text fields, credentials, client quality claims never persisted.
out[s['key']] = {'value':v, 'sourceUpdatedAt':at, 'valid':bool(good)}
return {'capturedAt':t, 'captureDurationSeconds':t-start, 'raw':out, 'invalidSources':issues}
def ingest_batch(con, plant, payload, now):
if not isinstance(payload,dict) or set(payload) != {'version','datasetId','records'} or type(payload.get('version')) is not int or payload['version'] != 1:
raise ValueError('Measurement batch version/fields invalid')
c = configuration(con,plant,payload['datasetId'])
if c.get('sourceDatasetId'):
raise ValueError('Derived dataset is read-only; append to the original measurement dataset')
records = payload['records']
if not isinstance(records,list) or not 1 <= len(records) <= 120:
raise ValueError('Batch requires 1..120 captures')
aliases = mapping_identity.approved(con, plant, c)
rows = [project(r,c,plant,now,aliases) for r in records]
if any(a['capturedAt'] >= b['capturedAt'] for a,b in zip(rows,rows[1:])):
raise ValueError('Batch must be in increasing capture order')
stored = duplicate = 0
con.execute('BEGIN IMMEDIATE')
try:
for original, r in zip(records, rows):
value = canonical(r); digest = sha256(value.encode()).hexdigest()
old = con.execute('SELECT fingerprint FROM planner_observations WHERE plant=? AND dataset=? AND captured_at=?', (plant,c['datasetId'],r['capturedAt'])).fetchone()
if old:
if old[0] != digest:
raise ValueError('Conflicting immutable observation')
duplicate += 1
else:
con.execute('INSERT INTO planner_observations VALUES(?,?,?,?,?,?)',(plant,c['datasetId'],r['capturedAt'],now,digest,value)); stored += 1
mapping_identity.save_origin(con, plant, c, original, r['capturedAt'], now, aliases)
con.commit()
except Exception:
con.rollback(); raise
return {'status':'stored' if stored else 'duplicate','stored':stored,'duplicates':duplicate,
'acceptedThrough':iso(rows[-1]['capturedAt']),'datasetId':c['datasetId'],'controlEnabled':False}
def physical_value(values, c):
"""Same explicit sign convention as configured acquisition. No virtual power in load."""
total = defaultdict(float)
sr = c['solarReference']
for s in c['sources']:
if s['role'] not in ('grid','pv','physical_storage','flexible_load'):
continue
if sr and s['key'] in (sr['pvKey'],sr['batteryKey']):
continue
val = values[s['key']]*s['factorToW']
if not numeric(val,1e9) or (s['role'] in ('pv','flexible_load') and val < 0):
raise ValueError('Invalid physical power')
total[s['role']] += val
solar = 0.
if sr:
raw, sf = values[sr['rawKey']], values[sr['scaleKey']]
if int(raw) != raw or not -32768 < raw <= 32767 or int(sf) != sf or not -6 <= sf <= 6:
raise ValueError('Invalid solar power/scaling sentinel')
solar = raw*10**int(sf)
load = total['grid']+total['pv']-total['physical_storage']-total['flexible_load']+solar
if not numeric(load,1e9) or load < 0:
raise ValueError('Negative/nonfinite physical load')
return load
def reconstruct(records, c, *, projection=None):
"""Bounded retrospective estimation, never a real-time feedback signal.
Missing observations split support. Source timestamps are not refreshed. Small
uncovered portions remain quantified and are never filled with zero.
"""
if len(records) < 2:
return []
if any(a['capturedAt'] >= b['capturedAt'] for a,b in zip(records,records[1:])):
raise ValueError('Capture sequence not ordered')
sr = c['solarReference']
primary = {s['key']:s for s in c['sources'] if s['role'] in ('grid','pv','physical_storage','flexible_load')}
if sr:
primary.pop(sr['pvKey']); primary.pop(sr['batteryKey'])
for s in c['sources']:
if s['key'] in (sr['rawKey'],sr['scaleKey']): primary[s['key']] = s
timing_policy = validate_policy(c.get('historyTimingPolicy'), c['sources'])
if projection is not None:
configured={s['key']:s for s in c['sources']}
if set(projection)!= {'keys','calculate'} or not callable(projection['calculate']) or not projection['keys'] or not set(projection['keys']) <= set(configured):
raise ValueError('Invalid internal projection')
primary={k:configured[k] for k in projection['keys']}
first,last = records[0]['capturedAt'],records[-1]['capturedAt']
series = {k:{} for k in primary}; blocks = {k:[] for k in primary}; gaps = []
first_observed = {k:{} for k in primary}
pending = {k:None for k in primary}; high = {k:0 for k in primary}
edges = {first,last}
for a,b in zip(records,records[1:]):
if b['capturedAt']-a['capturedAt'] > 45:
gaps.append((a['capturedAt'],b['capturedAt']));edges.update(gaps[-1])
for r in records:
at = r['capturedAt']
for k in primary:
v = r['raw'].get(k,{})
t = v.get('sourceUpdatedAt')
if not v.get('valid') or not numeric(v.get('value')) or type(t) is not int or t > at or t < high[k] or r['captureDurationSeconds'] > 5:
if pending[k] is None: pending[k] = at
continue
high[k] = max(high[k],t)
if pending[k] is not None:
blocks[k].append((pending[k],at));edges.update(blocks[k][-1]);pending[k] = None
if t in series[k] and series[k][t] != v['value']:
series[k][t] = None
else:
series[k].setdefault(t,v['value'])
first_observed[k].setdefault(t, max(at, r.get('_receivedAt', at)))
for k,s in primary.items():
if pending[k] is not None:
blocks[k].append((pending[k],last));edges.update(blocks[k][-1])
for t in series[k]: edges.update((t,t+s['maxAgeSeconds']))
edges.update(range(first//300*300+300,last,300))
edges = sorted(x for x in edges if first <= x <= last)
knots = {k:sorted(v) for k,v in series.items()}
bridges = endpoint_bridges(series, first_observed, blocks, gaps, timing_policy, primary)
bins = {}
for a,b in zip(edges,edges[1:]):
start = a//300*300; item = bins.setdefault(start,{'start':start,'seconds':0,'wattSeconds':0.,'maxGapSeconds':0,'currentGap':0,
'publicationEstimatedSeconds':0,'publicationEstimatedBySource':{},'knownAt':0,'missingSourceSeconds':{}})
vals = {}; usable = not any(x <= a < y for x,y in gaps)
extended = []; unavailable = []; known_at = 0
for k,s in primary.items():
pos = bisect_right(knots[k],a)-1
t = knots[k][pos] if pos >= 0 else None
bridge = bridges[k].get(t)
expired = t is None or a >= t+s['maxAgeSeconds']
supported_tail = bool(bridge and a < bridge['end'])
if t is None or (expired and not supported_tail) or series[k][t] is None or any(x <= a < y for x,y in blocks[k]):
usable = False; unavailable.append(k)
else:
vals[k] = series[k][t]
known_at = max(known_at, first_observed[k][t])
if expired:
extended.append(k); known_at = max(known_at, bridge['availableAt'])
load = None
if usable:
try:
load = physical_value(vals,c) if projection is None else projection['calculate'](vals,c)
if not numeric(load,1e9):raise ValueError('Invalid historical projection')
except ValueError: usable = False
if usable:
item['seconds'] += b-a; item['wattSeconds'] += load*(b-a);item['currentGap'] = 0
item['knownAt'] = max(item['knownAt'], known_at)
if extended: item['publicationEstimatedSeconds'] += b-a
for key in extended: item['publicationEstimatedBySource'][key] = item['publicationEstimatedBySource'].get(key,0)+b-a
else:
item['currentGap'] += b-a; item['maxGapSeconds'] = max(item['maxGapSeconds'],item['currentGap'])
for key in unavailable: item['missingSourceSeconds'][key] = item['missingSourceSeconds'].get(key,0)+b-a
out=[]
for t,item in sorted(bins.items()):
# Partial beginning/end bins remain diagnostic and cannot train.
complete_extent = first <= t and last >= t+300
coverage = item['seconds']/300
eligible = complete_extent and coverage >= c['minimumCoverage'] and item['maxGapSeconds'] <= c['maximumGapSeconds']
out.append({'start':t,'coverage':coverage,'coveredSeconds':item['seconds'],'maxGapSeconds':item['maxGapSeconds'],
'loadW':item['wattSeconds']/item['seconds'] if item['seconds'] else None,
'profileUsable':eligible,'estimated':True,'fullPhysicalIntervalMeasured':False,
'meterBoundaryVerified':False,'method':c['formula'],
'historyTimingMethod':(timing_policy or {}).get('method','strict_expiry'),
'publicationEstimatedSeconds':item['publicationEstimatedSeconds'],
'publicationEstimatedBySourceSeconds':item['publicationEstimatedBySource'],
'missingSourceSeconds':item['missingSourceSeconds'],
'availableNotBefore':max(t+300,item['knownAt'])})
return out
def slot(t):
local = datetime.fromtimestamp(t,UTC).astimezone(LOCAL)
return local.hour*12+local.minute//5
def build_profiles(rows):
samples=defaultdict(list); recent=defaultdict(list); weekend={False:defaultdict(list),True:defaultdict(list)}
anchor=max(r['start'] for r in rows)
for r in rows:
i=slot(r['start']); v=r['loadW']; samples[i].append(v)
if anchor-r['start'] < 86400: recent[i].append(v)
weekend[datetime.fromtimestamp(r['start'],UTC).astimezone(LOCAL).weekday()>=5][i].append(v)
overall = median([r['loadW'] for r in rows])
def profile(values):
# Missing calendar slots are a model estimate, not invented historical measurements.
result=[]
for i in range(288):
local=values.get(i,[])
if not local:
local=[v for j in ((i-2)%288,(i-1)%288,(i+1)%288,(i+2)%288) for v in values.get(j,[])]
result.append(float(median(local)) if local else float(overall))
return result
return {'3':profile(samples),'13':profile(recent),'23':{'weekday':profile(weekend[False] or samples),'weekend':profile(weekend[True] or samples)},
'slotCoverage':len(samples)/288}
def predict(model, family, t):
p=model['profiles'][family]
if family=='23': p=p['weekend' if datetime.fromtimestamp(t,UTC).astimezone(LOCAL).weekday()>=5 else 'weekday']
return p[slot(t)]
def advance(con, plant, dataset, settings, now):
"""Called by the existing worker; bounded data/model update once per five-minute tick."""
c=configuration(con,plant,dataset); tick=now//300
old=con.execute('SELECT tick FROM planner_pipeline_state WHERE plant=? AND dataset=?',(plant,dataset)).fetchone()
if old and old[0]==tick: return
observation_dataset=c.get('sourceDatasetId',dataset)
fetched=con.execute('SELECT value,received_at FROM planner_observations WHERE plant=? AND dataset=? AND captured_at>=? AND captured_at<=? AND received_at<=? ORDER BY captured_at',
(plant,observation_dataset,now-172800-300,now,now)).fetchall()
records=[{**json.loads(r[0]),'_receivedAt':r[1]} for r in fetched]
windows=reconstruct(records,c)
with con:
for w in windows:
if w['start']+300 > now-30 or w.get('availableNotBefore',0)>now: continue
con.execute('INSERT INTO planner_load_windows VALUES(?,?,?,?,?,?) ON CONFLICT(plant,dataset,start) DO UPDATE SET available_at=excluded.available_at,coverage=excluded.coverage,value=excluded.value',
(plant,dataset,w['start'],now,w['coverage'],canonical(w)))
rows=[json.loads(r[0]) for r in con.execute('SELECT value FROM planner_load_windows WHERE plant=? AND dataset=? AND start>=? AND start+300<=? ORDER BY start',(plant,dataset,now-c['historyDays']*86400,now))]
good=[r for r in rows if r['profileUsable']]
active=current_model(con,plant,dataset,now)
cadence=86400 if settings['trainingCadence']=='daily' else 604800
detail={'observationsInLast48h':len(records),'usableWindows':len(good),'requiredEquivalentHours':c['minimumTrainingHours'],
'usableEquivalentHours':sum(r['coverage'] for r in good)/12,'datasetId':dataset,'trainingCadence':settings['trainingCadence'],
'automaticTrainingConnected':True,'liveEnabled':False,'sourceIsConfiguredEstimate':True,
'observationDatasetId':observation_dataset,
'historyTimingMethod':c.get('historyTimingPolicy',{}).get('method','strict_expiry'),
'publicationEstimatedSeconds':sum(r.get('publicationEstimatedSeconds',0) for r in good),
'originalFreshnessLimitsChanged':False,
'lastCapture':iso(records[-1]['capturedAt']) if records else None}
state='collecting'
if sum(r['coverage'] for r in good) >= c['minimumTrainingHours']*12:
state='model_ready' if active else 'training'
attempted=con.execute('SELECT MAX(trained_at) FROM planner_load_models WHERE plant=? AND dataset=?',(plant,dataset)).fetchone()[0]
if attempted is None or now-attempted>=cadence:
profiles=build_profiles(good)
candidate={'profiles':profiles,'trainedAt':now,'trainedThrough':max(r['start']+300 for r in good),
'trainingWindowFrom':good[0]['start'],'sourceDataset':dataset,'formula':c['formula'],
'methodVersion':'physical-profile-v1','validation':{'status':'bootstrap_insufficient_holdout'},
'historyTimingMethod':c.get('historyTimingPolicy',{}).get('method','strict_expiry'),
'timingPolicySha256':sha256(canonical(c.get('historyTimingPolicy')).encode()).hexdigest(),
'observationDatasetId':observation_dataset,
'measurementBoundaryVerified':False}
# Causal held-out validation: build validation profiles without the final day.
split=good[-1]['start']-86400
train=[r for r in good if r['start']+300<=split and r.get('availableNotBefore',r['start']+300)<=split]; test=[r for r in good if r['start']>=split]
if len(train)>=288 and len(test)>=240:
val={'profiles':build_profiles(train)}
errors={f:sum(abs(predict(val,f,r['start'])-r['loadW'])*r['coverage'] for r in test)/sum(r['coverage'] for r in test) for f in FAMILIES}
candidate['validation']={'status':'causal_holdout','holdoutFrom':split,'holdoutWindows':len(test),'loadMaeWByFamily':errors}
# Initial model is labelled bootstrap, never a production measurement proof.
# Existing model can be replaced only with held-out evidence and no aggregate regression.
promote=active is None
if active and candidate['validation']['status']=='causal_holdout':
past_model_eligible=active['trainedThrough']<=split and active['trainedAt']<=split
if past_model_eligible:
incumbent=sum(abs(predict(active,f,r['start'])-r['loadW'])*r['coverage'] for f in FAMILIES for r in test)
challenger=sum(abs(predict(val,f,r['start'])-r['loadW'])*r['coverage'] for f in FAMILIES for r in test)
promote=challenger<=incumbent
candidate['validation']['incumbentCompared']=True
else:
candidate['validation']['status']='holdout_overlaps_active_training'
ident=sha256(canonical(candidate).encode()).hexdigest()
with con:
con.execute('INSERT OR IGNORE INTO planner_load_models VALUES(?,?,?,?,?,?)',(plant,dataset,ident,now,candidate['trainedThrough'],canonical(candidate)))
if promote:
con.execute('INSERT INTO planner_model_current VALUES(?,?,?) ON CONFLICT(plant,dataset) DO UPDATE SET model_id=excluded.model_id',(plant,dataset,ident))
detail['candidateModelId']=ident;detail['candidatePromoted']=promote
state='model_ready' if promote or active else 'candidate_pending'
active=current_model(con,plant,dataset,now)
if active: detail.update({'modelId':active['modelId'],'trainedAt':iso(active['trainedAt']),'trainedThrough':iso(active['trainedThrough']),'validation':active['validation']})
with con:
con.execute('INSERT INTO planner_pipeline_state VALUES(?,?,?,?,?) ON CONFLICT(plant,dataset) DO UPDATE SET tick=excluded.tick,status=excluded.status,detail=excluded.detail',
(plant,dataset,tick,state,canonical(detail)))
def current_model(con,plant,dataset,at):
row=con.execute('SELECT m.model_id,m.value FROM planner_load_models m JOIN planner_model_current c ON m.plant=c.plant AND m.dataset=c.dataset AND m.model_id=c.model_id WHERE m.plant=? AND m.dataset=? AND m.trained_at<=?',(plant,dataset,at)).fetchone()
return {**json.loads(row[1]),'modelId':row[0]} if row else None
def apply_load_forecast(con,plant,dataset,forecast,decision):
model=current_model(con,plant,dataset,decision)
if not model:
raise ValueError('Corrected profile is collecting data; legacy household forecast is not silently reused')
c=configuration(con,plant,dataset)
last=con.execute('SELECT value FROM planner_observations WHERE plant=? AND dataset=? AND captured_at<=? AND received_at<=? ORDER BY captured_at DESC LIMIT 1',(plant,c.get('sourceDatasetId',dataset),decision,decision)).fetchone()
if not last: raise ValueError('No recent corrected observation')
last=json.loads(last[0])
if decision-last['capturedAt']>120: raise ValueError('Corrected measurements older than 120 seconds')
sdl_sources=[s for s in c['sources'] if s['role']=='sdl_request']
if len(sdl_sources)!=1: raise ValueError('Explicit SDL request channel needed for the labelled persistence scenario')
s=sdl_sources[0];r=last['raw'][s['key']]
if not r['valid'] or decision-r['sourceUpdatedAt']>s['maxAgeSeconds']:
raise ValueError('No current external SDL request for the persistence scenario')
sdl=r['value']*s['factorToW']
result=json.loads(canonical(forecast)); result['families']={}
result['observedAt']=iso(max(epoch(forecast['observedAt']),model['trainedAt'],last['capturedAt']))
for family,old in forecast['families'].items():
if family not in FAMILIES: continue
points=[]
for p in old['points']:
t=epoch(p['time'])
points.append({**p,'loadW':predict(model,family,t),'externalW':sdl})
result['families'][family]={'loadBasis':'base_load','trainedUntil':iso(model['trainedThrough']),
'points':points,'dataPipeline':{'datasetId':dataset,'modelId':model['modelId'],
'loadMethodVersion':model['methodVersion'],'loadVariant':family,
'loadModelTrainedAt':iso(model['trainedAt']),'pvForecastEventId':forecast.get('eventId'),
'measurementBasis':'configured_physical_estimate','measurementBoundaryVerified':False,
'externalPolicy':'last_sdl_request_persistence_estimate','externalObservedAt':iso(r['sourceUpdatedAt']),
'externalPowerW':sdl,'futureSdlPublished':False,'validation':model['validation'],
'historyTimingMethod':model.get('historyTimingMethod','strict_expiry'),
'observationDatasetId':model.get('observationDatasetId',dataset)}}
return result
def pipeline_status(con,plant):
out=[]
for row in con.execute('SELECT dataset,config FROM planner_data_sets WHERE plant=? ORDER BY dataset',(plant,)):
ds=row[0]; c=json.loads(row[1]); state=con.execute('SELECT status,detail FROM planner_pipeline_state WHERE plant=? AND dataset=?',(plant,ds)).fetchone()
count=con.execute('SELECT COUNT(*),MIN(captured_at),MAX(captured_at),MAX(received_at) FROM planner_observations WHERE plant=? AND dataset=?',(plant,c.get('sourceDatasetId',ds))).fetchone()
out.append({'datasetId':ds,'formula':c['formula'],'mappingSha256':c['mappingSha256'],'records':count[0],
'firstCapture':iso(count[1]) if count[1] else None,'lastCapture':iso(count[2]) if count[2] else None,
'lastReceived':iso(count[3]) if count[3] else None,
'status':state[0] if state else 'awaiting_measurements','detail':json.loads(state[1]) if state else {},
'minimumCoverage':c['minimumCoverage'],'maximumGapSeconds':c['maximumGapSeconds'],
'observationDatasetId':c.get('sourceDatasetId',ds),
'historyTimingMethod':c.get('historyTimingPolicy',{}).get('method','strict_expiry')})
return {'datasets':out,'liveEnabled':False,'legacyHistoryModified':False,'historyTimingVersion':1}