From 750d947d7eadaf138e2b5ae2c839699579999bc8 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Thu, 3 Jul 2025 15:35:38 -0300 Subject: [PATCH] Update MongoDB connection handling in DataManager and enhance IngestorManager configuration - Updated DataManager to simplify MongoDB connection logic by removing retry loop and adding logging for connection attempts. - Enhanced IngestorManager to accept MongoDB connection parameters and export_to_kafka flag as part of its initialization. - Updated requirements.txt to use the latest version of the sientia-dataops-library. --- ingestor/managers/data_manager.py | 26 +++---- ingestor/managers/ingestor_manager.py | 3 +- requirements.txt | 2 +- tests/unit/managers/test_data_manager.py | 77 ++++++++++++++++++-- tests/unit/managers/test_ingestor_manager.py | 36 ++++++--- tests/unit/test_ingestor.py | 20 +++-- 6 files changed, 125 insertions(+), 39 deletions(-) diff --git a/ingestor/managers/data_manager.py b/ingestor/managers/data_manager.py index 6c29ffd..ad97a1a 100644 --- a/ingestor/managers/data_manager.py +++ b/ingestor/managers/data_manager.py @@ -75,26 +75,17 @@ class DataManager: logger.info( f"DataManager initialized with Kafka servers: {kafka_servers}") + logger.info( + f"Trying to initializing DataManager with MongoDB servers: {mongo_connection_string}" + ) + self.connection_string = mongo_connection_string self.database = mongo_database - for i in range(0, 3): - logger.info( - f"Trying ({i}) to initializing DataManager with MongoDB servers: {self.connection_string}" - ) - try: - self.mongo_client = MongoClient(self.connection_string) - self.mongo_client.server_info() + self.mongo_client = MongoClient(self.connection_string) + self.mongo_client.server_info() - self.mongo_db = self.mongo_client[self.database] - break - except Exception as e: - logger.error(f"Error connecting to MongoDB: {e}") - sleep(5) - else: - raise ValueError( - f"Failed to connect to MongoDB servers {self.connection_string} after 3 attempts." - ) + self.mongo_db = self.mongo_client[self.database] logger.info( f"DataManager initialized with MongoDB servers: {self.connection_string}" @@ -123,6 +114,9 @@ class DataManager: self.mongo_client.close() except Exception as e: self.logger.error(f"Error closing MongoDB client: {e}") + else: + self.logger.warning( + "MongoDB client is already closed or not initialized.") def __del__(self): self.shutdown() diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index dc56e55..8501e81 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -162,7 +162,8 @@ class IngestorManager(): _a = [self.managed_tags[slot].pop(server, None) for slot, _value in current_managed_tags.items()] - for server in self.opc_managers.keys(): + servers = list(self.opc_managers.keys()) + for server in servers: if server not in registered_servers: self.logger.warning( f"Server {server} not found in managed tags. " diff --git a/requirements.txt b/requirements.txt index f67d2f4..47e8725 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,5 +1,5 @@ asyncua==1.1.5 redis -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.1.14 +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.2.0 prometheus_client pymongo \ No newline at end of file diff --git a/tests/unit/managers/test_data_manager.py b/tests/unit/managers/test_data_manager.py index 80f75e0..8a79566 100644 --- a/tests/unit/managers/test_data_manager.py +++ b/tests/unit/managers/test_data_manager.py @@ -7,20 +7,28 @@ from ingestor.managers.data_manager import DataManager @fixture @patch("ingestor.managers.data_manager.KafkaProducer") -def data_manager(kafka): +@patch("ingestor.managers.data_manager.MongoClient") +def data_manager(mongo, kafka): return DataManager( kafka_servers="localhost:9092", + mongo_connection_string="mongodb://localhost:27017", + mongo_database="sientia", + export_to_kafka=True, logger=MagicMock(), notification_handler=MagicMock() ) @patch("ingestor.managers.data_manager.KafkaProducer") -def test___init___success(kafka): +@patch("ingestor.managers.data_manager.MongoClient") +def test___init___success(mongo, kafka): logger_mock = MagicMock() data_manager = DataManager( kafka_servers="localhost:9092", + mongo_connection_string="mongodb://localhost:27017", + mongo_database="sientia", + export_to_kafka=True, logger=logger_mock, notification_handler=MagicMock() ) @@ -38,16 +46,20 @@ def test___init___success(kafka): "DataManager initialized with Kafka servers: localhost:9092" ) logger_mock.error.assert_not_called() - assert logger_mock.info.call_count == 2 + assert logger_mock.info.call_count == 4 @patch("ingestor.managers.data_manager.KafkaProducer") -def test___init___second_attempt(kafka): +@patch("ingestor.managers.data_manager.MongoClient") +def test___init___second_attempt(mongo, kafka): kafka.side_effect = [NoBrokersAvailable, MagicMock()] logger_mock = MagicMock() data_manager = DataManager( kafka_servers="localhost:9092", + mongo_connection_string="mongodb://localhost:27017", + mongo_database="sientia", + export_to_kafka=True, logger=logger_mock, notification_handler=MagicMock() ) @@ -71,17 +83,21 @@ def test___init___second_attempt(kafka): logger_mock.error.assert_called_once_with( "Kafka servers localhost:9092 are not available. Retrying..." ) - assert logger_mock.info.call_count == 3 + assert logger_mock.info.call_count == 5 @patch("ingestor.managers.data_manager.KafkaProducer") -def test___init___failure_max_attempts(kafka): +@patch("ingestor.managers.data_manager.MongoClient") +def test___init___failure_max_attempts(mongo, kafka): kafka.side_effect = NoBrokersAvailable logger_mock = MagicMock() try: DataManager( kafka_servers="localhost:9092", + mongo_connection_string="mongodb://localhost:27017", + mongo_database="sientia", + export_to_kafka=True, logger=logger_mock, notification_handler=MagicMock() ) @@ -130,6 +146,16 @@ def test_shutdown_no_producer(data_manager): ) +def test_shutdown_no_mongo_client(data_manager): + data_manager.mongo_client = None + + data_manager.shutdown() + + data_manager.logger.warning.assert_any_call( + "MongoDB client is already closed or not initialized." + ) + + def test_shutdown_exception(data_manager): data_manager.kafka_producer.flush = MagicMock( side_effect=Exception("Test error")) @@ -141,6 +167,17 @@ def test_shutdown_exception(data_manager): ) +def test_shutdown_exception_mongo(data_manager): + data_manager.mongo_client.close = MagicMock( + side_effect=Exception("Test error")) + + data_manager.shutdown() + + data_manager.logger.error.assert_called_once_with( + "Error closing MongoDB client: Test error" + ) + + def test___del__(data_manager): data_manager.shutdown = MagicMock() data_manager.__del__() @@ -190,6 +227,18 @@ def test_publish(data_manager): data_manager.kafka_producer.flush.assert_called_once() +def test_publish_no_kafka(data_manager): + data_manager.export_to_kafka = False + topic = "test_topic" + data = {"key": "value"} + + data_manager.publish(topic, data) + + data_manager.logger.debug.assert_any_call( + f"Skipping message to topic {topic}: {data}" + ) + + @patch("ingestor.managers.data_manager.traceback") def test_publish_error(traceback, data_manager): topic = "test_topic" @@ -215,3 +264,19 @@ def test_publish_error(traceback, data_manager): level=NotificationLevel.ERROR, attachment_content=traceback.format_exc.return_value ) + + +def test_publish_error_mongo(data_manager): + data_manager.export_to_kafka = False + data_manager.mongo_db.__getitem__.return_value.insert_one = MagicMock( + side_effect=Exception("Test error")) + + data_manager.publish("test_topic", {"key": "value"}) + + data_manager.notification_handler.build_and_send_notification.assert_called_once_with( + notification_id="MONGO_PRODUCER_ERROR_test_topic", + message="Error inserting message to MongoDB: Test error", + block="mongo_producer", + level=NotificationLevel.ERROR, + attachment_content=ANY + ) diff --git a/tests/unit/managers/test_ingestor_manager.py b/tests/unit/managers/test_ingestor_manager.py index a1f2d12..9d27c55 100644 --- a/tests/unit/managers/test_ingestor_manager.py +++ b/tests/unit/managers/test_ingestor_manager.py @@ -10,12 +10,17 @@ from ingestor.managers.ingestor_manager import IngestorManager def ingestor_manager(data_manager_mock, resource_manager_mock): return IngestorManager( kafka_servers="localhost:9092", - redis_host="localhost", - redis_port=6379, + redis_data={ + "host": "localhost", + "port": 6379 + }, lease_ttl=60, heartbeat_ttl=60, pod_id="test_pod", poll_interval=5, + mongo_connection_string="mongodb://localhost:27017", + mongo_database="sientia", + export_to_kafka=False, logger=MagicMock(), notification_handler=MagicMock() ) @@ -29,19 +34,25 @@ def test___init__(notification_handler_mock, resource_manager_mock, data_manager ingestor = IngestorManager( kafka_servers="localhost:9092", - redis_host="localhost", - redis_port=6379, + redis_data={ + "host": "localhost", + "port": 6379 + }, lease_ttl=60, heartbeat_ttl=60, pod_id="test_pod", poll_interval=5, + mongo_connection_string="mongodb://localhost:27017", + mongo_database="sientia", + export_to_kafka=False, logger=MagicMock(), notification_handler=MagicMock() ) opc_manager_mock.assert_not_called() data_manager_mock.assert_called_once_with( - "localhost:9092", ingestor.logger, ingestor.notification_handler) + "localhost:9092", "mongodb://localhost:27017", "sientia", False, + ingestor.logger, ingestor.notification_handler) resource_manager_mock.assert_called_once_with( "localhost", 6379, 60, 60, "test_pod", None, None) assert ingestor.poll_interval == 5 @@ -176,8 +187,10 @@ def test_update_opc_servers(metrics, opc_manager, ingestor_manager): assert ingestor_manager.opc_managers['server2'] != mock - metrics.OPC_MANAGERS_ACTIVE.labels.assert_called_once_with(pod_id=ingestor_manager.pod_id) - metrics.OPC_MANAGERS_ACTIVE.labels.return_value.set.assert_called_once_with(len(ingestor_manager.opc_managers)) + metrics.OPC_MANAGERS_ACTIVE.labels.assert_called_once_with( + pod_id=ingestor_manager.pod_id) + metrics.OPC_MANAGERS_ACTIVE.labels.return_value.set.assert_called_once_with( + len(ingestor_manager.opc_managers)) def test_declare_active(ingestor_manager): @@ -359,7 +372,8 @@ def test_drop_slot_leases(metrics, ingestor_manager): ingestor_manager.resource_manager.drop_tag_lease.assert_any_call("1") ingestor_manager.resource_manager.drop_tag_lease.assert_any_call("2") - metrics.SLOTS_RELEASED.labels.assert_any_call(pod_id=ingestor_manager.pod_id) + metrics.SLOTS_RELEASED.labels.assert_any_call( + pod_id=ingestor_manager.pod_id) metrics.SLOTS_RELEASED.labels.return_value.inc.assert_any_call() @@ -540,8 +554,10 @@ def test_check_opc_servers_integrity_all_healthy(metrics, ingestor_manager): # Verify that no reinitialization was needed ingestor_manager.initialize_opc_from_config.assert_not_called() - metrics.OPC_MANAGERS_ACTIVE.labels.assert_called_once_with(pod_id=ingestor_manager.pod_id) - metrics.OPC_MANAGERS_ACTIVE.labels.return_value.set.assert_called_once_with(len(ingestor_manager.opc_managers)) + metrics.OPC_MANAGERS_ACTIVE.labels.assert_called_once_with( + pod_id=ingestor_manager.pod_id) + metrics.OPC_MANAGERS_ACTIVE.labels.return_value.set.assert_called_once_with( + len(ingestor_manager.opc_managers)) def test_check_opc_servers_integrity_server_lost(ingestor_manager): diff --git a/tests/unit/test_ingestor.py b/tests/unit/test_ingestor.py index cd39f23..189c0a6 100644 --- a/tests/unit/test_ingestor.py +++ b/tests/unit/test_ingestor.py @@ -11,6 +11,7 @@ from ingestor.ingestor import Ingestor def test___init__(notification_handler, init_logger, getenv): getenv.side_effect = [ "localhost:9092,localhost:35", # KAFKA_SERVERS + "true", # EXPORT_TO_KAFKA "localhost1", # REDIS_HOST '63790', # REDIS_PORT "user", # REDIS_USERNAME @@ -18,7 +19,11 @@ def test___init__(notification_handler, init_logger, getenv): '100', # LEASE_TTL '200', # HEARTBEAT_TTL "localhost1", # HOSTNAME - '50' # POLL_INTERVAL + '50', # POLL_INTERVAL + "mongodb://localhost:27017", # MONGODB_URL + "sientia", # MONGODB_USERNAME + "sientia", # MONGODB_PASSWORD + "sientia" # MONGODB_DATABASE ] ingestor = Ingestor() @@ -114,16 +119,21 @@ def test_prepare_ingestor(ingestor_manager_mock, ingestor): ingestor_manager_mock.assert_called_once_with( ingestor.kafka_servers, - ingestor.redis_host, - ingestor.redis_port, + { + "host": ingestor.redis_host, + "port": ingestor.redis_port, + "username": ingestor.redis_username, + "password": ingestor.redis_password + }, ingestor.lease_ttl, ingestor.heartbeat_ttl, ingestor.pod_id, ingestor.poll_interval, + ingestor.mongo_connection_string, + ingestor.mongo_database, ingestor.logger, ingestor.notification_handler, - ingestor.redis_username, - ingestor.redis_password + ingestor.export_to_kafka ) ingestor_manager.declare_active.assert_called_once() ingestor_manager.get_slot_leases.assert_called_once()