205 lines
11 KiB
Python
205 lines
11 KiB
Python
"""Synthetic publication gaps plus full source-reference/application integration."""
|
|
import copy
|
|
import json
|
|
import tempfile
|
|
import unittest
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from fastapi.testclient import TestClient
|
|
from netplan_v4 import measurement_pipeline as m
|
|
from netplan_v4.history_timing import validate_policy
|
|
from netplan_v4.store import PlannerStore
|
|
from netplan_v4.service import create_app, ingest, run_once
|
|
from test_measurement_pipeline import config, record, NOW, AID
|
|
from test_v4 import inputs, TOKEN
|
|
|
|
|
|
def timed_config(**kw):
|
|
return config(historyTimingPolicy={'version':1,'method':'equal_endpoint_v1',
|
|
'sources':{'pv':{'maxSpanSeconds':130,'evidenceId':'synthetic-publication-policy'}}}, **kw)
|
|
|
|
|
|
def projected(span=360, period=120, c=None):
|
|
c = c or timed_config()
|
|
result=[]
|
|
for delta in range(0,span+1,30):
|
|
r=record(NOW+delta,c)
|
|
r['raw']['pv']['sourceUpdatedAt']=NOW+(delta//period)*period
|
|
result.append(m.project(r,c,AID,NOW+span))
|
|
return result
|
|
|
|
|
|
class HistoricalTimingTest(unittest.TestCase):
|
|
def test_strict_policy_unchanged(self):
|
|
rows=projected(); result=m.reconstruct(rows,config())
|
|
self.assertEqual(result[0]['coveredSeconds'],180)
|
|
self.assertEqual(result[0]['publicationEstimatedSeconds'],0)
|
|
|
|
def test_equal_endpoint_bounded_tail_estimate(self):
|
|
w=m.reconstruct(projected(),timed_config())[0]
|
|
self.assertEqual(w['coveredSeconds'],300)
|
|
self.assertEqual(w['publicationEstimatedSeconds'],120)
|
|
self.assertEqual(w['publicationEstimatedBySourceSeconds'],{'pv':120})
|
|
self.assertEqual(w['loadW'],6000)
|
|
self.assertFalse(w['fullPhysicalIntervalMeasured'])
|
|
self.assertFalse(w['meterBoundaryVerified'])
|
|
|
|
def test_no_extension_at_open_end(self):
|
|
rows=projected(span=300)
|
|
for r in rows: r['raw']['pv']['sourceUpdatedAt']=NOW
|
|
w=m.reconstruct(rows,timed_config())[0]
|
|
self.assertEqual(w['coveredSeconds'],60)
|
|
|
|
def test_equal_zero_requires_new_original_timestamp(self):
|
|
rows=projected()
|
|
for r in rows:r['raw']['pv']['value']=0
|
|
self.assertEqual(m.reconstruct(rows,timed_config())[0]['coveredSeconds'],300)
|
|
for r in rows:r['raw']['pv']['sourceUpdatedAt']=NOW
|
|
self.assertEqual(m.reconstruct(rows,timed_config())[0]['coveredSeconds'],60)
|
|
|
|
def test_changed_endpoints_not_interpolated(self):
|
|
rows=projected()
|
|
for r in rows:r['raw']['pv']['value']+=r['raw']['pv']['sourceUpdatedAt']-NOW
|
|
self.assertEqual(m.reconstruct(rows,timed_config())[0]['publicationEstimatedSeconds'],0)
|
|
|
|
def test_no_floating_tolerance(self):
|
|
rows=projected()
|
|
for r in rows:r['raw']['pv']['value']+=1e-9*(r['raw']['pv']['sourceUpdatedAt']-NOW)
|
|
self.assertEqual(m.reconstruct(rows,timed_config())[0]['publicationEstimatedSeconds'],0)
|
|
|
|
def test_long_gap_not_recovered_by_equal_zero(self):
|
|
rows=projected(span=720,period=660)
|
|
for r in rows:r['raw']['pv']['value']=0
|
|
out=m.reconstruct(rows,timed_config())
|
|
self.assertEqual(sum(w['publicationEstimatedSeconds'] for w in out),0)
|
|
self.assertEqual(out[0]['coveredSeconds'],60)
|
|
|
|
def test_collector_gap_not_bridged(self):
|
|
rows=projected(); rows=[r for r in rows if r['capturedAt']!=NOW+90]
|
|
w=m.reconstruct(rows,timed_config())[0]
|
|
self.assertLess(w['coveredSeconds'],300)
|
|
self.assertEqual(w['publicationEstimatedSeconds'],60)
|
|
|
|
def test_invalid_source_blocks_equality_inference(self):
|
|
rows=projected();rows[3]['raw']['pv']['valid']=False
|
|
w=m.reconstruct(rows,timed_config())[0]
|
|
self.assertEqual(w['publicationEstimatedSeconds'],60)
|
|
self.assertLess(w['coverage'],1)
|
|
|
|
def test_conflicting_same_timestamp_blocks_equality(self):
|
|
rows=projected();rows[1]['raw']['pv']['value']+=10
|
|
w=m.reconstruct(rows,timed_config())[0]
|
|
self.assertEqual(w['publicationEstimatedSeconds'],60)
|
|
|
|
def test_invalid_other_physical_input_still_vetoes(self):
|
|
rows=projected();rows[3]['raw']['battery']['valid']=False
|
|
w=m.reconstruct(rows,timed_config())[0]
|
|
self.assertLess(w['coverage'],1)
|
|
|
|
def test_future_receipt_not_available_early(self):
|
|
rows=projected()
|
|
for r in rows:r['_receivedAt']=NOW+900
|
|
w=m.reconstruct(rows,timed_config())[0]
|
|
self.assertEqual(w['availableNotBefore'],NOW+900)
|
|
|
|
def test_unrelated_virtual_failure_does_not_veto(self):
|
|
rows=projected()
|
|
for r in rows:r['raw']['sdl']['valid']=False
|
|
self.assertEqual(m.reconstruct(rows,timed_config())[0]['coverage'],1)
|
|
|
|
def test_no_input_mutation_or_source_age_change(self):
|
|
rows=projected();c=timed_config();before=copy.deepcopy((rows,c))
|
|
m.reconstruct(rows,c)
|
|
self.assertEqual((rows,c),before)
|
|
self.assertEqual(c['sources'][1]['maxAgeSeconds'],60)
|
|
|
|
def test_policy_requires_trusted_explicit_bounds(self):
|
|
for change in ({'version':True},{'method':'always_fill'},{'sources':{'pv':{'maxSpanSeconds':600,'evidenceId':'test-policy'}}},
|
|
{'sources':{'grid':{'maxSpanSeconds':90,'evidenceId':'test-policy'}}},
|
|
{'sources':{'sdl':{'maxSpanSeconds':90,'evidenceId':'test-policy'}}}):
|
|
c=timed_config();c['historyTimingPolicy'].update(change)
|
|
with self.subTest(change=change),self.assertRaises(ValueError):m.validate_config(c)
|
|
|
|
def test_profile_values_stay_estimates(self):
|
|
w=m.reconstruct(projected(),timed_config())[0]
|
|
self.assertTrue(w['estimated'])
|
|
self.assertNotIn('controlEnabled',w)
|
|
self.assertNotIn('accountingEvidenceId',w)
|
|
|
|
|
|
class HistoryReferenceIntegrationTest(unittest.TestCase):
|
|
def setUp(self):
|
|
self.tmp=tempfile.TemporaryDirectory();self.store=PlannerStore(str(Path(self.tmp.name)/'test.sqlite'))
|
|
self.c=config();self.new=timed_config(datasetId='physical-v2',sourceDatasetId='physical-v1')
|
|
m.register_dataset(self.store.con,AID,self.c,NOW)
|
|
def tearDown(self):self.store.close();self.tmp.cleanup()
|
|
def register(self):m.register_dataset(self.store.con,AID,self.new,NOW)
|
|
def fill(self):
|
|
raw=[]
|
|
for t in range(NOW-7200,NOW+1,30):
|
|
r=record(t);r['raw']['pv']['sourceUpdatedAt']=t-(t-(NOW-7200))%120;raw.append(r)
|
|
for i in range(0,len(raw),120):m.ingest_batch(self.store.con,AID,{'version':1,'datasetId':'physical-v1','records':raw[i:i+120]},NOW)
|
|
return len(raw)
|
|
def test_reference_shares_original_immutable_observations(self):
|
|
n=self.fill();self.register()
|
|
self.assertEqual(m.pipeline_status(self.store.con,AID)['datasets'][1]['records'],n)
|
|
self.assertEqual(self.store.con.execute('SELECT count(*) FROM planner_observations').fetchone()[0],n)
|
|
self.assertEqual(m.configuration(self.store.con,AID,'physical-v1'),self.c)
|
|
def test_no_append_into_derived_dataset(self):
|
|
self.register()
|
|
with self.assertRaises(ValueError):m.ingest_batch(self.store.con,AID,{'version':1,'datasetId':'physical-v2','records':[record(NOW)]},NOW)
|
|
def test_source_mapping_cannot_change_in_reference(self):
|
|
for key,value in [('mappingSha256','c'*64),('formula','solar_terminal_v1'),('minimumCoverage',.9)]:
|
|
c=copy.deepcopy(self.new);c[key]=value
|
|
with self.subTest(key=key),self.assertRaises(ValueError):m.register_dataset(self.store.con,AID,c,NOW)
|
|
def test_no_reference_to_other_plant(self):
|
|
with self.assertRaises(ValueError):m.register_dataset(self.store.con,'00000000-0000-4000-8000-000000000099',self.new,NOW)
|
|
def test_no_cycles_or_reference_chains(self):
|
|
self.register(); c={**self.new,'datasetId':'physical-v3','sourceDatasetId':'physical-v2'}
|
|
with self.assertRaises(ValueError):m.register_dataset(self.store.con,AID,c,NOW)
|
|
def test_training_model_is_versioned_and_policy_annotated(self):
|
|
self.fill();self.register();settings=self.store.settings(AID)
|
|
m.advance(self.store.con,AID,'physical-v1',settings,NOW+60)
|
|
m.advance(self.store.con,AID,'physical-v2',settings,NOW+60)
|
|
self.assertIsNone(m.current_model(self.store.con,AID,'physical-v1',NOW+60))
|
|
model=m.current_model(self.store.con,AID,'physical-v2',NOW+60)
|
|
self.assertIsNotNone(model);self.assertEqual(model['historyTimingMethod'],'equal_endpoint_v1')
|
|
self.assertFalse(model['measurementBoundaryVerified'])
|
|
state=m.pipeline_status(self.store.con,AID)['datasets'][1]['detail']
|
|
self.assertGreater(state['publicationEstimatedSeconds'],0)
|
|
self.assertFalse(state['originalFreshnessLimitsChanged'])
|
|
def test_no_future_receipt_leaks_into_model(self):
|
|
self.fill();self.register()
|
|
m.advance(self.store.con,AID,'physical-v2',self.store.settings(AID),NOW-60)
|
|
self.assertIsNone(m.current_model(self.store.con,AID,'physical-v2',NOW-60))
|
|
def test_existing_sdl_freshness_not_relaxed_by_history(self):
|
|
self.fill();self.register();m.advance(self.store.con,AID,'physical-v2',self.store.settings(AID),NOW+30)
|
|
op,fc,tar=inputs(datetime.fromtimestamp(NOW+180,timezone.utc))
|
|
with self.assertRaises(ValueError):m.apply_load_forecast(self.store.con,AID,'physical-v2',fc,NOW+180)
|
|
def test_full_application_plan_uses_reference_model(self):
|
|
self.fill();self.register();settings=self.store.settings(AID)
|
|
m.advance(self.store.con,AID,'physical-v2',settings,NOW)
|
|
op,fc,tar=inputs(datetime.fromtimestamp(NOW,timezone.utc))
|
|
for kind,val in zip(('operation','forecast','tariffs'),(op,fc,tar)):ingest(self.store,AID,kind,val,datetime.fromtimestamp(NOW,timezone.utc))
|
|
self.store.save_settings(AID,{'forecastSource':'corrected_profile','measurementDataset':'physical-v2'},0,datetime.fromtimestamp(NOW,timezone.utc))
|
|
out=run_once(self.store,datetime.fromtimestamp(NOW,timezone.utc))
|
|
self.assertTrue(out['executable'],out);self.assertFalse(out['liveEnabled'])
|
|
self.assertEqual(out['inputQuality']['dataPipeline']['historyTimingMethod'],'equal_endpoint_v1')
|
|
def test_backfilled_or_late_confirmed_training_is_not_causal_holdout(self):
|
|
# Completed synthetic windows learned only NOW cannot validate a model at a past split.
|
|
with self.store.con:
|
|
for t in range(NOW-3*86400,NOW,300):
|
|
window={'start':t,'profileUsable':True,'coverage':1.,'loadW':5000.,'availableNotBefore':NOW}
|
|
self.store.con.execute('INSERT INTO planner_load_windows VALUES(?,?,?,?,?,?)',(AID,'physical-v1',t,NOW,1.,m.canonical(window)))
|
|
m.advance(self.store.con,AID,'physical-v1',self.store.settings(AID),NOW)
|
|
model=m.current_model(self.store.con,AID,'physical-v1',NOW)
|
|
self.assertEqual(model['validation']['status'],'bootstrap_insufficient_holdout')
|
|
|
|
def test_v1_raw_ingest_works_after_reference_added(self):
|
|
self.register()
|
|
r=m.ingest_batch(self.store.con,AID,{'version':1,'datasetId':'physical-v1','records':[record(NOW)]},NOW)
|
|
self.assertEqual(r['stored'],1)
|
|
|
|
|
|
if __name__=='__main__':unittest.main()
|