diff --git a/RELEASE_NOTES.md b/RELEASE_NOTES.md index 61ee6f2ad..41773815e 100644 --- a/RELEASE_NOTES.md +++ b/RELEASE_NOTES.md @@ -8,6 +8,8 @@ +- `ResamplerConfig` and `ResamplerConfig2`: `max_data_age_in_periods` has no default value anymore and must be set explicitly. To keep the previous behavior, pass `max_data_age_in_periods=3.0`. + ## New Features diff --git a/benchmarks/power_distribution/power_distributor.py b/benchmarks/power_distribution/power_distributor.py index 08880d173..de229c12f 100644 --- a/benchmarks/power_distribution/power_distributor.py +++ b/benchmarks/power_distribution/power_distributor.py @@ -139,7 +139,9 @@ async def run() -> None: """Create microgrid api and run tests.""" await microgrid.initialize( "grpc://microgrid.sandbox.api.frequenz.io:62060", - ResamplerConfig2(resampling_period=timedelta(seconds=1.0)), + ResamplerConfig2( + resampling_period=timedelta(seconds=1.0), max_data_age_in_periods=3.0 + ), ) all_batteries = connection_manager.get().component_graph.components( diff --git a/docs/tutorials/getting_started.md b/docs/tutorials/getting_started.md index ddeed9d00..e4f5499e3 100644 --- a/docs/tutorials/getting_started.md +++ b/docs/tutorials/getting_started.md @@ -49,7 +49,10 @@ async def run() -> None: # Initialize the microgrid await microgrid.initialize( server_url, - ResamplerConfig2(resampling_period=timedelta(seconds=1)), + ResamplerConfig2( + resampling_period=timedelta(seconds=1), + max_data_age_in_periods=3.0, + ), ) # Define your application logic here @@ -109,7 +112,10 @@ async def run() -> None: # Initialize the microgrid await microgrid.initialize( server_url, - ResamplerConfig2(resampling_period=timedelta(seconds=1)), + ResamplerConfig2( + resampling_period=timedelta(seconds=1), + max_data_age_in_periods=3.0, + ), ) # Define your application logic here diff --git a/examples/battery_pool.py b/examples/battery_pool.py index 971256c37..0babc8cf6 100644 --- a/examples/battery_pool.py +++ b/examples/battery_pool.py @@ -23,7 +23,9 @@ async def main() -> None: await microgrid.initialize( MICROGRID_API_URL, - resampler_config=ResamplerConfig2(resampling_period=timedelta(seconds=1.0)), + resampler_config=ResamplerConfig2( + resampling_period=timedelta(seconds=1.0), max_data_age_in_periods=3.0 + ), ) battery_pool = microgrid.new_battery_pool(priority=5) diff --git a/src/frequenz/sdk/timeseries/_resampling/_config.py b/src/frequenz/sdk/timeseries/_resampling/_config.py index b420a55ee..bc639547c 100644 --- a/src/frequenz/sdk/timeseries/_resampling/_config.py +++ b/src/frequenz/sdk/timeseries/_resampling/_config.py @@ -105,7 +105,7 @@ class ResamplerConfig: It must be a positive time span. """ - max_data_age_in_periods: float = 3.0 + max_data_age_in_periods: float """The maximum age a sample can have to be considered *relevant* for resampling. Expressed in number of periods, where period is the `resampling_period` @@ -115,6 +115,11 @@ class ResamplerConfig: It must be bigger than 1.0. + There is no default value, as the right value depends on the data being + resampled and on how the resampled data is used: bigger values smooth the + output and make it more robust against late or missing samples, but also + make it react slower to changes. + Example: If `resampling_period` is 3 seconds, the input sampling period is 1 and `max_data_age_in_periods` is 2, then data older than 3*2 @@ -323,7 +328,7 @@ class ResamplerConfig2(ResamplerConfig): It must be a positive time span. """ - max_data_age_in_periods: float = 3.0 + max_data_age_in_periods: float """The maximum age a sample can have to be considered *relevant* for resampling. Expressed in number of periods, where period is the `resampling_period` @@ -333,6 +338,11 @@ class ResamplerConfig2(ResamplerConfig): It must be bigger than 1.0. + There is no default value, as the right value depends on the data being + resampled and on how the resampled data is used: bigger values smooth the + output and make it more robust against late or missing samples, but also + make it react slower to changes. + Example: If `resampling_period` is 3 seconds, the input sampling period is 1 and `max_data_age_in_periods` is 2, then data older than 3*2 diff --git a/src/frequenz/sdk/timeseries/_voltage_streamer.py b/src/frequenz/sdk/timeseries/_voltage_streamer.py index 4759b362f..d76d5c5c8 100644 --- a/src/frequenz/sdk/timeseries/_voltage_streamer.py +++ b/src/frequenz/sdk/timeseries/_voltage_streamer.py @@ -41,7 +41,7 @@ class VoltageStreamer: await microgrid.initialize( "grpc://127.0.0.1:50051", - ResamplerConfig2(resampling_period=timedelta(seconds=1)) + ResamplerConfig2(resampling_period=timedelta(seconds=1), max_data_age_in_periods=3.0) ) # Get a receiver for the phase-to-neutral voltage. diff --git a/src/frequenz/sdk/timeseries/consumer.py b/src/frequenz/sdk/timeseries/consumer.py index 0376c0444..c53dc1c1e 100644 --- a/src/frequenz/sdk/timeseries/consumer.py +++ b/src/frequenz/sdk/timeseries/consumer.py @@ -40,7 +40,7 @@ class Consumer: await microgrid.initialize( "grpc://127.0.0.1:50051", - ResamplerConfig2(resampling_period=timedelta(seconds=1.0)) + ResamplerConfig2(resampling_period=timedelta(seconds=1.0), max_data_age_in_periods=3.0) ) consumer = microgrid.consumer() diff --git a/src/frequenz/sdk/timeseries/grid.py b/src/frequenz/sdk/timeseries/grid.py index 9640ee4bf..d5d81d122 100644 --- a/src/frequenz/sdk/timeseries/grid.py +++ b/src/frequenz/sdk/timeseries/grid.py @@ -46,7 +46,7 @@ class Grid: await microgrid.initialize( "grpc://127.0.0.1:50051", - ResamplerConfig2(resampling_period=timedelta(seconds=1)) + ResamplerConfig2(resampling_period=timedelta(seconds=1), max_data_age_in_periods=3.0) ) grid = microgrid.grid() diff --git a/src/frequenz/sdk/timeseries/logical_meter/_logical_meter.py b/src/frequenz/sdk/timeseries/logical_meter/_logical_meter.py index 44c45ad8b..a7bfc0481 100644 --- a/src/frequenz/sdk/timeseries/logical_meter/_logical_meter.py +++ b/src/frequenz/sdk/timeseries/logical_meter/_logical_meter.py @@ -40,7 +40,7 @@ class LogicalMeter: await microgrid.initialize( "grpc://microgrid.sandbox.api.frequenz.io:62060", - ResamplerConfig2(resampling_period=timedelta(seconds=1)), + ResamplerConfig2(resampling_period=timedelta(seconds=1), max_data_age_in_periods=3.0), ) logical_meter = ( diff --git a/src/frequenz/sdk/timeseries/producer.py b/src/frequenz/sdk/timeseries/producer.py index 0c82ae16c..54fe631b9 100644 --- a/src/frequenz/sdk/timeseries/producer.py +++ b/src/frequenz/sdk/timeseries/producer.py @@ -40,7 +40,7 @@ class Producer: await microgrid.initialize( "grpc://127.0.0.1:50051", - ResamplerConfig2(resampling_period=timedelta(seconds=1.0)) + ResamplerConfig2(resampling_period=timedelta(seconds=1.0), max_data_age_in_periods=3.0) ) producer = microgrid.producer() diff --git a/tests/config/test_actor.py b/tests/config/test_actor.py index 238ad15f3..cacead28f 100644 --- a/tests/config/test_actor.py +++ b/tests/config/test_actor.py @@ -135,7 +135,8 @@ async def test_update_multiple_files(self, config_file: pathlib.Path) -> None: config_receiver = config_channel.new_receiver() config_file2 = config_file.parent / "config2.toml" - config_file2.write_text(""" + config_file2.write_text( + """ logging_lvl = 'ERROR' var1 = "0" var2 = "15" @@ -149,7 +150,8 @@ async def test_update_multiple_files(self, config_file: pathlib.Path) -> None: b = 2 c = 4 d = 3 - """) + """ + ) async with ConfigManagingActor( [config_file, config_file2], @@ -177,10 +179,12 @@ async def test_update_multiple_files(self, config_file: pathlib.Path) -> None: } # We overwrite config_file with just two variables - config_file.write_text(""" + config_file.write_text( + """ logging_lvl = 'INFO' list_non_strict_bool = ["false", "0", "true"] - """) + """ + ) config = await config_receiver.receive() assert config is not None @@ -202,7 +206,8 @@ async def test_update_multiple_files(self, config_file: pathlib.Path) -> None: } # Now we only update logging_lvl in config_file2, it still takes precedence - config_file2.write_text(""" + config_file2.write_text( + """ logging_lvl = 'DEBUG' var1 = "0" var2 = "15" @@ -216,7 +221,8 @@ async def test_update_multiple_files(self, config_file: pathlib.Path) -> None: b = 2 c = 4 d = 3 - """) + """ + ) config = await config_receiver.receive() assert config is not None @@ -235,11 +241,14 @@ async def test_update_multiple_files(self, config_file: pathlib.Path) -> None: # Now add one variable to config_file not present in config_file2 and remove # a bunch of variables from config_file2 too, and update a few - config_file.write_text(""" + config_file.write_text( + """ logging_lvl = 'INFO' var10 = "10" - """) - config_file2.write_text(""" + """ + ) + config_file2.write_text( + """ logging_lvl = 'DEBUG' var1 = "3" var_off = "on" @@ -247,7 +256,8 @@ async def test_update_multiple_files(self, config_file: pathlib.Path) -> None: a = 1 b = 2 c = 4 - """) + """ + ) config = await config_receiver.receive() assert config is not None diff --git a/tests/config/test_manager.py b/tests/config/test_manager.py index 3a4341015..e7b1884fc 100644 --- a/tests/config/test_manager.py +++ b/tests/config/test_manager.py @@ -255,7 +255,8 @@ class TestConfigManagerIntegration: def config_file(self, tmp_path: pathlib.Path) -> pathlib.Path: """Create a temporary config file for testing.""" config_file = tmp_path / "config.toml" - config_file.write_text(""" + config_file.write_text( + """ [test] name = "test1" value = 42 @@ -263,7 +264,8 @@ def config_file(self, tmp_path: pathlib.Path) -> pathlib.Path: [logging.loggers.test] name = "test" level = "DEBUG" - """) + """ + ) return config_file async def test_full_config_flow(self, config_file: pathlib.Path) -> None: @@ -280,7 +282,8 @@ async def test_full_config_flow(self, config_file: pathlib.Path) -> None: assert logging.getLogger("test").level == logging.DEBUG # Update config file - config_file.write_text(""" + config_file.write_text( + """ [test] name = "test2" value = 43 @@ -288,7 +291,8 @@ async def test_full_config_flow(self, config_file: pathlib.Path) -> None: [logging.loggers.test] name = "test" level = "INFO" - """) + """ + ) # Check updated config config = await receiver.receive() @@ -316,7 +320,8 @@ async def test_full_config_flow_without_logging( assert logging.getLogger("test").level == logging.WARNING # Update config file - config_file.write_text(""" + config_file.write_text( + """ [test] name = "test2" value = 43 @@ -324,7 +329,8 @@ async def test_full_config_flow_without_logging( [logging.loggers.test] name = "test" level = "DEBUG" - """) + """ + ) # Check updated config config = await receiver.receive() diff --git a/tests/microgrid/fixtures.py b/tests/microgrid/fixtures.py index f1c6b0f12..6fa72940d 100644 --- a/tests/microgrid/fixtures.py +++ b/tests/microgrid/fixtures.py @@ -62,7 +62,9 @@ async def new( if microgrid._data_pipeline._DATA_PIPELINE is not None: microgrid._data_pipeline._DATA_PIPELINE = None await microgrid._data_pipeline.initialize( - ResamplerConfig2(resampling_period=timedelta(seconds=0.1)) + ResamplerConfig2( + resampling_period=timedelta(seconds=0.1), max_data_age_in_periods=3.0 + ) ) streamer = MockComponentDataStreamer(mockgrid.mock_client) diff --git a/tests/microgrid/test_datapipeline.py b/tests/microgrid/test_datapipeline.py index 343b8a7ba..25744d426 100644 --- a/tests/microgrid/test_datapipeline.py +++ b/tests/microgrid/test_datapipeline.py @@ -44,7 +44,9 @@ async def test_actors_started( ) -> None: """Test that the datasourcing, resampling and power distributing actors are started.""" datapipeline = _DataPipeline( - resampler_config=ResamplerConfig2(resampling_period=timedelta(seconds=1)) + resampler_config=ResamplerConfig2( + resampling_period=timedelta(seconds=1), max_data_age_in_periods=3.0 + ) ) await asyncio.sleep(1) diff --git a/tests/timeseries/_battery_pool/test_battery_pool.py b/tests/timeseries/_battery_pool/test_battery_pool.py index 1c56dba8c..3703fc324 100644 --- a/tests/timeseries/_battery_pool/test_battery_pool.py +++ b/tests/timeseries/_battery_pool/test_battery_pool.py @@ -148,7 +148,10 @@ async def setup_all_batteries(mocker: MockerFixture) -> AsyncIterator[SetupArgs] # pylint: disable=protected-access microgrid._data_pipeline._DATA_PIPELINE = None await microgrid._data_pipeline.initialize( - ResamplerConfig2(resampling_period=timedelta(seconds=min_update_interval)) + ResamplerConfig2( + resampling_period=timedelta(seconds=min_update_interval), + max_data_age_in_periods=3.0, + ) ) streamer = MockComponentDataStreamer(mock_microgrid) @@ -199,7 +202,10 @@ async def setup_batteries_pool(mocker: MockerFixture) -> AsyncIterator[SetupArgs # pylint: disable=protected-access microgrid._data_pipeline._DATA_PIPELINE = None await microgrid._data_pipeline.initialize( - ResamplerConfig2(resampling_period=timedelta(seconds=min_update_interval)) + ResamplerConfig2( + resampling_period=timedelta(seconds=min_update_interval), + max_data_age_in_periods=3.0, + ) ) # We don't use status channel from the sdk interface to limit diff --git a/tests/timeseries/_battery_pool/test_battery_pool_control_methods.py b/tests/timeseries/_battery_pool/test_battery_pool_control_methods.py index 0ab77ad22..3ab33c250 100644 --- a/tests/timeseries/_battery_pool/test_battery_pool_control_methods.py +++ b/tests/timeseries/_battery_pool/test_battery_pool_control_methods.py @@ -65,7 +65,9 @@ async def mocks(mocker: MockerFixture) -> AsyncIterator[Mocks]: if microgrid._data_pipeline._DATA_PIPELINE is not None: microgrid._data_pipeline._DATA_PIPELINE = None await microgrid._data_pipeline.initialize( - ResamplerConfig2(resampling_period=timedelta(seconds=0.1)) + ResamplerConfig2( + resampling_period=timedelta(seconds=0.1), max_data_age_in_periods=3.0 + ) ) streamer = MockComponentDataStreamer(mockgrid.mock_client) diff --git a/tests/timeseries/_ev_charger_pool/test_ev_charger_pool_control_methods.py b/tests/timeseries/_ev_charger_pool/test_ev_charger_pool_control_methods.py index 8c2a90d81..ac8f0080b 100644 --- a/tests/timeseries/_ev_charger_pool/test_ev_charger_pool_control_methods.py +++ b/tests/timeseries/_ev_charger_pool/test_ev_charger_pool_control_methods.py @@ -52,7 +52,9 @@ async def mocks(mocker: MockerFixture) -> AsyncIterator[_Mocks]: if microgrid._data_pipeline._DATA_PIPELINE is not None: microgrid._data_pipeline._DATA_PIPELINE = None await microgrid._data_pipeline.initialize( - ResamplerConfig2(resampling_period=timedelta(seconds=0.1)) + ResamplerConfig2( + resampling_period=timedelta(seconds=0.1), max_data_age_in_periods=3.0 + ) ) streamer = MockComponentDataStreamer(mockgrid.mock_client) diff --git a/tests/timeseries/_pv_pool/test_pv_pool_control_methods.py b/tests/timeseries/_pv_pool/test_pv_pool_control_methods.py index 1832e3e6f..3daecde0f 100644 --- a/tests/timeseries/_pv_pool/test_pv_pool_control_methods.py +++ b/tests/timeseries/_pv_pool/test_pv_pool_control_methods.py @@ -45,7 +45,9 @@ async def mocks(mocker: MockerFixture) -> typing.AsyncIterator[_Mocks]: if microgrid._data_pipeline._DATA_PIPELINE is not None: microgrid._data_pipeline._DATA_PIPELINE = None await microgrid._data_pipeline.initialize( - ResamplerConfig2(resampling_period=timedelta(seconds=0.1)) + ResamplerConfig2( + resampling_period=timedelta(seconds=0.1), max_data_age_in_periods=3.0 + ) ) streamer = MockComponentDataStreamer(mockgrid.mock_client) diff --git a/tests/timeseries/_steam_boiler_pool/test_steam_boiler_pool_control_methods.py b/tests/timeseries/_steam_boiler_pool/test_steam_boiler_pool_control_methods.py index 826392a01..e7b50a22b 100644 --- a/tests/timeseries/_steam_boiler_pool/test_steam_boiler_pool_control_methods.py +++ b/tests/timeseries/_steam_boiler_pool/test_steam_boiler_pool_control_methods.py @@ -45,7 +45,9 @@ async def mocks(mocker: MockerFixture) -> typing.AsyncIterator[_Mocks]: if microgrid._data_pipeline._DATA_PIPELINE is not None: microgrid._data_pipeline._DATA_PIPELINE = None await microgrid._data_pipeline.initialize( - ResamplerConfig2(resampling_period=timedelta(seconds=0.1)) + ResamplerConfig2( + resampling_period=timedelta(seconds=0.1), max_data_age_in_periods=3.0 + ) ) streamer = MockComponentDataStreamer(mockgrid.mock_client) diff --git a/tests/timeseries/mock_microgrid.py b/tests/timeseries/mock_microgrid.py index 248cace17..3d158660e 100644 --- a/tests/timeseries/mock_microgrid.py +++ b/tests/timeseries/mock_microgrid.py @@ -211,7 +211,10 @@ async def start(self, mocker: MockerFixture | None = None) -> None: self.init_mock_client(lambda mock_client: mock_client.initialize(local_mocker)) self.mock_resampler = MockResampler( mocker, - ResamplerConfig2(timedelta(seconds=self._sample_rate_s)), + ResamplerConfig2( + resampling_period=timedelta(seconds=self._sample_rate_s), + max_data_age_in_periods=3.0, + ), bat_inverter_ids=self.battery_inverter_ids, pv_inverter_ids=self.pv_inverter_ids, evc_ids=self.evc_ids, diff --git a/tests/timeseries/test_moving_window.py b/tests/timeseries/test_moving_window.py index 90751190e..ceeaa7363 100644 --- a/tests/timeseries/test_moving_window.py +++ b/tests/timeseries/test_moving_window.py @@ -463,7 +463,10 @@ async def test_wait_for_samples_with_resampling( ) -> None: """Test waiting for samples in a moving window with resampling.""" window, sender = init_moving_window( - timedelta(seconds=20), config_class(resampling_period=timedelta(seconds=2)) + timedelta(seconds=20), + config_class( + resampling_period=timedelta(seconds=2), max_data_age_in_periods=3.0 + ), ) async with window: task = asyncio.create_task(window.wait_for_samples(3)) @@ -537,7 +540,9 @@ async def test_resampling_window(fake_time: time_machine.Coordinates) -> None: window_size = timedelta(seconds=16) input_sampling = timedelta(seconds=1) output_sampling = timedelta(seconds=2) - resampler_config = ResamplerConfig(resampling_period=output_sampling) + resampler_config = ResamplerConfig( + resampling_period=output_sampling, max_data_age_in_periods=3.0 + ) async with MovingWindow( size=window_size, diff --git a/tests/timeseries/test_resampling.py b/tests/timeseries/test_resampling.py index 6d94a68bc..1af8125df 100644 --- a/tests/timeseries/test_resampling.py +++ b/tests/timeseries/test_resampling.py @@ -111,6 +111,7 @@ async def test_resampler_config_len_ok( """Test checks on the resampling buffer.""" config = config_class( resampling_period=timedelta(seconds=1.0), + max_data_age_in_periods=3.0, initial_buffer_len=init_len, ) assert config.initial_buffer_len == init_len @@ -129,6 +130,7 @@ async def test_resampler_config_len_warn( """Test checks on the resampling buffer.""" config = config_class( resampling_period=timedelta(seconds=1.0), + max_data_age_in_periods=3.0, initial_buffer_len=init_len, ) assert config.initial_buffer_len == init_len @@ -157,6 +159,7 @@ async def test_resampler_config_len_error( with pytest.raises(ValueError): _ = config_class( resampling_period=timedelta(seconds=1.0), + max_data_age_in_periods=3.0, initial_buffer_len=init_len, ) @@ -169,6 +172,7 @@ async def test_resampler_config_tick_delay_negative_error( with pytest.raises(ValueError, match="tick_delay"): _ = config_class( resampling_period=timedelta(seconds=1.0), + max_data_age_in_periods=3.0, tick_delay=timedelta(milliseconds=-1), ) @@ -182,6 +186,7 @@ async def test_resampler_config_tick_delay_too_big_error( with pytest.raises(ValueError, match="smaller than resampling_period"): _ = config_class( resampling_period=timedelta(seconds=1.0), + max_data_age_in_periods=3.0, tick_delay=tick_delay, ) @@ -587,6 +592,7 @@ async def test_calculate_window_end_trivial_cases( resampler = Resampler( ResamplerConfig( resampling_period=resampling_period, + max_data_age_in_periods=3.0, align_to=align_to, ) ) @@ -599,12 +605,14 @@ async def test_calculate_window_end_trivial_cases( resampler_now = Resampler( ResamplerConfig( resampling_period=resampling_period, + max_data_age_in_periods=3.0, align_to=now, ) ) resampler_none = Resampler( ResamplerConfig( resampling_period=resampling_period, + max_data_age_in_periods=3.0, align_to=None, ) )