- Updated the `sientia-dataops-library` dependency version from 1.4.3 to 1.4.6 in `requirements.txt`. - Modified the GitHub Actions workflow to install development and runtime dependencies separately, improving clarity and organization. - Added code formatting and linting checks using Ruff, along with type checking using mypy, to ensure code quality. - Updated `.gitignore` to include additional cache directories and log files. - Refactored code in various files for consistency in string formatting and improved logging messages.
225 lines
8.4 KiB
Python
225 lines
8.4 KiB
Python
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,
|
|
)
|