SIENTIAPDE-1083: add prometheus metrics to resource_manager.

This commit is contained in:
Bruno Domingues
2025-05-29 11:38:57 -03:00
parent f71b4c7937
commit 03ed5784b1
4 changed files with 205 additions and 69 deletions

View File

@@ -7,6 +7,6 @@
"projectKey": "Aignosi_sientia-dataops-opc-ingestor_1642c8bf-a148-4911-8362-0903d7fef99a" "projectKey": "Aignosi_sientia-dataops-opc-ingestor_1642c8bf-a148-4911-8362-0903d7fef99a"
}, },
"python.languageServer": "Pylance", "python.languageServer": "Pylance",
"python.analysis.typeCheckingMode": "basic", "python.analysis.typeCheckingMode": "standard",
"editor.suggestSelection": "first" "editor.suggestSelection": "first"
} }

View File

@@ -1,17 +1,59 @@
import json import json
from typing import List from typing import List
from redis import Redis from redis import Redis
from time import time
import ingestor.metrics as metrics
class ResourceManager: class ResourceManager:
def __init__(self, host: str, port: int, def __init__(
lease_ttl: int, heartbeat_ttl: int, pod_id: str, self,
username: str = None, password: str = None) -> None: host: str,
self.redis = Redis(host=host, port=port, decode_responses=True, port: int,
username=username, password=password) 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.lease_ttl = lease_ttl
self.heartbeat_ttl = heartbeat_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: def get(self, key: str) -> dict:
""" """
@@ -19,11 +61,11 @@ class ResourceManager:
Args: Args:
key (str): The key to look up in Redis. key (str): The key to look up in Redis.
Returns: 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. 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 return json.loads(history) if history else None
def get_tag_slot(self, id: str) -> dict: def get_tag_slot(self, id: str) -> dict:
@@ -48,14 +90,19 @@ class ResourceManager:
None None
""" """
self.redis.set( self._execute_redis_op(
f"heartbeat:ingestor:{self.pod_id}", 1, ex=self.heartbeat_ttl) "set",
self.redis.set,
f"heartbeat:ingestor:{self.pod_id}",
1,
ex=self.heartbeat_ttl,
)
def lease_tag(self, tag_id: str) -> bool: 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). 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 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 is only set if it does not already exist. The key is set with an expiration time
defined by `lease_ttl`. defined by `lease_ttl`.
Args: Args:
tag_id (str): The unique identifier of the tag to be leased. 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. bool: True if the lease was successfully acquired, False otherwise.
""" """
return self.redis.set( return self._execute_redis_op(
f"lease:opc_tags:{tag_id}", self.pod_id, nx=True, ex=self.lease_ttl) "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: 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. bool: True if the lease was successfully renewed, False otherwise.
""" """
current = self.redis.get( current = self._execute_redis_op(
f"lease:opc_tags:{tag_id}") "get", self.redis.get, f"lease:opc_tags:{tag_id}"
)
if current == self.pod_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 True
return False return False
def drop_tag_lease(self, tag_id: str) -> None: def drop_tag_lease(self, tag_id: str) -> None:
@@ -96,7 +151,7 @@ class ResourceManager:
None 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]: def get_all_ingestors(self) -> List[str]:
""" """
@@ -107,7 +162,7 @@ class ResourceManager:
list: A list of active ingestors. 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]: def get_all_slots(self) -> List[str]:
""" """
@@ -118,7 +173,7 @@ class ResourceManager:
int: The number of slots available. 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]: def get_all_leases(self) -> List[str]:
""" """
@@ -129,4 +184,4 @@ class ResourceManager:
list: A list of active leases. 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:*")

View File

View File

@@ -1,109 +1,190 @@
from unittest.mock import MagicMock, patch from unittest.mock import MagicMock, patch
from pytest import fixture from pytest import fixture, raises
from ingestor.managers.resource_manager import ResourceManager from ingestor.managers.resource_manager import ResourceManager
@fixture @fixture
@patch('ingestor.managers.resource_manager.Redis') @patch("ingestor.managers.resource_manager.Redis")
def resource_manager(redis): def resource_manager(redis):
return ResourceManager( return ResourceManager("localhost", 6379, 10, 10, "pod_id")
'localhost', 6379, 10, 10, 'pod_id'
)
def test_get_success(resource_manager): def test_get_success(resource_manager):
resource_manager.redis.get.return_value = '{"key": "value"}' resource_manager.redis.get.return_value = '{"key": "value"}'
result = resource_manager.get('key') result = resource_manager.get("key")
assert result == {"key": "value"} 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): def test_get_failure(resource_manager):
resource_manager.redis.get.return_value = None resource_manager.redis.get.return_value = None
result = resource_manager.get('key') result = resource_manager.get("key")
assert result is None 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): def test_get_tag_slot(resource_manager):
resource_manager.get = MagicMock(return_value={"tag": "slot"}) 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"} 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): 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:pod_id", 1, ex=10
) )
def test_lease_tag(resource_manager): def test_lease_tag(resource_manager):
resource_manager.redis.set.return_value = True 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 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", "pod_id", 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 = "pod_id"
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
resource_manager.redis.get.assert_called_once_with( resource_manager.redis.get.assert_called_once_with("lease:opc_tags:tag_id")
'lease:opc_tags:tag_id' resource_manager.redis.expire.assert_called_once_with("lease:opc_tags:tag_id", 10)
)
resource_manager.redis.expire.assert_called_once_with(
'lease:opc_tags:tag_id', 10
)
def test_renew_tag_lease_failure(resource_manager): 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 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 assert result is False
resource_manager.redis.get.assert_called_once_with( resource_manager.redis.get.assert_called_once_with("lease:opc_tags:tag_id")
'lease:opc_tags:tag_id'
)
resource_manager.redis.expire.assert_not_called() resource_manager.redis.expire.assert_not_called()
def test_drop_tag_lease(resource_manager): def test_drop_tag_lease(resource_manager):
resource_manager.redis.delete.return_value = True resource_manager.redis.delete.return_value = True
resource_manager.drop_tag_lease('tag_id') resource_manager.drop_tag_lease("tag_id")
resource_manager.redis.delete.assert_called_once_with( resource_manager.redis.delete.assert_called_once_with("lease:opc_tags:tag_id")
'lease:opc_tags:tag_id'
)
def test_get_all_ingestors(resource_manager): 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() result = resource_manager.get_all_ingestors()
assert result == ['ingestor1', 'ingestor2'] assert result == ["ingestor1", "ingestor2"]
resource_manager.redis.keys.assert_called_once_with( resource_manager.redis.keys.assert_called_once_with("heartbeat:ingestor:*")
'heartbeat:ingestor:*'
)
def test_get_all_slots(resource_manager): 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() result = resource_manager.get_all_slots()
assert result == ['slot1', 'slot2'] assert result == ["slot1", "slot2"]
resource_manager.redis.keys.assert_called_once_with( resource_manager.redis.keys.assert_called_once_with("slot:opc_tags:*")
'slot:opc_tags:*'
)
def test_get_all_leases(resource_manager): 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() result = resource_manager.get_all_leases()
assert result == ['lease1', 'lease2'] assert result == ["lease1", "lease2"]
resource_manager.redis.keys.assert_called_once_with( resource_manager.redis.keys.assert_called_once_with("lease:opc_tags:*")
'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"
)