-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathE_manager.py
More file actions
485 lines (402 loc) · 18.6 KB
/
Copy pathE_manager.py
File metadata and controls
485 lines (402 loc) · 18.6 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
#!/usr/bin/python3.11
from time import sleep
from datetime import datetime, timedelta
import csv
import json
import sqlite3
import systemd.daemon
import threading
from apscheduler.schedulers.background import BackgroundScheduler
from Enever_tarieven import priceData
from energietarieven import prices
from shared_store import get_key, set_key
from logger_setup import get_logger
from config import (
PV_MANAGER_CSV, LP_MANAGER_CSV, WP_MANAGER_CSV,
TARIEVEN_JSON, ENERGY_JSON, ENERGY_DB, LOG_MAX_SIZE,
READER_DATA_JSON, SETPOINTS_JSON,
MAX_AMP, MIN_HEAT, SET_DAY, SET_NIGHT, MAX_COOL, MIN_COOL,
WP_MAX_SETPOINT, LP_MIN_AMP,
PV_EXPORT_THRESHOLD, PV_LIMIT_MIN, PV_LIMIT_MAX,
PV_DECREASE_FACTOR, PV_DECREASE_BASE, PV_INCREASE_STEP,
WP_LEVERING_HIGH, WP_LEVERING_MID, WP_LEVERING_LOW,
WP_ADJUST_HIGH, WP_ADJUST_MID, WP_ADJUST_LOW, WP_FC_FACTOR,
WP_SUMMER_OUTSIDE_THRESHOLD, WP_SUMMER_RELAXED_MAX,
LP_START_AMP, LP_ADJUST_STEP, LP_HIGH_AMP_THRESHOLD, LP_LOW_AMP_THRESHOLD,
LP_LEVERING_HIGH, LP_LEVERING_LOW,
MAX_LP_RANK_THRESHOLD,
READER_DATA_STALE_THRESHOLD,
TARIEVEN_RELOAD_INTERVAL, MANAGER_INTERVAL, HOURDATA_SNAPSHOT_DELAY,
MAIN_LOOP_INTERVAL, WATCHDOG_INTERVAL,
SCHEDULER_HOURDATA_MINUTE, SCHEDULER_PRICEDATA_HOUR, SCHEDULER_PRICEDATA_MINUTE,
SCHEDULER_PRICES_HOUR, SCHEDULER_PRICES_MINUTE,
)
logger = get_logger("E_manager")
MAX_SIZE = LOG_MAX_SIZE # hergebruikt voor CSV-rotatie, zelfde grens als system.log
tarieven = None
block = False
block_solar = False
max_lp = False
wp_sp = None
def watchdog_task() -> None:
while True:
systemd.daemon.notify("WATCHDOG=1")
sleep(WATCHDOG_INTERVAL)
def load_json(keys, path=READER_DATA_JSON):
data = {}
for name in keys:
items = get_key(name, path=path) or {}
for item in items:
data[item] = items[item]
return data
def rotate_log(logfile):
if logfile.exists() and logfile.stat().st_size > MAX_SIZE:
logfile.replace(logfile.with_suffix(logfile.suffix + ".1"))
def log_data(table, FILE, timestamp, values):
stamp = timestamp.split(" ")
date = stamp[0]
time = stamp[1]
write_database(table, date, time, values)
write_csv (FILE, date, time, values)
def write_database(table, date, time, values):
try:
with sqlite3.connect(ENERGY_DB) as db:
cur = db.cursor()
cols = ["date", "time"] + list(values.keys())
vals = [date, time] + list(values.values())
placeholders = ", ".join(["?"] * len(vals))
colnames = ", ".join([f'"{c}"' for c in cols])
cur.execute(f'INSERT INTO {table} ({colnames}) VALUES ({placeholders})', vals)
db.commit()
except Exception as e:
logger.error(f"E_manager database error {e}")
def write_csv (FILE, date, time, values):
rotate_log(FILE)
FIELDNAMES = ["date", "time"] + list(values.keys())
FILE.parent.mkdir(parents=True, exist_ok=True)
with open(FILE, "a", newline="") as f:
data = {"date":date, "time":time, **values}
writer = csv.DictWriter(f, fieldnames = FIELDNAMES)
writer.writerow(data)
def tarieven_loader():
global tarieven
while True:
if TARIEVEN_JSON.exists():
try:
with open(TARIEVEN_JSON) as f:
tarieven = json.load(f)
logger.info(f"Tarieven bijgewerkt: {tarieven}")
except Exception as e:
logger.error(f"Fout bij laden tarieven: {e}")
sleep(TARIEVEN_RELOAD_INTERVAL)
def PVmanager():
PVdata = {}
while True:
# 1. Data ophalen
PVdata = load_json(["p1_reader"])
if not is_reader_data_fresh("p1_reader", READER_DATA_STALE_THRESHOLD):
logger.warning("PVmanager: p1_reader-data verouderd of ontbreekt, sla deze cyclus over")
sleep(MANAGER_INTERVAL)
continue
# 2. Basisvariabelen
try:
levering = PVdata["levering"]
logger.info(f"PV levering {levering}")
limit = get_key("PV_reader", default={"limit": PV_LIMIT_MAX}, path=SETPOINTS_JSON)["limit"]
new_limit = limit
# 3. Check of prijzen negatief zijn
if block_solar:
logger.info("solar block")
if levering > PV_EXPORT_THRESHOLD:
f = PV_DECREASE_BASE + (levering - PV_EXPORT_THRESHOLD) * PV_DECREASE_FACTOR
new_limit = max(PV_LIMIT_MIN, limit - f)
elif levering <= 0:
new_limit = min(PV_LIMIT_MAX, limit + PV_INCREASE_STEP)
else:
new_limit = PV_LIMIT_MAX
# 4. Wegschrijven (alleen de limiet - production is reader-data,
# wordt door het PV-uitleesscript zelf al in reader_data.json gezet)
production = get_key("PV_reader", default={"production": None}, path=READER_DATA_JSON)["production"]
set_key("PV_reader", {"limit": new_limit}, path=SETPOINTS_JSON)
# 5. Logging
if limit != new_limit:
logger.warning(f"PV limit aangepast van {limit} naar {new_limit}")
values = {"limit":limit, "new_limit":new_limit, "production": production, "block_solar": block_solar, "export":PVdata['levering']}
log_data("PVmanager", PV_MANAGER_CSV, PVdata["timestamp"], values)
except Exception as error:
logger.error (f"PV bepalen setpoint mislukt {error}")
# 6. Pauze
sleep(MANAGER_INTERVAL)
def is_reader_data_fresh(reader_key, threshold_seconds, path=READER_DATA_JSON):
"""
Controleert of de 'timestamp' van een specifieke reader-key (bijv.
'p1_reader') niet ouder is dan threshold_seconds. Leest de key apart
op (niet via de gemergde load_json-data) om te voorkomen dat een
andere gemergde bron per ongeluk zijn eigen 'timestamp'-veld laat
winnen (zelfde soort veldnaam-botsing als eerder bij WPoutsideT/Tmin).
Retourneert False bij een ontbrekende key of onleesbare timestamp -
dat is net zo onbetrouwbaar als een verouderde timestamp.
"""
data = get_key(reader_key, path=path)
if not data or "timestamp" not in data:
return False
try:
ts = datetime.strptime(data["timestamp"], "%Y-%m-%d %H:%M:%S")
except (ValueError, TypeError):
return False
return (datetime.now() - ts).total_seconds() < threshold_seconds
def clamp_step(current, delta, limit):
"""
Beweeg 'current' met 'delta' (positief of negatief), maar nooit
verder dan 'limit'. Als 'current' door een omslag al voorbij 'limit'
ligt (bijv. na een seizoens- of drempelovergang), wordt er hoogstens
|delta| per aanroep richting 'limit' bewogen, in plaats van er in
één stap naartoe te springen.
"""
target = current + delta
if delta >= 0:
# 'limit' is een plafond
if current > limit:
return max(current - abs(delta), limit)
return min(target, limit)
else:
# 'limit' is een bodem
if current < limit:
return min(current + abs(delta), limit)
return max(target, limit)
def WPmanager():
global wp_sp
WPdata = {}
while True:
# 1. Data ophalen
# KNMI levert alleen Tmin (voorspelling, via energietarieven.py's
# getTemp()) - de live buitentemperatuur (WPoutsideT) komt uitsluitend
# van WP_reader zelf. Geen veldoverlap meer tussen de twee bronnen.
WPdata = load_json(["p1_reader", "WP_reader", "KNMI"])
if not is_reader_data_fresh("p1_reader", READER_DATA_STALE_THRESHOLD):
logger.warning("WPmanager: p1_reader-data verouderd of ontbreekt, sla deze cyclus over")
sleep(MANAGER_INTERVAL)
continue
month = datetime.now().month
uur = datetime.now().hour
try:
# 2. Wintertime settings
if month < 4 or month > 9:
fc = tarieven[str(uur).zfill(2)]["fc"] # factor om voor te bereiden op hoge stroomprijzen
NORM = max(WPdata["WPsetpoint"], SET_DAY + (WP_FC_FACTOR * fc))
if uur >= 21 or uur < 9:
# Tmin ontbreekt zolang er geen echte weersverwachting-reader is
# aangesloten (zie DATA_CONTRACT.md) - conservatieve default:
# behandel als "kans op vorst" (Tmin<=0) zodat de veiligere
# SET_DAY-tak wordt gebruikt i.p.v. de hele cyclus te laten mislukken.
NORM = (SET_NIGHT + (WP_FC_FACTOR * fc)) if WPdata.get("Tmin", 0) > 0 else (SET_DAY + (WP_FC_FACTOR * fc))
setpoint = "WPsetpoint"
MIN = MIN_HEAT
f = 1
MAX = WP_MAX_SETPOINT
NEW = "WPnew"
# 3. Summertime settings
if month > 4 and month < 9:
setpoint = "WPcoolpoint"
MIN = MAX_COOL
NORM = MAX_COOL
f = -1
outsideT = round(WPdata["WPoutsideT"] * 2) / 2
if outsideT < WP_SUMMER_OUTSIDE_THRESHOLD:
MAX = WP_SUMMER_RELAXED_MAX
else:
MAX = MIN_COOL
NEW = "WPcool"
# 4. April/September settings (overgangsmaanden: geen actieve bijsturing)
if month == 4 or month == 9:
setpoint = "WPsetpoint"
MIN = MIN_HEAT
NORM = MIN_HEAT
f = 0
MAX = WP_MAX_SETPOINT
NEW = "WPnew"
# 5. Determine newSetpoint
newSetpoint = WPdata[setpoint]
if block: # Hoge stroomprijzen, WP naar min.
newSetpoint = MIN
elif WPdata["levering"] >= WP_LEVERING_HIGH:
newSetpoint = clamp_step(newSetpoint, f * WP_ADJUST_HIGH, MAX)
elif WPdata["levering"] >= WP_LEVERING_MID:
newSetpoint = clamp_step(newSetpoint, f * WP_ADJUST_MID, MAX)
elif WPdata["levering"] < WP_LEVERING_LOW:
newSetpoint = clamp_step(WPdata[setpoint], -(f * WP_ADJUST_LOW), NORM)
# 6. Wegschrijven
set_key(NEW, {"setpoint": newSetpoint}, path=SETPOINTS_JSON)
wp_sp = newSetpoint
set_key("hourData",{"wp_sp" : wp_sp}, ENERGY_JSON)
# 7. Logging
if newSetpoint != WPdata[setpoint]:
logger.warning(f"Temperatuur aanpassen: van {WPdata[setpoint]} naar {newSetpoint}, levering={WPdata['levering']}, WPpower={WPdata['WPpower']}")
values = {"setpoint":WPdata[setpoint], "new_setpoint":newSetpoint, "WPpower": WPdata["WPpower"], "export":WPdata['levering']}
log_data("WPmanager", WP_MANAGER_CSV, WPdata["timestamp"], values)
except Exception as error:
logger.error (f"WP bepalen setpoint mislukt {error}")
# 8. Pauze
sleep(MANAGER_INTERVAL)
def LPmanager():
LPdata = {}
while True:
# 1. Data ophalen
LPdata = load_json(["p1_reader", "LP_reader"])
if not is_reader_data_fresh("p1_reader", READER_DATA_STALE_THRESHOLD):
logger.warning("LPmanager: p1_reader-data verouderd of ontbreekt, sla deze cyclus over")
sleep(MANAGER_INTERVAL)
continue
# 2. Basisvariabelen
try:
minAmp = LP_MIN_AMP
loadAmp = LPdata["LPsetpoint"]
logger.info(f"loadAmp setpoint {loadAmp}")
LP_state = get_key("laadpaal", path=SETPOINTS_JSON)["state"]
laden = get_key("laadsessie", path=READER_DATA_JSON)["laden"]
newAmp = LP_START_AMP
uur = datetime.now().hour
# 3. Pauzevenster
if block: # Niet laden als tarieven hoog zijn
newAmp = 0
if LP_state != "pause":
set_key("laadpaal", {"state": "pause"}, path=SETPOINTS_JSON)
if LPdata["laadpower"] > 0:
logger.info(f"Pauzevenster: laden gestopt (uur={uur})")
# 4. Pauze voorbij -> starten
elif LP_state == "pause":
set_key("laadpaal", {"state": "active"}, path=SETPOINTS_JSON)
newAmp = LP_START_AMP
logger.info(f"Pauze voorbij: laden gestart (amps={LPdata['amps']})")
# 5. Hoge amps -> stoppen
elif LPdata["amps"][0] > LP_HIGH_AMP_THRESHOLD and laden:
newAmp = 0
logger.warning(f"Hoge amps: laden gestopt (amps={LPdata['amps']})")
# 6. Lage amps -> herstart
elif LPdata["amps"][0] < LP_LOW_AMP_THRESHOLD and LPdata["laadpower"] == 0 and laden:
logger.warning(f"laden hervat: Amps {LPdata['amps']} laadpower {LPdata['laadpower']} laden {laden}")
newAmp = LP_START_AMP
# 7. Dynamische load
elif LPdata["laadpower"] > 0:
if max_lp: # maximaal laden bij lage stroomprijs
newAmp = MAX_AMP
elif LPdata["levering"] > LP_LEVERING_HIGH:
newAmp = min(MAX_AMP, loadAmp + LP_ADJUST_STEP)
elif LPdata["levering"] >= LP_LEVERING_LOW:
newAmp = max(minAmp, loadAmp)
elif LPdata["levering"] < LP_LEVERING_LOW:
newAmp = max(minAmp, loadAmp - LP_ADJUST_STEP)
# 8. Begrenzen + wegschrijven
if newAmp != 0:
newAmp = max(minAmp, min(MAX_AMP, newAmp))
set_key("LPnew", {"setpoint": newAmp}, path=SETPOINTS_JSON)
# 9. Logging
if laden and loadAmp != newAmp:
logger.warning(f"Amperage aanpassen: van {loadAmp} naar {newAmp}, levering={LPdata['levering']}, laadpower={LPdata['laadpower']}")
values = {"amp":loadAmp, "new_amp":newAmp, "laadpower": LPdata["laadpower"], "export":LPdata['levering']}
log_data("LPmanager", LP_MANAGER_CSV, LPdata["timestamp"], values)
except Exception as error:
logger.error (f"LP bepalen setpoint mislukt {error}")
# 10. Pauze
sleep(MANAGER_INTERVAL)
def hourData():
global block, block_solar, max_lp, wp_sp
if tarieven is None:
logger.warning("hourData: tarieven nog niet geladen, sla deze run over")
return
uur = datetime.now().hour
try:
price = tarieven[str(uur).zfill(2)]["prijs"]
price_sell = tarieven[str(uur).zfill(2)]["prijs_terug"]
block = False
block_solar = False
max_lp = False
if tarieven[str(uur).zfill(2)]["block"]:
block = True
if tarieven[str(uur).zfill(2)]["block_solar"]:
block_solar = True
rank = tarieven[str(uur).zfill(2)]["rank"]
if block_solar or rank <= MAX_LP_RANK_THRESHOLD:
max_lp = True
set_key("hourData",{"price" : price, "price_sell": price_sell, "block" : block, "block_pv" : block_solar, "max_lp" : max_lp, "wp_sp" : wp_sp}, ENERGY_JSON)
except Exception as error:
logger.error(f"hourData mislukt: {error}")
return
sleep(HOURDATA_SNAPSHOT_DELAY)
store_hourly_snapshot(price, price_sell)
def store_hourly_snapshot(prijsverbruik, prijslevering):
try:
db = sqlite3.connect(ENERGY_DB)
db.row_factory = sqlite3.Row
cur = db.cursor()
ts = get_key("p1_reader", path=READER_DATA_JSON)["timestamp"]
dt = datetime.strptime(ts, "%Y-%m-%d %H:%M:%S")
hour = dt.strftime("%Y-%m-%d %H")
cur.execute("""UPDATE HourlyMeter SET prijsverbruik = ?, prijslevering = ? WHERE hour = ?""", (prijsverbruik, prijslevering, hour))
db.commit()
hour = {}
data = {}
stand_V = {}
stand_L = {}
verbruik = {}
levering = {}
# Ophalen van meterstanden
for i in ["0", "1", "24"]:
dt_i = dt - timedelta(hours=int(i))
hour[i] = dt_i.strftime("%Y-%m-%d %H")
cur.execute("SELECT * FROM HourlyMeter WHERE hour = ?", (hour[i],))
row = cur.fetchone()
if row is None:
logger.error(f"Hourly snapshot: geen data voor uur {hour[i]}")
else:
data[i] = dict(row)
stand_V[i] = data[i]["verbruik1"] + data[i]["verbruik2"]
stand_L[i] = data[i]["levering1"] + data[i]["levering2"]
# Berekenen verbruik laatste uur en laatste 24 uur
if i != "0" and i in data:
verbruik[i] = round(stand_V["0"] - stand_V[i], 1)
levering[i] = round(stand_L["0"] - stand_L[i], 1)
if i == "1" and i in data:
verbruik_EUR = verbruik[i] * float(data["0"]["prijsverbruik"])
levering_EUR = levering[i] * float(data["0"]["prijslevering"])
kosten = round(verbruik_EUR + levering_EUR, 2)
cur.execute("""UPDATE HourlyMeter SET verbruik = ?, levering = ?, kosten = ? WHERE hour = ?""", (verbruik[i], levering[i], kosten, hour["1"]))
set_key("cost", {i:{"use": verbruik[i], "sell": levering[i], "cost":kosten}},ENERGY_JSON)
db.commit()
if i == "24" and i in data:
cur.execute("SELECT kosten FROM HourlyMeter WHERE hour >= ?", (hour[i],))
rows = cur.fetchall()
kosten_24 = round(sum(r["kosten"] or 0 for r in rows), 2)
cur.execute("""UPDATE HourlyMeter SET verbruik_24 = ?, levering_24 = ?, kosten_24 = ? WHERE hour = ?""", (verbruik[i], levering[i], kosten_24, hour["1"]))
set_key("cost", {i:{"use": verbruik[i], "sell": levering[i], "cost":kosten_24}},ENERGY_JSON)
db.commit()
db.close()
except Exception as e:
logger.error(f"Hourly snapshot error: {e}")
def main():
systemd.daemon.notify("READY=1")
threading.Thread(target= tarieven_loader, daemon=True).start()
threading.Thread(target= watchdog_task, daemon=True).start()
threading.Thread(target= PVmanager, daemon=True).start()
threading.Thread(target= WPmanager, daemon=True).start()
threading.Thread(target= LPmanager, daemon=True).start()
# Wachten tot de eerste tarieven geladen zijn, zodat hourData() bij
# een koude start niet direct op tarieven=None crasht.
wait_ticks = 0
while tarieven is None and wait_ticks < 60:
sleep(1)
wait_ticks += 1
if tarieven is None:
logger.warning("main: tarieven.json nog niet gevonden na 60s, hourData start alsnog (guard vangt dit af)")
# start scheduler
scheduler = BackgroundScheduler()
scheduler.add_job(hourData, 'cron', minute=SCHEDULER_HOURDATA_MINUTE)
scheduler.add_job(priceData, 'cron', hour=SCHEDULER_PRICEDATA_HOUR, minute=SCHEDULER_PRICEDATA_MINUTE)
scheduler.add_job(prices, 'cron', hour=SCHEDULER_PRICES_HOUR, minute=SCHEDULER_PRICES_MINUTE)
scheduler.start()
hourData() # Direct herladen bij start
# hoofdthread blijft leven
while True:
sleep(MAIN_LOOP_INTERVAL)
if __name__ == "__main__":
main()