Loading manifests/metallb.yaml +1 −1 Viewed Changes for manifests/metallb.yaml: 1 added line, 1 removed line. Original line number Diff line number Diff line Loading @@ -17,7 +17,7 @@ apiVersion: metallb.io/v1beta1 kind: IPAddressPool metadata: name: my-ip-pool name: metallb-address-pool namespace: metallb-system spec: addresses: Loading src/telemetry/backend/service/collectors/gnmi_oc/GnmiOpenConfigCollector.py +1 −1 Viewed File changed.Contains only whitespace changes. Show changes src/telemetry/backend/service/collectors/intcollector/INTCollector.py→src/telemetry/backend/service/collectors/int_collector/INTCollector.py +162 −69 Viewed 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 Loading @@ -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 Loading @@ -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()}, Loading @@ -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.") Loading @@ -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) Loading @@ -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) Loading @@ -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, Loading @@ -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()) Loading @@ -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 Loading @@ -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 = [ Loading Loading @@ -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, Loading @@ -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)) src/telemetry/backend/service/collectors/intcollector/INTCollectorCommon.py→src/telemetry/backend/service/collectors/int_collector/INTCollectorCommon.py +1 −32 Viewed 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 Loading @@ -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), Loading @@ -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 = [ Loading src/telemetry/backend/service/collectors/intcollector/__init__.py→src/telemetry/backend/service/collectors/int_collector/__init__.py +0 −0 Viewed File moved. View file Loading
manifests/metallb.yaml +1 −1 Viewed Changes for manifests/metallb.yaml: 1 added line, 1 removed line. Original line number Diff line number Diff line Loading @@ -17,7 +17,7 @@ apiVersion: metallb.io/v1beta1 kind: IPAddressPool metadata: name: my-ip-pool name: metallb-address-pool namespace: metallb-system spec: addresses: Loading
src/telemetry/backend/service/collectors/gnmi_oc/GnmiOpenConfigCollector.py +1 −1 Viewed File changed.Contains only whitespace changes. Show changes
src/telemetry/backend/service/collectors/intcollector/INTCollector.py→src/telemetry/backend/service/collectors/int_collector/INTCollector.py +162 −69 Viewed 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 Loading @@ -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 Loading @@ -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()}, Loading @@ -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.") Loading @@ -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) Loading @@ -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) Loading @@ -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, Loading @@ -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()) Loading @@ -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 Loading @@ -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 = [ Loading Loading @@ -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, Loading @@ -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))
src/telemetry/backend/service/collectors/intcollector/INTCollectorCommon.py→src/telemetry/backend/service/collectors/int_collector/INTCollectorCommon.py +1 −32 Viewed 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 Loading @@ -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), Loading @@ -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 = [ Loading
src/telemetry/backend/service/collectors/intcollector/__init__.py→src/telemetry/backend/service/collectors/int_collector/__init__.py +0 −0 Viewed File moved. View file