Commit 3a936974 authored by Waleed Akbar's avatar Waleed Akbar
Browse files

Updated Kafka topic names for Telemetry service.

parent 56f63933
Loading
Loading
Loading
Loading
+4 −4
Original line number Diff line number Diff line
@@ -41,14 +41,14 @@ class KafkaConfig(Enum):

class KafkaTopic(Enum):
    # TODO: Later to be populated from ENV variable.
    REQUEST            = 'topic_request' 
    RESPONSE           = 'topic_response'
    TELEMETRY_REQUEST  = 'topic_telemetry_request' 
    TELEMETRY_RESPONSE = 'topic_telemetry_response'
    RAW                = 'topic_raw' 
    LABELED            = 'topic_labeled'
    VALUE              = 'topic_value'
    ALARMS             = 'topic_alarms'
    ANALYTICS_REQUEST  = 'topic_request_analytics'
    ANALYTICS_RESPONSE = 'topic_response_analytics'
    ANALYTICS_REQUEST  = 'topic_analytics_request'
    ANALYTICS_RESPONSE = 'topic_analytics_response'

    @staticmethod
    def create_all_topics() -> bool:
+6 −6
Original line number Diff line number Diff line
@@ -36,7 +36,7 @@ METRICS_POOL = MetricsPool('TelemetryBackend', 'backendService')
class TelemetryBackendService(GenericGrpcService):
    """
    Class listens for request on Kafka topic, fetches requested metrics from device.
    Produces metrics on both RESPONSE and VALUE kafka topics.
    Produces metrics on both TELEMETRY_RESPONSE and VALUE kafka topics.
    """
    def __init__(self, cls_name : str = __name__) -> None:
        LOGGER.info('Init TelemetryBackendService')
@@ -60,7 +60,7 @@ class TelemetryBackendService(GenericGrpcService):
        LOGGER.info('Telemetry backend request listener is running ...')
        # print      ('Telemetry backend request listener is running ...')
        consumer = self.kafka_consumer
        consumer.subscribe([KafkaTopic.REQUEST.value])
        consumer.subscribe([KafkaTopic.TELEMETRY_REQUEST.value])
        while True:
            receive_msg = consumer.poll(2.0)
            if receive_msg is None:
@@ -82,7 +82,7 @@ class TelemetryBackendService(GenericGrpcService):
                    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))
                LOGGER.warning("Unable to consumer message from topic: {:}. ERROR: {:}".format(KafkaTopic.TELEMETRY_REQUEST.value, e))

    def InitiateCollectorBackend(self, collector_id, collector):
        """
@@ -109,7 +109,7 @@ class TelemetryBackendService(GenericGrpcService):

    def GenerateCollectorResponse(self, collector_id: str, kpi_id: str, measured_kpi_value: Any):
        """
        Method to write kpi value on RESPONSE Kafka topic
        Method to write kpi value on TELEMETRY_RESPONSE Kafka topic
        """
        producer = self.kafka_producer
        kpi_value : Dict = {
@@ -138,7 +138,7 @@ class TelemetryBackendService(GenericGrpcService):

    def GenerateCollectorTerminationSignal(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 Termination signat on TELEMETRY_RESPONSE Kafka topic
        """
        producer = self.kafka_producer
        kpi_value : Dict = {
@@ -146,7 +146,7 @@ class TelemetryBackendService(GenericGrpcService):
            "kpi_value" : measured_kpi_value,
        }
        producer.produce(
            KafkaTopic.RESPONSE.value, # TODO: to the topic ...
            KafkaTopic.TELEMETRY_RESPONSE.value, # TODO: to the topic ...
            key      = collector_id,
            value    = json.dumps(kpi_value),
            callback = self.delivery_callback
+3 −3
Original line number Diff line number Diff line
@@ -74,7 +74,7 @@ class TelemetryFrontendServiceServicerImpl(TelemetryFrontendServiceServicer):
            "interval": collector_obj.interval_s
        }
        self.kafka_producer.produce(
            KafkaTopic.REQUEST.value,
            KafkaTopic.TELEMETRY_REQUEST.value,
            key      = collector_uuid,
            value    = json.dumps(collector_to_generate),
            callback = self.delivery_callback
@@ -110,7 +110,7 @@ class TelemetryFrontendServiceServicerImpl(TelemetryFrontendServiceServicer):
            "interval": -1
        }
        self.kafka_producer.produce(
            KafkaTopic.REQUEST.value,
            KafkaTopic.TELEMETRY_REQUEST.value,
            key      = collector_uuid,
            value    = json.dumps(collector_to_stop),
            callback = self.delivery_callback
@@ -168,7 +168,7 @@ class TelemetryFrontendServiceServicerImpl(TelemetryFrontendServiceServicer):
        """
        listener for response on Kafka topic.
        """
        self.kafka_consumer.subscribe([KafkaTopic.RESPONSE.value])
        self.kafka_consumer.subscribe([KafkaTopic.TELEMETRY_RESPONSE.value])
        while True:
            receive_msg = self.kafka_consumer.poll(2.0)
            if receive_msg is None: