diff --git a/.vscode/settings.json b/.vscode/settings.json index a3a53a3..a595a07 100644 --- a/.vscode/settings.json +++ b/.vscode/settings.json @@ -7,6 +7,6 @@ "projectKey": "Aignosi_sientia-dataops-opc-ingestor_1642c8bf-a148-4911-8362-0903d7fef99a" }, "python.languageServer": "Pylance", - "python.analysis.typeCheckingMode": "basic", + "python.analysis.typeCheckingMode": "standard", "editor.suggestSelection": "first" } diff --git a/ingestor/managers/resource_manager.py b/ingestor/managers/resource_manager.py index 2c30198..d958d85 100644 --- a/ingestor/managers/resource_manager.py +++ b/ingestor/managers/resource_manager.py @@ -1,17 +1,59 @@ import json from typing import List from redis import Redis +from time import time +import ingestor.metrics as metrics class ResourceManager: - def __init__(self, host: str, port: int, - lease_ttl: int, heartbeat_ttl: int, pod_id: str, - username: str = None, password: str = None) -> None: - self.redis = Redis(host=host, port=port, decode_responses=True, - username=username, password=password) + def __init__( + self, + host: str, + port: int, + lease_ttl: int, + heartbeat_ttl: int, + pod_id: str, + username: str | None = None, + password: str | None = None, + ) -> None: + self.pod_id = pod_id + try: + self.redis = Redis( + host=host, + port=port, + decode_responses=True, + username=username, + password=password, + ) + self.redis.ping() + metrics.REDIS_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(1) + except Exception as e: + print(f"Failed to connect to Redis: {e}") + metrics.REDIS_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(0) + raise + self.lease_ttl = lease_ttl self.heartbeat_ttl = heartbeat_ttl - self.pod_id = pod_id + + def _execute_redis_op(self, operation_name: str, func, *args, **kwargs): + """Wrapper to execute Redis operations and record metrics.""" + start_time = time() + try: + result = func(*args, **kwargs) + metrics.REDIS_OPERATIONS_TOTAL.labels( + pod_id=self.pod_id, operation=operation_name + ).inc() + duration = time() - start_time + metrics.REDIS_OPERATIONS_DURATION.labels( + pod_id=self.pod_id, operation=operation_name + ).observe(duration) + return result + except Exception as e: + metrics.REDIS_OPERATIONS_ERRORS.labels( + pod_id=self.pod_id, operation=operation_name + ).inc() + print(f"Error in Redis operation '{operation_name}': {e}") + raise def get(self, key: str) -> dict: """ @@ -19,11 +61,11 @@ class ResourceManager: Args: key (str): The key to look up in Redis. Returns: - dict: The value associated with the key, parsed as a dictionary, + dict: The value associated with the key, parsed as a dictionary, or None if the key does not exist or the value is empty. """ - history = self.redis.get(key) + history = self._execute_redis_op("get", self.redis.get, key) return json.loads(history) if history else None def get_tag_slot(self, id: str) -> dict: @@ -48,14 +90,19 @@ class ResourceManager: None """ - self.redis.set( - f"heartbeat:ingestor:{self.pod_id}", 1, ex=self.heartbeat_ttl) + self._execute_redis_op( + "set", + self.redis.set, + f"heartbeat:ingestor:{self.pod_id}", + 1, + ex=self.heartbeat_ttl, + ) def lease_tag(self, tag_id: str) -> bool: """ Attempts to lease a tag by setting a key in Redis with a specified TTL (time-to-live). - This method uses the Redis `SET` command with the `NX` option to ensure that the key - is only set if it does not already exist. The key is set with an expiration time + This method uses the Redis `SET` command with the `NX` option to ensure that the key + is only set if it does not already exist. The key is set with an expiration time defined by `lease_ttl`. Args: tag_id (str): The unique identifier of the tag to be leased. @@ -63,8 +110,14 @@ class ResourceManager: bool: True if the lease was successfully acquired, False otherwise. """ - return self.redis.set( - f"lease:opc_tags:{tag_id}", self.pod_id, nx=True, ex=self.lease_ttl) + return self._execute_redis_op( + "set_nx", + self.redis.set, + f"lease:opc_tags:{tag_id}", + self.pod_id, + nx=True, + ex=self.lease_ttl, + ) def renew_tag_lease(self, tag_id: str) -> bool: """ @@ -78,12 +131,14 @@ class ResourceManager: bool: True if the lease was successfully renewed, False otherwise. """ - current = self.redis.get( - f"lease:opc_tags:{tag_id}") + current = self._execute_redis_op( + "get", self.redis.get, f"lease:opc_tags:{tag_id}" + ) if current == self.pod_id: - self.redis.expire(f"lease:opc_tags:{tag_id}", self.lease_ttl) + self._execute_redis_op( + "expire", self.redis.expire, f"lease:opc_tags:{tag_id}", self.lease_ttl + ) return True - return False def drop_tag_lease(self, tag_id: str) -> None: @@ -96,7 +151,7 @@ class ResourceManager: None """ - self.redis.delete(f"lease:opc_tags:{tag_id}") + self._execute_redis_op("delete", self.redis.delete, f"lease:opc_tags:{tag_id}") def get_all_ingestors(self) -> List[str]: """ @@ -107,7 +162,7 @@ class ResourceManager: list: A list of active ingestors. """ - return self.redis.keys("heartbeat:ingestor:*") + return self._execute_redis_op("keys", self.redis.keys, "heartbeat:ingestor:*") def get_all_slots(self) -> List[str]: """ @@ -118,7 +173,7 @@ class ResourceManager: int: The number of slots available. """ - return self.redis.keys("slot:opc_tags:*") + return self._execute_redis_op("keys", self.redis.keys, "slot:opc_tags:*") def get_all_leases(self) -> List[str]: """ @@ -129,4 +184,4 @@ class ResourceManager: list: A list of active leases. """ - return self.redis.keys("lease:opc_tags:*") + return self._execute_redis_op("keys", self.redis.keys, "lease:opc_tags:*") diff --git a/tests/unit/managers/__init__.py b/tests/unit/managers/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/unit/managers/test_resource_manager.py b/tests/unit/managers/test_resource_manager.py index aaa8fd2..19ca732 100644 --- a/tests/unit/managers/test_resource_manager.py +++ b/tests/unit/managers/test_resource_manager.py @@ -1,109 +1,190 @@ from unittest.mock import MagicMock, patch -from pytest import fixture +from pytest import fixture, raises from ingestor.managers.resource_manager import ResourceManager @fixture -@patch('ingestor.managers.resource_manager.Redis') +@patch("ingestor.managers.resource_manager.Redis") def resource_manager(redis): - return ResourceManager( - 'localhost', 6379, 10, 10, 'pod_id' - ) + return ResourceManager("localhost", 6379, 10, 10, "pod_id") def test_get_success(resource_manager): resource_manager.redis.get.return_value = '{"key": "value"}' - result = resource_manager.get('key') + result = resource_manager.get("key") assert result == {"key": "value"} - resource_manager.redis.get.assert_called_once_with('key') + resource_manager.redis.get.assert_called_once_with("key") def test_get_failure(resource_manager): resource_manager.redis.get.return_value = None - result = resource_manager.get('key') + result = resource_manager.get("key") assert result is None - resource_manager.redis.get.assert_called_once_with('key') + resource_manager.redis.get.assert_called_once_with("key") def test_get_tag_slot(resource_manager): resource_manager.get = MagicMock(return_value={"tag": "slot"}) - result = resource_manager.get_tag_slot('id') + result = resource_manager.get_tag_slot("id") assert result == {"tag": "slot"} - resource_manager.get.assert_called_once_with('slot:opc_tags:id') + resource_manager.get.assert_called_once_with("slot:opc_tags:id") 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:pod_id", 1, ex=10 ) def test_lease_tag(resource_manager): resource_manager.redis.set.return_value = True - output = resource_manager.lease_tag('tag_id') + 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", "pod_id", 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 = "pod_id" 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 - resource_manager.redis.get.assert_called_once_with( - 'lease:opc_tags:tag_id' - ) - resource_manager.redis.expire.assert_called_once_with( - 'lease:opc_tags:tag_id', 10 - ) + resource_manager.redis.get.assert_called_once_with("lease:opc_tags:tag_id") + resource_manager.redis.expire.assert_called_once_with("lease:opc_tags:tag_id", 10) def test_renew_tag_lease_failure(resource_manager): - resource_manager.redis.get.return_value = 'other_pod_id' + resource_manager.redis.get.return_value = "other_pod_id" resource_manager.redis.expire.return_value = False - result = resource_manager.renew_tag_lease('tag_id') + result = resource_manager.renew_tag_lease("tag_id") assert result is False - resource_manager.redis.get.assert_called_once_with( - 'lease:opc_tags:tag_id' - ) + resource_manager.redis.get.assert_called_once_with("lease:opc_tags:tag_id") resource_manager.redis.expire.assert_not_called() def test_drop_tag_lease(resource_manager): resource_manager.redis.delete.return_value = True - resource_manager.drop_tag_lease('tag_id') - resource_manager.redis.delete.assert_called_once_with( - 'lease:opc_tags:tag_id' - ) + resource_manager.drop_tag_lease("tag_id") + resource_manager.redis.delete.assert_called_once_with("lease:opc_tags:tag_id") def test_get_all_ingestors(resource_manager): - resource_manager.redis.keys.return_value = ['ingestor1', 'ingestor2'] + resource_manager.redis.keys.return_value = ["ingestor1", "ingestor2"] result = resource_manager.get_all_ingestors() - assert result == ['ingestor1', 'ingestor2'] - resource_manager.redis.keys.assert_called_once_with( - 'heartbeat:ingestor:*' - ) + assert result == ["ingestor1", "ingestor2"] + resource_manager.redis.keys.assert_called_once_with("heartbeat:ingestor:*") def test_get_all_slots(resource_manager): - resource_manager.redis.keys.return_value = ['slot1', 'slot2'] + resource_manager.redis.keys.return_value = ["slot1", "slot2"] result = resource_manager.get_all_slots() - assert result == ['slot1', 'slot2'] - resource_manager.redis.keys.assert_called_once_with( - 'slot:opc_tags:*' - ) + assert result == ["slot1", "slot2"] + resource_manager.redis.keys.assert_called_once_with("slot:opc_tags:*") def test_get_all_leases(resource_manager): - resource_manager.redis.keys.return_value = ['lease1', 'lease2'] + resource_manager.redis.keys.return_value = ["lease1", "lease2"] result = resource_manager.get_all_leases() - assert result == ['lease1', 'lease2'] - resource_manager.redis.keys.assert_called_once_with( - 'lease:opc_tags:*' - ) + assert result == ["lease1", "lease2"] + resource_manager.redis.keys.assert_called_once_with("lease:opc_tags:*") + + +def test_init_connection_failure(monkeypatch): + # Mock Redis to raise an exception during initialization + mock_redis = MagicMock() + mock_redis.side_effect = Exception("Connection failed") + + monkeypatch.setattr("ingestor.managers.resource_manager.Redis", mock_redis) + + # Test that the exception is raised and metrics are set properly + with patch("ingestor.metrics.REDIS_CONNECTION_STATUS") as mock_metrics: + mock_status = MagicMock() + mock_metrics.labels.return_value = mock_status + + with raises(Exception, match="Connection failed"): + ResourceManager("localhost", 6379, 10, 10, "pod_id") + + mock_metrics.labels.assert_called_once_with(pod_id="pod_id") + mock_status.set.assert_called_once_with(0) + + +def test_init_ping_failure(monkeypatch): + # Mock Redis ping to raise an exception + mock_redis_instance = MagicMock() + mock_redis_instance.ping.side_effect = Exception("Ping failed") + + mock_redis_class = MagicMock(return_value=mock_redis_instance) + monkeypatch.setattr("ingestor.managers.resource_manager.Redis", mock_redis_class) + + # Test that the exception is raised and metrics are set properly + with patch("ingestor.metrics.REDIS_CONNECTION_STATUS") as mock_metrics: + mock_status = MagicMock() + mock_metrics.labels.return_value = mock_status + + with raises(Exception, match="Ping failed"): + ResourceManager("localhost", 6379, 10, 10, "pod_id") + + mock_metrics.labels.assert_called_once_with(pod_id="pod_id") + mock_status.set.assert_called_once_with(0) + + +def test_execute_redis_op_success(resource_manager): + # Mock the Redis operation and time function + mock_func = MagicMock(return_value="test_result") + + with patch("ingestor.managers.resource_manager.time", side_effect=[100, 100.5]): + with patch("ingestor.metrics.REDIS_OPERATIONS_TOTAL") as mock_total: + with patch("ingestor.metrics.REDIS_OPERATIONS_DURATION") as mock_duration: + mock_total_labels = MagicMock() + mock_duration_labels = MagicMock() + mock_total.labels.return_value = mock_total_labels + mock_duration.labels.return_value = mock_duration_labels + + # Execute the operation + result = resource_manager._execute_redis_op( + "test_op", mock_func, "arg1", kwarg1="value1" + ) + + # Verify the result and metrics + assert result == "test_result" + mock_func.assert_called_once_with("arg1", kwarg1="value1") + + mock_total.labels.assert_called_once_with( + pod_id="pod_id", operation="test_op" + ) + mock_total_labels.inc.assert_called_once() + + mock_duration.labels.assert_called_once_with( + pod_id="pod_id", operation="test_op" + ) + mock_duration_labels.observe.assert_called_once_with( + 0.5 + ) + + +def test_execute_redis_op_exception(resource_manager): + # Mock the Redis operation to raise an exception + mock_func = MagicMock(side_effect=Exception("Operation failed")) + + with patch("ingestor.managers.resource_manager.time", return_value=100): + with patch("ingestor.metrics.REDIS_OPERATIONS_ERRORS") as mock_errors: + with patch("builtins.print") as mock_print: + mock_errors_labels = MagicMock() + mock_errors.labels.return_value = mock_errors_labels + + # Execute the operation and expect an exception + with raises(Exception, match="Operation failed"): + resource_manager._execute_redis_op("test_op", mock_func, "arg1") + + # Verify metrics and error handling + mock_errors.labels.assert_called_once_with( + pod_id="pod_id", operation="test_op" + ) + mock_errors_labels.inc.assert_called_once() + mock_print.assert_called_once_with( + "Error in Redis operation 'test_op': Operation failed" + )