From 8cf3273714e7615cb1d5ca9209b2e12060b804e2 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 9 Jun 2025 15:09:53 -0300 Subject: [PATCH 1/2] SIENTIAPDE-1082 Implement shutdown methods for Activities and Kafka; enhance error handling in worker main function --- scouter/activities/activities.py | 4 ++++ scouter/activities/kafka.py | 8 ++++++++ scouter/worker/worker.py | 14 +++++++++++++- 3 files changed, 25 insertions(+), 1 deletion(-) diff --git a/scouter/activities/activities.py b/scouter/activities/activities.py index 1fba641..eec589e 100644 --- a/scouter/activities/activities.py +++ b/scouter/activities/activities.py @@ -65,3 +65,7 @@ class Activities(Postgres, Redis, Kafka, Gates): @activity.defn(name="prepare_activity") async def prepare_activity(self, input_data: dict[str, Any]): await super().prepare_activity(input_data) + + def shutdown(self): + Postgres.close(self) + Kafka.close(self) diff --git a/scouter/activities/kafka.py b/scouter/activities/kafka.py index 55cdee4..ee56368 100644 --- a/scouter/activities/kafka.py +++ b/scouter/activities/kafka.py @@ -25,6 +25,14 @@ class Kafka(BaseActivity): BaseActivity.__init__(self, logger, notification_handler) + def close(self): + """Closes the connector connection.""" + self.logger.info("Closing Kafka connector...") + self.kafka_connector.close() + + def __del__(self): + self.close() + @activity.defn(name="load_from_kafka") async def load_from_kafka(self, input_data: dict[str, Any]) -> dict[str, Any]: """ diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index c49ca64..321620a 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -2,6 +2,7 @@ from temporalio import workflow, client from temporalio.worker import Worker with workflow.unsafe.imports_passed_through(): + import sys import os from sientia_do.notifications.handlers import NotificationHandler from sientia_do.temporal.utils.logger import get_logger @@ -91,7 +92,18 @@ async def main(): logger.info('Workers started successfully') - await asyncio.gather(*handlers) + try: + await asyncio.gather(*handlers) + + except BaseException as e: + logger.error("An unhandled exception occurred: %s", e, exc_info=True) + finally: + if notification_handler: + notification_handler.shutdown() + if activities: + activities.shutdown() + # Exit with a non-zero status code to indicate failure to Kubernetes + sys.exit(1) if __name__ == '__main__': asyncio.run(main()) From f541e05e29023b0542bc3ba35b0f8ff63344e60d Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 9 Jun 2025 15:12:43 -0300 Subject: [PATCH 2/2] Update sientia-mlops-library version to 0.38.1 in requirements.txt --- requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/requirements.txt b/requirements.txt index 5604fc7..9f28c99 100644 --- a/requirements.txt +++ b/requirements.txt @@ -4,4 +4,4 @@ sqlalchemy asyncua redis git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.1.14 -git+ssh://git@github.com/Aignosi/sientia-mlops-library.git +git+ssh://git@github.com/Aignosi/sientia-mlops-library.git@0.38.1