Enhance IngestorManager and related tests with notification handling integration - Added notification_handler parameter to IngestorManager and OpcManager initialization for improved event handling. - Updated unit tests to include notification_handler in the Ingestor and IngestorManager instantiation, ensuring proper coverage of new functionality. - Refactored test cases to validate the integration of notification handling across components.
270 lines
10 KiB
Python
270 lines
10 KiB
Python
from unittest.mock import ANY, MagicMock, patch
|
|
|
|
from pytest import fixture
|
|
from ingestor.ingestor import Ingestor
|
|
|
|
|
|
@patch("ingestor.ingestor.getenv")
|
|
@patch("ingestor.ingestor.Ingestor.init_logger")
|
|
@patch("ingestor.ingestor.NotificationHandler")
|
|
def test___init__(notification_handler, init_logger, getenv):
|
|
getenv.side_effect = [
|
|
"localhost:9092,localhost:35", # KAFKA_SERVERS
|
|
"localhost1", # REDIS_HOST
|
|
'63790', # REDIS_PORT
|
|
"user", # REDIS_USERNAME
|
|
"password", # REDIS_PASSWORD
|
|
'100', # LEASE_TTL
|
|
'200', # HEARTBEAT_TTL
|
|
"localhost1", # HOSTNAME
|
|
'50' # POLL_INTERVAL
|
|
]
|
|
|
|
ingestor = Ingestor()
|
|
|
|
getenv.assert_any_call("KAFKA_SERVERS", "localhost:9092")
|
|
getenv.assert_any_call("REDIS_HOST", "localhost")
|
|
getenv.assert_any_call("REDIS_PORT", 6379)
|
|
getenv.assert_any_call("REDIS_USERNAME", None)
|
|
getenv.assert_any_call("REDIS_PASSWORD", None)
|
|
getenv.assert_any_call("LEASE_TTL", 10)
|
|
getenv.assert_any_call("HEARTBEAT_TTL", 20)
|
|
getenv.assert_any_call("HOSTNAME", "localhost")
|
|
getenv.assert_any_call("POLL_INTERVAL", 5)
|
|
|
|
assert ingestor.kafka_servers == ["localhost:9092", "localhost:35"]
|
|
assert ingestor.redis_host == "localhost1"
|
|
assert ingestor.redis_port == 63790
|
|
assert ingestor.redis_username == "user"
|
|
assert ingestor.redis_password == "password"
|
|
assert ingestor.lease_ttl == 100
|
|
assert ingestor.heartbeat_ttl == 200
|
|
assert ingestor.pod_id == "localhost1"
|
|
assert ingestor.poll_interval == 50
|
|
|
|
init_logger.assert_called_once()
|
|
notification_handler.assert_called_once_with(
|
|
servers=["localhost:9092", "localhost:35"],
|
|
logger=ingestor.logger,
|
|
project_name="OPC_INGESTOR",
|
|
pipeline_name="-",
|
|
trigger_name="-",
|
|
model_name="-",
|
|
model="-"
|
|
)
|
|
|
|
|
|
@fixture
|
|
@patch("ingestor.ingestor.getenv")
|
|
@patch("ingestor.ingestor.Ingestor.init_logger")
|
|
@patch("ingestor.ingestor.NotificationHandler")
|
|
def ingestor(notification_handler, init_logger, getenv):
|
|
ing = Ingestor()
|
|
ing.logger = MagicMock()
|
|
|
|
return ing
|
|
|
|
|
|
@fixture
|
|
def ingestor_manager_started(ingestor):
|
|
ingestor.ingestor_manager = MagicMock()
|
|
return ingestor
|
|
|
|
|
|
@patch("ingestor.ingestor.getLogger")
|
|
@patch("ingestor.ingestor.StreamHandler")
|
|
@patch("ingestor.ingestor.Formatter")
|
|
def test_init_logger(formatter, stream_handler, get_logger, ingestor):
|
|
ingestor.logger = None
|
|
ingestor.init_logger()
|
|
|
|
get_logger.assert_called_once_with('ingestor.ingestor')
|
|
stream_handler.assert_called_once()
|
|
formatter.assert_called_once_with(
|
|
'%(asctime)s - %(name)s - %(levelname)s - %(message)s')
|
|
ingestor.logger.setLevel.assert_called_once_with("INFO")
|
|
ingestor.logger.addHandler.assert_called_once_with(
|
|
stream_handler.return_value)
|
|
stream_handler.return_value.setFormatter.assert_called_once_with(
|
|
formatter.return_value)
|
|
|
|
|
|
def test_handle_acquired_tags_not_acquired(ingestor_manager_started):
|
|
ingestor_manager_started.handle_acquired_tags([])
|
|
|
|
ingestor_manager_started.logger.warning.assert_called_once_with(
|
|
"No slots available")
|
|
ingestor_manager_started.ingestor_manager.update_opc_servers.assert_not_called()
|
|
ingestor_manager_started.ingestor_manager.subscribe_to_tags.assert_not_called()
|
|
|
|
|
|
def test_handle_acquired_tags_success(ingestor_manager_started):
|
|
ingestor_manager_started.handle_acquired_tags(["tag1", "tag2"])
|
|
|
|
ingestor_manager_started.logger.warning.assert_not_called()
|
|
ingestor_manager_started.ingestor_manager.update_opc_servers.assert_called_once()
|
|
ingestor_manager_started.ingestor_manager.subscribe_to_tags.assert_called_once_with(
|
|
["tag1", "tag2"])
|
|
|
|
|
|
@patch("ingestor.ingestor.IngestorManager")
|
|
def test_prepare_ingestor(ingestor_manager_mock, ingestor):
|
|
ingestor_manager = ingestor_manager_mock.return_value
|
|
ingestor_manager.get_slot_leases.return_value = True
|
|
|
|
ingestor.prepare_ingestor()
|
|
|
|
ingestor_manager_mock.assert_called_once_with(
|
|
ingestor.kafka_servers,
|
|
ingestor.redis_host,
|
|
ingestor.redis_port,
|
|
ingestor.lease_ttl,
|
|
ingestor.heartbeat_ttl,
|
|
ingestor.pod_id,
|
|
ingestor.poll_interval,
|
|
ingestor.logger,
|
|
ingestor.notification_handler,
|
|
ingestor.redis_username,
|
|
ingestor.redis_password
|
|
)
|
|
ingestor_manager.declare_active.assert_called_once()
|
|
ingestor_manager.get_slot_leases.assert_called_once()
|
|
|
|
ingestor.handle_acquired_tags(
|
|
ingestor_manager.get_slot_leases.return_value)
|
|
|
|
|
|
def test_manage_slots_has_slots(ingestor_manager_started):
|
|
ingestor_manager_started.handle_acquired_tags = MagicMock()
|
|
ingestor_manager_started.ingestor_manager.managed_tags = True
|
|
|
|
ingestor_manager_started.manage_no_slots(5)
|
|
|
|
ingestor_manager_started.ingestor_manager.get_slot_leases.assert_not_called()
|
|
ingestor_manager_started.ingestor_manager.handle_acquired_tags.assert_not_called()
|
|
|
|
|
|
def test_manage_no_slots_has_slots_none_available(ingestor_manager_started):
|
|
ingestor_manager_started.handle_acquired_tags = MagicMock()
|
|
ingestor_manager_started.ingestor_manager.managed_tags = True
|
|
|
|
ingestor_manager_started.manage_no_slots(0)
|
|
|
|
ingestor_manager_started.ingestor_manager.get_slot_leases.assert_not_called()
|
|
ingestor_manager_started.ingestor_manager.handle_acquired_tags.assert_not_called()
|
|
|
|
|
|
def test_manage_slots_none_available(ingestor_manager_started):
|
|
ingestor_manager_started.handle_acquired_tags = MagicMock()
|
|
ingestor_manager_started.ingestor_manager.managed_tags = False
|
|
|
|
ingestor_manager_started.manage_no_slots(0)
|
|
|
|
ingestor_manager_started.ingestor_manager.get_slot_leases.assert_not_called()
|
|
ingestor_manager_started.ingestor_manager.handle_acquired_tags.assert_not_called()
|
|
|
|
|
|
def test_manage_slots_none_available_none_available(ingestor_manager_started):
|
|
ingestor_manager_started.handle_acquired_tags = MagicMock()
|
|
ingestor_manager_started.ingestor_manager.managed_tags = False
|
|
|
|
ingestor_manager_started.manage_no_slots(2)
|
|
|
|
ingestor_manager_started.ingestor_manager.get_slot_leases.assert_called_once_with(
|
|
1)
|
|
ingestor_manager_started.handle_acquired_tags.assert_called_once_with(
|
|
ingestor_manager_started.ingestor_manager.get_slot_leases.return_value)
|
|
|
|
|
|
def test_manage_leases_no_available_slots_no_extra_slots(ingestor_manager_started):
|
|
ingestor_manager_started.handle_acquired_tags = MagicMock()
|
|
|
|
ingestor_manager_started.manage_leases(0, 0, 0)
|
|
|
|
ingestor_manager_started.ingestor_manager.get_slot_leases.assert_not_called()
|
|
ingestor_manager_started.ingestor_manager.handle_acquired_tags.assert_not_called()
|
|
ingestor_manager_started.ingestor_manager.drop_slot_leases.assert_not_called()
|
|
|
|
|
|
def test_manage_leases_available_slots_innactive_ingestors(ingestor_manager_started):
|
|
ingestor_manager_started.handle_acquired_tags = MagicMock()
|
|
|
|
ingestor_manager_started.manage_leases(2, 2, 5)
|
|
|
|
ingestor_manager_started.ingestor_manager.get_slot_leases.assert_called_once_with(
|
|
2)
|
|
ingestor_manager_started.handle_acquired_tags.assert_called_once_with(
|
|
ingestor_manager_started.ingestor_manager.get_slot_leases.return_value)
|
|
|
|
ingestor_manager_started.ingestor_manager.drop_slot_leases.assert_not_called()
|
|
|
|
|
|
def test_manage_leases_no_available_slots_extra_sltos(ingestor_manager_started):
|
|
ingestor_manager_started.handle_acquired_tags = MagicMock()
|
|
ingestor_manager_started.ingestor_manager.managed_tags = {
|
|
"tag1": "server1",
|
|
"tag2": "server2",
|
|
"tag3": "server3"
|
|
}
|
|
|
|
ingestor_manager_started.manage_leases(0, 0, 2)
|
|
|
|
ingestor_manager_started.ingestor_manager.get_slot_leases.assert_not_called()
|
|
ingestor_manager_started.handle_acquired_tags.assert_not_called()
|
|
ingestor_manager_started.ingestor_manager.drop_slot_leases.assert_called_once_with(
|
|
["tag2", "tag3"])
|
|
|
|
|
|
def test_loop(ingestor_manager_started):
|
|
ingestor_manager_started.manage_no_slots = MagicMock()
|
|
ingestor_manager_started.manage_leases = MagicMock()
|
|
ingestor_manager_started.ingestor_manager.managed_tags = {
|
|
"slot1": "server1",
|
|
"slot2": "server2",
|
|
"slot3": "server3"
|
|
}
|
|
ingestor_manager_started.ingestor_manager.get_active_ingestors = MagicMock(
|
|
return_value=["ingestor1", "ingestor2"])
|
|
ingestor_manager_started.ingestor_manager.get_number_of_slots = MagicMock(
|
|
return_value=5)
|
|
ingestor_manager_started.ingestor_manager.get_number_of_leases = MagicMock(
|
|
return_value=1)
|
|
|
|
ingestor_manager_started.loop()
|
|
|
|
ingestor_manager_started.ingestor_manager.declare_active.assert_called_once()
|
|
ingestor_manager_started.ingestor_manager.get_active_ingestors.assert_called_once()
|
|
ingestor_manager_started.ingestor_manager.get_number_of_slots.assert_called_once()
|
|
|
|
ingestor_manager_started.manage_no_slots.assert_called_once_with(
|
|
ingestor_manager_started.ingestor_manager.get_number_of_slots.return_value)
|
|
# Explanation: 5 - 2 = 3, 3 - 1 = 2
|
|
ingestor_manager_started.manage_leases.assert_called_once_with(
|
|
4, 3, 2)
|
|
ingestor_manager_started.ingestor_manager.update_slot_config.assert_called_once()
|
|
|
|
|
|
def test_loop_no_managed(ingestor_manager_started):
|
|
ingestor_manager_started.manage_no_slots = MagicMock()
|
|
ingestor_manager_started.manage_leases = MagicMock()
|
|
ingestor_manager_started.ingestor_manager.managed_tags = {}
|
|
ingestor_manager_started.ingestor_manager.get_active_ingestors = MagicMock(
|
|
return_value=["ingestor1", "ingestor2"])
|
|
ingestor_manager_started.ingestor_manager.get_number_of_slots = MagicMock(
|
|
return_value=5)
|
|
|
|
ingestor_manager_started.loop()
|
|
|
|
ingestor_manager_started.ingestor_manager.declare_active.assert_called_once()
|
|
ingestor_manager_started.ingestor_manager.get_active_ingestors.assert_called_once()
|
|
ingestor_manager_started.ingestor_manager.get_number_of_slots.assert_called_once()
|
|
|
|
ingestor_manager_started.manage_no_slots.assert_called_once_with(
|
|
ingestor_manager_started.ingestor_manager.get_number_of_slots.return_value)
|
|
# Explanation: 5 - 2 = 3, 3 - 1 = 2
|
|
ingestor_manager_started.manage_leases.assert_called_once_with(
|
|
ANY, 3, -1)
|
|
ingestor_manager_started.ingestor_manager.update_slot_config.assert_called_once()
|
|
ingestor_manager_started.logger.info.assert_any_call(
|
|
"No slots acquired in this loop")
|