Merge pull request #6 from Aignosi/feature/SIENTIAPDE-1083

SIENTIAPDE-1083: Add prometheus metrics
This commit is contained in:
Matheus Demoner
2025-05-30 12:08:15 -03:00
committed by GitHub
17 changed files with 1415 additions and 268 deletions

View File

@@ -80,5 +80,5 @@ jobs:
-Dsonar.host.url=$SONAR_HOST_URL \ -Dsonar.host.url=$SONAR_HOST_URL \
-Dsonar.token=$SONAR_TOKEN \ -Dsonar.token=$SONAR_TOKEN \
-Dsonar.python.version=3.11 \ -Dsonar.python.version=3.11 \
-Dsonar.projectVersion=1.0.1 \ -Dsonar.projectVersion=1.2.0 \
-Dsonar.coverage.exclusions=ingestor/app.py -Dsonar.coverage.exclusions=ingestor/app.py

15
.vscode/settings.json vendored
View File

@@ -1,7 +1,12 @@
{ {
"python.testing.pytestArgs": [ "python.testing.pytestArgs": ["."],
"." "python.testing.unittestEnabled": false,
], "python.testing.pytestEnabled": true,
"python.testing.unittestEnabled": false, "sonarlint.connectedMode.project": {
"python.testing.pytestEnabled": true "connectionId": "sonardev-sientia-ai",
"projectKey": "Aignosi_sientia-dataops-opc-ingestor_1642c8bf-a148-4911-8362-0903d7fef99a"
},
"python.languageServer": "Pylance",
"python.analysis.typeCheckingMode": "standard",
"editor.suggestSelection": "first"
} }

View File

@@ -1,32 +1,53 @@
from threading import Event
import signal
import os import os
import signal
import traceback import traceback
from threading import Event
from time import sleep, time
from prometheus_client import start_http_server
import ingestor.metrics as metrics
from ingestor.ingestor import Ingestor from ingestor.ingestor import Ingestor
exit_signal = Event() exit_signal = Event()
POD_ID = os.getenv("HOSTNAME", "localhost")
def main(): def main():
start_prometheus_server()
ingestor = Ingestor() ingestor = Ingestor()
ingestor.prepare_ingestor() ingestor.prepare_ingestor()
ingestor.logger.info("Ingestor prepared. Starting main loop.") ingestor.logger.info("Ingestor prepared. Starting main loop.")
while not exit_signal.is_set(): while not exit_signal.is_set():
start_time = time() # Start loop timer
try: try:
ingestor.loop() ingestor.loop()
metrics.APP_LOOP_COUNT.labels(pod_id=POD_ID).inc() # Increment loop counter
exit_signal.wait(ingestor.poll_interval) exit_signal.wait(ingestor.poll_interval)
except KeyboardInterrupt: # Handle Ctrl+C gracefully
print("KeyboardInterrupt received. Setting exit_signal flag.")
exit_signal.set()
except Exception: except Exception:
print("Exception in main loop. Setting exit_signal flag.") print("Exception in main loop. Setting exit_signal flag.")
traceback.print_exc() traceback.print_exc()
metrics.APP_ERRORS_TOTAL.labels(pod_id=POD_ID).inc() # Increment errors
exit_signal.set() exit_signal.set()
finally:
# Record loop duration
duration = time() - start_time
metrics.APP_LOOP_DURATION.labels(pod_id=POD_ID).observe(duration)
ingestor.shutdown() ingestor.shutdown()
metrics.APP_UP.labels(pod_id=POD_ID).set(0) # Mark app as DOWN
ingestor.logger.info("Main loop exit_signaled.") ingestor.logger.info("Main loop exit_signaled.")
# Give Prometheus a chance to scrape one last time before exiting (optional)
sleep(5)
os._exit(0) os._exit(0)
@@ -35,6 +56,16 @@ def signal_handler(_signum, _frame):
exit_signal.set() exit_signal.set()
def start_prometheus_server():
try:
start_http_server(8000)
print("Prometheus server started on port 8000.")
metrics.APP_UP.labels(pod_id=POD_ID).set(1) # Mark app as UP
except Exception as e:
print(f"Failed to start Prometheus server: {e}")
os._exit(1)
if __name__ == "__main__": if __name__ == "__main__":
signal.signal(signal.SIGINT, signal_handler) signal.signal(signal.SIGINT, signal_handler)
signal.signal(signal.SIGTERM, signal_handler) signal.signal(signal.SIGTERM, signal_handler)

View File

@@ -7,6 +7,8 @@ from sientia_do.notifications.handlers import NotificationHandler
from ingestor.managers.ingestor_manager import IngestorManager from ingestor.managers.ingestor_manager import IngestorManager
import ingestor.metrics as metrics
class Ingestor: class Ingestor:
def __init__(self): def __init__(self):
@@ -33,13 +35,13 @@ class Ingestor:
kafka_servers = getenv("KAFKA_SERVERS", "localhost:9092") kafka_servers = getenv("KAFKA_SERVERS", "localhost:9092")
self.redis_host = getenv("REDIS_HOST", "localhost") 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_username = getenv("REDIS_USERNAME", None)
self.redis_password = getenv("REDIS_PASSWORD", None) self.redis_password = getenv("REDIS_PASSWORD", None)
self.lease_ttl = int(getenv("LEASE_TTL", '10')) self.lease_ttl = int(getenv("LEASE_TTL", "10"))
self.heartbeat_ttl = int(getenv("HEARTBEAT_TTL", '20')) self.heartbeat_ttl = int(getenv("HEARTBEAT_TTL", "20"))
self.pod_id = getenv("HOSTNAME", "localhost") 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.kafka_servers = kafka_servers.split(",")
self.logger = None self.logger = None
@@ -51,7 +53,7 @@ class Ingestor:
pipeline_name="-", pipeline_name="-",
trigger_name="-", trigger_name="-",
model_name="-", model_name="-",
model="-" model="-",
) )
self.ingestor_manager = None self.ingestor_manager = None
@@ -77,8 +79,7 @@ class Ingestor:
logger = getLogger(__name__) logger = getLogger(__name__)
logger.setLevel(getenv("LOG_LEVEL", "INFO")) logger.setLevel(getenv("LOG_LEVEL", "INFO"))
handler = StreamHandler() handler = StreamHandler()
formatter = Formatter( formatter = Formatter("%(asctime)s - %(name)s - %(levelname)s - %(message)s")
'%(asctime)s - %(name)s - %(levelname)s - %(message)s')
handler.setFormatter(formatter) handler.setFormatter(formatter)
logger.addHandler(handler) logger.addHandler(handler)
@@ -128,9 +129,17 @@ class Ingestor:
""" """
self.ingestor_manager = IngestorManager( self.ingestor_manager = IngestorManager(
self.kafka_servers, self.redis_host, self.redis_port, self.lease_ttl, self.kafka_servers,
self.heartbeat_ttl, self.pod_id, self.poll_interval, self.logger, self.redis_host,
self.notification_handler, self.redis_username, self.redis_password 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 # Declare ingestor ative
@@ -142,6 +151,10 @@ class Ingestor:
self.handle_acquired_tags(acquired) 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): def manage_no_slots(self, number_of_slots: int):
""" """
Manages the scenario where there are no slots assigned to the ingestor. Manages the scenario where there are no slots assigned to the ingestor.
@@ -159,7 +172,9 @@ class Ingestor:
# Get slot lease # Get slot lease
self.ingestor_manager.get_slot_leases(1) 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 Manages the allocation and deallocation of slot leases for ingestors based on
the number of available slots, lacking ingestors, and slot differences. 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.unsubscribe_slot(lease)
self.ingestor_manager.managed_tags.pop(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]): def update_ingestor_manager(self, old_managed_tags: Dict[str, Any]):
""" """
Updates the ingestor manager with the new managed tags. 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. old_managed_tags (Dict[str, Any]): The old managed tags.
""" """
self.logger.debug("Current managed tags: %s", self.logger.debug(
self.ingestor_manager.managed_tags) "Current managed tags: %s", self.ingestor_manager.managed_tags
)
self.ingestor_manager.update_opc_servers() self.ingestor_manager.update_opc_servers()
new_managed_tags = deepcopy(self.ingestor_manager.managed_tags) new_managed_tags = deepcopy(self.ingestor_manager.managed_tags)
self.logger.debug("Comparing new managed tags %s" self.logger.debug(
"with old managed tags %s", "Comparing new managed tags %s with old managed tags %s",
new_managed_tags, old_managed_tags) new_managed_tags,
old_managed_tags,
)
for slot, config in new_managed_tags.items(): for slot, config in new_managed_tags.items():
if slot not in old_managed_tags: if slot not in old_managed_tags:
@@ -228,6 +251,11 @@ class Ingestor:
if slot not in new_managed_tags: if slot not in new_managed_tags:
self.ingestor_manager.unsubscribe_slot(slot) 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): def loop(self):
""" """
Executes the main loop for managing ingestors and slots. 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_leases = self.ingestor_manager.get_number_of_leases()
number_of_slots = self.ingestor_manager.get_number_of_slots() number_of_slots = self.ingestor_manager.get_number_of_slots()
# Update active ingestors gauge
metrics.ACTIVE_INGESTORS.set(number_of_ingestors)
# Handle no slots # Handle no slots
self.logger.debug("Managing no slots...") self.logger.debug("Managing no slots...")
self.manage_no_slots(number_of_slots) self.manage_no_slots(number_of_slots)
@@ -270,6 +301,11 @@ class Ingestor:
self.logger.debug("Managing leases...") self.logger.debug("Managing leases...")
self.manage_leases(available_slots, lacking_ingestors, slot_diff) 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( self.logger.debug(
"Active ingestors: %s, " "Active ingestors: %s, "
"Number of slots: %s, " "Number of slots: %s, "
@@ -280,7 +316,7 @@ class Ingestor:
number_of_slots, number_of_slots,
number_of_leases, number_of_leases,
self.ingestor_manager.managed_tags, self.ingestor_manager.managed_tags,
self.ingestor_manager.opc_managers self.ingestor_manager.opc_managers,
) )
if not self.ingestor_manager.managed_tags: if not self.ingestor_manager.managed_tags:
# No slots acquired # No slots acquired
@@ -295,6 +331,4 @@ class Ingestor:
self.ingestor_manager.check_opc_servers_integrity() self.ingestor_manager.check_opc_servers_integrity()
self.logger.debug("Updating managed tags...") self.logger.debug("Updating managed tags...")
self.update_ingestor_manager( self.update_ingestor_manager(current_managed_tags)
current_managed_tags
)

View File

@@ -6,11 +6,17 @@ from kafka.errors import NoBrokersAvailable
from sientia_do.notifications.handlers import NotificationHandler from sientia_do.notifications.handlers import NotificationHandler
from sientia_do.notifications.models import NotificationLevel from sientia_do.notifications.models import NotificationLevel
import traceback import traceback
import ingestor.metrics as metrics
import os
class DataManager(): class DataManager:
def __init__(self, kafka_servers: str, logger: Logger, def __init__(
notification_handler: NotificationHandler) -> None: self,
kafka_servers: str,
logger: Logger,
notification_handler: NotificationHandler,
) -> None:
""" """
Initializes the DataManager instance with a Kafka producer. Initializes the DataManager instance with a Kafka producer.
This constructor attempts to establish a connection to the specified Kafka servers This constructor attempts to establish a connection to the specified Kafka servers
@@ -23,31 +29,39 @@ class DataManager():
NoBrokersAvailable: If the connection to Kafka servers fails after 3 attempts. NoBrokersAvailable: If the connection to Kafka servers fails after 3 attempts.
""" """
self.pod_id = os.getenv("HOSTNAME", "localhost")
self.kafka_producer = None self.kafka_producer = None
for i in range(0, 3): for i in range(0, 3):
logger.info( logger.info(
f"Trying ({i}) to initializing DataManager with Kafka servers: {kafka_servers}") f"Trying ({i}) to initializing DataManager with Kafka servers: {kafka_servers}"
)
try: try:
self.kafka_producer = KafkaProducer( self.kafka_producer = KafkaProducer(
bootstrap_servers=kafka_servers, bootstrap_servers=kafka_servers,
value_serializer=lambda v: json.dumps(v).encode( value_serializer=lambda v: json.dumps(v).encode(
'utf-8'), # Serialize JSON messages "utf-8"
key_serializer=lambda k: str( ), # Serialize JSON messages
k).encode('utf-8') if k else None, key_serializer=lambda k: str(k).encode("utf-8") if k else None,
) )
# Kafka connected
metrics.KAFKA_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(1)
break break
except NoBrokersAvailable: except NoBrokersAvailable:
logger.error( logger.error(
f"Kafka servers {kafka_servers} are not available. Retrying...") f"Kafka servers {kafka_servers} are not available. Retrying..."
)
sleep(5) sleep(5)
else: else:
# Kafka not connected
metrics.KAFKA_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(0)
logger.error( logger.error(
f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts.") f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts."
)
raise NoBrokersAvailable( raise NoBrokersAvailable(
f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts.") f"Failed to connect to Kafka servers {kafka_servers} after 3 attempts."
)
logger.info( logger.info(f"DataManager initialized with Kafka servers: {kafka_servers}")
f"DataManager initialized with Kafka servers: {kafka_servers}")
self.logger = logger self.logger = logger
self.notification_handler = notification_handler self.notification_handler = notification_handler
@@ -57,12 +71,12 @@ class DataManager():
try: try:
self.kafka_producer.flush(timeout=10) self.kafka_producer.flush(timeout=10)
self.kafka_producer.close() self.kafka_producer.close()
# Mark as disconnected
metrics.KAFKA_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(0)
except Exception as e: except Exception as e:
self.logger.error( self.logger.error(f"Error closing Kafka producer: {e}")
f"Error closing Kafka producer: {e}")
else: else:
self.logger.warning( self.logger.warning("Kafka producer is already closed or not initialized.")
"Kafka producer is already closed or not initialized.")
def __del__(self): def __del__(self):
self.shutdown() self.shutdown()
@@ -70,7 +84,8 @@ class DataManager():
def delivery_report(self, msg: str): def delivery_report(self, msg: str):
"""Callback for delivery reports from Kafka.""" """Callback for delivery reports from Kafka."""
self.logger.debug( self.logger.debug(
f"Record successfully produced to {msg.topic} [{msg.partition}] at offset {msg.offset}") f"Record successfully produced to {msg.topic} [{msg.partition}] at offset {msg.offset}"
)
def delivery_error(self, err: str): def delivery_error(self, err: str):
"""Callback for delivery reports from Kafka.""" """Callback for delivery reports from Kafka."""
@@ -93,22 +108,22 @@ class DataManager():
try: try:
self.logger.debug( self.logger.debug(f"Publishing message to topic {topic}: {data}")
f"Publishing message to topic {topic}: {data}") self.kafka_producer.send(topic=topic, value=data).add_callback(
self.kafka_producer.send( self.delivery_report
topic=topic, value=data).add_callback( ).add_errback(self.delivery_error)
self.delivery_report).add_errback(
self.delivery_error)
self.kafka_producer.flush(timeout=10) self.kafka_producer.flush(timeout=10)
metrics.KAFKA_MESSAGES_SENT.labels(pod_id=self.pod_id, topic=topic).inc()
except Exception as e: except Exception as e:
metrics.KAFKA_MESSAGES_ERRORS.labels(pod_id=self.pod_id, topic=topic).inc()
trace = traceback.format_exc() trace = traceback.format_exc()
self.notification_handler.build_and_send_notification( self.notification_handler.build_and_send_notification(
notification_id=f"KAFKA_PRODUCER_ERROR_{topic}", notification_id=f"KAFKA_PRODUCER_ERROR_{topic}",
message=f"Error publishing message to topic {topic}: {e}", message=f"Error publishing message to topic {topic}: {e}",
block="kafka_producer", block="kafka_producer",
level=NotificationLevel.ERROR, level=NotificationLevel.ERROR,
attachment_content=trace attachment_content=trace,
) )
self.logger.error(trace) self.logger.error(trace)

View File

@@ -7,6 +7,7 @@ from sientia_do.notifications.models import NotificationLevel
from ingestor.managers.data_manager import DataManager from ingestor.managers.data_manager import DataManager
from ingestor.managers.opc_manager import OpcManager from ingestor.managers.opc_manager import OpcManager
from ingestor.managers.resource_manager import ResourceManager from ingestor.managers.resource_manager import ResourceManager
import ingestor.metrics as metrics
class IngestorManager(): class IngestorManager():
@@ -29,6 +30,7 @@ class IngestorManager():
self.opc_servers = {} self.opc_servers = {}
self.notification_handler = notification_handler self.notification_handler = notification_handler
self.pod_id = pod_id
def initialize_opc_from_config(self, server_config: dict, def initialize_opc_from_config(self, server_config: dict,
data_manager: DataManager, logger: Logger) -> OpcManager | None: data_manager: DataManager, logger: Logger) -> OpcManager | None:
@@ -56,7 +58,7 @@ class IngestorManager():
manager = OpcManager( manager = OpcManager(
server_config['name'], server_config['url'], server_config['name'], server_config['url'],
data_manager, logger, server_config['server_uri'], data_manager, logger, server_config['server_uri'],
self.notification_handler, server_config.get('cert_path'), self.notification_handler, self.pod_id, server_config.get('cert_path'),
server_config.get('private_key_path'), server_config.get('private_key_path'),
server_config.get('server_cert_path') server_config.get('server_cert_path')
) )
@@ -161,6 +163,8 @@ class IngestorManager():
self.opc_managers[server].disconnect() self.opc_managers[server].disconnect()
self.opc_managers.pop(server, None) self.opc_managers.pop(server, None)
metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers))
def check_opc_servers_integrity(self): def check_opc_servers_integrity(self):
""" """
Checks the integrity of the OPC servers and updates the OPC servers if necessary. Checks the integrity of the OPC servers and updates the OPC servers if necessary.
@@ -178,6 +182,8 @@ class IngestorManager():
for slot, _config in self.managed_tags.items(): for slot, _config in self.managed_tags.items():
self.managed_tags[slot].pop(server, None) self.managed_tags[slot].pop(server, None)
metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers))
def declare_active(self): def declare_active(self):
""" """
Declares the ingestor as active by sending a heartbeat signal to the resource manager. Declares the ingestor as active by sending a heartbeat signal to the resource manager.
@@ -210,6 +216,7 @@ class IngestorManager():
leases = self.resource_manager.get_all_leases() leases = self.resource_manager.get_all_leases()
self.number_of_slots = len(leases) if leases else 0 self.number_of_slots = len(leases) if leases else 0
metrics.LEASES_TOTAL.set(self.number_of_slots)
return self.number_of_slots return self.number_of_slots
def get_number_of_slots(self) -> int: def get_number_of_slots(self) -> int:
@@ -223,6 +230,7 @@ class IngestorManager():
slots = self.resource_manager.get_all_slots() slots = self.resource_manager.get_all_slots()
self.number_of_slots = len(slots) if slots else 0 self.number_of_slots = len(slots) if slots else 0
metrics.SLOTS_TOTAL.set(self.number_of_slots)
return self.number_of_slots return self.number_of_slots
def get_slot_leases(self, max_slots: int = 1) -> Dict: def get_slot_leases(self, max_slots: int = 1) -> Dict:
@@ -254,9 +262,11 @@ class IngestorManager():
if slots is None: if slots is None:
continue continue
acquired[str(i)] = slots acquired[str(i)] = slots
metrics.SLOTS_ACQUIRED.labels(pod_id=self.pod_id).inc()
if len(acquired) >= max_slots: if len(acquired) >= max_slots:
self.managed_tags.update(acquired) self.managed_tags.update(acquired)
metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags))
return acquired return acquired
self.logger.warning( self.logger.warning(
@@ -264,6 +274,7 @@ class IngestorManager():
f"Only {acquired} slots were leased." f"Only {acquired} slots were leased."
) )
self.managed_tags.update(acquired) self.managed_tags.update(acquired)
metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags))
return acquired return acquired
def unsubscribe_slot(self, slot: str): def unsubscribe_slot(self, slot: str):
@@ -331,6 +342,7 @@ class IngestorManager():
for lease_id in ids: for lease_id in ids:
self.resource_manager.drop_tag_lease(lease_id) self.resource_manager.drop_tag_lease(lease_id)
metrics.SLOTS_RELEASED.labels(pod_id=self.pod_id).inc()
def manage_server(self, slot: str, server: str, server_config: dict, tags: dict) -> int: def manage_server(self, slot: str, server: str, server_config: dict, tags: dict) -> int:
""" """
@@ -389,6 +401,7 @@ class IngestorManager():
tags_to_sub tags_to_sub
) )
except Exception as e: except Exception as e:
metrics.OPC_SUBSCRIPTION_ERRORS.labels(pod_id=self.pod_id, server=server, slot=slot).inc()
trace = traceback.format_exc() trace = traceback.format_exc()
self.notification_handler.build_and_send_notification( self.notification_handler.build_and_send_notification(
notification_id=f'OPC_SUBSCRIPTION_ERROR_{slot}:{server}', notification_id=f'OPC_SUBSCRIPTION_ERROR_{slot}:{server}',

View File

@@ -7,11 +7,12 @@ from asyncua.sync import Client
from sientia_do.notifications.models import NotificationLevel from sientia_do.notifications.models import NotificationLevel
from sientia_do.notifications.handlers import NotificationHandler from sientia_do.notifications.handlers import NotificationHandler
from ingestor.managers.data_manager import DataManager from ingestor.managers.data_manager import DataManager
import ingestor.metrics as metrics
class OpcManager(): class OpcManager():
def __init__(self, name: str, url: str, data_manager: DataManager, def __init__(self, name: str, url: str, data_manager: DataManager,
logger: Logger, server_uri: str, notification_handler: NotificationHandler, logger: Logger, server_uri: str, notification_handler: NotificationHandler, pod_id: str,
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,8 +27,10 @@ class OpcManager():
self.nodes = {} self.nodes = {}
self.subscriptions = {} self.subscriptions = {}
self.data_manager = data_manager self.data_manager = data_manager
self.notification_handler = notification_handler self.notification_handler = notification_handler
self.pod_id = pod_id
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)
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" \
@@ -82,11 +85,20 @@ class OpcManager():
Exception: If the connection to the OPC server fails. Exception: If the connection to the OPC server fails.
""" """
self.client = Client(self.url) metrics.OPC_CONNECTIONS_TOTAL.labels(pod_id=self.pod_id, server_name=self.name).inc()
if self.cert_path: try:
self.set_security() self.client = Client(self.url)
self.logger.info('Starting connection...') if self.cert_path:
self.client.connect() self.set_security()
self.logger.info(f'Starting connection to {self.name}...')
self.client.connect()
metrics.OPC_CONNECTION_STATUS.labels(pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(1)
self.logger.info(f'Connection to {self.name} successful.')
except Exception as e:
metrics.OPC_CONNECTION_STATUS.labels(pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(0)
metrics.OPC_CONNECTIONS_FAILED.labels(pod_id=self.pod_id, server_name=self.name).inc()
self.logger.error(f"Failed to connect to {self.name}: {e}")
raise
def create_subscription(self, name: str, period: int = 500): def create_subscription(self, name: str, period: int = 500):
""" """
@@ -105,11 +117,14 @@ class OpcManager():
if not self.client: if not self.client:
raise ValueError("Client not connected. Call connect first.") raise ValueError("Client not connected. Call connect first.")
try:
p = period if period != None else 500 p = period if period is not None else 500
self.subscriptions[name] = self.client.create_subscription( self.subscriptions[name] = self.client.create_subscription(p, self)
p, self) self.logger.info(f'Subscription {name} created on {self.name}.')
self.logger.info('Subscription created.') metrics.OPC_SUBSCRIPTIONS_CREATED.labels(pod_id=self.pod_id, server_name=self.name, slot_name=name).inc()
except Exception as e:
self.logger.error(f"Failed to create subscription {name} on {self.name}: {e}")
raise
def subscribe(self, subscription: str, nodes: dict, collect_period: int): def subscribe(self, subscription: str, nodes: dict, collect_period: int):
""" """
@@ -129,14 +144,13 @@ class OpcManager():
""" """
if not self.subscriptions.get(subscription): if not self.subscriptions.get(subscription):
raise ValueError( raise ValueError("Subscription not created. Call create_subscription first.")
"Subscription not created. Call create_subscription first.")
self.logger.info(f"Subscribing to {subscription}...") self.logger.info(f"Subscribing to {subscription} on {self.name}...")
self.logger.info(f"Subscribing to nodes: {nodes}") self.logger.info(f"Subscribing to nodes: {nodes}")
self.addr_nodes = [self.client.get_node( self.addr_nodes = [self.client.get_node(n) for n in nodes if n not in self.nodes]
n) for n in nodes if n not in self.nodes]
self.nodes.update(nodes) self.nodes.update(nodes)
metrics.OPC_TAGS_SUBSCRIBED.labels(pod_id=self.pod_id, server_name=self.name).set(len(self.nodes))
self.collect_period = collect_period self.collect_period = collect_period
for node, config in self.nodes.items(): for node, config in self.nodes.items():
@@ -202,6 +216,8 @@ class OpcManager():
finally: finally:
del self.client del self.client
self.client = None self.client = None
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)
self.logger.warning("Disconnected from OPC UA server.") self.logger.warning("Disconnected from OPC UA server.")
def datachange_notification(self, node, _val, data): def datachange_notification(self, node, _val, data):
@@ -236,6 +252,7 @@ class OpcManager():
self.nodes[tag]['cycle_rule']['cycle_count'] = 0 self.nodes[tag]['cycle_rule']['cycle_count'] = 0
self.non_receive_count = 0 self.non_receive_count = 0
metrics.OPC_CYCLES_WITHOUT_DATA.labels(pod_id=self.pod_id, server_name=self.name).set(0)
data = { data = {
'tag': tag, 'tag': tag,
@@ -278,6 +295,7 @@ class OpcManager():
""" """
self.non_receive_count += 1 self.non_receive_count += 1
metrics.OPC_CYCLES_WITHOUT_DATA.labels(pod_id=self.pod_id, server_name=self.name).set(self.non_receive_count)
if self.non_receive_count >= 5: if self.non_receive_count >= 5:
self.notification_handler.build_and_send_notification( self.notification_handler.build_and_send_notification(
notification_id=f'OPC_LISTENNING_STOPPED__{self.name}', notification_id=f'OPC_LISTENNING_STOPPED__{self.name}',
@@ -287,6 +305,7 @@ class OpcManager():
level=NotificationLevel.ERROR level=NotificationLevel.ERROR
) )
if self.non_receive_count >= 15: if self.non_receive_count >= 15:
metrics.OPC_RECONNECTIONS_TOTAL.labels(pod_id=self.pod_id, server_name=self.name).inc()
self.notification_handler.build_and_send_notification( self.notification_handler.build_and_send_notification(
notification_id=f'OPC_CONNECTION_RETRY__{self.name}', notification_id=f'OPC_CONNECTION_RETRY__{self.name}',
message=f'Retrying to connect to server {self.name}', message=f'Retrying to connect to server {self.name}',

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:
""" """
@@ -23,7 +65,7 @@ class ResourceManager:
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,8 +90,13 @@ 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:
""" """
@@ -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:*")

148
ingestor/metrics.py Normal file
View File

@@ -0,0 +1,148 @@
from prometheus_client import Counter, Gauge, Histogram
POD_ID_LABEL = ["pod_id"]
SERVER_LABELS = ["pod_id", "server_name", "server_url"]
KAFKA_LABELS = ["pod_id", "topic"]
REDIS_LABELS = ["pod_id", "operation"]
NOTIFICATION_LABELS = ["pod_id", "level", "block"]
# --- General Application Metrics ---
APP_LOOP_COUNT = Counter(
"app_main_loop_total",
"Total number of times the application main loop has run",
POD_ID_LABEL,
)
APP_LOOP_DURATION = Histogram(
"app_main_loop_duration_seconds",
"Duration of the application main loop in seconds",
POD_ID_LABEL,
)
APP_ERRORS_TOTAL = Counter(
"app_errors_total",
"Total number of unhandled errors in the main loop",
POD_ID_LABEL,
)
APP_UP = Gauge(
"app_up",
"Indicates if the application is running (1) or shutting down (0)",
POD_ID_LABEL,
)
# --- Ingestor Manager Metrics ---
ACTIVE_INGESTORS = Gauge(
"ingestor_active_total",
"Number of active ingestors reported by Redis",
)
SLOTS_TOTAL = Gauge(
"ingestor_slots_total",
"Total number of slots configured in Redis",
)
LEASES_TOTAL = Gauge(
"ingestor_leases_total",
"Total number of leases (allocated slots) in Redis",
)
SLOTS_MANAGED = Gauge(
"ingestor_slots_managed_current",
"Number of slots currently managed by this ingestor instance",
POD_ID_LABEL,
)
SLOTS_ACQUIRED = Counter(
"ingestor_slots_acquired_total",
"Total number of slots acquired by this instance",
POD_ID_LABEL,
)
SLOTS_RELEASED = Counter(
"ingestor_slots_released_total",
"Total number of slots released by this instance",
POD_ID_LABEL,
)
OPC_MANAGERS_ACTIVE = Gauge(
"ingestor_opc_managers_active",
"Number of active OPC Managers in this instance",
POD_ID_LABEL,
)
OPC_SUBSCRIPTION_ERRORS = Counter(
"ingestor_opc_subscription_errors_total",
"Errors when trying to subscribe to OPC tags",
["pod_id", "server", "slot"],
)
# --- OPC Manager Metrics ---
OPC_CONNECTIONS_TOTAL = Counter(
"opc_connections_initiated_total",
"Total connection attempts to OPC servers",
["pod_id", "server_name"],
)
OPC_CONNECTIONS_FAILED = Counter(
"opc_connections_failed_total",
"Total failed connection attempts to OPC servers",
["pod_id", "server_name"],
)
OPC_CONNECTION_STATUS = Gauge(
"opc_connection_status",
"Connection status with the OPC server (1=connected, 0=disconnected)",
SERVER_LABELS,
)
OPC_SUBSCRIPTIONS_CREATED = Counter(
"opc_subscriptions_created_total",
"Total OPC subscriptions created",
["pod_id", "server_name", "slot_name"],
)
OPC_TAGS_SUBSCRIBED = Gauge(
"opc_tags_subscribed_current",
"Current number of OPC tags subscribed on a server",
["pod_id", "server_name"],
)
OPC_CYCLES_WITHOUT_DATA = Gauge(
"opc_cycles_without_data",
"Current number of cycles without receiving data from a server",
["pod_id", "server_name"],
)
OPC_RECONNECTIONS_TOTAL = Counter(
"opc_reconnections_tried_total",
"Reconnection attempts to an OPC server after a loss",
["pod_id", "server_name"],
)
# --- Data Manager (Kafka) Metrics ---
KAFKA_MESSAGES_SENT = Counter(
"kafka_messages_sent_total", "Total messages sent to Kafka", KAFKA_LABELS
)
KAFKA_MESSAGES_ERRORS = Counter(
"kafka_messages_errors_total",
"Total errors sending messages to Kafka",
KAFKA_LABELS,
)
KAFKA_CONNECTION_STATUS = Gauge(
"kafka_connection_status",
"Connection status with Kafka (1=connected, 0=disconnected)",
POD_ID_LABEL,
)
# --- Resource Manager (Redis) Metrics ---
REDIS_OPERATIONS_TOTAL = Counter(
"redis_operations_total", "Total number of Redis operations performed", REDIS_LABELS
)
REDIS_OPERATIONS_ERRORS = Counter(
"redis_operations_errors_total",
"Total number of errors in Redis operations",
REDIS_LABELS,
)
REDIS_OPERATIONS_DURATION = Histogram(
"redis_operations_duration_seconds",
"Duration of Redis operations in seconds",
REDIS_LABELS,
)
REDIS_CONNECTION_STATUS = Gauge(
"redis_connection_status",
"Connection status with Redis (1=connected, 0=disconnected)",
POD_ID_LABEL,
)
# --- Notification Metrics ---
NOTIFICATIONS_SENT = Counter(
"notifications_sent_total",
"Total number of notifications sent",
NOTIFICATION_LABELS,
)

View File

@@ -1,3 +1,4 @@
asyncua==1.1.5 asyncua==1.1.5
redis redis
git+ssh://git@github.com/Aignosi/sientia-dataops-library.git git+ssh://git@github.com/Aignosi/sientia-dataops-library.git
prometheus_client

0
tests/unit/__init__.py Normal file
View File

View File

View File

@@ -60,7 +60,8 @@ def test_initialize_opc_from_config(opc_manager, ingestor_manager):
'server_uri': 'http://opcua-server.simulator', 'server_uri': 'http://opcua-server.simulator',
'cert_path': '/path/to/cert', 'cert_path': '/path/to/cert',
'private_key_path': '/path/to/private_key', 'private_key_path': '/path/to/private_key',
'server_cert_path': '/path/to/server_cert' 'server_cert_path': '/path/to/server_cert',
'pod_id': 'test_pod'
} }
opc_manager.return_value = MagicMock() opc_manager.return_value = MagicMock()
@@ -70,6 +71,7 @@ def test_initialize_opc_from_config(opc_manager, ingestor_manager):
opc_manager.assert_called_once_with( opc_manager.assert_called_once_with(
server_config['name'], server_config['url'], ingestor_manager.data_manager, ingestor_manager.logger, server_config['name'], server_config['url'], ingestor_manager.data_manager, ingestor_manager.logger,
server_config['server_uri'], ingestor_manager.notification_handler, server_config['server_uri'], ingestor_manager.notification_handler,
server_config['pod_id'],
server_config['cert_path'], server_config['private_key_path'], server_config['cert_path'], server_config['private_key_path'],
server_config['server_cert_path'] server_config['server_cert_path']
) )
@@ -109,7 +111,8 @@ def test_initialize_opc_from_config_exception(traceback_mock, opc_manager, inges
@patch('ingestor.managers.ingestor_manager.OpcManager') @patch('ingestor.managers.ingestor_manager.OpcManager')
def test_update_opc_servers(opc_manager, ingestor_manager): @patch('ingestor.managers.ingestor_manager.metrics')
def test_update_opc_servers(metrics, opc_manager, ingestor_manager):
manager1 = MagicMock( manager1 = MagicMock(
config={"config": "config1"}) config={"config": "config1"})
manager2 = MagicMock( manager2 = MagicMock(
@@ -173,6 +176,9 @@ def test_update_opc_servers(opc_manager, ingestor_manager):
assert ingestor_manager.opc_managers['server2'] != mock assert ingestor_manager.opc_managers['server2'] != mock
metrics.OPC_MANAGERS_ACTIVE.labels.assert_called_once_with(pod_id=ingestor_manager.pod_id)
metrics.OPC_MANAGERS_ACTIVE.labels.return_value.set.assert_called_once_with(len(ingestor_manager.opc_managers))
def test_declare_active(ingestor_manager): def test_declare_active(ingestor_manager):
ingestor_manager.resource_manager.ingestor_heartbeat = MagicMock() ingestor_manager.resource_manager.ingestor_heartbeat = MagicMock()
@@ -194,36 +200,44 @@ def test_get_active_ingestors_empty(ingestor_manager):
ingestor_manager.resource_manager.get_all_ingestors.assert_called_once() ingestor_manager.resource_manager.get_all_ingestors.assert_called_once()
def test_get_number_of_leases_success(ingestor_manager): @patch('ingestor.managers.ingestor_manager.metrics')
def test_get_number_of_leases_success(metrics, ingestor_manager):
ingestor_manager.resource_manager.get_all_leases = MagicMock( ingestor_manager.resource_manager.get_all_leases = MagicMock(
return_value=["lease1", "lease2"]) return_value=["lease1", "lease2"])
result = ingestor_manager.get_number_of_leases() result = ingestor_manager.get_number_of_leases()
assert result == 2 assert result == 2
ingestor_manager.resource_manager.get_all_leases.assert_called_once() ingestor_manager.resource_manager.get_all_leases.assert_called_once()
metrics.LEASES_TOTAL.set.assert_called_once_with(2)
def test_get_number_of_leases_empty(ingestor_manager): @patch('ingestor.managers.ingestor_manager.metrics')
def test_get_number_of_leases_empty(metrics, ingestor_manager):
ingestor_manager.resource_manager.get_all_leases = MagicMock( ingestor_manager.resource_manager.get_all_leases = MagicMock(
return_value=None) return_value=None)
result = ingestor_manager.get_number_of_leases() result = ingestor_manager.get_number_of_leases()
assert result == 0 assert result == 0
ingestor_manager.resource_manager.get_all_leases.assert_called_once() ingestor_manager.resource_manager.get_all_leases.assert_called_once()
metrics.LEASES_TOTAL.set.assert_called_once_with(0)
def test_get_number_of_slots_success(ingestor_manager): @patch('ingestor.managers.ingestor_manager.metrics')
def test_get_number_of_slots_success(metrics, ingestor_manager):
ingestor_manager.resource_manager.get_all_slots = MagicMock( ingestor_manager.resource_manager.get_all_slots = MagicMock(
return_value=["slot1", "slot2"]) return_value=["slot1", "slot2"])
result = ingestor_manager.get_number_of_slots() result = ingestor_manager.get_number_of_slots()
assert result == 2 assert result == 2
ingestor_manager.resource_manager.get_all_slots.assert_called_once() ingestor_manager.resource_manager.get_all_slots.assert_called_once()
metrics.SLOTS_TOTAL.set.assert_called_once_with(2)
def test_get_number_of_slots_empty(ingestor_manager): @patch('ingestor.managers.ingestor_manager.metrics')
def test_get_number_of_slots_empty(metrics, ingestor_manager):
ingestor_manager.resource_manager.get_all_slots = MagicMock( ingestor_manager.resource_manager.get_all_slots = MagicMock(
return_value=None) return_value=None)
result = ingestor_manager.get_number_of_slots() result = ingestor_manager.get_number_of_slots()
assert result == 0 assert result == 0
ingestor_manager.resource_manager.get_all_slots.assert_called_once() ingestor_manager.resource_manager.get_all_slots.assert_called_once()
metrics.SLOTS_TOTAL.set.assert_called_once_with(0)
def test_get_slot_leases_1_success(ingestor_manager): def test_get_slot_leases_1_success(ingestor_manager):
@@ -337,13 +351,17 @@ def test_update_slot_config(ingestor_manager):
assert ingestor_manager.resource_manager.renew_tag_lease.call_count == 3 assert ingestor_manager.resource_manager.renew_tag_lease.call_count == 3
def test_drop_slot_leases(ingestor_manager): @patch('ingestor.managers.ingestor_manager.metrics')
def test_drop_slot_leases(metrics, ingestor_manager):
ingestor_manager.resource_manager.drop_tag_lease = MagicMock() ingestor_manager.resource_manager.drop_tag_lease = MagicMock()
ingestor_manager.drop_slot_leases(["1", "2"]) ingestor_manager.drop_slot_leases(["1", "2"])
ingestor_manager.resource_manager.drop_tag_lease.assert_any_call("1") ingestor_manager.resource_manager.drop_tag_lease.assert_any_call("1")
ingestor_manager.resource_manager.drop_tag_lease.assert_any_call("2") ingestor_manager.resource_manager.drop_tag_lease.assert_any_call("2")
metrics.SLOTS_RELEASED.labels.assert_any_call(pod_id=ingestor_manager.pod_id)
metrics.SLOTS_RELEASED.labels.return_value.inc.assert_any_call()
def test_manage_server_no_server(ingestor_manager): def test_manage_server_no_server(ingestor_manager):
ingestor_manager.opc_managers = { ingestor_manager.opc_managers = {
@@ -489,7 +507,8 @@ def test_subscribe_to_tags(ingestor_manager):
'server3', None) 'server3', None)
def test_check_opc_servers_integrity_all_healthy(ingestor_manager): @patch('ingestor.managers.ingestor_manager.metrics')
def test_check_opc_servers_integrity_all_healthy(metrics, ingestor_manager):
# Setup mock OPC managers # Setup mock OPC managers
opc_manager1 = MagicMock() opc_manager1 = MagicMock()
opc_manager1.check_cycles.return_value = None opc_manager1.check_cycles.return_value = None
@@ -521,6 +540,9 @@ def test_check_opc_servers_integrity_all_healthy(ingestor_manager):
# Verify that no reinitialization was needed # Verify that no reinitialization was needed
ingestor_manager.initialize_opc_from_config.assert_not_called() ingestor_manager.initialize_opc_from_config.assert_not_called()
metrics.OPC_MANAGERS_ACTIVE.labels.assert_called_once_with(pod_id=ingestor_manager.pod_id)
metrics.OPC_MANAGERS_ACTIVE.labels.return_value.set.assert_called_once_with(len(ingestor_manager.opc_managers))
def test_check_opc_servers_integrity_server_lost(ingestor_manager): def test_check_opc_servers_integrity_server_lost(ingestor_manager):
# Setup mock OPC manager that will be lost # Setup mock OPC manager that will be lost

View File

@@ -1,6 +1,7 @@
import json import json
from datetime import datetime from datetime import datetime
from unittest.mock import MagicMock, patch from unittest.mock import MagicMock, patch
from prometheus_client import Gauge
from pytest import fixture from pytest import fixture
from asyncua.crypto.security_policies import SecurityPolicyBasic256 from asyncua.crypto.security_policies import SecurityPolicyBasic256
@@ -9,57 +10,65 @@ from ingestor.managers.opc_manager import OpcManager
from sientia_do.notifications.models import NotificationLevel from sientia_do.notifications.models import NotificationLevel
tags = { tags = {
'ns=3;i=1001': { "ns=3;i=1001": {
'aggregation_function': 'LTS', "aggregation_function": "LTS",
'frequency': 1000, "frequency": 1000,
'max_value': 100, "max_value": 100,
'min_value': 0, "min_value": 0,
'tag_name': 'Counter' "tag_name": "Counter",
}, },
'ns=3;i=1003': { "ns=3;i=1003": {
'aggregation_function': 'AVG', "aggregation_function": "AVG",
'frequency': 1000, "frequency": 1000,
'max_value': 100, "max_value": 100,
'min_value': 0, "min_value": 0,
'tag_name': 'Random' "tag_name": "Random",
},
"ns=3;i=1004": {
"aggregation_function": "MDN",
"frequency": 1000,
"max_value": 100,
"min_value": 0,
"tag_name": "Sawtooth",
}, },
'ns=3;i=1004': {
'aggregation_function': 'MDN',
'frequency': 1000,
'max_value': 100,
'min_value': 0,
'tag_name': 'Sawtooth'
}
} }
@fixture @fixture
def raw_opc_manager(): def raw_opc_manager():
return OpcManager( return OpcManager(
'TestConnector', 'opc.tcp://localhost:4840', MagicMock(), "TestConnector",
MagicMock(), 'opc.tcp://localhost:4840', MagicMock() "opc.tcp://localhost:4840",
MagicMock(),
MagicMock(),
"opc.tcp://localhost:4840",
MagicMock(),
"localhost",
) )
@fixture @fixture
def opc_manager(raw_opc_manager): def opc_manager(raw_opc_manager):
raw_opc_manager.client = MagicMock() raw_opc_manager.client = MagicMock()
raw_opc_manager.cert_path = 'cert.pem' raw_opc_manager.cert_path = "cert.pem"
raw_opc_manager.private_key_path = 'private_key.pem' raw_opc_manager.private_key_path = "private_key.pem"
raw_opc_manager.server_cert_path = 'server_cert.pem' raw_opc_manager.server_cert_path = "server_cert.pem"
return raw_opc_manager return raw_opc_manager
@fixture @fixture
def opc_manager_subscribed(opc_manager): def opc_manager_subscribed(opc_manager):
opc_manager.subscriptions['sub1'] = MagicMock() opc_manager.subscriptions["sub1"] = MagicMock()
return opc_manager return opc_manager
def test___str__(opc_manager): def test___str__(opc_manager):
assert str(opc_manager) == 'OpcManager(name=TestConnector, url=opc.tcp://localhost:4840, server_uri=opc.tcp://localhost:4840)\nnodes={}, subscriptions={}' assert (
str(opc_manager)
== "OpcManager(name=TestConnector, url=opc.tcp://localhost:4840, server_uri=opc.tcp://localhost:4840)\nnodes={}, subscriptions={}"
)
def test_set_security_success(opc_manager): def test_set_security_success(opc_manager):
@@ -71,7 +80,7 @@ def test_set_security_success(opc_manager):
SecurityPolicyBasic256, SecurityPolicyBasic256,
certificate=opc_manager.cert_path, certificate=opc_manager.cert_path,
private_key=opc_manager.private_key_path, private_key=opc_manager.private_key_path,
server_certificate=opc_manager.server_cert_path server_certificate=opc_manager.server_cert_path,
) )
assert opc_manager.client.secure_channel_timeout == 10000000 assert opc_manager.client.secure_channel_timeout == 10000000
@@ -85,16 +94,19 @@ def test_set_security_no_cert(opc_manager):
try: try:
opc_manager.set_security() opc_manager.set_security()
except ValueError as e: except ValueError as e:
assert str( assert (
e) == "Certificate and private key paths must be provided for secure connection." str(e)
== "Certificate and private key paths must be provided for secure connection."
)
else: else:
assert False, "ValueError not raised" assert False, "ValueError not raised"
assert opc_manager.client.set_security.call_count == 0 assert opc_manager.client.set_security.call_count == 0
@patch('ingestor.managers.opc_manager.Client') @patch("ingestor.managers.opc_manager.metrics")
def test_connect_no_security(client, raw_opc_manager): @patch("ingestor.managers.opc_manager.Client")
def test_connect_no_security(client, mock_metrics, raw_opc_manager):
raw_opc_manager.set_security = MagicMock() raw_opc_manager.set_security = MagicMock()
raw_opc_manager.connect() raw_opc_manager.connect()
@@ -102,13 +114,25 @@ def test_connect_no_security(client, raw_opc_manager):
client.assert_called_once_with(raw_opc_manager.url) client.assert_called_once_with(raw_opc_manager.url)
raw_opc_manager.client.connect.assert_called_once() raw_opc_manager.client.connect.assert_called_once()
raw_opc_manager.set_security.assert_not_called() raw_opc_manager.set_security.assert_not_called()
mock_metrics.OPC_CONNECTIONS_TOTAL.labels.assert_called_once_with(
pod_id=raw_opc_manager.pod_id,
server_name=raw_opc_manager.name
)
mock_metrics.OPC_CONNECTIONS_TOTAL.labels.return_value.inc.assert_called_once()
mock_metrics.OPC_CONNECTION_STATUS.labels.assert_called_once_with(
pod_id=raw_opc_manager.pod_id,
server_name=raw_opc_manager.name,
server_url=raw_opc_manager.url
)
mock_metrics.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(1)
mock_metrics.OPC_CONNECTIONS_FAILED.labels.assert_not_called()
@patch('ingestor.managers.opc_manager.Client') @patch("ingestor.managers.opc_manager.Client")
def test_connect_with_security(client, raw_opc_manager): def test_connect_with_security(client, raw_opc_manager):
raw_opc_manager.cert_path = 'cert.pem' raw_opc_manager.cert_path = "cert.pem"
raw_opc_manager.private_key_path = 'private_key.pem' raw_opc_manager.private_key_path = "private_key.pem"
raw_opc_manager.server_cert_path = 'server_cert.pem' raw_opc_manager.server_cert_path = "server_cert.pem"
raw_opc_manager.set_security = MagicMock() raw_opc_manager.set_security = MagicMock()
raw_opc_manager.connect() raw_opc_manager.connect()
@@ -118,9 +142,48 @@ def test_connect_with_security(client, raw_opc_manager):
raw_opc_manager.set_security.assert_called_once() raw_opc_manager.set_security.assert_called_once()
@patch("ingestor.managers.opc_manager.Client")
@patch("ingestor.managers.opc_manager.metrics")
def test_connect_exception_handling_and_metrics(
mock_metrics_module, mock_opc_client_class, raw_opc_manager
):
mock_client_instance = mock_opc_client_class.return_value
simulated_error_message = "Erro de conexão simulado"
mock_client_instance.connect.side_effect = Exception(simulated_error_message)
opc_manager_instance = raw_opc_manager
opc_manager_instance.cert_path = None
with pytest.raises(Exception, match=simulated_error_message):
opc_manager_instance.connect()
mock_metrics_module.OPC_CONNECTIONS_TOTAL.labels.assert_called_once_with(
pod_id=opc_manager_instance.pod_id, server_name=opc_manager_instance.name
)
mock_metrics_module.OPC_CONNECTIONS_TOTAL.labels.return_value.inc.assert_called_once()
mock_metrics_module.OPC_CONNECTION_STATUS.labels.assert_called_once_with(
pod_id=opc_manager_instance.pod_id,
server_name=opc_manager_instance.name,
server_url=opc_manager_instance.url,
)
mock_metrics_module.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(
0
)
mock_metrics_module.OPC_CONNECTIONS_FAILED.labels.assert_called_once_with(
pod_id=opc_manager_instance.pod_id, server_name=opc_manager_instance.name
)
mock_metrics_module.OPC_CONNECTIONS_FAILED.labels.return_value.inc.assert_called_once()
opc_manager_instance.logger.error.assert_called_once_with(
f"Failed to connect to {opc_manager_instance.name}: {simulated_error_message}"
)
def test_create_subscription_no_client(raw_opc_manager): def test_create_subscription_no_client(raw_opc_manager):
try: try:
raw_opc_manager.create_subscription('sub1') raw_opc_manager.create_subscription("sub1")
except ValueError as e: except ValueError as e:
assert str(e) == "Client not connected. Call connect first." assert str(e) == "Client not connected. Call connect first."
else: else:
@@ -128,55 +191,95 @@ def test_create_subscription_no_client(raw_opc_manager):
def test_create_subscription_success_has_period(opc_manager): def test_create_subscription_success_has_period(opc_manager):
opc_manager.create_subscription('sub1', 1000) opc_manager.create_subscription("sub1", 1000)
opc_manager.client.create_subscription.assert_called_once_with( opc_manager.client.create_subscription.assert_called_once_with(1000, opc_manager)
1000, opc_manager) assert opc_manager.subscriptions["sub1"] is not None
assert opc_manager.subscriptions['sub1'] is not None
def test_create_subscription_success_no_period(opc_manager): def test_create_subscription_success_no_period(opc_manager):
opc_manager.create_subscription('sub1', None) opc_manager.create_subscription("sub1", None)
opc_manager.client.create_subscription.assert_called_once_with( opc_manager.client.create_subscription.assert_called_once_with(500, opc_manager)
500, opc_manager) assert opc_manager.subscriptions["sub1"] is not None
assert opc_manager.subscriptions['sub1'] is not None
def test_subscribe_no_subscription(opc_manager): @patch("ingestor.managers.opc_manager.metrics")
def test_create_subscription_with_metrics(metrics, opc_manager):
opc_manager.create_subscription("sub1", 1000)
metrics.OPC_SUBSCRIPTIONS_CREATED.labels.assert_called_once_with(
pod_id=opc_manager.pod_id, server_name=opc_manager.name, slot_name="sub1"
)
metrics.OPC_SUBSCRIPTIONS_CREATED.labels.return_value.inc.assert_called_once()
@patch("ingestor.managers.opc_manager.metrics")
def test_create_subscription_exception_during_client_call(
mock_metrics_module, raw_opc_manager
):
opc_manager_instance = raw_opc_manager
opc_manager_instance.client = MagicMock()
subscription_name = "test_sub_client_error"
simulated_period = 750
simulated_error_message = "Falha ao criar subscrição no cliente OPC"
opc_manager_instance.client.create_subscription.side_effect = Exception(
simulated_error_message
)
with pytest.raises(Exception, match=simulated_error_message):
opc_manager_instance.create_subscription(
subscription_name, period=simulated_period
)
opc_manager_instance.client.create_subscription.assert_called_once_with(
simulated_period, opc_manager_instance
)
opc_manager_instance.logger.error.assert_called_once_with(
f"Failed to create subscription {subscription_name} on {opc_manager_instance.name}: {simulated_error_message}"
)
mock_metrics_module.OPC_SUBSCRIPTIONS_CREATED.labels.assert_not_called()
@patch("ingestor.managers.opc_manager.metrics")
def test_subscribe_no_subscription(metrics, opc_manager):
try: try:
opc_manager.subscribe('sub1', tags, 1000) opc_manager.subscribe("sub1", tags, 1000)
except ValueError as e: except ValueError as e:
assert str( assert str(e) == "Subscription not created. Call create_subscription first."
e) == "Subscription not created. Call create_subscription first."
else: else:
assert False, "ValueError not raised" assert False, "ValueError not raised"
metrics.OPC_TAGS_SUBSCRIBED.labels.assert_not_called()
def test_subscribe_success(opc_manager_subscribed): def test_subscribe_success(opc_manager_subscribed):
opc_manager_subscribed.nodes = { opc_manager_subscribed.nodes = {"ns=3;i=1001": "data"}
'ns=3;i=1001': 'data'
}
opc_manager_subscribed.subscribe('sub1', tags, 1000) opc_manager_subscribed.subscribe("sub1", tags, 1000)
assert opc_manager_subscribed.nodes == tags assert opc_manager_subscribed.nodes == tags
assert opc_manager_subscribed.addr_nodes == [ assert opc_manager_subscribed.addr_nodes == [
opc_manager_subscribed.client.get_node(n) for n in tags if n != 'ns=3;i=1001'] opc_manager_subscribed.client.get_node(n) for n in tags if n != "ns=3;i=1001"
]
def test_unsubscribe_no_subscription(opc_manager): def test_unsubscribe_no_subscription(opc_manager):
opc_manager.unsubscribe('sub1') opc_manager.unsubscribe("sub1")
opc_manager.logger.warning.assert_called_once_with( opc_manager.logger.warning.assert_called_once_with(
"Subscription 'sub1' not found. Cannot unsubscribe.") "Subscription 'sub1' not found. Cannot unsubscribe."
assert opc_manager.subscriptions.get('sub1') is None )
assert opc_manager.subscriptions.get("sub1") is None
def test_unsubscribe_success(opc_manager_subscribed): def test_unsubscribe_success(opc_manager_subscribed):
opc_manager_subscribed.unsubscribe('sub1') opc_manager_subscribed.unsubscribe("sub1")
opc_manager_subscribed.subscriptions.get('sub1') is None opc_manager_subscribed.subscriptions.get("sub1") is None
def test_disconnect_success(opc_manager_subscribed): def test_disconnect_success(opc_manager_subscribed):
@@ -184,85 +287,123 @@ def test_disconnect_success(opc_manager_subscribed):
opc_manager_subscribed.disconnect() opc_manager_subscribed.disconnect()
opc_manager_subscribed.subscriptions['sub1'].delete.assert_called_once() opc_manager_subscribed.subscriptions["sub1"].delete.assert_called_once()
assert opc_manager_subscribed.client is None assert opc_manager_subscribed.client is None
def test_disconnect_error_unsubscribe(opc_manager_subscribed): def test_disconnect_error_unsubscribe(opc_manager_subscribed):
opc_manager_subscribed.client = MagicMock() opc_manager_subscribed.client = MagicMock()
opc_manager_subscribed.subscriptions['sub1'] = MagicMock( opc_manager_subscribed.subscriptions["sub1"] = MagicMock(
delete=MagicMock(side_effect=Exception("Test error")) delete=MagicMock(side_effect=Exception("Test error"))
) )
opc_manager_subscribed.disconnect() opc_manager_subscribed.disconnect()
opc_manager_subscribed.subscriptions['sub1'].delete.assert_called_once() opc_manager_subscribed.subscriptions["sub1"].delete.assert_called_once()
opc_manager_subscribed.client = None opc_manager_subscribed.client = None
opc_manager_subscribed.logger.error.assert_called_once_with( opc_manager_subscribed.logger.error.assert_called_once_with(
"Failed to clean up subscription: Test error") "Failed to clean up subscription: Test error"
)
def test_disconnect_error(opc_manager_subscribed): def test_disconnect_error(opc_manager_subscribed):
opc_manager_subscribed.client = MagicMock() opc_manager_subscribed.client = MagicMock()
opc_manager_subscribed.client.disconnect = MagicMock( opc_manager_subscribed.client.disconnect = MagicMock(
side_effect=Exception("Test error")) side_effect=Exception("Test error")
)
opc_manager_subscribed.disconnect() opc_manager_subscribed.disconnect()
opc_manager_subscribed.subscriptions['sub1'].delete.assert_called_once() opc_manager_subscribed.subscriptions["sub1"].delete.assert_called_once()
opc_manager_subscribed.client = None opc_manager_subscribed.client = None
opc_manager_subscribed.logger.error.assert_called_once_with( opc_manager_subscribed.logger.error.assert_called_once_with(
"Failed to disconnect from OPC UA server: Test error") "Failed to disconnect from OPC UA server: Test error"
)
def test_datachange_notification(opc_manager_subscribed): @patch('ingestor.managers.opc_manager.metrics')
def test_disconnect_metrics_on_successful_path(mock_metrics_module, raw_opc_manager):
mock_metrics_module.OPC_CONNECTION_STATUS.reset_mock()
mock_metrics_module.OPC_TAGS_SUBSCRIBED.reset_mock()
raw_opc_manager.client = MagicMock()
mock_sub1 = MagicMock()
mock_sub2 = MagicMock()
raw_opc_manager.subscriptions = {"sub1": mock_sub1, "sub2": mock_sub2}
raw_opc_manager.disconnect()
mock_metrics_module.OPC_CONNECTION_STATUS.labels.assert_called_once_with(
pod_id=raw_opc_manager.pod_id,
server_name=raw_opc_manager.name,
server_url=raw_opc_manager.url
)
mock_metrics_module.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(0)
mock_metrics_module.OPC_TAGS_SUBSCRIBED.labels.assert_called_once_with(
pod_id=raw_opc_manager.pod_id,
server_name=raw_opc_manager.name
)
mock_metrics_module.OPC_TAGS_SUBSCRIBED.labels.return_value.set.assert_called_once_with(0)
@patch('ingestor.managers.opc_manager.metrics')
def test_datachange_notification(metrics, opc_manager_subscribed):
data = MagicMock( data = MagicMock(
monitored_item=MagicMock( monitored_item=MagicMock(
Value=MagicMock( Value=MagicMock(
Value=MagicMock(Value=42), Value=MagicMock(Value=42),
SourceTimestamp=datetime.strptime( SourceTimestamp=datetime.strptime(
'2021-01-01T00:00:00', '%Y-%m-%dT%H:%M:%S') "2021-01-01T00:00:00", "%Y-%m-%dT%H:%M:%S"
))) ),
)
)
)
opc_manager_subscribed.nodes = { opc_manager_subscribed.nodes = {
'ns=3;i=1001': { "ns=3;i=1001": {
'tag_name': 'Counter', "tag_name": "Counter",
'cycle_rule': { "cycle_rule": {"cycle_increment": 1.0, "cycle_count": 2},
'cycle_increment': 1.0, "topics": ["topic1", "topic2"],
'cycle_count': 2
},
'topics': ['topic1', 'topic2']
} }
} }
opc_manager_subscribed.datachange_notification( metrics.OPC_CYCLES_WITHOUT_DATA.reset_mock()
'ns=3;i=1001', None, data)
opc_manager_subscribed.datachange_notification("ns=3;i=1001", None, data)
opc_manager_subscribed.data_manager.publish.assert_any_call( opc_manager_subscribed.data_manager.publish.assert_any_call(
'topic1', { "topic1",
'tag': 'ns=3;i=1001', {
'name': 'Counter', "tag": "ns=3;i=1001",
'timestamp': '2021-01-01 00:00:00', "name": "Counter",
'value': 42 "timestamp": "2021-01-01 00:00:00",
}) "value": 42,
},
)
opc_manager_subscribed.data_manager.publish.assert_any_call( opc_manager_subscribed.data_manager.publish.assert_any_call(
'topic2', { "topic2",
'tag': 'ns=3;i=1001', {
'name': 'Counter', "tag": "ns=3;i=1001",
'timestamp': '2021-01-01 00:00:00', "name": "Counter",
'value': 42 "timestamp": "2021-01-01 00:00:00",
}) "value": 42,
assert opc_manager_subscribed.nodes['ns=3;i=1001']['cycle_rule']['cycle_count'] == 0 },
)
assert opc_manager_subscribed.nodes["ns=3;i=1001"]["cycle_rule"]["cycle_count"] == 0
metrics.OPC_CYCLES_WITHOUT_DATA.labels.assert_called_once_with(
pod_id=opc_manager_subscribed.pod_id,
server_name=opc_manager_subscribed.name
)
metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with(0)
def test_check_cycles_no_notification(opc_manager): def test_check_cycles_no_notification(opc_manager):
# Setup: node with cycle_count just below threshold # Setup: node with cycle_count just below threshold
opc_manager.nodes = { opc_manager.nodes = {
'ns=3;i=1001': { "ns=3;i=1001": {
'tag_name': 'Counter', "tag_name": "Counter",
'cycle_rule': { "cycle_rule": {"cycle_increment": 1.0, "cycle_count": 3.0},
'cycle_increment': 1.0,
'cycle_count': 3.0
}
} }
} }
opc_manager.notification_handler.build_and_send_notification = MagicMock() opc_manager.notification_handler.build_and_send_notification = MagicMock()
@@ -270,20 +411,18 @@ def test_check_cycles_no_notification(opc_manager):
opc_manager.check_cycles() opc_manager.check_cycles()
# After one increment, cycle_count = 4.0, still below threshold # After one increment, cycle_count = 4.0, still below threshold
assert opc_manager.nodes['ns=3;i=1001']['cycle_rule']['cycle_count'] == pytest.approx( assert opc_manager.nodes["ns=3;i=1001"]["cycle_rule"][
4.0) "cycle_count"
] == pytest.approx(4.0)
opc_manager.notification_handler.build_and_send_notification.assert_not_called() opc_manager.notification_handler.build_and_send_notification.assert_not_called()
def test_check_cycles_triggers_notification(opc_manager): def test_check_cycles_triggers_notification(opc_manager):
# Setup: node with cycle_count just below threshold, increment will cross threshold # Setup: node with cycle_count just below threshold, increment will cross threshold
opc_manager.nodes = { opc_manager.nodes = {
'ns=3;i=1001': { "ns=3;i=1001": {
'tag_name': 'Counter', "tag_name": "Counter",
'cycle_rule': { "cycle_rule": {"cycle_increment": 2.5, "cycle_count": 3.0},
'cycle_increment': 2.5,
'cycle_count': 3.0
}
} }
} }
opc_manager.notification_handler.build_and_send_notification = MagicMock() opc_manager.notification_handler.build_and_send_notification = MagicMock()
@@ -291,19 +430,23 @@ def test_check_cycles_triggers_notification(opc_manager):
opc_manager.check_cycles() opc_manager.check_cycles()
# After increment, cycle_count = 5.5, should trigger notification # After increment, cycle_count = 5.5, should trigger notification
assert opc_manager.nodes['ns=3;i=1001']['cycle_rule']['cycle_count'] == pytest.approx( assert opc_manager.nodes["ns=3;i=1001"]["cycle_rule"][
5.5) "cycle_count"
] == pytest.approx(5.5)
opc_manager.notification_handler.build_and_send_notification.assert_called_once_with( opc_manager.notification_handler.build_and_send_notification.assert_called_once_with(
notification_id='TAG_ns=3;i=1001:Counter_LISTENNING_STOPPED', notification_id="TAG_ns=3;i=1001:Counter_LISTENNING_STOPPED",
message='5.5 cycles without receive from ns=3;i=1001:Counter', message="5.5 cycles without receive from ns=3;i=1001:Counter",
block="opc_manager", block="opc_manager",
level=NotificationLevel.WARNING level=NotificationLevel.WARNING,
) )
def test_check_opc_listenning_no_notification(opc_manager): @patch('ingestor.managers.opc_manager.metrics')
def test_check_opc_listenning_no_notification(metrics, opc_manager):
opc_manager.non_receive_count = 3 opc_manager.non_receive_count = 3
opc_manager.notification_handler.build_and_send_notification = MagicMock() opc_manager.notification_handler.build_and_send_notification = MagicMock()
metrics.OPC_CYCLES_WITHOUT_DATA.reset_mock()
metrics.OPC_RECONNECTIONS_TOTAL.reset_mock()
result = opc_manager.check_opc_listenning() result = opc_manager.check_opc_listenning()
@@ -311,6 +454,13 @@ def test_check_opc_listenning_no_notification(opc_manager):
opc_manager.notification_handler.build_and_send_notification.assert_not_called() opc_manager.notification_handler.build_and_send_notification.assert_not_called()
assert result is False assert result is False
metrics.OPC_CYCLES_WITHOUT_DATA.labels.assert_called_once_with(
pod_id=opc_manager.pod_id,
server_name=opc_manager.name
)
metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with(opc_manager.non_receive_count)
metrics.OPC_RECONNECTIONS_TOTAL.labels.assert_not_called()
def test_check_opc_listenning_warning_notification(opc_manager): def test_check_opc_listenning_warning_notification(opc_manager):
opc_manager.non_receive_count = 4 opc_manager.non_receive_count = 4
@@ -320,17 +470,20 @@ def test_check_opc_listenning_warning_notification(opc_manager):
assert opc_manager.non_receive_count == 5 assert opc_manager.non_receive_count == 5
opc_manager.notification_handler.build_and_send_notification.assert_called_once_with( opc_manager.notification_handler.build_and_send_notification.assert_called_once_with(
notification_id=f'OPC_LISTENNING_STOPPED__{opc_manager.name}', notification_id=f"OPC_LISTENNING_STOPPED__{opc_manager.name}",
message=f'5 cycles without receive from OPC {opc_manager.name}. Tags: {json.dumps(opc_manager.nodes)}', message=f"5 cycles without receive from OPC {opc_manager.name}. Tags: {json.dumps(opc_manager.nodes)}",
block="opc_manager", block="opc_manager",
level=NotificationLevel.ERROR level=NotificationLevel.ERROR,
) )
assert result is False assert result is False
def test_check_opc_listenning_error_notification_and_retry(opc_manager): @patch('ingestor.managers.opc_manager.metrics')
def test_check_opc_listenning_error_notification_and_retry(metrics, opc_manager):
opc_manager.non_receive_count = 14 opc_manager.non_receive_count = 14
opc_manager.notification_handler.build_and_send_notification = MagicMock() opc_manager.notification_handler.build_and_send_notification = MagicMock()
metrics.OPC_CYCLES_WITHOUT_DATA.reset_mock()
metrics.OPC_RECONNECTIONS_TOTAL.reset_mock()
result = opc_manager.check_opc_listenning() result = opc_manager.check_opc_listenning()
@@ -340,16 +493,56 @@ def test_check_opc_listenning_error_notification_and_retry(opc_manager):
calls = opc_manager.notification_handler.build_and_send_notification.call_args_list calls = opc_manager.notification_handler.build_and_send_notification.call_args_list
# First call: 5 cycles warning # First call: 5 cycles warning
assert calls[0].kwargs == dict( assert calls[0].kwargs == dict(
notification_id=f'OPC_LISTENNING_STOPPED__{opc_manager.name}', notification_id=f"OPC_LISTENNING_STOPPED__{opc_manager.name}",
message=f'15 cycles without receive from OPC {opc_manager.name}. Tags: {json.dumps(opc_manager.nodes)}', message=f"15 cycles without receive from OPC {opc_manager.name}. Tags: {json.dumps(opc_manager.nodes)}",
block="opc_manager", block="opc_manager",
level=NotificationLevel.ERROR level=NotificationLevel.ERROR,
) )
# Second call: 15 cycles retry # Second call: 15 cycles retry
assert calls[1].kwargs == dict( assert calls[1].kwargs == dict(
notification_id=f'OPC_CONNECTION_RETRY__{opc_manager.name}', notification_id=f"OPC_CONNECTION_RETRY__{opc_manager.name}",
message=f'Retrying to connect to server {opc_manager.name}', message=f"Retrying to connect to server {opc_manager.name}",
block="opc_manager", block="opc_manager",
level=NotificationLevel.ERROR level=NotificationLevel.ERROR,
) )
assert result is True assert result is True
metrics.OPC_CYCLES_WITHOUT_DATA.labels.assert_called_once_with(
pod_id=opc_manager.pod_id,
server_name=opc_manager.name
)
metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with(opc_manager.non_receive_count)
metrics.OPC_RECONNECTIONS_TOTAL.labels.assert_called_once_with(
pod_id=opc_manager.pod_id,
server_name=opc_manager.name
)
metrics.OPC_RECONNECTIONS_TOTAL.labels.return_value.inc.assert_called_once()
@patch("ingestor.managers.opc_manager.metrics")
def test_init_metrics_calls_correct_metric_methods(mock_metrics):
opc_manager = OpcManager(
name="TestInitConnector",
url="opc.tcp://init.test:4840",
data_manager=MagicMock(),
logger=MagicMock(),
server_uri="opc.tcp://init.test:4840/uri",
notification_handler=MagicMock(),
pod_id="init_pod_localhost",
)
mock_metrics.OPC_CONNECTION_STATUS.labels.assert_called_once_with(
pod_id=opc_manager.pod_id,
server_name=opc_manager.name,
server_url=opc_manager.url,
)
mock_metrics.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(0)
mock_metrics.OPC_TAGS_SUBSCRIBED.labels.assert_called_once_with(
pod_id=opc_manager.pod_id,
server_name=opc_manager.name
)
mock_metrics.OPC_TAGS_SUBSCRIBED.labels.return_value.set.assert_called_once_with(0)

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"
)

277
tests/unit/test_app.py Normal file
View File

@@ -0,0 +1,277 @@
import pytest
from unittest.mock import patch, MagicMock, call
from threading import Event
import signal as signal_module # To avoid conflict with mock names
import os
# Import the 'app' module to be tested
from ingestor import app
# Custom exception to catch os._exit calls
class OsExitCalled(Exception):
def __init__(self, code):
super().__init__(f"os._exit({code}) called")
self.code = code
# Helper function for the os_exit mock's side_effect
def raise_os_exit_with_code(exit_code):
raise OsExitCalled(exit_code)
@pytest.fixture
def mock_app_env(monkeypatch):
"""Fixture to mock dependencies of app.main and app.signal_handler."""
mocks = {
"start_http_server": MagicMock(),
"Ingestor": MagicMock(),
"metrics_APP_UP_labels_set": MagicMock(),
"metrics_APP_LOOP_COUNT_labels_inc": MagicMock(),
"metrics_APP_LOOP_DURATION_labels_observe": MagicMock(),
"metrics_APP_ERRORS_TOTAL_labels_inc": MagicMock(),
"os_exit": MagicMock(side_effect=raise_os_exit_with_code), # CORRECTED
"time_time": MagicMock(),
"time_sleep": MagicMock(),
"signal_signal": MagicMock(),
"traceback_print_exc": MagicMock(),
"mock_exit_signal": MagicMock(spec=Event),
}
monkeypatch.setattr(app, "start_http_server", mocks["start_http_server"])
monkeypatch.setattr(app, "Ingestor", mocks["Ingestor"])
monkeypatch.setattr(
app.metrics.APP_UP,
"labels",
MagicMock(return_value=MagicMock(set=mocks["metrics_APP_UP_labels_set"])),
)
monkeypatch.setattr(
app.metrics.APP_LOOP_COUNT,
"labels",
MagicMock(
return_value=MagicMock(inc=mocks["metrics_APP_LOOP_COUNT_labels_inc"])
),
)
monkeypatch.setattr(
app.metrics.APP_LOOP_DURATION,
"labels",
MagicMock(
return_value=MagicMock(
observe=mocks["metrics_APP_LOOP_DURATION_labels_observe"]
)
),
)
monkeypatch.setattr(
app.metrics.APP_ERRORS_TOTAL,
"labels",
MagicMock(
return_value=MagicMock(inc=mocks["metrics_APP_ERRORS_TOTAL_labels_inc"])
),
)
monkeypatch.setattr(app.os, "_exit", mocks["os_exit"])
monkeypatch.setattr(app, "time", mocks["time_time"])
monkeypatch.setattr(app, "sleep", mocks["time_sleep"])
monkeypatch.setattr(app.signal, "signal", mocks["signal_signal"])
monkeypatch.setattr(app.traceback, "print_exc", mocks["traceback_print_exc"])
monkeypatch.setattr(app, "exit_signal", mocks["mock_exit_signal"])
monkeypatch.setattr(app, "POD_ID", "test_pod")
mock_ingestor_instance = mocks["Ingestor"].return_value
mock_ingestor_instance.poll_interval = 0.01
mock_ingestor_instance.logger = MagicMock()
return mocks
def test_main_successful_run_one_loop(mock_app_env, capsys):
"""Test a successful run where the loop executes once and then exits gracefully."""
mock_ingestor_instance = mock_app_env["Ingestor"].return_value
mock_exit_signal = mock_app_env["mock_exit_signal"]
mock_exit_signal.is_set.side_effect = [False, True]
mock_app_env["time_time"].side_effect = [10.0, 11.5]
with pytest.raises(OsExitCalled) as excinfo:
app.main()
assert excinfo.value.code == 0
mock_app_env["start_http_server"].assert_called_once_with(8000)
app.metrics.APP_UP.labels.assert_any_call(pod_id="test_pod")
set_calls = mock_app_env["metrics_APP_UP_labels_set"].call_args_list
assert call(1) in set_calls
assert call(0) in set_calls
assert set_calls.index(call(1)) < set_calls.index(call(0))
mock_app_env["Ingestor"].assert_called_once_with()
mock_ingestor_instance.prepare_ingestor.assert_called_once()
mock_ingestor_instance.logger.info.assert_any_call(
"Ingestor prepared. Starting main loop."
)
mock_ingestor_instance.loop.assert_called_once()
app.metrics.APP_LOOP_COUNT.labels.assert_called_with(pod_id="test_pod")
mock_app_env["metrics_APP_LOOP_COUNT_labels_inc"].assert_called_once()
mock_exit_signal.wait.assert_called_once_with(mock_ingestor_instance.poll_interval)
app.metrics.APP_LOOP_DURATION.labels.assert_called_with(pod_id="test_pod")
mock_app_env["metrics_APP_LOOP_DURATION_labels_observe"].assert_called_once_with(
1.5
)
mock_ingestor_instance.shutdown.assert_called_once()
mock_app_env["time_sleep"].assert_called_once_with(5)
# CORREÇÃO APLICADA ABAIXO:
mock_ingestor_instance.logger.info.assert_any_call("Main loop exit_signaled.")
captured = capsys.readouterr()
assert "Prometheus server started on port 8000." in captured.out
def test_main_prometheus_server_fails_to_start(mock_app_env, capsys):
"""Test the scenario where starting the Prometheus server fails."""
mock_app_env["start_http_server"].side_effect = OSError("Port already in use")
with pytest.raises(OsExitCalled) as excinfo:
app.main()
assert excinfo.value.code == 1
# Verify that APP_UP.labels(...).set(1) was NOT called.
# The mock for .set is mock_app_env['metrics_APP_UP_labels_set']
# We need to check if it was called with 1.
# A more robust way is to check if the specific .labels(pod_id="test_pod") mock was ever called
# and then its .set(1) method.
# For simplicity here, we check if metrics_APP_UP_labels_set was ever called with 1.
# Check that .set(1) was not called. .set(0) definitely not called.
called_with_1 = False
for call_args in mock_app_env["metrics_APP_UP_labels_set"].call_args_list:
if call_args == call(1):
called_with_1 = True
break
assert (
not called_with_1
), "APP_UP.set(1) should not have been called if server start failed"
mock_app_env["Ingestor"].assert_not_called()
captured = capsys.readouterr()
assert "Failed to start Prometheus server: Port already in use" in captured.out
def test_main_loop_exception_handling(mock_app_env, capsys):
"""Test that an exception in ingestor.loop() is handled gracefully."""
mock_ingestor_instance = mock_app_env["Ingestor"].return_value
mock_exit_signal = mock_app_env["mock_exit_signal"]
mock_exit_signal.is_set.side_effect = [False, True]
mock_ingestor_instance.loop.side_effect = Exception("Test loop exception")
mock_app_env["time_time"].side_effect = [10.0, 10.1]
with pytest.raises(OsExitCalled) as excinfo:
app.main()
assert excinfo.value.code == 0
mock_ingestor_instance.loop.assert_called_once()
mock_app_env["traceback_print_exc"].assert_called_once()
app.metrics.APP_ERRORS_TOTAL.labels.assert_called_with(pod_id="test_pod")
mock_app_env["metrics_APP_ERRORS_TOTAL_labels_inc"].assert_called_once()
mock_exit_signal.set.assert_called_once()
mock_app_env["metrics_APP_LOOP_COUNT_labels_inc"].assert_not_called()
app.metrics.APP_LOOP_DURATION.labels.assert_called_with(pod_id="test_pod")
mock_app_env["metrics_APP_LOOP_DURATION_labels_observe"].assert_called_once_with(
pytest.approx(0.1)
)
mock_ingestor_instance.shutdown.assert_called_once()
captured = capsys.readouterr()
assert "Exception in main loop. Setting exit_signal flag." in captured.out
def test_main_keyboard_interrupt_handling(mock_app_env, capsys):
"""Test that KeyboardInterrupt in ingestor.loop() is handled."""
mock_ingestor_instance = mock_app_env["Ingestor"].return_value
mock_exit_signal = mock_app_env["mock_exit_signal"]
mock_exit_signal.is_set.side_effect = [False, True]
mock_ingestor_instance.loop.side_effect = KeyboardInterrupt()
mock_app_env["time_time"].side_effect = [10.0, 10.1]
with pytest.raises(OsExitCalled) as excinfo:
app.main()
assert excinfo.value.code == 0
mock_ingestor_instance.loop.assert_called_once()
mock_exit_signal.set.assert_called_once()
mock_ingestor_instance.shutdown.assert_called_once()
captured = capsys.readouterr()
assert "KeyboardInterrupt received. Setting exit_signal flag." in captured.out
mock_app_env["metrics_APP_ERRORS_TOTAL_labels_inc"].assert_not_called()
def test_signal_handler_sets_exit_signal(mock_app_env):
"""Test that the signal_handler function calls exit_signal.set()."""
mock_exit_signal_set = mock_app_env["mock_exit_signal"].set
app.signal_handler(signal_module.SIGINT, None)
mock_exit_signal_set.assert_called_once()
def test_main_multiple_loop_iterations(mock_app_env):
"""Test the main loop runs for a few iterations."""
mock_ingestor_instance = mock_app_env["Ingestor"].return_value
mock_exit_signal = mock_app_env["mock_exit_signal"]
mock_exit_signal.is_set.side_effect = [False, False, False, True]
mock_app_env["time_time"].side_effect = [10.0, 10.1, 10.2, 10.3, 10.4, 10.5]
with pytest.raises(OsExitCalled) as excinfo:
app.main()
assert excinfo.value.code == 0
assert mock_ingestor_instance.loop.call_count == 3
assert app.metrics.APP_LOOP_COUNT.labels.call_count == 3
app.metrics.APP_LOOP_COUNT.labels.assert_called_with(
pod_id="test_pod"
) # Checks last call or any call
assert mock_app_env["metrics_APP_LOOP_COUNT_labels_inc"].call_count == 3
assert mock_exit_signal.wait.call_count == 3
assert app.metrics.APP_LOOP_DURATION.labels.call_count == 3
app.metrics.APP_LOOP_DURATION.labels.assert_called_with(pod_id="test_pod")
duration_calls = mock_app_env[
"metrics_APP_LOOP_DURATION_labels_observe"
].call_args_list
assert duration_calls[0] == call(pytest.approx(0.1, abs=1e-9))
assert duration_calls[1] == call(pytest.approx(0.1, abs=1e-9))
assert duration_calls[2] == call(pytest.approx(0.1, abs=1e-9))
mock_ingestor_instance.shutdown.assert_called_once()
def test_main_pod_id_used_in_metrics(mock_app_env):
"""Test that the POD_ID from app module is used in metric labels."""
mock_exit_signal = mock_app_env["mock_exit_signal"]
mock_exit_signal.is_set.side_effect = [False, True]
mock_app_env["time_time"].side_effect = [10.0, 11.0]
with pytest.raises(OsExitCalled):
app.main()
app.metrics.APP_UP.labels.assert_any_call(pod_id="test_pod")
app.metrics.APP_LOOP_COUNT.labels.assert_any_call(pod_id="test_pod")
app.metrics.APP_LOOP_DURATION.labels.assert_any_call(pod_id="test_pod")
# APP_ERRORS_TOTAL would be checked similarly if it were called in this flow.
# Check the .set() / .inc() calls on the mocks returned by .labels()
mock_app_env["metrics_APP_UP_labels_set"].assert_any_call(1)
mock_app_env["metrics_APP_UP_labels_set"].assert_any_call(0)
mock_app_env["metrics_APP_LOOP_COUNT_labels_inc"].assert_called_once()

253
tests/unit/test_metrics.py Normal file
View File

@@ -0,0 +1,253 @@
# tests/unit/test_metrics.py
import pytest
from prometheus_client import Counter, Gauge, Histogram
import ingestor.metrics as metrics
# --- Test Functions for Each Metric (Corrected for v0.22.0 _name behavior) ---
def test_app_loop_count():
"""Verify the definition of APP_LOOP_COUNT."""
assert metrics.APP_LOOP_COUNT is not None
assert isinstance(metrics.APP_LOOP_COUNT, Counter)
assert metrics.APP_LOOP_COUNT._name == "app_main_loop" # REMOVED _total
assert set(metrics.APP_LOOP_COUNT._labelnames) == {"pod_id"}
def test_app_loop_duration():
"""Verify the definition of APP_LOOP_DURATION."""
assert metrics.APP_LOOP_DURATION is not None
assert isinstance(metrics.APP_LOOP_DURATION, Histogram)
assert (
metrics.APP_LOOP_DURATION._name == "app_main_loop_duration_seconds"
) # Histograms don't have _total
assert set(metrics.APP_LOOP_DURATION._labelnames) == {"pod_id"}
def test_app_errors_total():
"""Verify the definition of APP_ERRORS_TOTAL."""
assert metrics.APP_ERRORS_TOTAL is not None
assert isinstance(metrics.APP_ERRORS_TOTAL, Counter)
assert metrics.APP_ERRORS_TOTAL._name == "app_errors" # REMOVED _total
assert set(metrics.APP_ERRORS_TOTAL._labelnames) == {"pod_id"}
def test_app_up():
"""Verify the definition of APP_UP."""
assert metrics.APP_UP is not None
assert isinstance(metrics.APP_UP, Gauge)
assert metrics.APP_UP._name == "app_up" # Gauges don't have _total
assert set(metrics.APP_UP._labelnames) == {"pod_id"}
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_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_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():
"""Verify the definition of SLOTS_MANAGED."""
assert metrics.SLOTS_MANAGED is not None
assert isinstance(metrics.SLOTS_MANAGED, Gauge)
assert metrics.SLOTS_MANAGED._name == "ingestor_slots_managed_current"
assert set(metrics.SLOTS_MANAGED._labelnames) == {"pod_id"}
def test_slots_acquired():
"""Verify the definition of SLOTS_ACQUIRED."""
assert metrics.SLOTS_ACQUIRED is not None
assert isinstance(metrics.SLOTS_ACQUIRED, Counter)
assert metrics.SLOTS_ACQUIRED._name == "ingestor_slots_acquired" # REMOVED _total
assert set(metrics.SLOTS_ACQUIRED._labelnames) == {"pod_id"}
def test_slots_released():
"""Verify the definition of SLOTS_RELEASED."""
assert metrics.SLOTS_RELEASED is not None
assert isinstance(metrics.SLOTS_RELEASED, Counter)
assert metrics.SLOTS_RELEASED._name == "ingestor_slots_released" # REMOVED _total
assert set(metrics.SLOTS_RELEASED._labelnames) == {"pod_id"}
def test_opc_managers_active():
"""Verify the definition of OPC_MANAGERS_ACTIVE."""
assert metrics.OPC_MANAGERS_ACTIVE is not None
assert isinstance(metrics.OPC_MANAGERS_ACTIVE, Gauge)
assert metrics.OPC_MANAGERS_ACTIVE._name == "ingestor_opc_managers_active"
assert set(metrics.OPC_MANAGERS_ACTIVE._labelnames) == {"pod_id"}
def test_opc_subscription_errors():
"""Verify the definition of OPC_SUBSCRIPTION_ERRORS."""
assert metrics.OPC_SUBSCRIPTION_ERRORS is not None
assert isinstance(metrics.OPC_SUBSCRIPTION_ERRORS, Counter)
assert (
metrics.OPC_SUBSCRIPTION_ERRORS._name == "ingestor_opc_subscription_errors"
) # REMOVED _total
assert set(metrics.OPC_SUBSCRIPTION_ERRORS._labelnames) == {
"pod_id",
"server",
"slot",
}
def test_opc_connections_total():
"""Verify the definition of OPC_CONNECTIONS_TOTAL."""
assert metrics.OPC_CONNECTIONS_TOTAL is not None
assert isinstance(metrics.OPC_CONNECTIONS_TOTAL, Counter)
assert (
metrics.OPC_CONNECTIONS_TOTAL._name == "opc_connections_initiated"
) # REMOVED _total
assert set(metrics.OPC_CONNECTIONS_TOTAL._labelnames) == {"pod_id", "server_name"}
def test_opc_connections_failed():
"""Verify the definition of OPC_CONNECTIONS_FAILED."""
assert metrics.OPC_CONNECTIONS_FAILED is not None
assert isinstance(metrics.OPC_CONNECTIONS_FAILED, Counter)
assert (
metrics.OPC_CONNECTIONS_FAILED._name == "opc_connections_failed"
) # REMOVED _total
assert set(metrics.OPC_CONNECTIONS_FAILED._labelnames) == {"pod_id", "server_name"}
def test_opc_connection_status():
"""Verify the definition of OPC_CONNECTION_STATUS."""
assert metrics.OPC_CONNECTION_STATUS is not None
assert isinstance(metrics.OPC_CONNECTION_STATUS, Gauge)
assert metrics.OPC_CONNECTION_STATUS._name == "opc_connection_status"
assert set(metrics.OPC_CONNECTION_STATUS._labelnames) == {
"pod_id",
"server_name",
"server_url",
}
def test_opc_subscriptions_created():
"""Verify the definition of OPC_SUBSCRIPTIONS_CREATED."""
assert metrics.OPC_SUBSCRIPTIONS_CREATED is not None
assert isinstance(metrics.OPC_SUBSCRIPTIONS_CREATED, Counter)
assert (
metrics.OPC_SUBSCRIPTIONS_CREATED._name == "opc_subscriptions_created"
) # REMOVED _total
assert set(metrics.OPC_SUBSCRIPTIONS_CREATED._labelnames) == {
"pod_id",
"server_name",
"slot_name",
}
def test_opc_tags_subscribed():
"""Verify the definition of OPC_TAGS_SUBSCRIBED."""
assert metrics.OPC_TAGS_SUBSCRIBED is not None
assert isinstance(metrics.OPC_TAGS_SUBSCRIBED, Gauge)
assert metrics.OPC_TAGS_SUBSCRIBED._name == "opc_tags_subscribed_current"
assert set(metrics.OPC_TAGS_SUBSCRIBED._labelnames) == {"pod_id", "server_name"}
def test_opc_cycles_without_data():
"""Verify the definition of OPC_CYCLES_WITHOUT_DATA."""
assert metrics.OPC_CYCLES_WITHOUT_DATA is not None
assert isinstance(metrics.OPC_CYCLES_WITHOUT_DATA, Gauge)
assert metrics.OPC_CYCLES_WITHOUT_DATA._name == "opc_cycles_without_data"
assert set(metrics.OPC_CYCLES_WITHOUT_DATA._labelnames) == {"pod_id", "server_name"}
def test_opc_reconnections_total():
"""Verify the definition of OPC_RECONNECTIONS_TOTAL."""
assert metrics.OPC_RECONNECTIONS_TOTAL is not None
assert isinstance(metrics.OPC_RECONNECTIONS_TOTAL, Counter)
assert (
metrics.OPC_RECONNECTIONS_TOTAL._name == "opc_reconnections_tried"
) # REMOVED _total
assert set(metrics.OPC_RECONNECTIONS_TOTAL._labelnames) == {"pod_id", "server_name"}
def test_kafka_messages_sent():
"""Verify the definition of KAFKA_MESSAGES_SENT."""
assert metrics.KAFKA_MESSAGES_SENT is not None
assert isinstance(metrics.KAFKA_MESSAGES_SENT, Counter)
assert metrics.KAFKA_MESSAGES_SENT._name == "kafka_messages_sent" # REMOVED _total
assert set(metrics.KAFKA_MESSAGES_SENT._labelnames) == {"pod_id", "topic"}
def test_kafka_messages_errors():
"""Verify the definition of KAFKA_MESSAGES_ERRORS."""
assert metrics.KAFKA_MESSAGES_ERRORS is not None
assert isinstance(metrics.KAFKA_MESSAGES_ERRORS, Counter)
assert (
metrics.KAFKA_MESSAGES_ERRORS._name == "kafka_messages_errors"
) # REMOVED _total
assert set(metrics.KAFKA_MESSAGES_ERRORS._labelnames) == {"pod_id", "topic"}
def test_kafka_connection_status():
"""Verify the definition of KAFKA_CONNECTION_STATUS."""
assert metrics.KAFKA_CONNECTION_STATUS is not None
assert isinstance(metrics.KAFKA_CONNECTION_STATUS, Gauge)
assert metrics.KAFKA_CONNECTION_STATUS._name == "kafka_connection_status"
assert set(metrics.KAFKA_CONNECTION_STATUS._labelnames) == {"pod_id"}
def test_redis_operations_total():
"""Verify the definition of REDIS_OPERATIONS_TOTAL."""
assert metrics.REDIS_OPERATIONS_TOTAL is not None
assert isinstance(metrics.REDIS_OPERATIONS_TOTAL, Counter)
assert metrics.REDIS_OPERATIONS_TOTAL._name == "redis_operations" # REMOVED _total
assert set(metrics.REDIS_OPERATIONS_TOTAL._labelnames) == {"pod_id", "operation"}
def test_redis_operations_errors():
"""Verify the definition of REDIS_OPERATIONS_ERRORS."""
assert metrics.REDIS_OPERATIONS_ERRORS is not None
assert isinstance(metrics.REDIS_OPERATIONS_ERRORS, Counter)
assert (
metrics.REDIS_OPERATIONS_ERRORS._name == "redis_operations_errors"
) # REMOVED _total
assert set(metrics.REDIS_OPERATIONS_ERRORS._labelnames) == {"pod_id", "operation"}
def test_redis_operations_duration():
"""Verify the definition of REDIS_OPERATIONS_DURATION."""
assert metrics.REDIS_OPERATIONS_DURATION is not None
assert isinstance(metrics.REDIS_OPERATIONS_DURATION, Histogram)
assert (
metrics.REDIS_OPERATIONS_DURATION._name == "redis_operations_duration_seconds"
)
assert set(metrics.REDIS_OPERATIONS_DURATION._labelnames) == {"pod_id", "operation"}
def test_redis_connection_status():
"""Verify the definition of REDIS_CONNECTION_STATUS."""
assert metrics.REDIS_CONNECTION_STATUS is not None
assert isinstance(metrics.REDIS_CONNECTION_STATUS, Gauge)
assert metrics.REDIS_CONNECTION_STATUS._name == "redis_connection_status"
assert set(metrics.REDIS_CONNECTION_STATUS._labelnames) == {"pod_id"}
def test_notifications_sent():
"""Verify the definition of NOTIFICATIONS_SENT."""
assert metrics.NOTIFICATIONS_SENT is not None
assert isinstance(metrics.NOTIFICATIONS_SENT, Counter)
assert metrics.NOTIFICATIONS_SENT._name == "notifications_sent" # REMOVED _total
assert set(metrics.NOTIFICATIONS_SENT._labelnames) == {"pod_id", "level", "block"}