Commit 3f9c2b70 authored by Lluis Gifre Renom's avatar Lluis Gifre Renom
Browse files

E2E Orchestrator component:

- Updated Subscriptions Framework
- Corrected Controller Discovery mechanism
- Extended framework to support multiple Dispatchers
- Implemented Recommendations Dispatcher (being tested)
parent 66284d48
Loading
Loading
Loading
Loading
+33 −29
Original line number Diff line number Diff line
@@ -12,23 +12,20 @@
# See the License for the specific language governing permissions and
# limitations under the License.

import logging
import signal
import sys
import threading

import logging, signal, sys, threading
from prometheus_client import start_http_server

from common.Constants import ServiceNameEnum
from common.Settings import (ENVVAR_SUFIX_SERVICE_HOST,
                             ENVVAR_SUFIX_SERVICE_PORT_GRPC, get_env_var_name,
                             get_log_level, get_metrics_port,
                             wait_for_environment_variables)
from e2e_orchestrator.service.subscriptions.ControllerDiscovererThread import ControllerDiscoverer

from common.Settings import (
    ENVVAR_SUFIX_SERVICE_HOST, ENVVAR_SUFIX_SERVICE_PORT_GRPC, get_env_var_name,
    get_log_level, get_metrics_port, wait_for_environment_variables
)
from .subscriptions.ControllerDiscoverer import ControllerDiscoverer
from .subscriptions.Subscriptions import Subscriptions
from .subscriptions.dispatchers.Dispatchers import Dispatchers
from .subscriptions.dispatchers.recommendation.Dispatcher import RecommendationDispatcher
from .E2EOrchestratorService import E2EOrchestratorService

terminate = threading.Event()
TERMINATE = threading.Event()

LOG_LEVEL = get_log_level()
logging.basicConfig(level=LOG_LEVEL, format="[%(asctime)s] %(levelname)s:%(name)s:%(message)s")
@@ -36,40 +33,47 @@ LOGGER = logging.getLogger(__name__)


def signal_handler(signal, frame): # pylint: disable=redefined-outer-name
    LOGGER.warning("Terminate signal received")
    terminate.set()
    LOGGER.warning('Terminate signal received')
    TERMINATE.set()


def main():
    wait_for_environment_variables([
        get_env_var_name(ServiceNameEnum.CONTEXT, ENVVAR_SUFIX_SERVICE_HOST     ),
        get_env_var_name(ServiceNameEnum.CONTEXT, ENVVAR_SUFIX_SERVICE_PORT_GRPC),
    ])

    signal.signal(signal.SIGINT,  signal_handler)
    signal.signal(signal.SIGTERM, signal_handler)

    LOGGER.info("Starting...")
    LOGGER.info('Starting...')

    # Start metrics server
    metrics_port = get_metrics_port()
    start_http_server(metrics_port)

    # Starting service
    grpc_service = E2EOrchestratorService()
    grpc_service.start()

    controller_discoverer = ControllerDiscoverer(
        terminate=terminate
    )
    controller_discoverer.start()

    LOGGER.info("Running...")
    dispatchers   = Dispatchers(TERMINATE)
    dispatchers.add_dispatcher(RecommendationDispatcher)
    subscriptions = Subscriptions(dispatchers, TERMINATE)
    discoverer    = ControllerDiscoverer(subscriptions, TERMINATE)
    discoverer.start()

    LOGGER.info('Running...')
    # Wait for Ctrl+C or termination signal
    while not terminate.wait(timeout=1):
        pass
    while not TERMINATE.wait(timeout=1.0): pass

    LOGGER.info("Terminating...")
    controller_discoverer.stop()
    LOGGER.info('Terminating...')
    discoverer.stop()
    grpc_service.stop()

    LOGGER.info("Bye")
    LOGGER.info('Bye')
    return 0


if __name__ == "__main__":
if __name__ == '__main__':
    sys.exit(main())
+4 −12
Original line number Diff line number Diff line
@@ -68,24 +68,16 @@ class EventDispatcher(BaseEventDispatcher):

class ControllerDiscoverer:
    def __init__(
        self, terminate : Optional[threading.Event] = None
        self, subscriptions : Subscriptions, terminate : threading.Event
    ) -> None:
        self._context_client = ContextClient()

        self._event_collector = BaseEventCollector(
            terminate=terminate
        )
        self._event_collector = BaseEventCollector(terminate=terminate)
        self._event_collector.install_collector(
            self._context_client.GetDeviceEvents,
            Empty(), log_events_received=True
            self._context_client.GetDeviceEvents, Empty(), log_events_received=True
        )

        self._subscriptions = Subscriptions()

        self._event_dispatcher = EventDispatcher(
            self._event_collector.get_events_queue(),
            self._context_client,
            self._subscriptions,
            self._event_collector.get_events_queue(), self._context_client, subscriptions,
            terminate=terminate
        )

+9 −21
Original line number Diff line number Diff line
@@ -13,53 +13,41 @@
# limitations under the License.


import queue, socketio, threading
import socketio, threading
from common.Constants import ServiceNameEnum
from common.Settings import get_service_baseurl_http
from .RecommendationsClientNamespace import RecommendationsClientNamespace
from .dispatchers.Dispatchers import Dispatchers
from .TFSControllerSettings import TFSControllerSettings


NBI_SERVICE_PREFIX_URL = get_service_baseurl_http(ServiceNameEnum.NBI) or ''
CHILD_SOCKETIO_URL = 'http://{:s}:{:s}@{:s}:{:d}{:s}'
CHILD_SOCKETIO_URL = 'http://{:s}:{:s}@{:s}:{:d}' + NBI_SERVICE_PREFIX_URL


class Subscription(threading.Thread):
    def __init__(
        self, tfs_ctrl_settings : TFSControllerSettings,
        self, tfs_ctrl_settings : TFSControllerSettings, dispatchers : Dispatchers,
        terminate : threading.Event
    ) -> None:
        super().__init__(daemon=True)
        self._settings    = tfs_ctrl_settings
        self._dispatchers = dispatchers
        self._terminate   = terminate
        self._request_queue = queue.Queue()
        self._reply_queue   = queue.Queue()
        self._is_running  = threading.Event()

    @property
    def is_running(self): return self._is_running.is_set()

    @property
    def request_queue(self): return self._request_queue

    @property
    def reply_queue(self): return self._reply_queue

    def run(self) -> None:
        child_socketio_url = CHILD_SOCKETIO_URL.format(
            self._settings.nbi_username,
            self._settings.nbi_password,
            self._settings.nbi_address,
            self._settings.nbi_port,
            NBI_SERVICE_PREFIX_URL
        )

        namespace = RecommendationsClientNamespace(
            self._request_queue, self._reply_queue
        )

        sio = socketio.Client(logger=True, engineio_logger=True)
        sio.register_namespace(namespace)
        self._dispatchers.register(sio)
        sio.connect(child_socketio_url)

        while not self._terminate.is_set():
+7 −10
Original line number Diff line number Diff line
@@ -12,16 +12,18 @@
# See the License for the specific language governing permissions and
# limitations under the License.

import logging, queue, threading
import logging, threading
from typing import Dict
from .dispatchers.Dispatchers import Dispatchers
from .Subscription import Subscription
from .TFSControllerSettings import TFSControllerSettings

LOGGER = logging.getLogger(__name__)

class Subscriptions:
    def __init__(self) -> None:
        self._terminate = threading.Event()
    def __init__(self, dispatchers : Dispatchers, terminate : threading.Event) -> None:
        self._dispatchers = dispatchers
        self._terminate   = terminate
        self._lock        = threading.Lock()
        self._subscriptions : Dict[str, Subscription] = dict()

@@ -30,7 +32,7 @@ class Subscriptions:
        with self._lock:
            subscription = self._subscriptions.get(device_uuid)
            if (subscription is not None) and subscription.is_running: return
            subscription = Subscription(tfs_ctrl_settings, self._terminate)
            subscription = Subscription(tfs_ctrl_settings, self._dispatchers, self._terminate)
            self._subscriptions[device_uuid] = subscription
            subscription.start()

@@ -40,8 +42,3 @@ class Subscriptions:
            if subscription is None: return
            if subscription.is_running: subscription.stop()
            self._subscriptions.pop(device_uuid, None)

    def stop(self):
        self._terminate.set()
        for device_uuid in self._subscriptions:
            self.remove_subscription(device_uuid)
+14 −21
Original line number Diff line number Diff line
@@ -12,29 +12,22 @@
# See the License for the specific language governing permissions and
# limitations under the License.

import logging, queue, socketio
import logging, socketio, threading
from typing import List, Type
from ._Dispatcher import _Dispatcher

LOGGER = logging.getLogger(__name__)

class RecommendationsClientNamespace(socketio.ClientNamespace):
    def __init__(self, request_queue : queue.Queue, reply_queue : queue.Queue):
        self._request_queue = request_queue
        self._reply_queue   = reply_queue
        super().__init__(namespace='/recommendations')
class Dispatchers:
    def __init__(self, terminate : threading.Event) -> None:
        self._terminate = terminate
        self._dispatchers : List[_Dispatcher] = list()

    def on_connect(self):
        LOGGER.info('[on_connect] Connected')
    def add_dispatcher(self, dispatcher_class : Type[_Dispatcher]) -> None:
        dispatcher = dispatcher_class(self._terminate)
        self._dispatchers.append(dispatcher)
        dispatcher.start()

    def on_disconnect(self, reason):
        MSG = '[on_disconnect] Disconnected!, reason: {:s}'
        LOGGER.info(MSG.format(str(reason)))

    def on_recommendation(self, data):
        MSG = '[on_recommendation] data={:s}'
        LOGGER.info(MSG.format(str(data)))

        #MSG = '[on_recommendation] Recommendation: {:s}'
        #LOGGER.info(MSG.format(str(recommendation)))

        #request = (self._device_uuid, *sample)
        #self._request_queue.put_nowait(request)
    def register(self, sio_client : socketio.Client) -> None:
        for dispatcher in self._dispatchers:
            dispatcher.register(sio_client)
Loading