diff --git a/gbpservice/nfp/core/controller.py b/gbpservice/nfp/core/controller.py index b891762d3e..ffbc3f79fd 100644 --- a/gbpservice/nfp/core/controller.py +++ b/gbpservice/nfp/core/controller.py @@ -522,7 +522,12 @@ def load_nfp_modules(conf, controller): modules_dir = base_module.__path__[0] try: files = os.listdir(modules_dir) - for pyfile in set([f for f in files if f.endswith(".py")]): + pyfiles = set([f for f in files if f.endswith(".py")]) + for pyfile in pyfiles: + module_name = pyfile.strip('.py') + nsd_module_name = module_name + '_NSD.py' + if nsd_module_name in pyfiles: + continue try: pymodule = __import__(conf.nfp_modules_path, globals(), locals(), diff --git a/gbpservice/nfp/orchestrator/config_drivers/heat_driver.py b/gbpservice/nfp/orchestrator/config_drivers/heat_driver.py index c8df970260..433437fe7b 100644 --- a/gbpservice/nfp/orchestrator/config_drivers/heat_driver.py +++ b/gbpservice/nfp/orchestrator/config_drivers/heat_driver.py @@ -509,7 +509,7 @@ def _create_firewall_template(self, auth_token, subnets = consumer['subnets'] # Skip the stitching PTG - if ptg['proxied_group_id']: + if ptg.get('proxied_group_id'): continue fw_template_properties.update({'name': ptg['id'][:3]}) diff --git a/gbpservice/nfp/orchestrator/db/nfp_db_NSD.py b/gbpservice/nfp/orchestrator/db/nfp_db_NSD.py new file mode 100644 index 0000000000..98a090fe0a --- /dev/null +++ b/gbpservice/nfp/orchestrator/db/nfp_db_NSD.py @@ -0,0 +1,130 @@ +from gbpservice.nfp.orchestrator.db.nfp_db import NFPDbBase + +from gbpservice.nfp.orchestrator.db import common_db_mixin +from gbpservice.nfp.orchestrator.db import nfp_db_model + +from gbpservice.nfp.core import log as nfp_logging +LOG = nfp_logging.getLogger(__name__) + + +class NFPDbBaseNSD(NFPDbBase): + def __init__(self, *args, **kwargs): + super(NFPDbBaseNSD, self).__init__(*args, **kwargs) + + def _set_plugged_in_port_for_nfd_interface(self, session, nfd_interface_db, + interface, is_update=False): + plugged_in_port_id = interface.get('plugged_in_port_id') + if not plugged_in_port_id: + if not is_update: + nfd_interface_db.plugged_in_port_id = None + return + with session.begin(subtransactions=True): + port_info_db = nfp_db_model.PortInfo( + id=plugged_in_port_id['id'], + port_model=plugged_in_port_id['port_model'], + port_classification=plugged_in_port_id['port_classification'], + port_role=plugged_in_port_id['port_role']) + if is_update: + session.merge(port_info_db) + else: + session.add(port_info_db) + session.flush() + nfd_interface_db.plugged_in_port_id = port_info_db['id'] + del interface['plugged_in_port_id'] + + + def create_network_function_device_interface(self, session, + nfd_interface): + with session.begin(subtransactions=True): + mapped_real_port_id = nfd_interface.get('mapped_real_port_id') + nfd_interface_db = nfp_db_model.NetworkFunctionDeviceInterface( + id=(nfd_interface.get('id') or uuidutils.generate_uuid()), + tenant_id=nfd_interface['tenant_id'], + interface_position=nfd_interface['interface_position'], + mapped_real_port_id=mapped_real_port_id, + network_function_device_id=( + nfd_interface['network_function_device_id'])) + self._set_plugged_in_port_for_nfd_interface( + session, nfd_interface_db, nfd_interface) + session.add(nfd_interface_db) + + return self._make_network_function_device_interface_dict( + nfd_interface_db) + + def update_network_function_device_interface(self, session, + nfd_interface_id, + updated_nfd_interface): + with session.begin(subtransactions=True): + nfd_interface_db = self._get_network_function_device_interface( + session, nfd_interface_id) + self._set_plugged_in_port_for_nfd_interface( + session, nfd_interface_db, updated_nfd_interface, + is_update=True) + nfd_interface_db.update(updated_nfd_interface) + return self._make_network_function_device_interface_dict( + nfd_interface_db) + + def delete_network_function_device_interface( + self, session, network_function_device_interface_id): + with session.begin(subtransactions=True): + network_function_device_interface_db = ( + self._get_network_function_device_interface( + session, network_function_device_interface_id)) + if network_function_device_interface_db.plugged_in_port_id: + self.delete_port_info( + session, + network_function_device_interface_db.plugged_in_port_id) + session.delete(network_function_device_interface_db) + + def _get_network_function_device_interface(self, session, + network_function_device_id): + try: + return self._get_by_id( + session, + nfp_db_model.NetworkFunctionDeviceInterface, + network_function_device_id) + except exc.NoResultFound: + raise nfp_exc.NetworkFunctionDeviceNotFound( + network_function_device_id=network_function_device_id) + + def get_network_function_device_interface( + self, session, network_function_device_interface_id, + fields=None): + network_function_device_interface = ( + self._get_network_function_device_interface( + session, network_function_device_interface_id)) + return self._make_network_function_device_interface_dict( + network_function_device_interface, fields) + + def get_network_function_device_interfaces(self, session, filters=None, + fields=None, sorts=None, + limit=None, marker=None, + page_reverse=False): + marker_obj = self._get_marker_obj( + 'network_function_device_interfaces', limit, marker) + return self._get_collection( + session, + nfp_db_model.NetworkFunctionDeviceInterface, + self._make_network_function_device_interface_dict, + filters=filters, fields=fields, + sorts=sorts, limit=limit, + marker_obj=marker_obj, + page_reverse=page_reverse) + + + def _make_network_function_device_interface_dict(self, nfd_interface, + fields=None): + res = {'id': nfd_interface['id'], + 'tenant_id': nfd_interface['tenant_id'], + 'plugged_in_port_id': nfd_interface['plugged_in_port_id'], + 'interface_position': nfd_interface['interface_position'], + 'mapped_real_port_id': nfd_interface['mapped_real_port_id'], + 'network_function_device_id': ( + nfd_interface['network_function_device_id']), + } + return res + + + + + diff --git a/gbpservice/nfp/orchestrator/db/nfp_db_model_NSD.py b/gbpservice/nfp/orchestrator/db/nfp_db_model_NSD.py new file mode 100644 index 0000000000..4c2481d8a0 --- /dev/null +++ b/gbpservice/nfp/orchestrator/db/nfp_db_model_NSD.py @@ -0,0 +1,49 @@ +from neutron.db import model_base + + +class PortInfo(BASE, model_base.HasId, model_base.HasTenant): + """Represents the Port Information""" + __tablename__ = 'nfp_port_infos' + + port_model = sa.Column(sa.Enum(nfp_constants.NEUTRON_PORT, + nfp_constants.GBP_PORT, + name='port_model')) + port_classification = sa.Column(sa.Enum(nfp_constants.PROVIDER, + nfp_constants.CONSUMER, + nfp_constants.MANAGEMENT, + nfp_constants.MONITOR, + nfp_constants.ADVANCE_SHARING, + name='port_classification')) + port_role = sa.Column(sa.Enum(nfp_constants.ACTIVE_PORT, + nfp_constants.STANDBY_PORT, + nfp_constants.MASTER_PORT, + name='port_role'), + nullable=True) + + + +class NetworkFunctionDeviceInterface(BASE, model_base.HasId, model_base.HasTenant): + """Represents the Network Function Device""" + __tablename__ = 'nfp_network_function_device_interfaces' + + plugged_in_port_id = sa.Column(sa.String(36), + sa.ForeignKey('nfp_port_infos.id', + ondelete='SET NULL'), + nullable=True) + interface_position = sa.Column(sa.Integer(), nullable=False) + mapped_real_port_id = sa.Column(sa.String(36), + sa.ForeignKey('nfp_port_infos.id', + ondelete='SET NULL'), + nullable=True) + network_function_device_id = sa.Column( + sa.String(36), + sa.ForeignKey('nfp_network_function_devices.id', + ondelete='SET NULL'), + nullable=False) + + + + + + + diff --git a/gbpservice/nfp/orchestrator/drivers/orchestration_driver_NSD.py b/gbpservice/nfp/orchestrator/drivers/orchestration_driver_NSD.py new file mode 100644 index 0000000000..b65199b943 --- /dev/null +++ b/gbpservice/nfp/orchestrator/drivers/orchestration_driver_NSD.py @@ -0,0 +1,1244 @@ +from gbpservice.nfp.orchestrator.drivers.orchestration_driver import ( + OrchestrationDriver, _set_network_handler) + +from collections import defaultdict +from neutron._i18n import _LE +from neutron._i18n import _LI + +from gbpservice.nfp.common import constants as nfp_constants +from gbpservice.nfp.common import exceptions +from gbpservice.nfp.orchestrator.coal.networking import ( + nfp_gbp_network_driver +) +from gbpservice.nfp.orchestrator.coal.networking import ( + nfp_neutron_network_driver +) +from gbpservice.nfp.orchestrator.openstack import openstack_driver + +import ast +import operator + +from gbpservice.nfp.core import log as nfp_logging + +from gbpservice.nfp.core import executor as nfp_executor + +LOG = nfp_logging.getLogger(__name__) + + + +PROXY_PORT_PREFIX = "opflex_proxy:" +ADVANCE_SHARING_PTG_NAME = "Advance_Sharing_PTG" + + + +class OrchestrationDriverNSD(OrchestrationDriver): + """Generic Driver class for orchestration of virtual appliances + + Launches the VM with all the management and data ports and a new VM + is launched for each Network Service Instance + """ + + def __init__(self, config, supports_device_sharing=True, + supports_hotplug=True, max_interfaces=10): + super(OrchestrationDriverNSD, self).__init__(config, + supports_device_sharing,supports_hotplug, max_interfaces) + self.service_vendor = 'general' + self.supports_device_sharing = supports_device_sharing + self.supports_hotplug = supports_hotplug + self.maximum_interfaces = max_interfaces + + self.setup_mode = self._get_setup_mode(config) + self._advance_sharing_network_id = None + + # statistics available + # - instances + # - management_interfaces + # - keystone_token_get_failures + # - image_details_get_failures + # - port_details_get_failures + # - instance_launch_failures + # - instance_details_get_failures + # - instance_delete_failures + # - interface_plug_failures + # - interface_unplug_failures + self.stats = {} + + def _get_setup_mode(self, config): + # REVISIT(TODO): Removing sharing for cisco live demo + return {nfp_constants.NEUTRON_MODE: True} + # if nfp_constants.APIC_CONFIG_SECTION in config.list_all_sections(): + # return {nfp_constants.APIC_MODE: True} + # else: + # return {nfp_constants.NEUTRON_MODE: True} + + + def _get_token(self, device_data_token): + + try: + token = (device_data_token + if device_data_token + else self.identity_handler.get_admin_token()) + except Exception: + self._increment_stats_counter('keystone_token_get_failures') + LOG.error(_LE('Failed to get token')) + return None + return token + + def _increment_stats_counter(self, metric, by=1): + # TODO(RPM): create path and delete path have different driver objects. + # This will not work in case of increment and decrement. + # So, its no-operation now + return + try: + self.stats.update({metric: self.stats.get(metric, 0) + by}) + except Exception: + LOG.error(_LE("Statistics failure. Failed to increment" + " '%(metric)s' by %(by)d") + % {'metric': metric, 'by': by}) + + def _decrement_stats_counter(self, metric, by=1): + # TODO(RPM): create path and delete path have different driver objects. + # This will not work in case of increment and decrement. + # So, its no-operation now + return + try: + self.stats.update({metric: self.stats[metric] - by}) + except Exception: + LOG.error(_LE("Statistics failure. Failed to decrement" + " '%(metric)s' by %(by)d") + % {'metric': metric, 'by': by}) + + def _is_device_sharing_supported(self): + return self.supports_device_sharing + + + def _get_advance_sharing_network_id(self, admin_tenant_id, + network_handler=None): + if self._advance_sharing_network_id: + return self._advance_sharing_network_id + filters = {'tenant_id': admin_tenant_id, + 'name': ADVANCE_SHARING_PTG_NAME} + admin_token = self._get_token(None) + if not admin_token: + return None + sharing_networks = network_handler.get_networks( + admin_token, filters=filters) + if not sharing_networks: + LOG.error(_LE("Found empty network for tenant with Tenant ID: " + "%(admin_tenant_id)s for advance sharing"), + {'admin_tenant_id': admin_tenant_id}) + raise Exception() + elif len(sharing_networks) > 1: + LOG.error(_LE("Found more then one network for sharing with" + " Tenant ID: %(admin_tenant_id)s for " + "advance sharing"), + {'admin_tenant_id': admin_tenant_id}) + raise Exception() + else: + self._advance_sharing_network_id = sharing_networks[0]['id'] + return self._advance_sharing_network_id + + def _create_advance_sharing_interfaces(self, device_data, + network_handler=None): + token = self._get_token(device_data.get('token')) + if not token: + return None + + admin_tenant_id = self._get_admin_tenant_id(token=token) + port_infos = [] + port_model = (nfp_constants.GBP_PORT + if device_data['service_details'][ + 'network_mode'] == nfp_constants.GBP_MODE + else nfp_constants.NEUTRON_PORT) + advance_sharing_network_id = self._get_advance_sharing_network_id( + admin_tenant_id, + network_handler) + for i in range(self.maximum_interfaces): + port = network_handler.create_port(token, + admin_tenant_id, + advance_sharing_network_id) + port_infos.append({'id': port['id'], + 'port_model': port_model, + 'port_classification': ( + nfp_constants.ADVANCE_SHARING), + 'port_role': None, + 'plugged_in_pt_id': ( + network_handler.get_port_id(token, + port['id']))}) + return port_infos + + + + def _get_vendor_data(self, device_data, image_name): + token = self._get_token(device_data.get('token')) + if not token: + return None + try: + metadata = self.compute_handler_nova.get_image_metadata( + token, + self._get_admin_tenant_id(token=token), + image_name) + except Exception as e: + self._increment_stats_counter('image_details_get_failures') + LOG.error(_LE('Failed to get image metadata for image ' + 'name: %(image_name)s. Error: %(error)s'), + {'image_name': image_name, 'error': e}) + return None + vendor_data = self._verify_vendor_data(image_name, metadata) + if not vendor_data: + return None + return vendor_data + + def _get_vendor_data_fast(self, token, + admin_tenant_id, image_name, device_data): + try: + metadata = self.compute_handler_nova.get_image_metadata( + token, + admin_tenant_id, + image_name) + except Exception as e: + self._increment_stats_counter('image_details_get_failures') + LOG.error(_LE('Failed to get image metadata for image ' + 'name: %(image_name)s. Error: %(error)s'), + {'image_name': image_name, 'error': e}) + return None + vendor_data = self._verify_vendor_data(image_name, metadata) + if not vendor_data: + return None + return vendor_data + + + + def _update_vendor_data(self, device_data, token=None): + try: + image_name = self._get_image_name(device_data) + vendor_data = self._get_vendor_data(device_data, image_name) + LOG.info(_LI("Vendor data, specified in image: %(vendor_data)s"), + {'vendor_data': vendor_data}) + if vendor_data: + self._update_self_with_vendor_data( + vendor_data, + nfp_constants.MAXIMUM_INTERFACES) + self._update_self_with_vendor_data( + vendor_data, + nfp_constants.SUPPORTS_SHARING) + self._update_self_with_vendor_data( + vendor_data, + nfp_constants.SUPPORTS_HOTPLUG) + else: + LOG.info(_LI("No vendor data specified in image, " + "proceeding with default values")) + except Exception: + LOG.error(_LE("Error while getting metadata for image name:" + "%(image_name)s, proceeding with default values"), + {'image_name': image_name}) + + def _update_vendor_data_fast(self, token, admin_tenant_id, + image_name, device_data): + vendor_data = None + try: + vendor_data = self._get_vendor_data_fast( + token, admin_tenant_id, image_name, device_data) + LOG.info(_LI("Vendor data, specified in image: %(vendor_data)s"), + {'vendor_data': vendor_data}) + if vendor_data: + self._update_self_with_vendor_data( + vendor_data, + nfp_constants.MAXIMUM_INTERFACES) + self._update_self_with_vendor_data( + vendor_data, + nfp_constants.SUPPORTS_SHARING) + self._update_self_with_vendor_data( + vendor_data, + nfp_constants.SUPPORTS_HOTPLUG) + + else: + LOG.info(_LI("No vendor data specified in image, " + "proceeding with default values")) + except Exception: + LOG.error(_LE("Error while getting metadata for image name: " + "%(image_name)s, proceeding with default values"), + {'image_name': image_name}) + return vendor_data + + + + def get_network_function_device_sharing_info(self, device_data): + """ Get filters for NFD sharing + + :param device_data: NFD data + :type device_data: dict + + :returns: None -- when device sharing is not supported + :returns: dict -- It has the following scheme + { + 'filters': { + 'key': 'value', + ... + } + } + + :raises: exceptions.IncompleteData + """ + + if ( + any(key not in device_data + for key in ['tenant_id', + 'service_details']) or + + type(device_data['service_details']) is not dict or + + any(key not in device_data['service_details'] + for key in ['service_vendor']) + ): + raise exceptions.IncompleteData() + + if not self._is_device_sharing_supported(): + return None + + return { + 'filters': { + 'tenant_id': [device_data['tenant_id']], + 'service_vendor': [device_data['service_details'][ + 'service_vendor']], + 'status': [nfp_constants.ACTIVE] + } + } + + @_set_network_handler + def select_network_function_device(self, devices, device_data, + network_handler=None): + """ Select a NFD which is eligible for sharing + + :param devices: NFDs + :type devices: list + :param device_data: NFD data + :type device_data: dict + + :returns: None -- when device sharing is not supported, or + when no device is eligible for sharing + :return: dict -- NFD which is eligible for sharing + + :raises: exceptions.IncompleteData + """ + + if ( + any(key not in device_data + for key in ['ports']) or + + type(device_data['ports']) is not list or + + any(key not in port + for port in device_data['ports'] + for key in ['id', + 'port_classification', + 'port_model']) or + + type(devices) is not list or + + any(key not in device + for device in devices + for key in ['interfaces_in_use']) + ): + raise exceptions.IncompleteData() + + token = self._get_token(device_data.get('token')) + if not token: + return None + image_name = self._get_image_name(device_data) + if image_name: + self._update_vendor_data(device_data, + device_data.get('token')) + if not self._is_device_sharing_supported(): + return None + + hotplug_ports_count = 1 # for provider interface (default) + if any(port['port_classification'] == nfp_constants.CONSUMER + for port in device_data['ports']): + hotplug_ports_count = 2 + + device_service_types_map = ( + self._get_device_service_types_map(token, devices, + network_handler)) + service_type = device_data['service_details']['service_type'] + for device in devices: + if ( + (device['interfaces_in_use'] + hotplug_ports_count) <= + self.maximum_interfaces + ): + if (service_type.lower() == nfp_constants.VPN.lower() and + service_type in device_service_types_map[ + device['id']]): + # Restrict multiple VPN services to share same device + # If nfd request service type is VPN and current filtered + # device already has VPN service instantiated, ignore this + # device and checks for next one + continue + return device + return None + + + @_set_network_handler + def create_network_function_device(self, device_data, + network_handler=None): + """ Create a NFD + + :param device_data: NFD data + :type device_data: dict + + :returns: None -- when there is a failure in creating NFD + :return: dict -- NFD created + + :raises: exceptions.IncompleteData, + exceptions.ComputePolicyNotSupported + """ + if ( + any(key not in device_data + for key in ['service_details', + 'name', + 'management_network_info', + 'ports']) or + + type(device_data['service_details']) is not dict or + + any(key not in device_data['service_details'] + for key in ['service_vendor', + 'device_type', + 'network_mode']) or + + any(key not in device_data['management_network_info'] + for key in ['id']) or + + type(device_data['ports']) is not list or + + any(key not in port + for port in device_data['ports'] + for key in ['id', + 'port_classification', + 'port_model']) + ): + raise exceptions.IncompleteData() + + if ( + device_data['service_details']['device_type'] != + nfp_constants.NOVA_MODE + ): + raise exceptions.ComputePolicyNotSupported( + compute_policy=device_data['service_details']['device_type']) + + token = device_data['token'] + admin_tenant_id = device_data['admin_tenant_id'] + image_name = self._get_image_name(device_data) + + executor = nfp_executor.TaskExecutor(jobs=3) + + image_id_result = {} + vendor_data_result = {} + + executor.add_job('UPDATE_VENDOR_DATA', + self._update_vendor_data_fast, + token, admin_tenant_id, image_name, device_data, + result_store=vendor_data_result) + executor.add_job('GET_INTERFACES_FOR_DEVICE_CREATE', + self._get_interfaces_for_device_create, + token, admin_tenant_id, network_handler, device_data) + executor.add_job('GET_IMAGE_ID', + self.get_image_id, + self.compute_handler_nova, token, admin_tenant_id, + image_name, result_store=image_id_result) + + executor.fire() + + interfaces = device_data.pop('interfaces', None) + if not interfaces: + LOG.exception(_LE('Failed to get interfaces for device creation.')) + return None + else: + management_interface = interfaces[0] + self._increment_stats_counter('management_interfaces', + by=len(interfaces)) + + image_id = image_id_result.get('result', None) + if not image_id: + self._increment_stats_counter('image_details_get_failures') + LOG.error(_LE('Failed to get image id for device creation.')) + self._delete_interfaces(device_data, interfaces, + network_handler=network_handler) + self._decrement_stats_counter('management_interfaces', + by=len(interfaces)) + return None + + vendor_data = vendor_data_result.get('result', None) + if not vendor_data: + LOG.warn(_LE('Failed to get vendor data for device creation.')) + vendor_data = {} + + if device_data['service_details'].get('flavor'): + flavor = device_data['service_details']['flavor'] + else: + LOG.info(_LI("No Device flavor provided in service profile's " + "service flavor field, using default " + "flavor: m1.medium")) + flavor = 'm1.medium' + + interfaces_to_attach = [] + advance_sharing_interfaces = [] + try: + for interface in interfaces: + interfaces_to_attach.append({'port': interface['port_id']}) + + if vendor_data.get('supports_hotplug') == False: + if self.setup_mode.get(nfp_constants.NEUTRON_MODE): + # TODO(ashu): get neutron mode from conf + for port in device_data['ports']: + if (port['port_classification'] == + nfp_constants.PROVIDER): + if (device_data['service_details'][ + 'service_type'].lower() + in [nfp_constants.FIREWALL.lower(), + nfp_constants.VPN.lower()]): + network_handler.set_promiscuos_mode( + token, port['id']) + port_id = network_handler.get_port_id( + token, port['id']) + interfaces_to_attach.append({'port': port_id}) + for port in device_data['ports']: + if (port['port_classification'] == + nfp_constants.CONSUMER): + if (device_data['service_details'][ + 'service_type'].lower() + in [nfp_constants.FIREWALL.lower(), + nfp_constants.VPN.lower()]): + network_handler.set_promiscuos_mode( + token, port['id']) + port_id = network_handler.get_port_id( + token, port['id']) + interfaces_to_attach.append({'port': port_id}) + elif self.setup_mode.get(nfp_constants.APIC_MODE): + advance_sharing_interfaces = ( + self._create_advance_sharing_interfaces( + device_data, + network_handler)) + for interface in advance_sharing_interfaces: + port_id = network_handler.get_port_id(token, + interface['id']) + interfaces_to_attach.append({'port': port_id}) + + interfaces += advance_sharing_interfaces + + except Exception as e: + self._increment_stats_counter('port_details_get_failures') + LOG.error(_LE('Failed to fetch list of interfaces to attach' + ' for device creation %(error)s'), {'error': e}) + self._delete_interfaces(device_data, interfaces, + network_handler=network_handler) + self._decrement_stats_counter('management_interfaces', + by=len(interfaces)) + return None + + instance_name = device_data['name'] + instance_id_result = {} + port_details_result = {} + + executor.add_job('CREATE_INSTANCE', + self.create_instance, + self.compute_handler_nova, + token, admin_tenant_id, image_id, flavor, + interfaces_to_attach, instance_name, + result_store=instance_id_result) + + executor.add_job('GET_NEUTRON_PORT_DETAILS', + self.get_neutron_port_details, + network_handler, token, + management_interface['port_id'], + result_store=port_details_result) + + executor.fire() + + instance_id = instance_id_result.get('result', None) + if not instance_id: + self._increment_stats_counter('instance_launch_failures') + LOG.error(_LE('Failed to create %(device_type)s instance.')) + self._delete_interfaces(device_data, interfaces, + network_handler=network_handler) + self._decrement_stats_counter('management_interfaces', + by=len(interfaces)) + return None + else: + self._increment_stats_counter('instances') + + mgmt_ip_address = None + mgmt_neutron_port_info = port_details_result.get('result', None) + + if not mgmt_neutron_port_info: + self._increment_stats_counter('port_details_get_failures') + LOG.error(_LE('Failed to get management port details. ')) + try: + self.compute_handler_nova.delete_instance( + token, + admin_tenant_id, + instance_id) + except Exception as e: + self._increment_stats_counter('instance_delete_failures') + LOG.error(_LE('Failed to delete %(device_type)s instance.' + 'Error: %(error)s'), + {'device_type': ( + device_data['service_details']['device_type']), + 'error': e}) + self._decrement_stats_counter('instances') + self._delete_interfaces(device_data, interfaces, + network_handler=network_handler) + self._decrement_stats_counter('management_interfaces', + by=len(interfaces)) + return None + + mgmt_ip_address = mgmt_neutron_port_info['ip_address'] + return {'id': instance_id, + 'name': instance_name, + 'vendor_data': vendor_data, + 'mgmt_ip_address': mgmt_ip_address, + 'mgmt_port_id': interfaces[0], + 'mgmt_neutron_port_info': mgmt_neutron_port_info, + 'max_interfaces': self.maximum_interfaces, + 'interfaces_in_use': len(interfaces_to_attach), + 'advance_sharing_interfaces': advance_sharing_interfaces, + 'description': ''} # TODO(RPM): what should be the description + + + + @_set_network_handler + def delete_network_function_device(self, device_data, + network_handler=None): + """ Delete the NFD + + :param device_data: NFD + :type device_data: dict + + :returns: None -- Both on success and Failure + + :raises: exceptions.IncompleteData, + exceptions.ComputePolicyNotSupported + """ + if ( + any(key not in device_data + for key in ['service_details', + 'mgmt_port_id']) or + + type(device_data['service_details']) is not dict or + + any(key not in device_data['service_details'] + for key in ['service_vendor', + 'device_type', + 'network_mode']) or + + type(device_data['mgmt_port_id']) is not dict or + + any(key not in device_data['mgmt_port_id'] + for key in ['id', + 'port_classification', + 'port_model']) + ): + raise exceptions.IncompleteData() + + if ( + device_data['service_details']['device_type'] != + nfp_constants.NOVA_MODE + ): + raise exceptions.ComputePolicyNotSupported( + compute_policy=device_data['service_details']['device_type']) + + image_name = self._get_image_name(device_data) + if image_name: + self._update_vendor_data(device_data, + device_data.get('token')) + token = self._get_token(device_data.get('token')) + if not token: + return None + + if device_data.get('id'): + # delete the device instance + # + # this method will be invoked again + # once the device instance deletion is completed + try: + self.compute_handler_nova.delete_instance( + token, + self._get_admin_tenant_id( + token=token), + device_data['id']) + except Exception: + self._increment_stats_counter('instance_delete_failures') + LOG.error(_LE('Failed to delete %(instance)s instance'), + {'instance': + device_data['service_details']['device_type']}) + + else: + self._decrement_stats_counter('instances') + else: + # device instance deletion is done, delete remaining resources + try: + interfaces = [device_data['mgmt_port_id']] + interfaces.extend(device_data['advance_sharing_interfaces']) + self._delete_interfaces(device_data, + interfaces, + network_handler=network_handler) + except Exception as e: + LOG.error(_LE('Failed to delete the management data port(s). ' + 'Error: %(error)s'), {'error': e}) + else: + self._decrement_stats_counter('management_interfaces') + + + + def get_network_function_device_status(self, device_data, + ignore_failure=False): + """ Get the status of NFD + + :param device_data: NFD + :type device_data: dict + + :returns: None -- On failure + :return: str -- status string + + :raises: exceptions.IncompleteData, + exceptions.ComputePolicyNotSupported + """ + if ( + any(key not in device_data + for key in ['id', + 'service_details']) or + + type(device_data['service_details']) is not dict or + + any(key not in device_data['service_details'] + for key in ['service_vendor', + 'device_type', + 'network_mode']) + ): + raise exceptions.IncompleteData() + + if ( + device_data['service_details']['device_type'] != + nfp_constants.NOVA_MODE + ): + raise exceptions.ComputePolicyNotSupported( + compute_policy=device_data['service_details']['device_type']) + + try: + device = self.compute_handler_nova.get_instance( + device_data['token'], + device_data['tenant_id'], + device_data['id']) + except Exception: + if ignore_failure: + return None + self._increment_stats_counter('instance_details_get_failures') + LOG.error(_LE('Failed to get %(instance)s instance details'), + {device_data['service_details']['device_type']}) + return None # TODO(RPM): should we raise an Exception here? + + return device['status'] + + + + def get_network_function_device_status(self, device_data, + ignore_failure=False): + """ Get the status of NFD + + :param device_data: NFD + :type device_data: dict + + :returns: None -- On failure + :return: str -- status string + + :raises: exceptions.IncompleteData, + exceptions.ComputePolicyNotSupported + """ + if ( + any(key not in device_data + for key in ['id', + 'service_details']) or + + type(device_data['service_details']) is not dict or + + any(key not in device_data['service_details'] + for key in ['service_vendor', + 'device_type', + 'network_mode']) + ): + raise exceptions.IncompleteData() + + if ( + device_data['service_details']['device_type'] != + nfp_constants.NOVA_MODE + ): + raise exceptions.ComputePolicyNotSupported( + compute_policy=device_data['service_details']['device_type']) + + try: + device = self.compute_handler_nova.get_instance( + device_data['token'], + device_data['tenant_id'], + device_data['id']) + except Exception: + if ignore_failure: + return None + self._increment_stats_counter('instance_details_get_failures') + LOG.error(_LE('Failed to get %(instance)s instance details'), + {device_data['service_details']['device_type']}) + return None # TODO(RPM): should we raise an Exception here? + + return device['status'] + + + @_set_network_handler + def plug_network_function_device_interfaces(self, device_data, + network_handler=None): + """ Attach the network interfaces for NFD + + :param device_data: NFD + :type device_data: dict + + :returns: bool -- False on failure and True on Success + + :raises: exceptions.IncompleteData, + exceptions.ComputePolicyNotSupported + """ + + if ( + any(key not in device_data + for key in ['id', + 'service_details', + 'ports']) or + + type(device_data['service_details']) is not dict or + + any(key not in device_data['service_details'] + for key in ['service_vendor', + 'device_type', + 'network_mode']) or + + type(device_data['ports']) is not list or + + any(key not in port + for port in device_data['ports'] + for key in ['id', + 'port_classification', + 'port_model']) + ): + raise exceptions.IncompleteData() + + if ( + device_data['service_details']['device_type'] != + nfp_constants.NOVA_MODE + ): + raise exceptions.ComputePolicyNotSupported( + compute_policy=device_data['service_details']['device_type']) + + token = device_data['token'] + tenant_id = device_data['tenant_id'] + vendor_data = device_data['vendor_data'] + + if vendor_data: + self._update_self_with_vendor_data( + vendor_data, + nfp_constants.MAXIMUM_INTERFACES) + self._update_self_with_vendor_data( + vendor_data, + nfp_constants.SUPPORTS_SHARING) + self._update_self_with_vendor_data( + vendor_data, + nfp_constants.SUPPORTS_HOTPLUG) + else: + LOG.info(_LI("No vendor data specified in image, " + "proceeding with default values")) + + update_ifaces = [] + try: + if vendor_data.get('supports_hotplug') == False: + # configure interfaces instead of hotplug + if self.setup_mode.get(nfp_constants.APIC_MODE): + required_ports = len(device_data['ports']) + unused_ifaces = self._get_unused_interfaces( + device_data['advance_sharing_interfaces'], + required_ports) + data_port_ids = [] + + for port in device_data['ports']: + if (port['port_classification'] == + nfp_constants.PROVIDER): + data_port_ids.append(port['id']) + break + for port in device_data['ports']: + if (port['port_classification'] == + nfp_constants.CONSUMER): + data_port_ids.append(port['id']) + + for data_port_id, iface in zip(data_port_ids, + unused_ifaces): + self._update_attached_port_with_data_port( + token, + iface, + data_port_id, + network_handler, + stitch=True) + iface['mapped_real_port_id'] = data_port_id + update_ifaces = unused_ifaces + elif self.setup_mode.get(nfp_constants.NEUTRON_MODE): + pass + else: + executor = nfp_executor.TaskExecutor(jobs=10) + + for port in device_data['ports']: + if port['port_classification'] == nfp_constants.PROVIDER: + service_type = device_data[ + 'service_details']['service_type'].lower() + if service_type.lower() in [ + nfp_constants.FIREWALL.lower(), + nfp_constants.VPN.lower()]: + executor.add_job( + 'SET_PROMISCUOS_MODE', + network_handler.set_promiscuos_mode_fast, + token, port['id']) + executor.add_job( + 'ATTACH_INTERFACE', + self.compute_handler_nova.attach_interface, + token, tenant_id, device_data['id'], + port['id']) + break + + for port in device_data['ports']: + if port['port_classification'] == nfp_constants.CONSUMER: + service_type = device_data[ + 'service_details']['service_type'].lower() + if service_type.lower() in [ + nfp_constants.FIREWALL.lower(), + nfp_constants.VPN.lower()]: + executor.add_job( + 'SET_PROMISCUOS_MODE', + network_handler.set_promiscuos_mode_fast, + token, port['id']) + executor.add_job( + 'ATTACH_INTERFACE', + self.compute_handler_nova.attach_interface, + token, tenant_id, device_data['id'], + port['id']) + break + executor.fire() + + except Exception as e: + self._increment_stats_counter('interface_plug_failures') + LOG.error(_LE('Failed to plug interface(s) to the device.' + 'Error: %(error)s'), {'error': e}) + return None, [] + else: + return True, update_ifaces + + def _update_attached_port_with_data_port(self, token, + unused_interface, + data_port_id, + network_handler=None, + stitch=True): + if stitch: + # configure attached interface pt with real port id, + # to specify fabric/controller to stitch these interfaces + port_id = network_handler.get_port_id(token, data_port_id) + description = "%s%s" % ( + PROXY_PORT_PREFIX, + port_id) + # TODO(ashu): update attached port mac with data port mac. + else: + # configure attached interface pt with empty string, + # for further reuse of these interfaces + description = '' + port = {'description': description} + network_handler.update_port(token, unused_interface['id'], port) + + def _get_unused_interfaces(self, advance_sharing_ifaces, required_ports): + # sort the interfaces based on interface position + advance_sharing_ifaces.sort(key=operator.itemgetter( + 'interface_position')) + unused_interfaces = [] + for iface in advance_sharing_ifaces: + if not iface['mapped_real_port_id']: + current_position = iface['interface_position'] + unused_interfaces = advance_sharing_ifaces[ + current_position:required_ports] + break + + return unused_interfaces + + def _get_used_interfaces(self, advance_sharing_ifaces, data_port_ids): + used_interfaces = [] + for iface in advance_sharing_ifaces: + if iface['mapped_real_port_id'] in data_port_ids: + used_interfaces.append(iface) + return used_interfaces + + + @_set_network_handler + def unplug_network_function_device_interfaces(self, device_data, + network_handler=None): + """ Detach the network interfaces for NFD + + :param device_data: NFD + :type device_data: dict + + :returns: bool -- False on failure and True on Success + + :raises: exceptions.IncompleteData, + exceptions.ComputePolicyNotSupported + """ + + if ( + any(key not in device_data + for key in ['id', + 'service_details', + 'ports']) or + + type(device_data['service_details']) is not dict or + + any(key not in device_data['service_details'] + for key in ['service_vendor', + 'device_type', + 'network_mode']) or + + any(key not in port + for port in device_data['ports'] + for key in ['id', + 'port_classification', + 'port_model']) + ): + raise exceptions.IncompleteData() + + if ( + device_data['service_details']['device_type'] != + nfp_constants.NOVA_MODE + ): + raise exceptions.ComputePolicyNotSupported( + compute_policy=device_data['service_details']['device_type']) + + image_name = self._get_image_name(device_data) + if image_name: + self._update_vendor_data(device_data, + device_data.get('token')) + + token = self._get_token(device_data.get('token')) + if not token: + return None + + executor = nfp_executor.TaskExecutor(jobs=1) + vendor_data_result = {} + tenant_id = device_data.get('tenant_id') + + executor.add_job('UPDATE_VENDOR_DATA', + self._update_vendor_data_fast, + token, tenant_id, image_name, device_data, + result_store=vendor_data_result) + executor.fire() + + vendor_data = vendor_data_result.get('result', None) + if not vendor_data: + LOG.warn(_LE('Failed to get vendor data for device deletion.')) + vendor_data = {} + + update_ifaces = [] + try: + if vendor_data.get('supports_hotplug') == False: + if self.setup_mode.get(nfp_constants.APIC_MODE): + data_port_ids = [] + for port in device_data['ports']: + if (port['port_classification'] == + nfp_constants.PROVIDER): + data_port_ids.append(port['id']) + break + for port in device_data['ports']: + if (port['port_classification'] == + nfp_constants.CONSUMER): + data_port_ids.append(port['id']) + + used_ifaces = self._get_used_interfaces( + device_data['advance_sharing_interfaces'], + data_port_ids) + + for data_port_id, iface in zip(data_port_ids, used_ifaces): + self._update_attached_port_with_data_port( + token, + iface, + data_port_id, + network_handler, + stitch=False) + iface['mapped_real_port_id'] = '' + update_ifaces = used_ifaces + elif self.setup_mode.get(nfp_constants.NEUTRON_MODE): + pass + else: + for port in device_data['ports']: + port_id = network_handler.get_port_id(token, port['id']) + self.compute_handler_nova.detach_interface( + token, + self._get_admin_tenant_id(token=token), + device_data['id'], + port_id) + + except Exception as e: + self._increment_stats_counter('interface_unplug_failures') + LOG.error(_LE('Failed to unplug interface(s) from the device.' + 'Error: %(error)s'), {'error': e}) + return None, [] + else: + return True, update_ifaces + + + @_set_network_handler + def get_network_function_device_config_info(self, device_data, + network_handler=None): + """ Get the configuration information for NFD + + :param device_data: NFD + :type device_data: dict + + :returns: None -- On Failure + :returns: dict -- It has the following scheme + { + 'config': [ + { + 'resource': 'interfaces', + 'resource_data': { + ... + } + }, + { + 'resource': 'routes', + 'resource_data': { + ... + } + } + ] + } + + :raises: exceptions.IncompleteData + """ + if ( + any(key not in device_data + for key in ['service_details', + 'mgmt_ip_address', + 'ports']) or + + type(device_data['service_details']) is not dict or + + any(key not in device_data['service_details'] + for key in ['service_vendor', + 'device_type', + 'network_mode']) or + + type(device_data['ports']) is not list or + + any(key not in port + for port in device_data['ports'] + for key in ['id', + 'port_classification', + 'port_model']) + ): + raise exceptions.IncompleteData() + + token = self._get_token(device_data.get('token')) + if not token: + return None + + provider_ip = None + provider_mac = None + provider_cidr = None + consumer_ip = None + consumer_mac = None + consumer_cidr = None + consumer_gateway_ip = None + + for port in device_data['ports']: + if port['port_classification'] == nfp_constants.PROVIDER: + try: + (provider_ip, provider_mac, provider_cidr, dummy) = ( + network_handler.get_port_details(token, port['id']) + ) + except Exception: + self._increment_stats_counter('port_details_get_failures') + LOG.error(_LE('Failed to get provider port details' + ' for get device config info operation')) + return None + elif port['port_classification'] == nfp_constants.CONSUMER: + try: + (consumer_ip, consumer_mac, consumer_cidr, + consumer_gateway_ip) = ( + network_handler.get_port_details(token, port['id']) + ) + except Exception: + self._increment_stats_counter('port_details_get_failures') + LOG.error(_LE('Failed to get consumer port details' + ' for get device config info operation')) + return None + + return { + 'config': [ + { + 'resource': nfp_constants.INTERFACE_RESOURCE, + 'resource_data': { + 'mgmt_ip': device_data['mgmt_ip_address'], + 'provider_ip': provider_ip, + 'provider_cidr': provider_cidr, + 'provider_interface_index': 2, + 'stitching_ip': consumer_ip, + 'stitching_cidr': consumer_cidr, + 'stitching_interface_index': 3, + 'provider_mac': provider_mac, + 'stitching_mac': consumer_mac, + } + }, + { + 'resource': nfp_constants.ROUTES_RESOURCE, + 'resource_data': { + 'mgmt_ip': device_data['mgmt_ip_address'], + 'source_cidrs': ([provider_cidr, consumer_cidr] + if consumer_cidr + else [provider_cidr]), + 'destination_cidr': consumer_cidr, + 'provider_mac': provider_mac, + 'gateway_ip': consumer_gateway_ip, + 'provider_interface_index': 2 + } + } + ] + } + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/gbpservice/nfp/orchestrator/modules/device_orchestrator_NSD.py b/gbpservice/nfp/orchestrator/modules/device_orchestrator_NSD.py new file mode 100644 index 0000000000..1a0b865742 --- /dev/null +++ b/gbpservice/nfp/orchestrator/modules/device_orchestrator_NSD.py @@ -0,0 +1,304 @@ + +from neutron._i18n import _LE +from neutron._i18n import _LI +import oslo_messaging as messaging + +from gbpservice.nfp.common import constants as nfp_constants +from gbpservice.nfp.common import topics as nsf_topics +from gbpservice.nfp.core import event as nfp_event +from gbpservice.nfp.core.event import Event +from gbpservice.nfp.core import module as nfp_api +from gbpservice.nfp.core.rpc import RpcAgent +from gbpservice.nfp.lib import transport +from gbpservice.nfp.orchestrator.db import nfp_db_NSD as nfp_db +from gbpservice.nfp.orchestrator.drivers import orchestration_driver_NSD as orchestration_driver +from gbpservice.nfp.orchestrator.modules.device_orchestrator import ( + DeviceOrchestrator, events_init, rpc_init) +from gbpservice.nfp.orchestrator.openstack import openstack_driver +from neutron.common import rpc as n_rpc +from neutron import context as n_context +from neutron.db import api as db_api + +import sys +import traceback + +from gbpservice.nfp.core import log as nfp_logging +LOG = nfp_logging.getLogger(__name__) + +STOP_POLLING = {'poll': False} +CONTINUE_POLLING = {'poll': True} + + + +def nfp_module_init(controller, config): + events_init(controller, config, DeviceOrchestratorNSD(controller, config)) + rpc_init(controller, config) + LOG.debug("Device Orchestrator: module_init") + + + +class DeviceOrchestratorNSD(DeviceOrchestrator): + + def __init__(self, controller, config): + super(DeviceOrchestratorNSD, self).__init__(controller, config) + self.nsf_db = nfp_db.NFPDbBaseNSD() + self.orchestration_driver = ( + orchestration_driver.OrchestrationDriverNSD(self.config)) + + + def _create_advance_sharing_interfaces(self, device, interfaces_infos): + nfd_interfaces = [] + port_infos = [] + for position, interface in enumerate(interfaces_infos): + interface['network_function_device_id'] = device['id'] + interface['interface_position'] = position + interface['tenant_id'] = device['tenant_id'] + interface['plugged_in_port_id'] = {} + interface['plugged_in_port_id']['id'] = interface['id'] + interface['plugged_in_port_id']['port_model'] = ( + interface.get('port_model')) + interface['plugged_in_port_id']['port_classification'] = ( + interface.get('port_classification')) + interface['plugged_in_port_id']['port_role'] = ( + interface.get('port_role')) + + nfd_interfaces.append( + self.nsf_db.create_network_function_device_interface( + self.db_session, interface) + ) + LOG.debug("Created following entries in port_infos table : %s, " + " network function device interfaces table: %s." % + (port_infos, nfd_interfaces)) + + def _get_advance_sharing_interfaces(self, device_id): + filters = {'network_function_device_id': [device_id]} + network_function_device_interfaces = ( + self.nsf_db.get_network_function_device_interfaces( + self.db_session, + filters=filters) + ) + return network_function_device_interfaces + + def _update_advance_sharing_interfaces(self, device, nfd_ifaces): + for nfd_iface in nfd_ifaces: + for port in device['ports']: + if port['id'] == nfd_iface['mapped_real_port_id']: + nfd_iface['mapped_real_port_id'] = port['id'] + nfd_iface['plugged_in_port_id'] = ( + self.nsf_db.get_port_info( + self.db_session, + nfd_iface['plugged_in_port_id'])) + self.nsf_db.update_network_function_device_interface( + self.db_session, + nfd_iface['id'], + nfd_iface) + break + + def _delete_advance_sharing_interfaces(self, nfd_ifaces): + for nfd_iface in nfd_ifaces: + port_id = nfd_iface['id'] + self.nsf_db.delete_network_function_device_interface( + self.db_session, + port_id) + + def _create_network_function_device_db(self, device_info, state): + advance_sharing_interfaces = [] + + self._update_device_status(device_info, state) + # (ashu) driver should return device_id as vm_id + device_id = device_info.pop('id') + device_info['id'] = device_id + device_info['reference_count'] = 0 + if device_info.get('advance_sharing_interfaces'): + advance_sharing_interfaces = ( + device_info.pop('advance_sharing_interfaces')) + device_info['interfaces_in_use'] = 0 + device = self.nsf_db.create_network_function_device(self.db_session, + device_info) + if advance_sharing_interfaces: + self._create_advance_sharing_interfaces(device, + advance_sharing_interfaces) + return device + + def _delete_network_function_device_db(self, device_id, device): + advance_sharing_interfaces = device.get( + 'advance_sharing_interfaces', []) + if advance_sharing_interfaces: + self._delete_advance_sharing_interfaces( + advance_sharing_interfaces) + self.nsf_db.delete_network_function_device(self.db_session, device_id) + + + def _prepare_device_data(self, device_info): + network_function_id = device_info['network_function_id'] + network_function_device_id = device_info['network_function_device_id'] + network_function_instance_id = ( + device_info['network_function_instance_id']) + + network_function = self._get_nsf_db_resource( + 'network_function', + network_function_id) + network_function_device = self._get_nsf_db_resource( + 'network_function_device', + network_function_device_id) + network_function_instance = self._get_nsf_db_resource( + 'network_function_instance', + network_function_instance_id) + + admin_token = self.keystoneclient.get_admin_token() + service_profile = self.gbpclient.get_service_profile( + admin_token, network_function['service_profile_id']) + service_details = transport.parse_service_flavor_string( + service_profile['service_flavor']) + + device_info.update({ + 'network_function_instance': network_function_instance}) + device_info.update({'id': network_function_device_id}) + service_details.update({'service_type': self._get_service_type( + network_function['service_profile_id'])}) + device_info.update({'service_details': service_details}) + + device = self._get_device_data(device_info) + device = self._update_device_data(device, network_function_device) + + mgmt_port_id = network_function_device.pop('mgmt_port_id') + mgmt_port_id = self._get_port(mgmt_port_id) + device['mgmt_port_id'] = mgmt_port_id + device['network_function_id'] = network_function_id + + device['advance_sharing_interfaces'] = ( + self._get_advance_sharing_interfaces(device['id'])) + return device + + def _get_orchestration_driver(self, service_vendor): + return self.orchestration_driver + + def plug_interfaces(self, event, is_event_call=True): + if is_event_call: + device_info = event.data + else: + device_info = event + # Get event data, as configurator sends back only request_info, which + # contains nf_id, nfi_id, nfd_id. + device = self._prepare_device_data(device_info) + self._update_network_function_device_db(device, + 'HEALTH_CHECK_COMPLETED') + orchestration_driver = self._get_orchestration_driver( + device['service_details']['service_vendor']) + + _ifaces_plugged_in, advance_sharing_ifaces = ( + orchestration_driver.plug_network_function_device_interfaces( + device)) + if _ifaces_plugged_in: + if advance_sharing_ifaces: + self._update_advance_sharing_interfaces( + device, + advance_sharing_ifaces) + self._increment_device_interface_count(device) + self._create_event(event_id='CONFIGURE_DEVICE', + event_data=device, + is_internal_event=True) + else: + self._create_event(event_id='DEVICE_CONFIGURATION_FAILED', + event_data=device, + is_internal_event=True) + + + def plug_interfaces_fast(self, event): + + # In this case, the event will be + # happening in paralell with HEALTHMONITORIN, + # so, we should not generate CONFIGURE_DEVICE & should not update + # DB with HEALTH_CHECK_COMPLETED. + + nfp_context = event.data + + service_details = nfp_context['service_details'] + network_function_device = nfp_context['network_function_device'] + token = nfp_context['resource_owner_context']['admin_token'] + tenant_id = nfp_context['resource_owner_context']['admin_tenant_id'] + + consumer = nfp_context['consumer'] + provider = nfp_context['provider'] + + orchestration_driver = self._get_orchestration_driver( + service_details['service_vendor']) + + ports = self._make_ports_dict(consumer, provider, 'port') + + device = { + 'id': network_function_device['id'], + 'ports': ports, + 'service_details': service_details, + 'token': token, + 'tenant_id': tenant_id, + 'interfaces_in_use': network_function_device['interfaces_in_use'], + 'status': network_function_device['status'], + 'vendor_data': nfp_context['vendor_data']} + + _ifaces_plugged_in, advance_sharing_ifaces = ( + orchestration_driver.plug_network_function_device_interfaces( + device)) + if _ifaces_plugged_in: + if advance_sharing_ifaces: + self._update_advance_sharing_interfaces( + device, + advance_sharing_ifaces) + self._increment_device_interface_count(device) + # REVISIT(mak) - Check how incremented ref count can be updated in + # DB + self._controller.event_complete(event, result="SUCCESS") + else: + self._create_event(event_id="PLUG_INTERFACE_FAILED", + event_data=nfp_context, + is_internal_event=True) + self._controller.event_complete(event, result="FAILED") + + + def unplug_interfaces(self, event): + device_info = event.data + device = self._prepare_device_data(device_info) + orchestration_driver = self._get_orchestration_driver( + device['service_details']['service_vendor']) + + is_interface_unplugged, advance_sharing_ifaces = ( + orchestration_driver.unplug_network_function_device_interfaces( + device)) + if is_interface_unplugged: + if advance_sharing_ifaces: + self._update_advance_sharing_interfaces( + device, + advance_sharing_ifaces) + mgmt_port_id = device['mgmt_port_id'] + self._decrement_device_interface_count(device) + device['mgmt_port_id'] = mgmt_port_id + else: + # Ignore unplug error + pass + self._create_event(event_id='DELETE_DEVICE', + event_data=device, + is_internal_event=True) + + def delete_device(self, event): + # Update status in DB, send DEVICE_DELETED event to NSO. + device = event.data + orchestration_driver = self._get_orchestration_driver( + device['service_details']['service_vendor']) + + self._decrement_device_ref_count(device) + device_ref_count = device['reference_count'] + if device_ref_count <= 0: + orchestration_driver.delete_network_function_device(device) + self._create_event(event_id='DEVICE_BEING_DELETED', + event_data=device, + is_poll_event=True, + original_event=event) + else: + desc = 'Network Service Device can be reuse' + self._update_network_function_device_db(device, + device['status'], + desc) + # DEVICE_DELETED event for NSO + self._create_event(event_id='DEVICE_DELETED', + event_data=device) + diff --git a/gbpservice/nfp/orchestrator/modules/service_orchestrator_NSD.py b/gbpservice/nfp/orchestrator/modules/service_orchestrator_NSD.py new file mode 100644 index 0000000000..e151fc05c0 --- /dev/null +++ b/gbpservice/nfp/orchestrator/modules/service_orchestrator_NSD.py @@ -0,0 +1,186 @@ +# Licensed under the Apache License, Version 2.0 (the "License"); you may +# not use this file except in compliance with the License. You may obtain +# a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, 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. + +from neutron._i18n import _LE +from neutron._i18n import _LI + +from gbpservice.nfp.common import constants as nfp_constants +from gbpservice.nfp.common import exceptions as nfp_exc +from gbpservice.nfp.common import topics as nfp_rpc_topics +from gbpservice.nfp.core import context as nfp_core_context +from gbpservice.nfp.core.event import Event +from gbpservice.nfp.core import module as nfp_api +from gbpservice.nfp.core.rpc import RpcAgent +from gbpservice.nfp.lib import transport +from gbpservice.nfp.orchestrator.config_drivers import heat_driver +from gbpservice.nfp.orchestrator.db import nfp_db_NSD as nfp_db +from gbpservice.nfp.orchestrator.openstack import openstack_driver + +import sys +import traceback + +from gbpservice.nfp.core import log as nfp_logging + + +from gbpservice.nfp.orchestrator.modules.service_orchestrator import ( + ServiceOrchestrator, events_init, rpc_init) + +LOG = nfp_logging.getLogger(__name__) + +STOP_POLLING = {'poll': False} +CONTINUE_POLLING = {'poll': True} + + +def nfp_module_init(controller, config): + events_init(controller, config, ServiceOrchestratorNSD(controller, config)) + rpc_init(controller, config) + + +class ServiceOrchestratorNSD(ServiceOrchestrator): + + def __init__(self, controller, config): + super(ServiceOrchestratorNSD, self).__init__(controller, config) + self.db_handler = nfp_db.NFPDbBaseNSD() + + def _get_network_function_instance_for_multi_service_sharing(self, + port_info): + network_function_instances = ( + self.db_handler.get_network_function_instances(self.db_session, + filters={})) + provider_port_id = None + for port in port_info: + if port['port_classification'] == 'provider': + provider_port_id = port['id'] + break + for network_function_instance in network_function_instances: + if (provider_port_id in network_function_instance['port_info'] and + network_function_instance['network_function_device_id'] + is not None): + return network_function_instance + return None + + + def create_network_function_instance(self, event): + nfp_context = event.data + + network_function = nfp_context['network_function'] + # service_profile = nfp_context['service_profile'] + service_details = nfp_context['service_details'] + consumer = nfp_context['consumer'] + provider = nfp_context['provider'] + + port_info = [] + for ele in [consumer, provider]: + if ele['pt']: + # REVISIT(ashu): Only pick few chars from id + port_info.append( + {'id': ele['pt']['id'], + 'port_model': ele['port_model'], + 'port_classification': ele['port_classification'] + }) + + # REVISIT(ashu): Only pick few chars from id + name = '%s_%s' % (network_function['name'], + network_function['id']) + network_function_instance = ( + self._get_network_function_instance_for_multi_service_sharing( + port_info)) + if network_function_instance: + port_info = [] + create_nfi_request = { + 'name': name, + 'tenant_id': network_function['tenant_id'], + 'status': nfp_constants.PENDING_CREATE, + 'network_function_id': network_function['id'], + 'service_type': service_details['service_type'], + 'service_vendor': service_details['service_vendor'], + 'share_existing_device': nfp_context['share_existing_device'], + 'port_info': port_info, + } + nfi_db = self.db_handler.create_network_function_instance( + self.db_session, create_nfi_request) + if network_function_instance: + port_info = [] + for port_id in network_function_instance['port_info']: + port_info.append(self.db_handler.get_port_info(self.db_session, + port_id)) + nfi = { + 'port_info': port_info + } + nfi_db = self.db_handler.update_network_function_instance( + self.db_session, nfi_db['id'], nfi) + nfd_data = {} + nfd_data['network_function_instance_id'] = nfi_db['id'] + nfd_data['network_function_device_id'] = ( + network_function_instance['network_function_device_id']) + self._create_event('DEVICE_ACTIVE', + event_data=nfd_data) + + return + # Sending LogMeta Details to visibility + self._report_logging_info(network_function, + nfi_db, + service_details['service_type'], + service_details['service_vendor']) + + nfp_context['network_function_instance'] = nfi_db + + LOG.info(_LI("[Event:CreateService]")) + self._create_event('CREATE_NETWORK_FUNCTION_DEVICE', + event_data=nfp_context) + + + def delete_network_function_instance(self, event): + nfi_id = event.data + nfi = {'status': nfp_constants.PENDING_DELETE} + nfi = self.db_handler.update_network_function_instance( + self.db_session, nfi_id, nfi) + if nfi['network_function_device_id']: + + filters = { + 'network_function_device_id': [ + nfi['network_function_device_id']], + 'status': ['ACTIVE'] + } + network_function_instances = ( + self.db_handler.get_network_function_instances( + self.db_session, filters=filters)) + if network_function_instances: + device_deleted_event = { + 'network_function_instance_id': nfi['id'] + } + network_function = self.db_handler.get_network_function( + self.db_session, nfi['network_function_id']) + nf_id = network_function['id'] + self.db_handler.delete_network_function( + self.db_session, nfi['network_function_id']) + LOG.info(_LI("NSO: Deleted network function: %(nf_id)s"), + {'nf_id': nf_id}) + + return + delete_nfd_request = { + 'network_function_device_id': nfi[ + 'network_function_device_id'], + 'network_function_instance': nfi, + 'network_function_id': nfi['network_function_id'] + } + self._create_event('DELETE_NETWORK_FUNCTION_DEVICE', + event_data=delete_nfd_request) + else: + device_deleted_event = { + 'network_function_instance_id': nfi['id'] + } + self._create_event('DEVICE_DELETED', + event_data=device_deleted_event, + is_internal_event=True) + +