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

Changes to telemetry service

- TelemetryBackendService name and port is added in constants
- get_kafka_address() call added in TelemetryBackendService
- kafk topics test added in telemetrybackend test file
- get_kafka_address() call added in TelemetryBackendServiceImpl
- kafk topics test added in telemetryfrontend test file
- improvements in Telemetry gitlab-ci file
   + unit_test for telemetry backend and frontend added.
   + many other small changes
parent 1f495c78
Loading
Loading
Loading
Loading
+103 −45
Original line number Diff line number Diff line
@@ -21,6 +21,8 @@ build kpi-manager:
  before_script:
    - docker login -u "$CI_REGISTRY_USER" -p "$CI_REGISTRY_PASSWORD" $CI_REGISTRY
  script:
    # This first build tags the builder resulting image to prevent being removed by dangling image removal command
    - docker buildx build -t "${IMAGE_NAME}-backend:${IMAGE_TAG}-builder" --target builder -f ./src/$IMAGE_NAME/backend/Dockerfile .
    - docker buildx build -t "${IMAGE_NAME}-frontend:$IMAGE_TAG" -f ./src/$IMAGE_NAME/frontend/Dockerfile .
    - docker buildx build -t "${IMAGE_NAME}-backend:$IMAGE_TAG" -f ./src/$IMAGE_NAME/backend/Dockerfile .
    - docker tag "${IMAGE_NAME}-frontend:$IMAGE_TAG" "$CI_REGISTRY_IMAGE/${IMAGE_NAME}-frontend:$IMAGE_TAG"
@@ -35,6 +37,7 @@ build kpi-manager:
    - changes:
      - src/common/**/*.py
      - proto/*.proto
      - src/$IMAGE_NAME/.gitlab-ci.yml
      - src/$IMAGE_NAME/frontend/**/*.{py,in,yml}
      - src/$IMAGE_NAME/frontend/Dockerfile
      - src/$IMAGE_NAME/frontend/tests/*.py
@@ -45,7 +48,75 @@ build kpi-manager:
      - .gitlab-ci.yml

# Apply unit test to the component
unit_test telemetry:
unit_test telemetry-backend:
  variables:
    IMAGE_NAME: 'telemetry' # name of the microservice
    IMAGE_TAG: 'latest' # tag of the container image (production, development, etc)
  stage: unit_test
  needs:
    - build telemetry
  before_script:
    - docker login -u "$CI_REGISTRY_USER" -p "$CI_REGISTRY_PASSWORD" $CI_REGISTRY
    - if docker network list | grep teraflowbridge; then echo "teraflowbridge is already created"; else docker network create -d bridge teraflowbridge; fi
    - if docker container ls | grep kafka; then docker rm -f kafka; else echo "Kafka container is not in the system"; fi
    - if docker container ls | grep zookeeper; then docker rm -f zookeeper; else echo "Zookeeper container is not in the system"; fi
    # - if docker container ls | grep ${IMAGE_NAME}-frontend; then docker rm -f ${IMAGE_NAME}-frontend; else echo "${IMAGE_NAME}-frontend container is not in the system"; fi
    - if docker container ls | grep ${IMAGE_NAME}-backend; then docker rm -f ${IMAGE_NAME}-backend; else echo "${IMAGE_NAME}-backend container is not in the system"; fi
    - docker container prune -f
  script:
    - docker pull "$CI_REGISTRY_IMAGE/${IMAGE_NAME}-backend:$IMAGE_TAG"
    - docker pull "bitnami/zookeeper:latest"
    - docker pull "bitnami/kafka:latest"
    - >
      docker run --name zookeeper -d --network=teraflowbridge -p 2181:2181
      bitnami/zookeeper:latest
    - sleep 10 # Wait for Zookeeper to start
    - docker run --name kafka -d --network=teraflowbridge -p 9092:9092
      --env KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181
      --env ALLOW_PLAINTEXT_LISTENER=yes
      bitnami/kafka:latest
    - sleep 20 # Wait for Kafka to start
    - KAFKA_IP=$(docker inspect kafka --format "{{.NetworkSettings.Networks.teraflowbridge.IPAddress}}")
    - echo $KAFKA_IP    
    - >
      docker run --name $IMAGE_NAME -d -p 30060:30060
      --env "KFK_SERVER_ADDRESS=${KAFKA_IP}:9092"
      --volume "$PWD/src/$IMAGE_NAME/backend/tests:/opt/results"
      --network=teraflowbridge
      $CI_REGISTRY_IMAGE/${IMAGE_NAME}-backend:$IMAGE_TAG
    - docker ps -a
    - sleep 5
    - docker logs ${IMAGE_NAME}-backend
    - >
      docker exec -i ${IMAGE_NAME}-backend bash -c
      "coverage run -m pytest --log-level=INFO --verbose --junitxml=/opt/results/${IMAGE_NAME}-backend_report.xml $IMAGE_NAME/backend/tests/test_*.py"
    - docker exec -i ${IMAGE_NAME}-backend bash -c "coverage report --include='${IMAGE_NAME}/*' --show-missing"
  coverage: '/TOTAL\s+\d+\s+\d+\s+(\d+%)/'
  after_script:
    - docker network rm teraflowbridge
    - docker volume prune --force
    - docker image prune --force
    - docker rm -f ${IMAGE_NAME}-backend
    - docker rm -f zookeeper
    - docker rm -f kafka
  rules:
    - if: '$CI_PIPELINE_SOURCE == "merge_request_event" && ($CI_MERGE_REQUEST_TARGET_BRANCH_NAME == "develop" || $CI_MERGE_REQUEST_TARGET_BRANCH_NAME == $CI_DEFAULT_BRANCH)'
    - if: '$CI_PIPELINE_SOURCE == "push" && $CI_COMMIT_BRANCH == "develop"'
    - changes:
      - src/common/**/*.py
      - proto/*.proto
      - src/$IMAGE_NAME/backend/**/*.{py,in,yml}
      - src/$IMAGE_NAME/backend/Dockerfile
      - src/$IMAGE_NAME/backend/tests/*.py
      - manifests/${IMAGE_NAME}service.yaml
      - .gitlab-ci.yml
  artifacts:
      when: always
      reports:
        junit: src/$IMAGE_NAME/backend/tests/${IMAGE_NAME}-backend_report.xml

# Apply unit test to the component
unit_test telemetry-frontend:
  variables:
    IMAGE_NAME: 'telemetry' # name of the microservice
    IMAGE_TAG: 'latest' # tag of the container image (production, development, etc)
@@ -57,12 +128,14 @@ unit_test telemetry:
    - if docker network list | grep teraflowbridge; then echo "teraflowbridge is already created"; else docker network create -d bridge teraflowbridge; fi
    - if docker container ls | grep crdb; then docker rm -f crdb; else echo "CockroachDB container is not in the system"; fi
    - if docker volume ls | grep crdb; then docker volume rm -f crdb; else echo "CockroachDB volume is not in the system"; fi
    - if docker container ls | grep kafka; then docker rm -f kafka; else echo "Kafka container is not in the system"; fi
    - if docker container ls | grep zookeeper; then docker rm -f zookeeper; else echo "Zookeeper container is not in the system"; fi
    - if docker container ls | grep ${IMAGE_NAME}-frontend; then docker rm -f ${IMAGE_NAME}-frontend; else echo "${IMAGE_NAME}-frontend container is not in the system"; fi
    - if docker container ls | grep ${IMAGE_NAME}-backend; then docker rm -f ${IMAGE_NAME}-backend; else echo "${IMAGE_NAME}-backend container is not in the system"; fi
    - docker container prune -f
  script:
    - docker pull "$CI_REGISTRY_IMAGE/${IMAGE_NAME}-backend:$IMAGE_TAG"
    - docker pull "$CI_REGISTRY_IMAGE/${IMAGE_NAME}-frontend:$IMAGE_TAG"
    - docker pull "bitnami/zookeeper:latest"
    - docker pull "bitnami/kafka:latest"
    - docker pull "cockroachdb/cockroach:latest-v22.2"
    - docker volume create crdb
    - >
@@ -77,66 +150,51 @@ unit_test telemetry:
    - CRDB_ADDRESS=$(docker inspect crdb --format "{{.NetworkSettings.Networks.teraflowbridge.IPAddress}}")
    - echo $CRDB_ADDRESS
    - >
      docker run --name $IMAGE_NAME -d -p 30010:30010
      docker run --name zookeeper -d --network=teraflowbridge -p 2181:2181
      bitnami/zookeeper:latest
    - sleep 10 # Wait for Zookeeper to start
    - docker run --name kafka -d --network=teraflowbridge -p 9092:9092
      --env KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181
      --env ALLOW_PLAINTEXT_LISTENER=yes
      bitnami/kafka:latest
    - sleep 20 # Wait for Kafka to start
    - KAFKA_IP=$(docker inspect kafka --format "{{.NetworkSettings.Networks.teraflowbridge.IPAddress}}")
    - echo $KAFKA_IP    
    - >
      docker run --name $IMAGE_NAME -d -p 30050:30050
      --env "CRDB_URI=cockroachdb://tfs:tfs123@${CRDB_ADDRESS}:26257/tfs_test?sslmode=require"
      --volume "$PWD/src/$IMAGE_NAME/tests:/opt/results"
      --env "KFK_SERVER_ADDRESS=${KAFKA_IP}:9092"
      --volume "$PWD/src/$IMAGE_NAME/frontend/tests:/opt/results"
      --network=teraflowbridge
      $CI_REGISTRY_IMAGE/$IMAGE_NAME:$IMAGE_TAG
      $CI_REGISTRY_IMAGE/${IMAGE_NAME}-frontend:$IMAGE_TAG
    - docker ps -a
    - sleep 5
    - docker logs $IMAGE_NAME
    - docker logs ${IMAGE_NAME}-frontend
    - >
      docker exec -i $IMAGE_NAME bash -c
      "coverage run -m pytest --log-level=INFO --verbose --junitxml=/opt/results/${IMAGE_NAME}_report.xml $IMAGE_NAME/tests/test_*.py"
    - docker exec -i $IMAGE_NAME bash -c "coverage report --include='${IMAGE_NAME}/*' --show-missing"
      docker exec -i ${IMAGE_NAME}-frontend bash -c
      "coverage run -m pytest --log-level=INFO --verbose --junitxml=/opt/results/${IMAGE_NAME}-frontend_report.xml $IMAGE_NAME/frontend/tests/test_*.py"
    - docker exec -i ${IMAGE_NAME}-frontend bash -c "coverage report --include='${IMAGE_NAME}/*' --show-missing"
  coverage: '/TOTAL\s+\d+\s+\d+\s+(\d+%)/'
  after_script:
    - docker volume rm -f crdb
    - docker network rm teraflowbridge
    - docker volume prune --force
    - docker image prune --force
    - docker rm -f ${IMAGE_NAME}-frontend
    - docker rm -f zookeeper
    - docker rm -f kafka
  rules:
    - if: '$CI_PIPELINE_SOURCE == "merge_request_event" && ($CI_MERGE_REQUEST_TARGET_BRANCH_NAME == "develop" || $CI_MERGE_REQUEST_TARGET_BRANCH_NAME == $CI_DEFAULT_BRANCH)'
    - if: '$CI_PIPELINE_SOURCE == "push" && $CI_COMMIT_BRANCH == "develop"'
    - changes:
      - src/common/**/*.py
      - proto/*.proto
      - src/$IMAGE_NAME/**/*.{py,in,yml}
      - src/$IMAGE_NAME/frontend/**/*.{py,in,yml}
      - src/$IMAGE_NAME/frontend/Dockerfile
      - src/$IMAGE_NAME/frontend/tests/*.py
      - src/$IMAGE_NAME/frontend/tests/Dockerfile
      - src/$IMAGE_NAME/backend/Dockerfile
      - src/$IMAGE_NAME/backend/tests/*.py
      - src/$IMAGE_NAME/backend/tests/Dockerfile
      - manifests/${IMAGE_NAME}service.yaml
      - .gitlab-ci.yml
  # artifacts:
  #     when: always
  #     reports:
  #       junit: src/$IMAGE_NAME/tests/${IMAGE_NAME}_report.xml

## Deployment of the service in Kubernetes Cluster
#deploy context:
#  variables:
#    IMAGE_NAME: 'context' # name of the microservice
#    IMAGE_TAG: 'latest' # tag of the container image (production, development, etc)
#  stage: deploy
#  needs:
#    - unit test context
#    # - integ_test execute
#  script:
#    - 'sed -i "s/$IMAGE_NAME:.*/$IMAGE_NAME:$IMAGE_TAG/" manifests/${IMAGE_NAME}service.yaml'
#    - kubectl version
#    - kubectl get all
#    - kubectl apply -f "manifests/${IMAGE_NAME}service.yaml"
#    - kubectl get all
#  # environment:
#  #   name: test
#  #   url: https://example.com
#  #   kubernetes:
#  #     namespace: test
#  rules:
#    - if: '$CI_PIPELINE_SOURCE == "merge_request_event" && ($CI_MERGE_REQUEST_TARGET_BRANCH_NAME == "develop" || $CI_MERGE_REQUEST_TARGET_BRANCH_NAME == $CI_DEFAULT_BRANCH)'
#      when: manual    
#    - if: '$CI_PIPELINE_SOURCE == "push" && $CI_COMMIT_BRANCH == "develop"'
#      when: manual
  artifacts:
      when: always
      reports:
        junit: src/$IMAGE_NAME/frontend/tests/${IMAGE_NAME}-frontend_report.xml
 No newline at end of file
+2 −2
Original line number Diff line number Diff line
@@ -39,8 +39,8 @@ class TelemetryBackendService:

    def __init__(self):
        LOGGER.info('Init TelemetryBackendService')
        self.kafka_producer = KafkaProducer({'bootstrap.servers' : KafkaConfig.SERVER_ADDRESS.value})
        self.kafka_consumer = KafkaConsumer({'bootstrap.servers' : KafkaConfig.SERVER_ADDRESS.value,
        self.kafka_producer = KafkaProducer({'bootstrap.servers' : KafkaConfig.get_kafka_address()})
        self.kafka_consumer = KafkaConsumer({'bootstrap.servers' : KafkaConfig.get_kafka_address(),
                                            'group.id'           : 'backend',
                                            'auto.offset.reset'  : 'latest'})
        self.running_threads = {}
+7 −0
Original line number Diff line number Diff line
@@ -13,6 +13,7 @@
# limitations under the License.

import logging
from common.tools.kafka.Variables import KafkaTopic
from src.telemetry.backend.service.TelemetryBackendService import TelemetryBackendService


@@ -23,6 +24,12 @@ LOGGER = logging.getLogger(__name__)
# Tests Implementation of Telemetry Backend
###########################

# --- "test_validate_kafka_topics" should be run before the functionality tests ---
def test_validate_kafka_topics():
    LOGGER.debug(" >>> test_validate_kafka_topics: START <<< ")
    response = KafkaTopic.create_all_topics()
    assert isinstance(response, bool)

def test_RunRequestListener():
    LOGGER.info('test_RunRequestListener')
    TelemetryBackendServiceObj = TelemetryBackendService()
+2 −2
Original line number Diff line number Diff line
@@ -41,8 +41,8 @@ class TelemetryFrontendServiceServicerImpl(TelemetryFrontendServiceServicer):
    def __init__(self):
        LOGGER.info('Init TelemetryFrontendService')
        self.tele_db_obj = TelemetryDB()
        self.kafka_producer = KafkaProducer({'bootstrap.servers' : KafkaConfig.SERVER_ADDRESS.value})
        self.kafka_consumer = KafkaConsumer({'bootstrap.servers' : KafkaConfig.SERVER_ADDRESS.value,
        self.kafka_producer = KafkaProducer({'bootstrap.servers' : KafkaConfig.get_kafka_address()})
        self.kafka_consumer = KafkaConsumer({'bootstrap.servers' : KafkaConfig.get_kafka_address(),
                                            'group.id'           : 'frontend',
                                            'auto.offset.reset'  : 'latest'})

+9 −2
Original line number Diff line number Diff line
@@ -16,11 +16,10 @@ import os
import pytest
import logging

# from common.proto.context_pb2 import Empty
from common.Constants import ServiceNameEnum
from common.proto.telemetry_frontend_pb2 import CollectorId, CollectorList
from common.proto.context_pb2 import Empty

from common.tools.kafka.Variables import KafkaTopic
from common.Settings import ( 
    get_service_port_grpc, get_env_var_name, ENVVAR_SUFIX_SERVICE_HOST, ENVVAR_SUFIX_SERVICE_PORT_GRPC)

@@ -81,6 +80,13 @@ def telemetryFrontend_client(
###########################

# ------- Re-structuring Test ---------
# --- "test_validate_kafka_topics" should be run before the functionality tests ---
def test_validate_kafka_topics():
    LOGGER.debug(" >>> test_validate_kafka_topics: START <<< ")
    response = KafkaTopic.create_all_topics()
    assert isinstance(response, bool)

# ----- core funtionality test -----
def test_StartCollector(telemetryFrontend_client):
    LOGGER.info(' >>> test_StartCollector START: <<< ')
    response = telemetryFrontend_client.StartCollector(create_collector_request())
@@ -99,6 +105,7 @@ def test_SelectCollectors(telemetryFrontend_client):
    LOGGER.debug(str(response))
    assert isinstance(response, CollectorList)

# ----- Non-gRPC method tests ----- 
def test_RunResponseListener():
    LOGGER.info(' >>> test_RunResponseListener START: <<< ')
    TelemetryFrontendServiceObj = TelemetryFrontendServiceServicerImpl()