diff --git a/docs/NETPLAN_V4_AGENT_HANDOFF.md b/docs/NETPLAN_V4_AGENT_HANDOFF.md index 780137e..2bae37f 100644 --- a/docs/NETPLAN_V4_AGENT_HANDOFF.md +++ b/docs/NETPLAN_V4_AGENT_HANDOFF.md @@ -4,6 +4,8 @@ **Zeitbezug:** Die letzten hier belegten Betriebsberichte stammen vom **03.10.2026, ca. 14:02 Uhr Europe/Zurich**. Am 04.10. wurden für diese Übergabe die genannten Berichte und der Git-Stand erneut gelesen, aber keine neue Anlagenabnahme durchgeführt. Alte Messzahlen deshalb niemals als heutige Live-Werte ausgeben. +**Fortsetzung Prognose am 04.10.2026:** Der aktuelle Runtime-Stand wurde gelesen, ohne die alten Diagnose- oder Installationsschritte zu wiederholen. Der abgeleitete Datensatz `lihrenmoos-physical-published-v2` ist `model_ready` und erfüllt mit mehr als 24 nutzbaren äquivalenten Stunden die Trainingsschwelle. Die V4-Einstellungen stehen seit Revision 2 in `shadow` auf `forecastSource=corrected_profile` und `measurementDataset=lihrenmoos-physical-published-v2`; `liveEnabled` und Stellfreigabe bleiben aus. Der erste echte Lauf zeigte einen kompatibilitätsbedingten Stopp bei mikrosekundengenauen `observedAt`-Zeitstempeln. Die Korrektur rundet ausschliesslich die kausale Verfügbarkeit von Prognoseereignissen auf die nächste volle Sekunde auf; Rohmessungen bleiben streng ganzsekündlich. 310 Python- und 16 Portal-Tests sind erfolgreich. Repository und Runtime-Build-Kontext enthalten bytegleich den geprüften Fix; der laufende Container enthält ihn noch nicht, weil beide Agentzugänge weder Docker-Socket noch `sudo` erhalten. Der vorbereitete, syntaktisch geprüfte Root-Rollout liegt unter `/home/agent/services/netplan-v4-shadow/commissioning/finish_corrected_forecast_rollout.sh` und baut/testet vor dem Austausch, hält das vorige Image als Rollback fest und ändert keine Stellfreigabe. Daniel hat danach einen Aktivbetrieb angefragt; der aktuelle V4-Kern ist jedoch technisch `shadow-only`, und die weiterhin unbelegte Rückmeldung (`feedback_source_skew`, `usableForTrial=false`, `canDispatch=false`) verbietet ein blosses Umstellen. Keine Live-Freigabe wurde erteilt oder eingebaut. Git-Klarstellung durch Daniel: Commit-Autor `dh_Agent `, Gitea-Konto `dh`; Commit/Push auf `develop` und `beta` am 04.10.2026 ausdrücklich freigegeben. Der Agentzugang besitzt derzeit keine HTTPS-Schreibanmeldung für `dh`. + ## 1. Zuerst lesen: Wo wir tatsächlich stehen - Daniel möchte die **fertige, produktionsgeeignete Anwendung**, nicht weitere isolierte Sammler, Diagnosekategorien oder wiederholte Bestätigungsrunden. Er hat die fortlaufende Umsetzung mehrfach beauftragt. Probleme im Code selbst beheben, Tests und Auslieferung bündeln; ihn nur für wirklich notwendige Root-/Symcon-Ausführung oder echte Anlagenfreigaben einbeziehen. diff --git a/services/netplan-v4/commissioning/activate_corrected_forecast.py b/services/netplan-v4/commissioning/activate_corrected_forecast.py new file mode 100644 index 0000000..f3b251a --- /dev/null +++ b/services/netplan-v4/commissioning/activate_corrected_forecast.py @@ -0,0 +1,256 @@ +"""Commission the corrected Lihrenmoos forecast in shadow mode. + +The command is a dry run unless --apply is supplied. It never enables live control, +changes Symcon actuator flags, or reads/prints service credentials. +""" +from __future__ import annotations + +import argparse +import hashlib +import json +import os +import sqlite3 +import sys +import time +from datetime import datetime, timezone +from pathlib import Path +from uuid import UUID + +SERVICE_ROOT = Path(__file__).resolve().parents[1] +if str(SERVICE_ROOT) not in sys.path: + sys.path.insert(0, str(SERVICE_ROOT)) +RUNTIME_ROOT = Path("/home/agent/services/netplan-v4-shadow") +PLANT = "e3a08f9e-af12-4695-99bd-8b51c0520021" +DATASET = "lihrenmoos-physical-published-v2" +SOURCE_FILES = ( + "netplan_v4/store.py", + "netplan_v4/measurement_pipeline.py", + "netplan_v4/service.py", +) + + +def digest(path: Path) -> str: + if path.is_symlink() or not path.is_file(): + raise ValueError("Source file missing or is a symlink") + return hashlib.sha256(path.read_bytes()).hexdigest() + + +def verify_runtime_source(runtime_root: Path) -> None: + for name in SOURCE_FILES: + if digest(SERVICE_ROOT / name) != digest(runtime_root / name): + raise ValueError("Runtime source differs from reviewed repository: " + name) + + +def read_state(db: Path, plant: str, dataset: str, now: int) -> dict: + stamp = db.stat() + con = sqlite3.connect(f"file:{db}?mode=ro", uri=True) + con.row_factory = sqlite3.Row + try: + from netplan_v4.store import DEFAULT_SETTINGS + + row = con.execute( + "SELECT revision,value FROM planner_settings WHERE plant=?", (plant,) + ).fetchone() + settings = {**DEFAULT_SETTINGS, **(json.loads(row["value"]) if row else {})} + settings["revision"] = row["revision"] if row else 0 + ds = con.execute( + "SELECT config FROM planner_data_sets WHERE plant=? AND dataset=?", + (plant, dataset), + ).fetchone() + if not ds: + raise ValueError("Corrected dataset is not registered") + config = json.loads(ds["config"]) + state = con.execute( + "SELECT status,detail FROM planner_pipeline_state WHERE plant=? AND dataset=?", + (plant, dataset), + ).fetchone() + model = con.execute( + """SELECT m.model_id,m.trained_at,m.trained_through,m.value + FROM planner_model_current c JOIN planner_load_models m + ON m.plant=c.plant AND m.dataset=c.dataset AND m.model_id=c.model_id + WHERE c.plant=? AND c.dataset=?""", + (plant, dataset), + ).fetchone() + observation_dataset = config.get("sourceDatasetId", dataset) + observation = con.execute( + """SELECT COUNT(*) records,MAX(captured_at) last_capture, + MAX(received_at) last_received + FROM planner_observations WHERE plant=? AND dataset=?""", + (plant, observation_dataset), + ).fetchone() + plan = con.execute( + """SELECT p.value FROM planner_current c JOIN planner_plans p + ON p.plan_id=c.plan_id WHERE c.plant=? AND c.mode='shadow'""", + (plant,), + ).fetchone() + detail = json.loads(state["detail"]) if state else {} + model_value = json.loads(model["value"]) if model else None + return { + "settings": settings, + "dataset": dataset, + "observationDataset": observation_dataset, + "pipelineStatus": state["status"] if state else None, + "pipelineDetail": detail, + "records": observation["records"], + "lastCapture": observation["last_capture"], + "lastReceived": observation["last_received"], + "activeModel": { + "modelId": model["model_id"], + "trainedAt": model["trained_at"], + "trainedThrough": model["trained_through"], + "validation": model_value.get("validation"), + } if model else None, + "currentPlan": json.loads(plan["value"]) if plan else None, + "checkedAt": now, + "databaseMtimeNs": stamp.st_mtime_ns, + } + finally: + con.close() + + +def require_ready(state: dict, now: int) -> None: + if state["settings"]["runMode"] != "shadow": + raise ValueError("Only shadow mode may be commissioned") + if state["pipelineStatus"] != "model_ready" or not state["activeModel"]: + raise ValueError("Corrected profile model is not ready") + if state["records"] < 1 or not isinstance(state["lastCapture"], int): + raise ValueError("Corrected observations are missing") + if not 0 <= now - state["lastCapture"] <= 600: + raise ValueError("Corrected observations are not current") + if not 0 <= now - state["lastReceived"] <= 600: + raise ValueError("Corrected delivery is not current") + detail = state["pipelineDetail"] + if detail.get("modelId") != state["activeModel"]["modelId"]: + raise ValueError("Pipeline and active model disagree") + if detail.get("usableEquivalentHours", 0) < detail.get("requiredEquivalentHours", 24): + raise ValueError("Minimum usable equivalent hours not reached") + + +def backup_database(source: Path, destination: Path) -> None: + if destination.exists() or source.is_symlink(): + raise ValueError("Backup destination must be new and source must be regular") + with sqlite3.connect(f"file:{source}?mode=ro", uri=True) as src, sqlite3.connect(destination) as dst: + src.backup(dst) + if dst.execute("PRAGMA integrity_check").fetchone()[0] != "ok": + raise ValueError("Database backup integrity check failed") + destination.chmod(0o600) + + +def plan_uses_dataset(db: Path, plant: str, dataset: str, revision: int) -> dict | None: + con = sqlite3.connect(f"file:{db}?mode=ro", uri=True) + try: + row = con.execute( + """SELECT p.value FROM planner_current c JOIN planner_plans p + ON p.plan_id=c.plan_id WHERE c.plant=? AND c.mode='shadow'""", + (plant,), + ).fetchone() + if not row: + return None + plan = json.loads(row[0]) + pipeline = plan.get("inputQuality", {}).get("dataPipeline", {}) + if plan.get("configRevision") == revision and pipeline.get("datasetId") == dataset: + return plan + return None + finally: + con.close() + + +def activate(db: Path, plant: str, dataset: str, now: datetime) -> dict: + from netplan_v4.store import PlannerStore + + store = PlannerStore(str(db)) + try: + before = store.settings(plant) + if before["runMode"] != "shadow": + raise ValueError("Runtime left shadow mode") + if before["forecastSource"] == "corrected_profile" and before["measurementDataset"] == dataset: + store.request(plant, "manual", now) + return {"changed": False, "revision": before["revision"]} + result = store.save_settings( + plant, + {"forecastSource": "corrected_profile", "measurementDataset": dataset}, + before["revision"], + now, + ) + return {"changed": True, "revision": result["revision"]} + finally: + store.close() + + +def commission(runtime_root: Path, plant: str, dataset: str, apply: bool, wait_seconds: int = 90) -> dict: + plant = str(UUID(plant)) + if runtime_root.resolve() != RUNTIME_ROOT: + raise ValueError("Unexpected runtime root") + verify_runtime_source(runtime_root) + db = runtime_root / "data/netplan-v4.sqlite" + now = datetime.now(timezone.utc).replace(microsecond=0) + before = read_state(db, plant, dataset, int(now.timestamp())) + require_ready(before, int(now.timestamp())) + report = { + "scope": "corrected_forecast_shadow_commissioning", + "startedAt": now.isoformat(), + "installationId": plant, + "datasetId": dataset, + "liveEnabled": False, + "actuatorPermissionGranted": False, + "before": { + "settings": before["settings"], + "pipelineStatus": before["pipelineStatus"], + "records": before["records"], + "lastCapture": before["lastCapture"], + "activeModel": before["activeModel"], + "usableEquivalentHours": before["pipelineDetail"].get("usableEquivalentHours"), + }, + } + if not apply: + report["status"] = "ready_no_changes" + return report + + folder = runtime_root / "corrected-forecast-releases" / now.strftime("%Y%m%dT%H%M%SZ") + folder.mkdir(parents=True, mode=0o700) + backup_database(db, folder / "before.sqlite") + activated = activate(db, plant, dataset, now) + report["activation"] = activated + deadline = time.monotonic() + wait_seconds + plan = None + while time.monotonic() < deadline: + plan = plan_uses_dataset(db, plant, dataset, activated["revision"]) + if plan: + break + time.sleep(2) + after = read_state(db, plant, dataset, int(datetime.now(timezone.utc).timestamp())) + settings = after["settings"] + if settings["runMode"] != "shadow" or settings["forecastSource"] != "corrected_profile" or settings["measurementDataset"] != dataset: + raise ValueError("Corrected shadow settings were not retained") + report["after"] = { + "settings": settings, + "pipelineStatus": after["pipelineStatus"], + "activeModel": after["activeModel"], + "planId": plan.get("planId") if plan else None, + "planStatus": plan.get("status") if plan else None, + "sourceFamily": plan.get("sourceFamily") if plan else None, + "dataPipeline": plan.get("inputQuality", {}).get("dataPipeline") if plan else None, + } + report["status"] = "corrected_forecast_running_shadow" if plan else "corrected_forecast_activated_plan_pending" + report["finishedAt"] = datetime.now(timezone.utc).isoformat() + output = folder / "REPORT.json" + output.write_text(json.dumps(report, indent=2) + "\n") + output.chmod(0o640) + report["reportPath"] = str(output) + return report + + +if __name__ == "__main__": + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--plant", default=PLANT) + parser.add_argument("--dataset", default=DATASET) + parser.add_argument("--apply", action="store_true") + parser.add_argument("--wait-seconds", type=int, default=90) + args = parser.parse_args() + if not 1 <= args.wait_seconds <= 300: + raise SystemExit("wait-seconds must be 1..300") + try: + result = commission(RUNTIME_ROOT, args.plant, args.dataset, args.apply, args.wait_seconds) + except Exception as exc: + raise SystemExit("Stopped: " + type(exc).__name__ + ". No actuator permission was changed.") from exc + print(json.dumps(result, indent=2)) diff --git a/services/netplan-v4/netplan_v4/measurement_pipeline.py b/services/netplan-v4/netplan_v4/measurement_pipeline.py index 61ca5a9..e5af4da 100644 --- a/services/netplan-v4/netplan_v4/measurement_pipeline.py +++ b/services/netplan-v4/netplan_v4/measurement_pipeline.py @@ -9,7 +9,7 @@ from bisect import bisect_right from collections import defaultdict from datetime import datetime, timedelta, timezone from hashlib import sha256 -from math import isfinite +from math import ceil, isfinite from statistics import median from zoneinfo import ZoneInfo import json @@ -34,6 +34,16 @@ def epoch(value): return int(t.timestamp()) +def availability_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: + raise ValueError('Explicit UTC timestamp required') + # Never claim an input existed before its sub-second observation completed. + return int(ceil(t.timestamp())) + + def iso(t): return datetime.fromtimestamp(t, UTC).isoformat() @@ -479,7 +489,7 @@ def apply_load_forecast(con,plant,dataset,forecast,decision): 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'])) + result['observedAt']=iso(max(availability_epoch(forecast['observedAt']),model['trainedAt'],last['capturedAt'])) for family,old in forecast['families'].items(): if family not in FAMILIES: continue points=[] diff --git a/services/netplan-v4/tests/test_corrected_forecast_commissioning.py b/services/netplan-v4/tests/test_corrected_forecast_commissioning.py new file mode 100644 index 0000000..dda7d69 --- /dev/null +++ b/services/netplan-v4/tests/test_corrected_forecast_commissioning.py @@ -0,0 +1,116 @@ +import importlib.util +import json +import sqlite3 +import tempfile +import unittest +from datetime import datetime, timezone +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +SPEC = importlib.util.spec_from_file_location( + "activate_corrected_forecast", ROOT / "commissioning/activate_corrected_forecast.py" +) +module = importlib.util.module_from_spec(SPEC) +SPEC.loader.exec_module(module) + +from netplan_v4 import measurement_pipeline +from netplan_v4.store import PlannerStore + +PLANT = module.PLANT +DATASET = module.DATASET +NOW = 1791099600 + + +class CorrectedForecastCommissioningTest(unittest.TestCase): + def setUp(self): + self.tmp = tempfile.TemporaryDirectory() + self.db = Path(self.tmp.name) / "planner.sqlite" + self.store = PlannerStore(str(self.db)) + config = { + "datasetId": DATASET, + "mappingSha256": "a" * 64, + "inventorySha256": "b" * 64, + "formula": "physical_sum_v1", + "solarReference": None, + "minimumCoverage": 0.95, + "maximumGapSeconds": 10, + "minimumTrainingHours": 24, + "historyDays": 14, + "sources": [ + {"key": "grid", "variableId": 1, "role": "grid", "factorToW": 1, "maxAgeSeconds": 60}, + {"key": "pv", "variableId": 2, "role": "pv", "factorToW": 1, "maxAgeSeconds": 60}, + {"key": "battery", "variableId": 3, "role": "physical_storage", "factorToW": 1, "maxAgeSeconds": 60}, + ], + } + measurement_pipeline.register_dataset(self.store.con, PLANT, config, NOW - 100) + model = { + "version": 1, + "datasetId": DATASET, + "methodVersion": "profile_v1", + "trainedAt": NOW - 60, + "trainedThrough": NOW - 300, + "profiles": {}, + "validation": {"status": "bootstrap_insufficient_holdout"}, + } + model_id = "c" * 64 + with self.store.con: + self.store.con.execute( + "INSERT INTO planner_load_models VALUES(?,?,?,?,?,?)", + (PLANT, DATASET, model_id, NOW - 60, NOW - 300, json.dumps(model)), + ) + self.store.con.execute( + "INSERT INTO planner_model_current VALUES(?,?,?)", (PLANT, DATASET, model_id) + ) + self.store.con.execute( + "INSERT INTO planner_pipeline_state VALUES(?,?,?,?,?)", + ( + PLANT, + DATASET, + 1, + "model_ready", + json.dumps( + { + "modelId": model_id, + "usableEquivalentHours": 24.1, + "requiredEquivalentHours": 24, + } + ), + ), + ) + self.store.con.execute( + "INSERT INTO planner_observations VALUES(?,?,?,?,?,?)", + (PLANT, DATASET, NOW - 30, NOW - 20, "d" * 64, "{}"), + ) + self.store.close() + + def tearDown(self): + self.tmp.cleanup() + + def test_ready_state_can_activate_only_corrected_shadow_settings(self): + state = module.read_state(self.db, PLANT, DATASET, NOW) + module.require_ready(state, NOW) + result = module.activate( + self.db, PLANT, DATASET, datetime.fromtimestamp(NOW, timezone.utc) + ) + self.assertTrue(result["changed"]) + store = PlannerStore(str(self.db)) + settings = store.settings(PLANT) + self.assertEqual(settings["forecastSource"], "corrected_profile") + self.assertEqual(settings["measurementDataset"], DATASET) + self.assertEqual(settings["runMode"], "shadow") + store.close() + + def test_stale_observations_refused(self): + state = module.read_state(self.db, PLANT, DATASET, NOW + 700) + with self.assertRaises(ValueError): + module.require_ready(state, NOW + 700) + + def test_backup_is_consistent(self): + backup = Path(self.tmp.name) / "backup.sqlite" + module.backup_database(self.db, backup) + with sqlite3.connect(backup) as con: + self.assertEqual(con.execute("PRAGMA integrity_check").fetchone()[0], "ok") + + +if __name__ == "__main__": + unittest.main() diff --git a/services/netplan-v4/tests/test_measurement_pipeline.py b/services/netplan-v4/tests/test_measurement_pipeline.py index ab8cfa6..85b4f38 100644 --- a/services/netplan-v4/tests/test_measurement_pipeline.py +++ b/services/netplan-v4/tests/test_measurement_pipeline.py @@ -67,6 +67,9 @@ class MeasurementPipelineTest(unittest.TestCase): def test_no_naive_time(self): r=record(NOW);r['capturedAt']='2026-10-02T12:00:00' with self.assertRaises(ValueError):self.batch([r]) + def test_measurement_time_remains_whole_second(self): + r=record(NOW);r['capturedAt']='2026-10-02T12:00:00.000001+00:00' + with self.assertRaises(ValueError):self.batch([r]) def test_future_source_not_rejuvenated(self): r=record(NOW);r['raw']['grid']['sourceUpdatedAt']=NOW+1 self.assertFalse(m.project(r,self.c,AID,NOW)['raw']['grid']['valid']) @@ -115,7 +118,10 @@ class MeasurementPipelineTest(unittest.TestCase): with self.assertRaises(ValueError):m.apply_load_forecast(self.s.con,AID,'physical-v1',fc,NOW) def test_external_is_labelled_persistence_not_published(self): self.history();m.advance(self.s.con,AID,'physical-v1',self.s.settings(AID),NOW) - op,fc,tar=inputs(datetime.fromtimestamp(NOW,timezone.utc));out=m.apply_load_forecast(self.s.con,AID,'physical-v1',fc,NOW) + op,fc,tar=inputs(datetime.fromtimestamp(NOW,timezone.utc)) + fc['observedAt']=datetime.fromtimestamp(NOW-1,timezone.utc).replace(microsecond=1).isoformat() + out=m.apply_load_forecast(self.s.con,AID,'physical-v1',fc,NOW) + self.assertEqual(out['observedAt'],m.iso(NOW)) self.assertEqual(out['families']['3']['points'][0]['loadW'],6000) self.assertFalse(out['families']['3']['dataPipeline']['futureSdlPublished']) self.assertNotIn('accountingEvidenceId',out['families']['3'])