fix(history): model sparse publication without relaxing live freshness

This commit is contained in:
ENELIX Agent
2026-10-03 08:04:28 +00:00
parent 6d9aba74b0
commit 5890093b80
15 changed files with 1787 additions and 22 deletions
@@ -0,0 +1,58 @@
"""Explicit, retrospective equal-endpoint estimates for sparse state publication.
Never changes a source timestamp, real-time freshness limit, meter evidence or
actuator lease. A repeated value after a bounded gap supports a MODEL estimate,
not proof that the physical signal was constant between the observations.
"""
from math import isfinite
import re
def validate_policy(policy, sources):
if policy is None:
return None
if (not isinstance(policy, dict) or set(policy) != {'version', 'method', 'sources'}
or type(policy['version']) is not int or policy['version'] != 1
or policy['method'] != 'equal_endpoint_v1'
or not isinstance(policy['sources'], dict) or not policy['sources']):
raise ValueError('Explicit historical timing policy required')
defined = {s['key']: s for s in sources}
for key, rule in policy['sources'].items():
if key not in defined or defined[key]['role'] not in ('pv', 'physical_storage', 'solar_scale'):
raise ValueError('Timing estimates limited to explicit power/scale sources')
if not isinstance(rule, dict) or set(rule) != {'maxSpanSeconds', 'evidenceId'}:
raise ValueError('Timing bound and evidence reference required')
if type(rule['maxSpanSeconds']) is not int or not defined[key]['maxAgeSeconds'] <= rule['maxSpanSeconds'] <= 180:
raise ValueError('Historical endpoint span outside reviewed bounds')
if not isinstance(rule['evidenceId'], str) or not re.fullmatch(r'[A-Za-z0-9_-]{8,100}', rule['evidenceId']):
raise ValueError('Invalid publication evidence reference')
return policy
def endpoint_bridges(series, first_observed, blocks, capture_gaps, policy, source_specs):
"""Index strictly consecutive, equal original timestamps; never invent endpoints.
Only the stale tail of a normal hold is filled. A missing/invalid capture or a
source conflict forbids bridging. The right endpoint must have been observed;
its availability is retained by the caller for causal training and replay.
"""
result = {key: {} for key in series}
if policy is None:
return result
for key, rule in policy['sources'].items():
if key not in series:
continue # e.g. an unused DC channel in an AC-terminal formula
values = series[key]
times = sorted(values)
for left, right in zip(times, times[1:]):
value, next_value = values[left], values[right]
if (value is None or next_value is None or type(value) not in (int, float)
or not isfinite(value) or value != next_value):
continue
if not source_specs[key]['maxAgeSeconds'] < right-left <= rule['maxSpanSeconds']:
continue
if any(start < right and end > left for start, end in (*blocks[key], *capture_gaps)):
continue
result[key][left] = {'end': right, 'availableAt': first_observed[key][right],
'spanSeconds': right-left}
return result
@@ -13,6 +13,7 @@ from math import isfinite
from statistics import median
from zoneinfo import ZoneInfo
import json
from .history_timing import validate_policy, endpoint_bridges
UTC = timezone.utc
LOCAL = ZoneInfo('Europe/Zurich')
@@ -79,7 +80,8 @@ def validate_config(c):
fields = {'datasetId', 'mappingSha256', 'inventorySha256', 'sources',
'formula', 'solarReference', 'minimumCoverage', 'maximumGapSeconds',
'minimumTrainingHours', 'historyDays'}
if not isinstance(c, dict) or set(c) != fields:
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):
@@ -122,6 +124,10 @@ def validate_config(c):
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
@@ -129,6 +135,13 @@ def validate_config(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:
@@ -178,6 +191,8 @@ 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')
@@ -244,8 +259,10 @@ uncovered portions remain quantified and are never filled with zero.
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'])
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:]):
@@ -266,6 +283,7 @@ uncovered portions remain quantified and are never filled with zero.
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])
@@ -273,24 +291,38 @@ uncovered portions remain quantified and are never filled with zero.
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})
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
if t is None or a >= t+s['maxAgeSeconds'] or series[k][t] is None or any(x <= a < y for x,y in blocks[k]):
usable = False
else: vals[k] = series[k][t]
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)
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.
@@ -300,7 +332,12 @@ uncovered portions remain quantified and are never filled with zero.
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']})
'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
@@ -341,13 +378,14 @@ def advance(con, plant, dataset, settings, now):
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
fetched=con.execute('SELECT value FROM planner_observations WHERE plant=? AND dataset=? AND captured_at>=? AND captured_at<=? AND received_at<=? ORDER BY captured_at',
(plant,dataset,now-172800-300,now,now)).fetchall()
records=[json.loads(r[0]) for r in fetched]
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: continue
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))]
@@ -356,7 +394,12 @@ def advance(con, plant, dataset, settings, 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}
'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'
@@ -366,10 +409,13 @@ def advance(con, plant, dataset, settings, now):
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]; test=[r for r in good if r['start']>=split]
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}
@@ -378,7 +424,7 @@ def advance(con, plant, dataset, settings, now):
# 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
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)
@@ -409,9 +455,10 @@ 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')
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,dataset,decision,decision)).fetchone()
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]); c=configuration(con,plant,dataset)
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')
@@ -433,7 +480,9 @@ def apply_load_forecast(con,plant,dataset,forecast,decision):
'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']}}
'externalPowerW':sdl,'futureSdlPublished':False,'validation':model['validation'],
'historyTimingMethod':model.get('historyTimingMethod','strict_expiry'),
'observationDatasetId':model.get('observationDatasetId',dataset)}}
return result
@@ -441,9 +490,11 @@ 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) FROM planner_observations WHERE plant=? AND dataset=?',(plant,ds)).fetchone()
count=con.execute('SELECT COUNT(*),MIN(captured_at),MAX(captured_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,
'status':state[0] if state else 'awaiting_measurements','detail':json.loads(state[1]) if state else {},
'minimumCoverage':c['minimumCoverage'],'maximumGapSeconds':c['maximumGapSeconds']})
return {'datasets':out,'liveEnabled':False,'legacyHistoryModified':False}
'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}