Skip to content

Commit 66903e5

Browse files
authored
Add tests for some edge cases to prevent regressions (#1392)
Add targeted tests adapted to the v1.x.x branch that verify behavior for issues identified in the pull request #1371: - Data sourcing: Verify duplicate subscription and unknown components with same channel name does not hang. There is also a test for invalid metrics, but this one is currently failing (it hangs), but since we plan to remove the data sourcing actor, we delay the fix for now. - Resampling: Verify invalid metric request does not block later valid request. - Power status subscriptions (battery, EV charger, PV pools): Verify same-instance subscriptions share the registry channel correctly.
2 parents d2157e2 + b32620a commit 66903e5

5 files changed

Lines changed: 383 additions & 3 deletions

File tree

tests/actor/test_resampling.py

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -184,3 +184,48 @@ async def test_duplicate_request(
184184
)
185185

186186
await resampling_actor._resampler.stop() # pylint: disable=protected-access
187+
188+
189+
async def test_stalled_request_does_not_block_later_subscriptions() -> None:
190+
"""Ensure resampling does not head-of-line block."""
191+
channel_registry = ChannelRegistry(name="test")
192+
data_source_req_chan = Broadcast[ComponentMetricRequest](name="data-source-req")
193+
data_source_req_recv = data_source_req_chan.new_receiver()
194+
resampling_req_chan = Broadcast[ComponentMetricRequest](name="resample-req")
195+
resampling_req_sender = resampling_req_chan.new_sender()
196+
197+
async with ComponentMetricsResamplingActor(
198+
channel_registry=channel_registry,
199+
data_sourcing_request_sender=data_source_req_chan.new_sender(),
200+
resampling_request_receiver=resampling_req_chan.new_receiver(),
201+
config=ResamplerConfig2(
202+
resampling_period=timedelta(seconds=0.2),
203+
max_data_age_in_periods=2,
204+
),
205+
):
206+
first_request = ComponentMetricRequest(
207+
namespace="Resampling-A",
208+
component_id=ComponentId(9),
209+
metric=Metric.BATTERY_SOC_PCT,
210+
start_time=None,
211+
)
212+
second_request = ComponentMetricRequest(
213+
namespace="Resampling-B",
214+
component_id=ComponentId(10),
215+
metric=Metric.AC_ACTIVE_POWER,
216+
start_time=None,
217+
)
218+
219+
await resampling_req_sender.send(first_request)
220+
stalled_data_source_request = await data_source_req_recv.receive()
221+
assert stalled_data_source_request == dataclasses.replace(
222+
first_request, namespace="Resampling-A:Source"
223+
)
224+
225+
await resampling_req_sender.send(second_request)
226+
later_data_source_request = await asyncio.wait_for(
227+
data_source_req_recv.receive(), timeout=0.1
228+
)
229+
assert later_data_source_request == dataclasses.replace(
230+
second_request, namespace="Resampling-B:Source"
231+
)

tests/microgrid/test_data_sourcing.py

Lines changed: 105 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,8 @@
3838
)
3939
from frequenz.sdk.timeseries import Sample
4040

41+
from ..utils.receive_timeout import Timeout, receive_timeout
42+
4143
T = TypeVar("T", bound=ComponentData)
4244

4345
_MICROGRID_ID = MicrogridId(1)
@@ -168,6 +170,109 @@ async def test_data_sourcing_actor( # pylint: disable=too-many-locals
168170
assert -13.0 + i == sample.value.base_value
169171

170172

173+
async def test_duplicate_requests_do_not_block_new_receivers(
174+
mock_connection_manager: mock.Mock, # pylint: disable=redefined-outer-name,unused-argument
175+
) -> None:
176+
"""Ensure duplicate direct subscriptions still deliver samples."""
177+
req_chan = Broadcast[ComponentMetricRequest](name="data_sourcing_requests")
178+
req_sender = req_chan.new_sender()
179+
registry = ChannelRegistry(name="test-registry")
180+
181+
request = ComponentMetricRequest(
182+
"test-namespace",
183+
ComponentId(4),
184+
Metric.AC_ACTIVE_POWER,
185+
None,
186+
)
187+
188+
async with DataSourcingActor(req_chan.new_receiver(), registry):
189+
first_receiver = registry.get_or_create(
190+
Sample[Quantity], request.get_channel_name()
191+
).new_receiver()
192+
await req_sender.send(request)
193+
194+
first_sample = await receive_timeout(first_receiver, timeout=1.0)
195+
assert first_sample is not Timeout
196+
197+
second_receiver = registry.get_or_create(
198+
Sample[Quantity], request.get_channel_name()
199+
).new_receiver()
200+
await req_sender.send(request)
201+
202+
second_sample = await receive_timeout(second_receiver, timeout=1.0)
203+
assert second_sample is not Timeout
204+
205+
206+
async def test_unknown_component_request_does_not_block_later_valid_request(
207+
mock_connection_manager: mock.Mock, # pylint: disable=redefined-outer-name,unused-argument
208+
) -> None:
209+
"""Ensure unknown component requests don't block later valid streams."""
210+
req_chan = Broadcast[ComponentMetricRequest](name="data_sourcing_requests")
211+
req_sender = req_chan.new_sender()
212+
registry = ChannelRegistry(name="test-registry")
213+
214+
unknown_request = ComponentMetricRequest(
215+
"unknown-component",
216+
ComponentId(999),
217+
Metric.AC_ACTIVE_POWER,
218+
None,
219+
)
220+
valid_request = ComponentMetricRequest(
221+
"valid-component",
222+
ComponentId(4),
223+
Metric.AC_ACTIVE_POWER,
224+
None,
225+
)
226+
227+
async with DataSourcingActor(req_chan.new_receiver(), registry):
228+
valid_receiver = registry.get_or_create(
229+
Sample[Quantity], valid_request.get_channel_name()
230+
).new_receiver()
231+
232+
await req_sender.send(unknown_request)
233+
await req_sender.send(valid_request)
234+
235+
valid_sample = await receive_timeout(valid_receiver, timeout=1.0)
236+
assert valid_sample is not Timeout
237+
238+
239+
@pytest.mark.skip(
240+
reason="This is currently failing, but we are probably not going to "
241+
"fix it since the data sourcing actor will likely go away."
242+
)
243+
async def test_invalid_metric_request_does_not_block_later_valid_request(
244+
mock_connection_manager: mock.Mock, # pylint: disable=redefined-outer-name,unused-argument
245+
) -> None:
246+
"""Ensure a bad request doesn't poison later valid subscriptions."""
247+
req_chan = Broadcast[ComponentMetricRequest](name="data_sourcing_requests")
248+
req_sender = req_chan.new_sender()
249+
registry = ChannelRegistry(name="test-registry")
250+
251+
invalid_request = ComponentMetricRequest(
252+
"invalid-metric",
253+
ComponentId(4),
254+
Metric.BATTERY_SOC_PCT,
255+
None,
256+
)
257+
valid_request = ComponentMetricRequest(
258+
"valid-metric",
259+
ComponentId(4),
260+
Metric.AC_ACTIVE_POWER,
261+
None,
262+
)
263+
264+
async with DataSourcingActor(req_chan.new_receiver(), registry):
265+
valid_receiver = registry.get_or_create(
266+
Sample[Quantity], valid_request.get_channel_name()
267+
).new_receiver()
268+
269+
await req_sender.send(invalid_request)
270+
await req_sender.send(valid_request)
271+
272+
valid_sample = await receive_timeout(valid_receiver, timeout=1.0)
273+
assert valid_sample is not Timeout
274+
275+
171276
def _new_meter_data(
172277
component_id: ComponentId, timestamp: datetime, value: float
173278
) -> MeterData:

tests/timeseries/_battery_pool/test_battery_pool.py

Lines changed: 73 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,17 +15,19 @@
1515
from dataclasses import dataclass, is_dataclass, replace
1616
from datetime import datetime, timedelta, timezone
1717
from typing import Any, Generic, TypeVar
18+
from unittest.mock import MagicMock
1819

1920
import async_solipsism
2021
import pytest
2122
import time_machine
22-
from frequenz.channels import Receiver, Sender
23+
from frequenz.channels import Broadcast, Receiver, Sender
2324
from frequenz.client.common.microgrid.components import ComponentId
2425
from frequenz.client.microgrid.component import Battery, Component
2526
from frequenz.quantities import Energy, Percentage, Power, Temperature
2627
from pytest_mock import MockerFixture
2728

28-
from frequenz.sdk import microgrid
29+
from frequenz.sdk import microgrid, timeseries
30+
from frequenz.sdk._internal._channels import ChannelRegistry
2931
from frequenz.sdk._internal._constants import (
3032
MAX_BATTERY_DATA_AGE_SEC,
3133
WAIT_FOR_COMPONENT_DATA_SEC,
@@ -34,9 +36,13 @@
3436
from frequenz.sdk.microgrid._power_distributing._component_managers._battery_manager import (
3537
_get_battery_inverter_mappings,
3638
)
39+
from frequenz.sdk.microgrid._power_managing import ReportRequest, _Report
3740
from frequenz.sdk.timeseries import Bounds, ResamplerConfig2, Sample
3841
from frequenz.sdk.timeseries._base_types import SystemBounds
3942
from frequenz.sdk.timeseries.battery_pool import BatteryPool
43+
from frequenz.sdk.timeseries.battery_pool._battery_pool_reference_store import (
44+
BatteryPoolReferenceStore,
45+
)
4046

4147
from ...timeseries.mock_microgrid import MockMicrogrid
4248
from ...utils.component_data_streamer import MockComponentDataStreamer
@@ -50,6 +56,16 @@
5056
_logger = logging.getLogger(__name__)
5157

5258

59+
def _new_power_status_report(target_power_watts: float) -> _Report:
60+
"""Create a distinct report for power status assertions."""
61+
target_power = Power.from_watts(target_power_watts)
62+
return _Report(
63+
target_power=target_power,
64+
_inclusion_bounds=timeseries.Bounds(target_power, target_power),
65+
_exclusion_bounds=None,
66+
)
67+
68+
5369
@pytest.fixture(autouse=True)
5470
def event_loop_policy() -> async_solipsism.EventLoopPolicy:
5571
"""Return an event loop policy that uses the async solipsism event loop."""
@@ -1200,3 +1216,58 @@ async def run_temperature_test( # pylint: disable=too-many-locals
12001216
streamer.start_streaming(latest_data, sampling_rate=0.1)
12011217
msg = await asyncio.wait_for(receiver.receive(), timeout=waiting_time_sec)
12021218
compare_messages(msg, Sample(now, Temperature.from_celsius(15.0)))
1219+
1220+
1221+
async def test_power_status_same_instance_subscriptions_work(
1222+
mocker: MockerFixture,
1223+
) -> None:
1224+
"""Ensure same-instance power_status subscriptions share the same channel."""
1225+
mock_cm = MagicMock()
1226+
mock_graph = MagicMock()
1227+
mock_graph.components.return_value = [
1228+
MagicMock(id=ComponentId(8)),
1229+
MagicMock(id=ComponentId(18)),
1230+
]
1231+
mock_cm.component_graph = mock_graph
1232+
mocker.patch(
1233+
"frequenz.sdk.microgrid.connection_manager._CONNECTION_MANAGER",
1234+
mock_cm,
1235+
)
1236+
mocker.patch("frequenz.sdk.microgrid.connection_manager.get", return_value=mock_cm)
1237+
1238+
registry = ChannelRegistry(name="battery-pool-test")
1239+
requests_channel = Broadcast[ReportRequest](name="battery-pool-requests")
1240+
requests_rx = requests_channel.new_receiver()
1241+
component_ids = frozenset({ComponentId(8), ComponentId(18)})
1242+
pool = BatteryPool(
1243+
pool_ref_store=BatteryPoolReferenceStore(
1244+
channel_registry=registry,
1245+
resampler_subscription_sender=MagicMock(),
1246+
batteries_status_receiver=MagicMock(),
1247+
power_manager_requests_sender=MagicMock(),
1248+
power_manager_bounds_subscription_sender=requests_channel.new_sender(),
1249+
power_distribution_results_fetcher=MagicMock(),
1250+
min_update_interval=timedelta(seconds=1),
1251+
batteries_id=component_ids,
1252+
),
1253+
name="battery-pool",
1254+
priority=5,
1255+
)
1256+
1257+
first_status_rx = pool.power_status.new_receiver()
1258+
second_status_rx = pool.power_status.new_receiver()
1259+
1260+
await asyncio.sleep(0)
1261+
1262+
first_request = await asyncio.wait_for(requests_rx.receive(), timeout=1.0)
1263+
second_request = await asyncio.wait_for(requests_rx.receive(), timeout=1.0)
1264+
assert second_request.get_channel_name() == first_request.get_channel_name()
1265+
1266+
await registry.get_or_create(
1267+
_Report, first_request.get_channel_name()
1268+
).new_sender().send(_new_power_status_report(123.0))
1269+
1270+
first_report = await asyncio.wait_for(first_status_rx.receive(), timeout=1.0)
1271+
second_report = await asyncio.wait_for(second_status_rx.receive(), timeout=1.0)
1272+
assert first_report.target_power == Power.from_watts(123.0)
1273+
assert second_report.target_power == Power.from_watts(123.0)

tests/timeseries/_ev_charger_pool/test_ev_charger_pool.py

Lines changed: 76 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,13 +4,34 @@
44
"""Tests for the `EVChargerPool`."""
55

66

7+
import asyncio
8+
from unittest.mock import MagicMock
9+
10+
from frequenz.channels import Broadcast
11+
from frequenz.client.common.microgrid.components import ComponentId
712
from frequenz.quantities import Power
813
from pytest_mock import MockerFixture
914

10-
from frequenz.sdk import microgrid
15+
from frequenz.sdk import microgrid, timeseries
16+
from frequenz.sdk._internal._channels import ChannelRegistry
17+
from frequenz.sdk.microgrid._power_managing import ReportRequest, _Report
18+
from frequenz.sdk.timeseries.ev_charger_pool import EVChargerPool
19+
from frequenz.sdk.timeseries.ev_charger_pool._ev_charger_pool_reference_store import (
20+
EVChargerPoolReferenceStore,
21+
)
1122
from tests.timeseries.mock_microgrid import MockMicrogrid
1223

1324

25+
def _new_power_status_report(target_power_watts: float) -> _Report:
26+
"""Create a distinct report for power status assertions."""
27+
target_power = Power.from_watts(target_power_watts)
28+
return _Report(
29+
target_power=target_power,
30+
_inclusion_bounds=timeseries.Bounds(target_power, target_power),
31+
_exclusion_bounds=None,
32+
)
33+
34+
1435
class TestEVChargerPool:
1536
"""Tests for the `EVChargerPool`."""
1637

@@ -28,3 +49,57 @@ async def test_ev_power( # pylint: disable=too-many-locals
2849

2950
await mockgrid.mock_resampler.send_meter_power([16.0])
3051
assert (await power_receiver.receive()).value == Power.from_watts(16.0)
52+
53+
54+
async def test_power_status_same_instance_subscriptions_work(
55+
mocker: MockerFixture,
56+
) -> None:
57+
"""Ensure same-instance power_status subscriptions share the same channel."""
58+
mock_cm = MagicMock()
59+
mock_graph = MagicMock()
60+
mock_graph.components.return_value = [
61+
MagicMock(id=ComponentId(12)),
62+
MagicMock(id=ComponentId(22)),
63+
]
64+
mock_cm.component_graph = mock_graph
65+
mocker.patch(
66+
"frequenz.sdk.microgrid.connection_manager._CONNECTION_MANAGER",
67+
mock_cm,
68+
)
69+
mocker.patch("frequenz.sdk.microgrid.connection_manager.get", return_value=mock_cm)
70+
71+
registry = ChannelRegistry(name="ev-pool-test")
72+
requests_channel = Broadcast[ReportRequest](name="ev-pool-requests")
73+
requests_rx = requests_channel.new_receiver()
74+
component_ids = frozenset({ComponentId(12), ComponentId(22)})
75+
pool = EVChargerPool(
76+
pool_ref_store=EVChargerPoolReferenceStore(
77+
channel_registry=registry,
78+
resampler_subscription_sender=MagicMock(),
79+
status_receiver=MagicMock(),
80+
power_manager_requests_sender=MagicMock(),
81+
power_manager_bounds_subs_sender=requests_channel.new_sender(),
82+
power_distribution_results_fetcher=MagicMock(),
83+
component_ids=component_ids,
84+
),
85+
name="ev-pool",
86+
priority=5,
87+
)
88+
89+
first_status_rx = pool.power_status.new_receiver()
90+
second_status_rx = pool.power_status.new_receiver()
91+
92+
await asyncio.sleep(0)
93+
94+
first_request = await asyncio.wait_for(requests_rx.receive(), timeout=1.0)
95+
second_request = await asyncio.wait_for(requests_rx.receive(), timeout=1.0)
96+
assert second_request.get_channel_name() == first_request.get_channel_name()
97+
98+
await registry.get_or_create(
99+
_Report, first_request.get_channel_name()
100+
).new_sender().send(_new_power_status_report(123.0))
101+
102+
first_report = await asyncio.wait_for(first_status_rx.receive(), timeout=1.0)
103+
second_report = await asyncio.wait_for(second_status_rx.receive(), timeout=1.0)
104+
assert first_report.target_power == Power.from_watts(123.0)
105+
assert second_report.target_power == Power.from_watts(123.0)

0 commit comments

Comments
 (0)