Loading src/simap_connector/service/simap_updater/SimapUpdater.py +38 −25 Original line number Diff line number Diff line Loading @@ -17,7 +17,7 @@ import logging, queue, threading, uuid from typing import Any, Optional, Set from common.Constants import DEFAULT_TOPOLOGY_NAME from common.DeviceTypes import DeviceTypeEnum from common.proto.context_pb2 import ContextEvent, DeviceEvent, Empty, LinkEvent, ServiceEvent, TopologyEvent from common.proto.context_pb2 import ContextEvent, DeviceEvent, Empty, LinkEvent, ServiceEvent, SliceEvent, TopologyEvent from common.tools.grpc.BaseEventCollector import BaseEventCollector from common.tools.grpc.BaseEventDispatcher import BaseEventDispatcher from common.tools.grpc.Tools import grpc_message_to_json_string Loading Loading @@ -83,7 +83,12 @@ class EventDispatcher(BaseEventDispatcher): LOGGER.debug(MSG.format(grpc_message_to_json_string(context_event))) def _dispatch_topology_set(self, topology_event : TopologyEvent) -> None: def dispatch_slice(self, slice_event : SliceEvent) -> None: MSG = 'Skipping Slice Event: {:s}' LOGGER.debug(MSG.format(grpc_message_to_json_string(slice_event))) def _dispatch_topology_set(self, topology_event : TopologyEvent) -> bool: MSG = 'Processing Topology Event: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(topology_event))) Loading @@ -100,17 +105,18 @@ class EventDispatcher(BaseEventDispatcher): self._simap_client.network(topology_name).update( supporting_network_ids=supporting_network_ids ) return True def dispatch_topology_create(self, topology_event : TopologyEvent) -> None: self._dispatch_topology_set(topology_event) if not self._dispatch_topology_set(topology_event): return MSG = 'Topology Create: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(topology_event))) def dispatch_topology_update(self, topology_event : TopologyEvent) -> None: self._dispatch_topology_set(topology_event) if not self._dispatch_topology_set(topology_event): return MSG = 'Topology Updated: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(topology_event))) Loading @@ -132,7 +138,7 @@ class EventDispatcher(BaseEventDispatcher): LOGGER.info(MSG.format(grpc_message_to_json_string(topology_event))) def _dispatch_device_set(self, device_event : DeviceEvent) -> None: def _dispatch_device_set(self, device_event : DeviceEvent) -> bool: MSG = 'Processing Device Event: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(device_event))) Loading @@ -149,7 +155,7 @@ class EventDispatcher(BaseEventDispatcher): str_device_event = grpc_message_to_json_string(device_event) str_device = grpc_message_to_json_string(device) LOGGER.warning(MSG.format(str_device_event, str_device)) return return False #device_controller_uuid = device.controller_id.device_uuid.uuid #if len(device_controller_uuid) > 0: Loading @@ -170,7 +176,7 @@ class EventDispatcher(BaseEventDispatcher): str_device_event = grpc_message_to_json_string(device_event) str_device = grpc_message_to_json_string(device) LOGGER.warning(MSG.format(str_device_event, str_device)) return return False topology = self._object_cache.get(CachedEntities.TOPOLOGY, topology_uuid) topology_name = topology.name Loading @@ -186,17 +192,18 @@ class EventDispatcher(BaseEventDispatcher): te_device.termination_point(endpoint_name).update() #self._remove_skipped_device(device) return True def dispatch_device_create(self, device_event : DeviceEvent) -> None: self._dispatch_device_set(device_event) if not self._dispatch_device_set(device_event): return MSG = 'Device Created: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(device_event))) def dispatch_device_update(self, device_event : DeviceEvent) -> None: self._dispatch_device_set(device_event) if not self._dispatch_device_set(device_event): return MSG = 'Device Updated: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(device_event))) Loading Loading @@ -264,7 +271,7 @@ class EventDispatcher(BaseEventDispatcher): LOGGER.info(MSG.format(grpc_message_to_json_string(device_event))) def _dispatch_link_set(self, link_event : LinkEvent) -> None: def _dispatch_link_set(self, link_event : LinkEvent) -> bool: MSG = 'Processing Link Event: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(link_event))) Loading @@ -278,7 +285,7 @@ class EventDispatcher(BaseEventDispatcher): str_link_event = grpc_message_to_json_string(link_event) str_link = grpc_message_to_json_string(link) LOGGER.warning(MSG.format(str_link_event, str_link)) return return False topology = self._object_cache.get(CachedEntities.TOPOLOGY, topology_uuid) topology_name = topology.name Loading @@ -298,7 +305,7 @@ class EventDispatcher(BaseEventDispatcher): str_link_event = grpc_message_to_json_string(link_event) str_link = grpc_message_to_json_string(link) LOGGER.warning(MSG.format(str_link_event, str_link)) return return False # Skip links that connect to devices previously marked as skipped src_uuid = src_device.device_id.device_uuid.uuid Loading @@ -311,7 +318,7 @@ class EventDispatcher(BaseEventDispatcher): str_link_event = grpc_message_to_json_string(link_event) str_link = grpc_message_to_json_string(link) LOGGER.warning(MSG.format(str_link_event, str_link)) return return False try: if src_device is None: Loading @@ -332,17 +339,18 @@ class EventDispatcher(BaseEventDispatcher): te_link = te_topo.link(link_name) te_link.update(src_device.name, src_endpoint.name, dst_device.name, dst_endpoint.name) return True def dispatch_link_create(self, link_event : LinkEvent) -> None: self._dispatch_link_set(link_event) if not self._dispatch_link_set(link_event): return MSG = 'Link Created: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(link_event))) def dispatch_link_update(self, link_event : LinkEvent) -> None: self._dispatch_link_set(link_event) if not self._dispatch_link_set(link_event): return MSG = 'Link Updated: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(link_event))) Loading Loading @@ -400,7 +408,7 @@ class EventDispatcher(BaseEventDispatcher): LOGGER.info(MSG.format(grpc_message_to_json_string(link_event))) def _dispatch_service_set(self, service_event : ServiceEvent) -> None: def _dispatch_service_set(self, service_event : ServiceEvent) -> bool: MSG = 'Processing Service Event: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(service_event))) Loading @@ -411,7 +419,11 @@ class EventDispatcher(BaseEventDispatcher): try: uuid.UUID(hex=service_name) # skip it if properly parsed, means it is a service with a UUID-based name, i.e., a sub-service return MSG = 'ServiceEvent({:s}) skipped, it is a subservice: {:s}' str_service_event = grpc_message_to_json_string(service_event) str_service = grpc_message_to_json_string(service) LOGGER.warning(MSG.format(str_service_event, str_service)) return False except: # pylint: disable=bare-except pass Loading @@ -422,14 +434,14 @@ class EventDispatcher(BaseEventDispatcher): str_service_event = grpc_message_to_json_string(service_event) str_service = grpc_message_to_json_string(service) LOGGER.warning(MSG.format(str_service_event, str_service)) return return False if len(endpoint_uuids) < 2: MSG = 'ServiceEvent({:s}) skipped, not enough endpoint_ids to compose link: {:s}' str_service_event = grpc_message_to_json_string(service_event) str_service = grpc_message_to_json_string(service) LOGGER.warning(MSG.format(str_service_event, str_service)) return return False topologies = self._object_cache.get_all(CachedEntities.TOPOLOGY, fresh=False) topology_names = {t.name for t in topologies} Loading @@ -438,9 +450,9 @@ class EventDispatcher(BaseEventDispatcher): MSG = 'ServiceEvent({:s}) skipped, unable to identify on which topology to insert it' str_service_event = grpc_message_to_json_string(service_event) LOGGER.warning(MSG.format(str_service_event)) return domain_name = topology_names.pop() # trans-pkt/agg-net/e2e-net return False domain_name = topology_names.pop() # trans-pkt/agg-net/e2e-net domain_topo = self._simap_client.network(domain_name) domain_topo.update() Loading Loading @@ -490,17 +502,18 @@ class EventDispatcher(BaseEventDispatcher): ) dom_link = domain_topo.link(link_name) dom_link.update(src_dev_name, src_ep_name, dst_dev_name, dst_ep_name) return True def dispatch_service_created(self, service_event : ServiceEvent) -> None: self._dispatch_service_set(service_event) def dispatch_service_create(self, service_event : ServiceEvent) -> None: if not self._dispatch_service_set(service_event): return MSG = 'Logical Link Created for Service: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(service_event))) def dispatch_service_update(self, service_event : ServiceEvent) -> None: self._dispatch_service_set(service_event) if not self._dispatch_service_set(service_event): return MSG = 'Logical Link Updated for Service: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(service_event))) Loading Loading @@ -530,8 +543,8 @@ class EventDispatcher(BaseEventDispatcher): str_service_event = grpc_message_to_json_string(service_event) LOGGER.warning(MSG.format(str_service_event)) return domain_name = topology_names.pop() # trans-pkt/agg-net/e2e-net domain_name = topology_names.pop() # trans-pkt/agg-net/e2e-net domain_topo = self._simap_client.network(domain_name) domain_topo.update() Loading Loading
src/simap_connector/service/simap_updater/SimapUpdater.py +38 −25 Original line number Diff line number Diff line Loading @@ -17,7 +17,7 @@ import logging, queue, threading, uuid from typing import Any, Optional, Set from common.Constants import DEFAULT_TOPOLOGY_NAME from common.DeviceTypes import DeviceTypeEnum from common.proto.context_pb2 import ContextEvent, DeviceEvent, Empty, LinkEvent, ServiceEvent, TopologyEvent from common.proto.context_pb2 import ContextEvent, DeviceEvent, Empty, LinkEvent, ServiceEvent, SliceEvent, TopologyEvent from common.tools.grpc.BaseEventCollector import BaseEventCollector from common.tools.grpc.BaseEventDispatcher import BaseEventDispatcher from common.tools.grpc.Tools import grpc_message_to_json_string Loading Loading @@ -83,7 +83,12 @@ class EventDispatcher(BaseEventDispatcher): LOGGER.debug(MSG.format(grpc_message_to_json_string(context_event))) def _dispatch_topology_set(self, topology_event : TopologyEvent) -> None: def dispatch_slice(self, slice_event : SliceEvent) -> None: MSG = 'Skipping Slice Event: {:s}' LOGGER.debug(MSG.format(grpc_message_to_json_string(slice_event))) def _dispatch_topology_set(self, topology_event : TopologyEvent) -> bool: MSG = 'Processing Topology Event: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(topology_event))) Loading @@ -100,17 +105,18 @@ class EventDispatcher(BaseEventDispatcher): self._simap_client.network(topology_name).update( supporting_network_ids=supporting_network_ids ) return True def dispatch_topology_create(self, topology_event : TopologyEvent) -> None: self._dispatch_topology_set(topology_event) if not self._dispatch_topology_set(topology_event): return MSG = 'Topology Create: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(topology_event))) def dispatch_topology_update(self, topology_event : TopologyEvent) -> None: self._dispatch_topology_set(topology_event) if not self._dispatch_topology_set(topology_event): return MSG = 'Topology Updated: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(topology_event))) Loading @@ -132,7 +138,7 @@ class EventDispatcher(BaseEventDispatcher): LOGGER.info(MSG.format(grpc_message_to_json_string(topology_event))) def _dispatch_device_set(self, device_event : DeviceEvent) -> None: def _dispatch_device_set(self, device_event : DeviceEvent) -> bool: MSG = 'Processing Device Event: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(device_event))) Loading @@ -149,7 +155,7 @@ class EventDispatcher(BaseEventDispatcher): str_device_event = grpc_message_to_json_string(device_event) str_device = grpc_message_to_json_string(device) LOGGER.warning(MSG.format(str_device_event, str_device)) return return False #device_controller_uuid = device.controller_id.device_uuid.uuid #if len(device_controller_uuid) > 0: Loading @@ -170,7 +176,7 @@ class EventDispatcher(BaseEventDispatcher): str_device_event = grpc_message_to_json_string(device_event) str_device = grpc_message_to_json_string(device) LOGGER.warning(MSG.format(str_device_event, str_device)) return return False topology = self._object_cache.get(CachedEntities.TOPOLOGY, topology_uuid) topology_name = topology.name Loading @@ -186,17 +192,18 @@ class EventDispatcher(BaseEventDispatcher): te_device.termination_point(endpoint_name).update() #self._remove_skipped_device(device) return True def dispatch_device_create(self, device_event : DeviceEvent) -> None: self._dispatch_device_set(device_event) if not self._dispatch_device_set(device_event): return MSG = 'Device Created: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(device_event))) def dispatch_device_update(self, device_event : DeviceEvent) -> None: self._dispatch_device_set(device_event) if not self._dispatch_device_set(device_event): return MSG = 'Device Updated: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(device_event))) Loading Loading @@ -264,7 +271,7 @@ class EventDispatcher(BaseEventDispatcher): LOGGER.info(MSG.format(grpc_message_to_json_string(device_event))) def _dispatch_link_set(self, link_event : LinkEvent) -> None: def _dispatch_link_set(self, link_event : LinkEvent) -> bool: MSG = 'Processing Link Event: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(link_event))) Loading @@ -278,7 +285,7 @@ class EventDispatcher(BaseEventDispatcher): str_link_event = grpc_message_to_json_string(link_event) str_link = grpc_message_to_json_string(link) LOGGER.warning(MSG.format(str_link_event, str_link)) return return False topology = self._object_cache.get(CachedEntities.TOPOLOGY, topology_uuid) topology_name = topology.name Loading @@ -298,7 +305,7 @@ class EventDispatcher(BaseEventDispatcher): str_link_event = grpc_message_to_json_string(link_event) str_link = grpc_message_to_json_string(link) LOGGER.warning(MSG.format(str_link_event, str_link)) return return False # Skip links that connect to devices previously marked as skipped src_uuid = src_device.device_id.device_uuid.uuid Loading @@ -311,7 +318,7 @@ class EventDispatcher(BaseEventDispatcher): str_link_event = grpc_message_to_json_string(link_event) str_link = grpc_message_to_json_string(link) LOGGER.warning(MSG.format(str_link_event, str_link)) return return False try: if src_device is None: Loading @@ -332,17 +339,18 @@ class EventDispatcher(BaseEventDispatcher): te_link = te_topo.link(link_name) te_link.update(src_device.name, src_endpoint.name, dst_device.name, dst_endpoint.name) return True def dispatch_link_create(self, link_event : LinkEvent) -> None: self._dispatch_link_set(link_event) if not self._dispatch_link_set(link_event): return MSG = 'Link Created: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(link_event))) def dispatch_link_update(self, link_event : LinkEvent) -> None: self._dispatch_link_set(link_event) if not self._dispatch_link_set(link_event): return MSG = 'Link Updated: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(link_event))) Loading Loading @@ -400,7 +408,7 @@ class EventDispatcher(BaseEventDispatcher): LOGGER.info(MSG.format(grpc_message_to_json_string(link_event))) def _dispatch_service_set(self, service_event : ServiceEvent) -> None: def _dispatch_service_set(self, service_event : ServiceEvent) -> bool: MSG = 'Processing Service Event: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(service_event))) Loading @@ -411,7 +419,11 @@ class EventDispatcher(BaseEventDispatcher): try: uuid.UUID(hex=service_name) # skip it if properly parsed, means it is a service with a UUID-based name, i.e., a sub-service return MSG = 'ServiceEvent({:s}) skipped, it is a subservice: {:s}' str_service_event = grpc_message_to_json_string(service_event) str_service = grpc_message_to_json_string(service) LOGGER.warning(MSG.format(str_service_event, str_service)) return False except: # pylint: disable=bare-except pass Loading @@ -422,14 +434,14 @@ class EventDispatcher(BaseEventDispatcher): str_service_event = grpc_message_to_json_string(service_event) str_service = grpc_message_to_json_string(service) LOGGER.warning(MSG.format(str_service_event, str_service)) return return False if len(endpoint_uuids) < 2: MSG = 'ServiceEvent({:s}) skipped, not enough endpoint_ids to compose link: {:s}' str_service_event = grpc_message_to_json_string(service_event) str_service = grpc_message_to_json_string(service) LOGGER.warning(MSG.format(str_service_event, str_service)) return return False topologies = self._object_cache.get_all(CachedEntities.TOPOLOGY, fresh=False) topology_names = {t.name for t in topologies} Loading @@ -438,9 +450,9 @@ class EventDispatcher(BaseEventDispatcher): MSG = 'ServiceEvent({:s}) skipped, unable to identify on which topology to insert it' str_service_event = grpc_message_to_json_string(service_event) LOGGER.warning(MSG.format(str_service_event)) return domain_name = topology_names.pop() # trans-pkt/agg-net/e2e-net return False domain_name = topology_names.pop() # trans-pkt/agg-net/e2e-net domain_topo = self._simap_client.network(domain_name) domain_topo.update() Loading Loading @@ -490,17 +502,18 @@ class EventDispatcher(BaseEventDispatcher): ) dom_link = domain_topo.link(link_name) dom_link.update(src_dev_name, src_ep_name, dst_dev_name, dst_ep_name) return True def dispatch_service_created(self, service_event : ServiceEvent) -> None: self._dispatch_service_set(service_event) def dispatch_service_create(self, service_event : ServiceEvent) -> None: if not self._dispatch_service_set(service_event): return MSG = 'Logical Link Created for Service: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(service_event))) def dispatch_service_update(self, service_event : ServiceEvent) -> None: self._dispatch_service_set(service_event) if not self._dispatch_service_set(service_event): return MSG = 'Logical Link Updated for Service: {:s}' LOGGER.info(MSG.format(grpc_message_to_json_string(service_event))) Loading Loading @@ -530,8 +543,8 @@ class EventDispatcher(BaseEventDispatcher): str_service_event = grpc_message_to_json_string(service_event) LOGGER.warning(MSG.format(str_service_event)) return domain_name = topology_names.pop() # trans-pkt/agg-net/e2e-net domain_name = topology_names.pop() # trans-pkt/agg-net/e2e-net domain_topo = self._simap_client.network(domain_name) domain_topo.update() Loading