Commit 4cff2b80 authored by Carlos Natalino's avatar Carlos Natalino
Browse files

Merge branch 'feat/monitoring' into fix/ofc22-tests

parents 348bed26 92142ab1
Loading
Loading
Loading
Loading
+2 −1
Original line number Diff line number Diff line
@@ -9,13 +9,14 @@ Jinja2==3.0.3
ncclient==0.6.13
p4runtime==1.3.0
paramiko==2.9.2
influx-line-protocol==0.1.4
# influx-line-protocol==0.1.4
python-dateutil==2.8.2
python-json-logger==2.0.2
pytz==2021.3
redis==4.1.2
requests==2.27.1
xmltodict==0.12.0
questdb==1.0.1

# pip's dependency resolver does not take into account installed packages.
# p4runtime does not specify the version of grpcio/protobuf it needs, so it tries to install latest one
+3 −6
Original line number Diff line number Diff line
@@ -19,16 +19,13 @@ import grpc

from common.rpc_method_wrapper.ServiceExceptions import ServiceException
from context.client.ContextClient import ContextClient
#from common.proto import kpi_sample_types_pb2

from common.proto.context_pb2 import Empty, EventTypeEnum

from common.logger import getJSONLogger
from monitoring.client.MonitoringClient import MonitoringClient
from monitoring.service.MonitoringServiceServicerImpl import LOGGER
from common.proto import monitoring_pb2

LOGGER = getJSONLogger('monitoringservice-server')
LOGGER.setLevel('DEBUG')

class EventsDeviceCollector:
    def __init__(self) -> None: # pylint: disable=redefined-outer-name
        self._events_queue = Queue()
@@ -74,7 +71,7 @@ class EventsDeviceCollector:
            kpi_id_list = []

            while not self._events_queue.empty():
                LOGGER.info('getting Kpi by KpiID')
                # LOGGER.info('getting Kpi by KpiID')
                event = self.get_event(block=True)
                if event.event.event_type == EventTypeEnum.EVENTTYPE_CREATE:
                    device = self._context_client.GetDevice(event.device_id)
+26 −17
Original line number Diff line number Diff line
@@ -12,18 +12,16 @@
# See the License for the specific language governing permissions and
# limitations under the License.

from influx_line_protocol import Metric
import socket
from questdb.ingress import Sender, IngressError
import requests
import json
import sys
import logging
import datetime

LOGGER = logging.getLogger(__name__)

class MetricsDB():
  def __init__(self, host, ilp_port, rest_port, table):
    self.socket=socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    self.host=host
    self.ilp_port=int(ilp_port)
    self.rest_port=rest_port
@@ -31,19 +29,30 @@ class MetricsDB():
    self.create_table()

  def write_KPI(self,time,kpi_id,kpi_sample_type,device_id,endpoint_id,service_id,kpi_value):
    self.socket.connect((self.host,self.ilp_port))
    metric = Metric(self.table)
    metric.with_timestamp(time)
    metric.add_tag('kpi_id', kpi_id)
    metric.add_tag('kpi_sample_type', kpi_sample_type)
    metric.add_tag('device_id', device_id)
    metric.add_tag('endpoint_id', endpoint_id)
    metric.add_tag('service_id', service_id)
    metric.add_value('kpi_value', kpi_value)
    str_metric = str(metric)
    str_metric += "\n"
    self.socket.sendall((str_metric).encode())
    self.socket.close()
    counter=0
    number_of_retries=10
    while (counter<number_of_retries):
      try:
        with Sender(self.host, self.ilp_port) as sender:
          sender.row(
          self.table,
          symbols={
              'kpi_id': kpi_id,
              'kpi_sample_type': kpi_sample_type,
              'device_id': device_id,
              'endpoint_id': endpoint_id,
              'service_id': service_id},
          columns={
              'kpi_value': kpi_value},
          at=datetime.datetime.fromtimestamp(time))
          sender.flush()
        counter=number_of_retries
        LOGGER.info(f"KPI written")
      except IngressError as ierr:
        # LOGGER.info(ierr)
        # LOGGER.info(f"Ingress Retry number {counter}")
        counter=counter+1


  def run_query(self, sql_query):
    query_params = {'query': sql_query, 'fmt' : 'json'}
+7 −8
Original line number Diff line number Diff line
@@ -18,6 +18,7 @@ from typing import Iterator

from common.Constants import ServiceNameEnum
from common.Settings import get_setting, get_service_port_grpc, get_service_host
from common.logger import getJSONLogger
from common.proto.context_pb2 import Empty
from common.proto.device_pb2 import MonitoringSettings
from common.proto.kpi_sample_types_pb2 import KpiSampleType
@@ -26,14 +27,14 @@ from common.proto.monitoring_pb2 import AlarmResponse, AlarmDescriptor, AlarmIDL
    KpiDescriptor, KpiList, KpiQuery, SubsDescriptor, SubscriptionID, AlarmID, KpiDescriptorList, \
    MonitorKpiRequest, Kpi, AlarmSubscription
from common.rpc_method_wrapper.ServiceExceptions import ServiceException
from common.tools.timestamp.Converters import timestamp_float_to_string

from monitoring.service import SqliteTools, MetricsDBTools
from device.client.DeviceClient import DeviceClient

from prometheus_client import Counter, Summary

LOGGER = logging.getLogger(__name__)
LOGGER = getJSONLogger('monitoringservice-server')
LOGGER.setLevel('DEBUG')

MONITORING_GETINSTANTKPI_REQUEST_TIME = Summary(
    'monitoring_getinstantkpi_processing_seconds', 'Time spent processing monitoring instant kpi request')
@@ -57,7 +58,6 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
        self.sql_db = SqliteTools.SQLite('monitoring.db')
        self.deviceClient = DeviceClient(host=DEVICESERVICE_SERVICE_HOST, port=DEVICESERVICE_SERVICE_PORT_GRPC)  # instantiate the client

        # Set metrics_db client
        self.metrics_db = MetricsDBTools.MetricsDB(METRICSDB_HOSTNAME,METRICSDB_ILP_PORT,METRICSDB_REST_PORT,METRICSDB_TABLE)
        LOGGER.info('MetricsDB initialized')

@@ -81,7 +81,6 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
                kpi_description, kpi_sample_type, kpi_device_id, kpi_endpoint_id, kpi_service_id)

            kpi_id.kpi_id.uuid = str(data)

            # CREATEKPI_COUNTER_COMPLETED.inc()
            return kpi_id
        except ServiceException as e:
@@ -162,7 +161,7 @@ class MonitoringServiceServicerImpl(MonitoringServiceServicer):
            deviceId        = kpiDescriptor.device_id.device_uuid.uuid
            endpointId      = kpiDescriptor.endpoint_id.endpoint_uuid.uuid
            serviceId       = kpiDescriptor.service_id.service_uuid.uuid
            time_stamp      = timestamp_float_to_string(request.timestamp.timestamp)
            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
+6 −10
Original line number Diff line number Diff line
@@ -11,17 +11,13 @@
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
import datetime

from common.proto import monitoring_pb2
from common.proto.kpi_sample_types_pb2 import KpiSampleType
from common.tools.timestamp.Converters import timestamp_string_to_float
from common.tools.timestamp.Converters import timestamp_string_to_float, timestamp_utcnow_to_float


def kpi():
    _kpi                    = monitoring_pb2.Kpi()
    _kpi.kpi_id.kpi_id.uuid = 'KPIID0000'   # pylint: disable=maybe-no-member
    return _kpi

def kpi_id():
    _kpi_id             = monitoring_pb2.KpiId()
    _kpi_id.kpi_id.uuid = str(1)            # pylint: disable=maybe-no-member
@@ -43,9 +39,9 @@ def monitor_kpi_request(kpi_uuid, monitoring_window_s, sampling_rate_s):
    _monitor_kpi_request.sampling_rate_s     = sampling_rate_s
    return _monitor_kpi_request

def include_kpi_request():
def include_kpi_request(kpi_id):
    _include_kpi_request                        = monitoring_pb2.Kpi()
    _include_kpi_request.kpi_id.kpi_id.uuid     = str(1)    # pylint: disable=maybe-no-member
    _include_kpi_request.timestamp.timestamp    = timestamp_string_to_float("2021-10-12T13:14:42Z")
    _include_kpi_request.kpi_id.kpi_id.uuid     = kpi_id.kpi_id.uuid
    _include_kpi_request.timestamp.timestamp    = timestamp_utcnow_to_float()
    _include_kpi_request.kpi_value.int32Val     = 500       # pylint: disable=maybe-no-member
    return _include_kpi_request
Loading