Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions RELEASE_NOTES.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@

<!-- Here goes notes on how to upgrade from previous versions, including deprecations and what they should be replaced with -->

- `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

<!-- Here goes the main new features and examples or instructions on how to use them -->
Expand Down
4 changes: 3 additions & 1 deletion benchmarks/power_distribution/power_distributor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
10 changes: 8 additions & 2 deletions docs/tutorials/getting_started.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
4 changes: 3 additions & 1 deletion examples/battery_pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
14 changes: 12 additions & 2 deletions src/frequenz/sdk/timeseries/_resampling/_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand All @@ -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
Expand Down Expand Up @@ -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`
Expand All @@ -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
Expand Down
2 changes: 1 addition & 1 deletion src/frequenz/sdk/timeseries/_voltage_streamer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion src/frequenz/sdk/timeseries/consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
2 changes: 1 addition & 1 deletion src/frequenz/sdk/timeseries/grid.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 = (
Expand Down
2 changes: 1 addition & 1 deletion src/frequenz/sdk/timeseries/producer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
30 changes: 20 additions & 10 deletions tests/config/test_actor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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],
Expand Down Expand Up @@ -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
Expand All @@ -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"
Expand All @@ -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
Expand All @@ -235,19 +241,23 @@ 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"
[dict_str_int]
a = 1
b = 2
c = 4
""")
"""
)

config = await config_receiver.receive()
assert config is not None
Expand Down
18 changes: 12 additions & 6 deletions tests/config/test_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -255,15 +255,17 @@ 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

[logging.loggers.test]
name = "test"
level = "DEBUG"
""")
"""
)
return config_file

async def test_full_config_flow(self, config_file: pathlib.Path) -> None:
Expand All @@ -280,15 +282,17 @@ 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

[logging.loggers.test]
name = "test"
level = "INFO"
""")
"""
)

# Check updated config
config = await receiver.receive()
Expand Down Expand Up @@ -316,15 +320,17 @@ 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

[logging.loggers.test]
name = "test"
level = "DEBUG"
""")
"""
)

# Check updated config
config = await receiver.receive()
Expand Down
4 changes: 3 additions & 1 deletion tests/microgrid/fixtures.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
4 changes: 3 additions & 1 deletion tests/microgrid/test_datapipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
10 changes: 8 additions & 2 deletions tests/timeseries/_battery_pool/test_battery_pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
4 changes: 3 additions & 1 deletion tests/timeseries/_pv_pool/test_pv_pool_control_methods.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
5 changes: 4 additions & 1 deletion tests/timeseries/mock_microgrid.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading
Loading