From 7ebe6d714be5a8e265c0b5e411de4ed609391f28 Mon Sep 17 00:00:00 2001 From: dr-dimitri <87113560+dr-dimitri@users.noreply.github.com> Date: Wed, 16 Sep 2026 12:22:39 +0200 Subject: [PATCH] fix: detect recent Energy-Charts source cadence (#1315) Infer interval coverage from consistent recent original price spacings while preserving history and forecasting resolution policies. Reuse source data until a successful refresh and cover cadence transitions, fallback, and history repair. Fixes #1279 --- .../prediction/elecpriceenergycharts.py | 66 ++++-- tests/test_elecpriceenergycharts.py | 196 +++++++++++++++++- 2 files changed, 237 insertions(+), 25 deletions(-) diff --git a/src/akkudoktoreos/prediction/elecpriceenergycharts.py b/src/akkudoktoreos/prediction/elecpriceenergycharts.py index bc842314..c065edd3 100644 --- a/src/akkudoktoreos/prediction/elecpriceenergycharts.py +++ b/src/akkudoktoreos/prediction/elecpriceenergycharts.py @@ -98,6 +98,35 @@ class ElecPriceEnergyCharts(ElecPriceProvider): """Return the unique identifier for the Energy-Charts provider.""" return "ElecPriceEnergyCharts" + def _coverage_resolution_seconds(self, source_series: pd.Series) -> int: + """Infer interval coverage from original prices without changing forecast resolution. + + The caller excludes timestamps beyond ``highest_orig_datetime``. Within the last + 24 hours, four equal spacings among the final five differences establish a recent + cadence. This recognizes five consecutive points at a new resolution while tolerating + one exceptional gap. Only positive intervals dividing one hour are supported, as in + the shared resolution helper. + + Ambiguous or insufficient agreement falls back to the median resolution of these + 24 hours; fewer than two distinct timestamps fall back to hourly coverage. + """ + if source_series.empty: + return 3600 + + recent_series = source_series.sort_index() + recent_series = recent_series[ + recent_series.index >= recent_series.index[-1] - pd.Timedelta(hours=24) + ] + index = pd.DatetimeIndex(recent_series.index).drop_duplicates() + deltas = index.to_series().diff().dropna().dt.total_seconds().tail(5) + counts = deltas.value_counts() + if not counts.empty and counts.iloc[0] >= 4: + resolution = float(counts.index[0]) + if resolution > 0 and 3600 % resolution == 0: + return int(resolution) + + return self._resolution_seconds(recent_series) + def _has_complete_published_horizon( self, *, now: pd.Timestamp, resolution_seconds: int ) -> bool: @@ -247,12 +276,21 @@ class ElecPriceEnergyCharts(ElecPriceProvider): # Determine if update is needed and what start date is really necessary needs_update = False + raw_series: Optional[pd.Series] = None if self.highest_orig_datetime: - raw_history = await self.key_to_raw_series( + source_end = to_datetime(self.highest_orig_datetime).add(seconds=1) + raw_series = await self.key_to_raw_series( key="elecprice_marketprice_raw_wh", - start_datetime=start_datetime, - end_datetime=gross_start_datetime, + end_datetime=max(gross_start_datetime, source_end), ) + # Preserve the history window even during an outage when the latest original + # timestamp precedes EMS start. Reuse this read for coverage and forecasting, + # but exclude the extrapolated tail from both resolution estimates. + raw_history = raw_series[ + (raw_series.index >= pd.Timestamp(start_datetime)) + & (raw_series.index < pd.Timestamp(gross_start_datetime)) + ] + raw_series = raw_series[raw_series.index <= pd.Timestamp(self.highest_orig_datetime)] if raw_history.empty: # We need the default start date (35 days in past) @@ -271,16 +309,7 @@ class ElecPriceEnergyCharts(ElecPriceProvider): # Use default start date in case of forced update needs_update = True else: - # The latest source data may have a different resolution than - # the history before ems_start_datetime. Use its final 24 hours - # so older, finer intervals cannot dominate the median, and - # exclude the predicted tail. - source_series = await self.key_to_raw_series( - key="elecprice_marketprice_raw_wh", - start_datetime=to_datetime(self.highest_orig_datetime).subtract(hours=24), - end_datetime=to_datetime(self.highest_orig_datetime).add(seconds=1), - ) - source_resolution_seconds = self._resolution_seconds(source_series) + source_resolution_seconds = self._coverage_resolution_seconds(raw_series) if not self._has_complete_published_horizon( now=now, resolution_seconds=source_resolution_seconds ): @@ -311,6 +340,8 @@ class ElecPriceEnergyCharts(ElecPriceProvider): raise ValueError("No Energy-Charts electricity price data available") self.highest_orig_datetime = to_datetime(series_data.index.max()) await self.key_from_series("elecprice_marketprice_raw_wh", series_data) + # Reload after a successful fetch so prediction sees the new source data. + raw_series = None # Newly fetched data widens the window that needs its gross # (fee-inclusive) values recomputed. gross_start_datetime = to_datetime(series_data.index.min()) @@ -335,10 +366,11 @@ class ElecPriceEnergyCharts(ElecPriceProvider): logger.error(error_msg) raise ValueError(error_msg) - raw_series = await self.key_to_raw_series( - key="elecprice_marketprice_raw_wh", - end_datetime=to_datetime(self.highest_orig_datetime).add(seconds=1), - ) + if raw_series is None: + raw_series = await self.key_to_raw_series( + key="elecprice_marketprice_raw_wh", + end_datetime=to_datetime(self.highest_orig_datetime).add(seconds=1), + ) resolution_seconds = self._resolution_seconds(raw_series) slots_per_hour = 3600 // resolution_seconds diff --git a/tests/test_elecpriceenergycharts.py b/tests/test_elecpriceenergycharts.py index 7b8c63d6..77f47f4e 100644 --- a/tests/test_elecpriceenergycharts.py +++ b/tests/test_elecpriceenergycharts.py @@ -157,9 +157,66 @@ class TestElecPriceEnergyCharts: ) assert len(np_price_array) == provider.total_hours + @pytest.mark.parametrize( + ("old_interval_minutes", "recent_intervals_minutes", "expected_seconds"), + [ + pytest.param(60, [15, 15, 15, 15], 900, id="hourly-to-quarter-hourly"), + pytest.param(15, [60, 60, 60, 60], 3600, id="quarter-hourly-to-hourly"), + pytest.param(60, [15, 15, 60, 15, 15], 900, id="quarter-hourly-with-gap"), + pytest.param(15, [60, 60, 120, 60, 60], 3600, id="hourly-with-gap"), + pytest.param(60, [15, 15, 15, 15, 60], 900, id="quarter-hourly-with-last-gap"), + pytest.param(15, [60, 60, 60, 60, 15], 3600, id="hourly-with-last-outlier"), + pytest.param(60, [15, 15, 15], 3600, id="insufficient-quarter-hourly-run"), + pytest.param(15, [60, 60, 60], 900, id="insufficient-hourly-run"), + pytest.param(60, [15, 60, 15, 60, 15], 3600, id="ambiguous-latest-intervals"), + ], + ) + def test_coverage_resolution_recognizes_recent_cadence( + self, + provider: ElecPriceEnergyCharts, + old_interval_minutes: int, + recent_intervals_minutes: list[int], + expected_seconds: int, + ) -> None: + """A supported recent cadence survives one outlier; ambiguous samples use the median.""" + transition = pd.Timestamp("2026-01-15 18:00:00", tz="Europe/Berlin") + older_index = pd.date_range( + start=transition - pd.Timedelta(days=1), + end=transition, + freq=f"{old_interval_minutes}min", + ) + recent_index = transition + pd.to_timedelta(np.cumsum(recent_intervals_minutes), unit="min") + source = pd.Series(0.0001, index=older_index.append(recent_index)) + + assert provider._coverage_resolution_seconds(source) == expected_seconds + + @pytest.mark.parametrize( + ("offsets_minutes", "expected_seconds"), + [ + pytest.param([], 3600, id="empty"), + pytest.param([0], 3600, id="one-timestamp"), + pytest.param([0, 0], 3600, id="duplicate-only"), + pytest.param([0, 15], 900, id="two-timestamps"), + pytest.param([60, 15, 0, 45, 30, 30], 900, id="unordered-and-duplicated"), + pytest.param([0, 45, 90, 135, 180], 3600, id="unsupported-interval"), + ], + ) + def test_coverage_resolution_handles_sparse_or_irregular_source( + self, + provider: ElecPriceEnergyCharts, + offsets_minutes: list[int], + expected_seconds: int, + ) -> None: + """Normalize source ordering and retain the shared helper's sparse-data fallback.""" + start = pd.Timestamp("2026-01-15 00:00:00", tz="Europe/Berlin") + source = pd.Series(0.0001, index=start + pd.to_timedelta(offsets_minutes, unit="min")) + + assert provider._coverage_resolution_seconds(source) == expected_seconds + @pytest.mark.asyncio @pytest.mark.parametrize("host_timezone", ["UTC", "Europe/Berlin"]) @pytest.mark.parametrize("history_interval_minutes", [15, 60]) + @pytest.mark.parametrize("recent_source_points", [None, 5]) @pytest.mark.parametrize( ("now", "last_price", "interval_minutes", "needs_update"), [ @@ -182,6 +239,7 @@ class TestElecPriceEnergyCharts: set_other_timezone: Callable[[str], str], host_timezone: str, history_interval_minutes: int, + recent_source_points: int | None, now: str, last_price: str, interval_minutes: int, @@ -198,16 +256,19 @@ class TestElecPriceEnergyCharts: pd.Timestamp(last_price, tz="Europe/Berlin"), in_timezone="Europe/Berlin" ) get_ems().set_start_datetime(start) - history_index = pd.date_range( - start=start.subtract(days=35), - end=start, - freq=f"{history_interval_minutes}min", - inclusive="left", + source_start = ( + start + if recent_source_points is None + else last_original.subtract(minutes=(recent_source_points - 1) * interval_minutes) ) source_index = pd.date_range( - start=start, - end=last_original, - freq=f"{interval_minutes}min", + start=source_start, end=last_original, freq=f"{interval_minutes}min" + ) + history_index = pd.date_range( + start=start.subtract(days=35), + end=source_index[0], + freq=f"{history_interval_minutes}min", + inclusive="left", ) await provider.key_from_series( "elecprice_marketprice_raw_wh", @@ -271,6 +332,125 @@ class TestElecPriceEnergyCharts: request.assert_not_called() assert provider.highest_orig_datetime == last_original + @pytest.mark.asyncio + @pytest.mark.parametrize( + ( + "history_days", + "history_end_days_ago", + "force_update", + "complete_horizon", + "fetch_fails", + "needs_update", + ), + [ + pytest.param(35, 0, False, True, False, False, id="unchanged"), + pytest.param(35, 0, False, False, False, True, id="successful-refresh"), + pytest.param(35, 0, False, False, True, True, id="failed-refresh"), + pytest.param(35, 0, True, True, False, True, id="forced-refresh"), + pytest.param(0, 0, False, True, False, True, id="missing-history"), + pytest.param(2, 0, False, True, False, True, id="insufficient-history"), + pytest.param(14, 0, False, True, False, True, id="history-threshold"), + pytest.param(15, 0, False, True, False, False, id="sufficient-history"), + pytest.param(35, 35, False, True, False, True, id="history-outside-window"), + ], + ) + async def test_update_data_reuses_source_snapshot_until_successful_fetch( + self, + provider: ElecPriceEnergyCharts, + set_other_timezone: Callable[[str], str], + history_days: int, + history_end_days_ago: int, + force_update: bool, + complete_horizon: bool, + fetch_fails: bool, + needs_update: bool, + ) -> None: + """Retain history repair and fallback while rereading source data only after a fetch.""" + set_other_timezone("Europe/Berlin") + provider.config.merge_settings_from_dict( + {"general": {"latitude": 52.52, "longitude": 13.405}} + ) + start = to_datetime("2026-01-15 00:00:00", in_timezone="Europe/Berlin") + fixed_now = pd.Timestamp("2026-01-15 13:59:59", tz="Europe/Berlin") + get_ems().set_start_datetime(start) + last_original = start.add(hours=23 if complete_horizon else 22) + history_end = start.subtract(days=history_end_days_ago) + history_index = pd.date_range( + start=history_end.subtract(days=history_days), periods=history_days * 24, freq="1h" + ) + source_index = pd.date_range(start=start, end=last_original, freq="1h") + await provider.key_from_series( + "elecprice_marketprice_raw_wh", + pd.Series(0.0001, index=history_index.append(source_index)), + ) + provider.highest_orig_datetime = last_original + predicted_index = pd.date_range( + start=last_original.add(minutes=15), periods=120, freq="15min" + ) + await provider.key_from_series( + "elecprice_marketprice_raw_wh", pd.Series(0.00005, index=predicted_index) + ) + + repair_history = history_days <= 14 or history_end_days_ago > 0 or force_update + fetch_start = start.subtract(days=35) if repair_history else start + response_index = pd.date_range( + start=fetch_start, end=start.add(days=1), freq="1h", inclusive="left" + ) + response = EnergyChartsElecPrice( + license_info="", + unix_seconds=[int(timestamp.timestamp()) for timestamp in response_index], + price=[200.0] * len(response_index), + unit="EUR/MWh", + deprecated=False, + ) + + def predict(history: np.ndarray, hours: int, slots_per_hour: int = 1) -> np.ndarray: + return np.full(hours, 0.00005) + + # Isolate source lookups from resampling and gross-price derivation, which + # perform their own raw-series reads for different purposes. + with ( + patch("akkudoktoreos.prediction.elecpriceenergycharts.pd", wraps=pd) as pandas, + patch.object( + provider, + "_request_forecast", + return_value=response, + side_effect=requests.exceptions.ReadTimeout("unavailable") if fetch_fails else None, + ) as request, + patch.object( + ElecPriceEnergyCharts, "key_to_raw_series", wraps=provider.key_to_raw_series + ) as raw_reads, + patch.object(ElecPriceEnergyCharts, "key_to_array", return_value=np.full(48, 0.0001)), + patch.object(provider, "_store_gross_series"), + patch.object(provider, "_resolution_seconds", wraps=provider._resolution_seconds) + as resolution, + patch.object(provider, "_predict", side_effect=predict) as prediction, + ): + pandas.Timestamp.now.return_value = fixed_now + await provider._update_data(force_update=force_update) + + if needs_update: + request.assert_called_once_with( + start_date=fetch_start.format("YYYY-MM-DD"), force_update=force_update + ) + else: + request.assert_not_called() + + fetched = needs_update and not fetch_fails + assert raw_reads.await_count == (2 if fetched else 1) + assert raw_reads.await_args_list[0].kwargs["end_datetime"] == last_original.add(seconds=1) + last_source = to_datetime(response_index[-1]) if fetched else last_original + assert provider.highest_orig_datetime == last_source + if fetched: + assert raw_reads.await_args_list[1].kwargs["end_datetime"] == last_source.add(seconds=1) + + # Forecast resolution must use the refreshed snapshot after success and + # retain the original one when a request fails, without predicted points. + forecast_source = resolution.call_args.args[0] + assert forecast_source.index.max() == pd.Timestamp(last_source) + assert forecast_source.iloc[-1] == pytest.approx(0.0002 if fetched else 0.0001) + prediction.assert_called_once() + @pytest.mark.asyncio @patch("requests.get") async def test_update_data_with_incomplete_forecast(self, mock_get, caplog, provider):