Commit d25c52c0 authored by Javi Moreno's avatar Javi Moreno
Browse files

Monitoring component updated to new information model

parent 7ca58a42
Loading
Loading
Loading
Loading
+0 −6
Original line number Diff line number Diff line
@@ -29,12 +29,6 @@ class MonitoringClient:
        LOGGER.info('MonitorKpi result: {}'.format(response))
        return context_pb2.Empty()

    def MonitorDeviceKpi(self, request):
        LOGGER.info('MonitorDeviceKpi: {}'.format(request))
        response = self.server.MonitorDeviceKpi(request)
        LOGGER.info('MonitorDeviceKpi result: {}'.format(response))
        return context_pb2.Empty()

    def IncludeKpi(self, request):
        LOGGER.info('IncludeKpi: {}'.format(request))
        response = self.server.IncludeKpi(request)
+44 −32
Original line number Diff line number Diff line
import os

import device
from device.Config import GRPC_SERVICE_PORT
from device.client.DeviceClient import DeviceClient
from device.proto import device_pb2
from monitoring.proto import context_pb2
from monitoring.service import sqlite_tools, influx_tools

@@ -21,28 +25,31 @@ INFLUXDB_USER = os.environ.get("INFLUXDB_USER")
INFLUXDB_PASSWORD = os.environ.get("INFLUXDB_PASSWORD")
INFLUXDB_DATABASE = os.environ.get("INFLUXDB_DATABASE")


class MonitoringServiceServicerImpl(monitoring_pb2_grpc.MonitoringServiceServicer):
    def __init__(self):
        LOGGER.info('Init monitoringService')

        # Init sqlite monitoring db
        self.sql_db = sqlite_tools.SQLite('monitoring.db')
        self.sql_db = sqlite_tools.SQLite('monitoring2.db')

        # Create influx_db client
        self.influx_db = influx_tools.Influx(INFLUXDB_HOSTNAME,"8086",INFLUXDB_USER,INFLUXDB_PASSWORD,INFLUXDB_DATABASE)

    # CreateKpi (CreateKpiRequest) returns (KpiId) {}
    def CreateKpi(self, request : monitoring_pb2.CreateKpiRequest, context) -> monitoring_pb2.KpiId :
    def CreateKpi(self, request : monitoring_pb2.KpiDescriptor, context) -> monitoring_pb2.KpiId :
        LOGGER.info('CreateKpi')

        # Here the code to create a sqlite query to crete a KPI and return a KpiID
        kpi_id = monitoring_pb2.KpiId()

        kpi_description = request.kpiDescription
        kpi_device_id = request.device_id.device_id.uuid
        kpi_description = request.kpi_description
        kpi_sample_type = request.kpi_sample_type
        kpi_device_id   = request.device_id.device_uuid.uuid
        kpi_endpoint_id = request.endpoint_id.endpoint_uuid.uuid
        kpi_service_id  = request.service_id.service_uuid.uuid

        data = self.sql_db.insert_KPI(kpi_description, kpi_device_id, kpi_sample_type)
        data = self.sql_db.insert_KPI(kpi_description, kpi_sample_type, kpi_device_id, kpi_endpoint_id, kpi_service_id)
        kpi_id.kpi_id.uuid = str(data)

        return kpi_id
@@ -53,39 +60,42 @@ class MonitoringServiceServicerImpl(monitoring_pb2_grpc.MonitoringServiceService
        LOGGER.info('MonitorKpi')

        # Creates the request to send to the device service
        monitor_device_request = monitoring_pb2.MonitorDeviceKpiRequest()
        kpi = self.get_Kpi(request.kpi_id)
        monitor_device_request = device_pb2.MonitoringSettings()
        kpiDescriptor = self.get_KpiDescriptor(request.kpi_id)

        monitor_device_request.kpi.kpi_id.kpi_id.uuid  = kpi.kpi_id.kpi_id.uuid
        monitor_device_request.connexion_time_s = request.connexion_time_s
        monitor_device_request.sample_rate_ms = request.sample_rate_ms
        monitor_device_request.kpi_id.kpi_id.uuid  = request.kpi_id.kpi_id.uuid
        monitor_device_request.kpi_descriptor.kpi_description                   = kpiDescriptor.kpi_description
        monitor_device_request.kpi_descriptor.kpi_sample_type                   = kpiDescriptor.kpi_sample_type
        monitor_device_request.kpi_descriptor.device_id.device_uuid.uuid        = kpiDescriptor.device_id.device_uuid.uuid
        monitor_device_request.kpi_descriptor.endpoint_id.endpoint_uuid.uuid    = kpiDescriptor.endpoint_id.endpoint_uuid.uuid
        monitor_device_request.kpi_descriptor.service_id.service_uuid.uuid      = kpiDescriptor.service_id.service_uuid.uuid
        monitor_device_request.sampling_duration_s                              = request.sampling_duration_s
        monitor_device_request.sampling_interval_s                              = request.sampling_interval_s

        self.MonitorDeviceKpi(monitor_device_request,context)
        # deviceClient = DeviceClient(address="localhost", port=GRPC_SERVICE_PORT )  # instantiate the client
        # deviceClient.MonitorDeviceKpi(monitor_device_request)

        return context_pb2.Empty()

    # rpc MonitorDeviceKpi(MonitorDeviceKpiRequest) returns(context.Empty) {}
    def MonitorDeviceKpi ( self, request : monitoring_pb2.MonitorDeviceKpiRequest, context) -> context_pb2.Empty:

        # Some code device to perform its actions

        LOGGER.info('MonitorDeviceKpi')

        # Notify device about monitoring (device client to add)

        return context_pb2.Empty()

    # rpc IncludeKpi(IncludeKpiRequest)  returns(context.Empty)    {}
    def IncludeKpi(self, request : monitoring_pb2.IncludeKpiRequest, context) -> context_pb2.Empty:
    def IncludeKpi(self, request : monitoring_pb2.Kpi, context) -> context_pb2.Empty:

        LOGGER.info('IncludeKpi')

        kpi = self.get_Kpi(request.kpi_id)
        time_stamp = request.time_stamp
        kpiDescriptor = self.get_KpiDescriptor(request.kpi_id)

        kpiSampleType = kpiDescriptor.kpi_sample_type
        kpiId = request.kpi_id.kpi_id.uuid
        deviceId = kpiDescriptor.device_id.device_uuid.uuid
        endpointId = kpiDescriptor.endpoint_id.endpoint_uuid.uuid
        serviceId = kpiDescriptor.service_id.service_uuid.uuid

        time_stamp = request.timestamp
        kpi_value = request.kpi_value.intVal

        # Build the structure to be included as point in the influxDB
        self.influx_db.write_KPI(time_stamp,kpi.kpi_id.kpi_id.uuid,kpi.device_id.device_id.uuid,kpi.kpi_sample_type,kpi_value)
        self.influx_db.write_KPI(time_stamp,kpiId,kpiSampleType,deviceId,endpointId,serviceId,kpi_value)

        self.influx_db.read_KPI_points()

@@ -103,15 +113,17 @@ class MonitoringServiceServicerImpl(monitoring_pb2_grpc.MonitoringServiceService
        LOGGER.info('GetInstantKpi')
        return monitoring_pb2.Kpi()

    def get_Kpi(self, kpiId):
    def get_KpiDescriptor(self, kpiId):
        LOGGER.info('getting Kpi by KpiID')

        kpi_db = self.sql_db.get_KPI(int(kpiId.kpi_id.uuid))

        kpi = monitoring_pb2.Kpi()
        kpi.kpi_id.kpi_id.uuid = str(kpi_db[0])
        kpi.kpiDescription = kpi_db[1]
        kpi.device_id.device_id.uuid = kpi_db[2]
        kpi.kpi_sample_type = kpi_db[3]
        kpiDescriptor = monitoring_pb2.KpiDescriptor()

        kpiDescriptor.kpi_description = kpi_db[1]
        kpiDescriptor.kpi_sample_type = kpi_db[2]
        kpiDescriptor.device_id.device_uuid.uuid = str(kpi_db[3])
        kpiDescriptor.endpoint_id.endpoint_uuid.uuid = str(kpi_db[4])
        kpiDescriptor.service_id.service_uuid.uuid = str(kpi_db[5])

        return kpi
 No newline at end of file
        return kpiDescriptor
 No newline at end of file
+4 −2
Original line number Diff line number Diff line
@@ -4,14 +4,16 @@ class Influx():
  def __init__(self, host, port, username, password, database):
      self.client = InfluxDBClient(host=host, port=port, username=username, password=password, database=database)

  def write_KPI(self,time,kpi_id,device_id,kpi_sample_type,kpi_value):
  def write_KPI(self,time,kpi_id,kpi_sample_type,device_id,endpoint_id,service_id,kpi_value):
    data = [{
      "measurement": "samples",
      "time": time,
      "tags": {
          "kpi_id" : kpi_id,
          "kpi_sample_type": kpi_sample_type,
          "device_id"  : device_id,
          "kpi_sample_type": kpi_sample_type
          "endpoint_id" : endpoint_id,
          "service_id" : service_id
      },
      "fields": {
          "kpi_value": kpi_value
+6 −4
Original line number Diff line number Diff line
@@ -6,18 +6,20 @@ class SQLite():
        self.client.execute("""
            CREATE TABLE IF NOT EXISTS KPI(
                kpi_id INTEGER PRIMARY KEY AUTOINCREMENT,
                kpiDescription TEXT,
                kpi_description TEXT,
                kpi_sample_type INTEGER,
                device_id INTEGER,
                kpi_sample_type INTEGER
                endpoint_id INTEGER,
                service_id INTEGER
            );
        """)

    def insert_KPI(self,kpiDescription,device_id,kpi_sample_type):
    def insert_KPI(self,kpi_description,kpi_sample_type,device_id,endpoint_id,service_id ):
        c = self.client.cursor()
        c.execute("SELECT kpi_id FROM KPI WHERE device_id is ? AND kpi_sample_type is ?",(device_id,kpi_sample_type))
        data=c.fetchone()
        if data is None:
            c.execute("INSERT INTO KPI (kpiDescription, device_id,kpi_sample_type) VALUES (?,?,?)", (kpiDescription,device_id,kpi_sample_type))
            c.execute("INSERT INTO KPI (kpi_description,kpi_sample_type,device_id,endpoint_id,service_id) VALUES (?,?,?,?,?)", (kpi_description,kpi_sample_type,device_id,endpoint_id,service_id))
            self.client.commit()
            return c.lastrowid
        else:
+11 −24
Original line number Diff line number Diff line
import logging, os
import pytest

from monitoring.proto import context_pb2
from monitoring.proto import context_pb2, kpi_sample_types_pb2
from monitoring.proto import monitoring_pb2
from monitoring.client.monitoring_client import MonitoringClient
from monitoring.Config import GRPC_SERVICE_PORT, GRPC_MAX_WORKERS, GRPC_GRACE_PERIOD, LOG_LEVEL, METRICS_PORT
@@ -71,10 +71,12 @@ def kpi_id():
def create_kpi_request():
    LOGGER.warning('test_include_kpi begin')

    create_kpi_request = monitoring_pb2.CreateKpiRequest()
    create_kpi_request.device_id.device_id.uuid = 'DEV1'  # pylint: disable=maybe-no-member
    create_kpi_request.kpiDescription = 'KPI Description'
    create_kpi_request.kpi_sample_type = monitoring_pb2.KpiSampleType.PACKETS_TRANSMITTED
    create_kpi_request = monitoring_pb2.KpiDescriptor()
    create_kpi_request.kpi_description = 'KPI Description Test'
    create_kpi_request.kpi_sample_type = kpi_sample_types_pb2.KpiSampleType.PACKETS_TRANSMITTED
    create_kpi_request.device_id.device_uuid.uuid = 'DEV1'  # pylint: disable=maybe-no-member
    create_kpi_request.service_id.service_uuid.uuid = "SERV1"
    create_kpi_request.endpoint_id.endpoint_uuid.uuid = "END1"

    return create_kpi_request

@@ -84,28 +86,19 @@ def monitor_kpi_request():

    monitor_kpi_request = monitoring_pb2.MonitorKpiRequest()
    monitor_kpi_request.kpi_id.kpi_id.uuid = str(1)
    monitor_kpi_request.connexion_time_s = 120
    monitor_kpi_request.sample_rate_ms = 5
    monitor_kpi_request.sampling_duration_s = 120
    monitor_kpi_request.sampling_interval_s = 5

    return monitor_kpi_request

@pytest.fixture(scope='session')
def monitor_device_kpi_request():
    LOGGER.warning('test_monitor_kpi begin')

    monitor_device_kpi_request = monitoring_pb2.MonitorDeviceKpiRequest()
    monitor_device_kpi_request.connexion_time_s = 120
    monitor_device_kpi_request.sample_rate_ms = 5

    return monitor_device_kpi_request

@pytest.fixture(scope='session')
def include_kpi_request():
    LOGGER.warning('test_include_kpi begin')

    include_kpi_request = monitoring_pb2.IncludeKpiRequest()
    include_kpi_request = monitoring_pb2.Kpi()
    include_kpi_request.kpi_id.kpi_id.uuid = str(1)
    include_kpi_request.time_stamp = "2021-10-12T13:14:42Z"
    include_kpi_request.timestamp = "2021-10-12T13:14:42Z"
    include_kpi_request.kpi_value.intVal = 500

    return include_kpi_request
@@ -129,12 +122,6 @@ def test_monitor_kpi(monitoring_client,monitor_kpi_request):
    LOGGER.debug(str(response))
    assert isinstance(response, context_pb2.Empty)

# Test case that makes use of client fixture to test server's MonitorDeviceKpi method
def test_monitor_device_kpi(monitoring_client,monitor_device_kpi_request):
    LOGGER.warning('test_monitor_device_kpi begin')
    response = monitoring_client.MonitorDeviceKpi(monitor_device_kpi_request)
    LOGGER.debug(str(response))
    assert isinstance(response, context_pb2.Empty)

# Test case that makes use of client fixture to test server's IncludeKpi method
def test_include_kpi(monitoring_client,include_kpi_request):