from unittest.mock import MagicMock, patch from pytest import fixture, raises from sientia_do.notifications.models import NotificationLevel from ingestor.managers.resource_manager import ResourceManager metadata = { 'metadata': { 'model_id': 'test_model', 'model_name': 'test_model', 'workflow_name': 'test_workflow', 'schema_name': 'test_schedule', }, } @fixture @patch('ingestor.managers.resource_manager.Redis') def resource_manager(redis): resource_manager = ResourceManager( host='localhost', port=6379, lease_ttl=10, heartbeat_ttl=10, metadata=metadata['metadata'], logger=MagicMock(), notification_handler=MagicMock(), ) resource_manager.send_notification = MagicMock() return resource_manager def test_get_success(resource_manager): resource_manager.redis.get.return_value = '{"key": "value"}' result = resource_manager.get('key') assert result == {'key': 'value'} resource_manager.redis.get.assert_called_once_with('key') def test_get_failure(resource_manager): resource_manager.redis.get.return_value = None result = resource_manager.get('key') assert result is None resource_manager.redis.get.assert_called_once_with('key') def test_get_tag_slot(resource_manager): resource_manager.get = MagicMock(return_value={'tag': 'slot'}) result = resource_manager.get_tag_slot('id') assert result == {'tag': 'slot'} resource_manager.get.assert_called_once_with('slot:opc_tags:id') def test_ingestor_heartbeat(resource_manager): resource_manager.redis.set.return_value = True resource_manager.ingestor_heartbeat() resource_manager.redis.set.assert_called_once_with('heartbeat:ingestor:localhost', 1, ex=10) def test_lease_tag(resource_manager): resource_manager.redis.set.return_value = True output = resource_manager.lease_tag('tag_id') assert output is True resource_manager.redis.set.assert_called_once_with( 'lease:opc_tags:tag_id', 'localhost', nx=True, ex=10 ) def test_renew_tag_lease_success(resource_manager): resource_manager.redis.get.return_value = 'localhost' resource_manager.redis.expire.return_value = True result = resource_manager.renew_tag_lease('tag_id') assert result is True resource_manager.redis.get.assert_called_once_with('lease:opc_tags:tag_id') resource_manager.redis.expire.assert_called_once_with('lease:opc_tags:tag_id', 10) def test_renew_tag_lease_failure(resource_manager): resource_manager.redis.get.return_value = 'other_pod_id' resource_manager.redis.expire.return_value = False result = resource_manager.renew_tag_lease('tag_id') assert result is False resource_manager.redis.get.assert_called_once_with('lease:opc_tags:tag_id') resource_manager.redis.expire.assert_not_called() def test_drop_tag_lease(resource_manager): resource_manager.redis.delete.return_value = True resource_manager.drop_tag_lease('tag_id') resource_manager.redis.delete.assert_called_once_with('lease:opc_tags:tag_id') def test_get_all_ingestors(resource_manager): resource_manager.redis.keys.return_value = ['ingestor1', 'ingestor2'] result = resource_manager.get_all_ingestors() assert result == ['ingestor1', 'ingestor2'] resource_manager.redis.keys.assert_called_once_with('heartbeat:ingestor:*') def test_get_all_slots(resource_manager): resource_manager.redis.keys.return_value = ['slot1', 'slot2'] result = resource_manager.get_all_slots() assert result == ['slot1', 'slot2'] resource_manager.redis.keys.assert_called_once_with('slot:opc_tags:*') def test_get_all_leases(resource_manager): resource_manager.redis.keys.return_value = ['lease1', 'lease2'] result = resource_manager.get_all_leases() assert result == ['lease1', 'lease2'] resource_manager.redis.keys.assert_called_once_with('lease:opc_tags:*') def test_init_connection_failure(monkeypatch): # Mock Redis to raise an exception during initialization mock_redis = MagicMock() mock_redis.side_effect = Exception('Connection failed') monkeypatch.setattr('ingestor.managers.resource_manager.Redis', mock_redis) # Test that the exception is raised and metrics are set properly with patch('ingestor.metrics.REDIS_CONNECTION_STATUS') as mock_metrics: mock_status = MagicMock() mock_metrics.labels.return_value = mock_status with raises(Exception, match='Connection failed'): ResourceManager( host='localhost', port=6379, lease_ttl=10, heartbeat_ttl=10, metadata=metadata['metadata'], logger=MagicMock(), notification_handler=MagicMock(), ) mock_metrics.labels.assert_called_once_with(pod_id='localhost') mock_status.set.assert_called_once_with(0) def test_init_ping_failure(monkeypatch): # Mock Redis ping to raise an exception mock_redis_instance = MagicMock() mock_redis_instance.ping.side_effect = Exception('Ping failed') mock_redis_class = MagicMock(return_value=mock_redis_instance) monkeypatch.setattr('ingestor.managers.resource_manager.Redis', mock_redis_class) # Test that the exception is raised and metrics are set properly with patch('ingestor.metrics.REDIS_CONNECTION_STATUS') as mock_metrics: mock_status = MagicMock() mock_metrics.labels.return_value = mock_status with raises(Exception, match='Ping failed'): ResourceManager( host='localhost', port=6379, lease_ttl=10, heartbeat_ttl=10, metadata=metadata['metadata'], logger=MagicMock(), notification_handler=MagicMock(), ) mock_metrics.labels.assert_called_once_with(pod_id='localhost') mock_status.set.assert_called_once_with(0) def test_execute_redis_op_success(resource_manager): # Mock the Redis operation and time function mock_func = MagicMock(return_value='test_result') with patch('ingestor.managers.resource_manager.time', side_effect=[100, 100.5]): with patch('ingestor.metrics.REDIS_OPERATIONS_TOTAL') as mock_total: with patch('ingestor.metrics.REDIS_OPERATIONS_DURATION') as mock_duration: mock_total_labels = MagicMock() mock_duration_labels = MagicMock() mock_total.labels.return_value = mock_total_labels mock_duration.labels.return_value = mock_duration_labels # Execute the operation result = resource_manager._execute_redis_op( 'test_op', mock_func, 'arg1', kwarg1='value1' ) # Verify the result and metrics assert result == 'test_result' mock_func.assert_called_once_with('arg1', kwarg1='value1') mock_total.labels.assert_called_once_with(pod_id='localhost', operation='test_op') mock_total_labels.inc.assert_called_once() mock_duration.labels.assert_called_once_with( pod_id='localhost', operation='test_op' ) mock_duration_labels.observe.assert_called_once_with(0.5) def test_execute_redis_op_exception(resource_manager): # Mock the Redis operation to raise an exception mock_func = MagicMock(side_effect=Exception('Operation failed')) with patch('ingestor.managers.resource_manager.time', return_value=100): with patch('ingestor.metrics.REDIS_OPERATIONS_ERRORS') as mock_errors: mock_errors_labels = MagicMock() mock_errors.labels.return_value = mock_errors_labels # Execute the operation and expect an exception with raises(Exception, match='Operation failed'): resource_manager._execute_redis_op('test_op', mock_func, 'arg1') # Verify metrics and error handling mock_errors.labels.assert_called_once_with(pod_id='localhost', operation='test_op') mock_errors_labels.inc.assert_called_once() resource_manager.send_notification.assert_called_once_with( metadata=metadata['metadata'], notification_id='REDIS_OPERATION_ERROR_test_op', message="Error in Redis operation 'test_op': Operation failed", block='redis_manager', level=NotificationLevel.ERROR, )