From 6672005a99d1f3823074626991143d35db9f4657 Mon Sep 17 00:00:00 2001 From: Mark Vrijlandt Date: Wed, 23 Sep 2026 18:50:56 +0200 Subject: [PATCH 1/2] cleanup job resources --- src/omotes_optimizer_worker/prefect_flow.py | 46 +++++++++ tests/test_prefect_flow.py | 107 +++++++++++++++++++- 2 files changed, 152 insertions(+), 1 deletion(-) diff --git a/src/omotes_optimizer_worker/prefect_flow.py b/src/omotes_optimizer_worker/prefect_flow.py index 0effeaf..bc0659e 100644 --- a/src/omotes_optimizer_worker/prefect_flow.py +++ b/src/omotes_optimizer_worker/prefect_flow.py @@ -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 @@ -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, @@ -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: + 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. diff --git a/tests/test_prefect_flow.py b/tests/test_prefect_flow.py index 917e99b..6e97e79 100644 --- a/tests/test_prefect_flow.py +++ b/tests/test_prefect_flow.py @@ -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", @@ -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_database_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, + } From e9c1ebd9f117ac31e4a29f86f9196b6c63eec522 Mon Sep 17 00:00:00 2001 From: Mark Vrijlandt Date: Fri, 25 Sep 2026 21:19:00 +0200 Subject: [PATCH 2/2] to sdk 5.1.0 --- dev.Dockerfile | 3 +++ pyproject.toml | 2 +- tests/test_prefect_flow.py | 2 +- uv.lock | 8 ++++---- 4 files changed, 9 insertions(+), 6 deletions(-) diff --git a/dev.Dockerfile b/dev.Dockerfile index d0e3f74..abe7e0c 100644 --- a/dev.Dockerfile +++ b/dev.Dockerfile @@ -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" diff --git a/pyproject.toml b/pyproject.toml index a3394f9..75628f4 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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", ] diff --git a/tests/test_prefect_flow.py b/tests/test_prefect_flow.py index 6e97e79..9c2325d 100644 --- a/tests/test_prefect_flow.py +++ b/tests/test_prefect_flow.py @@ -100,7 +100,7 @@ def run_mesido(*args: object, **kwargs: object) -> SimpleNamespace: 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_database_cleanup_resource") as publish_resource, + patch("omotes_optimizer_worker.prefect_flow.publish_optimizer_timeseries_cleanup_resource") as publish_resource, ): result = optimizer_flow.fn( input_esdl=input_esdl, diff --git a/uv.lock b/uv.lock index 4e66046..3fd29ce 100644 --- a/uv.lock +++ b/uv.lock @@ -1102,7 +1102,7 @@ dev = [ [package.metadata] requires-dist = [ { name = "mesido", specifier = "~=0.1.22" }, - { name = "omotes-sdk-python", specifier = "~=5.0.4" }, + { name = "omotes-sdk-python", specifier = "~=5.1.0" }, { name = "pyesdl", specifier = "~=26.7.1" }, { name = "python-dotenv", specifier = "~=1.0.0" }, ] @@ -1121,7 +1121,7 @@ dev = [ [[package]] name = "omotes-sdk-python" -version = "5.0.4" +version = "5.1.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "prefect" }, @@ -1129,9 +1129,9 @@ dependencies = [ { name = "streamcapture" }, { name = "typing-extensions" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/03/eb/b067c8f84aac11a20fbdfa5e10946654ac2894213e00d174113a88f54c7b/omotes_sdk_python-5.0.4.tar.gz", hash = "sha256:47ff5f310ad85dbeb014291d7e29f23c02b20a29724552fb32bca73ca6ea3bd1", size = 26908, upload-time = "2026-08-28T19:47:17.389Z" } +sdist = { url = "https://files.pythonhosted.org/packages/fc/ea/b1a20b8f47367a84f6f34ac65a71ce461ba4e8dc59168822c42bee7b6611/omotes_sdk_python-5.1.0.tar.gz", hash = "sha256:a963656b323a3668ff3777bc31412ba72d37117eef09bfc781a96d5987a77f6d", size = 28538, upload-time = "2026-09-25T16:40:19.169Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/09/89/c6ca93ac2201ab6e88cf41dc5f22edca155c23199362770650d6927033cc/omotes_sdk_python-5.0.4-py3-none-any.whl", hash = "sha256:8653a33fc5a0f6ef838b496b832a26655ee3de3891852d3d3244f6a2d5501c1b", size = 23044, upload-time = "2026-08-28T19:47:15.759Z" }, + { url = "https://files.pythonhosted.org/packages/85/f4/b33ddc4644e96754fd513cda03e977e517ddefb5fecff3f321f95fac496b/omotes_sdk_python-5.1.0-py3-none-any.whl", hash = "sha256:3ac932d37d0bc3e00e877b86f7c4f9ce8a8770bb8218843f93d57e162d05cf0a", size = 23691, upload-time = "2026-09-25T16:40:17.952Z" }, ] [[package]]