Commit 59437ffe authored by Waleed Akbar's avatar Waleed Akbar
Browse files

Implement idle telemetry handling and connection count ramping in SIMAP synthesizer

parent a28d50cd
Loading
Loading
Loading
Loading
+34 −3
Changes for src/simap_connector/service/simap_updater/SimapUpdater.py: 34 added lines, 3 removed lines.
Original line number Diff line number Diff line
@@ -865,6 +865,36 @@ class EventDispatcher(BaseEventDispatcher):
        return active_count


    def _stop_link_synthesizer(self, worker_name : str, link_name : str) -> None:
        """Stop a link's synthesizer, leaving its SIMAP telemetry describing an idle link.

        Stopping the worker only ends the sampling loop. The last values it wrote stay in
        the SIMAP datastore, so the link keeps advertising the utilization of the
        connection that was just removed, and every collector polling that link keeps
        reading the same stale value for as long as the datastore lives. Write one idle
        sample once the sampling thread has joined, so the two cannot interleave.
        """
        worker = self._telemetry_pool.get_worker(WorkerTypeEnum.SYNTHESIZER, worker_name)
        self._telemetry_pool.stop_worker(WorkerTypeEnum.SYNTHESIZER, worker_name)

        if worker is None:
            MSG = 'Worker {:s} vanished before idle telemetry could be written'
            LOGGER.warning(MSG.format(worker_name))
            return

        assert isinstance(worker, SynthesizerWorker), \
            'Expected SynthesizerWorker, got {:s}'.format(type(worker).__name__)

        try:
            worker.write_idle_telemetry()
            LOGGER.info('Reset SIMAP telemetry of link {:s} to idle'.format(link_name))
        except Exception:
            # The worker is already stopped; a failed idle write must not abort the
            # removal of the remaining links.
            MSG = 'Failed to write idle telemetry for link {:s}'
            LOGGER.exception(MSG.format(link_name))


    def dispatch_connection_create(self, connection_event : ConnectionEvent) -> None:
        if not self.dispatch_connection_set(connection_event): return

@@ -923,8 +953,9 @@ class EventDispatcher(BaseEventDispatcher):
                remaining_conn_count = self._count_active_connections(link_uuid, domain_name)
                
                if remaining_conn_count == 0:
                    # No other connections use this link, stop the worker
                    self._telemetry_pool.stop_worker(WorkerTypeEnum.SYNTHESIZER, worker_name)
                    # No other connections use this link, stop the worker and leave the
                    # link reporting idle telemetry instead of its last busy sample.
                    self._stop_link_synthesizer(worker_name, link_name)
                    LOGGER.info('Stopped telemetry worker for link {:s}, no connections remain'.format(link_name))

                    # ---- TEMPORARY: Stop triggered links (L3 and L13 when L6 is removed from trans-pkt) ----
@@ -939,7 +970,7 @@ class EventDispatcher(BaseEventDispatcher):
                                    trig_worker_name = '{:s}:{:s}'.format(trig_link_topology_name, trig_link_name)
                                    
                                    if self._telemetry_pool.has_worker(WorkerTypeEnum.SYNTHESIZER, trig_worker_name):
                                        self._telemetry_pool.stop_worker(WorkerTypeEnum.SYNTHESIZER, trig_worker_name)
                                        self._stop_link_synthesizer(trig_worker_name, trig_link_name)
                                        LOGGER.info('Stopped triggered telemetry worker for link {:s}'.format(trig_link_name))
                                    else:
                                        LOGGER.warning('Triggered worker {:s} not found during cleanup'.format(trig_worker_name))
+59 −1
Changes for src/simap_connector/service/telemetry/worker/SynthesizerWorker.py: 59 added lines, 1 removed line.
Original line number Diff line number Diff line
@@ -15,6 +15,7 @@

import math, threading, time
from typing import Optional
from common.Settings import get_setting
from simap_connector.service.simap_updater.SimapClient import SimapClient
from .data.Resources import Resources
from ._Worker import _Worker, WorkerTypeEnum
@@ -22,6 +23,15 @@ from ._Worker import _Worker, WorkerTypeEnum

WAIT_LOOP_GRANULARITY = 0.5

IDLE_CONNECTION_COUNT = 0

# Seconds a link takes to move between congestion levels when its connection count
# changes, so that provisioning or tearing down a service shows congestion building up
# and draining away instead of stepping instantly. Set to 0 for immediate steps.
# Keep it comfortably above the telemetry subscription period (10s in the MWC scenario),
# otherwise subscribers sample the transition too coarsely to see it as a ramp.
TELEMETRY_RAMP_SECONDS = float(get_setting('TELEMETRY_RAMP_SECONDS', default='30.0'))


class SynthesizerWorker(_Worker):
    def __init__(
@@ -34,10 +44,58 @@ class SynthesizerWorker(_Worker):
        self._resources = resources
        self._sampling_interval = sampling_interval

        # The ramp is configured in seconds but applied in samples, and only this worker
        # knows how often it samples.
        self._ramp_samples = (
            0 if TELEMETRY_RAMP_SECONDS <= 0.0 or sampling_interval <= 0.0 else
            int(math.ceil(TELEMETRY_RAMP_SECONDS / sampling_interval))
        )

    def change_resources(self, connection_count: int) -> None:
        with self._lock:
            self._set_connection_count(connection_count)

    def _set_connection_count(self, connection_count : int, force_reset : bool = False) -> None:
        # Caller must hold self._lock.
        for link in self._resources.links:
                link.metrics_sampler.connection_count = connection_count
            sampler = link.metrics_sampler
            changed = sampler.connection_count != connection_count
            sampler.connection_count = connection_count

            if force_reset:
                # Adopt the new range at once, with no transition. Used for the idle
                # write, which happens after the sampling loop has stopped and therefore
                # only ever produces one sample: a ramp there would never be played out.
                sampler.reset()
            elif changed:
                # Each sample is derived from the previous one, so carrying the history
                # across a change of connection count leaves the value clamped to the
                # near edge of the new range, and it then has to random walk towards the
                # average of that range. Ramping up hides this (the clamp jumps the value
                # straight up), but tearing a service down does not: the link would keep
                # reporting the old congestion level for several minutes. Retargeting
                # replaces that with a bounded, symmetric transition onto the average of
                # the new range, so the value a link reports depends on its connection
                # count alone and not on the order it got there.
                sampler.retarget(self._ramp_samples)

    def write_idle_telemetry(self) -> None:
        """Publish one sample describing a link with no connections on it.

        Stopping the worker only stops the sampling loop; whatever it wrote last stays
        in the SIMAP datastore, so the link keeps advertising the utilization of the
        connection that was just removed, and collectors polling it keep reading that
        stale value indefinitely. Writing one idle sample leaves the datastore in a
        state that matches reality.
        """
        with self._lock:
            # force_reset: the idle write must land on the idle average whatever the
            # worker was doing before, including when it is already marked idle.
            self._set_connection_count(IDLE_CONNECTION_COUNT, force_reset=True)
            self._resources.generate_samples(self._simap_client)

        MSG = '[write_idle_telemetry] Wrote idle telemetry for {:d} link(s)'
        self._logger.info(MSG.format(len(self._resources.links)))

    def run(self) -> None:
        self._logger.info('[run] Starting...')
+106 −9
Changes for src/simap_connector/service/telemetry/worker/data/SyntheticSamplers.py: 106 added lines, 9 removed lines.
Original line number Diff line number Diff line
@@ -27,31 +27,47 @@ class SyntheticSampler:
    Bandwidth ranges based on connection count:
      0 conns: avg=3%,   range 1-10%
      1 conn:  avg=25%,  range 15-30%
      2 conns: avg=45%,  range 35-55%
      2 conns: avg=40%,  range 35-50%
      3 conns: avg=65%,  range 60-80%
      4+ conns: avg=85%, range 80-95%
    
    Latency uses bandwidth ranges divided by 10 (0-10ms):
      0 conns: avg=0.3ms, range 0.1-1.0ms
      1 conn:  avg=2.5ms, range 1.5-3.0ms
      2 conns: avg=4.5ms, range 3.5-5.5ms
      3 conns: avg=6.5ms, range 6.0-8.0ms
      4+ conns: avg=8.5ms, range 8.0-9.5ms
      0 conns: avg=0.4ms, range 0.1-0.8ms
      1 conn:  avg=1.4ms, range 1.0-1.8ms
      2 conns: avg=2.4ms, range 2.0-2.8ms
      3 conns: avg=3.4ms, range 3.0-3.8ms
      4+ conns: avg=4.4ms, range 4.0-4.8ms

    Values vary by ±1% between consecutive samples for temporal continuity.

    A change of connection_count moves the sampler to a different range. Call reset() to
    adopt the new range at once, or retarget(n) to glide into it over the next n samples
    so that congestion visibly builds up and drains away instead of stepping. Either way
    the sampler settles on the average of the new range, so the value a link reports
    depends on its connection count alone and not on the order it got there.
    """
    connection_count : int             = field(default = 0)
    link_capacity    : float           = field(default = 100.0)
    prev_bw          : Optional[float] = field(default = None)
    prev_latency     : Optional[float] = field(default = None)

    # Transition state: glide from (ramp_from_*) to the average of the current range over
    # ramp_total samples. ramp_left == 0 means no transition is in progress.
    ramp_left        : int             = field(default = 0)
    ramp_total       : int             = field(default = 0)
    ramp_from_bw     : Optional[float] = field(default = None)
    ramp_from_latency: Optional[float] = field(default = None)
    
    # Connection count to (avg, min, max) percentage mapping
    # Latency uses same ranges divided by 10 (0-10ms range)
    BW_RANGES = {
           0: (3,  5,  10),
           # NOTE: every row must satisfy min <= avg <= max. get_sample() clamps each
           # value into [min, max], so a row whose min exceeds its avg can never report
           # that average: the link silently sits at the minimum instead.
           0: (3,  1,  10),
           1: (25, 15, 30),
           2: (40, 35, 50),
           3: (60, 65, 80),
           3: (65, 60, 80),
           4: (85, 80, 95),
    }
    LAT_RANGES = {
@@ -71,6 +87,56 @@ class SyntheticSampler:
        """Factory method for compatibility (ignores unused parameters)."""
        return cls(connection_count=connection_count, link_capacity=link_capacity)

    def reset(self) -> None:
        """Drop temporal continuity with the previously sampled regime.

        Samples are derived from the previous value, so after a change of
        connection_count the value only converges to the edge of the new range and
        stays pinned there. Clearing the history makes the next sample the average of
        the new range, which is what an observer expects to see after the change.
        """
        self.prev_bw           = None
        self.prev_latency      = None
        self.ramp_left         = 0
        self.ramp_total        = 0
        self.ramp_from_bw      = None
        self.ramp_from_latency = None

    def retarget(self, ramp_samples : int = 0) -> None:
        """Converge on the range matching the current connection_count.

        With ramp_samples <= 0 this is reset(): the new range is adopted immediately.
        Otherwise the next ramp_samples samples interpolate from the current value to
        the average of the new range, so an operator provisioning or tearing down a
        service sees congestion build up and drain away rather than jump.

        Retargeting again mid-transition simply starts a new one from wherever the
        value has reached, so rapid changes of connection_count stay continuous.
        """
        if ramp_samples <= 0 or self.prev_bw is None or self.prev_latency is None:
            # Nothing to glide from (a freshly created sampler), or ramping disabled.
            self.reset()
            return

        self.ramp_from_bw      = self.prev_bw
        self.ramp_from_latency = self.prev_latency
        self.ramp_total        = ramp_samples
        self.ramp_left         = ramp_samples

    @staticmethod
    def _interpolate(
        value_from : float, value_to : float, progress : float, noise : float
    ) -> float:
        """Blend value_from into value_to, jittered, held inside the interval they span.

        Intermediate values deliberately fall outside the ranges of both the old and the
        new connection count, so they are bounded by the transition itself rather than
        clamped into either range.
        """
        value = value_from + (value_to - value_from) * progress
        value = value * (1.0 + noise)
        return max(min(value_from, value_to), min(max(value_from, value_to), value))

    def get_sample(self) -> Tuple[Sample, Sample]:
        """Generate bandwidth and latency samples with temporal continuity.
        
@@ -81,6 +147,38 @@ class SyntheticSampler:
        conn_key  = min(self.connection_count, 4)

        avg, min_bw, max_bw = self.BW_RANGES[conn_key]
        avg_lat, min_lat, max_lat = self.LAT_RANGES[conn_key]

        if self.ramp_left > 0:
            # Transition in progress: advance first, so the very first sample after a
            # change of connection_count already moves, and the last one lands exactly
            # on the average of the new range.
            self.ramp_left -= 1
            progress = 1.0 - (self.ramp_left / float(self.ramp_total))

            # The closing sample is taken without jitter, so a transition always ends
            # exactly on the average of the new range.
            final          = (self.ramp_left == 0)
            bw_noise       = 0.0 if final else random.uniform(-0.01, 0.01)
            latency_noise  = 0.0 if final else random.uniform(-0.05, 0.05)

            bw_utilization = self._interpolate(
                self.ramp_from_bw, avg, progress, bw_noise
            )
            latency = self._interpolate(
                self.ramp_from_latency, avg_lat, progress, latency_noise
            )

            self.prev_bw      = bw_utilization
            self.prev_latency = latency

            if self.ramp_left == 0:
                self.ramp_total        = 0
                self.ramp_from_bw      = None
                self.ramp_from_latency = None

            return (Sample(timestamp, 0, bw_utilization), Sample(timestamp, 0, latency))

        if self.prev_bw is None:
            bw_utilization = avg
        else:
@@ -90,7 +188,6 @@ class SyntheticSampler:
        bw_utilization = max(min_bw, min(max_bw, bw_utilization))
        self.prev_bw   = bw_utilization

        avg_lat, min_lat, max_lat = self.LAT_RANGES[conn_key]
        if self.prev_latency is None:
            latency = avg_lat
        else: