diff --git a/docs/netplan-v4-history-timing.md b/docs/netplan-v4-history-timing.md new file mode 100644 index 0000000..b96f8d4 --- /dev/null +++ b/docs/netplan-v4-history-timing.md @@ -0,0 +1,90 @@ +# Historical publication timing (2026-10-03) + +## Implemented scope + +`services/netplan-v4/netplan_v4/history_timing.py` implements optional, bounded +retrospective estimates between consecutive identical source values at DIFFERENT +original publication timestamps. It is used by the existing measurement pipeline, +model training and forecast source substitution. It is NOT a real-time estimator, +a measurement certificate, an actuator authority or a repair of device communications. + +No source `maxAgeSeconds` is changed. When a regular hold expires, only its tail +until the next qualifying original publication may be estimated. Strict processing +remains the default. The following forbid the estimate: unequal endpoints, +nonfinite/missing observations, conflicting same-timestamp values, collector gaps, +invalid source intervals or a span beyond the explicit source-specific bound. +There is no extrapolation past the last observation and no fabricated zero. + +Each window reports `publicationEstimatedSeconds`, per-source estimated portions, +uncovered-source seconds and `availableNotBefore`. Training/validation must not use +a confirmation or receipt before it was actually available. Backfilled windows +cannot establish a supposedly causal holdout at a date before their availability. +All physical windows remain configured estimates; even full support is not proof +that the physical waveform stayed constant between the two endpoint readings. + +## Evidence and model assumptions + +Read-only source configuration projection on the test server, snapshot +2026-10-03T07:47:24Z: GoodWe 1 instance 19742 (ModBus Device) has Poller 60000 ms; +GoodWe 2 instance 57658 has Poller 5000 ms; SolarEdge scale instance 48996 +(ModBus Address) has Poller 20000 ms. No polling setting was changed. +Projection program: /srv/agent/netplan-v4-application-build/inspect_source_update_policy.py. + +Official Symcon documentation accessed on 2026-10-03: +https://www.symcon.de/de/service/dokumentation/modulreferenz/geraete/modbus-rtu-tcp/vorlagen/ +It documents value publication for ModBus Device on changes or when the variable is +older than 60 seconds, including ordinary and virtual addresses. VariableUpdated +is therefore not automatically the last successful device poll. The documentation +is not a proof that a particular installed device link was healthy. + +Explicit Lihrenmoos HISTORICAL estimate bounds: +- GoodWe 1 PV/physical storage: equal original publications at most 130 s apart + (60 s publication suppression + configured 60 s poll + 10 s scheduling allowance). +- GoodWe 2 PV/physical storage: at most 75 s (60 + 5 + 10). +- SolarEdge scale: equal scale publications at most 90 s apart. This is a modelling + bound based on the configured 20 s polling and observed sparse scale publication, + not a manufacturer guarantee and not proof of the same internal implementation + as ModBus Device. Raw AC power remains under the original strict limit. +The 10 s margins are explicit scheduling assumptions, NOT measured timing guarantees. +Gaps of many minutes (e.g. the previously observed GoodWe gaps above 600 s) are NOT +bridged. Nonzero changes are never interpolated by this policy. + +## Versioned dataset without restarting collection + +`lihrenmoos-physical-published-v2` references immutable observations in +`lihrenmoos-physical-v1`. No observations are copied, deleted, relabelled or resent. +Registration requires the same installation and identical original mapping, +formula, coverage and training configuration. Cycles/reference chains and device +uploads directly into a derived dataset are rejected. Original ingestion continues +unchanged in the Manager. The 24 usable-hour minimum stays in force. + +Models/windows are stored separately for v1/v2 and record their timing policy. An +existing selected model/family is not automatically switched by this release. +The current SDL request still uses the strict current-state requirements; no +historical bridge grants a future SDL schedule or a live battery permission. + +## Tests and release + +254 isolated Python tests passed (219 previous + 35 new including simulated +installation/recovery). Includes the actual service pipeline from referenced +observations to model and optimizer using synthetic records. No new target-container +or field-history evaluation has been run with this release yet; do not infer an +improved real coverage percentage or production acceptance from unit tests. +Evidence: /home/agent/services/qa/history-timing-20261003/ALL_TEST_RESULTS.txt. + +Prepared command on enelix-services (root required for Docker): + + python3 /home/agent/services/netplan-v4-shadow/commissioning/deploy_history_timing.py \ + --plant e3a08f9e-af12-4695-99bd-8b51c0520021 --apply + +Default without --apply validates staged source only. The apply path backs up source +files, builds and runs the exact target tests, backs up the V4 SQLite database, and +recreates ONLY the V4 service. A failed test restores source without a service +restart; failed post-deploy validation attempts the previous image and preserves +additive data. Concurrently modified files are not silently overwritten. + +No portal/forecast/tariff/Symcon restart, no new sensor requests, no change to +161.44 kWh / 39 kW, SOC accounting, old raw samplers or actuator permissions. +The new dataset is registered but not selected for the active forecast. A root +execution and successful report are still required to install it. Previous release +manifests remain historical records; do not bypass them to reapply an older build. diff --git a/services/netplan-v4/Dockerfile b/services/netplan-v4/Dockerfile index b98928f..ceff456 100644 --- a/services/netplan-v4/Dockerfile +++ b/services/netplan-v4/Dockerfile @@ -10,6 +10,7 @@ COPY gui ./gui COPY release_preflight.py forecast_acceptance.py deploy_integrated_shadow.py approved_previous_assets.json ./ COPY acceptance ./acceptance COPY commissioning/deploy_application.py ./commissioning/deploy_application.py +COPY commissioning/deploy_history_timing.py ./commissioning/deploy_history_timing.py # Host sources can be 0600/0700. COPY makes them root-owned. # Normalize only packaged application code; never change host secrets or sockets. RUN find /app -type d -exec chmod 0755 {} + \ diff --git a/services/netplan-v4/SOURCE_MANIFEST.json b/services/netplan-v4/SOURCE_MANIFEST.json index f32a7a7..34ed6eb 100644 --- a/services/netplan-v4/SOURCE_MANIFEST.json +++ b/services/netplan-v4/SOURCE_MANIFEST.json @@ -1,7 +1,7 @@ { "sourceHashes": { ".dockerignore": "fab6861d98f34e54646ae966b237e792fc0a0b4f95f1df3b62a1f062cbb8f790", - "Dockerfile": "655c600a0364e91d47bfc80faaf27e26362bfc2683c8e57b653d913133865713", + "Dockerfile": "f61404d65d9c031d0ddd5db9d2fa52279e39b1bd6015249a375b9167c744ae90", "acceptance/Dockerfile.php": "715a2d232d7f901bd6ca1f2e453b7b7f48fc1d7f49bff8944a7d277a8dc30062", "acceptance/check_forecast.py": "ccee58ec1078767a580f15f895506eceeb15e13b54e8eb9b4905cc03ba14ccd5", "acceptance/forecast-src/SOURCE_MANIFEST.json": "8103775396c58d82921e2e2a1513ee185ae2431200763717544ff9b1da181fe5", @@ -58,7 +58,7 @@ "netplan_v4/controlled_trial.py": "4d4bd8ed97050b3bf5872c839e4bcc2fc9dcb9f2e0a8c66a4611bbd50bf9e1c4", "netplan_v4/domain.py": "03aadfb6d69f349041c872e105b3ae31e880b507a630f82bdd774948a618b39b", "netplan_v4/forecast_quality.py": "b656868538cb822dc1dec97c7726bcc47d78744a1db3ce02a11d259ba9c2a840", - "netplan_v4/measurement_pipeline.py": "caab05a69ac28c086f06b10b20460330b4c3721a6467417542d7666ce82ac743", + "netplan_v4/measurement_pipeline.py": "17e59c8e8d7403f107e95620260c7b71d9f51e8165c0c29944bb931c55c4001d", "netplan_v4/meter_runtime.py": "43275072a212351fa35d34ed2f4932a259112d516d594fad19a7e2b3fd44c547", "netplan_v4/metering.py": "9ed4d3747de76f646d7be603cbdf37a3fc971109605705017c6644c0baa80987", "netplan_v4/optimizer.py": "bf1ad3dcadf10c84e76f6758525fc1309b9f665853c660a1107c1347f1ce9063", @@ -87,7 +87,13 @@ "tests/test_peak_release.py": "223f4374111cfab99ef352a71a7b91a7f71abc4ed126f2c180b138d1bd390d3d", "tests/test_receiver_contract.py": "21daccb358fa62a983d2292d5de2b9b3c7300a08e314e850749316a24d3d254e", "tests/test_release_preflight.py": "64104923bea0890e5010471466de87c289a97bca7dd29d56734d0009e79f24d8", - "tests/test_v4.py": "c35db22a864aa0afc4d0abc357ae854f038e86de8e3001760c0e63ce036a7763" + "tests/test_v4.py": "c35db22a864aa0afc4d0abc357ae854f038e86de8e3001760c0e63ce036a7763", + "netplan_v4/history_timing.py": "d2992e15d8d1a4f79d068405eefba4d4b37a4326afb09064cfd7c098822ec63f", + "tests/test_history_timing.py": "7a11d52c779cd25f44bc2bb56a758b862154e2c8df6987046e18eace46ed51a2", + "tests/test_history_timing_deployment.py": "7508d24b70c1fa4c33e21fbad4e99d5bff23e262d91222df5b0405e937c96d2e", + "commissioning/deploy_history_timing.py": "7769a473927bb10acf1ba1aef902e6b14fdf26b548d58056b557e6e4c33f62eb", + "commissioning/history-timing-source/dataset.json": "671216e1cc04148f4f553103077bcb4500bf0120cb74780703bfcf26cd1121d3", + "commissioning/history-timing-source/RELEASE.json": "14cb690a357c40f698c2e9b4190933318396b9c08e15cc4ad9f9f4439d6054ca" }, "packagedAt": "2026-10-02T21:00:35.231591+00:00", "runtimeChanged": false diff --git a/services/netplan-v4/commissioning/deploy_history_timing.py b/services/netplan-v4/commissioning/deploy_history_timing.py new file mode 100644 index 0000000..c4c5914 --- /dev/null +++ b/services/netplan-v4/commissioning/deploy_history_timing.py @@ -0,0 +1,146 @@ +"""Install the tested historical-publication correction in the V4 service only. + +Default is source validation. --apply requires root. No Symcon changes, polling, +portal restart, old-data rewrite, forecast-source switch or actuator permission. +""" +from datetime import datetime, timezone +from pathlib import Path +import argparse +import hashlib +import json +import os +import subprocess + +ROOT = Path(__file__).resolve().parents[1] +PACKAGE = Path(__file__).with_name('history-timing-source') +TARGETS = ('Dockerfile', 'netplan_v4/history_timing.py', 'netplan_v4/measurement_pipeline.py', + 'tests/test_history_timing.py', 'tests/test_history_timing_deployment.py') + + +def digest(path): + if path.is_symlink(): + raise ValueError('Source symlink refused') + return hashlib.sha256(path.read_bytes()).hexdigest() if path.is_file() else None + + +def verify(): + manifest = json.loads((PACKAGE/'RELEASE.json').read_text()) + if manifest.get('scope') != 'historical_publication_only' or tuple(manifest.get('files', {})) != TARGETS: + raise ValueError('Unexpected timing release scope') + for name, versions in manifest['files'].items(): + if digest(PACKAGE/'source'/name) != versions['after']: + raise ValueError('Timing candidate changed: '+name) + if digest(ROOT/name) not in (versions['before'], versions['after']): + raise ValueError('Concurrent application change: '+name) + if digest(PACKAGE/'dataset.json') != manifest['datasetSha256']: + raise ValueError('Dataset candidate changed') + if digest(ROOT/'commissioning/deploy_application.py') != manifest['existingDeploymentHelperSha256']: + raise ValueError('Existing deployment helper changed') + if digest(Path(__file__).resolve()) != manifest['installerSha256']: + raise ValueError('Timing installer changed') + return manifest + + +BOOTSTRAP = r'''import json,os,sys,time,urllib.request +payload=json.load(sys.stdin) +base='http://127.0.0.1:9100/internal/v2/prognosis/'+payload['plant']+'/planner' +headers={'X-Enelix-Service-Token':os.environ['PROGNOSIS_SERVICE_TOKEN'],'Content-Type':'application/json'} +def get(): + return json.load(urllib.request.urlopen(urllib.request.Request(base,headers=headers),timeout=10)) +before=get() +assert before['dataPipeline'].get('historyTimingVersion')==1 +assert before['liveEnabled'] is False +req=urllib.request.Request(base+'/datasets/'+payload['dataset']['datasetId'], + data=json.dumps(payload['dataset']).encode(),headers=headers,method='PUT') +receipt=json.load(urllib.request.urlopen(req,timeout=10)) +after=get() +assert before['settings']==after['settings'] and after['liveEnabled'] is False +view=next(d for d in after['dataPipeline']['datasets'] if d['datasetId']==payload['dataset']['datasetId']) +print(json.dumps({'dataset':receipt,'originalObservationDataset':view['observationDatasetId'], + 'availableRecords':view['records'],'pipelineStatus':view['status'], + 'settingsUnchanged':True,'forecastSource':after['settings']['forecastSource'], + 'liveEnabled':False})) +''' + + +def run(plant, apply=False): + manifest = verify() + if plant != manifest['installationId']: + raise ValueError('Use the reviewed installation') + if not apply: + print('HISTORY TIMING SOURCE CHECK PASSED; no runtime changes.') + return + if os.geteuid() != 0: + raise ValueError('Run as root; do not change Docker permissions') + # Reuse the reviewed atomic-write and consistent-backup implementation. + from deploy_application import atomic, backup_database + folder = ROOT/'history-timing-releases'/datetime.now(timezone.utc).strftime('%Y%m%dT%H%M%S.%fZ') + folder.mkdir(parents=True, mode=0o700) + result={'scope':'historical_publication_only','startedAt':datetime.now(timezone.utc).isoformat(), + 'steps':[],'liveEnabled':False,'symconChanged':False,'pollingChanged':False, + 'sourceFreshnessLimitsChanged':False,'forecastSelectionChanged':False} + compose=['docker','compose','-f',str(ROOT/'compose.yaml')] + env={**os.environ,'NETPLAN_V4_PLANTS':plant} + def cmd(args, timeout=180, capture=False, data=None): + return subprocess.run(args,cwd=ROOT,env=env,check=True,timeout=timeout,text=True,input=data, + stdout=subprocess.PIPE if capture else None,stderr=subprocess.PIPE if capture else None) + changed={}; previous_image=None; replaced=False + try: + ids=cmd(compose+['ps','-q','netplan-v4'],capture=True).stdout.split() + if len(ids)!=1: raise ValueError('Expected one running V4 container') + previous_image=cmd(['docker','inspect','--format','{{.Image}}',ids[0]],capture=True).stdout.strip() + result['previousImage']=previous_image + for name, versions in manifest['files'].items(): + target=ROOT/name; content=(PACKAGE/'source'/name).read_bytes() + before=target.read_bytes() if target.is_file() else None + if digest(target)==versions['after']: continue + if digest(target)!=versions['before']: raise ValueError('Concurrent source change') + backup=folder/'source-before'/name; backup.parent.mkdir(parents=True,exist_ok=True) + if before is not None: backup.write_bytes(before) + atomic(target,content); changed[name]=(before,content) + verify() + cmd(compose+['build','netplan-v4'],timeout=900) + cmd(compose+['run','--rm','--no-deps','--entrypoint','python','netplan-v4','/app/run_tests.py'],timeout=240) + result['steps'].append('target_tests_passed'); verify() + backup_database(ROOT/'data/netplan-v4.sqlite',folder/'before.sqlite') + result['steps'].append('consistent_database_backup') + rollback=folder/'rollback.yaml' + rollback.write_text('services:\n netplan-v4:\n image: '+previous_image+'\n') + replaced=True + cmd(compose+['up','-d','--no-deps','--no-build','--wait','netplan-v4']) + payload=json.dumps({'plant':plant,'dataset':json.loads((PACKAGE/'dataset.json').read_text())}) + response=cmd(compose+['exec','-T','netplan-v4','python','-c',BOOTSTRAP],capture=True,data=payload) + result['application']=json.loads(response.stdout) + result['status']='historical_timing_installed_no_actuation' + result['steps'].append('versioned_dataset_registered_original_observations_retained') + except Exception as exc: + result['status']='needs_review';result['errorType']=type(exc).__name__ + conflicts=[] + for name,(before,after) in changed.items(): + path=ROOT/name + if digest(path)!=hashlib.sha256(after).hexdigest(): conflicts.append(name);continue + if before is None:path.unlink() + else:atomic(path,before) + result['concurrentFilesNotOverwritten']=conflicts + if replaced and previous_image: + try: + cmd(compose+['-f',str(folder/'rollback.yaml'),'up','-d','--no-deps','--no-build','--pull','never','--wait','netplan-v4']) + result['rollback']='previous_image_restored_additive_dataset_retained' + except Exception: result['rollback']='manual_review_required' + raise + finally: + result['finishedAt']=datetime.now(timezone.utc).isoformat() + report=folder/'REPORT.json';report.write_text(json.dumps(result,indent=2)+'\n') + uid,gid=ROOT.stat().st_uid,ROOT.stat().st_gid + for path in (folder.parent,folder,report):os.chown(path,uid,gid) + report.chmod(0o640) + print('HISTORY TIMING REPORT:',report) + print('Only V4 updated. Original data, manager, polling, forecast selection and actuator permissions unchanged.') + + +if __name__=='__main__': + parser=argparse.ArgumentParser(description=__doc__) + parser.add_argument('--plant',required=True);parser.add_argument('--apply',action='store_true') + args=parser.parse_args() + try:run(args.plant,args.apply) + except Exception as exc:raise SystemExit('Stopped: '+type(exc).__name__+'. See HISTORY TIMING REPORT; do not enable actuators.') diff --git a/services/netplan-v4/commissioning/history-timing-source/RELEASE.json b/services/netplan-v4/commissioning/history-timing-source/RELEASE.json new file mode 100644 index 0000000..65f5a18 --- /dev/null +++ b/services/netplan-v4/commissioning/history-timing-source/RELEASE.json @@ -0,0 +1,29 @@ +{ + "scope": "historical_publication_only", + "installationId": "e3a08f9e-af12-4695-99bd-8b51c0520021", + "files": { + "Dockerfile": { + "before": "655c600a0364e91d47bfc80faaf27e26362bfc2683c8e57b653d913133865713", + "after": "f61404d65d9c031d0ddd5db9d2fa52279e39b1bd6015249a375b9167c744ae90" + }, + "netplan_v4/history_timing.py": { + "before": null, + "after": "d2992e15d8d1a4f79d068405eefba4d4b37a4326afb09064cfd7c098822ec63f" + }, + "netplan_v4/measurement_pipeline.py": { + "before": "caab05a69ac28c086f06b10b20460330b4c3721a6467417542d7666ce82ac743", + "after": "17e59c8e8d7403f107e95620260c7b71d9f51e8165c0c29944bb931c55c4001d" + }, + "tests/test_history_timing.py": { + "before": null, + "after": "7a11d52c779cd25f44bc2bb56a758b862154e2c8df6987046e18eace46ed51a2" + }, + "tests/test_history_timing_deployment.py": { + "before": null, + "after": "7508d24b70c1fa4c33e21fbad4e99d5bff23e262d91222df5b0405e937c96d2e" + } + }, + "datasetSha256": "671216e1cc04148f4f553103077bcb4500bf0120cb74780703bfcf26cd1121d3", + "installerSha256": "7769a473927bb10acf1ba1aef902e6b14fdf26b548d58056b557e6e4c33f62eb", + "existingDeploymentHelperSha256": "8fabbbbf41e00be677a6f109690758035bbe098f9161d9c44865b28ab7f1bd4e" +} diff --git a/services/netplan-v4/commissioning/history-timing-source/dataset.json b/services/netplan-v4/commissioning/history-timing-source/dataset.json new file mode 100644 index 0000000..bf153a3 --- /dev/null +++ b/services/netplan-v4/commissioning/history-timing-source/dataset.json @@ -0,0 +1,206 @@ +{ + "datasetId": "lihrenmoos-physical-published-v2", + "mappingSha256": "517d1907631c7f152096e21759bb3452cdd0b4b1b73f91f1c7cd3bd62d8aaa8b", + "inventorySha256": "0a9dde190d81bc263da222ad8bfb6268295ea03aec239f6c1fcd249b3d34c1bc", + "sources": [ + { + "key": "grid", + "variableId": 40348, + "factorToW": 1000, + "maxAgeSeconds": 60, + "role": "grid" + }, + { + "key": "pv_goodwe1", + "variableId": 48459, + "factorToW": 1, + "maxAgeSeconds": 60, + "role": "pv" + }, + { + "key": "pv_goodwe2", + "variableId": 53802, + "factorToW": 1, + "maxAgeSeconds": 60, + "role": "pv" + }, + { + "key": "pv_solaredge", + "variableId": 20335, + "factorToW": 1, + "maxAgeSeconds": 60, + "role": "pv" + }, + { + "key": "physical_goodwe1", + "variableId": 47725, + "factorToW": -1, + "maxAgeSeconds": 60, + "role": "physical_storage" + }, + { + "key": "physical_goodwe2", + "variableId": 35724, + "factorToW": -1, + "maxAgeSeconds": 60, + "role": "physical_storage" + }, + { + "key": "physical_solaredge", + "variableId": 21447, + "factorToW": 1, + "maxAgeSeconds": 60, + "role": "physical_storage" + }, + { + "key": "ev_account", + "variableId": 52020, + "factorToW": 1, + "maxAgeSeconds": 60, + "role": "reference" + }, + { + "key": "sdl_account", + "variableId": 25085, + "factorToW": 1, + "maxAgeSeconds": 60, + "role": "reference" + }, + { + "key": "grid_display", + "variableId": 49301, + "role": "reference", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "soc_goodwe1", + "variableId": 27361, + "role": "reference", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "soc_goodwe2", + "variableId": 23109, + "role": "reference", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "soc_solaredge", + "variableId": 51938, + "role": "reference", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "soc_ev", + "variableId": 32871, + "role": "reference", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "soc_sdl", + "variableId": 23879, + "role": "reference", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "ev_available_charge", + "variableId": 50230, + "role": "reference", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "ev_available_discharge", + "variableId": 43899, + "role": "reference", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "energy_t1", + "variableId": 59607, + "role": "reference", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "energy_t2", + "variableId": 26620, + "role": "reference", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "ev_requested", + "variableId": 19651, + "role": "reference", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "sdl_requested", + "variableId": 38943, + "role": "sdl_request", + "factorToW": 1, + "maxAgeSeconds": 120 + }, + { + "key": "solar_ac_value", + "variableId": 37975, + "role": "solar_raw", + "factorToW": 1, + "maxAgeSeconds": 60 + }, + { + "key": "solar_ac_scale", + "variableId": 41853, + "role": "solar_scale", + "factorToW": 1, + "maxAgeSeconds": 60 + } + ], + "formula": "solar_terminal_v1", + "solarReference": { + "pvKey": "pv_solaredge", + "batteryKey": "physical_solaredge", + "rawKey": "solar_ac_value", + "scaleKey": "solar_ac_scale" + }, + "minimumCoverage": 0.95, + "maximumGapSeconds": 10, + "minimumTrainingHours": 24, + "historyDays": 28, + "sourceDatasetId": "lihrenmoos-physical-v1", + "historyTimingPolicy": { + "version": 1, + "method": "equal_endpoint_v1", + "sources": { + "pv_goodwe1": { + "maxSpanSeconds": 130, + "evidenceId": "symcon-modbus-publish60-poll60-jitter10" + }, + "physical_goodwe1": { + "maxSpanSeconds": 130, + "evidenceId": "symcon-modbus-publish60-poll60-jitter10" + }, + "pv_goodwe2": { + "maxSpanSeconds": 75, + "evidenceId": "symcon-modbus-publish60-poll5-jitter10" + }, + "physical_goodwe2": { + "maxSpanSeconds": 75, + "evidenceId": "symcon-modbus-publish60-poll5-jitter10" + }, + "solar_ac_scale": { + "maxSpanSeconds": 90, + "evidenceId": "configured-scale-poll20-bounded-equal-value-estimate" + } + } + } +} diff --git a/services/netplan-v4/commissioning/history-timing-source/source/Dockerfile b/services/netplan-v4/commissioning/history-timing-source/source/Dockerfile new file mode 100644 index 0000000..ceff456 --- /dev/null +++ b/services/netplan-v4/commissioning/history-timing-source/source/Dockerfile @@ -0,0 +1,24 @@ +FROM python:3.11-slim +WORKDIR /app +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt +COPY netplan_v4 ./netplan_v4 +COPY tests ./tests +COPY run_tests.py install_hooks.py runtime_preflight.py Dockerfile . +COPY integrations ./integrations +COPY gui ./gui +COPY release_preflight.py forecast_acceptance.py deploy_integrated_shadow.py approved_previous_assets.json ./ +COPY acceptance ./acceptance +COPY commissioning/deploy_application.py ./commissioning/deploy_application.py +COPY commissioning/deploy_history_timing.py ./commissioning/deploy_history_timing.py +# Host sources can be 0600/0700. COPY makes them root-owned. +# Normalize only packaged application code; never change host secrets or sockets. +RUN find /app -type d -exec chmod 0755 {} + \ + && find /app -type f -exec chmod 0644 {} + +ENV PYTHONDONTWRITEBYTECODE=1 \ + PYTHONUNBUFFERED=1 \ + NETPLAN_V4_TEST_REPORT_DIR=/tmp/test-results +USER 1000:1000 +# Fail the build early if the actual unprivileged runtime cannot read the code. +RUN python /app/runtime_preflight.py +CMD ["uvicorn", "netplan_v4.service:from_environment", "--factory", "--host", "0.0.0.0", "--port", "9100", "--workers", "1"] diff --git a/services/netplan-v4/commissioning/history-timing-source/source/netplan_v4/history_timing.py b/services/netplan-v4/commissioning/history-timing-source/source/netplan_v4/history_timing.py new file mode 100644 index 0000000..6205d84 --- /dev/null +++ b/services/netplan-v4/commissioning/history-timing-source/source/netplan_v4/history_timing.py @@ -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 diff --git a/services/netplan-v4/commissioning/history-timing-source/source/netplan_v4/measurement_pipeline.py b/services/netplan-v4/commissioning/history-timing-source/source/netplan_v4/measurement_pipeline.py new file mode 100644 index 0000000..0579467 --- /dev/null +++ b/services/netplan-v4/commissioning/history-timing-source/source/netplan_v4/measurement_pipeline.py @@ -0,0 +1,500 @@ +"""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 + +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): + 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): + 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') + if record.get('mappingSha256') != c['mappingSha256'] 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') + rows = [project(r,c,plant,now) 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 r in 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 + 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): + """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']) + 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) + 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) 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'], + 'observationDatasetId':c.get('sourceDatasetId',ds), + 'historyTimingMethod':c.get('historyTimingPolicy',{}).get('method','strict_expiry')}) + return {'datasets':out,'liveEnabled':False,'legacyHistoryModified':False,'historyTimingVersion':1} diff --git a/services/netplan-v4/commissioning/history-timing-source/source/tests/test_history_timing.py b/services/netplan-v4/commissioning/history-timing-source/source/tests/test_history_timing.py new file mode 100644 index 0000000..b05b046 --- /dev/null +++ b/services/netplan-v4/commissioning/history-timing-source/source/tests/test_history_timing.py @@ -0,0 +1,204 @@ +"""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() diff --git a/services/netplan-v4/commissioning/history-timing-source/source/tests/test_history_timing_deployment.py b/services/netplan-v4/commissioning/history-timing-source/source/tests/test_history_timing_deployment.py new file mode 100644 index 0000000..2acf194 --- /dev/null +++ b/services/netplan-v4/commissioning/history-timing-source/source/tests/test_history_timing_deployment.py @@ -0,0 +1,94 @@ +"""Installer ordering/recovery tests; subprocesses, privileges and DB backup are simulated.""" +import hashlib +import importlib.util +import json +import subprocess +import sys +import tempfile +import types +import unittest +from pathlib import Path +from unittest.mock import patch + +SCRIPT=Path(__file__).resolve().parents[1]/'commissioning/deploy_history_timing.py' +spec=importlib.util.spec_from_file_location('history_deploy_tests',SCRIPT) +d=importlib.util.module_from_spec(spec);spec.loader.exec_module(d) + + +class TimingDeploymentTest(unittest.TestCase): + def setUp(self): + self.tmp=tempfile.TemporaryDirectory();self.root=Path(self.tmp.name) + self.package=self.root/'commissioning/history-timing-source';self.package.mkdir(parents=True) + files={} + for index,name in enumerate(d.TARGETS): + original=(b'old:'+name.encode()) if index%2==0 else None + target=self.root/name;target.parent.mkdir(parents=True,exist_ok=True) + if original is not None:target.write_bytes(original) + candidate=self.package/'source'/name;candidate.parent.mkdir(parents=True,exist_ok=True);candidate.write_bytes(b'new:'+name.encode()) + files[name]={'before':hashlib.sha256(original).hexdigest() if original else None,'after':d.digest(candidate)} + (self.package/'dataset.json').write_text('{}') + (self.root/'commissioning/deploy_application.py').write_text('synthetic helper') + self.manifest={'scope':'historical_publication_only','installationId':'plant','files':files, + 'existingDeploymentHelperSha256':d.digest(self.root/'commissioning/deploy_application.py'), + 'datasetSha256':d.digest(self.package/'dataset.json'),'installerSha256':d.digest(SCRIPT)} + (self.package/'RELEASE.json').write_text(json.dumps(self.manifest)) + self.context=patch.multiple(d,ROOT=self.root,PACKAGE=self.package);self.context.start() + self.calls=[] + def tearDown(self):self.context.stop();self.tmp.cleanup() + def fake_run(self,args,**kwargs): + self.calls.append(args) + if 'inspect' in args:out='sha256:test-old-image' + elif 'ps' in args:out='test-container' + elif 'exec' in args:out=json.dumps({'settingsUnchanged':True,'liveEnabled':False}) + else:out='' + return subprocess.CompletedProcess(args,0,out,'') + def execute(self,fail_at=None): + def run(args,**kwargs): + out=self.fake_run(args,**kwargs) + if fail_at and fail_at in args:raise subprocess.CalledProcessError(1,args) + return out + def atomic(path,data,*args,**kwargs):path.write_bytes(data) + def backup(source,dest):self.calls.append(['test-backup']);dest.write_text('synthetic backup') + helpers=types.SimpleNamespace(atomic=atomic,backup_database=backup) + with patch.object(d.os,'geteuid',return_value=0),patch.object(d.os,'chown'),patch.object(d.subprocess,'run',side_effect=run),patch.dict(sys.modules,{'deploy_application':helpers}): + d.run('plant',True) + def report(self):return json.loads(next(self.root.glob('history-timing-releases/*/REPORT.json')).read_text()) + def test_default_checks_only(self): + d.run('plant');self.assertEqual(self.calls,[]) + self.assertFalse((self.root/'history-timing-releases').exists()) + def test_order_tests_backup_then_only_v4_replacement(self): + self.execute() + test=next(i for i,c in enumerate(self.calls) if '/app/run_tests.py' in c) + backup=next(i for i,c in enumerate(self.calls) if c==['test-backup']) + up=next(i for i,c in enumerate(self.calls) if 'up' in c) + self.assertLess(test,backup);self.assertLess(backup,up) + self.assertTrue(all('license-portal' not in c and 'forecast-engine' not in c for c in self.calls)) + self.assertEqual(self.report()['status'],'historical_timing_installed_no_actuation') + def test_failed_tests_restore_host_source_without_restart(self): + with self.assertRaises(subprocess.CalledProcessError):self.execute('/app/run_tests.py') + self.assertFalse(any('up' in c for c in self.calls)) + for name,v in self.manifest['files'].items():self.assertEqual(d.digest(self.root/name),v['before']) + def test_failed_registration_rolls_image_back(self): + with self.assertRaises(subprocess.CalledProcessError):self.execute('exec') + self.assertIn('previous_image_restored',self.report()['rollback']) + self.assertTrue(any('--pull' in c and 'never' in c for c in self.calls)) + def test_source_drift_refused_before_actions(self): + (self.root/d.TARGETS[0]).write_text('parallel change') + with self.assertRaises(ValueError):d.verify() + self.assertEqual(self.calls,[]) + def test_source_already_installed_idempotent(self): + for name in d.TARGETS:(self.root/name).write_bytes((self.package/'source'/name).read_bytes()) + self.assertEqual(d.verify()['scope'],'historical_publication_only') + self.execute() + self.assertEqual(self.report()['status'],'historical_timing_installed_no_actuation') + def test_no_production_selection_or_permission_in_bootstrap(self): + self.assertNotIn("'/settings'",d.BOOTSTRAP) + self.assertNotIn('/trial/arm',d.BOOTSTRAP) + self.assertIn("before['settings']==after['settings']",d.BOOTSTRAP) + self.assertIn("historyTimingVersion",d.BOOTSTRAP) + def test_target_image_contains_installer_test_dependency(self): + dockerfile=(SCRIPT.parents[1]/'Dockerfile').read_text() + self.assertIn('COPY commissioning/deploy_history_timing.py ./commissioning/deploy_history_timing.py',dockerfile) + + +if __name__=='__main__':unittest.main() diff --git a/services/netplan-v4/netplan_v4/history_timing.py b/services/netplan-v4/netplan_v4/history_timing.py new file mode 100644 index 0000000..6205d84 --- /dev/null +++ b/services/netplan-v4/netplan_v4/history_timing.py @@ -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 diff --git a/services/netplan-v4/netplan_v4/measurement_pipeline.py b/services/netplan-v4/netplan_v4/measurement_pipeline.py index ead16bc..0579467 100644 --- a/services/netplan-v4/netplan_v4/measurement_pipeline.py +++ b/services/netplan-v4/netplan_v4/measurement_pipeline.py @@ -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} diff --git a/services/netplan-v4/tests/test_history_timing.py b/services/netplan-v4/tests/test_history_timing.py new file mode 100644 index 0000000..b05b046 --- /dev/null +++ b/services/netplan-v4/tests/test_history_timing.py @@ -0,0 +1,204 @@ +"""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() diff --git a/services/netplan-v4/tests/test_history_timing_deployment.py b/services/netplan-v4/tests/test_history_timing_deployment.py new file mode 100644 index 0000000..2acf194 --- /dev/null +++ b/services/netplan-v4/tests/test_history_timing_deployment.py @@ -0,0 +1,94 @@ +"""Installer ordering/recovery tests; subprocesses, privileges and DB backup are simulated.""" +import hashlib +import importlib.util +import json +import subprocess +import sys +import tempfile +import types +import unittest +from pathlib import Path +from unittest.mock import patch + +SCRIPT=Path(__file__).resolve().parents[1]/'commissioning/deploy_history_timing.py' +spec=importlib.util.spec_from_file_location('history_deploy_tests',SCRIPT) +d=importlib.util.module_from_spec(spec);spec.loader.exec_module(d) + + +class TimingDeploymentTest(unittest.TestCase): + def setUp(self): + self.tmp=tempfile.TemporaryDirectory();self.root=Path(self.tmp.name) + self.package=self.root/'commissioning/history-timing-source';self.package.mkdir(parents=True) + files={} + for index,name in enumerate(d.TARGETS): + original=(b'old:'+name.encode()) if index%2==0 else None + target=self.root/name;target.parent.mkdir(parents=True,exist_ok=True) + if original is not None:target.write_bytes(original) + candidate=self.package/'source'/name;candidate.parent.mkdir(parents=True,exist_ok=True);candidate.write_bytes(b'new:'+name.encode()) + files[name]={'before':hashlib.sha256(original).hexdigest() if original else None,'after':d.digest(candidate)} + (self.package/'dataset.json').write_text('{}') + (self.root/'commissioning/deploy_application.py').write_text('synthetic helper') + self.manifest={'scope':'historical_publication_only','installationId':'plant','files':files, + 'existingDeploymentHelperSha256':d.digest(self.root/'commissioning/deploy_application.py'), + 'datasetSha256':d.digest(self.package/'dataset.json'),'installerSha256':d.digest(SCRIPT)} + (self.package/'RELEASE.json').write_text(json.dumps(self.manifest)) + self.context=patch.multiple(d,ROOT=self.root,PACKAGE=self.package);self.context.start() + self.calls=[] + def tearDown(self):self.context.stop();self.tmp.cleanup() + def fake_run(self,args,**kwargs): + self.calls.append(args) + if 'inspect' in args:out='sha256:test-old-image' + elif 'ps' in args:out='test-container' + elif 'exec' in args:out=json.dumps({'settingsUnchanged':True,'liveEnabled':False}) + else:out='' + return subprocess.CompletedProcess(args,0,out,'') + def execute(self,fail_at=None): + def run(args,**kwargs): + out=self.fake_run(args,**kwargs) + if fail_at and fail_at in args:raise subprocess.CalledProcessError(1,args) + return out + def atomic(path,data,*args,**kwargs):path.write_bytes(data) + def backup(source,dest):self.calls.append(['test-backup']);dest.write_text('synthetic backup') + helpers=types.SimpleNamespace(atomic=atomic,backup_database=backup) + with patch.object(d.os,'geteuid',return_value=0),patch.object(d.os,'chown'),patch.object(d.subprocess,'run',side_effect=run),patch.dict(sys.modules,{'deploy_application':helpers}): + d.run('plant',True) + def report(self):return json.loads(next(self.root.glob('history-timing-releases/*/REPORT.json')).read_text()) + def test_default_checks_only(self): + d.run('plant');self.assertEqual(self.calls,[]) + self.assertFalse((self.root/'history-timing-releases').exists()) + def test_order_tests_backup_then_only_v4_replacement(self): + self.execute() + test=next(i for i,c in enumerate(self.calls) if '/app/run_tests.py' in c) + backup=next(i for i,c in enumerate(self.calls) if c==['test-backup']) + up=next(i for i,c in enumerate(self.calls) if 'up' in c) + self.assertLess(test,backup);self.assertLess(backup,up) + self.assertTrue(all('license-portal' not in c and 'forecast-engine' not in c for c in self.calls)) + self.assertEqual(self.report()['status'],'historical_timing_installed_no_actuation') + def test_failed_tests_restore_host_source_without_restart(self): + with self.assertRaises(subprocess.CalledProcessError):self.execute('/app/run_tests.py') + self.assertFalse(any('up' in c for c in self.calls)) + for name,v in self.manifest['files'].items():self.assertEqual(d.digest(self.root/name),v['before']) + def test_failed_registration_rolls_image_back(self): + with self.assertRaises(subprocess.CalledProcessError):self.execute('exec') + self.assertIn('previous_image_restored',self.report()['rollback']) + self.assertTrue(any('--pull' in c and 'never' in c for c in self.calls)) + def test_source_drift_refused_before_actions(self): + (self.root/d.TARGETS[0]).write_text('parallel change') + with self.assertRaises(ValueError):d.verify() + self.assertEqual(self.calls,[]) + def test_source_already_installed_idempotent(self): + for name in d.TARGETS:(self.root/name).write_bytes((self.package/'source'/name).read_bytes()) + self.assertEqual(d.verify()['scope'],'historical_publication_only') + self.execute() + self.assertEqual(self.report()['status'],'historical_timing_installed_no_actuation') + def test_no_production_selection_or_permission_in_bootstrap(self): + self.assertNotIn("'/settings'",d.BOOTSTRAP) + self.assertNotIn('/trial/arm',d.BOOTSTRAP) + self.assertIn("before['settings']==after['settings']",d.BOOTSTRAP) + self.assertIn("historyTimingVersion",d.BOOTSTRAP) + def test_target_image_contains_installer_test_dependency(self): + dockerfile=(SCRIPT.parents[1]/'Dockerfile').read_text() + self.assertIn('COPY commissioning/deploy_history_timing.py ./commissioning/deploy_history_timing.py',dockerfile) + + +if __name__=='__main__':unittest.main()