Refactor Ingestor and Metrics Classes for Enhanced Asynchronous Operations - Updated IngestorManager, OpcManager, and ResourceManager to utilize asynchronous methods for improved performance. - Integrated MetricsController into various classes for better observability and monitoring. - Adjusted unit tests to accommodate the new asynchronous behavior, ensuring proper mocking of async methods. - Removed deprecated Redis metrics and streamlined resource management logic.
179 lines
6.0 KiB
Python
179 lines
6.0 KiB
Python
from unittest.mock import AsyncMock, MagicMock, patch
|
|
|
|
from pytest import fixture, mark, raises
|
|
|
|
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',
|
|
},
|
|
}
|
|
|
|
|
|
@patch('ingestor.managers.resource_manager.RedisRepository')
|
|
def test___init__(redis_repository):
|
|
logger = MagicMock()
|
|
notification_handler = MagicMock()
|
|
metrics_controller = MagicMock()
|
|
resource_manager = ResourceManager(
|
|
host='localhost',
|
|
port=6379,
|
|
lease_ttl=10,
|
|
heartbeat_ttl=10,
|
|
metadata=metadata['metadata'],
|
|
logger=logger,
|
|
notification_handler=notification_handler,
|
|
metrics_controller=metrics_controller,
|
|
)
|
|
assert resource_manager.redis_repository == redis_repository.return_value
|
|
redis_repository.assert_called_once_with(
|
|
host='localhost',
|
|
port=6379,
|
|
username=None,
|
|
password=None,
|
|
logger=logger,
|
|
notification_handler=notification_handler,
|
|
metrics_controller=metrics_controller,
|
|
)
|
|
redis_repository.return_value.redis_client.ping.assert_called_once()
|
|
|
|
|
|
@patch('ingestor.managers.resource_manager.RedisRepository')
|
|
def test___init__connection_failure(redis_repository):
|
|
redis_repository.return_value.redis_client.ping.side_effect = Exception('Connection failed')
|
|
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(),
|
|
metrics_controller=MagicMock(),
|
|
)
|
|
redis_repository.return_value.redis_client.ping.assert_called_once()
|
|
redis_repository.logger.error.assert_called_once_with(
|
|
'Failed to connect to Redis: Connection failed'
|
|
)
|
|
redis_repository.return_value.redis_client.ping.assert_called_once()
|
|
|
|
|
|
@fixture
|
|
@patch('ingestor.managers.resource_manager.RedisRepository')
|
|
def resource_manager(redis_repository):
|
|
resource_manager = ResourceManager(
|
|
host='localhost',
|
|
port=6379,
|
|
lease_ttl=10,
|
|
heartbeat_ttl=10,
|
|
metadata=metadata['metadata'],
|
|
logger=MagicMock(),
|
|
notification_handler=MagicMock(),
|
|
metrics_controller=MagicMock(),
|
|
)
|
|
|
|
resource_manager.send_notification = MagicMock()
|
|
resource_manager.send_notification_async = AsyncMock()
|
|
resource_manager.emit_metric = AsyncMock()
|
|
|
|
resource_manager.redis_repository = AsyncMock()
|
|
|
|
return resource_manager
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_get_tag_slot(resource_manager):
|
|
result = await resource_manager.get_tag_slot('id')
|
|
assert result == resource_manager.redis_repository.get.return_value
|
|
|
|
resource_manager.redis_repository.get.assert_called_once_with(
|
|
'slot:opc_tags:id', metadata=metadata['metadata']
|
|
)
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_ingestor_heartbeat(resource_manager):
|
|
await resource_manager.ingestor_heartbeat()
|
|
resource_manager.redis_repository.set.assert_called_once_with(
|
|
'heartbeat:ingestor:localhost', 1, ttl=10, metadata=metadata['metadata']
|
|
)
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_lease_tag(resource_manager):
|
|
output = await resource_manager.lease_tag('tag_id')
|
|
assert output is resource_manager.redis_repository.set.return_value
|
|
resource_manager.redis_repository.set.assert_called_once_with(
|
|
'lease:opc_tags:tag_id', 'localhost', ttl=10, nx=True, metadata=metadata['metadata']
|
|
)
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_renew_tag_lease_success(resource_manager):
|
|
resource_manager.redis_repository.get.return_value = 'localhost'
|
|
resource_manager.redis_repository.expire.return_value = True
|
|
result = await resource_manager.renew_tag_lease('tag_id')
|
|
assert result is resource_manager.redis_repository.expire.return_value
|
|
resource_manager.redis_repository.get.assert_called_once_with(
|
|
'lease:opc_tags:tag_id', metadata=metadata['metadata']
|
|
)
|
|
resource_manager.redis_repository.expire.assert_called_once_with(
|
|
'lease:opc_tags:tag_id', 10, metadata=metadata['metadata']
|
|
)
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_renew_tag_lease_failure(resource_manager):
|
|
resource_manager.redis_repository.get.return_value = 'other_pod_id'
|
|
|
|
result = await resource_manager.renew_tag_lease('tag_id')
|
|
assert result is False
|
|
|
|
resource_manager.redis_repository.get.assert_called_once_with(
|
|
'lease:opc_tags:tag_id', metadata=metadata['metadata']
|
|
)
|
|
resource_manager.redis_repository.expire.assert_not_called()
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_drop_tag_lease(resource_manager):
|
|
await resource_manager.drop_tag_lease('tag_id')
|
|
resource_manager.redis_repository.delete.assert_called_once_with(
|
|
'lease:opc_tags:tag_id', metadata=metadata['metadata']
|
|
)
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_get_all_ingestors(resource_manager):
|
|
resource_manager.redis_repository.keys.return_value = ['ingestor1', 'ingestor2']
|
|
result = await resource_manager.get_all_ingestors()
|
|
assert result == ['ingestor1', 'ingestor2']
|
|
resource_manager.redis_repository.keys.assert_called_once_with(
|
|
'heartbeat:ingestor:*', metadata=metadata['metadata']
|
|
)
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_get_all_slots(resource_manager):
|
|
resource_manager.redis_repository.keys.return_value = ['slot1', 'slot2']
|
|
result = await resource_manager.get_all_slots()
|
|
assert result == ['slot1', 'slot2']
|
|
resource_manager.redis_repository.keys.assert_called_once_with(
|
|
'slot:opc_tags:*', metadata=metadata['metadata']
|
|
)
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_get_all_leases(resource_manager):
|
|
resource_manager.redis_repository.keys.return_value = ['lease1', 'lease2']
|
|
result = await resource_manager.get_all_leases()
|
|
assert result == ['lease1', 'lease2']
|
|
resource_manager.redis_repository.keys.assert_called_once_with(
|
|
'lease:opc_tags:*', metadata=metadata['metadata']
|
|
)
|