Loading src/simap_connector/service/SimapConnectorServiceServicerImpl.py +20 −12 Original line number Diff line number Diff line Loading @@ -21,6 +21,7 @@ from common.proto.simap_connector_pb2_grpc import SimapConnectorServiceServicer from common.tools.rest_conf.client.RestConfClient import RestConfClient from common.method_wrappers.Decorator import MetricsPool, safe_and_metered_rpc_method from device.client.DeviceClient import DeviceClient from simap_connector.service.telemetry.worker._Worker import WorkerTypeEnum from .database.Subscription import subscription_get, subscription_set, subscription_delete from .database.SubSubscription import ( sub_subscription_list, sub_subscription_set, sub_subscription_delete Loading Loading @@ -77,8 +78,8 @@ class SimapConnectorServiceServicerImpl(SimapConnectorServiceServicer): sup_link_xpath_filters.append(sup_link_xpath_filter) if controller_id is None: collector_name = 'SIMAP:{:s}:{:s}'.format( str(supporting_link.network_id), str(supporting_link.link_id) collector_name = '{:d}:SIMAP:{:s}:{:s}'.format( parent_subscription_id, str(supporting_link.network_id), str(supporting_link.link_id) ) target_uri = sup_link_xpath_filter underlay_subscription_id = 0 Loading @@ -86,8 +87,8 @@ class SimapConnectorServiceServicerImpl(SimapConnectorServiceServicer): underlay_sub_id = establish_underlay_subscription( device_client, controller_id, sup_link_xpath_filter, period ) collector_name = '{:s}:{:s}'.format( controller_id, str(underlay_sub_id.subscription_id) collector_name = '{:d}:{:s}:{:s}'.format( parent_subscription_id, controller_id, str(underlay_sub_id.subscription_id) ) target_uri = underlay_sub_id.subscription_uri underlay_subscription_id = underlay_sub_id.subscription_id Loading @@ -103,7 +104,8 @@ class SimapConnectorServiceServicerImpl(SimapConnectorServiceServicer): sub_request.period = period sub_subscription_set( self._db_engine, parent_subscription_uuid, controller_id, datastore, sup_link_xpath_filter, period, underlay_subscription_id, target_uri sup_link_xpath_filter, period, underlay_subscription_id, target_uri, collector_name ) topic = 'subscription.{:d}'.format(parent_subscription_id) Loading @@ -123,20 +125,26 @@ class SimapConnectorServiceServicerImpl(SimapConnectorServiceServicer): subscription = subscription_get(self._db_engine, parent_subscription_id) if subscription is None: return Empty() # TODO: desactivate subscription aggregator and collectors topic = 'subscription.{:d}'.format(parent_subscription_id) delete_kafka_topic(topic) parent_subscription_uuid = subscription['subscription_uuid'] aggregator_name = str(parent_subscription_id) self._telemetry_pool.stop_worker(WorkerTypeEnum.AGGREGATOR, aggregator_name) device_client = DeviceClient() parent_subscription_uuid = subscription['subscription_uuid'] sub_subscriptions = sub_subscription_list(self._db_engine, parent_subscription_uuid) for sub_subscription in sub_subscriptions: sub_subscription_id = sub_subscription['sub_subscription_id'] controller_id = sub_subscription['controller_uuid' ] collector_name = sub_subscription['collector_name' ] self._telemetry_pool.stop_worker(WorkerTypeEnum.COLLECTOR, collector_name) if controller_id is not None and len(controller_id) > 0: delete_underlay_subscription(device_client, controller_id, sub_subscription_id) sub_subscription_delete(self._db_engine, parent_subscription_uuid, sub_subscription_id) topic = 'subscription.{:d}'.format(parent_subscription_id) delete_kafka_topic(topic) subscription_delete(self._db_engine, parent_subscription_id) return Empty() src/simap_connector/service/database/SubSubscription.py +4 −1 Original line number Diff line number Diff line Loading @@ -57,7 +57,8 @@ def sub_subscription_get( def sub_subscription_set( db_engine : Engine, parent_subscription_uuid : str, controller_uuid : str, datastore : str, xpath_filter : str, period : float, sub_subscription_id : int, sub_subscription_uri : str xpath_filter : str, period : float, sub_subscription_id : int, sub_subscription_uri : str, collector_name : str ) -> str: now = datetime.datetime.now(datetime.timezone.utc) if controller_uuid is None: controller_uuid = '' Loading @@ -69,6 +70,7 @@ def sub_subscription_set( 'period' : period, 'sub_subscription_id' : sub_subscription_id, 'sub_subscription_uri': sub_subscription_uri, 'collector_name' : collector_name, 'created_at' : now, 'updated_at' : now, } Loading @@ -87,6 +89,7 @@ def sub_subscription_set( period = stmt.excluded.period, sub_subscription_id = stmt.excluded.sub_subscription_id, sub_subscription_uri = stmt.excluded.sub_subscription_uri, collector_name = stmt.excluded.collector_name, updated_at = stmt.excluded.updated_at, ) ) Loading src/simap_connector/service/database/models/SubSubscriptionModel.py +2 −0 Original line number Diff line number Diff line Loading @@ -30,6 +30,7 @@ class SubSubscriptionModel(_Base): period = Column(Float, nullable=False, unique=False) sub_subscription_id = Column(BigInteger, nullable=False, unique=False) sub_subscription_uri = Column(String, nullable=False, unique=False) collector_name = Column(String, nullable=False, unique=False) created_at = Column(DateTime, nullable=False) updated_at = Column(DateTime, nullable=False) Loading @@ -45,4 +46,5 @@ class SubSubscriptionModel(_Base): 'period' : self.period, 'sub_subscription_id' : self.sub_subscription_id, 'sub_subscription_uri' : self.sub_subscription_uri, 'collector_name' : self.collector_name, } src/simap_connector/service/simap_updater/SimapUpdater.py +4 −4 Original line number Diff line number Diff line Loading @@ -533,9 +533,9 @@ class EventDispatcher(BaseEventDispatcher): #dom_link.update(src_dev_name, src_ep_name, dst_dev_name, dst_ep_name) resources = Resources() sampling_interval = 1.0 self._telemetry_pool.start_synthesizer(domain_name, resources, sampling_interval) #resources = Resources() #sampling_interval = 1.0 #self._telemetry_pool.start_synthesizer(domain_name, resources, sampling_interval) return True Loading Loading @@ -631,7 +631,7 @@ class EventDispatcher(BaseEventDispatcher): #self._object_cache.delete(CachedEntities.SERVICE, service_uuid) #self._object_cache.delete(CachedEntities.SERVICE, service_name) self._telemetry_pool.stop_worker(WorkerTypeEnum.SYNTHESIZER, domain_name) #self._telemetry_pool.stop_worker(WorkerTypeEnum.SYNTHESIZER, domain_name) MSG = 'Logical Link Removed for Service: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(service_event))) Loading Loading
src/simap_connector/service/SimapConnectorServiceServicerImpl.py +20 −12 Original line number Diff line number Diff line Loading @@ -21,6 +21,7 @@ from common.proto.simap_connector_pb2_grpc import SimapConnectorServiceServicer from common.tools.rest_conf.client.RestConfClient import RestConfClient from common.method_wrappers.Decorator import MetricsPool, safe_and_metered_rpc_method from device.client.DeviceClient import DeviceClient from simap_connector.service.telemetry.worker._Worker import WorkerTypeEnum from .database.Subscription import subscription_get, subscription_set, subscription_delete from .database.SubSubscription import ( sub_subscription_list, sub_subscription_set, sub_subscription_delete Loading Loading @@ -77,8 +78,8 @@ class SimapConnectorServiceServicerImpl(SimapConnectorServiceServicer): sup_link_xpath_filters.append(sup_link_xpath_filter) if controller_id is None: collector_name = 'SIMAP:{:s}:{:s}'.format( str(supporting_link.network_id), str(supporting_link.link_id) collector_name = '{:d}:SIMAP:{:s}:{:s}'.format( parent_subscription_id, str(supporting_link.network_id), str(supporting_link.link_id) ) target_uri = sup_link_xpath_filter underlay_subscription_id = 0 Loading @@ -86,8 +87,8 @@ class SimapConnectorServiceServicerImpl(SimapConnectorServiceServicer): underlay_sub_id = establish_underlay_subscription( device_client, controller_id, sup_link_xpath_filter, period ) collector_name = '{:s}:{:s}'.format( controller_id, str(underlay_sub_id.subscription_id) collector_name = '{:d}:{:s}:{:s}'.format( parent_subscription_id, controller_id, str(underlay_sub_id.subscription_id) ) target_uri = underlay_sub_id.subscription_uri underlay_subscription_id = underlay_sub_id.subscription_id Loading @@ -103,7 +104,8 @@ class SimapConnectorServiceServicerImpl(SimapConnectorServiceServicer): sub_request.period = period sub_subscription_set( self._db_engine, parent_subscription_uuid, controller_id, datastore, sup_link_xpath_filter, period, underlay_subscription_id, target_uri sup_link_xpath_filter, period, underlay_subscription_id, target_uri, collector_name ) topic = 'subscription.{:d}'.format(parent_subscription_id) Loading @@ -123,20 +125,26 @@ class SimapConnectorServiceServicerImpl(SimapConnectorServiceServicer): subscription = subscription_get(self._db_engine, parent_subscription_id) if subscription is None: return Empty() # TODO: desactivate subscription aggregator and collectors topic = 'subscription.{:d}'.format(parent_subscription_id) delete_kafka_topic(topic) parent_subscription_uuid = subscription['subscription_uuid'] aggregator_name = str(parent_subscription_id) self._telemetry_pool.stop_worker(WorkerTypeEnum.AGGREGATOR, aggregator_name) device_client = DeviceClient() parent_subscription_uuid = subscription['subscription_uuid'] sub_subscriptions = sub_subscription_list(self._db_engine, parent_subscription_uuid) for sub_subscription in sub_subscriptions: sub_subscription_id = sub_subscription['sub_subscription_id'] controller_id = sub_subscription['controller_uuid' ] collector_name = sub_subscription['collector_name' ] self._telemetry_pool.stop_worker(WorkerTypeEnum.COLLECTOR, collector_name) if controller_id is not None and len(controller_id) > 0: delete_underlay_subscription(device_client, controller_id, sub_subscription_id) sub_subscription_delete(self._db_engine, parent_subscription_uuid, sub_subscription_id) topic = 'subscription.{:d}'.format(parent_subscription_id) delete_kafka_topic(topic) subscription_delete(self._db_engine, parent_subscription_id) return Empty()
src/simap_connector/service/database/SubSubscription.py +4 −1 Original line number Diff line number Diff line Loading @@ -57,7 +57,8 @@ def sub_subscription_get( def sub_subscription_set( db_engine : Engine, parent_subscription_uuid : str, controller_uuid : str, datastore : str, xpath_filter : str, period : float, sub_subscription_id : int, sub_subscription_uri : str xpath_filter : str, period : float, sub_subscription_id : int, sub_subscription_uri : str, collector_name : str ) -> str: now = datetime.datetime.now(datetime.timezone.utc) if controller_uuid is None: controller_uuid = '' Loading @@ -69,6 +70,7 @@ def sub_subscription_set( 'period' : period, 'sub_subscription_id' : sub_subscription_id, 'sub_subscription_uri': sub_subscription_uri, 'collector_name' : collector_name, 'created_at' : now, 'updated_at' : now, } Loading @@ -87,6 +89,7 @@ def sub_subscription_set( period = stmt.excluded.period, sub_subscription_id = stmt.excluded.sub_subscription_id, sub_subscription_uri = stmt.excluded.sub_subscription_uri, collector_name = stmt.excluded.collector_name, updated_at = stmt.excluded.updated_at, ) ) Loading
src/simap_connector/service/database/models/SubSubscriptionModel.py +2 −0 Original line number Diff line number Diff line Loading @@ -30,6 +30,7 @@ class SubSubscriptionModel(_Base): period = Column(Float, nullable=False, unique=False) sub_subscription_id = Column(BigInteger, nullable=False, unique=False) sub_subscription_uri = Column(String, nullable=False, unique=False) collector_name = Column(String, nullable=False, unique=False) created_at = Column(DateTime, nullable=False) updated_at = Column(DateTime, nullable=False) Loading @@ -45,4 +46,5 @@ class SubSubscriptionModel(_Base): 'period' : self.period, 'sub_subscription_id' : self.sub_subscription_id, 'sub_subscription_uri' : self.sub_subscription_uri, 'collector_name' : self.collector_name, }
src/simap_connector/service/simap_updater/SimapUpdater.py +4 −4 Original line number Diff line number Diff line Loading @@ -533,9 +533,9 @@ class EventDispatcher(BaseEventDispatcher): #dom_link.update(src_dev_name, src_ep_name, dst_dev_name, dst_ep_name) resources = Resources() sampling_interval = 1.0 self._telemetry_pool.start_synthesizer(domain_name, resources, sampling_interval) #resources = Resources() #sampling_interval = 1.0 #self._telemetry_pool.start_synthesizer(domain_name, resources, sampling_interval) return True Loading Loading @@ -631,7 +631,7 @@ class EventDispatcher(BaseEventDispatcher): #self._object_cache.delete(CachedEntities.SERVICE, service_uuid) #self._object_cache.delete(CachedEntities.SERVICE, service_name) self._telemetry_pool.stop_worker(WorkerTypeEnum.SYNTHESIZER, domain_name) #self._telemetry_pool.stop_worker(WorkerTypeEnum.SYNTHESIZER, domain_name) MSG = 'Logical Link Removed for Service: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(service_event))) Loading