Commit bbf9e2eb authored by Lluis Gifre Renom's avatar Lluis Gifre Renom
Browse files

SIMAP Connector:

- Implemented creation/removal of workers for basic links
parent 48698c8d
Loading
Loading
Loading
Loading
+29 −2
Original line number Diff line number Diff line
@@ -25,7 +25,7 @@ from common.tools.grpc.BaseEventDispatcher import BaseEventDispatcher
from common.tools.grpc.Tools import grpc_message_to_json_string
from context.client.ContextClient import ContextClient
from simap_connector.service.simap_updater.MockSimaps import delete_mock_simap, set_mock_simap
from simap_connector.service.telemetry.Resources import Resources
from simap_connector.service.telemetry.Resources import ResourceLink, Resources, SyntheticSampler
from simap_connector.service.telemetry.TelemetryPool import TelemetryPool
from .ObjectCache import CachedEntities, ObjectCache
from .SimapClient import SimapClient
@@ -335,6 +335,30 @@ class EventDispatcher(BaseEventDispatcher):

        te_link = te_topo.link(link_name)
        te_link.update(src_device.name, src_endpoint.name, dst_device.name, dst_endpoint.name)

        worker_name = '{:s}:{:s}'.format(topology_name, link_name)
        resources = Resources()
        resources.links.append(ResourceLink(
            domain_name=topology_name, link_name=link_name,
            bandwidth_utilization_sampler=SyntheticSampler.create_random(
                amplitude_scale = 2.0,
                phase_scale     = 1e-7,
                period_scale    = 86_400,
                offset_scale    = 10_000_000,
                noise_ratio     = 0.05,
            ),
            latency_sampler=SyntheticSampler.create_random(
                amplitude_scale = 0.5,
                phase_scale     = 1e-7,
                period_scale    = 60.0,
                offset_scale    = 10.0,
                noise_ratio     = 0.05,
            ),
            related_service_ids=[],
        ))
        sampling_interval = 1.0
        self._telemetry_pool.start_worker(worker_name, resources, sampling_interval)

        return True


@@ -400,7 +424,10 @@ class EventDispatcher(BaseEventDispatcher):
        self._object_cache.delete(CachedEntities.LINK, link_uuid)
        self._object_cache.delete(CachedEntities.LINK, link_name)

        MSG = 'Link Remove: {:s}'
        worker_name = '{:s}:{:s}'.format(topology_name, link_name)
        self._telemetry_pool.stop_worker(worker_name)

        MSG = 'Link Removed: {:s}'
        LOGGER.info(MSG.format(grpc_message_to_json_string(link_event)))


+9 −7
Original line number Diff line number Diff line
@@ -21,28 +21,30 @@ from simap_connector.service.telemetry.SyntheticSamplers import SyntheticSampler

@dataclass
class ResourceNode:
    domain_name             : str
    node_name               : str
    cpu_utilization_sampler : SyntheticSampler
    related_service_ids     : List[str] = field(default_factory=list)

    def generate_samples(self, simap_client : SimapClient, domain_name : str) -> None:
    def generate_samples(self, simap_client : SimapClient) -> None:
        cpu_utilization = self.cpu_utilization_sampler.get_sample()
        simap_node = simap_client.network(domain_name).node(self.node_name)
        simap_node = simap_client.network(self.domain_name).node(self.node_name)
        simap_node.telemetry.update(
            cpu_utilization.value, related_service_ids=self.related_service_ids
        )

@dataclass
class ResourceLink:
    domain_name                   : str
    link_name                     : str
    bandwidth_utilization_sampler : SyntheticSampler
    latency_sampler               : SyntheticSampler
    related_service_ids           : List[str] = field(default_factory=list)

    def generate_samples(self, simap_client : SimapClient, domain_name : str) -> None:
    def generate_samples(self, simap_client : SimapClient) -> None:
        bandwidth_utilization = self.bandwidth_utilization_sampler.get_sample()
        latency               = self.latency_sampler.get_sample()
        simap_link = simap_client.network(domain_name).link(self.link_name)
        simap_link = simap_client.network(self.domain_name).link(self.link_name)
        simap_link.telemetry.update(
            bandwidth_utilization.value, latency.value,
            related_service_ids=self.related_service_ids
@@ -54,9 +56,9 @@ class Resources:
    nodes : List[ResourceNode] = field(default_factory=list)
    links : List[ResourceLink] = field(default_factory=list)

    def generate_samples(self, simap_client : SimapClient, domain_name : str) -> None:
    def generate_samples(self, simap_client : SimapClient) -> None:
        for resource in self.nodes:
            resource.generate_samples(simap_client, domain_name)
            resource.generate_samples(simap_client)

        for resource in self.links:
            resource.generate_samples(simap_client, domain_name)
            resource.generate_samples(simap_client)
+17 −14
Original line number Diff line number Diff line
@@ -32,41 +32,44 @@ class TelemetryPool:
        self._lock = threading.Lock()
        self._terminate = threading.Event() if terminate is None else terminate

    def has_worker(self, worker_name : str) -> bool:
        with self._lock:
            return worker_name in self._workers

    def start_worker(
        self, domain_name : str, resources : Resources, sampling_interval : float
        self, worker_name : str, resources : Resources, sampling_interval : float
    ) -> None:
        with self._lock:
            if domain_name in self._workers:
                MSG = '[start_worker] Worker already running for Domain({:s})'
                LOGGER.debug(MSG.format(str(domain_name)))
            if worker_name in self._workers:
                MSG = '[start_worker] Worker({:s}) already exists'
                LOGGER.debug(MSG.format(str(worker_name)))
                return

            worker = TelemetryWorker(
                domain_name, self._simap_client, resources, sampling_interval,
                worker_name, self._simap_client, resources, sampling_interval,
                terminate=self._terminate
            )
            worker.start()

            MSG = '[start_worker] Started worker for Domain({:s})'
            LOGGER.info(MSG.format(str(domain_name)))
            MSG = '[start_worker] Started Worker({:s})'
            LOGGER.info(MSG.format(str(worker_name)))

            self._workers[domain_name] = worker
            self._workers[worker_name] = worker


    def stop_worker(self, domain_name : str) -> None:
    def stop_worker(self, worker_name : str) -> None:
        with self._lock:
            worker = self._workers.pop(domain_name, None)
            worker = self._workers.pop(worker_name, None)

        if worker is None:
            MSG = '[stop_worker] No worker found for Domain({:s})'
            LOGGER.debug(MSG.format(str(domain_name)))
            MSG = '[stop_worker] Worker({:s}) not found'
            LOGGER.debug(MSG.format(str(worker_name)))
            return

        worker.stop()

        MSG = '[stop_worker] Stopped worker for Domain({:s})'
        LOGGER.info(MSG.format(str(domain_name)))
        MSG = '[stop_worker] Stopped Worker({:s})'
        LOGGER.info(MSG.format(str(worker_name)))


    def stop_all(self) -> None:
+9 −9
Original line number Diff line number Diff line
@@ -24,12 +24,12 @@ LOGGER = logging.getLogger(__name__)

class TelemetryWorker(threading.Thread):
    def __init__(
        self, domain_name : str, simap_client : SimapClient, resources : Resources,
        self, worker_name : str, simap_client : SimapClient, resources : Resources,
        sampling_interval : float, terminate : Optional[threading.Event] = None
    ) -> None:
        name = 'TelemetryWorker({:s})'.format(str(domain_name))
        name = 'TelemetryWorker({:s})'.format(str(worker_name))
        super().__init__(name=name, daemon=True)
        self._domain_name = domain_name
        self._worker_name = worker_name
        self._simap_client = simap_client
        self._resources = resources
        self._sampling_interval = sampling_interval
@@ -38,20 +38,20 @@ class TelemetryWorker(threading.Thread):

    def stop(self) -> None:
        MSG = '[stop][{:s}] Stopping...'
        LOGGER.info(MSG.format(str(self._domain_name)))
        LOGGER.info(MSG.format(str(self._worker_name)))
        self._stop_event.set()
        self.join()

    def run(self) -> None:
        MSG = '[run][{:s}] Starting...'
        LOGGER.info(MSG.format(str(self._domain_name)))
        LOGGER.info(MSG.format(str(self._worker_name)))

        try:
            while not self._stop_event.is_set() and not self._terminate.is_set():
                MSG = '[run][{:s}] Sampling...'
                LOGGER.info(MSG.format(str(self._domain_name)))
                LOGGER.info(MSG.format(str(self._worker_name)))

                self._resources.generate_samples(self._simap_client, self._domain_name)
                self._resources.generate_samples(self._simap_client)

                # Make wait responsible to terminations
                iterations = self._sampling_interval / 0.1
@@ -62,7 +62,7 @@ class TelemetryWorker(threading.Thread):

        except Exception:
            MSG = '[run][{:s}] Unhandled Exception'
            LOGGER.info(MSG.format(str(self._domain_name)))
            LOGGER.info(MSG.format(str(self._worker_name)))
        finally:
            MSG = '[run][{:s}] Terminated'
            LOGGER.info(MSG.format(str(self._domain_name)))
            LOGGER.info(MSG.format(str(self._worker_name)))