Update MongoDB configuration in values.yaml and enhance Ingestor class to utilize MongoDB connection parameters

This commit is contained in:
vitor-aignosi
2025-07-01 12:31:46 -03:00
parent d79f882018
commit c5b23da373
3 changed files with 34 additions and 10 deletions

View File

@@ -14,11 +14,12 @@ class IngestorManager():
def __init__(self,
kafka_servers: str, redis_host: str, redis_port: int,
lease_ttl: int, heartbeat_ttl: int, pod_id: str,
poll_interval: int, logger: Logger, notification_handler: NotificationHandler,
poll_interval: int, mongo_connection_string: str, mongo_database: str,
logger: Logger, notification_handler: NotificationHandler,
redis_username: str = None, redis_password: str = None):
self.data_manager = DataManager(
kafka_servers, logger, notification_handler)
kafka_servers, mongo_connection_string, mongo_database, logger, notification_handler)
self.opc_managers = {}
self.resource_manager = ResourceManager(
redis_host, redis_port, lease_ttl, heartbeat_ttl, pod_id, redis_username, redis_password
@@ -58,7 +59,8 @@ class IngestorManager():
manager = OpcManager(
server_config['name'], server_config['url'],
data_manager, logger, server_config['server_uri'],
self.notification_handler, self.pod_id, 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('server_cert_path')
)
@@ -163,7 +165,8 @@ class IngestorManager():
self.opc_managers[server].disconnect()
self.opc_managers.pop(server, None)
metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers))
metrics.OPC_MANAGERS_ACTIVE.labels(
pod_id=self.pod_id).set(len(self.opc_managers))
def check_opc_servers_integrity(self):
"""
@@ -182,7 +185,8 @@ class IngestorManager():
for slot, _config in self.managed_tags.items():
self.managed_tags[slot].pop(server, None)
metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers))
metrics.OPC_MANAGERS_ACTIVE.labels(
pod_id=self.pod_id).set(len(self.opc_managers))
def declare_active(self):
"""
@@ -266,7 +270,8 @@ class IngestorManager():
if len(acquired) >= max_slots:
self.managed_tags.update(acquired)
metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags))
metrics.SLOTS_MANAGED.labels(
pod_id=self.pod_id).set(len(self.managed_tags))
return acquired
self.logger.warning(
@@ -274,7 +279,8 @@ class IngestorManager():
f"Only {acquired} slots were leased."
)
self.managed_tags.update(acquired)
metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags))
metrics.SLOTS_MANAGED.labels(
pod_id=self.pod_id).set(len(self.managed_tags))
return acquired
def unsubscribe_slot(self, slot: str):
@@ -400,7 +406,8 @@ class IngestorManager():
tags_to_sub
)
except Exception as e:
metrics.OPC_SUBSCRIPTION_ERRORS.labels(pod_id=self.pod_id, server=server, slot=slot).inc()
metrics.OPC_SUBSCRIPTION_ERRORS.labels(
pod_id=self.pod_id, server=server, slot=slot).inc()
trace = traceback.format_exc()
self.notification_handler.build_and_send_notification(
notification_id=f'OPC_SUBSCRIPTION_ERROR_{slot}:{server}',