Commit 56f63933 authored by Waleed Akbar's avatar Waleed Akbar
Browse files

Additions to Backend Telemetry:

- Added Emulated Metric Collector.
- Implemented a few other changes.
parent 88b4d7c7
Loading
Loading
Loading
Loading
+2 −5
Original line number Diff line number Diff line
@@ -18,15 +18,12 @@ PROJECTDIR=`pwd`

cd $PROJECTDIR/src
# RCFILE=$PROJECTDIR/coverage/.coveragerc
# coverage run --rcfile=$RCFILE --append -m pytest --log-level=INFO --verbose \
#     kpi_manager/tests/test_unitary.py

# python3 kpi_manager/tests/test_unitary.py
export KFK_SERVER_ADDRESS='127.0.0.1:9092'
CRDB_SQL_ADDRESS=$(kubectl get service cockroachdb-public --namespace crdb -o jsonpath='{.spec.clusterIP}')
export CRDB_URI="cockroachdb://tfs:tfs123@${CRDB_SQL_ADDRESS}:26257/tfs_telemetry?sslmode=require"
RCFILE=$PROJECTDIR/coverage/.coveragerc


python3 -m pytest --log-level=INFO --log-cli-level=debug --verbose \
    telemetry/backend/tests/test_TelemetryBackend.py
python3 -m pytest --log-level=debug --log-cli-level=debug --verbose \
    telemetry/backend/tests/test_backend.py
+75 −0
Original line number Diff line number Diff line
import numpy as np
import random
import threading
import time
import logging
import queue

LOGGER = logging.getLogger(__name__)

class NetworkMetricsEmulator(threading.Thread):
    def __init__(self, interval=1, duration=10, metric_queue=None, network_state="moderate"):
        LOGGER.info("Initiaitng Emulator")
        super().__init__()
        self.interval            = interval
        self.duration            = duration
        self.metric_queue        = metric_queue if metric_queue is not None else queue.Queue()
        self.network_state       = network_state
        self.running             = True
        self.base_utilization    = None
        self.states              = None
        self.state_probabilities = None
        self.set_inital_parameter_values()

    def set_inital_parameter_values(self):
        self.states              = ["good", "moderate", "poor"]
        self.state_probabilities = {
            "good"    : [0.9, 0.1, 0.0],
            "moderate": [0.2, 0.7, 0.1],
            "poor"    : [0.0, 0.3, 0.7]
        }
        if self.network_state     == "good":
            self.base_utilization = random.uniform(700, 900)
        elif self.network_state   == "moderate":
            self.base_utilization = random.uniform(300, 700)
        else:
            self.base_utilization = random.uniform(100, 300)

    def generate_synthetic_data_point(self):
        if self.network_state   == "good":
            variance = random.uniform(-5, 5)  
        elif self.network_state == "moderate":
            variance = random.uniform(-50, 50)
        elif self.network_state == "poor":
            variance = random.uniform(-100, 100)
        else:
            raise ValueError("Invalid network state. Must be 'good', 'moderate', or 'poor'.")
        self.base_utilization += variance

        period       = 60 * 60 * random.uniform(10, 100)
        amplitude    = random.uniform(50, 100) 
        sin_wave     = amplitude * np.sin(2 * np.pi * 100 / period) + self.base_utilization
        random_noise = random.uniform(-10, 10)
        utilization  = sin_wave + random_noise 

        state_prob = self.state_probabilities[self.network_state]
        self.network_state = random.choices(self.states, state_prob)[0]

        return utilization

    def run(self):
        while self.running and (self.duration == -1 or self.duration > 0):
            utilization = self.generate_synthetic_data_point()
            self.metric_queue.put(round(utilization,3))
            time.sleep(self.interval)  
            if self.duration > 0:
                self.duration -= self.interval
                if self.duration == -1:
                    self.duration = 0
        LOGGER.debug("Emulator collector is stopped.")
        self.stop()

    def stop(self):
        self.running = False
        if not self.is_alive():
            LOGGER.debug("Emulator Collector is Termintated.")
+57 −68
Original line number Diff line number Diff line
@@ -12,22 +12,23 @@
# See the License for the specific language governing permissions and
# limitations under the License.

import queue
import json
import time
import random
import logging
import threading
from typing           import Any, Dict
from datetime         import datetime, timezone
# from common.proto.context_pb2 import Empty
from confluent_kafka  import Producer as KafkaProducer
from confluent_kafka  import Consumer as KafkaConsumer
from confluent_kafka  import KafkaError
from numpy import info
from common.Constants import ServiceNameEnum
from common.Settings  import get_service_port_grpc
from common.tools.kafka.Variables import KafkaConfig, KafkaTopic
from common.method_wrappers.Decorator        import MetricsPool
from common.tools.kafka.Variables            import KafkaConfig, KafkaTopic
from common.tools.service.GenericGrpcService import GenericGrpcService
from telemetry.backend.service.EmulatedCollector import NetworkMetricsEmulator

LOGGER             = logging.getLogger(__name__)
METRICS_POOL       = MetricsPool('TelemetryBackend', 'backendService')
@@ -46,6 +47,8 @@ class TelemetryBackendService(GenericGrpcService):
                                            'group.id'           : 'backend',
                                            'auto.offset.reset'  : 'latest'})
        self.running_threads   = {}
        self.emulatorCollector = None
        self.metric_queue      = queue.Queue()

    def install_servicers(self):
        threading.Thread(target=self.RequestListener).start()
@@ -66,93 +69,84 @@ class TelemetryBackendService(GenericGrpcService):
                if receive_msg.error().code() == KafkaError._PARTITION_EOF:
                    continue
                else:
                    # print("Consumer error: {}".format(receive_msg.error()))
                    LOGGER.error("Consumer error: {}".format(receive_msg.error()))
                    break
            try: 
                collector = json.loads(receive_msg.value().decode('utf-8'))
                collector_id = receive_msg.key().decode('utf-8')
                LOGGER.debug('Recevied Collector: {:} - {:}'.format(collector_id, collector))
                # print('Recevied Collector: {:} - {:}'.format(collector_id, collector))

                if collector['duration'] == -1 and collector['interval'] == -1:
                    self.TerminateCollectorBackend(collector_id)
                else:
                    self.RunInitiateCollectorBackend(collector_id, collector)
                    threading.Thread(target=self.InitiateCollectorBackend, 
                                  args=(collector_id, collector)).start()
            except Exception as e:
                LOGGER.warning("Unable to consumer message from topic: {:}. ERROR: {:}".format(KafkaTopic.REQUEST.value, e))
                # print         ("Unable to consumer message from topic: {:}. ERROR: {:}".format(KafkaTopic.REQUEST.value, e))

    def TerminateCollectorBackend(self, collector_id):
        if collector_id in self.running_threads:
            thread, stop_event = self.running_threads[collector_id]
            stop_event.set()
            thread.join()
            # print ("Terminating backend (by StopCollector): Collector Id: ", collector_id)
            del self.running_threads[collector_id]
            self.GenerateCollectorTerminationSignal(collector_id, "-1", -1)          # Termination confirmation to frontend.
        else:
            # print ('Backend collector {:} not found'.format(collector_id))
            LOGGER.warning('Backend collector {:} not found'.format(collector_id))

    def RunInitiateCollectorBackend(self, collector_id: str, collector: str):
        stop_event = threading.Event()
        thread = threading.Thread(target=self.InitiateCollectorBackend, 
                                  args=(collector_id, collector, stop_event))
        self.running_threads[collector_id] = (thread, stop_event)
        thread.start()

    def InitiateCollectorBackend(self, collector_id, collector, stop_event):
    def InitiateCollectorBackend(self, collector_id, collector):
        """
        Method receives collector request and initiates collecter backend.
        """
        # print("Initiating backend for collector: ", collector_id)
        LOGGER.info("Initiating backend for collector: {:s}".format(str(collector_id)))
        start_time = time.time()
        while not stop_event.is_set():
            if int(collector['duration']) != -1 and time.time() - start_time >= collector['duration']:            # condition to terminate backend
                print("Execuation duration completed: Terminating backend: Collector Id: ", collector_id, " - ", time.time() - start_time)
                self.GenerateCollectorTerminationSignal(collector_id, "-1", -1)       # Termination confirmation to frontend.
                break
            self.ExtractKpiValue(collector_id, collector['kpi_id'])
            time.sleep(collector['interval'])
        self.emulatorCollector = NetworkMetricsEmulator(
            duration           = collector['duration'],
            interval           = collector['interval'],
            metric_queue       = self.metric_queue
        )
        self.emulatorCollector.start()
        self.running_threads[collector_id] = self.emulatorCollector

    def GenerateCollectorTerminationSignal(self, collector_id: str, kpi_id: str, measured_kpi_value: Any):
        while self.emulatorCollector.is_alive():
            if not self.metric_queue.empty():
                metric_value = self.metric_queue.get()
                LOGGER.debug("Metric: {:} - Value : {:}".format(collector['kpi_id'], metric_value))
                self.GenerateCollectorResponse(collector_id, collector['kpi_id'] , metric_value)
            time.sleep(1)
        self.TerminateCollectorBackend(collector_id)

    def GenerateCollectorResponse(self, collector_id: str, kpi_id: str, measured_kpi_value: Any):
        """
        Method to write kpi Termination signat on RESPONSE Kafka topic
        Method to write kpi value on RESPONSE Kafka topic
        """
        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,
            "kpi_value" : measured_kpi_value
        }
        producer.produce(
            KafkaTopic.RESPONSE.value, # TODO: to  the topic ...
            KafkaTopic.VALUE.value, # TODO: to  the topic ...
            key      = collector_id,
            value    = json.dumps(kpi_value),
            callback = self.delivery_callback
        )
        producer.flush()

    def ExtractKpiValue(self, collector_id: str, kpi_id: str):
        """
        Method to extract kpi value.
        """
        measured_kpi_value = random.randint(1,100)                      # TODO: To be extracted from a device
        # print ("Measured Kpi value: {:}".format(measured_kpi_value))
        self.GenerateCollectorResponse(collector_id, kpi_id , measured_kpi_value)
    def TerminateCollectorBackend(self, collector_id):
        LOGGER.debug("Terminating collector backend...")
        if collector_id in self.running_threads:
            thread = self.running_threads[collector_id]
            thread.stop()
            del self.running_threads[collector_id]
            LOGGER.debug("Collector backend terminated. Collector ID: {:}".format(collector_id))
            self.GenerateCollectorTerminationSignal(collector_id, "-1", -1)          # Termination confirmation to frontend.
        else:
            LOGGER.warning('Backend collector {:} not found'.format(collector_id))

    def GenerateCollectorResponse(self, collector_id: str, kpi_id: str, measured_kpi_value: Any):
    def GenerateCollectorTerminationSignal(self, collector_id: str, kpi_id: str, measured_kpi_value: Any):
        """
        Method to write kpi value on RESPONSE Kafka topic
        Method to write kpi Termination signat on RESPONSE Kafka topic
        """
        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
            "kpi_value" : measured_kpi_value,
        }
        producer.produce(
            KafkaTopic.VALUE.value, # TODO: to  the topic ...
            KafkaTopic.RESPONSE.value, # TODO: to the topic ...
            key      = collector_id,
            value    = json.dumps(kpi_value),
            callback = self.delivery_callback
@@ -160,14 +154,9 @@ class TelemetryBackendService(GenericGrpcService):
        producer.flush()

    def delivery_callback(self, err, msg):
        """
        Callback function to handle message delivery status.
        Args: err (KafkaError): Kafka error object.
              msg (Message): Kafka message object.
        """
        if err: 
            LOGGER.error('Message delivery failed: {:}'.format(err))
            LOGGER.error('Message delivery failed: {:s}'.format(str(err)))
            # print(f'Message delivery failed: {err}')
        # else:
        #    LOGGER.debug('Message delivered to topic {:}'.format(msg.topic()))
        #    # print(f'Message delivered to topic {msg.topic()}')
        #     LOGGER.info('Message delivered to topic {:}'.format(msg.topic()))
            # print(f'Message delivered to topic {msg.topic()}')
+16 −0
Original line number Diff line number Diff line
@@ -12,4 +12,20 @@
# See the License for the specific language governing permissions and
# limitations under the License.

import uuid
import random
from common.proto import telemetry_frontend_pb2
# from common.proto.kpi_sample_types_pb2 import KpiSampleType
# from common.proto.kpi_manager_pb2 import KpiId

def create_collector_request():
    _create_collector_request                                = telemetry_frontend_pb2.Collector()
    _create_collector_request.collector_id.collector_id.uuid = str(uuid.uuid4()) 
    # _create_collector_request.collector_id.collector_id.uuid = "efef4d95-1cf1-43c4-9742-95c283dddddd"
    _create_collector_request.kpi_id.kpi_id.uuid             = str(uuid.uuid4())
    # _create_collector_request.kpi_id.kpi_id.uuid             = "6e22f180-ba28-4641-b190-2287bf448888"
    _create_collector_request.duration_s                     = float(random.randint(8, 16))
    # _create_collector_request.duration_s                     = -1
    _create_collector_request.interval_s                     = float(random.randint(2, 4)) 
    return _create_collector_request
+19 −3
Original line number Diff line number Diff line
@@ -13,10 +13,11 @@
# limitations under the License.

import logging
import threading
from typing import Dict
from .messages import create_collector_request
from common.tools.kafka.Variables import KafkaTopic
from telemetry.backend.service.TelemetryBackendService import TelemetryBackendService

import time

LOGGER = logging.getLogger(__name__)

@@ -35,3 +36,18 @@ def test_validate_kafka_topics():
#     LOGGER.info('test_RunRequestListener')
#     TelemetryBackendServiceObj = TelemetryBackendService()
#     threading.Thread(target=TelemetryBackendServiceObj.RequestListener).start()

def test_RunInitiateCollectorBackend():
    LOGGER.debug(">>> RunInitiateCollectorBackend <<<")
    collector_obj = create_collector_request()
    collector_id = collector_obj.collector_id.collector_id.uuid
    collector_dict :  Dict = {
        "kpi_id"  : collector_obj.kpi_id.kpi_id.uuid,
        "duration": collector_obj.duration_s,
        "interval": collector_obj.interval_s
    }
    TeleObj = TelemetryBackendService()
    TeleObj.InitiateCollectorBackend(collector_id, collector_dict)
    time.sleep(20)

    LOGGER.debug("--- Execution Finished Sucessfully---")
Loading