Commit 4713ece2 authored by Luis de la Cal's avatar Luis de la Cal
Browse files

Merge branch 'feat/monitoring' into feat/l3-components

parents a8e18718 ae0398a5
Loading
Loading
Loading
Loading
+1 −0
Original line number Diff line number Diff line
@@ -48,6 +48,7 @@ message KpiDescriptor {
  context.EndPointId             endpoint_id     = 6;
  context.ServiceId              service_id      = 7;
  context.SliceId                slice_id        = 8;
  context.ConnectionId           connection_id   = 9;
}

message MonitorKpiRequest {
+10 −3
Original line number Diff line number Diff line
import pytz
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.executors.pool import ProcessPoolExecutor
from apscheduler.jobstores.base import JobLookupError
@@ -19,10 +20,16 @@ class AlarmManager():
        end_date=None
        if subscription_timeout_s:
            start_timestamp=time.time()
            start_date=datetime.fromtimestamp(start_timestamp)
            end_date=datetime.fromtimestamp(start_timestamp+subscription_timeout_s)
        self.scheduler.add_job(self.metrics_db.get_alarm_data, args=(alarm_queue,kpi_id, kpiMinValue, kpiMaxValue, inRange, includeMinValue, includeMaxValue, subscription_frequency_ms),trigger='interval', seconds=(subscription_frequency_ms/1000), start_date=start_date, end_date=end_date, id=alarm_id)
            end_timestamp = start_timestamp + subscription_timeout_s
            start_date = datetime.utcfromtimestamp(start_timestamp).isoformat()
            end_date = datetime.utcfromtimestamp(end_timestamp).isoformat()

        job = self.scheduler.add_job(self.metrics_db.get_alarm_data,
                               args=(alarm_queue,kpi_id, kpiMinValue, kpiMaxValue, inRange, includeMinValue, includeMaxValue, subscription_frequency_ms),
                               trigger='interval', seconds=(subscription_frequency_ms/1000), start_date=start_date,
                               end_date=end_date,timezone=pytz.utc, id=str(alarm_id))
        LOGGER.debug(f"Alarm job {alarm_id} succesfully created")
        job.remove()

    def delete_alarm(self, alarm_id):
        try:
+10 −6
Original line number Diff line number Diff line
@@ -40,7 +40,9 @@ class ManagementDB:
                    kpi_sample_type INTEGER,
                    device_id INTEGER,
                    endpoint_id INTEGER,
                    service_id INTEGER
                    service_id INTEGER,
                    slice_id INTEGER,
                    connection_id INTEGER
                );
            """
            )
@@ -91,18 +93,20 @@ class ManagementDB:
            LOGGER.debug(f"Alarm table cannot be created in the ManagementDB. {e}")
            raise Exception

    def insert_KPI(self, kpi_description, kpi_sample_type, device_id, endpoint_id, service_id):
    def insert_KPI(
        self, kpi_description, kpi_sample_type, device_id, endpoint_id, service_id, slice_id, connection_id
    ):
        try:
            c = self.client.cursor()
            c.execute(
                "SELECT kpi_id FROM kpi WHERE device_id is ? AND kpi_sample_type is ? AND endpoint_id is ? AND service_id is ?",
                (device_id, kpi_sample_type, endpoint_id, service_id),
                "SELECT kpi_id FROM kpi WHERE device_id is ? AND kpi_sample_type is ? AND endpoint_id is ? AND service_id is ? AND slice_id is ? AND connection_id is ?",
                (device_id, kpi_sample_type, endpoint_id, service_id, slice_id, connection_id),
            )
            data = c.fetchone()
            if data is None:
                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),
                    "INSERT INTO kpi (kpi_description,kpi_sample_type,device_id,endpoint_id,service_id,slice_id,connection_id) VALUES (?,?,?,?,?,?,?)",
                    (kpi_description, kpi_sample_type, device_id, endpoint_id, service_id, slice_id, connection_id),
                )
                self.client.commit()
                kpi_id = c.lastrowid
+13 −7
Original line number Diff line number Diff line
@@ -87,6 +87,8 @@ class MetricsDB():
                    'device_id SYMBOL,' \
                    'endpoint_id SYMBOL,' \
                    'service_id SYMBOL,' \
                    'slice_id SYMBOL,' \
                    'connection_id SYMBOL,' \
                    'timestamp TIMESTAMP,' \
                    'kpi_value DOUBLE)' \
                    'TIMESTAMP(timestamp);'
@@ -97,7 +99,7 @@ class MetricsDB():
            LOGGER.debug(f"Table {self.table} cannot be created. {e}")
            raise Exception

    def write_KPI(self, time, kpi_id, kpi_sample_type, device_id, endpoint_id, service_id, kpi_value):
    def write_KPI(self, time, kpi_id, kpi_sample_type, device_id, endpoint_id, service_id, slice_id, connection_id, kpi_value):
        counter = 0
        while (counter < self.retries):
            try:
@@ -109,7 +111,9 @@ class MetricsDB():
                            'kpi_sample_type': kpi_sample_type,
                            'device_id': device_id,
                            'endpoint_id': endpoint_id,
                            'service_id': service_id},
                            'service_id': service_id,
                            'slice_id': slice_id,
                            'connection_id': connection_id,},
                        columns={
                            'kpi_value': kpi_value},
                        at=datetime.datetime.fromtimestamp(time))
@@ -201,6 +205,8 @@ class MetricsDB():
                kpi_list = self.run_query(query)
            if kpi_list:
                LOGGER.debug(f"New data received for alarm of KPI {kpi_id}")
                LOGGER.info(kpi_list)
                valid_kpi_list = []
                for kpi in kpi_list:
                    alarm = False
                    kpi_value = kpi[2]
@@ -263,8 +269,8 @@ class MetricsDB():
                        if (kpi_value >= kpiMaxValue):
                            alarm = True
                    if alarm:
                        # queue.append[kpi]
                        alarm_queue.put_nowait(kpi)
                        valid_kpi_list.append(kpi)
                alarm_queue.put_nowait(valid_kpi_list)
                LOGGER.debug(f"Alarm of KPI {kpi_id} triggered -> kpi_value:{kpi[2]}, timestamp:{kpi[1]}")
            else:
                LOGGER.debug(f"No new data for the alarm of KPI {kpi_id}")
+62 −47
Original line number Diff line number Diff line
@@ -28,7 +28,7 @@ from common.proto.monitoring_pb2 import AlarmResponse, AlarmDescriptor, AlarmLis
    KpiDescriptor, KpiList, KpiQuery, SubsDescriptor, SubscriptionID, AlarmID, KpiDescriptorList, \
    MonitorKpiRequest, Kpi, AlarmSubscription, SubsResponse
from common.rpc_method_wrapper.ServiceExceptions import ServiceException
from common.tools.timestamp.Converters import timestamp_string_to_float
from common.tools.timestamp.Converters import timestamp_string_to_float, timestamp_utcnow_to_float

from monitoring.service import ManagementDBTools, MetricsDBTools
from device.client.DeviceClient import DeviceClient
@@ -85,13 +85,16 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
            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
            kpi_slice_id = request.slice_id.slice_uuid.uuid
            kpi_connection_id = request.connection_id.connection_uuid.uuid


            if request.kpi_id.kpi_id.uuid is not "":
                response.kpi_id.uuid = request.kpi_id.kpi_id.uuid
            #     Here the code to modify an existing kpi
            else:
                data = self.management_db.insert_KPI(
                    kpi_description, kpi_sample_type, kpi_device_id, kpi_endpoint_id, kpi_service_id)
                    kpi_description, kpi_sample_type, kpi_device_id, kpi_endpoint_id, kpi_service_id, kpi_slice_id, kpi_connection_id)
                response.kpi_id.uuid = str(data)

            return response
@@ -136,6 +139,8 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
                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])
                kpiDescriptor.slice_id.slice_uuid.uuid              = str(kpi_db[6])
                kpiDescriptor.connection_id.connection_uuid.uuid    = str(kpi_db[7])
            return kpiDescriptor
        except ServiceException as e:
            LOGGER.exception('GetKpiDescriptor exception')
@@ -160,6 +165,8 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
                kpi_descriptor.device_id.device_uuid.uuid           = str(item[3])
                kpi_descriptor.endpoint_id.endpoint_uuid.uuid       = str(item[4])
                kpi_descriptor.service_id.service_uuid.uuid         = str(item[5])
                kpi_descriptor.slice_id.slice_uuid.uuid             = str(item[6])
                kpi_descriptor.connection_id.connection_uuid.uuid   = str(item[7])

                kpi_descriptor_list.kpi_descriptor_list.append(kpi_descriptor)

@@ -186,11 +193,13 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
                deviceId = kpiDescriptor.device_id.device_uuid.uuid
                endpointId = kpiDescriptor.endpoint_id.endpoint_uuid.uuid
                serviceId = kpiDescriptor.service_id.service_uuid.uuid
                sliceId   = kpiDescriptor.slice_id.slice_uuid.uuid
                connectionId = kpiDescriptor.connection_id.connection_uuid.uuid
                time_stamp = request.timestamp.timestamp
                kpi_value = getattr(request.kpi_value, request.kpi_value.WhichOneof('value'))

                # Build the structure to be included as point in the MetricsDB
                self.metrics_db.write_KPI(time_stamp, kpiId, kpiSampleType, deviceId, endpointId, serviceId, kpi_value)
                self.metrics_db.write_KPI(time_stamp, kpiId, kpiSampleType, deviceId, endpointId, serviceId, sliceId, connectionId, kpi_value)

            return Empty()
        except ServiceException as e:
@@ -250,9 +259,7 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):

        LOGGER.info('SubscribeKpi')
        try:

            subs_queue = Queue()
            subs_response = SubsResponse()

            kpi_id = request.kpi_id.kpi_id.uuid
            sampling_duration_s = request.sampling_duration_s
@@ -268,7 +275,9 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
                                                  start_timestamp, end_timestamp)

            # parse queue to append kpis into the list
            while True:
                while not subs_queue.empty():
                    subs_response = SubsResponse()
                    list = subs_queue.get_nowait()
                    for item in list:
                        kpi = Kpi()
@@ -276,10 +285,11 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
                        kpi.timestamp.timestamp = timestamp_string_to_float(item[1])
                        kpi.kpi_value.floatVal = item[2]  # This must be improved
                        subs_response.kpi_list.kpi.append(kpi)

                    subs_response.subs_id.subs_id.uuid = str(subs_id)

                    yield subs_response
                if timestamp_utcnow_to_float() > end_timestamp:
                    break
            # yield subs_response
        except ServiceException as e:
            LOGGER.exception('SubscribeKpi exception')
            grpc_context.abort(e.code, e.details)
@@ -424,6 +434,7 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
        LOGGER.info('GetAlarmDescriptor')
        try:
            alarm_id = request.alarm_id.uuid
            LOGGER.debug(alarm_id)
            alarm = self.management_db.get_alarm(alarm_id)
            response = AlarmDescriptor()

@@ -454,15 +465,13 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
        LOGGER.info('GetAlarmResponseStream')
        try:
            alarm_id = request.alarm_id.alarm_id.uuid
            alarm = self.management_db.get_alarm(alarm_id)
            alarm_response = AlarmResponse()

            if alarm:
            alarm_data = self.management_db.get_alarm(alarm_id)
            real_start_time = timestamp_utcnow_to_float()

            if alarm_data:
                LOGGER.debug(f"{alarm_data}")
                alarm_queue = Queue()

                alarm_data = self.management_db.get_alarm(alarm)

                alarm_id = request.alarm_id.alarm_id.uuid
                kpi_id = alarm_data[3]
                kpiMinValue = alarm_data[4]
@@ -473,24 +482,30 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
                subscription_frequency_ms = request.subscription_frequency_ms
                subscription_timeout_s = request.subscription_timeout_s

                end_timestamp = real_start_time + subscription_timeout_s

                self.alarm_manager.create_alarm(alarm_queue, alarm_id, kpi_id, kpiMinValue, kpiMaxValue, inRange,
                                                includeMinValue, includeMaxValue, subscription_frequency_ms,
                                                subscription_timeout_s)

                while True:
                    while not alarm_queue.empty():
                        alarm_response = AlarmResponse()
                        list = alarm_queue.get_nowait()
                        size = len(list)
                        for item in list:
                            kpi = Kpi()
                            kpi.kpi_id.kpi_id.uuid = str(item[0])
                            kpi.timestamp.timestamp = timestamp_string_to_float(item[1])
                            kpi.kpi_value.floatVal = item[2]  # This must be improved
                            alarm_response.kpi_list.kpi.append(kpi)

                        alarm_response.alarm_id.alarm_id.uuid = alarm_id

                        yield alarm_response
                    if timestamp_utcnow_to_float() > end_timestamp:
                        break
            else:
                LOGGER.info('GetAlarmResponseStream error: AlarmID({:s}): not found in database'.format(str(alarm_id)))
                alarm_response = AlarmResponse()
                alarm_response.alarm_id.alarm_id.uuid = "NoID"
                return alarm_response
        except ServiceException as e:
@@ -527,7 +542,7 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
        kpi_db = self.management_db.get_KPI(int(kpi_id))
        response = Kpi()
        if kpi_db is None:
            LOGGER.info('GetInstantKpi error: KpiID({:s}): not found in database'.format(str(kpi_id)))
            LOGGER.info('GetStreamKpi error: KpiID({:s}): not found in database'.format(str(kpi_id)))
            response.kpi_id.kpi_id.uuid = "NoID"
            return response
        else:
Loading