Loading manifests/e2e_orchestratorservice.yaml +3 −0 Original line number Diff line number Diff line Loading @@ -25,6 +25,9 @@ spec: metadata: annotations: config.linkerd.io/skip-outbound-ports: "8761" config.linkerd.io/skip-inbound-ports: "8761" labels: app: e2e-orchestratorservice spec: Loading src/e2e_orchestrator/service/E2EOrchestratorServiceServicerImpl.py +84 −73 Original line number Diff line number Diff line Loading @@ -12,23 +12,23 @@ # See the License for the specific language governing permissions and # limitations under the License. import logging import networkx as nx import grpc import copy from websockets.sync.client import connect import time from common.method_wrappers.Decorator import MetricsPool, safe_and_metered_rpc_method from common.proto.e2eorchestrator_pb2 import E2EOrchestratorRequest, E2EOrchestratorReply from common.proto.context_pb2 import Empty, Connection, EndPointId, Link, LinkId from common.proto.e2eorchestrator_pb2_grpc import E2EOrchestratorServiceServicer from context.client.ContextClient import ContextClient from context.service.database.uuids.EndPoint import endpoint_get_uuid from common.proto.vnt_manager_pb2 import VNTSubscriptionRequest, VNTSubscriptionReply from common.proto.vnt_manager_pb2 import VNTSubscriptionRequest from common.tools.grpc.Tools import grpc_message_to_json_string from websockets.sync.server import serve import grpc import json import logging import networkx as nx from threading import Thread import time from websockets.sync.client import connect from websockets.sync.server import serve LOGGER = logging.getLogger(__name__) Loading @@ -37,6 +37,78 @@ METRICS_POOL = MetricsPool("E2EOrchestrator", "RPC") context_client: ContextClient = ContextClient() EXT_HOST = "nbiservice.tfs-ip.svc.cluster.local" EXT_PORT = "8762" OWN_HOST = "e2e-orchestratorservice.tfs-e2e.svc.cluster.local" OWN_PORT = "8761" def _event_received(websocket): for message in websocket: LOGGER.info("Message received!!!: {}".format(message)) message_json = json.loads(message) if 'event_type' in message_json: pass elif 'link_id' in message_json: obj = Link(**message_json) if _check_policies(obj): pass elif 'link_uuid' in message_json: obj = LinkId(**message_json) websocket.send(message) def _check_policies(link): return True def requestSubscription(): url = "ws://" + EXT_HOST + ":" + EXT_PORT request = VNTSubscriptionRequest() request.host = OWN_HOST request.port = OWN_PORT LOGGER.debug("Trying to connect to {}".format(url)) try: websocket = connect(url) LOGGER.debug("Connected to {}".format(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) LOGGER.debug("Sending {}".format(send)) websocket.send(send) try: message = websocket.recv() LOGGER.debug("Received message from WebSocket: {}".format(message)) except Exception as ex: LOGGER.info('Exception receiving from WebSocket: {}'.format(ex)) events_server() LOGGER.info('Subscription requested') def events_server(): all_hosts = "0.0.0.0" try: server = serve(_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)) with server: LOGGER.info("Running events server...: {}:{}".format(all_hosts, OWN_PORT)) server.serve_forever() LOGGER.info("Exiting events server...") class E2EOrchestratorServiceServicerImpl(E2EOrchestratorServiceServicer): Loading @@ -47,9 +119,10 @@ class E2EOrchestratorServiceServicerImpl(E2EOrchestratorServiceServicer): time.sleep(5) try: LOGGER.info("Requesting subscription") self.RequestSubscription() except Exception as E: LOGGER.info("Exception0!: {}".format(E)) subscription_thread = Thread(target=requestSubscription) subscription_thread.start() except Exception as ex: LOGGER.info("Exception!: {}".format(ex)) Loading Loading @@ -104,66 +177,4 @@ class E2EOrchestratorServiceServicerImpl(E2EOrchestratorServiceServicer): path.connections.append(conn) def RequestSubscription(self): LOGGER.info("Trying to connect...!!!") EXT_HOST = "nbiservice.tfs-ip.svc.cluster.local" EXT_PORT = "8762" OWN_HOST = "e2e-orchestratorservice.tfs-e2e.svc.cluster.local" OWN_PORT = "8761" url = "ws://" + EXT_HOST + ":" + EXT_PORT request = VNTSubscriptionRequest() request.host = OWN_HOST request.port = OWN_PORT LOGGER.info("Trying to connect... to {}".format(url)) with connect(url) as websocket: LOGGER.info("CONNECTED!!! {}") send = grpc_message_to_json_string(request) LOGGER.info("Sending {}".format(send)) websocket.send(send) try: message = websocket.recv() except Exception as e: LOGGER.info('Exception1!: {}'.format(e)) try: LOGGER.info("Received ws: {}".format(message)) except Exception as e: LOGGER.info('Exception2!: {}'.format(e)) with serve(self._event_received, "0.0.0.0", OWN_PORT, logger=LOGGER) as server: LOGGER.info("Running subscription server...: {}:{}".format("0.0.0.0", OWN_PORT)) server.serve_forever() LOGGER.info("Exiting subscription server...") def _event_received(self, websocket): for message in websocket: LOGGER.info("Message received!!!: {}".format(message)) message_json = json.loads(message) if 'event_type' in message_json: pass elif 'link_id' in message_json: obj = Link(**message_json) if self._check_policies(obj): pass elif 'link_uuid' in message_json: obj = LinkId(**message_json) websocket.send(message) def _check_policies(self, link): return True src/nbi/Dockerfile +2 −3 Original line number Diff line number Diff line Loading @@ -89,9 +89,8 @@ COPY src/service/__init__.py service/__init__.py COPY src/service/client/. service/client/ COPY src/slice/__init__.py slice/__init__.py COPY src/slice/client/. slice/client/ # COPY src/vnt_manager/__init__.py vnt_manager/__init__.py # COPY src/vnt_manager/client/. vnt_manager/client/ COPY --chown=teraflow:teraflow ./src/vnt_manager/. vnt_manager COPY src/vnt_manager/__init__.py vnt_manager/__init__.py COPY src/vnt_manager/client/. vnt_manager/client/ RUN mkdir -p /var/teraflow/tests/tools COPY src/tests/tools/mock_osm/. tests/tools/mock_osm/ Loading src/nbi/service/__main__.py +1 −1 Original line number Diff line number Diff line Loading @@ -26,7 +26,7 @@ from .rest_server.nbi_plugins.ietf_l3vpn import register_ietf_l3vpn from .rest_server.nbi_plugins.ietf_network import register_ietf_network from .rest_server.nbi_plugins.ietf_network_slice import register_ietf_nss from .rest_server.nbi_plugins.tfs_api import register_tfs_api from .rest_server.nbi_plugins.context_subscription import register_context_subscription from .context_subscription import register_context_subscription terminate = threading.Event() LOGGER = None Loading src/nbi/service/rest_server/nbi_plugins/context_subscription/__init__.py→src/nbi/service/context_subscription/__init__.py +4 −12 Original line number Diff line number Diff line Loading @@ -50,25 +50,17 @@ def register_context_subscription(): def subcript_to_vnt_manager(websocket): for message in websocket: LOGGER.info("Message received: {}".format(message)) LOGGER.debug("Message received: {}".format(message)) message_json = json.loads(message) request = VNTSubscriptionRequest() request.host = message_json['host'] request.port = message_json['port'] LOGGER.info("Received gRPC from ws: {}".format(request)) LOGGER.debug("Received gRPC from ws: {}".format(request)) reply = VNTSubscriptionReply() try: vntm_reply = vnt_manager_client.VNTSubscript(request) LOGGER.info("Received gRPC from vntm: {}".format(vntm_reply)) LOGGER.debug("Received gRPC from vntm: {}".format(vntm_reply)) except Exception as e: LOGGER.error('Could not subscript to VTNManager: {}'.format(e)) reply.subscription = "NOT OK" else: reply.subscription = "OK" websocket.send(reply.subscription) websocket.send(vntm_reply.subscription) Loading
manifests/e2e_orchestratorservice.yaml +3 −0 Original line number Diff line number Diff line Loading @@ -25,6 +25,9 @@ spec: metadata: annotations: config.linkerd.io/skip-outbound-ports: "8761" config.linkerd.io/skip-inbound-ports: "8761" labels: app: e2e-orchestratorservice spec: Loading
src/e2e_orchestrator/service/E2EOrchestratorServiceServicerImpl.py +84 −73 Original line number Diff line number Diff line Loading @@ -12,23 +12,23 @@ # See the License for the specific language governing permissions and # limitations under the License. import logging import networkx as nx import grpc import copy from websockets.sync.client import connect import time from common.method_wrappers.Decorator import MetricsPool, safe_and_metered_rpc_method from common.proto.e2eorchestrator_pb2 import E2EOrchestratorRequest, E2EOrchestratorReply from common.proto.context_pb2 import Empty, Connection, EndPointId, Link, LinkId from common.proto.e2eorchestrator_pb2_grpc import E2EOrchestratorServiceServicer from context.client.ContextClient import ContextClient from context.service.database.uuids.EndPoint import endpoint_get_uuid from common.proto.vnt_manager_pb2 import VNTSubscriptionRequest, VNTSubscriptionReply from common.proto.vnt_manager_pb2 import VNTSubscriptionRequest from common.tools.grpc.Tools import grpc_message_to_json_string from websockets.sync.server import serve import grpc import json import logging import networkx as nx from threading import Thread import time from websockets.sync.client import connect from websockets.sync.server import serve LOGGER = logging.getLogger(__name__) Loading @@ -37,6 +37,78 @@ METRICS_POOL = MetricsPool("E2EOrchestrator", "RPC") context_client: ContextClient = ContextClient() EXT_HOST = "nbiservice.tfs-ip.svc.cluster.local" EXT_PORT = "8762" OWN_HOST = "e2e-orchestratorservice.tfs-e2e.svc.cluster.local" OWN_PORT = "8761" def _event_received(websocket): for message in websocket: LOGGER.info("Message received!!!: {}".format(message)) message_json = json.loads(message) if 'event_type' in message_json: pass elif 'link_id' in message_json: obj = Link(**message_json) if _check_policies(obj): pass elif 'link_uuid' in message_json: obj = LinkId(**message_json) websocket.send(message) def _check_policies(link): return True def requestSubscription(): url = "ws://" + EXT_HOST + ":" + EXT_PORT request = VNTSubscriptionRequest() request.host = OWN_HOST request.port = OWN_PORT LOGGER.debug("Trying to connect to {}".format(url)) try: websocket = connect(url) LOGGER.debug("Connected to {}".format(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) LOGGER.debug("Sending {}".format(send)) websocket.send(send) try: message = websocket.recv() LOGGER.debug("Received message from WebSocket: {}".format(message)) except Exception as ex: LOGGER.info('Exception receiving from WebSocket: {}'.format(ex)) events_server() LOGGER.info('Subscription requested') def events_server(): all_hosts = "0.0.0.0" try: server = serve(_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)) with server: LOGGER.info("Running events server...: {}:{}".format(all_hosts, OWN_PORT)) server.serve_forever() LOGGER.info("Exiting events server...") class E2EOrchestratorServiceServicerImpl(E2EOrchestratorServiceServicer): Loading @@ -47,9 +119,10 @@ class E2EOrchestratorServiceServicerImpl(E2EOrchestratorServiceServicer): time.sleep(5) try: LOGGER.info("Requesting subscription") self.RequestSubscription() except Exception as E: LOGGER.info("Exception0!: {}".format(E)) subscription_thread = Thread(target=requestSubscription) subscription_thread.start() except Exception as ex: LOGGER.info("Exception!: {}".format(ex)) Loading Loading @@ -104,66 +177,4 @@ class E2EOrchestratorServiceServicerImpl(E2EOrchestratorServiceServicer): path.connections.append(conn) def RequestSubscription(self): LOGGER.info("Trying to connect...!!!") EXT_HOST = "nbiservice.tfs-ip.svc.cluster.local" EXT_PORT = "8762" OWN_HOST = "e2e-orchestratorservice.tfs-e2e.svc.cluster.local" OWN_PORT = "8761" url = "ws://" + EXT_HOST + ":" + EXT_PORT request = VNTSubscriptionRequest() request.host = OWN_HOST request.port = OWN_PORT LOGGER.info("Trying to connect... to {}".format(url)) with connect(url) as websocket: LOGGER.info("CONNECTED!!! {}") send = grpc_message_to_json_string(request) LOGGER.info("Sending {}".format(send)) websocket.send(send) try: message = websocket.recv() except Exception as e: LOGGER.info('Exception1!: {}'.format(e)) try: LOGGER.info("Received ws: {}".format(message)) except Exception as e: LOGGER.info('Exception2!: {}'.format(e)) with serve(self._event_received, "0.0.0.0", OWN_PORT, logger=LOGGER) as server: LOGGER.info("Running subscription server...: {}:{}".format("0.0.0.0", OWN_PORT)) server.serve_forever() LOGGER.info("Exiting subscription server...") def _event_received(self, websocket): for message in websocket: LOGGER.info("Message received!!!: {}".format(message)) message_json = json.loads(message) if 'event_type' in message_json: pass elif 'link_id' in message_json: obj = Link(**message_json) if self._check_policies(obj): pass elif 'link_uuid' in message_json: obj = LinkId(**message_json) websocket.send(message) def _check_policies(self, link): return True
src/nbi/Dockerfile +2 −3 Original line number Diff line number Diff line Loading @@ -89,9 +89,8 @@ COPY src/service/__init__.py service/__init__.py COPY src/service/client/. service/client/ COPY src/slice/__init__.py slice/__init__.py COPY src/slice/client/. slice/client/ # COPY src/vnt_manager/__init__.py vnt_manager/__init__.py # COPY src/vnt_manager/client/. vnt_manager/client/ COPY --chown=teraflow:teraflow ./src/vnt_manager/. vnt_manager COPY src/vnt_manager/__init__.py vnt_manager/__init__.py COPY src/vnt_manager/client/. vnt_manager/client/ RUN mkdir -p /var/teraflow/tests/tools COPY src/tests/tools/mock_osm/. tests/tools/mock_osm/ Loading
src/nbi/service/__main__.py +1 −1 Original line number Diff line number Diff line Loading @@ -26,7 +26,7 @@ from .rest_server.nbi_plugins.ietf_l3vpn import register_ietf_l3vpn from .rest_server.nbi_plugins.ietf_network import register_ietf_network from .rest_server.nbi_plugins.ietf_network_slice import register_ietf_nss from .rest_server.nbi_plugins.tfs_api import register_tfs_api from .rest_server.nbi_plugins.context_subscription import register_context_subscription from .context_subscription import register_context_subscription terminate = threading.Event() LOGGER = None Loading
src/nbi/service/rest_server/nbi_plugins/context_subscription/__init__.py→src/nbi/service/context_subscription/__init__.py +4 −12 Original line number Diff line number Diff line Loading @@ -50,25 +50,17 @@ def register_context_subscription(): def subcript_to_vnt_manager(websocket): for message in websocket: LOGGER.info("Message received: {}".format(message)) LOGGER.debug("Message received: {}".format(message)) message_json = json.loads(message) request = VNTSubscriptionRequest() request.host = message_json['host'] request.port = message_json['port'] LOGGER.info("Received gRPC from ws: {}".format(request)) LOGGER.debug("Received gRPC from ws: {}".format(request)) reply = VNTSubscriptionReply() try: vntm_reply = vnt_manager_client.VNTSubscript(request) LOGGER.info("Received gRPC from vntm: {}".format(vntm_reply)) LOGGER.debug("Received gRPC from vntm: {}".format(vntm_reply)) except Exception as e: LOGGER.error('Could not subscript to VTNManager: {}'.format(e)) reply.subscription = "NOT OK" else: reply.subscription = "OK" websocket.send(reply.subscription) websocket.send(vntm_reply.subscription)