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())