Loading src/service/service/service_handler_api/ServiceHandlerFactory.py +4 −6 Original line number Diff line number Diff line Loading @@ -100,9 +100,6 @@ class ServiceHandlerFactory: candidate_service_handler_classes.items(), key=operator.itemgetter(1), reverse=True) return candidate_service_handler_classes[0][0] def get_device_supported_drivers(device : Device) -> Set[int]: return {device_driver for device_driver in device.device_drivers} def get_common_device_drivers(drivers_per_device : List[Set[int]]) -> Set[int]: common_device_drivers = None for device_drivers in drivers_per_device: Loading @@ -114,15 +111,16 @@ def get_common_device_drivers(drivers_per_device : List[Set[int]]) -> Set[int]: return common_device_drivers def get_service_handler_class( service_handler_factory : ServiceHandlerFactory, service : Service, connection_devices : Dict[str, Device] service_handler_factory : ServiceHandlerFactory, service : Service, device_and_drivers: Dict[str, Tuple[Device, Set[int]]] ) -> Optional['_ServiceHandler']: str_service_key = grpc_message_to_json_string(service.service_id) # Assume all devices involved in the service's connection must support at least one driver in common common_device_drivers = get_common_device_drivers([ get_device_supported_drivers(device) for device in connection_devices.values() device_drivers for _,device_drivers in device_and_drivers.values() ]) filter_fields = { Loading src/service/service/task_scheduler/TaskExecutor.py +54 −11 Original line number Diff line number Diff line Loading @@ -14,7 +14,7 @@ import json, logging from enum import Enum from typing import TYPE_CHECKING, Any, Dict, List, Optional, Tuple, Union from typing import TYPE_CHECKING, Any, Dict, List, Optional, Set, Tuple, Union from common.DeviceTypes import DeviceTypeEnum from common.method_wrappers.ServiceExceptions import NotFoundException from common.proto.qkd_app_pb2 import QKDAppStatusEnum Loading Loading @@ -333,6 +333,38 @@ class TaskExecutor: # return devices return devices def get_device_type_drivers_for_connection( self, connection : Connection ) -> Dict[DeviceTypeEnum, Dict[str, Tuple[Device, Set[int]]]]: devices : Dict[DeviceTypeEnum, Dict[str, Tuple[Device, Set[int]]]] = dict() for endpoint_id in connection.path_hops_endpoint_ids: device = self.get_device(endpoint_id.device_id) device_uuid = endpoint_id.device_id.device_uuid.uuid if device is None: raise Exception('Device({:s}) not found'.format(str(device_uuid))) controller = self.get_device_controller(device) if controller is None: device_type = DeviceTypeEnum._value2member_map_[device.device_type] device_drivers = set(driver for driver in device.device_drivers) devices.setdefault(device_type, dict())[device_uuid] = (device, device_drivers) else: # Controller device types for those underlying path is needed by service handler device_type = DeviceTypeEnum._value2member_map_[controller.device_type] controller_drivers = set(driver for driver in controller.device_drivers) if device_type not in EXPANSION_CONTROLLER_DEVICE_TYPES: devices.setdefault(device_type, dict())[device_uuid] = (device, controller_drivers) else: controller_uuid = controller.device_id.device_uuid.uuid devices.setdefault(device_type, dict())[controller_uuid] = (controller, controller_drivers) LOGGER.debug('[get_devices_from_connection] devices = {:s}'.format(str(devices))) return devices # ----- Service-related methods ------------------------------------------------------------------------------------ def get_service(self, service_id : ServiceId) -> Service: Loading Loading @@ -374,18 +406,22 @@ class TaskExecutor: #LOGGER.debug('connection_device_types_included = {:s}'.format(str(connection_device_types_included))) # ================================================================================================ connection_device_types : Dict[DeviceTypeEnum, Dict[str, Device]] = self.get_devices_from_connection( connection, exclude_managed_by_controller=False ) device_type_to_device_and_drivers : Dict[DeviceTypeEnum, Dict[str, Tuple[Device, Set[int]]]] = \ self.get_device_type_drivers_for_connection(connection) service_handlers : Dict[DeviceTypeEnum, Tuple['_ServiceHandler', Dict[str, Device]]] = dict() # ===== Ryu original test ======================================================================== #for device_type, connection_devices in connection_device_types_excluded.items(): # ================================================================================================ for device_type, connection_devices in connection_device_types.items(): for device_type, device_and_drivers in device_type_to_device_and_drivers.items(): try: service_handler_class = get_service_handler_class( self._service_handler_factory, service, connection_devices self._service_handler_factory, service, device_and_drivers ) connection_devices = { device_uuid : device for device_uuid, (device, _) in device_and_drivers.items() } # ===== Ryu original test ======================================================================== #LOGGER.debug('service_handler_class IN CONNECTION DEVICE TYPE EXCLUDED = {:s}'.format(str(service_handler_class.__name__))) #service_handler = service_handler_class(service, self, **service_handler_settings) Loading @@ -402,11 +438,18 @@ class TaskExecutor: UnsupportedFilterFieldValueException ): dict_connection_devices = { cd_data.name : (cd_uuid, cd_data.name, { cd_data.name : ( cd_uuid, cd_data.name, { (device_driver, DeviceDriverEnum.Name(device_driver)) for device_driver in cd_data.device_drivers }) for cd_uuid,cd_data in connection_devices.items() }, { (device_driver, DeviceDriverEnum.Name(device_driver)) for device_driver in drivers } ) for cd_uuid,(cd_data, drivers) in device_and_drivers.items() } MSG = 'Unable to select service handler. service={:s} connection={:s} connection_devices={:s}' LOGGER.exception(MSG.format( Loading Loading
src/service/service/service_handler_api/ServiceHandlerFactory.py +4 −6 Original line number Diff line number Diff line Loading @@ -100,9 +100,6 @@ class ServiceHandlerFactory: candidate_service_handler_classes.items(), key=operator.itemgetter(1), reverse=True) return candidate_service_handler_classes[0][0] def get_device_supported_drivers(device : Device) -> Set[int]: return {device_driver for device_driver in device.device_drivers} def get_common_device_drivers(drivers_per_device : List[Set[int]]) -> Set[int]: common_device_drivers = None for device_drivers in drivers_per_device: Loading @@ -114,15 +111,16 @@ def get_common_device_drivers(drivers_per_device : List[Set[int]]) -> Set[int]: return common_device_drivers def get_service_handler_class( service_handler_factory : ServiceHandlerFactory, service : Service, connection_devices : Dict[str, Device] service_handler_factory : ServiceHandlerFactory, service : Service, device_and_drivers: Dict[str, Tuple[Device, Set[int]]] ) -> Optional['_ServiceHandler']: str_service_key = grpc_message_to_json_string(service.service_id) # Assume all devices involved in the service's connection must support at least one driver in common common_device_drivers = get_common_device_drivers([ get_device_supported_drivers(device) for device in connection_devices.values() device_drivers for _,device_drivers in device_and_drivers.values() ]) filter_fields = { Loading
src/service/service/task_scheduler/TaskExecutor.py +54 −11 Original line number Diff line number Diff line Loading @@ -14,7 +14,7 @@ import json, logging from enum import Enum from typing import TYPE_CHECKING, Any, Dict, List, Optional, Tuple, Union from typing import TYPE_CHECKING, Any, Dict, List, Optional, Set, Tuple, Union from common.DeviceTypes import DeviceTypeEnum from common.method_wrappers.ServiceExceptions import NotFoundException from common.proto.qkd_app_pb2 import QKDAppStatusEnum Loading Loading @@ -333,6 +333,38 @@ class TaskExecutor: # return devices return devices def get_device_type_drivers_for_connection( self, connection : Connection ) -> Dict[DeviceTypeEnum, Dict[str, Tuple[Device, Set[int]]]]: devices : Dict[DeviceTypeEnum, Dict[str, Tuple[Device, Set[int]]]] = dict() for endpoint_id in connection.path_hops_endpoint_ids: device = self.get_device(endpoint_id.device_id) device_uuid = endpoint_id.device_id.device_uuid.uuid if device is None: raise Exception('Device({:s}) not found'.format(str(device_uuid))) controller = self.get_device_controller(device) if controller is None: device_type = DeviceTypeEnum._value2member_map_[device.device_type] device_drivers = set(driver for driver in device.device_drivers) devices.setdefault(device_type, dict())[device_uuid] = (device, device_drivers) else: # Controller device types for those underlying path is needed by service handler device_type = DeviceTypeEnum._value2member_map_[controller.device_type] controller_drivers = set(driver for driver in controller.device_drivers) if device_type not in EXPANSION_CONTROLLER_DEVICE_TYPES: devices.setdefault(device_type, dict())[device_uuid] = (device, controller_drivers) else: controller_uuid = controller.device_id.device_uuid.uuid devices.setdefault(device_type, dict())[controller_uuid] = (controller, controller_drivers) LOGGER.debug('[get_devices_from_connection] devices = {:s}'.format(str(devices))) return devices # ----- Service-related methods ------------------------------------------------------------------------------------ def get_service(self, service_id : ServiceId) -> Service: Loading Loading @@ -374,18 +406,22 @@ class TaskExecutor: #LOGGER.debug('connection_device_types_included = {:s}'.format(str(connection_device_types_included))) # ================================================================================================ connection_device_types : Dict[DeviceTypeEnum, Dict[str, Device]] = self.get_devices_from_connection( connection, exclude_managed_by_controller=False ) device_type_to_device_and_drivers : Dict[DeviceTypeEnum, Dict[str, Tuple[Device, Set[int]]]] = \ self.get_device_type_drivers_for_connection(connection) service_handlers : Dict[DeviceTypeEnum, Tuple['_ServiceHandler', Dict[str, Device]]] = dict() # ===== Ryu original test ======================================================================== #for device_type, connection_devices in connection_device_types_excluded.items(): # ================================================================================================ for device_type, connection_devices in connection_device_types.items(): for device_type, device_and_drivers in device_type_to_device_and_drivers.items(): try: service_handler_class = get_service_handler_class( self._service_handler_factory, service, connection_devices self._service_handler_factory, service, device_and_drivers ) connection_devices = { device_uuid : device for device_uuid, (device, _) in device_and_drivers.items() } # ===== Ryu original test ======================================================================== #LOGGER.debug('service_handler_class IN CONNECTION DEVICE TYPE EXCLUDED = {:s}'.format(str(service_handler_class.__name__))) #service_handler = service_handler_class(service, self, **service_handler_settings) Loading @@ -402,11 +438,18 @@ class TaskExecutor: UnsupportedFilterFieldValueException ): dict_connection_devices = { cd_data.name : (cd_uuid, cd_data.name, { cd_data.name : ( cd_uuid, cd_data.name, { (device_driver, DeviceDriverEnum.Name(device_driver)) for device_driver in cd_data.device_drivers }) for cd_uuid,cd_data in connection_devices.items() }, { (device_driver, DeviceDriverEnum.Name(device_driver)) for device_driver in drivers } ) for cd_uuid,(cd_data, drivers) in device_and_drivers.items() } MSG = 'Unable to select service handler. service={:s} connection={:s} connection_devices={:s}' LOGGER.exception(MSG.format( Loading