Resolve "[UBI] INT Collector invocation and processing crashes after integrating gNMI"

Closes issue #356 (closed)

Edited by Georgios P. Katsikas

Merge request reports

Loading
+1 −1
Changes for manifests/metallb.yaml: 1 added line, 1 removed line.
Original line number Diff line number Diff line
@@ -17,7 +17,7 @@
apiVersion: metallb.io/v1beta1
kind: IPAddressPool
metadata:
  name: my-ip-pool
  name: metallb-address-pool
  namespace: metallb-system
spec:
  addresses:
+1 −1

File changed.

Contains only whitespace changes.

+162 −69
Changes for src/telemetry/backend/service/collectors/int_collector/INTCollector.py: 162 added lines, 69 removed lines.
Original line number Diff line number Diff line
@@ -27,17 +27,16 @@ import struct
import socket
import ipaddress

from .INTCollectorCommon import IntDropReport, IntLocalReport, IntFixedReport, FlowInfo
from .INTCollectorCommon import IntDropReport, IntLocalReport, IntFixedReport, FlowInfo, IPPacket, UDPPacket
from common.proto.kpi_manager_pb2 import KpiId, KpiDescriptor
from confluent_kafka import Producer as KafkaProducer
from common.tools.kafka.Variables import KafkaConfig, KafkaTopic
from uuid import uuid4
from typing import Dict
from typing import Dict, List, Tuple
from datetime import datetime, timezone
import json

from kpi_manager.client.KpiManagerClient import KpiManagerClient
from common.proto.analytics_frontend_pb2 import Analyzer, AnalyzerId
from context.client.ContextClient import ContextClient
from analytics.frontend.client.AnalyticsFrontendClient import AnalyticsFrontendClient
from common.proto.kpi_sample_types_pb2 import KpiSampleType
@@ -46,20 +45,37 @@ import logging

LOGGER = logging.getLogger(__name__)

class INTCollector(_Collector):
DEF_SW_NUM = 10

    last_packet_time = time.time() # Track last packet time
class INTCollector(_Collector):

    max_idle_time = 5  # for how long we tolerate inactivity
    sniff_timeout = 3   # how often we stop sniffing to check for inactivity
    last_packet_time = 0.0 # Track the timestamp of the last packet
    max_idle_time = 5      # For how long we tolerate inactivity
    sniff_timeout = 3      # How often we stop sniffing to check for inactivity

    """
    INTCollector is a class that simulates a network collector for testing purposes.
    It provides functionalities to manage configurations, state subscriptions, and synthetic data generation.
    INTCollector spawns a packet sniffer at the interface of the Telemetry service,
    which is mapped to the interface of the TFS host.
    INT packets arriving there are:
    - picked up by the INT collector
    - parsed, and
    - telemetry KPI metric values are extracted and reported to Kafka as KPIDescriptors
    """
    def __init__(self, collector_id: str , address: str, interface: str, port: str, kpi_id: str, service_id: str, context_id: str, **settings):
        super().__init__('int_collector', address, port, **settings)
        self._out_samples    = queue.Queue()                # Queue to hold synthetic state samples

    def __init__(self, address: str, port: int, **settings) -> None:
        super().__init__('INTCollector', address, port, **settings)
        self.collector_id = settings.pop('collector_id', None)
        self.interface = settings.pop('interface', None)
        self.kpi_id = settings.pop('kpi_id', None)
        self.service_id = settings.pop('service_id', None)
        self.context_id = settings.pop('context_id', None)

        if any(item is None for item in [
                self.collector_id, self.interface, self.kpi_id, self.service_id, self.context_id]):
            LOGGER.error("INT collector not instantiated properly: Bad input")
            return

        self._out_samples = queue.Queue()
        self._scheduler   = BackgroundScheduler(daemon=True)
        self._scheduler.configure(
            jobstores = {'default': MemoryJobStore()},
@@ -67,29 +83,30 @@ class INTCollector(_Collector):
            timezone  = pytz.utc
        )
        self.kafka_producer = KafkaProducer({'bootstrap.servers': KafkaConfig.get_kafka_address()})
        self.collector_id    = collector_id
        self.interface    = interface
        self.kpi_manager_client = KpiManagerClient()
        self.analytics_frontend_client = AnalyticsFrontendClient()
        self.context_client = ContextClient()
        self.kpi_id     = kpi_id
        self.service_id = service_id
        self.context_id = context_id
        self.table = {}
        self.connected = False          # To track connection state
        LOGGER.info("INT Collector initialized")
        self.connected = False
        LOGGER.info("=== INT Collector initialized")

    def Connect(self) -> bool:
        LOGGER.info("=== INT Collector Connect()")
        LOGGER.info(f"Connecting to {self.interface}:{self.port}")
        self.connected = True

        self._scheduler.add_job(self.sniff_with_restarts_on_idle, id=self.kpi_id ,args=[self.interface , self.port , self.service_id, self.context_id])
        self._scheduler.add_job(
            self.sniff_with_restarts_on_idle,
            id=self.kpi_id,
            args=[self.interface, self.port, self.service_id]
        )

        self._scheduler.start()
        self.connected = True
        LOGGER.info(f"Successfully connected to {self.interface}:{self.port}")
        return True

    def Disconnect(self) -> bool:
        LOGGER.info("=== INT Collector Disconnect()")
        LOGGER.info(f"Disconnecting from {self.interface}:{self.port}")
        if not self.connected:
            LOGGER.warning("INT Collector is not connected. Nothing to disconnect.")
@@ -102,12 +119,52 @@ class INTCollector(_Collector):
        LOGGER.info(f"Successfully disconnected from {self.interface}:{self.port}")
        return True

    def require_connection(self):
        if not self.connected:
            raise RuntimeError("INT collector is not connected. Please connect before performing operations.")

    def SubscribeState(self, subscriptions: List[Tuple[str, dict, float, float, str, int, str, str]]) -> bool:
        LOGGER.info("=== INT Collector SubscribeState()")
        self.require_connection()
        try:
            _, _, _, _, interface, port, service_id, _ = subscriptions
        except:
            LOGGER.exception(f"Invalid subscription format: {subscriptions}")
            return False

        if self.kpi_id:
            self._scheduler.add_job(
                self.sniff_with_restarts_on_idle,
                id=self.kpi_id,
                args=[interface, port, service_id]
            )

        return True

    def UnsubscribeState(self, resource_key: str) -> bool:
        LOGGER.info("=== INT Collector UnsubscribeState()")
        self.require_connection()
        try: 
            # Check if job exists
            job_ids = [job.id for job in self._scheduler.get_jobs() if resource_key in job.id]
            if not job_ids:
                LOGGER.warning(f"No active jobs found for {resource_key}. It might have already been terminated.")
                return False
            for job_id in job_ids:
                self._scheduler.remove_job(job_id)
            LOGGER.info(f"Unsubscribed from {resource_key} with job IDs: {job_ids}")
            return True
        except:
            LOGGER.exception(f"Failed to unsubscribe from {resource_key}")
            return False

    def on_idle_timeout(self):
        LOGGER.info(f"Sniffer idle for more than {self.max_idle_time} seconds.")
        LOGGER.info(f"=== INT Collector IDLE() - No INT packets arrived during the last {self.max_idle_time}")
        LOGGER.debug(f"last_packet_time {self.last_packet_time} seconds.")

        # Report a zero value for the P4 switch KPIs
        values = [0]
        for sw_id in range(1, 6):
        for sw_id in range(1, DEF_SW_NUM+1):
            sw = self.table.get(sw_id)
            self.overwrite_switch_values(sw, values)

@@ -122,38 +179,52 @@ class INTCollector(_Collector):
        for key, value in switch.items():
            self.send_message_to_kafka(key, value)

    def process_packet(self , packet, port, service_id , context_id):
        # global last_packet_time
    def process_packet(self, packet, port, service_id):
        LOGGER.debug("=== INT Collector Packet-In()")
        LOGGER.debug(packet)

        # Check for IP layer
        # Check for IP header
        if IP not in packet:
            return None
        ip_layer = packet[IP]
        # ip_pkt = IPPacket(ip_layer[:20])

        # IP parsing
        try:
            ihl = ip_layer.ihl * 4
            raw_ip = bytes(ip_layer)
            ip_pkt = IPPacket(raw_ip[:ihl]) # exclude options if any
            src_ip_str = str(ipaddress.IPv4Address(ip_pkt.ip_src))
            dst_ip_str = str(ipaddress.IPv4Address(ip_pkt.ip_dst))
            # ip_pkt.show()
        except Exception as ex:
            LOGGER.exception(f"Failed to parse IP packet: {ex}")
            return None

        # Check for UDP
        LOGGER.debug(f"ip src: {src_ip_str}")
        LOGGER.debug(f"ip dst: {dst_ip_str}")
        LOGGER.debug(f"ip-proto: {ip_pkt.ip_proto}")

        # Check for UDP header
        if UDP not in ip_layer:
            return None
        udp_layer = ip_layer[UDP]

        # Only the INT port
        # We care about datagrams arriving on the INT port
        if udp_layer.dport != port:
            LOGGER.warning(f"Expected UDP INT packet on port {udp_layer.dport}. Received packet on port {port}")
            return None
        # udp_dgram = UDPPacket(bytes(udp_layer))
        # udp_dgram.show()

        src_ip = socket.ntohl(struct.unpack('<I', socket.inet_aton(ip_layer.src))[0])
        src_ip_str = str(ipaddress.IPv4Address(src_ip))
        LOGGER.debug("ip src: {}".format(src_ip_str))

        dst_ip = socket.ntohl(struct.unpack('<I', socket.inet_aton(ip_layer.dst))[0])
        dst_ip_str = str(ipaddress.IPv4Address(dst_ip))
        LOGGER.debug("ip dst: {}".format(dst_ip_str))
        LOGGER.debug("ip-proto: {}".format(ip_layer.proto))
        try:
            raw_udp = bytes(udp_layer)
            udp_dgram = UDPPacket(raw_udp)
            # udp_dgram.show()
        except Exception as ex:
            LOGGER.exception(f"Failed to parse UDP datagram: {ex}")
            return None

        LOGGER.debug("port src: {}".format(udp_layer.sport))
        LOGGER.debug("port dst: {}".format(udp_layer.dport))
        # UDP parsing
        LOGGER.debug(f"port src: {udp_dgram.udp_port_src}")
        LOGGER.debug(f"port dst: {udp_dgram.udp_port_dst}")

        # Get the INT report data (after UDP header)
        int_data = bytes(udp_layer.payload)
@@ -167,21 +238,23 @@ class INTCollector(_Collector):
        local_report = None
        lat = 0

        # Drop report
        if fixed_report.d == 1:
            drop_report = IntDropReport(int_data[offset:offset + 4])
            offset += 4
            # drop_report.show()
        # Regular report
        elif fixed_report.f == 1 or fixed_report.q == 1:
            local_report = IntLocalReport(int_data[offset:offset + 8])
            offset += 8
            lat = local_report.egress_timestamp - fixed_report.ingress_timestamp
            assert lat > 0, "Egress timestamp must be > ingress timestamp"
            assert lat > 0, f"Egress timestamp must be > ingress timestamp. Got a diff: {lat}"
            # local_report.show()

        # Create flow info
        flow_info = FlowInfo(
            src_ip=src_ip,
            dst_ip=dst_ip,
            src_ip=ip_pkt.ip_src,
            dst_ip=ip_pkt.ip_dst,
            src_port=udp_layer.sport,
            dst_port=udp_layer.dport,
            ip_proto=ip_layer.proto,
@@ -201,25 +274,26 @@ class INTCollector(_Collector):
        )
        LOGGER.debug(f"Flow info: {flow_info}")

        self.create_descriptors_and_send_to_kafka(flow_info , service_id , context_id)

        self.create_descriptors_and_send_to_kafka(flow_info, service_id)
        self.last_packet_time = time.time()

        return flow_info

    def set_kpi_descriptor(self , kpi_uuid , service_id , device_id , endpoint_id , sample_type):
    def set_kpi_descriptor(self, kpi_uuid, service_id, sample_type):
        kpi_descriptor = KpiDescriptor()
        kpi_descriptor.kpi_sample_type = sample_type
        kpi_descriptor.service_id.service_uuid.uuid = service_id
        # kpi_descriptor.device_id.device_uuid.uuid = device_id
        # kpi_descriptor.endpoint_id.endpoint_uuid.uuid = endpoint_id
        kpi_descriptor.kpi_id.kpi_id.uuid = kpi_uuid

        try:
            kpi_id: KpiId = self.kpi_manager_client.SetKpiDescriptor(kpi_descriptor)
        except Exception as ex:
            LOGGER.exception(f"Failed to set KPI descriptor {kpi_uuid}: {ex}")

        return kpi_id

    def create_descriptors_and_send_to_kafka(self, flow_info , service_id , context_id):
        LOGGER.debug(f"PACKET FROM SWITCH: {flow_info.switch_id} LATENCY: {flow_info.hop_latency}")
    def create_descriptors_and_send_to_kafka(self, flow_info, service_id):
        LOGGER.debug(f"Packet from switch: {flow_info.switch_id} with latency: {flow_info.hop_latency}")
        if(self.table.get(flow_info.switch_id) == None):
            seq_num_kpi_id     = str(uuid4())
            ingress_ts_kpi_id  = str(uuid4())
@@ -241,14 +315,14 @@ class INTCollector(_Collector):
            LOGGER.debug(f"is_drop_kpi_id     for switch {flow_info.switch_id}: {is_drop_kpi_id}")
            LOGGER.debug(f"sw_lat_kpi_id      for switch {flow_info.switch_id}: {sw_lat_kpi_id}")

            seq_num_kpi           = self.set_kpi_descriptor(seq_num_kpi_id,     service_id ,'', '', KpiSampleType.KPISAMPLETYPE_INT_SEQ_NUM)
            ingress_timestamp_kpi = self.set_kpi_descriptor(ingress_ts_kpi_id,  service_id, '', '', KpiSampleType.KPISAMPLETYPE_INT_TS_ING)
            egress_timestamp_kpi  = self.set_kpi_descriptor(egress_ts_kpi_id,   service_id, '', '', KpiSampleType.KPISAMPLETYPE_INT_TS_EGR)
            hop_latency_kpi       = self.set_kpi_descriptor(hop_lat_kpi_id,     service_id, '', '', KpiSampleType.KPISAMPLETYPE_INT_HOP_LAT)
            ingress_port_id_kpi   = self.set_kpi_descriptor(ing_port_id_kpi_id, service_id, '', '', KpiSampleType.KPISAMPLETYPE_INT_PORT_ID_ING)
            egress_port_id_kpi    = self.set_kpi_descriptor(egr_port_id_kpi_id, service_id, '', '', KpiSampleType.KPISAMPLETYPE_INT_PORT_ID_EGR)
            queue_occup_kpi       = self.set_kpi_descriptor(queue_occup_kpi_id, service_id, '', '', KpiSampleType.KPISAMPLETYPE_INT_QUEUE_OCCUP)
            is_drop_kpi           = self.set_kpi_descriptor(is_drop_kpi_id,     service_id, '', '', KpiSampleType.KPISAMPLETYPE_INT_IS_DROP)
            seq_num_kpi           = self.set_kpi_descriptor(seq_num_kpi_id,     service_id ,KpiSampleType.KPISAMPLETYPE_INT_SEQ_NUM)
            ingress_timestamp_kpi = self.set_kpi_descriptor(ingress_ts_kpi_id,  service_id, KpiSampleType.KPISAMPLETYPE_INT_TS_ING)
            egress_timestamp_kpi  = self.set_kpi_descriptor(egress_ts_kpi_id,   service_id, KpiSampleType.KPISAMPLETYPE_INT_TS_EGR)
            hop_latency_kpi       = self.set_kpi_descriptor(hop_lat_kpi_id,     service_id, KpiSampleType.KPISAMPLETYPE_INT_HOP_LAT)
            ingress_port_id_kpi   = self.set_kpi_descriptor(ing_port_id_kpi_id, service_id, KpiSampleType.KPISAMPLETYPE_INT_PORT_ID_ING)
            egress_port_id_kpi    = self.set_kpi_descriptor(egr_port_id_kpi_id, service_id, KpiSampleType.KPISAMPLETYPE_INT_PORT_ID_EGR)
            queue_occup_kpi       = self.set_kpi_descriptor(queue_occup_kpi_id, service_id, KpiSampleType.KPISAMPLETYPE_INT_QUEUE_OCCUP)
            is_drop_kpi           = self.set_kpi_descriptor(is_drop_kpi_id,     service_id, KpiSampleType.KPISAMPLETYPE_INT_IS_DROP)

            # Set a dedicated KPI descriptor for every switch
            sw_lat_kpi = None
@@ -259,10 +333,10 @@ class INTCollector(_Collector):
                KpiSampleType.KPISAMPLETYPE_INT_HOP_LAT_SW07, KpiSampleType.KPISAMPLETYPE_INT_HOP_LAT_SW08,
                KpiSampleType.KPISAMPLETYPE_INT_HOP_LAT_SW09, KpiSampleType.KPISAMPLETYPE_INT_HOP_LAT_SW10
            ]
            for i, sw_id in enumerate(range(1, 11)):
            for i, sw_id in enumerate(range(1, DEF_SW_NUM+1)):
                if flow_info.switch_id == sw_id:
                    LOGGER.debug(f"SET KPI : seq_num_kpi_id for switch {flow_info.switch_id}: {sw_lat_kpi_id}")
                    sw_lat_kpi = self.set_kpi_descriptor(sw_lat_kpi_id, service_id, '', '', sw_sample_types[i])
                    LOGGER.debug(f"Set latency KPI for switch {flow_info.switch_id}: {sw_lat_kpi_id}")
                    sw_lat_kpi = self.set_kpi_descriptor(sw_lat_kpi_id, service_id, sw_sample_types[i])

            # Gather keys and values
            keys   = [
@@ -313,12 +387,14 @@ class INTCollector(_Collector):
            self.overwrite_switch_values(switch, values)

    def send_message_to_kafka(self, kpi_id, measured_kpi_value):
        LOGGER.debug("=== INT Collector Kafka Writer()")
        producer = self.kafka_producer
        kpi_value: Dict = {
            "time_stamp": datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"),
            "kpi_id": kpi_id,
            "kpi_value": measured_kpi_value
        }
        try:
            producer.produce(
                KafkaTopic.VALUE.value,
                key=self.collector_id,
@@ -326,32 +402,49 @@ class INTCollector(_Collector):
                callback=self.delivery_callback
            )
            producer.flush()
        LOGGER.debug(f"Message with kpi_id: {kpi_id} was send to kafka!")
        except Exception as ex:
            LOGGER.error(f"Message with kpi_id: {kpi_id} is NOT sent to kafka!")
            LOGGER.exception(f"{ex}")
            return
        LOGGER.debug(f"Message with kpi_id: {kpi_id} is sent to kafka!")

    def packet_callback(self, packet, port , service_id,context_id):
        flow_info = self.process_packet(packet , port , service_id, context_id)
    def packet_callback(self, packet, port, service_id):
        flow_info = self.process_packet(packet, port, service_id)
        if flow_info:
            LOGGER.debug(f"Flow info: {flow_info}")

    def sniff_with_restarts_on_idle(self, interface, port, service_id , context_id):
    def sniff_with_restarts_on_idle(self, interface, port, service_id):
        LOGGER.info("=== INT Collector Sniffer Start")
        while True:
            # Run sniff for a short period to periodically check for idle timeout
            sniff(
            try:
                sniff( # type: ignore
                    iface=interface,
                    filter=f"udp port {port}",
                prn=lambda pkt: self.packet_callback(pkt, port, service_id , context_id),
                    prn=lambda pkt: self.packet_callback(pkt, port, service_id),
                    timeout=self.sniff_timeout
                )
            except Exception as ex:
                LOGGER.exception(ex)
                self.Disconnect()

            if not self.connected:
                break

            # Check if idle period has been exceeded
            now = time.time()
            LOGGER.debug(f"Time now: {self.epoch_to_day_time(now)}")
            LOGGER.debug(f"Time last pkt: {self.epoch_to_day_time(self.last_packet_time)}")
            diff = now - self.last_packet_time
            assert diff > 0, f"Time diff: {diff} sec must be positive"
            if (now - self.last_packet_time) > self.max_idle_time:
                self.on_idle_timeout()
                self.last_packet_time = now  # Reset timer after action
        LOGGER.info("=== INT Collector Sniffer End")

    def delivery_callback(self, err, msg):
        if err:
            LOGGER.error('Message delivery failed: {:s}'.format(str(err)))
            LOGGER.error(f"Kafka message delivery failed: {str(err)}")

    def epoch_to_day_time(self, ep_time : float):
        return time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(ep_time))
+1 −32
Changes for src/telemetry/backend/service/collectors/int_collector/INTCollectorCommon.py: 1 added line, 32 removed lines.
Original line number Diff line number Diff line
@@ -35,7 +35,7 @@ class IPPacket(Packet):
    ]

class UDPPacket(Packet):
    name = "UDPPacket"
    name = "UDPDatagram"
    fields_desc = [
        BitField("udp_port_src", 0, 16),
        BitField("udp_port_dst", 0, 16),
@@ -43,37 +43,6 @@ class UDPPacket(Packet):
        BitField("udp_csum", 0, 16)
    ]

"""
    Private source repo: https://github.com/SNVLab-WUT/int-collector/blob/main/collector/BPFCollector_v0_5.c
    fabric-tna repo: https://github.com/stratum/fabric-tna/blob/main/p4src/shared/header.p4

    Little Endian order
        BitField("nproto", 0, 4),
        BitField("ver", 0, 4),
        BitField("rsvd1", 0, 5),
        BitField("f", 0, 1),
        BitField("q", 0, 1),
        BitField("d", 0, 1),
        BitField("hw_id", 0, 6),
        BitField("rsvd2", 0, 10),

        BitField("queue_id", 0, 8),
        BitField("queue_occupancy", 0, 24),

    Big Endian order
        BitField("ver", 0, 4),
        BitField("nproto", 0, 4),
        BitField("d", 0, 1),
        BitField("q", 0, 1),
        BitField("f", 0, 1),
        BitField("rsvd1", 0, 5),
        BitField("rsvd2", 0, 10),
        BitField("hw_id", 0, 6),

        BitField("queue_occupancy", 0, 24),
        BitField("queue_id", 0, 8),
"""

class IntFixedReport(Packet):
    name = "IntFixedReport"
    fields_desc = [
+0 −0

File moved.

Loading
Loading