Commit a77e8488 authored by Waleed Akbar's avatar Waleed Akbar
Browse files

feat(collectors): Enhance NETCONF collector connection handling and improve...

feat(collectors): Enhance NETCONF collector connection handling and improve error logging for driver preloading
parent 75d4ecf2
Loading
Loading
Loading
Loading
+2 −1
Original line number Diff line number Diff line
@@ -17,8 +17,9 @@ APScheduler>=3.10.4
confluent-kafka==2.3.*
deepdiff==6.7.*
kafka-python==2.0.6
ncclient>=0.6.15
ncclient==0.6.15
numpy==2.0.1
paramiko==2.11.*
pygnmi==0.8.14
pytz>=2025.2
scapy==2.6.1    # TODO: UBI need to confirm the version (This depencdency was missing)
+14 −9
Original line number Diff line number Diff line
@@ -57,6 +57,7 @@ class TelemetryBackendService(GenericGrpcService):
        self.context_client        = ContextClient()
        self.kpi_manager_client    = KpiManagerClient()
        self.active_jobs = {}
        self.active_collectors = {}

    def install_servicers(self):
        threading.Thread(target=self.RequestListener).start()
@@ -138,9 +139,10 @@ class TelemetryBackendService(GenericGrpcService):
        Method to handle collector request.
        """
        LOGGER.info(f"Starting Collector Handler for Collector ID: {collector_id} - KPI ID: {kpi_id}")
        device_collector = None
        # INT collector invocation
        if interface:
            self.device_collector = get_node_level_int_collector(
            device_collector = get_node_level_int_collector(
                collector_id=collector_id,
                kpi_id=kpi_id,
                address="127.0.0.1",
@@ -149,6 +151,7 @@ class TelemetryBackendService(GenericGrpcService):
                service_id=service_id,
                context_id=context_id
            )
            self.active_collectors[collector_id] = device_collector
            return
        # Rest of the collectors
        # elif context_id == "43813baf-195e-5da6-af20-b3d0922e71a7":
@@ -165,18 +168,19 @@ class TelemetryBackendService(GenericGrpcService):
            device_type   = collector.get('device_type',   None)  # str: e.g. 'packet-router'
            if device_driver and device_type:
                LOGGER.info(f"Getting collector by meta info for KPI ID: {kpi_id} - Device Type: {device_type} - Device Driver: {device_driver}")
                self.device_collector = get_collector_by_meta_info(
                device_collector = get_collector_by_meta_info(
                    kpi_id, device_type, device_driver,
                    self.kpi_manager_client, self.context_client, self.driver_instance_cache
                )
            else:
                self.device_collector = get_collector_by_kpi_id(
                device_collector = get_collector_by_kpi_id(
                    kpi_id, self.kpi_manager_client, self.context_client, self.driver_instance_cache
                )

        if not self.device_collector:
        if not device_collector:
            LOGGER.warning(f"KPI ID: {kpi_id} - Collector not found. Skipping...")
            raise Exception(f"KPI ID: {kpi_id} - Collector not found.")
        self.active_collectors[collector_id] = device_collector

        resource_to_subscribe = get_subscription_parameters(
            kpi_id, self.kpi_manager_client, self.context_client, duration, interval
@@ -186,7 +190,7 @@ class TelemetryBackendService(GenericGrpcService):
            LOGGER.warning(f"KPI ID: {kpi_id} - Resource to subscribe not found. Skipping...")
            raise Exception(f"KPI ID: {kpi_id} - Resource to subscribe not found.")

        responses = self.device_collector.SubscribeState(resource_to_subscribe)
        responses = device_collector.SubscribeState(resource_to_subscribe)
        for status in responses:
            if isinstance(status, Exception):
                LOGGER.error(f"Subscription failed for KPI ID: {kpi_id} - Error: {status}")
@@ -194,11 +198,11 @@ class TelemetryBackendService(GenericGrpcService):
            else:
                LOGGER.info(f"Subscription successful for KPI ID: {kpi_id} - Status: {status}")
                
        for (timestamp, resource_key, value) in self.device_collector.GetState(duration=duration, blocking=True):
        for (timestamp, resource_key, value) in device_collector.GetState(duration=duration, blocking=True):
            LOGGER.info(f"KPI ID: {kpi_id} - resource={resource_key} value={value}")
            self.GenerateKpiValue(collector_id, kpi_id, value)
            if stop_event.is_set():
                self.device_collector.Disconnect()
                device_collector.Disconnect()
                break

    def GenerateKpiValue(self, collector_id: str, kpi_id: str, measured_kpi_value: Any):
@@ -235,8 +239,9 @@ class TelemetryBackendService(GenericGrpcService):
                if stop_event:
                    stop_event.set()
                    LOGGER.info(f"Job {job_id} terminated.")
                    if self.device_collector is not None:
                        if self.device_collector.UnsubscribeState(job_id):
                    device_collector = self.active_collectors.pop(job_id, None)
                    if device_collector is not None:
                        if device_collector.UnsubscribeState(job_id):
                            LOGGER.info(f"Unsubscribed from collector: {job_id}")
                        else:
                            LOGGER.warning(f"Failed to unsubscribe from collector: {job_id}")
+6 −1
Original line number Diff line number Diff line
@@ -140,4 +140,9 @@ def get_driver(driver_instance_cache : DriverInstanceCache, device : Device) ->
def preload_drivers(driver_instance_cache : DriverInstanceCache) -> None:
    context_client = ContextClient()
    devices = context_client.ListDevices(Empty())
    for device in devices.devices: get_driver(driver_instance_cache, device)
    for device in devices.devices:
        device_uuid = device.device_id.device_uuid.uuid
        try:
            get_driver(driver_instance_cache, device)
        except Exception:
            LOGGER.warning("Skipping driver preload for device(%s)", device_uuid, exc_info=True)
+14 −1
Original line number Diff line number Diff line
@@ -74,6 +74,18 @@ class NetconfOpenConfigCollector(_Collector):

    def Connect(self) -> bool:
        """Open a NETCONF-over-SSH session and start the background scheduler."""
        if self._session is not None and getattr(self._session, 'connected', True):
            if not self._scheduler.running:
                self._scheduler.start()
            LOGGER.info(
                "NetconfOpenConfigCollector already connected to %s:%s",
                self._address,
                self._port,
            )
            return True
        if self._session is not None:
            self._session = None

        try:
            self._session = connect_ssh(
                host           = self._address,
@@ -84,6 +96,7 @@ class NetconfOpenConfigCollector(_Collector):
                allow_agent    = False,
                look_for_keys  = False,
            )
            if not self._scheduler.running:
                self._scheduler.start()
            LOGGER.info("NetconfOpenConfigCollector connected to %s:%s", self._address, self._port)
            return True
+2 −2
Original line number Diff line number Diff line
@@ -32,8 +32,8 @@ LOGGER = logging.getLogger(__name__)

# OpenConfig namespace for platform (components/component)
NS_PLATFORM = 'http://openconfig.net/yang/platform'
# OpenConfig namespace for terminal-device-digital-subcarriers (optical-channel augments platform component)
NS_TERMINAL = 'http://openconfig.net/yang/terminal-device'
# OpenConfig namespace for DSCM optical-channel augments on platform component.
NS_TERMINAL = 'http://openconfig.net/yang/terminal-device-digital-subcarriers'

_NAMESPACES = {
    'ocp' : NS_PLATFORM,