diff --git a/gbpservice/neutron/services/servicechain/plugins/ncp/node_drivers/nfp_node_driver.py b/gbpservice/neutron/services/servicechain/plugins/ncp/node_drivers/nfp_node_driver.py index bf48502647..15422bc064 100644 --- a/gbpservice/neutron/services/servicechain/plugins/ncp/node_drivers/nfp_node_driver.py +++ b/gbpservice/neutron/services/servicechain/plugins/ncp/node_drivers/nfp_node_driver.py @@ -14,6 +14,7 @@ # limitations under the License. import eventlet +from eventlet import greenpool from keystoneclient import exceptions as k_exceptions from keystoneclient.v2_0 import client as keyclient from neutron._i18n import _LE @@ -38,6 +39,7 @@ from gbpservice.nfp.common import constants as nfp_constants from gbpservice.nfp.common import topics as nfp_rpc_topics +from gbpservice.neutron.services.grouppolicy.common import constants as gp_constants NFP_NODE_DRIVER_OPTS = [ cfg.BoolOpt('is_service_admin_owned', @@ -59,7 +61,6 @@ LOG = logging.getLogger(__name__) - class InvalidServiceType(exc.NodeCompositionPluginBadRequest): message = _("The NFP Node driver only supports the services " "VPN, Firewall and LB in a Service Chain") @@ -227,6 +228,9 @@ class NFPNodeDriver(driver_base.NodeDriverBase): def __init__(self): super(NFPNodeDriver, self).__init__() self._lbaas_plugin = None + self.thread_pool = greenpool.GreenPool(10) + self.active_threads = [] + self.sc_node_count = 0 @property def name(self): @@ -336,8 +340,18 @@ def create(self, context): self._set_node_instance_network_function_map( context.plugin_session, context.current_node['id'], context.instance['id'], network_function_id) - self._wait_for_network_function_operation_completion( - context, network_function_id, operation='create') + + # Check for NF status in a separate thread + gth = self.thread_pool.spawn(self._wait_for_network_function_operation_completion, + context, network_function_id, operation='create') + + self.active_threads.append(gth) + + # At last wait for the threads to complete, success/failure/timeout + if len(self.active_threads) == self.sc_node_count: + for gth in self.active_threads: + gth.wait() + self.active_threads = [] def update(self, context): context._plugin_context = self._get_resource_owner_context( @@ -591,23 +605,49 @@ def _get_service_targets(self, context): {'service_type': service_type}) raise Exception("Service Targets are not created for the Node") - service_target_info = {'provider_ports': [], 'provider_pts': [], - 'consumer_ports': [], 'consumer_pts': []} + service_target_info = { + 'provider_ports': [], + 'provider_subnet': None, + 'provider_pts': [], + 'provider_pt_objs': [], + 'provider_ptg': [], + 'consumer_ports': [], + 'consumer_subnet': None, + 'consumer_pts': [], + 'consumer_pt_objs': [], + 'consumer_ptg': []} + for service_target in provider_service_targets: policy_target = context.gbp_plugin.get_policy_target( context.plugin_context, service_target.policy_target_id) + policy_target_group = context.gbp_plugin.get_policy_target_group( + context.plugin_context, policy_target['policy_target_group_id']) port = context.core_plugin.get_port( context.plugin_context, policy_target['port_id']) + port['ip_address'] = port['fixed_ips'][0]['ip_address'] + subnet = context.core_plugin.get_subnet( + context.plugin_context, port['fixed_ips'][0]['subnet_id']) service_target_info['provider_ports'].append(port) + service_target_info['provider_subnet'] = subnet service_target_info['provider_pts'].append(policy_target['id']) + service_target_info['provider_pt_objs'].append(policy_target) + service_target_info['provider_ptg'].append(policy_target_group) for service_target in consumer_service_targets: policy_target = context.gbp_plugin.get_policy_target( context.plugin_context, service_target.policy_target_id) + policy_target_group = context.gbp_plugin.get_policy_target_group( + context.plugin_context, policy_target['policy_target_group_id']) port = context.core_plugin.get_port( context.plugin_context, policy_target['port_id']) + port['ip_address'] = port['fixed_ips'][0]['ip_address'] + subnet = context.core_plugin.get_subnet( + context.plugin_context, port['fixed_ips'][0]['subnet_id']) service_target_info['consumer_ports'].append(port) + service_target_info['consumer_subnet'] = subnet service_target_info['consumer_pts'].append(policy_target['id']) + service_target_info['consumer_pt_objs'].append(policy_target) + service_target_info['consumer_ptg'].append(policy_target_group) return service_target_info @@ -619,6 +659,7 @@ def _is_node_order_in_spec_supported(self, context): for spec in current_specs: node_list.extend(spec['nodes']) + self.sc_node_count = len(node_list) for node_id in node_list: node_info = context.sc_plugin.get_servicechain_node( context.plugin_context, node_id) @@ -641,9 +682,65 @@ def _is_node_order_in_spec_supported(self, context): raise InvalidNodeOrderInChain( node_order=allowed_chain_combinations) + def _get_consumers_for_provider(self, context, provider): + ''' + { + consuming_ptgs_details: [{'ptg': <>, 'subnets': <>}] + consuming_eps_details: [] + } + ''' + + consuming_ptgs_details = [] + consuming_eps_details = [] + + provided_prs_id = provider['provided_policy_rule_sets'][0] + provided_prs = context.gbp_plugin.get_policy_rule_set( + context.plugin_context, provided_prs_id) + consuming_ptg_ids = provided_prs['consuming_policy_target_groups'] + consuming_ep_ids = provided_prs['consuming_external_policies'] + + consuming_ptgs = context.gbp_plugin.get_policy_target_groups( + context.plugin_context, filters={'id':consuming_ptg_ids}) + consuming_eps_details = context.gbp_plugin.get_external_policies( + context.plugin_context, filters={'id': consuming_ep_ids}) + + for ptg in consuming_ptgs: + subnet_ids = ptg['subnets'] + subnets = context.core_plugin.get_subnets(context.plugin_context, filters={'id':subnet_ids}) + consuming_ptgs_details.append({'ptg':ptg, 'subnets':subnets}) + + return consuming_ptgs_details, consuming_eps_details + + def _create_network_function(self, context): + """ + nfp_create_nf_data :- + + {'resource_owner_context': <>, + 'service_chain_instance': <>, + 'service_chain_node': <>, + 'service_profile': <>, + 'service_config': context.current_node.get('config'), + 'provider': {'pt':<>, 'ptg':<>, 'port':<>, 'subnet':<>}, + 'consumer': {'pt':<>, 'ptg':<>, 'port':<>, 'subnet':<>}, + 'management': {'pt':<>, 'ptg':<>, 'port':<>, 'subnet':<>}, + 'management_ptg_id': <>, + 'network_function_mode': nfp_constants.GBP_MODE, + 'tenant_id': <>, + 'consuming_ptgs_details': [], + 'consuming_eps_details': [] + } + + """ + nfp_create_nf_data = {} + sc_instance = context.instance service_targets = self._get_service_targets(context) + + consuming_ptgs_details, consuming_eps_details = \ + self._get_consumers_for_provider(context, + service_targets['provider_ptg'][0]) + if context.current_profile['service_type'] == pconst.LOADBALANCER: config_param_values = sc_instance.get('config_param_values', {}) if config_param_values: @@ -661,35 +758,58 @@ def _create_network_function(self, context): context.core_plugin.update_port( context.plugin_context, provider_port['id'], port) - port_info = [] - if service_targets.get('provider_pts'): - # Device case, for Base mode ports won't be available. - port_info = [ - { - 'id': service_targets['provider_pts'][0], - 'port_model': nfp_constants.GBP_PORT, - 'port_classification': nfp_constants.PROVIDER, - } - ] - if service_targets.get('consumer_ports'): - port_info.append({ - 'id': service_targets['consumer_pts'][0], - 'port_model': nfp_constants.GBP_PORT, - 'port_classification': nfp_constants.CONSUMER, - }) - network_function = { - 'tenant_id': context.provider['tenant_id'], - 'service_chain_id': sc_instance['id'], - 'service_id': context.current_node['id'], - 'service_profile_id': context.current_profile['id'], - 'management_ptg_id': sc_instance['management_ptg_id'], + provider = { + 'pt': service_targets.get('provider_pt_objs', [None])[0], + 'ptg': service_targets.get('provider_ptg', [None])[0], + 'port': service_targets.get('provider_ports', [None])[0], + 'subnet': service_targets.get('provider_subnet', None), + 'port_model': nfp_constants.GBP_PORT, + 'port_classification': nfp_constants.PROVIDER} + + consumer_pt = None + consumer_ptg = None + consumer_ports = None + + if service_targets['consumer_pt_objs']: + consumer_pt = service_targets.get('consumer_pt_objs', [None])[0] + if service_targets['consumer_ptg']: + consumer_ptg = service_targets.get('consumer_ptg', [None])[0] + if service_targets['consumer_ports']: + consumer_ports = service_targets.get('consumer_ports', [None])[0] + + consumer = { + 'pt': consumer_pt, + 'ptg': consumer_ptg, + 'port': consumer_ports, + 'subnet': service_targets.get('consumer_subnet', None), + 'port_model': nfp_constants.GBP_PORT, + 'port_classification': nfp_constants.CONSUMER} + + management = { + 'pt': None, + 'ptg': None, + 'port': None, + 'subnet': None, + 'port_model': nfp_constants.GBP_NETWORK, + 'port_classification': nfp_constants.MANAGEMENT} + + nfp_create_nf_data = { + 'resource_owner_context': context._plugin_context.to_dict(), + 'service_chain_instance': sc_instance, + 'service_chain_node': context.current_node, + 'service_profile': context.current_profile, 'service_config': context.current_node.get('config'), - 'port_info': port_info, + 'provider': provider, + 'consumer': consumer, + 'management': management, + 'management_ptg_id': sc_instance['management_ptg_id'], 'network_function_mode': nfp_constants.GBP_MODE, - } + 'tenant_id': context.provider['tenant_id'], + 'consuming_ptgs_details': consuming_ptgs_details, + 'consuming_eps_details': consuming_eps_details} return self.nfp_notifier.create_network_function( - context.plugin_context, network_function=network_function)['id'] + context.plugin_context, network_function=nfp_create_nf_data)['id'] def _set_node_instance_network_function_map( self, session, sc_node_id, sc_instance_id, network_function_id): diff --git a/gbpservice/nfp/bin/nfp_configurator.ini b/gbpservice/nfp/bin/nfp_configurator.ini index f6f52abfb0..bcad20de30 100644 --- a/gbpservice/nfp/bin/nfp_configurator.ini +++ b/gbpservice/nfp/bin/nfp_configurator.ini @@ -12,10 +12,10 @@ kombu_reconnect_delay=1.0 rabbit_use_ssl=False rabbit_virtual_host=/ -workers=1 +workers=2 modules_dir=gbpservice.nfp.configurator.modules reportstate_interval=10 -periodic_interval=9 +periodic_interval=2 log_forward_ip_address= log_forward_port=514 diff --git a/gbpservice/nfp/bin/nfp_orch_agent.ini b/gbpservice/nfp/bin/nfp_orch_agent.ini index 2d20594cba..9ddfbccb0b 100644 --- a/gbpservice/nfp/bin/nfp_orch_agent.ini +++ b/gbpservice/nfp/bin/nfp_orch_agent.ini @@ -1,5 +1,5 @@ [DEFAULT] -workers=1 +workers=2 debug=False kombu_reconnect_delay=1.0 rabbit_use_ssl=False diff --git a/gbpservice/nfp/bin/nfp_proxy_agent.ini b/gbpservice/nfp/bin/nfp_proxy_agent.ini index 62792a47ab..8698c2a0fa 100644 --- a/gbpservice/nfp/bin/nfp_proxy_agent.ini +++ b/gbpservice/nfp/bin/nfp_proxy_agent.ini @@ -6,4 +6,4 @@ rabbit_use_ssl=False rabbit_virtual_host=/ modules_dir=gbpservice.nfp.proxy_agent.modules backend=unix_rest -periodic_interval=10 +periodic_interval=2 diff --git a/gbpservice/nfp/bin/proxy.ini b/gbpservice/nfp/bin/proxy.ini index 564cc9dda0..18babe5eea 100644 --- a/gbpservice/nfp/bin/proxy.ini +++ b/gbpservice/nfp/bin/proxy.ini @@ -5,7 +5,8 @@ max_connections=10 rest_server_address= 11.0.0.3 ##for docker ## rest_server_port= 8070 -worker_threads=40 +#[Note: worker threads should not be less than connect_max_wait_timeout/{periodic_interval or spacing for pull_notification}] +worker_threads=100 connect_max_wait_timeout=120 idle_max_wait_timeout=120 idle_min_wait_timeout=0.1 diff --git a/gbpservice/nfp/config_orchestrator/agent/otc_service_events.py b/gbpservice/nfp/config_orchestrator/agent/otc_service_events.py index 21b6f060e3..10b7693867 100644 --- a/gbpservice/nfp/config_orchestrator/agent/otc_service_events.py +++ b/gbpservice/nfp/config_orchestrator/agent/otc_service_events.py @@ -67,7 +67,7 @@ def _delete_service(self, context, resource): "DELETE", network_function_event=True) - @core_pt.poll_event_desc(event='SERVICE_CREATE_PENDING', spacing=5) + @core_pt.poll_event_desc(event='SERVICE_CREATE_PENDING', spacing=2) def create_sevice_pending_event(self, ev): event_data = ev.data ctxt = n_context.Context.from_dict(event_data['context']) diff --git a/gbpservice/nfp/configurator/agents/generic_config.py b/gbpservice/nfp/configurator/agents/generic_config.py index 7c6c53f899..b27ee0cdeb 100644 --- a/gbpservice/nfp/configurator/agents/generic_config.py +++ b/gbpservice/nfp/configurator/agents/generic_config.py @@ -21,6 +21,9 @@ from gbpservice.nfp.core import event as nfp_event from gbpservice.nfp.core import poll as nfp_poll +STOP_POLLING = {'poll': False} +CONTINUE_POLLING = {'poll': True} + LOG = nfp_logging.getLogger(__name__) """Implements APIs invoked by configurator for processing RPC messages. @@ -65,6 +68,7 @@ def _send_event(self, context, resource_data, event_id, event_key=None): arg_dict = {'context': context, 'resource_data': resource_data} ev = self.sc.new_event(id=event_id, data=arg_dict, key=event_key) + self.sc.post_event(ev) def configure_interfaces(self, context, resource_data): @@ -268,8 +272,9 @@ def _process_event(self, ev): if (resource_data.get('periodicity') == gen_cfg_const.INITIAL and result == common_const.SUCCESS): notification_data = self._prepare_notification_data(ev, result) - self.sc.poll_event_done(ev) + # self.sc.poll_event_done(ev) self.notify._notification(notification_data) + return STOP_POLLING elif resource_data.get('periodicity') == gen_cfg_const.FOREVER: if result == common_const.FAILED: """If health monitoring fails continuously for 5 times @@ -282,8 +287,9 @@ def _process_event(self, ev): notification_data = self._prepare_notification_data( ev, result) - self.sc.poll_event_done(ev) + # self.sc.poll_event_done(ev) self.notify._notification(notification_data) + return STOP_POLLING elif result == common_const.SUCCESS: """set fail_count to 0 if it had failed earlier even once """ @@ -293,8 +299,9 @@ def _process_event(self, ev): that particular service vm's health monitor """ notification_data = self._prepare_notification_data(ev, result) - self.sc.poll_event_done(ev) + # self.sc.poll_event_done(ev) self.notify._notification(notification_data) + return STOP_POLLING else: """For other events, irrespective of result send notification""" notification_data = self._prepare_notification_data(ev, result) @@ -357,7 +364,7 @@ def poll_event_cancel(self, ev): @nfp_poll.poll_event_desc( event=gen_cfg_const.EVENT_CONFIGURE_HEALTHMONITOR, - spacing=5) + spacing=2) def handle_configure_healthmonitor(self, ev): """Decorator method called for poll event CONFIGURE_HEALTHMONITOR Finally it Enqueues response into notification queue. @@ -367,7 +374,7 @@ def handle_configure_healthmonitor(self, ev): Returns: None """ - self._process_event(ev) + return self._process_event(ev) def events_init(sc, drivers, rpcmgr): diff --git a/gbpservice/nfp/configurator/api/config.py b/gbpservice/nfp/configurator/api/config.py index 28498d2070..f0d45757c6 100644 --- a/gbpservice/nfp/configurator/api/config.py +++ b/gbpservice/nfp/configurator/api/config.py @@ -48,7 +48,8 @@ 'logfile': { 'class': 'logging.FileHandler', 'filename': '/var/log/nfp/nfp_pecan.log', - 'level': 'INFO' + 'level': 'INFO', + 'formatter': 'simple' } }, 'formatters': { diff --git a/gbpservice/nfp/configurator/lib/generic_config_constants.py b/gbpservice/nfp/configurator/lib/generic_config_constants.py index 8018a7a10f..3ea05a7d59 100644 --- a/gbpservice/nfp/configurator/lib/generic_config_constants.py +++ b/gbpservice/nfp/configurator/lib/generic_config_constants.py @@ -22,4 +22,4 @@ MAX_FAIL_COUNT = 12 # 5 secs delay * 12 = 60 secs INITIAL = 'initial' FOREVER = 'forever' -INITIAL_HM_RETRIES = 24 # 5 secs delay * 24 = 120 secs +INITIAL_HM_RETRIES = 90 # 5 secs delay * 24 = 120 secs diff --git a/gbpservice/nfp/core/context.py b/gbpservice/nfp/core/context.py new file mode 100644 index 0000000000..e002ca322b --- /dev/null +++ b/gbpservice/nfp/core/context.py @@ -0,0 +1,34 @@ +# 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. + +import threading + +nfp_context_store = threading.local() + +class NfpContext(object): + def __init__(self, context): + self.context = context + + def get_context(self): + return self.context + +def store_nfp_context(context): + nfp_context_store.context = NfpContext(context) + +def clear_nfp_context(): + nfp_context_store.context = None + +def get_nfp_context(): + context = getattr(nfp_context_store, 'context', None) + if context: + return context.get_context() + return {} diff --git a/gbpservice/nfp/core/controller.py b/gbpservice/nfp/core/controller.py index 36764cccc0..0e0d15be7b 100644 --- a/gbpservice/nfp/core/controller.py +++ b/gbpservice/nfp/core/controller.py @@ -345,15 +345,20 @@ def get_stashed_events(self): """ events = [] if self._process_name == 'distributor-process': + # return at max 5 events + maxx = 5 # wait sometime for first event in the queue timeout = 0.1 - try: - event = self._stashq.get(timeout=timeout) - self.decompress(event) - events.append(event) - timeout = 0 - except Queue.Empty: - pass + while maxx: + try: + event = self._stashq.get(timeout=timeout) + self.decompress(event) + events.append(event) + timeout = 0 + maxx -= 1 + except Queue.Empty: + maxx = 0 + pass else: LOG.error("worker cannot pull stashed events") return events diff --git a/gbpservice/nfp/core/task.py b/gbpservice/nfp/core/task.py new file mode 100644 index 0000000000..efaaab2661 --- /dev/null +++ b/gbpservice/nfp/core/task.py @@ -0,0 +1,73 @@ +# 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 gbpservice.nfp.core import threadpool as core_tp + +class InUse(Exception): + + """Exception raised when same executor is fired twice or jobs + added after executor is fired. + """ + pass + +def check_in_use(f): + def wrapped(self, *args, **kwargs): + if self.fired: + raise InUse("Executor in use") + return f(self, *args, **kwargs) + return wrapped + + +class TaskExecutor(object): + def __init__(self, jobs=0): + if not jobs: + self.thread_pool = core_tp.ThreadPool() + else: + self.thread_pool = core_tp.ThreadPool(thread_pool_size=jobs) + + self.pipe_line = [] + self.fired = False + + @check_in_use + def add_job(self, id, func, *args, **kwargs): + result_store = kwargs.pop('result_store', None) + + job = { + 'id':id, 'method': func, 'args': args, + 'kwargs': kwargs} + + if result_store is not None: + job.update({'result_store': result_store}) + + self.pipe_line.append(job) + + def _complete(self): + self.pipe_line = [] + self.fired = False + + @check_in_use + def fire(self): + self.fired = True + for job in self.pipe_line: + th = self.thread_pool.dispatch(job['method'], *job['args'], **job['kwargs']) + job['thread'] = th + + for job in self.pipe_line: + result = job['thread'].wait() + job.pop('thread') + job['result'] = result + if 'result_store' in job.keys(): + job['result_store']['result'] = result + + done_jobs = self.pipe_line[:] + self._complete() + return done_jobs diff --git a/gbpservice/nfp/lib/RestClientOverUnix.py b/gbpservice/nfp/lib/RestClientOverUnix.py index adf0059e91..15a3729db2 100644 --- a/gbpservice/nfp/lib/RestClientOverUnix.py +++ b/gbpservice/nfp/lib/RestClientOverUnix.py @@ -82,7 +82,6 @@ def send_request(self, path, method_type, request_method='http', if method_type.upper() != 'GET': body = jsonutils.dumps(body) body = zlib.compress(body) - path = '/v1/nfp/' + path url = urlparse.urlunsplit(( request_method, diff --git a/gbpservice/nfp/lib/transport.py b/gbpservice/nfp/lib/transport.py index 575aa1a5de..373bd4f2e9 100644 --- a/gbpservice/nfp/lib/transport.py +++ b/gbpservice/nfp/lib/transport.py @@ -203,6 +203,7 @@ def send_request_to_configurator(conf, context, body, elif conf.backend == UNIX_REST: try: + LOG.info("#############sends post_notification##############") resp, content = unix_rc.post(method_name, body=body) LOG.info( @@ -257,6 +258,7 @@ def get_response_from_configurator(conf): elif conf.backend == UNIX_REST: try: + LOG.info("#############sends get_notification##############") resp, content = unix_rc.get('get_notifications') content = jsonutils.loads(content) if content: diff --git a/gbpservice/nfp/orchestrator/coal/networking/nfp_gbp_network_driver.py b/gbpservice/nfp/orchestrator/coal/networking/nfp_gbp_network_driver.py index ca00982b7b..0cfb91756b 100644 --- a/gbpservice/nfp/orchestrator/coal/networking/nfp_gbp_network_driver.py +++ b/gbpservice/nfp/orchestrator/coal/networking/nfp_gbp_network_driver.py @@ -15,7 +15,6 @@ nfp_neutron_network_driver as neutron_nd ) - class NFPGBPNetworkDriver(neutron_nd.NFPNeutronNetworkDriver): def __init__(self, config): self.config = config @@ -42,12 +41,19 @@ def update_port(self, token, port_id, port): port) return pt['port_id'] + def get_neutron_port_details(self, token, port_id): + #self.network_handler = openstack_driver.NeutronClient(self.config) + port_details = super(NFPGBPNetworkDriver, self).get_port_and_subnet_details( + token, port_id) + #self.network_handler = openstack_driver.GBPClient(self.config) + return port_details + def get_port_details(self, token, port_id): _port_id = self.get_port_id(token, port_id) - self.network_handler = openstack_driver.NeutronClient(self.config) + #self.network_handler = openstack_driver.NeutronClient(self.config) port_details = super(NFPGBPNetworkDriver, self).get_port_details( token, _port_id) - self.network_handler = openstack_driver.GBPClient(self.config) + #self.network_handler = openstack_driver.GBPClient(self.config) return port_details def get_networks(self, token, filters): @@ -56,6 +62,11 @@ def get_networks(self, token, filters): def set_promiscuos_mode(self, token, port_id): port_id = self.get_port_id(token, port_id) - self.network_handler = openstack_driver.NeutronClient(self.config) + #self.network_handler = openstack_driver.NeutronClient(self.config) + super(NFPGBPNetworkDriver, self).set_promiscuos_mode(token, port_id) + #self.network_handler = openstack_driver.GBPClient(self.config) + + def set_promiscuos_mode_v1(self, token, port_id): + #self.network_handler = openstack_driver.NeutronClient(self.config) super(NFPGBPNetworkDriver, self).set_promiscuos_mode(token, port_id) - self.network_handler = openstack_driver.GBPClient(self.config) + #self.network_handler = openstack_driver.GBPClient(self.config) diff --git a/gbpservice/nfp/orchestrator/coal/networking/nfp_neutron_network_driver.py b/gbpservice/nfp/orchestrator/coal/networking/nfp_neutron_network_driver.py index 0ee00dcf84..10c871dd8e 100644 --- a/gbpservice/nfp/orchestrator/coal/networking/nfp_neutron_network_driver.py +++ b/gbpservice/nfp/orchestrator/coal/networking/nfp_neutron_network_driver.py @@ -18,28 +18,46 @@ class NFPNeutronNetworkDriver(ndb.NFPNetworkDriverBase): def __init__(self, config): - self.network_handler = openstack_driver.NeutronClient(config) + # self.network_handler = openstack_driver.NeutronClient(config) + self.neutron_client = openstack_driver.NeutronClient(config) def setup_traffic_steering(self): pass def create_port(self, token, tenant_id, net_id, name=None): - port = self.network_handler.create_port(token, tenant_id, net_id, + port = self.neutron_client.create_port(token, tenant_id, net_id, attrs={'name': name}) return port def delete_port(self, token, port_id): - self.network_handler.delete_port(token, port_id) + self.neutron_client.delete_port(token, port_id) def get_port_id(self, token, port_id): return port_id def update_port(self, token, port_id, port): - port = self.network_handler.update_port(token, port_id, port) + port = self.neutron_client.update_port(token, port_id, port) return port['port'] + def get_port_and_subnet_details(self, token, port_id): + port = self.neutron_client.get_port(token, port_id) + + # ip + ip = port['port']['fixed_ips'][0]['ip_address'] + + # mac + mac = port['port']['mac_address'] + + # gateway ip + subnet_id = port['port']['fixed_ips'][0]['subnet_id'] + subnet = self.neutron_client.get_subnet(token, subnet_id) + cidr = subnet['subnet']['cidr'] + gateway_ip = subnet['subnet']['gateway_ip'] + + return (ip, mac, cidr, gateway_ip, port, subnet) + def get_port_details(self, token, port_id): - port = self.network_handler.get_port(token, port_id) + port = self.neutron_client.get_port(token, port_id) # ip ip = port['port']['fixed_ips'][0]['ip_address'] @@ -49,13 +67,13 @@ def get_port_details(self, token, port_id): # gateway ip subnet_id = port['port']['fixed_ips'][0]['subnet_id'] - subnet = self.network_handler.get_subnet(token, subnet_id) + subnet = self.neutron_client.get_subnet(token, subnet_id) cidr = subnet['subnet']['cidr'] gateway_ip = subnet['subnet']['gateway_ip'] return (ip, mac, cidr, gateway_ip) def set_promiscuos_mode(self, token, port_id): - self.network_handler.update_port(token, port_id, + self.neutron_client.update_port(token, port_id, security_groups=[], port_security_enabled=False) diff --git a/gbpservice/nfp/orchestrator/config_drivers/heat_client.py b/gbpservice/nfp/orchestrator/config_drivers/heat_client.py index b0b241d75c..f243197c8e 100644 --- a/gbpservice/nfp/orchestrator/config_drivers/heat_client.py +++ b/gbpservice/nfp/orchestrator/config_drivers/heat_client.py @@ -13,7 +13,7 @@ # under the License. from heatclient import client as heat_client from heatclient import exc as heat_exc -from neutron._i18n import _LW +from neutron.i18n import _LW from gbpservice.nfp.core import log as nfp_logging LOG = nfp_logging.getLogger(__name__) diff --git a/gbpservice/nfp/orchestrator/config_drivers/heat_driver.py b/gbpservice/nfp/orchestrator/config_drivers/heat_driver.py index 7d4e92f704..c6548c70db 100644 --- a/gbpservice/nfp/orchestrator/config_drivers/heat_driver.py +++ b/gbpservice/nfp/orchestrator/config_drivers/heat_driver.py @@ -32,9 +32,9 @@ from heatclient import exc as heat_exc from keystoneclient import exceptions as k_exceptions -from neutron._i18n import _LE -from neutron._i18n import _LI -from neutron._i18n import _LW +from neutron.i18n import _LE +from neutron.i18n import _LI +from neutron.i18n import _LW from neutron.plugins.common import constants as pconst from oslo_config import cfg from oslo_serialization import jsonutils @@ -43,6 +43,8 @@ from gbpservice.nfp.core import log as nfp_logging + +from gbpservice.nfp.core import threadpool as core_tp HEAT_DRIVER_OPTS = [ cfg.StrOpt('svc_management_ptg_name', default='svc_management_ptg', @@ -96,6 +98,14 @@ def __init__(self, config): self.neutron_client = NeutronClient(config) # self.resource_owner_tenant_id = None + keystone_conf = cfg.CONF.keystone_authtoken + keystone_version = keystone_conf.auth_version + self.v2client = self.keystoneclient._get_v2_keystone_admin_client() + self.admin_id = self.v2client.users.find(name=keystone_conf.admin_user).id + self.admin_role = self._get_role_by_name(self.v2client, "admin", keystone_version) + self.heat_role = self._get_role_by_name(self.v2client, "heat_stack_owner", keystone_version) + + ''' @property def resource_owner_tenant_id(self): @@ -133,7 +143,9 @@ def _get_resource_owner_context(self): self.keystoneclient.get_keystone_creds() auth_token = self.keystoneclient.get_scoped_keystone_token( user, pwd, tenant_name, tenant_id) - return auth_token, tenant_id + + tenant_id = self.keystoneclient.get_tenant_id(auth_token, tenant_name) + return auth_token, tenant_id def _get_role_by_name(self, keystone_client, name, keystone_version): if keystone_version == 'v2.0': @@ -155,29 +167,23 @@ def get_allocated_roles(self, v2client, user, tenant_id=None): allocated_role_names.append(role.name) return allocated_role_names + def _assign_admin_user_to_project_v2(self, project_id): + allocated_role_names = self.get_allocated_roles(self.v2client, self.admin_id, project_id) + if self.admin_role: + if self.admin_role.name not in allocated_role_names: + self.v2client.roles.add_user_role( + self.admin_id, self.admin_role.id, tenant=project_id) + if self.heat_role: + if self.heat_role.name not in allocated_role_names: + self.v2client.roles.add_user_role(self.admin_id, self.heat_role.id, + tenant=project_id) + def _assign_admin_user_to_project(self, project_id): keystone_conf = cfg.CONF.keystone_authtoken keystone_version = keystone_conf.auth_version if keystone_version == 'v2.0': - v2client = self.keystoneclient._get_v2_keystone_admin_client() - admin_id = v2client.users.find(name=keystone_conf.admin_user).id - admin_role = self._get_role_by_name(v2client, "admin", - keystone_version) - allocated_role_names = self.get_allocated_roles( - v2client, admin_id, project_id) - - if admin_role: - if admin_role.name not in allocated_role_names: - v2client.roles.add_user_role( - admin_id, admin_role.id, tenant=project_id) - - heat_role = self._get_role_by_name(v2client, "heat_stack_owner", - keystone_version) - if heat_role: - if heat_role.name not in allocated_role_names: - v2client.roles.add_user_role(admin_id, heat_role.id, - tenant=project_id) + return self._assign_admin_user_to_project_v2(project_id) else: v3client = self.keystoneclient._get_v3_keystone_admin_client() admin_id = v3client.users.find(name=keystone_conf.admin_user).id @@ -200,6 +206,37 @@ def keystone(self, user, pwd, tenant_name, tenant_id=None): return self.keystoneclient.get_scoped_keystone_token( user, pwd, tenant_name) + def _get_heat_client_v1(self, tenant_id, assign_admin=False): + if assign_admin: + try: + self._assign_admin_user_to_project(tenant_id) + except Exception: + LOG.exception(_LE("Failed to assign admin user to project")) + return None + + user, password, tenant, auth_url =\ + self.keystoneclient.get_keystone_creds() + + auth_token = self.keystone(user, password, tenant, tenant_id=tenant_id) + + timeout_mins, timeout_seconds = divmod(STACK_ACTION_WAIT_TIME, 60) + if timeout_seconds: + timeout_mins = timeout_mins + 1 + try: + heat_client = HeatClient( + user, + tenant_id, + cfg.CONF.heat_driver.heat_uri, + password, + auth_token=auth_token, + timeout_mins=timeout_mins) + except Exception: + LOG.exception(_LE("Failed to create heatclient object")) + return None + + return heat_client + + def _get_heat_client(self, resource_owner_tenant_id, tenant_id=None): user_tenant_id = tenant_id or resource_owner_tenant_id try: @@ -207,6 +244,7 @@ def _get_heat_client(self, resource_owner_tenant_id, tenant_id=None): except Exception: LOG.exception(_LE("Failed to assign admin user to project")) return None + user, password, tenant, auth_url =\ self.keystoneclient.get_keystone_creds() admin_token = self.keystone( @@ -445,6 +483,70 @@ def _get_all_heat_resource_keys(self, template_resource_dict, resource_keys.append(key) return resource_keys + def _create_firewall_template(self, auth_token, service_details, stack_template): + provider = service_details['provider_ptg'] + + consuming_ptgs_details = service_details['consuming_ptgs_details'] + consumer_eps = service_details['consuming_external_policies'] + + if (not consuming_ptgs_details) and (not consumer_eps): + return None + + is_template_aws_version = stack_template.get( + 'AWSTemplateFormatVersion', False) + resources_key = 'Resources' if is_template_aws_version else 'resources' + properties_key = ('Properties' if is_template_aws_version + else 'properties') + fw_rule_keys = self._get_all_heat_resource_keys( + stack_template[resources_key], is_template_aws_version, + 'OS::Neutron::FirewallRule') + fw_policy_key = self._get_all_heat_resource_keys( + stack_template['resources'], is_template_aws_version, + 'OS::Neutron::FirewallPolicy')[0] + + provider_subnet = service_details['provider_subnet'] + provider_cidr = provider_subnet['cidr'] + + fw_template_properties = dict( + resources_key=resources_key, properties_key=properties_key, + is_template_aws_version=is_template_aws_version, + fw_rule_keys=fw_rule_keys, + fw_policy_key=fw_policy_key) + + for consumer in consuming_ptgs_details: + ptg = consumer['ptg'] + subnets = consumer['subnets'] + + # Skip the stitching PTG + if ptg['proxied_group_id']: + continue + + fw_template_properties.update({'name': ptg['id'][:3]}) + for subnet in subnets: + if subnet['name'].startswith(APIC_OWNED_RES): + continue + + consumer_cidr = subnet['cidr'] + self._append_firewall_rule(stack_template, + provider_cidr, consumer_cidr, + fw_template_properties, ptg['id']) + + for consumer_ep in consumer_eps: + fw_template_properties.update({'name': consumer_ep['id'][:3]}) + self._append_firewall_rule(stack_template, provider_cidr, + "0.0.0.0/0", fw_template_properties, + consumer_ep['id']) + + for rule_key in fw_rule_keys: + del stack_template[resources_key][rule_key] + stack_template[resources_key][fw_policy_key][ + properties_key]['firewall_rules'].remove( + {'get_resource': rule_key}) + + return stack_template + + + def _update_firewall_template(self, auth_token, provider, stack_template): consumer_ptgs, consumer_eps = self._get_consumers_for_chain( auth_token, provider) @@ -606,6 +708,121 @@ def _get_site_conn_keys(self, template_resource_dict, keys.append(key) return keys + def _create_node_config_data(self, auth_token, tenant_id, service_chain_node, service_chain_instance, + provider, provider_port, consumer, consumer_port, network_function, + mgmt_ip, service_details): + + nf_desc = None + common_desc = {'network_function_id': network_function['id']} + + service_type = service_details['service_details']['service_type'] + service_vendor = service_details['service_details']['service_vendor'] + device_type = service_details['service_details']['device_type'] + base_mode_support = (True if device_type == 'None' + else False) + + + _, stack_template_str = self.parse_template_config_string( + service_chain_node.get('config')) + try: + stack_template = (jsonutils.loads(stack_template_str) if + stack_template_str.startswith('{') else + yaml.load(stack_template_str)) + except Exception: + LOG.error(_LE( + "Unable to load stack template for service chain " + "node: %(node_id)s") % {'node_id': service_chain_node}) + return None, None + config_param_values = service_chain_instance.get( + 'config_param_values', '{}') + stack_params = {} + try: + config_param_values = jsonutils.loads(config_param_values) + except Exception: + LOG.error(_LE("Unable to load config parameters")) + return None, None + + is_template_aws_version = stack_template.get( + 'AWSTemplateFormatVersion', False) + resources_key = ('Resources' if is_template_aws_version + else 'resources') + parameters_key = ('Parameters' if is_template_aws_version + else 'parameters') + properties_key = ('Properties' if is_template_aws_version + else 'properties') + + if not base_mode_support: + provider_port_mac = provider_port['mac_address'] + provider_cidr = service_details['provider_subnet']['cidr'] + provider_subnet = service_details['provider_subnet'] + else: + provider_port_mac = '' + provider_cidr = '' + standby_provider_port_mac = None + + if service_type == pconst.LOADBALANCER: + self._generate_pool_members( + auth_token, stack_template, config_param_values, + provider, is_template_aws_version) + config_param_values['Subnet'] = provider_subnet['id'] + config_param_values['service_chain_metadata'] = "" + if not base_mode_support: + config_param_values[ + 'service_chain_metadata'] = str(common_desc) + nf_desc = str((SC_METADATA % (service_chain_instance['id'], + mgmt_ip, + provider_port_mac, + standby_provider_port_mac, + network_function['id'], + service_vendor))) + + lb_pool_key = self._get_heat_resource_key( + stack_template[resources_key], + is_template_aws_version, + 'OS::Neutron::Pool') + stack_template[resources_key][lb_pool_key][properties_key][ + 'description'] = str(common_desc) + elif service_type == pconst.FIREWALL: + stack_template = self._create_firewall_template(auth_token, service_details, stack_template) + + if not stack_template: + return None, None + self._modify_fw_resources_name( + stack_template, provider, is_template_aws_version) + if not base_mode_support: + firewall_desc = {'vm_management_ip': mgmt_ip, + 'provider_ptg_info': [provider_port_mac], + 'provider_cidr': provider_cidr, + 'service_vendor': service_vendor, + 'network_function_id': network_function[ + 'id']} + + fw_key = self._get_heat_resource_key( + stack_template[resources_key], + is_template_aws_version, + 'OS::Neutron::Firewall') + stack_template[resources_key][fw_key][properties_key][ + 'description'] = str(common_desc) + + nf_desc = str(firewall_desc) + + if nf_desc: + network_function['description'] = network_function[ + 'description'] + '\n' + nf_desc + + for parameter in stack_template.get(parameters_key) or []: + if parameter in config_param_values: + stack_params[parameter] = config_param_values[parameter] + + LOG.info(_LI('Final stack_template : %(stack_data)s, ' + 'stack_params : %(params)s') % + {'stack_data': stack_template, 'params': stack_params}) + return (stack_template, stack_params) + + + + + def _update_node_config(self, auth_token, tenant_id, service_profile, service_chain_node, service_chain_instance, provider, consumer_port, network_function, @@ -855,6 +1072,82 @@ def parse_template_config_string(self, config_str): tag_str = nfp_constants.HEAT_CONFIG_TAG return tag_str, service_config + def get_servicechain_node(self, gbp, admin_token, service_id, result): + servicechain_node = gbp.get_servicechain_node(admin_token, service_id) + result['result'] = servicechain_node + + def get_servicechain_instance(self, gbp, admin_token, service_chain_id, result): + servicechain_instance = gbp.get_servicechain_instance(admin_token, service_chain_id) + result['result'] = servicechain_instance + + def get_ptg(self, gbp, admin_token, pt_id, result): + policy_target = gbp.get_policy_target( + admin_token, pt_id) + policy_target_group = gbp.get_policy_target_group( + admin_token, + policy_target['policy_target_group_id']) + result['result'] = policy_target_group + + def get_provider_details(self, gbp, admin_token, pt_id, result): + l_result = {} + self.get_ptg(gbp, admin_token, pt_id, l_result) + ptg = l_result['result'] + result['ptg'] = ptg + _,consuming_eps = self._get_consumers_for_chain(admin_token, ptg) + result['consuming_eps'] = consuming_eps + + def get_provider_consumer_details(self, gbp, admin_token, nfp_context, thread_pool, result): + nfp_device_data = nfp_context['nfp_device_data'] + consumer_th = None + consumer_result = {} + consumer_port = None + consumer_subnet = None + for port_info in nfp_device_data.get('ports'): + port_classification = None + if port_info['port_model'] == nfp_constants.GBP_PORT: + policy_target_id = port_info['id'] + port_classification = port_info['port_classification'] + port_id = port_info['port_id'] + else: + port_id = port_info['id'] + + if port_classification == nfp_constants.CONSUMER: + consumer_port = port_info['neutron_info']['port'] + consumer_subnet = port_info['neutron_info']['subnet'] + + # consumer_result = {} + consumer_th = thread_pool.dispatch(self.get_ptg, self.gbp_client, admin_token, policy_target_id, consumer_result) + + elif port_classification == nfp_constants.PROVIDER: + LOG.info(_LI("provider info: %s") % (port_id)) + provider_port = port_info['neutron_info']['port'] + provider_subnet = port_info['neutron_info']['subnet'] + provider_policy_target_group = None + + provider_result = {} + provider_th = thread_pool.dispatch(self.get_provider_details, self.gbp_client, admin_token, policy_target_id, provider_result) + + if consumer_th: + consumer_th.wait() + provider_th.wait() + + consumer_policy_target_group = consumer_result.get('result', None) + provider_policy_target_group = provider_result.get('ptg', None) + consuming_external_policies = provider_result.get('consuming_eps', None) + + result['consumer'] = {} + result['provider'] = {} + + result['consumer']['ptg'] = consumer_policy_target_group + result['consumer']['port'] = consumer_port + result['consumer']['subnet'] = consumer_subnet + + result['provider']['ptg'] = provider_policy_target_group + result['provider']['port'] = provider_port + result['provider']['subnet'] = provider_subnet + result['provider']['consuming_eps'] = consuming_external_policies + + def get_service_details(self, network_function_details): db_handler = nfp_db.NFPDbBase() db_session = nfp_db_api.get_session() @@ -1071,6 +1364,42 @@ def is_config_complete(self, stack_id, tenant_id, LOG.exception(_LE("Retrieving the stack %(stack)s failed."), {'stack': stack_id}) return failure_status + + def check_config_complete(self, nfp_context): + success_status = "COMPLETED" + failure_status = "ERROR" + intermediate_status = "IN_PROGRESS" + + provider_tenant_id = nfp_context['tenant_id'] + stack_id = nfp_context['heat_stack_id'] + + heatclient = self._get_heat_client_v1(provider_tenant_id) + if not heatclient: + return failure_status + try: + stack = heatclient.get(stack_id) + if stack.stack_status == 'DELETE_FAILED': + return failure_status + elif stack.stack_status == 'CREATE_COMPLETE': + return success_status + elif stack.stack_status == 'UPDATE_COMPLETE': + return success_status + elif stack.stack_status == 'DELETE_COMPLETE': + LOG.info(_LI("Stack %(stack)s is deleted"), + {'stack': stack_id}) + return failure_status + elif stack.stack_status == 'CREATE_FAILED': + return failure_status + elif stack.stack_status == 'UPDATE_FAILED': + return failure_status + elif stack.stack_status not in [ + 'UPDATE_IN_PROGRESS', 'CREATE_IN_PROGRESS', + 'DELETE_IN_PROGRESS']: + return intermediate_status + except Exception: + LOG.exception(_LE("Retrieving the stack %(stack)s failed."), + {'stack': stack_id}) + return failure_status def is_config_delete_complete(self, stack_id, tenant_id): success_status = "COMPLETED" @@ -1105,6 +1434,43 @@ def is_config_delete_complete(self, stack_id, tenant_id): {'stack': stack_id}) return failure_status + def get_service_details_from_nfp_context(self, nfp_context): + network_function = nfp_context['network_function'] + network_function_instance = nfp_context['network_function_instance'] + service_details = nfp_context['service_details'] + mgmt_ip = nfp_context['management']['port']['ip_address'] + heat_stack_id = network_function['heat_stack_id'] + service_id = network_function['service_id'] + service_chain_id = network_function['service_chain_id'] + servicechain_instance = nfp_context['service_chain_instance'] + servicechain_node = nfp_context['service_chain_node'] + + consumer_policy_target_group = nfp_context['consumer']['ptg'] + provider_policy_target_group = nfp_context['provider']['ptg'] + provider_port = nfp_context['provider']['port'] + provider_subnet = nfp_context['provider']['subnet'] + consumer_port = nfp_context['consumer']['port'] + consumer_subnet = nfp_context['consumer']['subnet'] + service_details['consuming_external_policies'] = nfp_context['consuming_eps_details'] + service_details['consuming_ptgs_details'] = nfp_context['consuming_ptgs_details'] + + return { + 'service_profile': None, + 'service_details': service_details, + 'servicechain_node': servicechain_node, + 'servicechain_instance': servicechain_instance, + 'consumer_port': consumer_port, + 'consumer_subnet': consumer_subnet, + 'provider_port': provider_port, + 'provider_subnet': provider_subnet, + 'mgmt_ip': mgmt_ip, + 'heat_stack_id': heat_stack_id, + 'provider_ptg': provider_policy_target_group, + 'consumer_ptg': consumer_policy_target_group, + 'consuming_external_policies': service_details['consuming_external_policies'], + 'consuming_ptgs_details': service_details['consuming_ptgs_details'] + } + def apply_config(self, network_function_details): service_details = self.get_service_details(network_function_details) service_profile = service_details['service_profile'] @@ -1161,6 +1527,65 @@ def apply_config(self, network_function_details): return stack_id + + def apply_heat_config(self, nfp_context): + service_details = self.get_service_details_from_nfp_context(nfp_context) + + network_function = nfp_context['network_function'] + service_profile = service_details['service_profile'] + service_chain_node = service_details['servicechain_node'] + service_chain_instance = service_details['servicechain_instance'] + provider = service_details['provider_ptg'] + consumer = service_details['consumer_ptg'] + consumer_port = service_details['consumer_port'] + provider_port = service_details['provider_port'] + mgmt_ip = service_details['mgmt_ip'] + + auth_token = nfp_context['resource_owner_context']['admin_token'] + provider_tenant_id = nfp_context['tenant_id'] + heatclient = self._get_heat_client_v1(provider_tenant_id, + assign_admin=True) + if not heatclient: + return None + + stack_template, stack_params = self._create_node_config_data( + auth_token, provider_tenant_id, + service_chain_node, service_chain_instance, + provider, provider_port, consumer, consumer_port, + network_function, mgmt_ip, service_details) + + if not stack_template and not stack_params: + return None + + if not heatclient: + return None + + stack_name = ("stack_" + service_chain_instance['name'] + + service_chain_node['name'] + + service_chain_instance['id'][:8] + + service_chain_node['id'][:8] + '-' + + time.strftime("%Y%m%d%H%M%S")) + # Heat does not accept space in stack name + stack_name = stack_name.replace(" ", "") + + try: + stack = heatclient.create(stack_name, stack_template, stack_params) + except Exception as err: + LOG.error(_LE("Heat stack creation failed for template : " + "%(template)s and stack parameters : %(params)s " + "with Error: %(error)s") % + {'template': stack_template, 'params': stack_params, + 'error': err}) + return None + + stack_id = stack['stack']['id'] + LOG.info(_LI("Created stack with ID %(stack_id)s and " + "name %(stack_name)s for provider PTG %(provider)s"), + {'stack_id': stack_id, 'stack_name': stack_name, + 'provider': provider['id']}) + + return stack_id, heatclient + def delete_config(self, stack_id, tenant_id): auth_token, resource_owner_tenant_id =\ self._get_resource_owner_context() diff --git a/gbpservice/nfp/orchestrator/db/nfp_db.py b/gbpservice/nfp/orchestrator/db/nfp_db.py index 775b2b87f0..1fe6bcaafe 100644 --- a/gbpservice/nfp/orchestrator/db/nfp_db.py +++ b/gbpservice/nfp/orchestrator/db/nfp_db.py @@ -202,7 +202,7 @@ def _set_mgmt_port_for_nfd(self, session, network_function_device_db, session.add(port_info_db) session.flush() nfd_db.mgmt_port_id = port_info_db['id'] - del network_function_device['mgmt_port_id'] + # del network_function_device['mgmt_port_id'] def _set_plugged_in_port_for_nfd_interface(self, session, nfd_interface_db, interface, is_update=False): @@ -324,7 +324,11 @@ def update_network_function_device(self, session, session, network_function_device_db, updated_network_function_device) + mgmt_port_id = updated_network_function_device.pop('mgmt_port_id', None) + if mgmt_port_id: + updated_network_function_device['mgmt_port_id'] = mgmt_port_id['id'] network_function_device_db.update(updated_network_function_device) + updated_network_function_device['mgmt_port_id'] = mgmt_port_id return self._make_network_function_device_dict( network_function_device_db) diff --git a/gbpservice/nfp/orchestrator/drivers/orchestration_driver.py b/gbpservice/nfp/orchestrator/drivers/orchestration_driver.py index f3dcdee59b..b3bcbc369d 100644 --- a/gbpservice/nfp/orchestrator/drivers/orchestration_driver.py +++ b/gbpservice/nfp/orchestrator/drivers/orchestration_driver.py @@ -1,4 +1,5 @@ # 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 # @@ -27,6 +28,9 @@ import operator from gbpservice.nfp.core import log as nfp_logging + +from gbpservice.nfp.core import task as core_task + LOG = nfp_logging.getLogger(__name__) @@ -91,20 +95,7 @@ def _get_setup_mode(self, config): return {nfp_constants.NEUTRON_MODE: True} def _get_admin_tenant_id(self, token=None): - try: - (dummy, - dummy, - admin_tenant_name, - dummy) = self.identity_handler.get_keystone_creds() - if not token: - token = self.identity_handler.get_admin_token() - admin_tenant_id = self.identity_handler.get_tenant_id( - token, - admin_tenant_name) - return admin_tenant_id - except Exception: - LOG.error(_LE("Failed to get admin's tenant ID")) - raise + return self.identity_handler.get_admin_tenant_id(token) def _get_token(self, device_data_token): @@ -145,19 +136,16 @@ def _decrement_stats_counter(self, metric, by=1): def _is_device_sharing_supported(self): return self.supports_device_sharing - def _create_management_interface(self, device_data, network_handler=None): - token = self._get_token(device_data.get('token')) - if not token: - return None - + def _create_management_interface(self, token, admin_tenant_id, device_data, network_handler): name = nfp_constants.MANAGEMENT_INTERFACE_NAME mgmt_interface = network_handler.create_port( token, - self._get_admin_tenant_id(token=token), + admin_tenant_id, device_data['management_network_info']['id'], name=name) return {'id': mgmt_interface['id'], + 'port_id': mgmt_interface['port_id'], 'port_model': (nfp_constants.GBP_PORT if device_data['service_details'][ 'network_mode'] == @@ -227,12 +215,17 @@ def _create_advance_sharing_interfaces(self, device_data, port['id'])}) return port_infos - def _get_interfaces_for_device_create(self, device_data, - network_handler=None): - mgmt_interface = self._create_management_interface( - device_data, - network_handler=network_handler) - return [mgmt_interface] + def _get_interfaces_for_device_create(self, token, admin_tenant_id, network_handler, device_data): + try: + mgmt_interface = self._create_management_interface( + token, + admin_tenant_id, + device_data, + network_handler) + device_data['interfaces'] = [mgmt_interface] + except Exception as e: + LOG.exception(_LE('Failed to get interfaces for device creation.' + 'Error: %(error)s'), {'error', e}) def _delete_interfaces(self, device_data, interfaces, network_handler=None): @@ -275,6 +268,24 @@ def _get_vendor_data(self, device_data, image_name): return None return vendor_data + def _get_vendor_data_v1(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_self_with_vendor_data(self, vendor_data, attr): attr_value = getattr(self, attr) if attr in vendor_data: @@ -306,6 +317,27 @@ def _update_vendor_data(self, device_data, token=None): " proceeding with default values") % (image_name)) + + def _update_vendor_data_v1(self, token, admin_tenant_id, image_name, device_data): + try: + vendor_data = self._get_vendor_data_v1(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: %s," + " proceeding with default values") + % (image_name)) + def _get_image_name(self, device_data): if device_data['service_details'].get('image_name'): image_name = device_data['service_details']['image_name'] @@ -316,6 +348,7 @@ def _get_image_name(self, device_data): % (device_data['service_details']['service_vendor'])) image_name = device_data['service_details']['service_vendor'] image_name = '%s' % image_name.lower() + device_data['service_details']['image_name'] = image_name return image_name def get_network_function_device_sharing_info(self, device_data): @@ -348,12 +381,9 @@ def get_network_function_device_sharing_info(self, device_data): ): raise exceptions.IncompleteData() - 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 + return { 'filters': { 'tenant_id': [device_data['tenant_id']], @@ -418,6 +448,55 @@ def select_network_function_device(self, devices, device_data): return device return None + def get_image_id(self, nova, token, admin_tenant_id, image_name): + try: + image_id = nova.get_image_id(token, admin_tenant_id, image_name) + return image_id + except Exception as e: + LOG.error(_LE('Failed to get image id for device creation.' + ' image name: %(image_name)s. Error: %(error)s'), + {'image_name': image_name, 'error': e}) + + + def create_instance(self, nova, token, admin_tenant_id, + image_id, flavor, interfaces_to_attach, + instance_name): + try: + instance_id = nova.create_instance(token, admin_tenant_id, + image_id, flavor, interfaces_to_attach, instance_name) + return instance_id + except Exception as e: + LOG.error(_LE('Failed to create %(device_type)s instance.' + 'Error: %(error)s'), + {'device_type': ( + device_data['service_details']['device_type']), + 'error': e}) + + def get_neutron_port_details(self, network_handler, token, port_id): + try: + (mgmt_ip_address, + mgmt_mac, mgmt_cidr, gateway_ip, + mgmt_port, mgmt_subnet) = network_handler.get_neutron_port_details(token, port_id) + + result = {'neutron_port': mgmt_port['port'], + 'neutron_subnet': mgmt_subnet['subnet'], + 'ip_address': mgmt_ip_address, + 'mac': mgmt_mac, + 'cidr': mgmt_cidr, + 'gateway_ip': gateway_ip} + return result + except Exception as e: + import sys + import traceback + exc_type, exc_value, exc_traceback = sys.exc_info() + print traceback.format_exception(exc_type, exc_value, + exc_traceback) + LOG.error(traceback.format_exception(exc_type, exc_value, + exc_traceback)) + LOG.error(_LE('Failed to get management port details. ' + 'Error: %(error)s'), {'error': e}) + + @_set_network_handler def create_network_function_device(self, device_data, network_handler=None): @@ -466,46 +545,41 @@ def create_network_function_device(self, device_data, 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) - if image_name: - self._update_vendor_data(device_data, - device_data.get('token')) - try: - interfaces = self._get_interfaces_for_device_create( - device_data, - network_handler=network_handler - ) - except Exception as e: - LOG.exception(_LE('Failed to get interfaces for device creation.' - 'Error: %(error)s'), {'error', e}) + + executor = core_task.TaskExecutor(jobs=3) + + image_id_result = {} + + executor.add_job('UPDATE_VENDOR_DATA', + self._update_vendor_data_v1, + token, admin_tenant_id, image_name, device_data) + 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) + + + completed = 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)) - token = self._get_token(device_data.get('token')) - if not token: - return None - - if device_data['service_details'].get('image_name'): - image_name = device_data['service_details']['image_name'] - else: - LOG.info(_LI("No image name provided in service profile's " - "service flavor field, image will be selected " - "based on service vendor's name : %s") - % (device_data['service_details']['service_vendor'])) - image_name = device_data['service_details']['service_vendor'] - image_name = '%s' % image_name.lower() - try: - image_id = self.compute_handler_nova.get_image_id( - token, - self._get_admin_tenant_id(token=token), - image_name) - except Exception as e: + 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.' - ' image name: %(image_name)s. Error: %(error)s'), - {'image_name': image_name, 'error': e}) + 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', @@ -524,8 +598,7 @@ def create_network_function_device(self, device_data, advance_sharing_interfaces = [] try: for interface in interfaces: - port_id = network_handler.get_port_id(token, interface['id']) - interfaces_to_attach.append({'port': port_id}) + interfaces_to_attach.append({'port': interface['port_id']}) if not self.supports_hotplug: if self.setup_mode.get(nfp_constants.NEUTRON_MODE): @@ -575,18 +648,29 @@ def create_network_function_device(self, device_data, return None instance_name = device_data['name'] - try: - instance_id = self.compute_handler_nova.create_instance( - token, self._get_admin_tenant_id(token=token), - image_id, flavor, - interfaces_to_attach, instance_name) - except Exception as e: + 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) + + completed = 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.' - 'Error: %(error)s'), - {'device_type': ( - device_data['service_details']['device_type']), - 'error': e}) + 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', @@ -596,23 +680,15 @@ def create_network_function_device(self, device_data, self._increment_stats_counter('instances') mgmt_ip_address = None - try: - for interface in interfaces: - if interface['port_classification'] == ( - nfp_constants.MANAGEMENT): - (mgmt_ip_address, - dummy, dummy, - dummy) = network_handler.get_port_details( - token, interface['id']) - except Exception as e: + 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. ' - 'Error: %(error)s'), {'error': e}) + LOG.error(_LE('Failed to get management port details. ')) try: self.compute_handler_nova.delete_instance( token, - self._get_admin_tenant_id( - token=token), + admin_tenant_id, instance_id) except Exception as e: self._increment_stats_counter('instance_delete_failures') @@ -628,10 +704,12 @@ def create_network_function_device(self, device_data, by=len(interfaces)) return None + mgmt_ip_address = mgmt_neutron_port_info['ip_address'] return {'id': instance_id, 'name': instance_name, '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, @@ -751,14 +829,10 @@ def get_network_function_device_status(self, device_data, raise exceptions.ComputePolicyNotSupported( compute_policy=device_data['service_details']['device_type']) - token = self._get_token(device_data.get('token')) - if not token: - return None - try: device = self.compute_handler_nova.get_instance( - token, - self._get_admin_tenant_id(token=token), + device_data['token'], + device_data['tenant_id'], device_data['id']) except Exception: if ignore_failure: @@ -814,13 +888,8 @@ def plug_network_function_device_interfaces(self, device_data, raise exceptions.ComputePolicyNotSupported( compute_policy=device_data['service_details']['device_type']) - 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) + token = device_data['token'] + tenant_id = device_data['tenant_id'] update_ifaces = [] try: @@ -855,40 +924,37 @@ def plug_network_function_device_interfaces(self, device_data, elif self.setup_mode.get(nfp_constants.NEUTRON_MODE): pass else: + executor = core_task.TaskExecutor(jobs=10) + 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()] - ): - network_handler.set_promiscuos_mode(token, - port['id']) - port_id = network_handler.get_port_id(token, - port['id']) - self.compute_handler_nova.attach_interface( - token, - self._get_admin_tenant_id(token=token), - device_data['id'], - port_id) + service_type = device_data['service_details']['service_type'].lower() + if service_type == nfp_constants.FIREWALL.lower(): + executor.add_job('SET_PROMISCUOS_MODE', + network_handler.set_promiscuos_mode_v1, + token, port['id']) + executor.add_job('ATTACH_INTERFACE', + self.compute_handler_nova.attach_interface, + token, tenant_id, device_data['id'], + port['id']) break + # Configurator expects interface to attach in order + # executor.fire() + 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()] - ): - network_handler.set_promiscuos_mode(token, - port['id']) - port_id = network_handler.get_port_id(token, - port['id']) - self.compute_handler_nova.attach_interface( - token, - self._get_admin_tenant_id(token=token), - device_data['id'], - port_id) + service_type = device_data['service_details']['service_type'].lower() + if service_type == nfp_constants.FIREWALL.lower(): + executor.add_job('SET_PROMISCUOS_MODE', + network_handler.set_promiscuos_mode_v1, + 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.' @@ -1104,6 +1170,7 @@ def get_network_function_device_healthcheck_info(self, device_data): ] } + @_set_network_handler def get_network_function_device_config_info(self, device_data, network_handler=None): @@ -1221,3 +1288,77 @@ def get_network_function_device_config_info(self, device_data, } ] } + + + + @_set_network_handler + def get_create_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 + """ + + mgmt_ip = device_data.get('mgmt_ip', None) + provider_ip = device_data.get('provider_ip', None) + provider_mac = device_data.get('provider_mac', None) + provider_cidr = device_data.get('provider_cidr', None) + provider_gateway_ip = device_data.get('provider_gateway_ip', None) + consumer_ip = device_data.get('consumer_ip', None) + consumer_mac = device_data.get('consumer_mac', None) + consumer_cidr = device_data.get('consumer_cidr', None) + consumer_gateway_ip = device_data.get('consumer_gateway_ip', None) + + return { + 'config': [ + { + 'resource': nfp_constants.INTERFACE_RESOURCE, + 'resource_data': { + 'mgmt_ip': mgmt_ip, + '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': mgmt_ip, + 'source_cidrs': ([provider_cidr, consumer_cidr] + if consumer_cidr + else [provider_cidr]), + 'destination_cidr': consumer_cidr, + 'gateway_ip': consumer_gateway_ip, + 'provider_interface_index': 2 + } + } + ] + } diff --git a/gbpservice/nfp/orchestrator/modules/device_orchestrator.py b/gbpservice/nfp/orchestrator/modules/device_orchestrator.py index 0f5d6f7634..0a25622a11 100644 --- a/gbpservice/nfp/orchestrator/modules/device_orchestrator.py +++ b/gbpservice/nfp/orchestrator/modules/device_orchestrator.py @@ -10,8 +10,8 @@ # License for the specific language governing permissions and limitations # under the License. -from neutron._i18n import _LE -from neutron._i18n import _LI +from neutron.i18n import _LE +from neutron.i18n import _LI import oslo_messaging as messaging from gbpservice.nfp.common import constants as nfp_constants @@ -32,6 +32,8 @@ import traceback from gbpservice.nfp.core import log as nfp_logging + + LOG = nfp_logging.getLogger(__name__) STOP_POLLING = {'poll': False} @@ -50,11 +52,12 @@ def rpc_init(controller, config): def events_init(controller, config, device_orchestrator): events = ['CREATE_NETWORK_FUNCTION_DEVICE', 'DEVICE_SPAWNING', - 'DEVICE_HEALTHY', 'CONFIGURE_DEVICE', + 'DEVICE_HEALTHY', 'HEALTH_MONITOR_COMPLETE', + 'CONFIGURE_DEVICE', 'CREATE_DEVICE_CONFIGURATION', 'DEVICE_CONFIGURED', "DELETE_CONFIGURATION", 'DELETE_NETWORK_FUNCTION_DEVICE', 'DELETE_CONFIGURATION_COMPLETED', 'DEVICE_BEING_DELETED', - 'DEVICE_NOT_REACHABLE', 'DEVICE_CONFIGURATION_FAILED'] + 'DEVICE_NOT_REACHABLE', 'DEVICE_CONFIGURATION_FAILED', 'DEVICE_UP'] events_to_register = [] for event in events: events_to_register.append( @@ -76,7 +79,7 @@ def __init__(self, conf, controller): self.conf = conf self._controller = controller self.rpc_event_mapping = { - 'healthmonitor': ['DEVICE_HEALTHY', + 'healthmonitor': ['HEALTH_MONITOR_COMPLETE', 'DEVICE_NOT_REACHABLE', 'DEVICE_NOT_REACHABLE'], 'interfaces': ['DEVICE_CONFIGURED', @@ -221,8 +224,11 @@ def event_method_mapping(self, event_id): "CREATE_NETWORK_FUNCTION_DEVICE": ( self.create_network_function_device), "DEVICE_UP": self.perform_health_check, + "PLUG_INTERFACES": self.plug_interfaces_v1, "DEVICE_HEALTHY": self.plug_interfaces, - "CONFIGURE_DEVICE": self.create_device_configuration, + "HEALTH_MONITOR_COMPLETE": self.device_healthy, + "CONFIGURE_DEVICE": self.configure_device, + "CREATE_DEVICE_CONFIGURATION": self.create_device_configuration, "DEVICE_CONFIGURED": self.device_configuration_complete, "DELETE_NETWORK_FUNCTION_DEVICE": ( @@ -409,10 +415,6 @@ def _create_network_function_device_db(self, device_info, state): device_info['interfaces_in_use'] = 0 device = self.nsf_db.create_network_function_device(self.db_session, device_info) - mgmt_port_id = device.pop('mgmt_port_id') - mgmt_port_id = self._get_port(mgmt_port_id) - device['mgmt_port_id'] = mgmt_port_id - if advance_sharing_interfaces: self._create_advance_sharing_interfaces(device, advance_sharing_interfaces) @@ -470,6 +472,7 @@ def _get_device_to_reuse(self, device_data, dev_sharing_info): return device def _get_device_data(self, nfd_request): + device_data = {} network_function = nfd_request.get('network_function') network_function_instance = nfd_request['network_function_instance'] @@ -504,6 +507,7 @@ def _get_device_data(self, nfd_request): device_data['service_details']['network_mode'] = ( nfp_constants.NEUTRON_MODE) device_data['service_vendor'] = service_details['service_vendor'] + return device_data def _get_nsf_db_resource(self, resource_name, resource_id): @@ -514,26 +518,78 @@ def _update_device_data(self, device, device_data): device.update(device_data) return device - # Create path + def _prepare_device_data_from_nfp_context(self, nfp_context): + device_data = {} + + network_function = nfp_context['network_function'] + network_function_instance = nfp_context['network_function_instance'] + service_details = nfp_context['service_details'] + + device_data['token'] = nfp_context['resource_owner_context']['admin_token'] + device_data['admin_tenant_id'] = nfp_context['resource_owner_context']['admin_tenant_id'] + device_data['name'] = network_function_instance['name'] + device_data['share_existing_device'] = nfp_context['share_existing_device'] + + management_network_info = { + 'id': nfp_context['management_ptg_id'], + 'port_model': nfp_constants.GBP_NETWORK + } + + consumer = nfp_context['consumer'] + provider = nfp_context['provider'] + ports = [] + + if consumer['port']: + ports.append({'id': consumer['port']['id'], + 'port_classification': consumer['port_classification'], + 'port_model': consumer['port_model']}) + + if provider['port']: + ports.append({'id': provider['port']['id'], + 'port_classification': provider['port_classification'], + 'port_model': provider['port_model']}) + + device_data['management_network_info'] = management_network_info + + device_data['network_function_id'] = network_function['id'] + device_data['service_chain_id'] = network_function['service_chain_id'] + device_data['network_function_instance_id'] = network_function_instance['id'] + device_data['tenant_id'] = network_function_instance['tenant_id'] + device_data['ports'] = ports + device_data['service_details'] = service_details + device_data['service_details']['network_mode'] = nfp_constants.GBP_MODE + device_data['service_vendor'] = service_details['service_vendor'] + + return device_data + + # Create path def create_network_function_device(self, event): """ Returns device instance for a new service This method either returns existing device which could be reused for a new service or it creates new device instance """ - nfd_request = event.data + device = None + nfd_request = event.data + nfp_context = event.data + + service_details = nfp_context['service_details'] + LOG.info(_LI("Device Orchestrator received create network service " "device request with data %(data)s"), {'data': nfd_request}) - device_data = self._get_device_data(nfd_request) orchestration_driver = self._get_orchestration_driver( - device_data['service_details']['service_vendor']) + service_details['service_vendor']) + + device_data = self._prepare_device_data_from_nfp_context(nfp_context) + dev_sharing_info = ( orchestration_driver.get_network_function_device_sharing_info( device_data)) + if dev_sharing_info: device = self._get_device_to_reuse(device_data, dev_sharing_info) if device: @@ -552,6 +608,8 @@ def create_network_function_device(self, event): else: LOG.info(_LI("No Device exists for sharing, Creating new device," "device request: %(device)s"), {'device': nfd_request}) + LOG.info("Prepared device_data %s" %(device_data)) + driver_device_info = ( orchestration_driver.create_network_function_device( device_data)) @@ -561,42 +619,83 @@ def create_network_function_device(self, event): event_data=nfd_request, is_internal_event=True) return None + + management = nfp_context['management'] + management['port'] = driver_device_info['mgmt_neutron_port_info']['neutron_port'] + management['port']['ip_address'] = management['port']['fixed_ips'][0]['ip_address'] + management['subnet'] = driver_device_info['mgmt_neutron_port_info']['neutron_subnet'] # Update newly created device with required params device = self._update_device_data(driver_device_info, device_data) device['network_function_device_id'] = device['id'] # Create DB entry with status as DEVICE_SPAWNING - self._create_network_function_device_db(device, + network_function_device = self._create_network_function_device_db(device, 'DEVICE_SPAWNING') + + #[mak: TODO] Wrong by nfp_db method needs in this format + network_function_device['mgmt_port_id'] = device['mgmt_port_id'] + nfp_context['network_function_device'] = network_function_device + # Create an event to NSO, to give device_id device_created_data = { 'network_function_instance_id': ( nfd_request['network_function_instance']['id']), - 'network_function_device_id': device['id'] + 'network_function_device_id': device['id'], } - self._create_event(event_id='DEVICE_CREATED', - event_data=device_created_data) + self._create_event(event_id='DEVICE_SPAWNING', - event_data=device, + event_data=nfp_context, is_poll_event=True, original_event=event) - @poll_event_desc(event='DEVICE_SPAWNING', spacing=20) + self._create_event(event_id='DEVICE_CREATED', + event_data=device_created_data) + + @poll_event_desc(event='DEVICE_SPAWNING', spacing=2) def check_device_is_up(self, event): - device = event.data + 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']['tenant_id'] + + device = { + 'token': token, + 'tenant_id': tenant_id, + 'id': network_function_device['id'], + 'service_details': service_details} orchestration_driver = self._get_orchestration_driver( - device['service_details']['service_vendor']) + service_details['service_vendor']) + is_device_up = ( orchestration_driver.get_network_function_device_status(device)) + if is_device_up == nfp_constants.ACTIVE: + # [(mak)TODO] - Update interfaces count here before + # sending health monitor rpc in DEVICE_UP event. + # [HACK] to handle a very corner case where + # PLUG_INTERFACES completes later than HEALTHMONITOR. + # till proper fix is identified. + provider = nfp_context['provider']['ptg'] + consumer = nfp_context['consumer']['ptg'] + network_function_device = nfp_context['network_function_device'] + + if provider: + network_function_device['interfaces_in_use'] += 1 + if consumer: + network_function_device['interfaces_in_use'] += 1 + # create event DEVICE_UP self._create_event(event_id='DEVICE_UP', - event_data=device, + event_data=nfp_context) + + self._create_event(event_id='PLUG_INTERFACES', + event_data=nfp_context, is_internal_event=True) - self._update_network_function_device_db(device, - 'DEVICE_UP') + return STOP_POLLING elif is_device_up == nfp_constants.ERROR: # create event DEVICE_NOT_UP @@ -611,10 +710,27 @@ def check_device_is_up(self, event): return CONTINUE_POLLING def perform_health_check(self, event): + nfp_context = event.data + + service_details = nfp_context['service_details'] + network_function_device = nfp_context['network_function_device'] + network_function = nfp_context['network_function'] + network_function_instance = nfp_context['network_function_instance'] + mgmt_ip_address = nfp_context['management']['port']['ip_address'] + # The driver tells which protocol / port to monitor ?? - device = event.data orchestration_driver = self._get_orchestration_driver( - device['service_details']['service_vendor']) + service_details['service_vendor']) + + device ={ + 'id': network_function_device['id'], + 'mgmt_ip_address': mgmt_ip_address, + 'service_details': service_details, + 'network_function_id': network_function['id'], + 'network_function_instance_id': network_function_instance['id'], + 'nfp_context': nfp_context + } + hm_req = ( orchestration_driver.get_network_function_device_healthcheck_info( device)) @@ -623,11 +739,14 @@ def perform_health_check(self, event): event_data=device, is_internal_event=True) return None + self.configurator_rpc.create_network_function_device_config(device, hm_req) LOG.debug("Health Check RPC sent to configurator for device: " "%s with health check parameters: %s" % ( device['id'], hm_req)) + + device['status'] = 'HEALTH_CHECK_PENDING' self._update_network_function_device_db(device, 'HEALTH_CHECK_PENDING') @@ -678,6 +797,18 @@ def _prepare_device_data(self, device_info): device['advance_sharing_interfaces'] = ( self._get_advance_sharing_interfaces(device['id'])) return device + + def device_healthy(self, event): + nfp_context = event.data['nfp_context'] + + device = nfp_context['network_function_device'] + + self._update_network_function_device_db(device, 'HEALTH_CHECK_COMPLETED') + + self._create_event(event_id='DEVICE_ACTIVE', event_data=nfp_context) + + self._create_event(event_id='CREATE_DEVICE_CONFIGURATION', + event_data=nfp_context, is_internal_event=True) def plug_interfaces(self, event, is_event_call=True): if is_event_call: @@ -691,15 +822,10 @@ def plug_interfaces(self, event, is_event_call=True): 'HEALTH_CHECK_COMPLETED') orchestration_driver = self._get_orchestration_driver( device['service_details']['service_vendor']) - - _ifaces_plugged_in, advance_sharing_ifaces = ( + _ifaces_plugged_in = ( 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, @@ -709,7 +835,57 @@ def plug_interfaces(self, event, is_event_call=True): event_data=device, is_internal_event=True) - def create_device_configuration(self, event): + + def plug_interfaces_v1(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 = [] + if consumer['port']: + ports.append({'id': consumer['port']['id'], + 'port_classification': consumer['port_classification'], + 'port_model': consumer['port_model']}) + if provider['port']: + ports.append({'id': provider['port']['id'], + 'port_classification': provider['port_classification'], + 'port_model': provider['port_model']}) + 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']} + + _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) + #[mak: TODO] - Check how incremented ref count can be updated in DB + + def configure_device(self, event): device = event.data orchestration_driver = self._get_orchestration_driver( device['service_details']['service_vendor']) @@ -725,9 +901,59 @@ def create_device_configuration(self, event): self.configurator_rpc.create_network_function_device_config( device, config_params) + def create_device_configuration(self, event): + nfp_context = event.data + + service_details = nfp_context['service_details'] + token = nfp_context['resource_owner_context']['admin_token'] + tenant_id = nfp_context['resource_owner_context']['tenant_id'] + consumer = nfp_context['consumer'] + provider = nfp_context['provider'] + management = nfp_context['management'] + network_function = nfp_context['network_function'] + network_function_instance = nfp_context['network_function_instance'] + network_function_device = nfp_context['network_function_device'] + + orchestration_driver = self._get_orchestration_driver( + service_details['service_vendor']) + device = { + 'mgmt_ip': management['port']['ip_address'], + 'provider_ip': provider['port']['ip_address'], + 'provider_cidr': provider['subnet']['cidr'], + 'provider_mac': provider['port']['mac_address'], + 'provider_gateway_ip': provider['subnet']['gateway_ip']} + + if consumer['port'] and consumer['subnet']: + device.update({'consumer_ip': consumer['port']['ip_address'], + 'consumer_cidr': consumer['subnet']['cidr'], + 'consumer_mac': consumer['port']['mac_address'], + 'consumer_gateway_ip': consumer['subnet']['gateway_ip']}) + + config_params = ( + orchestration_driver.get_create_network_function_device_config_info( + device)) + device.update({ + 'id': network_function_device['id'], + 'mgmt_ip_address': management['port']['ip_address'], + 'service_details': service_details, + 'network_function_id': network_function['id'], + 'network_function_instance_id': network_function_instance['id'], + 'nfp_context': nfp_context}) + + if not config_params: + self._create_event(event_id='DRIVER_ERROR', + event_data=device, + is_internal_event=True) + return None + # Sends RPC to configurator to create generic config + self.configurator_rpc.create_network_function_device_config( + device, config_params) + def device_configuration_complete(self, event): - device_info = event.data - device = self._prepare_device_data(device_info) + nfp_context = event.data.get('nfp_context', None) + + device = nfp_context['network_function_device'] + # Change status to active in DB and generate an event DEVICE_ACTIVE # to inform NSO self._increment_device_ref_count(device) @@ -737,17 +963,6 @@ def device_configuration_complete(self, event): "reference count for %(device)s"), {'device_id': device['id'], 'device': device}) - device_created_data = { - 'network_function_id': ( - device['network_function_id']), - 'network_function_instance_id': ( - device['network_function_instance_id']), - 'network_function_device_id': device['id'] - } - # DEVICE_ACTIVE event for NSO. - self._create_event(event_id='DEVICE_ACTIVE', - event_data=device_created_data) - # Delete path def delete_network_function_device(self, event): delete_nfd_request = event.data @@ -928,7 +1143,9 @@ def _get_request_info(self, device, operation): 'nfd_id': device['id'], 'requester': nfp_constants.DEVICE_ORCHESTRATOR, 'operation': operation, - 'logging_context': nfp_logging.get_logging_context() + 'logging_context': nfp_logging.get_logging_context(), + # So that notification callbacks can work on cached data + 'nfp_context': device.get('nfp_context', None) } nfd_ip = device['mgmt_ip_address'] request_info.update({'device_ip': nfd_ip}) diff --git a/gbpservice/nfp/orchestrator/modules/service_orchestrator.py b/gbpservice/nfp/orchestrator/modules/service_orchestrator.py index 3dd9b325a6..1f17ede2a2 100644 --- a/gbpservice/nfp/orchestrator/modules/service_orchestrator.py +++ b/gbpservice/nfp/orchestrator/modules/service_orchestrator.py @@ -10,8 +10,8 @@ # License for the specific language governing permissions and limitations # under the License. -from neutron._i18n import _LE -from neutron._i18n import _LI +from neutron.i18n import _LE +from neutron.i18n import _LI from neutron.common import rpc as n_rpc from neutron import context as n_context from oslo_log import helpers as log_helpers @@ -32,6 +32,8 @@ import traceback from gbpservice.nfp.core import log as nfp_logging +from gbpservice.nfp.core import context as nfp_core_context + LOG = nfp_logging.getLogger(__name__) STOP_POLLING = {'poll': False} @@ -69,7 +71,7 @@ def events_init(controller, config, service_orchestrator): 'DELETE_USER_CONFIG_IN_PROGRESS', 'CONFIG_APPLIED', 'USER_CONFIG_APPLIED', 'USER_CONFIG_DELETED', 'USER_CONFIG_DELETE_FAILED', 'USER_CONFIG_UPDATE_FAILED', - 'USER_CONFIG_FAILED'] + 'USER_CONFIG_FAILED', 'CHECK_USER_CONFIG_COMPLETE'] events_to_register = [] for event in events: events_to_register.append( @@ -402,13 +404,15 @@ def event_method_mapping(self, event_id): "CONSUMER_ADD": self.consumer_ptg_add_user_config, "CONSUMER_REMOVE": self.consumer_ptg_remove_user_config, "APPLY_USER_CONFIG_IN_PROGRESS": ( + self.apply_user_config_in_progress), + "CHECK_USER_CONFIG_COMPLETE": ( self.check_for_user_config_complete), "UPDATE_USER_CONFIG_PREPARING_TO_START": ( self.check_for_user_config_deleted), "UPDATE_USER_CONFIG_IN_PROGRESS": ( self.handle_continue_update_user_config), "UPDATE_USER_CONFIG_STILL_IN_PROGRESS": ( - self.check_for_user_config_complete), + self.apply_user_config_in_progress), "DELETE_USER_CONFIG_IN_PROGRESS": ( self.check_for_user_config_deleted), "CONFIG_APPLIED": self.handle_config_applied, @@ -519,9 +523,11 @@ def update_network_function_user_config(self, network_function_id, if tag_str != nfp_constants.CONFIG_INIT_TAG: network_function_details = self.get_network_function_details( network_function_id) - service_type = self._get_service_type( - network_function_details['network_function'][ - 'service_profile_id']) + service_type = network_function_details.pop('service_type') + if not service_type: + service_type = self._get_service_type( + network_function_details['network_function'][ + 'service_profile_id']) network_function_data = { 'network_function_details': network_function_details, 'service_type': service_type @@ -607,16 +613,24 @@ def _report_logging_info(self, nf, nfi, service_type, def create_network_function(self, context, network_function_info): self._validate_create_service_input(context, network_function_info) + + admin_token = self.keystoneclient.get_admin_token() + admin_tenant_id = self.keystoneclient.get_admin_tenant_id(admin_token) + + network_function_info['resource_owner_context']['admin_token'] = admin_token + network_function_info['resource_owner_context']['admin_tenant_id'] = admin_tenant_id + + tenant_id = network_function_info['tenant_id'] + # GBP or Neutron mode = network_function_info['network_function_mode'] - service_profile_id = network_function_info['service_profile_id'] - service_id = network_function_info['service_id'] - admin_token = self.keystoneclient.get_admin_token() - service_profile = self.gbpclient.get_service_profile( - admin_token, service_profile_id) - service_chain_id = network_function_info.get('service_chain_id') + service_profile = network_function_info['service_profile'] + service_profile_id = service_profile['id'] + service_id = network_function_info['service_chain_node']['id'] + service_chain_id = network_function_info['service_chain_instance']['id'] service_details = transport.parse_service_flavor_string( service_profile['service_flavor']) + base_mode_support = (True if service_details['device_type'] == 'None' else False) service_vendor = service_details['service_vendor'] @@ -627,13 +641,14 @@ def create_network_function(self, context, network_function_info): network_function = { 'name': name, 'description': '', - 'tenant_id': network_function_info['tenant_id'], + 'tenant_id': tenant_id, 'service_id': service_id, # GBP Service Node or Neutron Service ID 'service_chain_id': service_chain_id, # GBP SC instance ID 'service_profile_id': service_profile_id, 'service_config': service_config_str, 'status': nfp_constants.PENDING_CREATE } + network_function = self.db_handler.create_network_function( self.db_session, network_function) @@ -658,26 +673,19 @@ def create_network_function(self, context, network_function_info): service_config_str) return network_function - if mode == nfp_constants.GBP_MODE: - management_network_info = { - 'id': network_function_info['management_ptg_id'], - 'port_model': nfp_constants.GBP_NETWORK - } - else: - management_network_info = {} - create_network_function_instance_request = { - 'network_function': network_function, - 'network_function_port_info': network_function_info['port_info'], - 'management_network_info': management_network_info, - 'service_type': service_profile['service_type'], - 'service_details': service_details, - 'share_existing_device': False # Extend service profile if needed - } + nfp_context = network_function_info + + service_details['service_type'] = service_profile['service_type'] + service_details['network_mode'] = nfp_context['network_function_mode'] + nfp_context['network_function'] = network_function + nfp_context['service_details'] = service_details + nfp_context['share_existing_device'] = False # Create and event to perform Network service instance self._create_event('CREATE_NETWORK_FUNCTION_INSTANCE', - event_data=create_network_function_instance_request, + event_data=nfp_context, is_internal_event=True) + nfp_logging.clear_logging_context() return network_function @@ -751,42 +759,49 @@ def delete_user_config(self, event): is_poll_event=True, original_event=event) def create_network_function_instance(self, event): - request_data = event.data - name = '%s_%s' % (request_data['network_function']['name'], - request_data['network_function']['id']) + 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']: + port_info.append({'id': ele['pt']['id'], + 'port_model': ele['port_model'], + 'port_classification': ele['port_classification'] + }) + + name = '%s_%s' % (network_function['name'], + network_function['id']) create_nfi_request = { 'name': name, - 'tenant_id': request_data['network_function']['tenant_id'], + 'tenant_id': network_function['tenant_id'], 'status': nfp_constants.PENDING_CREATE, - 'network_function_id': request_data['network_function']['id'], - 'service_type': request_data['service_type'], - 'service_vendor': ( - request_data['service_details']['service_vendor']), - 'share_existing_device': request_data['share_existing_device'], - 'port_info': request_data['network_function_port_info'], + '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) + + network_function_instance = \ + self.db_handler.create_network_function_instance( + self.db_session, create_nfi_request) # Sending LogMeta Details to visibility - self._report_logging_info(request_data['network_function'], nfi_db, - request_data['service_type'], - request_data['service_details'][ - 'service_vendor']) - - request_data['service_details'].update( - service_type=request_data['service_type']) - create_nfd_request = { - 'network_function': request_data['network_function'], - 'network_function_instance': nfi_db, - 'management_network_info': request_data['management_network_info'], - 'service_vendor': ( - request_data['service_details']['service_vendor']), - 'service_details': request_data['service_details'], - 'share_existing_device': request_data['share_existing_device'], - } + self._report_logging_info(network_function, + network_function_instance, + service_details['service_type'], + service_details['service_vendor']) + + nfp_context['network_function_instance'] = network_function_instance + LOG.info(_LI("[Event:CreateService]")) self._create_event('CREATE_NETWORK_FUNCTION_DEVICE', - event_data=create_nfd_request) + event_data=nfp_context) def handle_device_created(self, event): request_data = event.data @@ -799,7 +814,15 @@ def handle_device_created(self, event): return def handle_device_active(self, event): - request_data = event.data + nfp_context = event.data + + network_function_instance = nfp_context['network_function_instance'] + network_function_device = nfp_context['network_function_device'] + network_function = nfp_context['network_function'] + + request_data = {'network_function_device_id': network_function_device['id'], + 'network_function_instance_id': network_function_instance['id']} + nfi = { 'status': nfp_constants.ACTIVE, 'network_function_device_id': request_data[ @@ -807,37 +830,48 @@ def handle_device_active(self, event): } nfi = self.db_handler.update_network_function_instance( self.db_session, request_data['network_function_instance_id'], nfi) - network_function = self.db_handler.get_network_function( - self.db_session, nfi['network_function_id']) + network_function_instance['status'] = nfp_constants.ACTIVE + network_function_instance['network_function_device_id'] = network_function_device['id'] + service_config = network_function['service_config'] + nfp_core_context.store_nfp_context(nfp_context) self.create_network_function_user_config(network_function['id'], service_config) def apply_user_config(self, event): request_data = event.data + nfp_context = event.data['nfp_context'] + + network_function = nfp_context['network_function'] network_function_details = self.get_network_function_details( - request_data['network_function_id']) - request_data['heat_stack_id'] = self.config_driver.apply_config( - network_function_details) # Heat driver to launch stack - network_function = network_function_details['network_function'] + network_function['id']) + request_data['heat_stack_id'], heat_client = self.config_driver.apply_heat_config( + nfp_context) # Heat driver to launch stack request_data['network_function_id'] = network_function['id'] + if not request_data['heat_stack_id']: self._create_event('USER_CONFIG_FAILED', event_data=request_data, is_internal_event=True) return - request_data['tenant_id'] = network_function['tenant_id'] - request_data['network_function_details'] = network_function_details + LOG.debug("handle_device_active heat_stack_id: %s" % (request_data['heat_stack_id'])) + + nfp_context['heat_stack_id'] = request_data['heat_stack_id'] + nfp_context['network_function'].update({ + 'heat_stack_id': request_data['heat_stack_id'], + 'description': network_function['description']}) + + self._create_event('CHECK_USER_CONFIG_COMPLETE', + event_data=nfp_context, + is_poll_event=True, + original_event=event) + self.db_handler.update_network_function( self.db_session, network_function['id'], {'heat_stack_id': request_data['heat_stack_id'], 'description': network_function['description']}) - self._create_event('APPLY_USER_CONFIG_IN_PROGRESS', - event_data=request_data, - is_poll_event=True, - original_event=event) def handle_update_user_config(self, event): request_data = event.data @@ -980,8 +1014,10 @@ def delete_network_function_instance(self, event): # FIXME: Add all possible validations here def _validate_create_service_input(self, context, create_service_request): - required_attributes = ["tenant_id", "service_id", "service_chain_id", - "service_profile_id", "network_function_mode"] + required_attributes = ["resource_owner_context","service_chain_instance", + "service_chain_node", "service_profile", + "service_config", "provider", "consumer", + "network_function_mode"] if (set(required_attributes) & set(create_service_request.keys()) != set(required_attributes)): missing_keys = (set(required_attributes) - @@ -990,8 +1026,7 @@ def _validate_create_service_input(self, context, create_service_request): required_data=", ".join(missing_keys), request="Create Network Function") if create_service_request['network_function_mode'].lower() == "gbp": - gbp_required_attributes = ["port_info", "service_chain_id", - "management_ptg_id"] + gbp_required_attributes = ["management_ptg_id"] if (set(gbp_required_attributes) & set(create_service_request.keys()) != set(gbp_required_attributes)): @@ -1001,7 +1036,7 @@ def _validate_create_service_input(self, context, create_service_request): required_data=", ".join(missing_keys), request="Create Network Function") - def check_for_user_config_complete(self, event): + def apply_user_config_in_progress(self, event): request_data = event.data config_status = self.config_driver.is_config_complete( request_data['heat_stack_id'], request_data['tenant_id'], @@ -1036,6 +1071,42 @@ def check_for_user_config_complete(self, event): elif config_status == nfp_constants.IN_PROGRESS: return CONTINUE_POLLING + def check_for_user_config_complete(self, event): + nfp_context = event.data + + network_function = nfp_context['network_function'] + config_status = self.config_driver.check_config_complete(nfp_context) + + if config_status == nfp_constants.ERROR: + LOG.info(_LI("NSO: applying user config failed for " + "network function %(network_function_id)s data " + "%(data)s"), {'data': nfp_context, + 'network_function_id': + network_function['id']}) + updated_network_function = {'status': nfp_constants.ERROR} + self.db_handler.update_network_function( + self.db_session, + network_function['id'], + updated_network_function) + self._controller.event_done(event) + return STOP_POLLING + # Trigger RPC to notify the Create_Service caller with status + elif config_status == nfp_constants.COMPLETED: + updated_network_function = {'status': nfp_constants.ACTIVE} + LOG.info(_LI("NSO: applying user config is successfull moving " + "network function %(network_function_id)s to ACTIVE"), + {'network_function_id': + network_function['id']}) + self.db_handler.update_network_function( + self.db_session, + network_function['id'], + updated_network_function) + self._controller.event_done(event) + return STOP_POLLING + # Trigger RPC to notify the Create_Service caller with status + elif config_status == nfp_constants.IN_PROGRESS: + return CONTINUE_POLLING + def check_for_user_config_deleted(self, event): request_data = event.data event_data = { @@ -1477,24 +1548,44 @@ def get_port_info(self, port_id): return None def get_network_function_details(self, network_function_id): - network_function = self.db_handler.get_network_function( - self.db_session, network_function_id) + network_function = None + network_function_instance = None + network_function_device = None + service_type = None + + nfp_context = nfp_core_context.get_nfp_context() + if nfp_context: + network_function = nfp_context['network_function'] + network_function_instance = nfp_context['network_function_instance'] + network_function_device = nfp_context['network_function_device'] + service_type = nfp_context['service_details']['service_type'] + + if not network_function: + network_function = self.db_handler.get_network_function( + self.db_session, network_function_id) + network_function_details = { - 'network_function': network_function + 'network_function': network_function, + 'service_type': service_type } - network_function_instances = network_function[ - 'network_function_instances'] - if not network_function_instances: - return network_function_details - nfi = self.db_handler.get_network_function_instance( - self.db_session, network_function_instances[0]) - network_function_details['network_function_instance'] = nfi - if nfi['network_function_device_id']: - network_function_device = ( - self.db_handler.get_network_function_device( - self.db_session, nfi['network_function_device_id'])) - network_function_details['network_function_device'] = ( - network_function_device) + + if not network_function_instance: + network_function_instances = network_function[ + 'network_function_instances'] + if not network_function_instances: + return network_function_details + network_function_instance = self.db_handler.get_network_function_instance( + self.db_session, network_function_instances[0]) + + network_function_details['network_function_instance'] = network_function_instance + + if not network_function_device: + if network_function_instance['network_function_device_id']: + network_function_device = ( + self.db_handler.get_network_function_device( + self.db_session, network_function_instance['network_function_device_id'])) + network_function_details['network_function_device'] = ( + network_function_device) return network_function_details @@ -1525,7 +1616,8 @@ def _get_request_info(self, user_config_data, operation): 'nfd_id': None, 'requester': nfp_constants.SERVICE_ORCHESTRATOR, 'operation': operation, - 'logging_context': nfp_logging.get_logging_context() + 'logging_context': nfp_logging.get_logging_context(), + 'nfp_context': nfp_core_context.get_nfp_context() } if operation in ['consumer_add', 'consumer_remove']: request_info.update({'consumer_ptg': user_config_data[ diff --git a/gbpservice/nfp/orchestrator/openstack/openstack_driver.py b/gbpservice/nfp/orchestrator/openstack/openstack_driver.py index 869983c1b4..24a3091297 100644 --- a/gbpservice/nfp/orchestrator/openstack/openstack_driver.py +++ b/gbpservice/nfp/orchestrator/openstack/openstack_driver.py @@ -41,6 +41,7 @@ def __init__(self, config, username=None, self.tenant_name = (tenant_name or config.keystone_authtoken.admin_tenant_name) self.token = None + self.admin_tenant_id = None class KeystoneClient(OpenstackApi): @@ -103,6 +104,13 @@ def get_scoped_keystone_token(self, user, password, tenant_name, else: return scoped_token + def get_admin_tenant_id(self, token): + if not self.admin_tenant_id: + _,_,name,_ = self.get_keystone_creds() + self.admin_tenant_id = self.get_tenant_id(token, name) + + return self.admin_tenant_id + def get_tenant_id(self, token, tenant_name): """ Get the tenant UUID associated to tenant name diff --git a/gbpservice/nfp/proxy_agent/proxy/proxy.py b/gbpservice/nfp/proxy_agent/proxy/proxy.py index 2c5bd73d8a..793b287e77 100644 --- a/gbpservice/nfp/proxy_agent/proxy/proxy.py +++ b/gbpservice/nfp/proxy_agent/proxy/proxy.py @@ -167,7 +167,9 @@ def idle_reset(self): def _wait(self, timeout): if self.type == 'unix': eventlet.sleep(timeout) - self._socket.settimeout(timeout) + self._socket.setblocking(0) + else: + self._socket.settimeout(timeout) def recv(self): self._wait(self._idle_wait) @@ -184,7 +186,9 @@ def recv(self): return None def send(self, data): - self._socket.send(data) + self._socket.setblocking(1) + self._socket.sendall(data) + self._socket.setblocking(0) def close(self): LOG.debug("Closing Socket - %d" % (self.identify()))