"""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()