from unittest.mock import ANY, MagicMock, patch from pytest import fixture from kafka.errors import NoBrokersAvailable from sientia_do.notifications.models import NotificationLevel from ingestor.managers.data_manager import DataManager metadata = { "metadata": { "model_id": "test_model", "model_name": "test_model", "workflow_name": "test_workflow", "schema_name": "test_schedule", }, } @fixture @patch("ingestor.managers.data_manager.KafkaProducer") @patch("ingestor.managers.data_manager.MongoClient") def data_manager(mongo, kafka): data_manager = DataManager( kafka_servers="localhost:9092", mongo_connection_string="mongodb://localhost:27017", mongo_database="sientia", export_to_kafka=True, logger=MagicMock(), notification_handler=MagicMock(), metadata=metadata["metadata"], ) data_manager.send_notification = MagicMock() return data_manager @patch("ingestor.managers.data_manager.KafkaProducer") @patch("ingestor.managers.data_manager.MongoClient") def test___init___success(mongo, kafka): logger_mock = MagicMock() data_manager = DataManager( metadata=metadata["metadata"], kafka_servers="localhost:9092", mongo_connection_string="mongodb://localhost:27017", mongo_database="sientia", export_to_kafka=True, logger=logger_mock, notification_handler=MagicMock() ) kafka.assert_called_once_with( bootstrap_servers="localhost:9092", value_serializer=ANY, key_serializer=ANY ) assert data_manager.kafka_producer is not None logger_mock.info.assert_any_call( "Trying (0) to initializing DataManager with Kafka servers: localhost:9092" ) logger_mock.info.assert_any_call( "DataManager initialized with Kafka servers: localhost:9092" ) logger_mock.error.assert_not_called() assert logger_mock.info.call_count == 4 @patch("ingestor.managers.data_manager.KafkaProducer") @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(), metadata=metadata["metadata"], ) kafka.assert_any_call( bootstrap_servers="localhost:9092", value_serializer=ANY, key_serializer=ANY ) assert kafka.call_count == 2 assert data_manager.kafka_producer is not None logger_mock.info.assert_any_call( "Trying (0) to initializing DataManager with Kafka servers: localhost:9092" ) logger_mock.info.assert_any_call( "Trying (1) to initializing DataManager with Kafka servers: localhost:9092" ) logger_mock.info.assert_any_call( "DataManager initialized with Kafka servers: localhost:9092" ) logger_mock.error.assert_called_once_with( "Kafka servers localhost:9092 are not available. Retrying..." ) assert logger_mock.info.call_count == 5 @patch("ingestor.managers.data_manager.KafkaProducer") @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(), metadata=metadata["metadata"], ) except NoBrokersAvailable as e: assert str( e) == "NoBrokersAvailable: Failed to connect to Kafka servers localhost:9092 after 3 attempts." assert kafka.call_count == 3 logger_mock.info.assert_any_call( "Trying (0) to initializing DataManager with Kafka servers: localhost:9092" ) logger_mock.info.assert_any_call( "Trying (1) to initializing DataManager with Kafka servers: localhost:9092" ) logger_mock.info.assert_any_call( "Trying (2) to initializing DataManager with Kafka servers: localhost:9092" ) logger_mock.error.assert_called_with( "Failed to connect to Kafka servers localhost:9092 after 3 attempts." ) assert logger_mock.info.call_count == 3 else: assert False, "Expected NoBrokersAvailable exception was not raised." def test_shutdown_has_producer(data_manager): flush_mock = MagicMock() close_mock = MagicMock() data_manager.kafka_producer.flush = flush_mock data_manager.kafka_producer.close = close_mock data_manager.shutdown() flush_mock.assert_called_once() close_mock.assert_called_once() def test_shutdown_no_producer(data_manager): data_manager.kafka_producer = None data_manager.shutdown() data_manager.logger.warning.assert_any_call( "Kafka producer is already closed or not initialized." ) 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")) data_manager.kafka_producer.close = MagicMock() data_manager.shutdown() data_manager.logger.error.assert_called_once_with( "Error closing Kafka producer: Test error" ) 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__() data_manager.shutdown.assert_called_once() def test_delivery_report(data_manager): msg = MagicMock() msg.topic = "test_topic" msg.partition = 0 msg.offset = 1 data_manager.delivery_report(msg) data_manager.logger.debug.assert_called_once_with( f"Record successfully produced to {msg.topic} [{msg.partition}] at offset {msg.offset}" ) def test_delivery_error(data_manager): err = "Test error" data_manager.delivery_error(err) data_manager.logger.error.assert_called_once_with( f"Delivery failed for record : {err}" ) def test_publish(data_manager): topic = "test_topic" data = {"key": "value"} # Mock the send method of the Kafka producer send_mock = MagicMock() data_manager.kafka_producer.send = send_mock # Call the publish method data_manager.publish(topic, data) # Check if the send method was called with the correct arguments send_mock.assert_called_once_with( topic=topic, value=data ) send_mock.return_value.add_callback.assert_called_once() 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.kafka_producer.send.assert_not_called() @patch("ingestor.managers.data_manager.traceback") def test_publish_error(traceback, data_manager): topic = "test_topic" data = {"key": "value"} # Mock the send method of the Kafka producer to raise an exception send_mock = MagicMock(side_effect=Exception("Test error")) data_manager.kafka_producer.send = send_mock # Call the publish method data_manager.publish(topic, data) # Check if the send method was called with the correct arguments send_mock.assert_called_once_with( topic=topic, value=data ) # Check if the error was logged data_manager.send_notification.assert_called_once_with( notification_id=f"KAFKA_PRODUCER_ERROR_{topic}", message=f"Error publishing message to topic {topic}: Test error", block="kafka_producer", level=NotificationLevel.ERROR, attachment_content=traceback.format_exc.return_value, metadata=metadata["metadata"] ) 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.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, metadata=metadata["metadata"] )