514 lines
31 KiB
Python
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}
|