Skip to content
Snippets Groups Projects
E2EOrchestratorServiceServicerImpl.py 11.3 KiB
Newer Older
# Copyright 2022-2024 ETSI SDG TeraFlowSDN (TFS) (https://tfs.etsi.org/)
#
# 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 copy
Carlos Manso's avatar
Carlos Manso committed
import requests
Carlos Manso's avatar
Carlos Manso committed
from common.method_wrappers.Decorator import MetricsPool, safe_and_metered_rpc_method
from common.proto.e2eorchestrator_pb2 import E2EOrchestratorRequest, E2EOrchestratorReply
Carlos Manso's avatar
Carlos Manso committed
from common.proto.context_pb2 import (
Carlos Manso's avatar
Carlos Manso committed
    Empty, Connection, EndPointId, Link, LinkId, TopologyDetails, Topology, Context, Service, ServiceId,
    ServiceTypeEnum, ServiceStatusEnum)
from common.proto.e2eorchestrator_pb2_grpc import E2EOrchestratorServiceServicer
Carlos Manso's avatar
Carlos Manso committed
from common.Settings import get_setting
from context.client.ContextClient import ContextClient
Carlos Manso's avatar
Carlos Manso committed
from service.client.ServiceClient import ServiceClient
from context.service.database.uuids.EndPoint import endpoint_get_uuid
Carlos Manso's avatar
Carlos Manso committed
from context.service.database.uuids.Device import device_get_uuid
Carlos Manso's avatar
Carlos Manso committed
from common.proto.vnt_manager_pb2 import VNTSubscriptionRequest
Carlos Manso's avatar
Carlos Manso committed
from common.tools.grpc.Tools import grpc_message_to_json_string
Carlos Manso's avatar
Carlos Manso committed
import grpc
Carlos Manso's avatar
Carlos Manso committed
import json
Carlos Manso's avatar
Carlos Manso committed
import logging
import networkx as nx
from threading import Thread
from websockets.sync.client import connect
from websockets.sync.server import serve
Carlos Manso's avatar
Carlos Manso committed
from common.Constants import DEFAULT_CONTEXT_NAME, DEFAULT_TOPOLOGY_NAME
Carlos Manso's avatar
Carlos Manso committed


LOGGER = logging.getLogger(__name__)
Carlos Manso's avatar
Carlos Manso committed
logging.getLogger("websockets").propagate = True
Carlos Manso's avatar
Carlos Manso committed
logging.getLogger("requests.packages.urllib3").propagate = True

METRICS_POOL = MetricsPool("E2EOrchestrator", "RPC")

context_client: ContextClient = ContextClient()
Carlos Manso's avatar
Carlos Manso committed
service_client: ServiceClient = ServiceClient()
Carlos Manso's avatar
Carlos Manso committed
EXT_HOST = str(get_setting('WS_IP_HOST'))
EXT_PORT = str(get_setting('WS_IP_PORT'))

OWN_HOST = str(get_setting('WS_E2E_HOST'))
OWN_PORT = str(get_setting('WS_E2E_PORT'))
Carlos Manso's avatar
Carlos Manso committed


Carlos Manso's avatar
Carlos Manso committed
ALL_HOSTS = "0.0.0.0"
Carlos Manso's avatar
Carlos Manso committed

Carlos Manso's avatar
Carlos Manso committed
class SubscriptionServer(Thread):
    def __init__(self):
        Thread.__init__(self)
            
    def run(self):
        url = "ws://" + EXT_HOST + ":" + EXT_PORT
        request = VNTSubscriptionRequest()
        request.host = OWN_HOST
        request.port = OWN_PORT
        try: 
            LOGGER.debug("Trying to connect to {}".format(url))
            websocket = connect(url)
        except Exception as ex:
            LOGGER.error('Error connecting to {}'.format(url))
        else:
            with websocket:
                LOGGER.debug("Connected to {}".format(url))
                send = grpc_message_to_json_string(request)
                websocket.send(send)
                LOGGER.debug("Sent: {}".format(send))
                try:
                    message = websocket.recv()
                    LOGGER.debug("Received message from WebSocket: {}".format(message))
                except Exception as ex:
Carlos Manso's avatar
Carlos Manso committed
                    LOGGER.error('Exception receiving from WebSocket: {}'.format(ex))
Carlos Manso's avatar
Carlos Manso committed
            self._events_server()
Carlos Manso's avatar
Carlos Manso committed

Carlos Manso's avatar
Carlos Manso committed

Carlos Manso's avatar
Carlos Manso committed
    def _events_server(self):
        all_hosts = "0.0.0.0"
Carlos Manso's avatar
Carlos Manso committed

Carlos Manso's avatar
Carlos Manso committed
        try:
            server = serve(self._event_received, all_hosts, int(OWN_PORT))
        except Exception as ex:
            LOGGER.error('Error starting server on {}:{}'.format(all_hosts, OWN_PORT))
            LOGGER.error('Exception!: {}'.format(ex))
Carlos Manso's avatar
Carlos Manso committed
            with server:
                LOGGER.info("Running events server...: {}:{}".format(all_hosts, OWN_PORT))
                server.serve_forever()
Carlos Manso's avatar
Carlos Manso committed
    def _event_received(self, connection):
Carlos Manso's avatar
Carlos Manso committed
        LOGGER.debug('Event received')
Carlos Manso's avatar
Carlos Manso committed
        for message in connection:
            message_json = json.loads(message)
Carlos Manso's avatar
Carlos Manso committed
            # Link creation
Carlos Manso's avatar
Carlos Manso committed
            if 'link_id' in message_json:
Carlos Manso's avatar
Carlos Manso committed
                LOGGER.debug('Link creation')
Carlos Manso's avatar
Carlos Manso committed
                link = Link(**message_json)
Carlos Manso's avatar
Carlos Manso committed
                service = Service()
                service.service_id.service_uuid.uuid = link.link_id.link_uuid.uuid
                service.service_id.context_id.context_uuid.uuid = DEFAULT_CONTEXT_NAME
                service.service_type = ServiceTypeEnum.SERVICETYPE_OPTICAL_CONNECTIVITY
                service.service_status.service_status = ServiceStatusEnum.SERVICESTATUS_PLANNED
                service_client.CreateService(service)
Carlos Manso's avatar
Carlos Manso committed
                a_device_uuid = device_get_uuid(link.link_endpoint_ids[0].device_id)
                a_endpoint_uuid = endpoint_get_uuid(link.link_endpoint_ids[0])[2]
                z_device_uuid = device_get_uuid(link.link_endpoint_ids[1].device_id)
                z_endpoint_uuid = endpoint_get_uuid(link.link_endpoint_ids[1])[2]
Carlos Manso's avatar
Carlos Manso committed

Carlos Manso's avatar
Carlos Manso committed
                links = context_client.ListLinks(Empty()).links
Carlos Manso's avatar
Carlos Manso committed
                for _link in links:
                    for _endpoint_id in _link.link_endpoint_ids:
                        if _endpoint_id.device_id.device_uuid.uuid == a_device_uuid and \
                        _endpoint_id.endpoint_uuid.uuid == a_endpoint_uuid:
                            a_ep_id = _endpoint_id
                        elif _endpoint_id.device_id.device_uuid.uuid == z_device_uuid and \
                        _endpoint_id.endpoint_uuid.uuid == z_endpoint_uuid:
                            z_ep_id = _endpoint_id
Carlos Manso's avatar
Carlos Manso committed
                if (not 'a_ep_id' in locals()) or (not 'z_ep_id' in locals()):
Carlos Manso's avatar
Carlos Manso committed
                    error_msg = f'Could not get VNT link endpoints\
                                    \n\ta_endpoint_uuid= {a_endpoint_uuid}\
                                    \n\tz_endpoint_uuid= {z_device_uuid}'
Carlos Manso's avatar
Carlos Manso committed
                    LOGGER.error(error_msg)
                    connection.send(error_msg)
                    return
Carlos Manso's avatar
Carlos Manso committed

Carlos Manso's avatar
Carlos Manso committed
                service.service_endpoint_ids.append(copy.deepcopy(a_ep_id))
                service.service_endpoint_ids.append(copy.deepcopy(z_ep_id))
Carlos Manso's avatar
Carlos Manso committed

Carlos Manso's avatar
Carlos Manso committed
                service_client.UpdateService(service)
                re_svc = context_client.GetService(service.service_id)
Carlos Manso's avatar
Carlos Manso committed
                connection.send(grpc_message_to_json_string(link))
Carlos Manso's avatar
Carlos Manso committed
                context_client.SetLink(link)
Carlos Manso's avatar
Carlos Manso committed
            elif 'link_uuid' in message_json:
Carlos Manso's avatar
Carlos Manso committed
                LOGGER.debug('Link removal')
Carlos Manso's avatar
Carlos Manso committed
                link_id = LinkId(**message_json)

                service_id = ServiceId()
                service_id.service_uuid.uuid = link_id.link_uuid.uuid
                service_id.context_id.context_uuid.uuid = DEFAULT_CONTEXT_NAME
Carlos Manso's avatar
Carlos Manso committed
                service_client.DeleteService(service_id)
Carlos Manso's avatar
Carlos Manso committed
                connection.send(grpc_message_to_json_string(link_id))
Carlos Manso's avatar
Carlos Manso committed
                context_client.RemoveLink(link_id)
Carlos Manso's avatar
Carlos Manso committed
            else:
Carlos Manso's avatar
Carlos Manso committed
                LOGGER.debug('Topology received')
Carlos Manso's avatar
Carlos Manso committed
                topology_details = TopologyDetails(**message_json)
Carlos Manso's avatar
Carlos Manso committed

Carlos Manso's avatar
Carlos Manso committed
                context = Context()
                context.context_id.context_uuid.uuid = topology_details.topology_id.context_id.context_uuid.uuid
                context_client.SetContext(context)
Carlos Manso's avatar
Carlos Manso committed

Carlos Manso's avatar
Carlos Manso committed
                topology = Topology()
Carlos Manso's avatar
Carlos Manso committed
                topology.topology_id.context_id.CopyFrom(context.context_id)
Carlos Manso's avatar
Carlos Manso committed
                topology.topology_id.topology_uuid.uuid = topology_details.topology_id.topology_uuid.uuid
                context_client.SetTopology(topology)
Carlos Manso's avatar
Carlos Manso committed

Carlos Manso's avatar
Carlos Manso committed
                for device in topology_details.devices:
                    context_client.SetDevice(device)

                for link in topology_details.links:
                    context_client.SetLink(link)
Carlos Manso's avatar
Carlos Manso committed

class E2EOrchestratorServiceServicerImpl(E2EOrchestratorServiceServicer):
    def __init__(self):
        LOGGER.debug("Creating Servicer...")
Carlos Manso's avatar
Carlos Manso committed
        try:
Carlos Manso's avatar
Carlos Manso committed
            LOGGER.debug("Requesting subscription")
Carlos Manso's avatar
Carlos Manso committed
            sub_server = SubscriptionServer()
            sub_server.start()
            LOGGER.debug("Servicer Created")
Carlos Manso's avatar
Carlos Manso committed
            self.retrieve_external_topologies()
Carlos Manso's avatar
Carlos Manso committed
        except Exception as ex:
            LOGGER.info("Exception!: {}".format(ex))
Carlos Manso's avatar
Carlos Manso committed
    def retrieve_external_topologies(self):
        i = 1
        while True:
            try:
                ADD = str(get_setting(f'EXT_CONTROLLER{i}_ADD'))
                PORT = str(get_setting(f'EXT_CONTROLLER{i}_PORT'))
            except Exception as e:
                break
            try:
Carlos Manso's avatar
Carlos Manso committed
                LOGGER.info(f'Retrieving external controller #{i}')
Carlos Manso's avatar
Carlos Manso committed
                url = f'http://{ADD}:{PORT}/tfs-api/context/{DEFAULT_CONTEXT_NAME}/topology_details/{DEFAULT_TOPOLOGY_NAME}'
Carlos Manso's avatar
Carlos Manso committed
                LOGGER.info(f'url= {url}')
Carlos Manso's avatar
Carlos Manso committed
                topo = requests.get(url).json()
Carlos Manso's avatar
Carlos Manso committed
                LOGGER.info(f'Retrieved external controller #{i}')
Carlos Manso's avatar
Carlos Manso committed
            except Exception as e:
                LOGGER.info(f'Exception retrieven topology from external controler #{i}: {e}')
            topology_details = TopologyDetails(**topo)
            context = Context()
            context.context_id.context_uuid.uuid = topology_details.topology_id.context_id.context_uuid.uuid
            context_client.SetContext(context)

            topology = Topology()
            topology.topology_id.context_id.CopyFrom(context.context_id)
            topology.topology_id.topology_uuid.uuid = topology_details.topology_id.topology_uuid.uuid
            context_client.SetTopology(topology)

            for device in topology_details.devices:
                context_client.SetDevice(device)

            for link in topology_details.links:
                context_client.SetLink(link)

            i+=1
    @safe_and_metered_rpc_method(METRICS_POOL, LOGGER)
    def Compute(self, request: E2EOrchestratorRequest, context: grpc.ServicerContext) -> E2EOrchestratorReply:
        endpoints_ids = []
        for endpoint_id in request.service.service_endpoint_ids:
            endpoints_ids.append(endpoint_get_uuid(endpoint_id)[2])

        graph = nx.Graph()

        devices = context_client.ListDevices(Empty()).devices

        for device in devices:
            endpoints_uuids = [endpoint.endpoint_id.endpoint_uuid.uuid
                               for endpoint in device.device_endpoints]
            for ep in endpoints_uuids:
                graph.add_node(ep)

            for ep in endpoints_uuids:
                for ep_i in endpoints_uuids:
                    if ep == ep_i:
                        continue
                    graph.add_edge(ep, ep_i)

        links = context_client.ListLinks(Empty()).links
        for link in links:
            eps = []
            for endpoint_id in link.link_endpoint_ids:
                eps.append(endpoint_id.endpoint_uuid.uuid)
            graph.add_edge(eps[0], eps[1])


        shortest = nx.shortest_path(graph, endpoints_ids[0], endpoints_ids[1])

        path = E2EOrchestratorReply()
        path.services.append(copy.deepcopy(request.service))
        for i in range(0, int(len(shortest)/2)):
            conn = Connection()
            ep_a_uuid = str(shortest[i*2])
            ep_z_uuid = str(shortest[i*2+1])

            conn.connection_id.connection_uuid.uuid = str(ep_a_uuid) + '_->_' + str(ep_z_uuid)

            ep_a_id = EndPointId()
            ep_a_id.endpoint_uuid.uuid = ep_a_uuid
            conn.path_hops_endpoint_ids.append(ep_a_id)

            ep_z_id = EndPointId()
            ep_z_id.endpoint_uuid.uuid = ep_z_uuid
            conn.path_hops_endpoint_ids.append(ep_z_id)

            path.connections.append(conn)

Carlos Manso's avatar
Carlos Manso committed
        return path