mirror of
https://github.com/Akkudoktor-EOS/EOS.git
synced 2026-10-09 16:06:40 +00:00
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
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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):
|
||||
|
||||
Reference in New Issue
Block a user