diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index a9303da..eec7116 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -14,7 +14,7 @@ import ingestor.metrics as metrics class IngestorManager(BaseActivity): def __init__(self, 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, metadata: dict, logger: Logger, notification_handler: NotificationHandler, @@ -40,7 +40,6 @@ class IngestorManager(BaseActivity): port=redis_port, lease_ttl=lease_ttl, heartbeat_ttl=heartbeat_ttl, - pod_id=pod_id, metadata=metadata, logger=logger, notification_handler=notification_handler, @@ -52,7 +51,6 @@ class IngestorManager(BaseActivity): self.managed_tags = {} self.opc_servers = {} - self.pod_id = pod_id self.metadata = metadata BaseActivity.__init__(self, logger=logger, @@ -88,7 +86,6 @@ class IngestorManager(BaseActivity): logger=self.logger, server_uri=server_config['server_uri'], notification_handler=self.notification_handler, - pod_id=self.pod_id, metadata=self.metadata, cert_path=server_config.get('cert_path'), private_key_path=server_config.get('private_key_path'), diff --git a/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index 332a2ca..b5fa903 100644 --- a/ingestor/managers/opc_manager.py +++ b/ingestor/managers/opc_manager.py @@ -11,8 +11,8 @@ import ingestor.metrics as metrics class OpcManager(BaseActivity): - def __init__(self, name: str, url: str, data_manager: DataManager, - logger: Logger, server_uri: str, notification_handler: NotificationHandler, pod_id: str, metadata: dict, + def __init__(self, name: str, url: str, data_manager: DataManager, logger: Logger, + server_uri: str, notification_handler: NotificationHandler, metadata: dict, cert_path: str = None, private_key_path: str = None, server_cert_path: str = None): self.url = url self.name = name @@ -26,18 +26,17 @@ class OpcManager(BaseActivity): self.nodes = {} self.subscriptions = {} self.data_manager = data_manager - self.pod_id = pod_id self.metadata = metadata + BaseActivity.__init__(self, logger=logger, + notification_handler=notification_handler, + set_error_counter=True) + metrics.OPC_CONNECTION_STATUS.labels( pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(0) metrics.OPC_TAGS_SUBSCRIBED.labels( 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): return f"OpcManager(name={self.name}, url={self.url}, server_uri={self.server_uri})\n" \ f"nodes={self.nodes}, subscriptions={self.subscriptions}" diff --git a/ingestor/managers/resource_manager.py b/ingestor/managers/resource_manager.py index de2d3b2..d2db351 100644 --- a/ingestor/managers/resource_manager.py +++ b/ingestor/managers/resource_manager.py @@ -16,14 +16,15 @@ class ResourceManager(BaseActivity): port: int, lease_ttl: int, heartbeat_ttl: int, - pod_id: str, metadata: dict, logger: Logger, notification_handler: NotificationHandler, username: str | None = None, password: str | None = None, ) -> None: - self.pod_id = pod_id + BaseActivity.__init__(self, logger=logger, + notification_handler=notification_handler, + set_error_counter=True) try: self.redis = Redis( host=host, @@ -35,7 +36,7 @@ class ResourceManager(BaseActivity): self.redis.ping() metrics.REDIS_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(1) 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) raise @@ -43,10 +44,6 @@ class ResourceManager(BaseActivity): self.heartbeat_ttl = heartbeat_ttl 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): """Wrapper to execute Redis operations and record metrics.""" start_time = time() diff --git a/requirements.txt b/requirements.txt index efdfaf2..5687b65 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,6 +1,5 @@ asyncua==1.1.5 redis -# git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.3.5 -../sientia-dataops-library +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.3.8 prometheus_client pymongo \ No newline at end of file diff --git a/tests/unit/managers/test_ingestor_manager.py b/tests/unit/managers/test_ingestor_manager.py index f6ace9f..df994b7 100644 --- a/tests/unit/managers/test_ingestor_manager.py +++ b/tests/unit/managers/test_ingestor_manager.py @@ -25,7 +25,6 @@ def ingestor_manager(data_manager_mock, resource_manager_mock): }, lease_ttl=60, heartbeat_ttl=60, - pod_id="test_pod", poll_interval=5, mongo_connection_string="mongodb://localhost:27017", mongo_database="sientia", @@ -54,7 +53,6 @@ def test___init__(notification_handler_mock, resource_manager_mock, data_manager }, lease_ttl=60, heartbeat_ttl=60, - pod_id="test_pod", poll_interval=5, mongo_connection_string="mongodb://localhost:27017", mongo_database="sientia", @@ -79,7 +77,6 @@ def test___init__(notification_handler_mock, resource_manager_mock, data_manager port=6379, lease_ttl=60, heartbeat_ttl=60, - pod_id="test_pod", metadata=metadata["metadata"], logger=ingestor.logger, notification_handler=ingestor.notification_handler, @@ -117,7 +114,6 @@ def test_initialize_opc_from_config(opc_manager, ingestor_manager): logger=ingestor_manager.logger, server_uri=server_config['server_uri'], notification_handler=ingestor_manager.notification_handler, - pod_id=server_config['pod_id'], cert_path=server_config['cert_path'], private_key_path=server_config['private_key_path'], server_cert_path=server_config['server_cert_path'], diff --git a/tests/unit/managers/test_opc_manager.py b/tests/unit/managers/test_opc_manager.py index 386f2ae..7cdfb80 100644 --- a/tests/unit/managers/test_opc_manager.py +++ b/tests/unit/managers/test_opc_manager.py @@ -46,13 +46,12 @@ metadata = { @fixture def raw_opc_manager(): return OpcManager( - "TestConnector", - "opc.tcp://localhost:4840", - MagicMock(), - MagicMock(), - "opc.tcp://localhost:4840", - MagicMock(), - "localhost", + name="TestConnector", + url="opc.tcp://localhost:4840", + data_manager=MagicMock(), + logger=MagicMock(), + server_uri="opc.tcp://localhost:4840", + notification_handler=MagicMock(), metadata=metadata["metadata"], ) @@ -552,7 +551,6 @@ def test_init_metrics_calls_correct_metric_methods(mock_metrics): logger=MagicMock(), server_uri="opc.tcp://init.test:4840/uri", notification_handler=MagicMock(), - pod_id="init_pod_localhost", metadata=metadata["metadata"], ) diff --git a/tests/unit/managers/test_resource_manager.py b/tests/unit/managers/test_resource_manager.py index 0e528a2..00238ea 100644 --- a/tests/unit/managers/test_resource_manager.py +++ b/tests/unit/managers/test_resource_manager.py @@ -22,7 +22,6 @@ def resource_manager(redis): port=6379, lease_ttl=10, heartbeat_ttl=10, - pod_id="pod_id", metadata=metadata["metadata"], logger=MagicMock(), notification_handler=MagicMock(), @@ -58,7 +57,7 @@ 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: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") assert output is True 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): - resource_manager.redis.get.return_value = "pod_id" + 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 @@ -136,13 +135,12 @@ def test_init_connection_failure(monkeypatch): port=6379, lease_ttl=10, heartbeat_ttl=10, - pod_id="pod_id", metadata=metadata["metadata"], logger=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) @@ -166,13 +164,12 @@ def test_init_ping_failure(monkeypatch): port=6379, lease_ttl=10, heartbeat_ttl=10, - pod_id="pod_id", metadata=metadata["metadata"], logger=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) @@ -198,12 +195,12 @@ def test_execute_redis_op_success(resource_manager): mock_func.assert_called_once_with("arg1", kwarg1="value1") 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_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( 0.5 @@ -226,7 +223,7 @@ def test_execute_redis_op_exception(resource_manager): # Verify metrics and error handling 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() resource_manager.send_notification.assert_called_once_with(