Skip to content
Merged
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
3 changes: 3 additions & 0 deletions dev.Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,9 @@ RUN apt update && \
apt install -y wget && \
apt-get clean

COPY --from=docker:27-cli /usr/local/bin/docker /usr/local/bin/docker
COPY --from=docker:27-cli /usr/local/libexec/docker/cli-plugins/docker-buildx /usr/local/libexec/docker/cli-plugins/docker-buildx

ENV GUROBI_HOME=/app/gurobi/
ENV PATH="/app/gurobi/bin:$PATH"
ENV LD_LIBRARY_PATH="/app/gurobi/lib"
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ requires-python = "==3.11.*"
dependencies = [
"python-dotenv ~= 1.0.0",
"mesido ~= 0.1.22",
"omotes-sdk-python ~= 5.0.4",
"omotes-sdk-python ~= 5.1.0",
"pyesdl ~= 26.7.1",
]

Expand Down
46 changes: 46 additions & 0 deletions src/omotes_optimizer_worker/prefect_flow.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,11 @@
)
from omotes_sdk.log_forwarding import StdCaptureToLogSession
from omotes_sdk.prefect_util import (
TimeseriesResource,
create_flow_progress_updater,
in_prefect_flow_context,
load_gurobi_license,
publish_job_cleanup_resource,
write_flow_return_artifact_to_minio,
)
from prefect import flow
Expand All @@ -42,6 +44,42 @@ class OptimizerFlowResult(BaseModel):
esdl_messages: list[dict[str, Any]] = Field(default_factory=list, json_schema_extra={"file_extension": ".json"})


def publish_optimizer_timeseries_cleanup_resource(
db_host: str,
db_port: int,
output_energy_system_id: str,
output_profiles_type: ESDLOutputProfilesType,
pg_database: str | None = None,
) -> None:
"""Attach the timeseries resource created by Mesido to the Prefect flow run, for cleanup purposes.

Raises:
ValueError: If PostgreSQL output is missing its database name.

"""
if output_profiles_type == ESDLOutputProfilesType.POSTGRESQL:
if pg_database is None:
raise ValueError("PostgreSQL resource metadata requires a database name")
resource = TimeseriesResource(
type="postgresql",
host=db_host,
port=db_port,
database=pg_database,
schema_name=output_energy_system_id,
)
elif output_profiles_type == ESDLOutputProfilesType.INFLUXDB:
resource = TimeseriesResource(
type="influxdb",
host=db_host,
port=db_port,
database=output_energy_system_id,
)
else:
return

publish_job_cleanup_resource(resource)


@flow(timeout_seconds=EnvSettings.prefect_flow_timeout_seconds())
def optimizer_flow(
input_esdl: str,
Expand Down Expand Up @@ -151,6 +189,14 @@ def optimizer_flow(
output_esh = EnergySystemHandler()
output_esh.load_from_string(output_esdl)
output_energy_system: EnergySystem = output_esh.energy_system
if esdl_output_profiles_type is not None:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Minor: I will move this below line 200. Just for a better readability.

publish_optimizer_timeseries_cleanup_resource(
db_host=db_host,
db_port=db_port,
output_energy_system_id=output_energy_system.id,
output_profiles_type=esdl_output_profiles_type,
pg_database=pg_db_timeseries,
)
output_energy_system.name = output_esdl_name

# TODO get esdl_messages from successful run after mesido update.
Expand Down
107 changes: 106 additions & 1 deletion tests/test_prefect_flow.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,17 @@
from os import environ
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import patch

from esdl.esdl_handler import EnergySystemHandler
from mesido.esdl.esdl_mixin import ESDLOutputProfilesType
from prefect.states import State

from omotes_optimizer_worker.prefect_flow import OptimizerFlowResult, optimizer_flow
from omotes_optimizer_worker.prefect_flow import (
OptimizerFlowResult,
optimizer_flow,
publish_optimizer_timeseries_cleanup_resource,
)

MINIO_TEST_ENV = {
"MINIO_HOST": "minio",
Expand Down Expand Up @@ -68,3 +75,101 @@ def test_optimizer_flow_returns_delft_feedback_messages() -> None:
assert feedback_result.esdl_messages
assert all(message["technical_message"] for message in feedback_result.esdl_messages)
assert all(message["severity"] == "ERROR" for message in feedback_result.esdl_messages)


def test_optimizer_flow_configures_influxdb_output() -> None:
"""Pass InfluxDB output settings to Mesido and publish its generated database."""
fixture_path = Path(__file__).parent / "data" / "esdl" / "Delft_T.esdl"
input_esdl = fixture_path.read_text()
output_esh = EnergySystemHandler()
output_esh.load_from_string(input_esdl)
mesido_arguments: dict = {}

def run_mesido(*args: object, **kwargs: object) -> SimpleNamespace:
mesido_arguments.update(kwargs)
return SimpleNamespace(optimized_esdl_string=input_esdl)

influx_env = MINIO_TEST_ENV | {
"ESDL_OUTPUT_PROFILES_TYPE": "INFLUXDB",
"DB_HOSTNAME": "omotes_influxdb",
"DB_PORT": "8096",
}
with (
patch.dict(environ, influx_env, clear=False),
patch("omotes_optimizer_worker.prefect_flow.get_problem_function", return_value=run_mesido),
patch("omotes_optimizer_worker.prefect_flow.get_problem_type"),
patch("omotes_optimizer_worker.prefect_flow.get_solver_class"),
patch("omotes_optimizer_worker.prefect_flow.write_flow_return_artifact_to_minio"),
patch("omotes_optimizer_worker.prefect_flow.publish_optimizer_timeseries_cleanup_resource") as publish_resource,
):
result = optimizer_flow.fn(
input_esdl=input_esdl,
workflow_config={},
workflow_type_name="grow_optimizer_no_heat_losses",
)

assert isinstance(result, OptimizerFlowResult)
assert mesido_arguments["esdl_output_profiles_type"] == ESDLOutputProfilesType.INFLUXDB
assert mesido_arguments["database_connections"] == [
{
"access_type": "read_write",
"host": "omotes_influxdb",
"port": 8096,
"username": "user",
"password": "password",
"ssl": False,
"verify_ssl": False,
}
]
publish_resource.assert_called_once_with(
db_host="omotes_influxdb",
db_port=8096,
output_energy_system_id=output_esh.energy_system.id,
output_profiles_type=ESDLOutputProfilesType.INFLUXDB,
pg_database=None,
)


def test_publish_optimizer_database_cleanup_resource_for_postgresql() -> None:
"""Publish deletable PostgreSQL resource coordinates without credentials."""
with (
patch("omotes_optimizer_worker.prefect_flow.publish_job_cleanup_resource") as publish_resource,
):
publish_optimizer_timeseries_cleanup_resource(
db_host="postgres",
db_port=5432,
output_energy_system_id="output-esdl-id",
output_profiles_type=ESDLOutputProfilesType.POSTGRESQL,
pg_database="omotes_timeseries",
)

assert publish_resource.call_args.args[0].model_dump(mode="json", by_alias=True) == {
"type": "postgresql",
"host": "postgres",
"port": 5432,
"database": "omotes_timeseries",
"schema": "output-esdl-id",
}
assert "username" not in str(publish_resource.call_args)
assert "password" not in str(publish_resource.call_args)


def test_publish_optimizer_database_cleanup_resource_for_influxdb() -> None:
"""Use the output ESDL ID as the Mesido-created InfluxDB database name."""
with (
patch("omotes_optimizer_worker.prefect_flow.publish_job_cleanup_resource") as publish_resource,
):
publish_optimizer_timeseries_cleanup_resource(
db_host="influxdb",
db_port=8086,
output_energy_system_id="output-esdl-id",
output_profiles_type=ESDLOutputProfilesType.INFLUXDB,
)

assert publish_resource.call_args.args[0].model_dump(mode="json", by_alias=True) == {
"type": "influxdb",
"host": "influxdb",
"port": 8086,
"database": "output-esdl-id",
"schema": None,
}
8 changes: 4 additions & 4 deletions uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading