Loading src/monitoring/client/monitoring_client.py +4 −10 Original line number Diff line number Diff line Loading @@ -21,19 +21,19 @@ class MonitoringClient: LOGGER.info('CreateKpi: {}'.format(request)) response = self.server.CreateKpi(request) LOGGER.info('CreateKpi result: {}'.format(response)) return monitoring_pb2.KpiId() return response def MonitorKpi(self, request): LOGGER.info('MonitorKpi: {}'.format(request)) response = self.server.MonitorKpi(request) LOGGER.info('MonitorKpi result: {}'.format(response)) return context_pb2.Empty() return response def IncludeKpi(self, request): LOGGER.info('IncludeKpi: {}'.format(request)) response = self.server.IncludeKpi(request) LOGGER.info('IncludeKpi result: {}'.format(response)) return context_pb2.Empty() return response def GetStreamKpi(self, request): LOGGER.info('GetStreamKpi: {}'.format(request)) Loading @@ -51,13 +51,7 @@ class MonitoringClient: LOGGER.info('GetKpiDescriptor: {}'.format(request)) response = self.server.GetKpiDescriptor(request) LOGGER.info('GetKpiDescriptor result: {}'.format(response)) return monitoring_pb2.KpiDescriptor() def ListenEvents(self, ): LOGGER.info('ListenEvents: {}'.format()) response = self.server.ListenEvents() LOGGER.info('ListenEvents result: {}'.format(response)) return monitoring_pb2.KpiDescriptor() return response if __name__ == '__main__': # get port Loading src/monitoring/service/EventTools.py +21 −19 Original line number Diff line number Diff line Loading @@ -47,16 +47,18 @@ class EventsDeviceCollector: def listen_events(self): LOGGER.info('getting Kpi by KpiID') qsize = self._events_queue.qsize() try: kpi_id_list = [] for i in range(self._events_queue.qsize()): if qsize > 0: for i in range(qsize): print("Queue size: "+str(qsize)) event = self.get_event(block=True) if event.event.event_type == EventTypeEnum.EVENTTYPE_CREATE: device = self._context_client.GetDevice(event.device_id) print("Endpoints value: " + str(len(device.device_endpoints))) for j,end_point in enumerate(device.device_endpoints): # for k,rule in enumerate(device.device_config.config_rules): kpi_descriptor = monitoring_pb2.KpiDescriptor() Loading src/monitoring/service/SqliteTools.py +1 −1 Original line number Diff line number Diff line Loading @@ -16,7 +16,7 @@ class SQLite(): 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)) c.execute("SELECT kpi_id FROM KPI WHERE device_id is ? AND kpi_sample_type is ? AND endpoint_id is ?",(device_id,kpi_sample_type,endpoint_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)) Loading src/monitoring/service/__main__.py +38 −1 Original line number Diff line number Diff line import logging, os, signal, sys, threading import time from context.client.ContextClient import ContextClient from monitoring.Config import GRPC_SERVICE_PORT, GRPC_MAX_WORKERS, GRPC_GRACE_PERIOD, LOG_LEVEL, METRICS_PORT from common.logger import getJSONLogger from monitoring.client.monitoring_client import MonitoringClient from monitoring.proto import monitoring_pb2 from monitoring.service.EventTools import EventsDeviceCollector from monitoring.service.MonitoringService import MonitoringService LOGGER = getJSONLogger('monitoringservice-server') Loading @@ -16,9 +20,40 @@ logger = None def signal_handler(signal, frame): global terminate, logger logger.warning('Terminate signal received') LOGGER.warning('Terminate signal received') terminate.set() def start_monitoring(): LOGGER.info('Start Monitoring...') context_client_grpc = ContextClient(address='localhost', port='2020') monitoring_client = MonitoringClient(server='localhost', port='7070') # instantiate the client while True: if terminate.is_set(): LOGGER.warning("Stopping execution...") break # Start Listening Events events_collector = EventsDeviceCollector(context_client_grpc, monitoring_client) events_collector.start() list_new_kpi_ids = events_collector.listen_events() # Monitor Kpis if bool(list_new_kpi_ids): for kpi_id in list_new_kpi_ids: # Create Monitor Kpi Requests monitor_kpi_request = monitoring_pb2.MonitorKpiRequest() monitor_kpi_request.kpi_id.CopyFrom(kpi_id) monitor_kpi_request.sampling_duration_s = 120 monitor_kpi_request.sampling_interval_s = 5 # MonitorKpi(monitor_kpi_request) def main(): global terminate, logger Loading @@ -42,6 +77,8 @@ def main(): grpc_service = MonitoringService(port=service_port, max_workers=max_workers, grace_period=grace_period) grpc_service.start() # start_monitoring() # Wait for Ctrl+C or termination signal while not terminate.wait(timeout=0.1): pass Loading src/monitoring/tests/test_unitary.py +4 −1 Original line number Diff line number Diff line Loading @@ -263,4 +263,7 @@ def test_listen_events(monitoring_client: MonitoringClient, populate('localhost', GRPC_PORT_CONTEXT) # place this call in the appropriate line, according to your tests events_collector.listen_events() kpi_id_list = events_collector.listen_events() assert bool(kpi_id_list) == True Loading
src/monitoring/client/monitoring_client.py +4 −10 Original line number Diff line number Diff line Loading @@ -21,19 +21,19 @@ class MonitoringClient: LOGGER.info('CreateKpi: {}'.format(request)) response = self.server.CreateKpi(request) LOGGER.info('CreateKpi result: {}'.format(response)) return monitoring_pb2.KpiId() return response def MonitorKpi(self, request): LOGGER.info('MonitorKpi: {}'.format(request)) response = self.server.MonitorKpi(request) LOGGER.info('MonitorKpi result: {}'.format(response)) return context_pb2.Empty() return response def IncludeKpi(self, request): LOGGER.info('IncludeKpi: {}'.format(request)) response = self.server.IncludeKpi(request) LOGGER.info('IncludeKpi result: {}'.format(response)) return context_pb2.Empty() return response def GetStreamKpi(self, request): LOGGER.info('GetStreamKpi: {}'.format(request)) Loading @@ -51,13 +51,7 @@ class MonitoringClient: LOGGER.info('GetKpiDescriptor: {}'.format(request)) response = self.server.GetKpiDescriptor(request) LOGGER.info('GetKpiDescriptor result: {}'.format(response)) return monitoring_pb2.KpiDescriptor() def ListenEvents(self, ): LOGGER.info('ListenEvents: {}'.format()) response = self.server.ListenEvents() LOGGER.info('ListenEvents result: {}'.format(response)) return monitoring_pb2.KpiDescriptor() return response if __name__ == '__main__': # get port Loading
src/monitoring/service/EventTools.py +21 −19 Original line number Diff line number Diff line Loading @@ -47,16 +47,18 @@ class EventsDeviceCollector: def listen_events(self): LOGGER.info('getting Kpi by KpiID') qsize = self._events_queue.qsize() try: kpi_id_list = [] for i in range(self._events_queue.qsize()): if qsize > 0: for i in range(qsize): print("Queue size: "+str(qsize)) event = self.get_event(block=True) if event.event.event_type == EventTypeEnum.EVENTTYPE_CREATE: device = self._context_client.GetDevice(event.device_id) print("Endpoints value: " + str(len(device.device_endpoints))) for j,end_point in enumerate(device.device_endpoints): # for k,rule in enumerate(device.device_config.config_rules): kpi_descriptor = monitoring_pb2.KpiDescriptor() Loading
src/monitoring/service/SqliteTools.py +1 −1 Original line number Diff line number Diff line Loading @@ -16,7 +16,7 @@ class SQLite(): 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)) c.execute("SELECT kpi_id FROM KPI WHERE device_id is ? AND kpi_sample_type is ? AND endpoint_id is ?",(device_id,kpi_sample_type,endpoint_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)) Loading
src/monitoring/service/__main__.py +38 −1 Original line number Diff line number Diff line import logging, os, signal, sys, threading import time from context.client.ContextClient import ContextClient from monitoring.Config import GRPC_SERVICE_PORT, GRPC_MAX_WORKERS, GRPC_GRACE_PERIOD, LOG_LEVEL, METRICS_PORT from common.logger import getJSONLogger from monitoring.client.monitoring_client import MonitoringClient from monitoring.proto import monitoring_pb2 from monitoring.service.EventTools import EventsDeviceCollector from monitoring.service.MonitoringService import MonitoringService LOGGER = getJSONLogger('monitoringservice-server') Loading @@ -16,9 +20,40 @@ logger = None def signal_handler(signal, frame): global terminate, logger logger.warning('Terminate signal received') LOGGER.warning('Terminate signal received') terminate.set() def start_monitoring(): LOGGER.info('Start Monitoring...') context_client_grpc = ContextClient(address='localhost', port='2020') monitoring_client = MonitoringClient(server='localhost', port='7070') # instantiate the client while True: if terminate.is_set(): LOGGER.warning("Stopping execution...") break # Start Listening Events events_collector = EventsDeviceCollector(context_client_grpc, monitoring_client) events_collector.start() list_new_kpi_ids = events_collector.listen_events() # Monitor Kpis if bool(list_new_kpi_ids): for kpi_id in list_new_kpi_ids: # Create Monitor Kpi Requests monitor_kpi_request = monitoring_pb2.MonitorKpiRequest() monitor_kpi_request.kpi_id.CopyFrom(kpi_id) monitor_kpi_request.sampling_duration_s = 120 monitor_kpi_request.sampling_interval_s = 5 # MonitorKpi(monitor_kpi_request) def main(): global terminate, logger Loading @@ -42,6 +77,8 @@ def main(): grpc_service = MonitoringService(port=service_port, max_workers=max_workers, grace_period=grace_period) grpc_service.start() # start_monitoring() # Wait for Ctrl+C or termination signal while not terminate.wait(timeout=0.1): pass Loading
src/monitoring/tests/test_unitary.py +4 −1 Original line number Diff line number Diff line Loading @@ -263,4 +263,7 @@ def test_listen_events(monitoring_client: MonitoringClient, populate('localhost', GRPC_PORT_CONTEXT) # place this call in the appropriate line, according to your tests events_collector.listen_events() kpi_id_list = events_collector.listen_events() assert bool(kpi_id_list) == True