SIENTIAPDE-1174
Update requirements and refactor manager classes to remove pod_id - Updated the sientia-dataops-library dependency version to 1.3.8 in requirements.txt. - Refactored IngestorManager, OpcManager, and ResourceManager classes to remove pod_id from their initialization and internal handling, enhancing code clarity and consistency. - Adjusted unit tests to reflect the removal of pod_id, ensuring they remain functional and accurate.
This commit is contained in:
@@ -14,7 +14,7 @@ import ingestor.metrics as metrics
|
|||||||
class IngestorManager(BaseActivity):
|
class IngestorManager(BaseActivity):
|
||||||
def __init__(self,
|
def __init__(self,
|
||||||
kafka_servers: str, redis_data: dict,
|
kafka_servers: str, redis_data: dict,
|
||||||
lease_ttl: int, heartbeat_ttl: int, pod_id: str,
|
lease_ttl: int, heartbeat_ttl: int,
|
||||||
poll_interval: int, mongo_connection_string: str, mongo_database: str,
|
poll_interval: int, mongo_connection_string: str, mongo_database: str,
|
||||||
metadata: dict,
|
metadata: dict,
|
||||||
logger: Logger, notification_handler: NotificationHandler,
|
logger: Logger, notification_handler: NotificationHandler,
|
||||||
@@ -40,7 +40,6 @@ class IngestorManager(BaseActivity):
|
|||||||
port=redis_port,
|
port=redis_port,
|
||||||
lease_ttl=lease_ttl,
|
lease_ttl=lease_ttl,
|
||||||
heartbeat_ttl=heartbeat_ttl,
|
heartbeat_ttl=heartbeat_ttl,
|
||||||
pod_id=pod_id,
|
|
||||||
metadata=metadata,
|
metadata=metadata,
|
||||||
logger=logger,
|
logger=logger,
|
||||||
notification_handler=notification_handler,
|
notification_handler=notification_handler,
|
||||||
@@ -52,7 +51,6 @@ class IngestorManager(BaseActivity):
|
|||||||
self.managed_tags = {}
|
self.managed_tags = {}
|
||||||
self.opc_servers = {}
|
self.opc_servers = {}
|
||||||
|
|
||||||
self.pod_id = pod_id
|
|
||||||
self.metadata = metadata
|
self.metadata = metadata
|
||||||
|
|
||||||
BaseActivity.__init__(self, logger=logger,
|
BaseActivity.__init__(self, logger=logger,
|
||||||
@@ -88,7 +86,6 @@ class IngestorManager(BaseActivity):
|
|||||||
logger=self.logger,
|
logger=self.logger,
|
||||||
server_uri=server_config['server_uri'],
|
server_uri=server_config['server_uri'],
|
||||||
notification_handler=self.notification_handler,
|
notification_handler=self.notification_handler,
|
||||||
pod_id=self.pod_id,
|
|
||||||
metadata=self.metadata,
|
metadata=self.metadata,
|
||||||
cert_path=server_config.get('cert_path'),
|
cert_path=server_config.get('cert_path'),
|
||||||
private_key_path=server_config.get('private_key_path'),
|
private_key_path=server_config.get('private_key_path'),
|
||||||
|
|||||||
@@ -11,8 +11,8 @@ import ingestor.metrics as metrics
|
|||||||
|
|
||||||
|
|
||||||
class OpcManager(BaseActivity):
|
class OpcManager(BaseActivity):
|
||||||
def __init__(self, name: str, url: str, data_manager: DataManager,
|
def __init__(self, name: str, url: str, data_manager: DataManager, logger: Logger,
|
||||||
logger: Logger, server_uri: str, notification_handler: NotificationHandler, pod_id: str, metadata: dict,
|
server_uri: str, notification_handler: NotificationHandler, metadata: dict,
|
||||||
cert_path: str = None, private_key_path: str = None, server_cert_path: str = None):
|
cert_path: str = None, private_key_path: str = None, server_cert_path: str = None):
|
||||||
self.url = url
|
self.url = url
|
||||||
self.name = name
|
self.name = name
|
||||||
@@ -26,18 +26,17 @@ class OpcManager(BaseActivity):
|
|||||||
self.nodes = {}
|
self.nodes = {}
|
||||||
self.subscriptions = {}
|
self.subscriptions = {}
|
||||||
self.data_manager = data_manager
|
self.data_manager = data_manager
|
||||||
self.pod_id = pod_id
|
|
||||||
self.metadata = metadata
|
self.metadata = metadata
|
||||||
|
|
||||||
|
BaseActivity.__init__(self, logger=logger,
|
||||||
|
notification_handler=notification_handler,
|
||||||
|
set_error_counter=True)
|
||||||
|
|
||||||
metrics.OPC_CONNECTION_STATUS.labels(
|
metrics.OPC_CONNECTION_STATUS.labels(
|
||||||
pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(0)
|
pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(0)
|
||||||
metrics.OPC_TAGS_SUBSCRIBED.labels(
|
metrics.OPC_TAGS_SUBSCRIBED.labels(
|
||||||
pod_id=self.pod_id, server_name=self.name).set(0)
|
pod_id=self.pod_id, server_name=self.name).set(0)
|
||||||
|
|
||||||
BaseActivity.__init__(self, logger=logger,
|
|
||||||
notification_handler=notification_handler,
|
|
||||||
set_error_counter=True)
|
|
||||||
|
|
||||||
def __str__(self):
|
def __str__(self):
|
||||||
return f"OpcManager(name={self.name}, url={self.url}, server_uri={self.server_uri})\n" \
|
return f"OpcManager(name={self.name}, url={self.url}, server_uri={self.server_uri})\n" \
|
||||||
f"nodes={self.nodes}, subscriptions={self.subscriptions}"
|
f"nodes={self.nodes}, subscriptions={self.subscriptions}"
|
||||||
|
|||||||
@@ -16,14 +16,15 @@ class ResourceManager(BaseActivity):
|
|||||||
port: int,
|
port: int,
|
||||||
lease_ttl: int,
|
lease_ttl: int,
|
||||||
heartbeat_ttl: int,
|
heartbeat_ttl: int,
|
||||||
pod_id: str,
|
|
||||||
metadata: dict,
|
metadata: dict,
|
||||||
logger: Logger,
|
logger: Logger,
|
||||||
notification_handler: NotificationHandler,
|
notification_handler: NotificationHandler,
|
||||||
username: str | None = None,
|
username: str | None = None,
|
||||||
password: str | None = None,
|
password: str | None = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
self.pod_id = pod_id
|
BaseActivity.__init__(self, logger=logger,
|
||||||
|
notification_handler=notification_handler,
|
||||||
|
set_error_counter=True)
|
||||||
try:
|
try:
|
||||||
self.redis = Redis(
|
self.redis = Redis(
|
||||||
host=host,
|
host=host,
|
||||||
@@ -35,7 +36,7 @@ class ResourceManager(BaseActivity):
|
|||||||
self.redis.ping()
|
self.redis.ping()
|
||||||
metrics.REDIS_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(1)
|
metrics.REDIS_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(1)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Failed to connect to Redis: {e}")
|
self.logger.error(f"Failed to connect to Redis: {e}")
|
||||||
metrics.REDIS_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(0)
|
metrics.REDIS_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(0)
|
||||||
raise
|
raise
|
||||||
|
|
||||||
@@ -43,10 +44,6 @@ class ResourceManager(BaseActivity):
|
|||||||
self.heartbeat_ttl = heartbeat_ttl
|
self.heartbeat_ttl = heartbeat_ttl
|
||||||
self.metadata = metadata
|
self.metadata = metadata
|
||||||
|
|
||||||
BaseActivity.__init__(self, logger=logger,
|
|
||||||
notification_handler=notification_handler,
|
|
||||||
set_error_counter=True)
|
|
||||||
|
|
||||||
def _execute_redis_op(self, operation_name: str, func, *args, **kwargs):
|
def _execute_redis_op(self, operation_name: str, func, *args, **kwargs):
|
||||||
"""Wrapper to execute Redis operations and record metrics."""
|
"""Wrapper to execute Redis operations and record metrics."""
|
||||||
start_time = time()
|
start_time = time()
|
||||||
|
|||||||
@@ -1,6 +1,5 @@
|
|||||||
asyncua==1.1.5
|
asyncua==1.1.5
|
||||||
redis
|
redis
|
||||||
# git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.3.5
|
git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.3.8
|
||||||
../sientia-dataops-library
|
|
||||||
prometheus_client
|
prometheus_client
|
||||||
pymongo
|
pymongo
|
||||||
@@ -25,7 +25,6 @@ def ingestor_manager(data_manager_mock, resource_manager_mock):
|
|||||||
},
|
},
|
||||||
lease_ttl=60,
|
lease_ttl=60,
|
||||||
heartbeat_ttl=60,
|
heartbeat_ttl=60,
|
||||||
pod_id="test_pod",
|
|
||||||
poll_interval=5,
|
poll_interval=5,
|
||||||
mongo_connection_string="mongodb://localhost:27017",
|
mongo_connection_string="mongodb://localhost:27017",
|
||||||
mongo_database="sientia",
|
mongo_database="sientia",
|
||||||
@@ -54,7 +53,6 @@ def test___init__(notification_handler_mock, resource_manager_mock, data_manager
|
|||||||
},
|
},
|
||||||
lease_ttl=60,
|
lease_ttl=60,
|
||||||
heartbeat_ttl=60,
|
heartbeat_ttl=60,
|
||||||
pod_id="test_pod",
|
|
||||||
poll_interval=5,
|
poll_interval=5,
|
||||||
mongo_connection_string="mongodb://localhost:27017",
|
mongo_connection_string="mongodb://localhost:27017",
|
||||||
mongo_database="sientia",
|
mongo_database="sientia",
|
||||||
@@ -79,7 +77,6 @@ def test___init__(notification_handler_mock, resource_manager_mock, data_manager
|
|||||||
port=6379,
|
port=6379,
|
||||||
lease_ttl=60,
|
lease_ttl=60,
|
||||||
heartbeat_ttl=60,
|
heartbeat_ttl=60,
|
||||||
pod_id="test_pod",
|
|
||||||
metadata=metadata["metadata"],
|
metadata=metadata["metadata"],
|
||||||
logger=ingestor.logger,
|
logger=ingestor.logger,
|
||||||
notification_handler=ingestor.notification_handler,
|
notification_handler=ingestor.notification_handler,
|
||||||
@@ -117,7 +114,6 @@ def test_initialize_opc_from_config(opc_manager, ingestor_manager):
|
|||||||
logger=ingestor_manager.logger,
|
logger=ingestor_manager.logger,
|
||||||
server_uri=server_config['server_uri'],
|
server_uri=server_config['server_uri'],
|
||||||
notification_handler=ingestor_manager.notification_handler,
|
notification_handler=ingestor_manager.notification_handler,
|
||||||
pod_id=server_config['pod_id'],
|
|
||||||
cert_path=server_config['cert_path'],
|
cert_path=server_config['cert_path'],
|
||||||
private_key_path=server_config['private_key_path'],
|
private_key_path=server_config['private_key_path'],
|
||||||
server_cert_path=server_config['server_cert_path'],
|
server_cert_path=server_config['server_cert_path'],
|
||||||
|
|||||||
@@ -46,13 +46,12 @@ metadata = {
|
|||||||
@fixture
|
@fixture
|
||||||
def raw_opc_manager():
|
def raw_opc_manager():
|
||||||
return OpcManager(
|
return OpcManager(
|
||||||
"TestConnector",
|
name="TestConnector",
|
||||||
"opc.tcp://localhost:4840",
|
url="opc.tcp://localhost:4840",
|
||||||
MagicMock(),
|
data_manager=MagicMock(),
|
||||||
MagicMock(),
|
logger=MagicMock(),
|
||||||
"opc.tcp://localhost:4840",
|
server_uri="opc.tcp://localhost:4840",
|
||||||
MagicMock(),
|
notification_handler=MagicMock(),
|
||||||
"localhost",
|
|
||||||
metadata=metadata["metadata"],
|
metadata=metadata["metadata"],
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -552,7 +551,6 @@ def test_init_metrics_calls_correct_metric_methods(mock_metrics):
|
|||||||
logger=MagicMock(),
|
logger=MagicMock(),
|
||||||
server_uri="opc.tcp://init.test:4840/uri",
|
server_uri="opc.tcp://init.test:4840/uri",
|
||||||
notification_handler=MagicMock(),
|
notification_handler=MagicMock(),
|
||||||
pod_id="init_pod_localhost",
|
|
||||||
metadata=metadata["metadata"],
|
metadata=metadata["metadata"],
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -22,7 +22,6 @@ def resource_manager(redis):
|
|||||||
port=6379,
|
port=6379,
|
||||||
lease_ttl=10,
|
lease_ttl=10,
|
||||||
heartbeat_ttl=10,
|
heartbeat_ttl=10,
|
||||||
pod_id="pod_id",
|
|
||||||
metadata=metadata["metadata"],
|
metadata=metadata["metadata"],
|
||||||
logger=MagicMock(),
|
logger=MagicMock(),
|
||||||
notification_handler=MagicMock(),
|
notification_handler=MagicMock(),
|
||||||
@@ -58,7 +57,7 @@ def test_ingestor_heartbeat(resource_manager):
|
|||||||
resource_manager.redis.set.return_value = True
|
resource_manager.redis.set.return_value = True
|
||||||
resource_manager.ingestor_heartbeat()
|
resource_manager.ingestor_heartbeat()
|
||||||
resource_manager.redis.set.assert_called_once_with(
|
resource_manager.redis.set.assert_called_once_with(
|
||||||
"heartbeat:ingestor:pod_id", 1, ex=10
|
"heartbeat:ingestor:localhost", 1, ex=10
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@@ -67,12 +66,12 @@ def test_lease_tag(resource_manager):
|
|||||||
output = resource_manager.lease_tag("tag_id")
|
output = resource_manager.lease_tag("tag_id")
|
||||||
assert output is True
|
assert output is True
|
||||||
resource_manager.redis.set.assert_called_once_with(
|
resource_manager.redis.set.assert_called_once_with(
|
||||||
"lease:opc_tags:tag_id", "pod_id", nx=True, ex=10
|
"lease:opc_tags:tag_id", "localhost", nx=True, ex=10
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def test_renew_tag_lease_success(resource_manager):
|
def test_renew_tag_lease_success(resource_manager):
|
||||||
resource_manager.redis.get.return_value = "pod_id"
|
resource_manager.redis.get.return_value = "localhost"
|
||||||
resource_manager.redis.expire.return_value = True
|
resource_manager.redis.expire.return_value = True
|
||||||
result = resource_manager.renew_tag_lease("tag_id")
|
result = resource_manager.renew_tag_lease("tag_id")
|
||||||
assert result is True
|
assert result is True
|
||||||
@@ -136,13 +135,12 @@ def test_init_connection_failure(monkeypatch):
|
|||||||
port=6379,
|
port=6379,
|
||||||
lease_ttl=10,
|
lease_ttl=10,
|
||||||
heartbeat_ttl=10,
|
heartbeat_ttl=10,
|
||||||
pod_id="pod_id",
|
|
||||||
metadata=metadata["metadata"],
|
metadata=metadata["metadata"],
|
||||||
logger=MagicMock(),
|
logger=MagicMock(),
|
||||||
notification_handler=MagicMock(),
|
notification_handler=MagicMock(),
|
||||||
)
|
)
|
||||||
|
|
||||||
mock_metrics.labels.assert_called_once_with(pod_id="pod_id")
|
mock_metrics.labels.assert_called_once_with(pod_id="localhost")
|
||||||
mock_status.set.assert_called_once_with(0)
|
mock_status.set.assert_called_once_with(0)
|
||||||
|
|
||||||
|
|
||||||
@@ -166,13 +164,12 @@ def test_init_ping_failure(monkeypatch):
|
|||||||
port=6379,
|
port=6379,
|
||||||
lease_ttl=10,
|
lease_ttl=10,
|
||||||
heartbeat_ttl=10,
|
heartbeat_ttl=10,
|
||||||
pod_id="pod_id",
|
|
||||||
metadata=metadata["metadata"],
|
metadata=metadata["metadata"],
|
||||||
logger=MagicMock(),
|
logger=MagicMock(),
|
||||||
notification_handler=MagicMock(),
|
notification_handler=MagicMock(),
|
||||||
)
|
)
|
||||||
|
|
||||||
mock_metrics.labels.assert_called_once_with(pod_id="pod_id")
|
mock_metrics.labels.assert_called_once_with(pod_id="localhost")
|
||||||
mock_status.set.assert_called_once_with(0)
|
mock_status.set.assert_called_once_with(0)
|
||||||
|
|
||||||
|
|
||||||
@@ -198,12 +195,12 @@ def test_execute_redis_op_success(resource_manager):
|
|||||||
mock_func.assert_called_once_with("arg1", kwarg1="value1")
|
mock_func.assert_called_once_with("arg1", kwarg1="value1")
|
||||||
|
|
||||||
mock_total.labels.assert_called_once_with(
|
mock_total.labels.assert_called_once_with(
|
||||||
pod_id="pod_id", operation="test_op"
|
pod_id="localhost", operation="test_op"
|
||||||
)
|
)
|
||||||
mock_total_labels.inc.assert_called_once()
|
mock_total_labels.inc.assert_called_once()
|
||||||
|
|
||||||
mock_duration.labels.assert_called_once_with(
|
mock_duration.labels.assert_called_once_with(
|
||||||
pod_id="pod_id", operation="test_op"
|
pod_id="localhost", operation="test_op"
|
||||||
)
|
)
|
||||||
mock_duration_labels.observe.assert_called_once_with(
|
mock_duration_labels.observe.assert_called_once_with(
|
||||||
0.5
|
0.5
|
||||||
@@ -226,7 +223,7 @@ def test_execute_redis_op_exception(resource_manager):
|
|||||||
|
|
||||||
# Verify metrics and error handling
|
# Verify metrics and error handling
|
||||||
mock_errors.labels.assert_called_once_with(
|
mock_errors.labels.assert_called_once_with(
|
||||||
pod_id="pod_id", operation="test_op"
|
pod_id="localhost", operation="test_op"
|
||||||
)
|
)
|
||||||
mock_errors_labels.inc.assert_called_once()
|
mock_errors_labels.inc.assert_called_once()
|
||||||
resource_manager.send_notification.assert_called_once_with(
|
resource_manager.send_notification.assert_called_once_with(
|
||||||
|
|||||||
Reference in New Issue
Block a user