Skip to content

Commit 9f0bb8e

Browse files
fix(fetch): publish rate_import/rate_export atomically to remove a read race
update_pred rebuilds rate_import/rate_export on an AppDaemon worker thread while async components (e.g. the Teslemetry tariff builder) read them, with no lock. The build reset them to {}, reassigned per step, and mutated them in place (load_saving_slot/load_free_slot/load_axle_slot), so a concurrent reader could see an empty or half-merged dict. Build each side into a local (import_rates/export_rates) through the whole fetch + replicate + slot pipeline and assign self.rate_import/self.rate_export exactly once at the end, so a reader always sees either the previous complete dict or the new complete one - never an intermediate. The three in-place slot mutators now take the working dict as a parameter instead of touching self.rate_*. Behaviour-preserving: derived attrs (rate_import_base/_replicated/_no_io, rate_export_base/_replicated) are unchanged, and the full --quick suite (octopus rates, saving/free/axle slots, overrides, manual rates, IO) passes. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
1 parent 263c5df commit 9f0bb8e

6 files changed

Lines changed: 73 additions & 74 deletions

File tree

apps/predbat/axle.py

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -647,9 +647,9 @@ def fetch_axle_sessions(base):
647647
return axle_events_deduplicated
648648

649649

650-
def load_axle_slot(base, axle_sessions, export, rate_replicate=None):
650+
def load_axle_slot(base, axle_sessions, rate_dict, export, rate_replicate=None):
651651
"""
652-
Load Axle VPP session slot
652+
Load Axle VPP session slot into rate_dict (in place)
653653
"""
654654
if rate_replicate is None:
655655
rate_replicate = {}
@@ -679,11 +679,11 @@ def load_axle_slot(base, axle_sessions, export, rate_replicate=None):
679679
base.log("Setting Axle VPP session in range {} - {} export {} pence_per_kwh {}".format(base.time_abs_str(start_minutes), base.time_abs_str(end_minutes), export, pence_per_kwh))
680680
for minute in range(start_minutes, end_minutes):
681681
if export:
682-
base.rate_export[minute] = base.rate_export.get(minute, 0) + pence_per_kwh
682+
rate_dict[minute] = rate_dict.get(minute, 0) + pence_per_kwh
683683
base.load_scaling_dynamic[minute] = base.load_scaling_saving
684684
rate_replicate[minute] = "saving"
685685
else:
686-
base.rate_import[minute] = base.rate_import.get(minute, 0) - pence_per_kwh
686+
rate_dict[minute] = rate_dict.get(minute, 0) - pence_per_kwh
687687
base.load_scaling_dynamic[minute] = base.load_scaling_free
688688
rate_replicate[minute] = "saving"
689689

apps/predbat/fetch.py

Lines changed: 45 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -708,9 +708,9 @@ def fetch_sensor_data(self, save=True):
708708
prev_octopus_free_slots = self.octopus_free_slots.copy()
709709
prev_axle_sessions = self.axle_sessions.copy()
710710

711-
self.rate_import = {}
711+
import_rates = {}
712712
self.rate_import_replicated = {}
713-
self.rate_export = {}
713+
export_rates = {}
714714
self.rate_export_replicated = {}
715715
self.rate_slots = []
716716
self.io_adjusted = {}
@@ -849,36 +849,36 @@ def fetch_sensor_data(self, save=True):
849849
# Fixed URL for rate import
850850
self.log("Downloading import rates directly from URL {}".format(self.get_arg("rates_import_octopus_url", indirect=False)))
851851
# Need to take a copy, as saving sessions will repeatedly increment cached import rates
852-
self.rate_import = copy.deepcopy(self.download_octopus_rates(self.get_arg("rates_import_octopus_url", indirect=False)))
852+
import_rates = copy.deepcopy(self.download_octopus_rates(self.get_arg("rates_import_octopus_url", indirect=False)))
853853
elif "metric_octopus_import" in self.args:
854854
# Octopus import rates
855855
entity_id = self.get_arg("metric_octopus_import", None, indirect=False)
856-
self.rate_import = self.fetch_octopus_rates(entity_id, adjust_key="is_intelligent_adjusted")
857-
if not self.rate_import:
856+
import_rates = self.fetch_octopus_rates(entity_id, adjust_key="is_intelligent_adjusted")
857+
if not import_rates:
858858
self.log("Error: metric_octopus_import is not set correctly in apps.yaml, or no energy rates can be read")
859859
self.record_status(message="Error: metric_octopus_import not set correctly in apps.yaml, or no energy rates can be read", had_errors=True)
860860
elif "metric_energidataservice_import" in self.args:
861861
# Energi Data Service import rates
862862
entity_id = self.get_arg("metric_energidataservice_import", None, indirect=False)
863-
self.rate_import = self.fetch_energidataservice_rates(entity_id, adjust_key="is_intelligent_adjusted")
864-
if not self.rate_import:
863+
import_rates = self.fetch_energidataservice_rates(entity_id, adjust_key="is_intelligent_adjusted")
864+
if not import_rates:
865865
self.log("Error: metric_energidataservice_import is not set correctly in apps.yaml, or no energy rates can be read")
866866
self.record_status(message="Error: metric_energidataservice_import not set correctly in apps.yaml, or no energy rates can be read", had_errors=True)
867867
elif "metric_stromligning_import_today" in self.args or "metric_stromligning_import_tomorrow" in self.args:
868868
# Strømligning import rates
869869
entity_id_today = self.get_arg("metric_stromligning_import_today", None, indirect=False)
870870
entity_id_tomorrow = self.get_arg("metric_stromligning_import_tomorrow", None, indirect=False)
871-
self.rate_import = self.fetch_stromligning_rates(entity_id_today, entity_id_tomorrow, adjust_key="is_intelligent_adjusted")
872-
if not self.rate_import:
871+
import_rates = self.fetch_stromligning_rates(entity_id_today, entity_id_tomorrow, adjust_key="is_intelligent_adjusted")
872+
if not import_rates:
873873
self.log("Error: metric_stromligning_import sensors are not set correctly or no energy rates can be read")
874874
self.record_status(message="Error: metric_stromligning_import sensors not set correctly or no energy rates can be read", had_errors=True)
875875

876876
# Fallback if no other rate types are set
877-
if not self.rate_import:
877+
if not import_rates:
878878
# Basic rates defined by user over time
879879
rate_import_dict = self.get_arg("rates_import", [], indirect=False)
880880
if rate_import_dict:
881-
self.rate_import = self.basic_rates(rate_import_dict, "rates_import")
881+
import_rates = self.basic_rates(rate_import_dict, "rates_import")
882882

883883
# Gas rates if set
884884
if "metric_octopus_gas" in self.args:
@@ -925,36 +925,36 @@ def fetch_sensor_data(self, save=True):
925925
# Fixed URL for rate export
926926
self.log("Downloading export rates directly from URL {}".format(self.get_arg("rates_export_octopus_url", indirect=False)))
927927
# Need to take a copy, as saving sessions will repeatedly increment cached export rates
928-
self.rate_export = copy.deepcopy(self.download_octopus_rates(self.get_arg("rates_export_octopus_url", indirect=False)))
928+
export_rates = copy.deepcopy(self.download_octopus_rates(self.get_arg("rates_export_octopus_url", indirect=False)))
929929
elif "metric_octopus_export" in self.args:
930930
# Octopus export rates
931931
entity_id = self.get_arg("metric_octopus_export", None, indirect=False)
932-
self.rate_export = self.fetch_octopus_rates(entity_id)
933-
if not self.rate_export:
932+
export_rates = self.fetch_octopus_rates(entity_id)
933+
if not export_rates:
934934
self.log("Warning: metric_octopus_export is not set correctly in apps.yaml, or no energy rates can be read")
935935
self.record_status(message="Error: metric_octopus_export not set correctly in apps.yaml, or no energy rates can be read", had_errors=True)
936936
elif "metric_energidataservice_export" in self.args:
937937
# Energi Data Service export rates
938938
entity_id = self.get_arg("metric_energidataservice_export", None, indirect=False)
939-
self.rate_export = self.fetch_energidataservice_rates(entity_id)
940-
if not self.rate_export:
939+
export_rates = self.fetch_energidataservice_rates(entity_id)
940+
if not export_rates:
941941
self.log("Warning: metric_energidataservice_export is not set correctly in apps.yaml, or no energy rates can be read")
942942
self.record_status(message="Error: metric_energidataservice_export not set correctly in apps.yaml, or no energy rates can be read", had_errors=True)
943943
elif "metric_stromligning_export_today" in self.args or "metric_stromligning_export_tomorrow" in self.args:
944944
# Strømligning export rates
945945
entity_id_today = self.get_arg("metric_stromligning_export_today", None, indirect=False)
946946
entity_id_tomorrow = self.get_arg("metric_stromligning_export_tomorrow", None, indirect=False)
947-
self.rate_export = self.fetch_stromligning_rates(entity_id_today, entity_id_tomorrow)
948-
if not self.rate_export:
947+
export_rates = self.fetch_stromligning_rates(entity_id_today, entity_id_tomorrow)
948+
if not export_rates:
949949
self.log("Warning: metric_stromligning_export sensors are not set correctly or no energy rates can be read")
950950
self.record_status(message="Error: metric_stromligning_export sensors not set correctly or no energy rates can be read", had_errors=True)
951951

952952
# Fallback if no other rate types are set
953-
if not self.rate_export:
953+
if not export_rates:
954954
# Basic rates defined by user over time
955955
rate_export_dict = self.get_arg("rates_export", [], indirect=False)
956956
# Allow all zero export rates, as some users have a feed-in tariff that is zero
957-
self.rate_export = self.basic_rates(rate_export_dict, "rates_export")
957+
export_rates = self.basic_rates(rate_export_dict, "rates_export")
958958

959959
# Fetch Axle sessions first so Octopus auto-join can skip saving sessions that overlap an Axle VPP session
960960
self.axle_sessions = fetch_axle_sessions(self)
@@ -967,44 +967,47 @@ def fetch_sensor_data(self, save=True):
967967

968968
# futurerate data
969969
futurerate = FutureRate(self)
970-
self.future_energy_rates_import, self.future_energy_rates_export = futurerate.futurerate_analysis(self.rate_import, self.rate_export)
970+
self.future_energy_rates_import, self.future_energy_rates_export = futurerate.futurerate_analysis(import_rates, export_rates)
971971

972972
# Replicate and scan import rates
973-
if self.rate_import:
974-
self.rate_scan(self.rate_import, print=False)
973+
if import_rates:
974+
self.rate_scan(import_rates, print=False)
975975
self.rate_max_base = self.rate_max # True peak rate before saving sessions / overrides inflate it
976976
self.rate_min_base = self.rate_min # True off-peak rate before free sessions / overrides deflate it
977-
self.rate_import_base, _ = self.rate_replicate(self.rate_import.copy(), {}, is_import=True) # True import rates, gap-filled but without IO/saving/override distortion
978-
self.rate_import, self.rate_import_replicated = self.rate_replicate(self.rate_import, self.io_adjusted, is_import=True)
979-
self.rate_import_no_io = self.rate_import.copy()
977+
self.rate_import_base, _ = self.rate_replicate(import_rates.copy(), {}, is_import=True) # True import rates, gap-filled but without IO/saving/override distortion
978+
import_rates, self.rate_import_replicated = self.rate_replicate(import_rates, self.io_adjusted, is_import=True)
979+
self.rate_import_no_io = import_rates.copy()
980980
for car_n in range(self.num_cars):
981-
self.rate_import = self.rate_add_io_slots(car_n, self.rate_import, self.octopus_slots[car_n])
982-
self.load_saving_slot(self.octopus_saving_slots, export=False, rate_replicate=self.rate_import_replicated)
983-
self.load_free_slot(self.octopus_free_slots, export=False, rate_replicate=self.rate_import_replicated)
984-
load_axle_slot(self, self.axle_sessions, export=False, rate_replicate=self.rate_import_replicated)
985-
self.rate_import = self.basic_rates(self.get_arg("rates_import_override", [], indirect=False), "rates_import_override", self.rate_import, self.rate_import_replicated)
986-
self.rate_import = self.apply_manual_rates(self.rate_import, self.manual_import_rates, is_import=True, rate_replicate=self.rate_import_replicated)
987-
self.rate_scan(self.rate_import, print=True)
981+
import_rates = self.rate_add_io_slots(car_n, import_rates, self.octopus_slots[car_n])
982+
self.load_saving_slot(self.octopus_saving_slots, import_rates, export=False, rate_replicate=self.rate_import_replicated)
983+
self.load_free_slot(self.octopus_free_slots, import_rates, export=False, rate_replicate=self.rate_import_replicated)
984+
load_axle_slot(self, self.axle_sessions, import_rates, export=False, rate_replicate=self.rate_import_replicated)
985+
import_rates = self.basic_rates(self.get_arg("rates_import_override", [], indirect=False), "rates_import_override", import_rates, self.rate_import_replicated)
986+
import_rates = self.apply_manual_rates(import_rates, self.manual_import_rates, is_import=True, rate_replicate=self.rate_import_replicated)
987+
self.rate_scan(import_rates, print=True)
988988
else:
989989
self.rate_import_no_io = {}
990990
self.log("Warning: No import rate data provided")
991991
self.record_status(message="Error: No import rate data provided", had_errors=True)
992+
# Atomic publish: readers (e.g. async components) never see a half-built or empty rate_import.
993+
self.rate_import = import_rates
992994

993995
# Replicate and scan export rates
994-
if self.rate_export:
995-
self.rate_scan_export(self.rate_export, print=False)
996-
self.rate_export, self.rate_export_replicated = self.rate_replicate(self.rate_export, is_import=False)
997-
self.rate_export_base = self.rate_export.copy()
996+
if export_rates:
997+
self.rate_scan_export(export_rates, print=False)
998+
export_rates, self.rate_export_replicated = self.rate_replicate(export_rates, is_import=False)
999+
self.rate_export_base = export_rates.copy()
9981000
# For export tariff only load the saving session if enabled
9991001
if self.rate_export_max > 0:
1000-
self.load_saving_slot(self.octopus_saving_slots, export=True, rate_replicate=self.rate_export_replicated)
1001-
load_axle_slot(self, self.axle_sessions, export=True, rate_replicate=self.rate_export_replicated)
1002-
self.rate_export = self.basic_rates(self.get_arg("rates_export_override", [], indirect=False), "rates_export_override", self.rate_export, self.rate_export_replicated)
1003-
self.rate_export = self.apply_manual_rates(self.rate_export, self.manual_export_rates, is_import=False, rate_replicate=self.rate_export_replicated)
1004-
self.rate_scan_export(self.rate_export, print=True)
1002+
self.load_saving_slot(self.octopus_saving_slots, export_rates, export=True, rate_replicate=self.rate_export_replicated)
1003+
load_axle_slot(self, self.axle_sessions, export_rates, export=True, rate_replicate=self.rate_export_replicated)
1004+
export_rates = self.basic_rates(self.get_arg("rates_export_override", [], indirect=False), "rates_export_override", export_rates, self.rate_export_replicated)
1005+
export_rates = self.apply_manual_rates(export_rates, self.manual_export_rates, is_import=False, rate_replicate=self.rate_export_replicated)
1006+
self.rate_scan_export(export_rates, print=True)
10051007
else:
10061008
self.log("Warning: No export rate data provided")
10071009
self.record_status(message="Error: No export rate data provided", had_errors=True)
1010+
self.rate_export = export_rates
10081011

10091012
# Set rate thresholds
10101013
if self.rate_import or self.rate_export:

apps/predbat/octopus.py

Lines changed: 10 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -2405,9 +2405,9 @@ def _parse_slot_time(self, value):
24052405
except (ValueError, TypeError):
24062406
return None
24072407

2408-
def load_free_slot(self, octopus_free_slots, export=False, rate_replicate=None):
2408+
def load_free_slot(self, octopus_free_slots, rate_dict, export=False, rate_replicate=None):
24092409
"""
2410-
Load octopus free session slot
2410+
Load octopus free session slot into rate_dict (in place)
24112411
"""
24122412
if rate_replicate is None:
24132413
rate_replicate = {}
@@ -2436,15 +2436,15 @@ def load_free_slot(self, octopus_free_slots, export=False, rate_replicate=None):
24362436
self.log("Setting Octopus free session in range {} - {} export {} rate {}".format(self.time_abs_str(start_minutes), self.time_abs_str(end_minutes), export, rate))
24372437
for minute in range(start_minutes, end_minutes):
24382438
if export:
2439-
self.rate_export[minute] = rate
2439+
rate_dict[minute] = rate
24402440
else:
2441-
self.rate_import[minute] = min(rate, self.rate_import[minute])
2441+
rate_dict[minute] = min(rate, rate_dict[minute])
24422442
self.load_scaling_dynamic[minute] = self.load_scaling_free
24432443
rate_replicate[minute] = "saving"
24442444

2445-
def load_saving_slot(self, octopus_saving_slots, export=False, rate_replicate=None):
2445+
def load_saving_slot(self, octopus_saving_slots, rate_dict, export=False, rate_replicate=None):
24462446
"""
2447-
Load octopus saving session slot
2447+
Load octopus saving session slot into rate_dict (in place)
24482448
"""
24492449
if rate_replicate is None:
24502450
rate_replicate = {}
@@ -2476,15 +2476,11 @@ def load_saving_slot(self, octopus_saving_slots, export=False, rate_replicate=No
24762476
if start_minutes < (self.forecast_minutes + self.minutes_now):
24772477
self.log("Octopus: Setting Octopus saving session in range {} - {} export {} rate {}".format(self.time_abs_str(start_minutes), self.time_abs_str(end_minutes), export, rate))
24782478
for minute in range(start_minutes, end_minutes):
2479-
if export:
2480-
if minute in self.rate_export:
2481-
self.rate_export[minute] += rate
2482-
rate_replicate[minute] = "saving"
2483-
else:
2484-
if minute in self.rate_import:
2485-
self.rate_import[minute] += rate
2479+
if minute in rate_dict:
2480+
rate_dict[minute] += rate
2481+
rate_replicate[minute] = "saving"
2482+
if not export:
24862483
self.load_scaling_dynamic[minute] = self.load_scaling_saving
2487-
rate_replicate[minute] = "saving"
24882484

24892485
def decode_octopus_slot(self, car_n, slot, raw=False):
24902486
"""

apps/predbat/tests/test_axle.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1072,7 +1072,7 @@ def get_arg(self, name, indirect=True):
10721072

10731073
# Load the Axle export slot
10741074
rate_replicate = {}
1075-
load_axle_slot(base, axle_sessions, export=True, rate_replicate=rate_replicate)
1075+
load_axle_slot(base, axle_sessions, base.rate_export, export=True, rate_replicate=rate_replicate)
10761076

10771077
# Verify rates were increased by 100p/kWh during the event period
10781078
for minute in range(start_minutes, end_minutes):
@@ -1155,7 +1155,7 @@ def get_arg(self, name, indirect=True):
11551155
original_rate = base.rate_import[start_minutes]
11561156

11571157
rate_replicate = {}
1158-
load_axle_slot(base, axle_sessions, export=False, rate_replicate=rate_replicate)
1158+
load_axle_slot(base, axle_sessions, base.rate_import, export=False, rate_replicate=rate_replicate)
11591159

11601160
# Verify import rates were decreased by pence_per_kwh during the event period
11611161
for minute in range(start_minutes, end_minutes):

0 commit comments

Comments
 (0)