feat(forecast): commission corrected Lihrenmoos profile
This commit is contained in:
@@ -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))
|
||||
Reference in New Issue
Block a user