from __future__ import annotations import json import sqlite3 from datetime import datetime,timedelta from uuid import uuid4 from .domain import default_registry,month_key,number,quarter_start,utc from . import meter_runtime, controlled_trial, economic_replay, retention, measurement_pipeline def canonical(value): return json.dumps(value,sort_keys=True,separators=(',',':'),allow_nan=False) DEFAULT_SETTINGS={'family':'3','autoLookbackDays':14,'autoMinimumDays':7,'autoMinimumCoverage':.9,'autoSwitchMarginChf':1.,'tariffPolicy':'published_only','trainingCadence':'daily','trainingPromotion':'validated_only','runMode':'shadow','roundtripEfficiency':.90,'measurementPolicy':'verified_only','peakOutlookPolicy':'empirical_if_available','peakOutlookReliance':.5,'forecastSource':'legacy','measurementDataset':''} class PlannerStore: """Own SQLite file, no mutation of legacy application databases.""" def __init__(self,path,registry=None): self.registry=registry or default_registry();self.con=sqlite3.connect(path,timeout=10) self.con.row_factory=sqlite3.Row;self.con.execute('PRAGMA foreign_keys=ON') self.con.executescript(''' CREATE TABLE IF NOT EXISTS planner_settings(plant TEXT PRIMARY KEY,revision INTEGER NOT NULL,value TEXT NOT NULL); CREATE TABLE IF NOT EXISTS planner_audit(id INTEGER PRIMARY KEY,plant TEXT NOT NULL,at TEXT NOT NULL,kind TEXT NOT NULL,detail TEXT NOT NULL); CREATE TABLE IF NOT EXISTS planner_work(plant TEXT PRIMARY KEY,sequence INTEGER NOT NULL,revision INTEGER NOT NULL,reasons TEXT NOT NULL,requested_at TEXT NOT NULL,lease_until TEXT,lease_token TEXT); CREATE TABLE IF NOT EXISTS planner_measurements(plant TEXT NOT NULL,start TEXT NOT NULL,import_kwh REAL NOT NULL,PRIMARY KEY(plant,start)); CREATE TABLE IF NOT EXISTS planner_month_peaks(plant TEXT NOT NULL,month TEXT NOT NULL,peak_kw REAL NOT NULL,source TEXT NOT NULL,updated_at TEXT NOT NULL,PRIMARY KEY(plant,month)); CREATE TABLE IF NOT EXISTS planner_snapshots(id TEXT PRIMARY KEY,plant TEXT NOT NULL,issued_at TEXT NOT NULL,value TEXT NOT NULL); CREATE TABLE IF NOT EXISTS planner_plans(plan_id TEXT PRIMARY KEY,plant TEXT NOT NULL,revision INTEGER NOT NULL,mode TEXT NOT NULL,value TEXT NOT NULL,created_at TEXT NOT NULL); CREATE TABLE IF NOT EXISTS planner_current(plant TEXT NOT NULL,mode TEXT NOT NULL,plan_id TEXT NOT NULL REFERENCES planner_plans(plan_id),PRIMARY KEY(plant,mode)); CREATE TABLE IF NOT EXISTS planner_ack(plant TEXT PRIMARY KEY,plan_id TEXT NOT NULL,revision INTEGER NOT NULL,received_at TEXT NOT NULL,applied_step TEXT,status TEXT NOT NULL); CREATE TABLE IF NOT EXISTS planner_inputs(plant TEXT NOT NULL,kind TEXT NOT NULL,event_id TEXT NOT NULL,observed_at TEXT NOT NULL,value TEXT NOT NULL,PRIMARY KEY(plant,kind,event_id)); CREATE TABLE IF NOT EXISTS planner_input_current(plant TEXT NOT NULL,kind TEXT NOT NULL,event_id TEXT NOT NULL,PRIMARY KEY(plant,kind)); CREATE TABLE IF NOT EXISTS planner_run_status(plant TEXT PRIMARY KEY,updated_at TEXT NOT NULL,status TEXT NOT NULL,detail TEXT NOT NULL); CREATE TABLE IF NOT EXISTS planner_ticks(plant TEXT PRIMARY KEY,tick INTEGER NOT NULL); ''') meter_runtime.schema(self.con) controlled_trial.schema(self.con) economic_replay.schema(self.con) retention.schema(self.con) measurement_pipeline.schema(self.con) def close(self):self.con.close() def settings(self,plant): row=self.con.execute('SELECT revision,value FROM planner_settings WHERE plant=?',(plant,)).fetchone() return {'revision':row['revision'],**DEFAULT_SETTINGS,**json.loads(row['value'])} if row else {'revision':0,**DEFAULT_SETTINGS} def _validate_settings(self,value): if set(value)!=set(DEFAULT_SETTINGS):raise ValueError('Unknown or missing setting') if value['family']!='auto':self.registry.get(value['family']) for key in ('autoLookbackDays','autoMinimumDays'): if type(value[key]) is not int:raise ValueError('Days must be integers') number(value['autoLookbackDays'],'lookback',7,90);number(value['autoMinimumDays'],'minimum days',1,value['autoLookbackDays']) number(value['autoMinimumCoverage'],'coverage',.5,1);number(value['autoSwitchMarginChf'],'margin',0) if value['tariffPolicy']!='published_only':raise ValueError('Only published-price policy implemented') if value['trainingCadence'] not in ('daily','weekly') or value['trainingPromotion']!='validated_only':raise ValueError('Training must use validated promotion') if value['runMode']!='shadow':raise ValueError('Shadow-only: live release requires separate validation') number(value['roundtripEfficiency'],'roundtrip efficiency',.01,1) if value['measurementPolicy'] not in ('verified_only','allow_estimates'):raise ValueError('Invalid measurement policy') if value['peakOutlookPolicy'] not in ('full_incremental','empirical_if_available'):raise ValueError('Invalid peak outlook policy') number(value['peakOutlookReliance'],'peak outlook reliance',0,1) if value['forecastSource'] not in ('legacy','corrected_profile'):raise ValueError('Unknown forecast source') if not isinstance(value['measurementDataset'],str) or len(value['measurementDataset'])>80:raise ValueError('Invalid measurement dataset') if value['forecastSource']=='corrected_profile' and not value['measurementDataset']:raise ValueError('Corrected forecast requires an explicit dataset') def _request(self,plant,revision,reason,now): row=self.con.execute('SELECT * FROM planner_work WHERE plant=?',(plant,)).fetchone() reasons=set(json.loads(row['reasons'])) if row else set();reasons.add(reason) seq=row['sequence']+1 if row else 1 self.con.execute('''INSERT INTO planner_work(plant,sequence,revision,reasons,requested_at) VALUES(?,?,?,?,?) ON CONFLICT(plant) DO UPDATE SET sequence=excluded.sequence,revision=excluded.revision,reasons=excluded.reasons,requested_at=excluded.requested_at''',(plant,seq,revision,canonical(sorted(reasons)),utc(now).isoformat())) def save_settings(self,plant,changes,expected_revision,now): if type(expected_revision) is not int or expected_revision<0:raise ValueError('Invalid expected revision') self.con.execute('BEGIN IMMEDIATE') try: current=self.settings(plant) if current.pop('revision')!=expected_revision:raise ValueError('Revision conflict; reload before saving') current.update(changes);self._validate_settings(current);revision=expected_revision+1 self.con.execute('''INSERT INTO planner_settings VALUES(?,?,?) ON CONFLICT(plant) DO UPDATE SET revision=excluded.revision,value=excluded.value''',(plant,revision,canonical(current))) self._request(plant,revision,'configuration_changed',now) self.con.execute('INSERT INTO planner_audit(plant,at,kind,detail) VALUES(?,?,?,?)',(plant,utc(now).isoformat(),'settings',canonical({'revision':revision,'changes':changes}))) self.con.commit();return {'revision':revision,**current} except Exception:self.con.rollback();raise def request(self,plant,reason,now): if reason not in ('prices_changed','telemetry_changed','five_minute_tick','manual','model_promoted','forecast_changed','operation_changed','tariffs_changed'):raise ValueError('Unknown trigger') with self.con:self._request(plant,self.settings(plant)['revision'],reason,now) def claim(self,now,lease_seconds=120): at=utc(now);self.con.execute('BEGIN IMMEDIATE') try: row=self.con.execute('SELECT * FROM planner_work WHERE lease_until IS NULL OR lease_until < ? ORDER BY requested_at LIMIT 1',(at.isoformat(),)).fetchone() if not row:self.con.commit();return None token=str(uuid4());self.con.execute('UPDATE planner_work SET lease_until=?,lease_token=? WHERE plant=?',((at+timedelta(seconds=lease_seconds)).isoformat(),token,row['plant'])) self.con.commit();return {**dict(row),'lease_token':token} except Exception:self.con.rollback();raise def finish(self,claim): self.con.execute('BEGIN IMMEDIATE') try: row=self.con.execute('SELECT sequence,lease_token FROM planner_work WHERE plant=?',(claim['plant'],)).fetchone() if not row or row['lease_token']!=claim['lease_token']:self.con.commit();return False if row['sequence']==claim['sequence']:self.con.execute('DELETE FROM planner_work WHERE plant=?',(claim['plant'],)) else:self.con.execute('UPDATE planner_work SET lease_until=NULL,lease_token=NULL WHERE plant=?',(claim['plant'],)) self.con.commit();return True except Exception:self.con.rollback();raise def initialize_peak(self,plant,month,peak_kw,source,now): number(peak_kw,'authoritative measured peak',0) if source not in ('meter_month_register','verified_month_history','verified_new_month'):raise ValueError('Configured cap is NOT measured peak') datetime.strptime(month,'%Y-%m') if month>month_key(now):raise ValueError('Future month cannot have a measured peak') with self.con: old=self.con.execute('SELECT peak_kw FROM planner_month_peaks WHERE plant=? AND month=?',(plant,month)).fetchone() if old and peak_kwreceived_at:raise ValueError('Completed aligned intervals required') with self.con: old=self.con.execute('SELECT import_kwh FROM planner_measurements WHERE plant=? AND start=?',(plant,start.isoformat())).fetchone() if old and abs(old[0]-import_kwh)>1e-9:raise ValueError('Conflicting metering fact') self.con.execute('INSERT OR IGNORE INTO planner_measurements VALUES(?,?,?)',(plant,start.isoformat(),import_kwh)) q=quarter_start(start);rows=self.con.execute('SELECT start,import_kwh FROM planner_measurements WHERE plant=? AND start>=? AND start