diff --git a/.github/workflows/quality-gate.yml b/.github/workflows/quality-gate.yml index e13ecb2..cb52e6d 100644 --- a/.github/workflows/quality-gate.yml +++ b/.github/workflows/quality-gate.yml @@ -80,5 +80,5 @@ jobs: -Dsonar.host.url=$SONAR_HOST_URL \ -Dsonar.token=$SONAR_TOKEN \ -Dsonar.python.version=3.11 \ - -Dsonar.projectVersion=1.0.1 \ + -Dsonar.projectVersion=1.2.0 \ -Dsonar.coverage.exclusions=ingestor/app.py diff --git a/.vscode/settings.json b/.vscode/settings.json index 78e50ab..bdd01e8 100644 --- a/.vscode/settings.json +++ b/.vscode/settings.json @@ -1,11 +1,12 @@ { - "python.testing.pytestArgs": [ - "." - ], - "python.testing.unittestEnabled": false, - "python.testing.pytestEnabled": true, - "sonarlint.connectedMode.project": { - "connectionId": "sonardev-sientia-ai", - "projectKey": "Aignosi_sientia-dataops-opc-ingestor_1642c8bf-a148-4911-8362-0903d7fef99a" - } -} \ No newline at end of file + "python.testing.pytestArgs": ["."], + "python.testing.unittestEnabled": false, + "python.testing.pytestEnabled": true, + "sonarlint.connectedMode.project": { + "connectionId": "sonardev-sientia-ai", + "projectKey": "Aignosi_sientia-dataops-opc-ingestor_1642c8bf-a148-4911-8362-0903d7fef99a" + }, + "python.languageServer": "Pylance", + "python.analysis.typeCheckingMode": "off", + "editor.suggestSelection": "first" +} diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index b508b56..883af54 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -7,6 +7,8 @@ from sientia_do.notifications.handlers import NotificationHandler from ingestor.managers.ingestor_manager import IngestorManager +import ingestor.metrics as metrics + class Ingestor: def __init__(self): @@ -33,13 +35,13 @@ class Ingestor: kafka_servers = getenv("KAFKA_SERVERS", "localhost:9092") self.redis_host = getenv("REDIS_HOST", "localhost") - self.redis_port = int(getenv("REDIS_PORT", '6379')) + self.redis_port = int(getenv("REDIS_PORT", "6379")) self.redis_username = getenv("REDIS_USERNAME", None) self.redis_password = getenv("REDIS_PASSWORD", None) - self.lease_ttl = int(getenv("LEASE_TTL", '10')) - self.heartbeat_ttl = int(getenv("HEARTBEAT_TTL", '20')) + self.lease_ttl = int(getenv("LEASE_TTL", "10")) + self.heartbeat_ttl = int(getenv("HEARTBEAT_TTL", "20")) self.pod_id = getenv("HOSTNAME", "localhost") - self.poll_interval = int(getenv("POLL_INTERVAL", '5')) + self.poll_interval = int(getenv("POLL_INTERVAL", "5")) self.kafka_servers = kafka_servers.split(",") self.logger = None @@ -51,7 +53,7 @@ class Ingestor: pipeline_name="-", trigger_name="-", model_name="-", - model="-" + model="-", ) self.ingestor_manager = None @@ -77,8 +79,7 @@ class Ingestor: logger = getLogger(__name__) logger.setLevel(getenv("LOG_LEVEL", "INFO")) handler = StreamHandler() - formatter = Formatter( - '%(asctime)s - %(name)s - %(levelname)s - %(message)s') + formatter = Formatter("%(asctime)s - %(name)s - %(levelname)s - %(message)s") handler.setFormatter(formatter) logger.addHandler(handler) @@ -128,9 +129,17 @@ class Ingestor: """ self.ingestor_manager = IngestorManager( - self.kafka_servers, self.redis_host, self.redis_port, self.lease_ttl, - self.heartbeat_ttl, self.pod_id, self.poll_interval, self.logger, - self.notification_handler, self.redis_username, self.redis_password + self.kafka_servers, + self.redis_host, + self.redis_port, + self.lease_ttl, + self.heartbeat_ttl, + self.pod_id, + self.poll_interval, + self.logger, + self.notification_handler, + self.redis_username, + self.redis_password, ) # Declare ingestor ative @@ -142,6 +151,10 @@ class Ingestor: self.handle_acquired_tags(acquired) + metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set( + len(self.ingestor_manager.managed_tags) + ) # Set initial + def manage_no_slots(self, number_of_slots: int): """ Manages the scenario where there are no slots assigned to the ingestor. @@ -159,7 +172,9 @@ class Ingestor: # Get slot lease self.ingestor_manager.get_slot_leases(1) - def manage_leases(self, available_slots: int, lacking_ingestors: int, slot_diff: int): + def manage_leases( + self, available_slots: int, lacking_ingestors: int, slot_diff: int + ): """ Manages the allocation and deallocation of slot leases for ingestors based on the number of available slots, lacking ingestors, and slot differences. @@ -198,6 +213,11 @@ class Ingestor: self.ingestor_manager.unsubscribe_slot(lease) self.ingestor_manager.managed_tags.pop(lease) + # Update metric after removal + metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set( + len(self.ingestor_manager.managed_tags) + ) + def update_ingestor_manager(self, old_managed_tags: Dict[str, Any]): """ Updates the ingestor manager with the new managed tags. @@ -205,15 +225,18 @@ class Ingestor: old_managed_tags (Dict[str, Any]): The old managed tags. """ - self.logger.debug("Current managed tags: %s", - self.ingestor_manager.managed_tags) + self.logger.debug( + "Current managed tags: %s", self.ingestor_manager.managed_tags + ) self.ingestor_manager.update_opc_servers() new_managed_tags = deepcopy(self.ingestor_manager.managed_tags) - self.logger.debug("Comparing new managed tags %s" - "with old managed tags %s", - new_managed_tags, old_managed_tags) + self.logger.debug( + "Comparing new managed tags %s with old managed tags %s", + new_managed_tags, + old_managed_tags, + ) for slot, config in new_managed_tags.items(): if slot not in old_managed_tags: @@ -228,6 +251,11 @@ class Ingestor: if slot not in new_managed_tags: self.ingestor_manager.unsubscribe_slot(slot) + # Ensure the gauge is updated after any potential changes here + metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set( + len(self.ingestor_manager.managed_tags) + ) + def loop(self): """ Executes the main loop for managing ingestors and slots. @@ -259,6 +287,9 @@ class Ingestor: number_of_leases = self.ingestor_manager.get_number_of_leases() number_of_slots = self.ingestor_manager.get_number_of_slots() + # Update active ingestors gauge + metrics.ACTIVE_INGESTORS.set(number_of_ingestors) + # Handle no slots self.logger.debug("Managing no slots...") self.manage_no_slots(number_of_slots) @@ -270,6 +301,11 @@ class Ingestor: self.logger.debug("Managing leases...") self.manage_leases(available_slots, lacking_ingestors, slot_diff) + # Update managed slots gauge + metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set( + len(self.ingestor_manager.managed_tags) + ) + self.logger.debug( "Active ingestors: %s, " "Number of slots: %s, " @@ -280,7 +316,7 @@ class Ingestor: number_of_slots, number_of_leases, self.ingestor_manager.managed_tags, - self.ingestor_manager.opc_managers + self.ingestor_manager.opc_managers, ) if not self.ingestor_manager.managed_tags: # No slots acquired @@ -295,6 +331,4 @@ class Ingestor: self.ingestor_manager.check_opc_servers_integrity() self.logger.debug("Updating managed tags...") - self.update_ingestor_manager( - current_managed_tags - ) + self.update_ingestor_manager(current_managed_tags) diff --git a/ingestor/metrics.py b/ingestor/metrics.py index c605497..bb0976d 100644 --- a/ingestor/metrics.py +++ b/ingestor/metrics.py @@ -34,17 +34,17 @@ APP_UP = Gauge( ) # --- Ingestor Manager Metrics --- -INGESTORS_ATIVOS = Gauge( +ACTIVE_INGESTORS = Gauge( "ingestor_active_total", "Number of active ingestors reported by Redis", # No pod_id here, as it's a global view from Redis ) -SLOTS_TOTAIS = Gauge( +SLOTS_TOTAL = Gauge( "ingestor_slots_total", "Total number of slots configured in Redis", # No pod_id here, as it's a global view from Redis ) -LEASES_TOTAIS = Gauge( +LEASES_TOTAL = Gauge( "ingestor_leases_total", "Total number of leases (allocated slots) in Redis", # No pod_id here, as it's a global view from Redis diff --git a/tests/unit/test_metrics.py b/tests/unit/test_metrics.py index eb2670a..8d9ad7c 100644 --- a/tests/unit/test_metrics.py +++ b/tests/unit/test_metrics.py @@ -41,28 +41,28 @@ def test_app_up(): assert set(metrics.APP_UP._labelnames) == {"pod_id"} -def test_ingestors_ativos(): - """Verify the definition of INGESTORS_ATIVOS.""" - assert metrics.INGESTORS_ATIVOS is not None - assert isinstance(metrics.INGESTORS_ATIVOS, Gauge) - assert metrics.INGESTORS_ATIVOS._name == "ingestor_active_total" - assert set(metrics.INGESTORS_ATIVOS._labelnames) == set() +def test_active_ingestors(): + """Verify the definition of ACTIVE_INGESTORS.""" + assert metrics.ACTIVE_INGESTORS is not None + assert isinstance(metrics.ACTIVE_INGESTORS, Gauge) + assert metrics.ACTIVE_INGESTORS._name == "ingestor_active_total" + assert set(metrics.ACTIVE_INGESTORS._labelnames) == set() -def test_slots_totais(): - """Verify the definition of SLOTS_TOTAIS.""" - assert metrics.SLOTS_TOTAIS is not None - assert isinstance(metrics.SLOTS_TOTAIS, Gauge) - assert metrics.SLOTS_TOTAIS._name == "ingestor_slots_total" - assert set(metrics.SLOTS_TOTAIS._labelnames) == set() +def test_slots_total(): + """Verify the definition of SLOTS_TOTAL.""" + assert metrics.SLOTS_TOTAL is not None + assert isinstance(metrics.SLOTS_TOTAL, Gauge) + assert metrics.SLOTS_TOTAL._name == "ingestor_slots_total" + assert set(metrics.SLOTS_TOTAL._labelnames) == set() -def test_leases_totais(): - """Verify the definition of LEASES_TOTAIS.""" - assert metrics.LEASES_TOTAIS is not None - assert isinstance(metrics.LEASES_TOTAIS, Gauge) - assert metrics.LEASES_TOTAIS._name == "ingestor_leases_total" - assert set(metrics.LEASES_TOTAIS._labelnames) == set() +def test_leases_total(): + """Verify the definition of LEASES_TOTAL.""" + assert metrics.LEASES_TOTAL is not None + assert isinstance(metrics.LEASES_TOTAL, Gauge) + assert metrics.LEASES_TOTAL._name == "ingestor_leases_total" + assert set(metrics.LEASES_TOTAL._labelnames) == set() def test_slots_managed():