Refactor OpcManager connection error handling and enhance unit tests - Simplified error handling during connection attempts in OpcManager by removing redundant disconnection logic. - Updated unit tests to include subscription_period_ms in server configuration and adjusted connection assertions to include a timeout parameter. - Added new tests for disconnection fallback functionality to ensure robust error handling during disconnect attempts.
548 lines
21 KiB
Python
548 lines
21 KiB
Python
from unittest.mock import AsyncMock, MagicMock, patch
|
|
|
|
from pytest import fixture, mark
|
|
from sientia_do.notifications.models import NotificationLevel
|
|
|
|
from ingestor.managers.ingestor_manager import IngestorManager
|
|
|
|
metadata = {
|
|
'metadata': {
|
|
'model_id': 'test_model',
|
|
'model_name': 'test_model',
|
|
'workflow_name': 'test_workflow',
|
|
'schema_name': 'test_schedule',
|
|
},
|
|
}
|
|
|
|
|
|
@fixture
|
|
@patch('ingestor.managers.ingestor_manager.DataManager')
|
|
@patch('ingestor.managers.ingestor_manager.ResourceManager')
|
|
def ingestor_manager(data_manager_mock, resource_manager_mock):
|
|
ingestor = IngestorManager(
|
|
kafka_servers='localhost:9092',
|
|
redis_data={'host': 'localhost', 'port': 6379},
|
|
lease_ttl=60,
|
|
heartbeat_ttl=60,
|
|
poll_interval=5,
|
|
mongo_connection_string='mongodb://localhost:27017',
|
|
mongo_database='sientia',
|
|
export_to_kafka=False,
|
|
logger=MagicMock(),
|
|
notification_handler=MagicMock(),
|
|
metadata=metadata['metadata'],
|
|
)
|
|
|
|
ingestor.send_notification = MagicMock()
|
|
|
|
return ingestor
|
|
|
|
|
|
@patch('ingestor.managers.ingestor_manager.OpcManager')
|
|
@patch('ingestor.managers.ingestor_manager.DataManager')
|
|
@patch('ingestor.managers.ingestor_manager.ResourceManager')
|
|
@patch('ingestor.managers.ingestor_manager.NotificationHandler')
|
|
def test___init__(
|
|
notification_handler_mock, resource_manager_mock, data_manager_mock, opc_manager_mock
|
|
):
|
|
ingestor = IngestorManager(
|
|
kafka_servers='localhost:9092',
|
|
redis_data={'host': 'localhost', 'port': 6379},
|
|
lease_ttl=60,
|
|
heartbeat_ttl=60,
|
|
poll_interval=5,
|
|
mongo_connection_string='mongodb://localhost:27017',
|
|
mongo_database='sientia',
|
|
export_to_kafka=False,
|
|
logger=MagicMock(),
|
|
notification_handler=MagicMock(),
|
|
metadata=metadata['metadata'],
|
|
)
|
|
|
|
opc_manager_mock.assert_not_called()
|
|
data_manager_mock.assert_called_once_with(
|
|
kafka_servers='localhost:9092',
|
|
mongo_connection_string='mongodb://localhost:27017',
|
|
mongo_database='sientia',
|
|
export_to_kafka=False,
|
|
metadata=metadata['metadata'],
|
|
logger=ingestor.logger,
|
|
notification_handler=ingestor.notification_handler,
|
|
)
|
|
resource_manager_mock.assert_called_once_with(
|
|
host='localhost',
|
|
port=6379,
|
|
lease_ttl=60,
|
|
heartbeat_ttl=60,
|
|
metadata=metadata['metadata'],
|
|
logger=ingestor.logger,
|
|
notification_handler=ingestor.notification_handler,
|
|
username=None,
|
|
password=None,
|
|
)
|
|
assert ingestor.poll_interval == 5
|
|
assert ingestor.managed_tags == {}
|
|
assert ingestor.opc_servers == {}
|
|
assert ingestor.opc_managers == {}
|
|
assert ingestor.data_manager == data_manager_mock.return_value
|
|
assert ingestor.resource_manager == resource_manager_mock.return_value
|
|
|
|
|
|
@mark.asyncio
|
|
@patch('ingestor.managers.ingestor_manager.OpcManager')
|
|
async def test_initialize_opc_from_config(opc_manager, ingestor_manager):
|
|
server_config = {
|
|
'name': 'server1',
|
|
'url': 'opc.tcp://localhost:4840',
|
|
'subscription_period_ms': 1000,
|
|
'server_uri': 'http://opcua-server.simulator',
|
|
'cert_path': '/path/to/cert',
|
|
'private_key_path': '/path/to/private_key',
|
|
'server_cert_path': '/path/to/server_cert',
|
|
'pod_id': 'test_pod',
|
|
}
|
|
|
|
opc_manager.return_value = MagicMock(connect=AsyncMock())
|
|
result = await ingestor_manager.initialize_opc_from_config(server_config)
|
|
|
|
opc_manager.assert_called_once_with(
|
|
name=server_config['name'],
|
|
url=server_config['url'],
|
|
data_manager=ingestor_manager.data_manager,
|
|
logger=ingestor_manager.logger,
|
|
subscription_period_ms=server_config['subscription_period_ms'],
|
|
server_uri=server_config['server_uri'],
|
|
notification_handler=ingestor_manager.notification_handler,
|
|
cert_path=server_config['cert_path'],
|
|
private_key_path=server_config['private_key_path'],
|
|
server_cert_path=server_config['server_cert_path'],
|
|
metadata=metadata['metadata'],
|
|
)
|
|
|
|
assert result == opc_manager.return_value
|
|
result.connect.assert_called_once()
|
|
|
|
|
|
@mark.asyncio
|
|
@patch('ingestor.managers.ingestor_manager.OpcManager')
|
|
@patch('ingestor.managers.ingestor_manager.traceback')
|
|
async def test_initialize_opc_from_config_exception(traceback_mock, opc_manager, ingestor_manager):
|
|
server_config = {
|
|
'name': 'server1',
|
|
'url': 'opc.tcp://localhost:4840',
|
|
'subscription_period_ms': 1000,
|
|
'server_uri': 'http://opcua-server.simulator',
|
|
'cert_path': '/path/to/cert',
|
|
'private_key_path': '/path/to/private_key',
|
|
'server_cert_path': '/path/to/server_cert',
|
|
}
|
|
|
|
ingestor_manager.logger.error = MagicMock()
|
|
opc_manager.side_effect = Exception('Initialization error')
|
|
|
|
result = await ingestor_manager.initialize_opc_from_config(server_config)
|
|
|
|
assert result is None
|
|
|
|
traceback_mock.format_exc.assert_called_once()
|
|
ingestor_manager.send_notification.assert_called_once_with(
|
|
metadata=metadata['metadata'],
|
|
notification_id=f'OPC_CONNECTION_ERROR_{server_config["name"]}',
|
|
message='Error initializing OPC manager: Initialization error',
|
|
block='opc_manager',
|
|
level=NotificationLevel.ERROR,
|
|
attachment_content=traceback_mock.format_exc.return_value,
|
|
)
|
|
|
|
|
|
@mark.asyncio
|
|
@patch('ingestor.managers.ingestor_manager.OpcManager')
|
|
@patch('ingestor.managers.ingestor_manager.metrics')
|
|
async def test_update_opc_servers(metrics, opc_manager, ingestor_manager):
|
|
manager1 = MagicMock(config={'config': 'config1'})
|
|
manager2 = MagicMock(config={'config': 'config2'})
|
|
manager3 = MagicMock(config={'config': 'config3'})
|
|
|
|
async def mock_initialize_from_config(config):
|
|
if config == {'config': 'config1'}:
|
|
return manager1
|
|
elif config == {'config': 'config2'}:
|
|
return manager2
|
|
elif config == {'config': 'config3'}:
|
|
return manager3
|
|
else:
|
|
return None
|
|
|
|
ingestor_manager.initialize_opc_from_config = AsyncMock(side_effect=mock_initialize_from_config)
|
|
|
|
ingestor_manager.managed_tags = {
|
|
'slot1': {
|
|
'server1': {'config': 'config1'},
|
|
'server2': {'config': 'config2'},
|
|
'server5': {'config': 'config5'},
|
|
},
|
|
'slot2': {'server3': {'config': 'config3'}, 'server1': {'config': 'config1'}},
|
|
}
|
|
|
|
mock = MagicMock(config={'config': 'old_config2'})
|
|
ingestor_manager.opc_managers['server3'] = AsyncMock(config={'config': 'config3'})
|
|
ingestor_manager.opc_managers['server2'] = mock
|
|
ingestor_manager.opc_managers['server4'] = AsyncMock()
|
|
|
|
await ingestor_manager.update_opc_servers()
|
|
|
|
assert len(ingestor_manager.opc_managers) == 3
|
|
|
|
ingestor_manager.initialize_opc_from_config.assert_any_call({'config': 'config1'})
|
|
ingestor_manager.initialize_opc_from_config.assert_any_call({'config': 'config2'})
|
|
ingestor_manager.initialize_opc_from_config.assert_any_call({'config': 'config5'})
|
|
assert ingestor_manager.initialize_opc_from_config.call_count == 3
|
|
|
|
assert ingestor_manager.opc_managers['server1'].config == {'config': 'config1'}
|
|
assert ingestor_manager.opc_managers['server2'].config == {'config': 'config2'}
|
|
assert ingestor_manager.opc_managers['server3'].config == {'config': 'config3'}
|
|
assert 'server4' not in ingestor_manager.opc_managers
|
|
assert 'server5' not in ingestor_manager.opc_managers
|
|
|
|
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)
|
|
)
|
|
|
|
|
|
def test_declare_active(ingestor_manager):
|
|
ingestor_manager.resource_manager.ingestor_heartbeat = MagicMock()
|
|
ingestor_manager.declare_active()
|
|
ingestor_manager.resource_manager.ingestor_heartbeat.assert_called_once()
|
|
|
|
|
|
def test_get_active_ingestors(ingestor_manager):
|
|
ingestor_manager.resource_manager.get_all_ingestors = MagicMock()
|
|
ingestor_manager.get_active_ingestors()
|
|
ingestor_manager.resource_manager.get_all_ingestors.assert_called_once()
|
|
|
|
|
|
def test_get_active_ingestors_empty(ingestor_manager):
|
|
ingestor_manager.resource_manager.get_all_ingestors = MagicMock(return_value=None)
|
|
result = ingestor_manager.get_active_ingestors()
|
|
assert result == []
|
|
ingestor_manager.resource_manager.get_all_ingestors.assert_called_once()
|
|
|
|
|
|
@patch('ingestor.managers.ingestor_manager.metrics')
|
|
def test_get_number_of_leases_success(metrics, ingestor_manager):
|
|
ingestor_manager.resource_manager.get_all_leases = MagicMock(return_value=['lease1', 'lease2'])
|
|
result = ingestor_manager.get_number_of_leases()
|
|
assert result == 2
|
|
ingestor_manager.resource_manager.get_all_leases.assert_called_once()
|
|
metrics.LEASES_TOTAL.set.assert_called_once_with(2)
|
|
|
|
|
|
@patch('ingestor.managers.ingestor_manager.metrics')
|
|
def test_get_number_of_leases_empty(metrics, ingestor_manager):
|
|
ingestor_manager.resource_manager.get_all_leases = MagicMock(return_value=None)
|
|
result = ingestor_manager.get_number_of_leases()
|
|
assert result == 0
|
|
ingestor_manager.resource_manager.get_all_leases.assert_called_once()
|
|
metrics.LEASES_TOTAL.set.assert_called_once_with(0)
|
|
|
|
|
|
@patch('ingestor.managers.ingestor_manager.metrics')
|
|
def test_get_number_of_slots_success(metrics, ingestor_manager):
|
|
ingestor_manager.resource_manager.get_all_slots = MagicMock(return_value=['slot1', 'slot2'])
|
|
result = ingestor_manager.get_number_of_slots()
|
|
assert result == 2
|
|
ingestor_manager.resource_manager.get_all_slots.assert_called_once()
|
|
metrics.SLOTS_TOTAL.set.assert_called_once_with(2)
|
|
|
|
|
|
@patch('ingestor.managers.ingestor_manager.metrics')
|
|
def test_get_number_of_slots_empty(metrics, ingestor_manager):
|
|
ingestor_manager.resource_manager.get_all_slots = MagicMock(return_value=None)
|
|
result = ingestor_manager.get_number_of_slots()
|
|
assert result == 0
|
|
ingestor_manager.resource_manager.get_all_slots.assert_called_once()
|
|
metrics.SLOTS_TOTAL.set.assert_called_once_with(0)
|
|
|
|
|
|
def test_get_slot_leases_1_success(ingestor_manager):
|
|
ingestor_manager.resource_manager.lease_tag = MagicMock(return_value=True)
|
|
ingestor_manager.resource_manager.get_tag_slot = MagicMock(return_value={'tags': ['tag1']})
|
|
|
|
ingestor_manager.number_of_slots = 1
|
|
result = ingestor_manager.get_slot_leases()
|
|
|
|
assert result == {'1': {'tags': ['tag1']}}
|
|
|
|
|
|
def test_get_slot_leases_2_success(ingestor_manager):
|
|
ingestor_manager.resource_manager.lease_tag = MagicMock(side_effect=[True, True])
|
|
ingestor_manager.resource_manager.get_tag_slot = MagicMock(
|
|
side_effect=[{'tags': ['tag1']}, {'tags': ['tag2']}]
|
|
)
|
|
|
|
ingestor_manager.number_of_slots = 2
|
|
result = ingestor_manager.get_slot_leases(max_slots=2)
|
|
|
|
assert result == {'1': {'tags': ['tag1']}, '2': {'tags': ['tag2']}}
|
|
|
|
|
|
def test_get_slot_leases_2_1_none(ingestor_manager):
|
|
ingestor_manager.resource_manager.lease_tag = MagicMock(side_effect=[True, True])
|
|
ingestor_manager.resource_manager.get_tag_slot = MagicMock(
|
|
side_effect=[None, {'tags': ['tag1']}]
|
|
)
|
|
|
|
ingestor_manager.number_of_slots = 1
|
|
result = ingestor_manager.get_slot_leases(max_slots=1)
|
|
|
|
assert result == {}
|
|
|
|
|
|
def test_get_slot_leases_1_failure(ingestor_manager):
|
|
ingestor_manager.resource_manager.lease_tag = MagicMock(return_value=False)
|
|
ingestor_manager.resource_manager.get_tag_slot = MagicMock(return_value={'tags': ['tag1']})
|
|
|
|
result = ingestor_manager.get_slot_leases()
|
|
ingestor_manager.resource_manager.get_tag_slot.assert_not_called()
|
|
|
|
assert result == {}
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_unsubscribe_slot(ingestor_manager):
|
|
ingestor_manager.managed_tags = {
|
|
'slot1': {'server1': {'tags': 'config1'}, 'server2': {'tags': 'config2'}},
|
|
'slot2': {'server3': {'tags': 'config3'}, 'server1': {'tags': 'config1'}},
|
|
}
|
|
ingestor_manager.opc_managers = {
|
|
'server1': AsyncMock(),
|
|
'server2': AsyncMock(),
|
|
'server3': AsyncMock(),
|
|
}
|
|
await ingestor_manager.unsubscribe_slot('slot1')
|
|
|
|
ingestor_manager.opc_managers['server1'].unsubscribe.assert_called_once_with('slot1')
|
|
ingestor_manager.opc_managers['server2'].unsubscribe.assert_called_once_with('slot1')
|
|
ingestor_manager.opc_managers['server3'].unsubscribe.assert_not_called()
|
|
|
|
|
|
def test_update_slot_config(ingestor_manager):
|
|
ingestor_manager.managed_tags = {
|
|
'slot1': {'config': 'old_config'},
|
|
'slot2': {'config': 'new_config'},
|
|
'slot3': {'config': 'old_config'},
|
|
}
|
|
|
|
ingestor_manager.resource_manager.get_tag_slot = MagicMock(
|
|
side_effect=[{'config': 'updated_config'}, {'config': 'new_config'}, None]
|
|
)
|
|
|
|
ingestor_manager.update_opc_servers = MagicMock()
|
|
ingestor_manager.subscribe_to_tags = MagicMock()
|
|
ingestor_manager.unsubscribe_slot = MagicMock()
|
|
|
|
ingestor_manager.update_slot_config()
|
|
|
|
assert ingestor_manager.managed_tags['slot1'] == {'config': 'updated_config'}
|
|
assert ingestor_manager.managed_tags['slot2'] == {'config': 'new_config'}
|
|
assert 'slot3' not in ingestor_manager.managed_tags
|
|
|
|
ingestor_manager.resource_manager.renew_tag_lease.assert_any_call('slot1')
|
|
ingestor_manager.resource_manager.renew_tag_lease.assert_any_call('slot3')
|
|
ingestor_manager.resource_manager.renew_tag_lease.assert_any_call('slot2')
|
|
assert ingestor_manager.resource_manager.renew_tag_lease.call_count == 3
|
|
|
|
|
|
@patch('ingestor.managers.ingestor_manager.metrics')
|
|
def test_drop_slot_leases(metrics, ingestor_manager):
|
|
ingestor_manager.resource_manager.drop_tag_lease = MagicMock()
|
|
ingestor_manager.drop_slot_leases(['1', '2'])
|
|
|
|
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.return_value.inc.assert_any_call()
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_manage_server_no_server(ingestor_manager):
|
|
ingestor_manager.opc_managers = {'server1': MagicMock(), 'server2': MagicMock()}
|
|
server_config = {'tags': 'config1'}
|
|
|
|
result = await ingestor_manager.manage_server('slot1', 'server3', server_config, server_config)
|
|
|
|
assert result == 1
|
|
ingestor_manager.opc_managers['server1'].create_subscription.assert_not_called()
|
|
ingestor_manager.opc_managers['server1'].subscribe.assert_not_called()
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_manage_server_create_subscription_failure(ingestor_manager):
|
|
ingestor_manager.opc_managers = {'server1': MagicMock(), 'server2': MagicMock()}
|
|
ingestor_manager.subscriptions = {'server1': MagicMock()}
|
|
server_config = {'tags': 'config1'}
|
|
|
|
ingestor_manager.opc_managers['server1'].create_subscription.side_effect = Exception(
|
|
'Subscription error'
|
|
)
|
|
|
|
result = await ingestor_manager.manage_server('slot1', 'server1', server_config, server_config)
|
|
|
|
assert result == 2
|
|
ingestor_manager.opc_managers['server1'].create_subscription.assert_called_once_with('slot1')
|
|
ingestor_manager.opc_managers['server1'].subscribe.assert_not_called()
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_manage_server(ingestor_manager):
|
|
ingestor_manager.opc_managers = {'server1': AsyncMock(), 'server2': AsyncMock()}
|
|
ingestor_manager.subscriptions = {'server1': AsyncMock()}
|
|
server_config = {'tags': 'config1'}
|
|
|
|
result = await ingestor_manager.manage_server('slot1', 'server1', server_config, server_config)
|
|
|
|
assert result == 0
|
|
ingestor_manager.opc_managers['server1'].create_subscription.assert_called_once_with('slot1')
|
|
ingestor_manager.opc_managers['server1'].subscribe.assert_called_once_with(
|
|
'slot1', 'config1', ingestor_manager.poll_interval
|
|
)
|
|
|
|
|
|
@patch('ingestor.managers.ingestor_manager.traceback')
|
|
@mark.asyncio
|
|
async def test_manage_server_subscribe_failure(traceback_mock, ingestor_manager):
|
|
ingestor_manager.opc_managers = {'server1': AsyncMock(), 'server2': AsyncMock()}
|
|
ingestor_manager.subscriptions = {'server1': AsyncMock()}
|
|
server_config = {'tags': 'config1'}
|
|
|
|
ingestor_manager.opc_managers['server1'].subscribe.side_effect = Exception('Subscription error')
|
|
|
|
result = await ingestor_manager.manage_server('slot1', 'server1', server_config, server_config)
|
|
|
|
assert result == 2
|
|
ingestor_manager.opc_managers['server1'].create_subscription.assert_called_once_with('slot1')
|
|
ingestor_manager.opc_managers['server1'].subscribe.assert_called_once_with(
|
|
'slot1', 'config1', ingestor_manager.poll_interval
|
|
)
|
|
ingestor_manager.opc_managers['server1'].unsubscribe.assert_called_once_with('slot1')
|
|
|
|
traceback_mock.format_exc.assert_called_once()
|
|
|
|
ingestor_manager.send_notification.assert_called_once_with(
|
|
metadata=metadata['metadata'],
|
|
notification_id='OPC_SUBSCRIPTION_ERROR_slot1:server1',
|
|
message="Failed to subscribe to tags from slot1:server1\n{'tags': 'config1'}: Subscription error",
|
|
block='opc_manager',
|
|
level=NotificationLevel.ERROR,
|
|
attachment_content=traceback_mock.format_exc.return_value,
|
|
)
|
|
|
|
ingestor_manager.logger.warning.assert_any_call(
|
|
'Removing subscription from server server1 for slot slot1'
|
|
)
|
|
|
|
|
|
@mark.asyncio
|
|
async def test_subscribe_to_tags(ingestor_manager):
|
|
ingestor_manager.manage_server = AsyncMock(side_effect=[0, 1, 2])
|
|
ingestor_manager.managed_tags = {'slot1': MagicMock(), 'slot2': MagicMock()}
|
|
|
|
ingestor_manager.opc_managers = {'server1': AsyncMock(), 'server2': AsyncMock()}
|
|
ingestor_manager.subscriptions = {'server1': AsyncMock()}
|
|
tags = {
|
|
'slot1': {
|
|
'server1': {'tags': 'config1'},
|
|
'server2': {'tags': 'config2'},
|
|
'server3': {'tags': 'config3'},
|
|
}
|
|
}
|
|
|
|
await ingestor_manager.subscribe_to_tags(tags)
|
|
|
|
ingestor_manager.manage_server.assert_any_call('slot1', 'server1', {'tags': 'config1'}, tags)
|
|
ingestor_manager.manage_server.assert_any_call('slot1', 'server2', {'tags': 'config2'}, tags)
|
|
ingestor_manager.manage_server.assert_any_call('slot1', 'server3', {'tags': 'config3'}, tags)
|
|
|
|
assert ingestor_manager.manage_server.call_count == 3
|
|
|
|
ingestor_manager.managed_tags['slot1'].pop.assert_called_once_with('server3', None)
|
|
|
|
|
|
@patch('ingestor.managers.ingestor_manager.metrics')
|
|
def test_check_opc_servers_integrity_all_healthy(metrics, ingestor_manager):
|
|
# Setup mock OPC managers
|
|
opc_manager1 = MagicMock()
|
|
opc_manager1.check_cycles.return_value = None
|
|
opc_manager1.check_opc_listenning.return_value = False
|
|
opc_manager1.config = {'config': 'config1'}
|
|
|
|
opc_manager2 = MagicMock()
|
|
opc_manager2.check_cycles.return_value = None
|
|
opc_manager2.check_opc_listenning.return_value = False
|
|
opc_manager2.config = {'config': 'config2'}
|
|
|
|
ingestor_manager.opc_managers = {'server1': opc_manager1, 'server2': opc_manager2}
|
|
|
|
# Mock the initialize_opc_from_config method
|
|
ingestor_manager.initialize_opc_from_config = MagicMock()
|
|
|
|
# Call the method
|
|
ingestor_manager.check_opc_servers_integrity()
|
|
|
|
# Verify that check_cycles and check_opc_listenning were called for each server
|
|
opc_manager1.check_cycles.assert_called_once()
|
|
opc_manager1.check_opc_listenning.assert_called_once()
|
|
opc_manager2.check_cycles.assert_called_once()
|
|
opc_manager2.check_opc_listenning.assert_called_once()
|
|
|
|
# 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)
|
|
)
|
|
|
|
|
|
def test_check_opc_servers_integrity_server_lost(ingestor_manager):
|
|
# Setup mock OPC manager that will be lost
|
|
opc_manager = MagicMock()
|
|
opc_manager.check_cycles.return_value = None
|
|
opc_manager.check_opc_listenning.return_value = True # Server is lost
|
|
opc_manager.config = {'config': 'config1'}
|
|
|
|
ingestor_manager.opc_managers = {'server1': opc_manager}
|
|
|
|
# Mock the initialize_opc_from_config method to return a new manager
|
|
new_manager = MagicMock()
|
|
ingestor_manager.initialize_opc_from_config = MagicMock(return_value=new_manager)
|
|
|
|
# Call the method
|
|
ingestor_manager.check_opc_servers_integrity()
|
|
|
|
|
|
def test_check_opc_servers_integrity_server_lost_with_tags(ingestor_manager):
|
|
# Setup mock OPC manager that will be lost
|
|
opc_manager = MagicMock()
|
|
opc_manager.check_cycles.return_value = None
|
|
opc_manager.check_opc_listenning.return_value = True # Server is lost
|
|
opc_manager.config = {'config': 'config1'}
|
|
|
|
ingestor_manager.opc_managers = {'server1': opc_manager}
|
|
|
|
# Setup managed tags
|
|
ingestor_manager.managed_tags = {
|
|
'slot1': {'server1': {'config': 'config1', 'tags': {'tag1': 'value1'}}}
|
|
}
|
|
|
|
# Mock the initialize_opc_from_config method to return a new manager
|
|
new_manager = MagicMock()
|
|
ingestor_manager.initialize_opc_from_config = MagicMock(return_value=new_manager)
|
|
|
|
# Call the method
|
|
ingestor_manager.check_opc_servers_integrity()
|