Resolve "(CTTC) End-to-End integration test with ETSI OSM RO IETF L2VPN"

Closes #367 (closed)

Merge request reports

Loading
+72 −0
Changes for src/common/tests/test_base_event_collector_retry.py: 72 added lines, 0 removed lines.
Original line number Diff line number Diff line
# Copyright 2022-2025 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 grpc
import common.tools.grpc.BaseEventCollector as base_event_collector_module
from common.proto.context_pb2 import DeviceEvent, EventTypeEnum


class _FakeRpcError(grpc.RpcError):
    def __init__(self, status_code):
        self._status_code = status_code

    def code(self):
        return self._status_code


class _FakeStream:
    def __init__(self, events=None):
        self._events = iter(events or [])

    def __iter__(self):
        return self

    def __next__(self):
        return next(self._events)

    def cancel(self):
        pass


def _create_device_event():
    event = DeviceEvent()
    event.event.event_type = EventTypeEnum.EVENTTYPE_CREATE
    event.event.timestamp.timestamp = 1.0
    event.device_id.device_uuid.uuid = 'dev1'
    return event


def test_base_event_collector_retries_if_subscription_creation_fails(monkeypatch):
    monkeypatch.setattr(base_event_collector_module.time, 'sleep', lambda _seconds: None)

    state = {'calls': 0}

    def subscription_method(_request):
        state['calls'] += 1
        if state['calls'] == 1:
            raise _FakeRpcError(grpc.StatusCode.UNAVAILABLE)
        if state['calls'] == 2:
            return _FakeStream(events=[_create_device_event()])
        raise _FakeRpcError(grpc.StatusCode.CANCELLED)

    collector = base_event_collector_module.BaseEventCollector()
    collector.install_collector(subscription_method, object())

    collector.start()
    try:
        event = collector.get_event(block=True, timeout=1.0)
    finally:
        collector.stop()

    assert event.device_id.device_uuid.uuid == 'dev1'
+1 −1
Changes for src/common/tools/grpc/BaseEventCollector.py: 1 added line, 1 removed line.
Original line number Diff line number Diff line
@@ -41,8 +41,8 @@ class CollectorThread(threading.Thread):

    def run(self) -> None:
        while not self._terminate.is_set():
            self._stream = self._subscription_func()
            try:
                self._stream = self._subscription_func()
                for event in self._stream:
                    if self._log_events_received:
                        str_event = grpc_message_to_json_string(event)
+19 −130
Changes for src/context/client/EventsCollector.py: 19 added lines, 130 removed lines.
Original line number Diff line number Diff line
@@ -12,52 +12,15 @@
# See the License for the specific language governing permissions and
# limitations under the License.

import grpc, logging, queue, threading, time
from typing import Callable
import logging
from common.proto.context_pb2 import Empty
from common.tools.grpc.Tools import grpc_message_to_json_string
from common.tools.grpc.BaseEventCollector import BaseEventCollector
from context.client.ContextClient import ContextClient

LOGGER = logging.getLogger(__name__)
LOGGER.setLevel(logging.DEBUG)

class _Collector(threading.Thread):
    def __init__(
        self, subscription_func : Callable, events_queue = queue.PriorityQueue,
        terminate = threading.Event, log_events_received: bool = False
    ) -> None:
        super().__init__(daemon=False)
        self._subscription_func = subscription_func
        self._events_queue = events_queue
        self._terminate = terminate
        self._log_events_received = log_events_received
        self._stream = None

    def cancel(self) -> None:
        if self._stream is None: return
        self._stream.cancel()

    def run(self) -> None:
        while not self._terminate.is_set():
            self._stream = self._subscription_func()
            try:
                for event in self._stream:
                    if self._log_events_received:
                        str_event = grpc_message_to_json_string(event)
                        LOGGER.info('[_collect] event: {:s}'.format(str_event))
                    timestamp = event.event.timestamp.timestamp
                    self._events_queue.put_nowait((timestamp, event))
            except grpc.RpcError as e:
                if e.code() == grpc.StatusCode.UNAVAILABLE:
                    LOGGER.info('[_collect] UNAVAILABLE... retrying...')
                    time.sleep(0.5)
                    continue
                elif e.code() == grpc.StatusCode.CANCELLED:
                    break
                else:
                    raise # pragma: no cover

class EventsCollector:
class EventsCollector(BaseEventCollector):
    def __init__(
        self, context_client          : ContextClient,
        log_events_received           : bool = False,
@@ -69,93 +32,19 @@ class EventsCollector:
        activate_slice_collector      : bool = True,
        activate_connection_collector : bool = True,
    ) -> None:
        self._events_queue = queue.PriorityQueue()
        self._terminate = threading.Event()
        self._log_events_received = log_events_received

        self._context_thread = _Collector(
                lambda: context_client.GetContextEvents(Empty()),
                self._events_queue, self._terminate, self._log_events_received
            ) if activate_context_collector else None

        self._topology_thread = _Collector(
                lambda: context_client.GetTopologyEvents(Empty()),
                self._events_queue, self._terminate, self._log_events_received
            ) if activate_topology_collector else None

        self._device_thread = _Collector(
                lambda: context_client.GetDeviceEvents(Empty()),
                self._events_queue, self._terminate, self._log_events_received
            ) if activate_device_collector else None

        self._link_thread = _Collector(
                lambda: context_client.GetLinkEvents(Empty()),
                self._events_queue, self._terminate, self._log_events_received
            ) if activate_link_collector else None

        self._service_thread = _Collector(
                lambda: context_client.GetServiceEvents(Empty()),
                self._events_queue, self._terminate, self._log_events_received
            ) if activate_service_collector else None

        self._slice_thread = _Collector(
                lambda: context_client.GetSliceEvents(Empty()),
                self._events_queue, self._terminate, self._log_events_received
            ) if activate_slice_collector else None

        self._connection_thread = _Collector(
                lambda: context_client.GetConnectionEvents(Empty()),
                self._events_queue, self._terminate, self._log_events_received
            ) if activate_connection_collector else None

    def start(self):
        self._terminate.clear()

        if self._context_thread    is not None: self._context_thread.start()
        if self._topology_thread   is not None: self._topology_thread.start()
        if self._device_thread     is not None: self._device_thread.start()
        if self._link_thread       is not None: self._link_thread.start()
        if self._service_thread    is not None: self._service_thread.start()
        if self._slice_thread      is not None: self._slice_thread.start()
        if self._connection_thread is not None: self._connection_thread.start()

    def get_event(self, block : bool = True, timeout : float = 0.1):
        try:
            _,event = self._events_queue.get(block=block, timeout=timeout)
            return event
        except queue.Empty: # pylint: disable=catching-non-exception
            return None

    def get_events(self, block : bool = True, timeout : float = 0.1, count : int = None):
        events = []
        if count is None:
            while not self._terminate.is_set():
                event = self.get_event(block=block, timeout=timeout)
                if event is None: break
                events.append(event)
        else:
            while len(events) < count:
                if self._terminate.is_set(): break
                event = self.get_event(block=block, timeout=timeout)
                if event is None: continue
                events.append(event)
        return sorted(events, key=lambda e: e.event.timestamp.timestamp)

    def stop(self):
        self._terminate.set()

        if self._context_thread    is not None: self._context_thread.cancel()
        if self._topology_thread   is not None: self._topology_thread.cancel()
        if self._device_thread     is not None: self._device_thread.cancel()
        if self._link_thread       is not None: self._link_thread.cancel()
        if self._service_thread    is not None: self._service_thread.cancel()
        if self._slice_thread      is not None: self._slice_thread.cancel()
        if self._connection_thread is not None: self._connection_thread.cancel()

        if self._context_thread    is not None: self._context_thread.join()
        if self._topology_thread   is not None: self._topology_thread.join()
        if self._device_thread     is not None: self._device_thread.join()
        if self._link_thread       is not None: self._link_thread.join()
        if self._service_thread    is not None: self._service_thread.join()
        if self._slice_thread      is not None: self._slice_thread.join()
        if self._connection_thread is not None: self._connection_thread.join()
        super().__init__()

        if activate_context_collector:
            self.install_collector(context_client.GetContextEvents, Empty(), log_events_received)
        if activate_topology_collector:
            self.install_collector(context_client.GetTopologyEvents, Empty(), log_events_received)
        if activate_device_collector:
            self.install_collector(context_client.GetDeviceEvents, Empty(), log_events_received)
        if activate_link_collector:
            self.install_collector(context_client.GetLinkEvents, Empty(), log_events_received)
        if activate_service_collector:
            self.install_collector(context_client.GetServiceEvents, Empty(), log_events_received)
        if activate_slice_collector:
            self.install_collector(context_client.GetSliceEvents, Empty(), log_events_received)
        if activate_connection_collector:
            self.install_collector(context_client.GetConnectionEvents, Empty(), log_events_received)
+84 −0
Changes for src/context/tests/test_events_collector_retry.py: 84 added lines, 0 removed lines.
Original line number Diff line number Diff line
# Copyright 2022-2025 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 grpc
import common.tools.grpc.BaseEventCollector as base_event_collector_module
import context.client.EventsCollector as context_events_module
from common.proto.context_pb2 import DeviceEvent, EventTypeEnum


class _FakeRpcError(grpc.RpcError):
    def __init__(self, status_code):
        self._status_code = status_code

    def code(self):
        return self._status_code


class _FakeStream:
    def __init__(self, events=None):
        self._events = iter(events or [])

    def __iter__(self):
        return self

    def __next__(self):
        return next(self._events)

    def cancel(self):
        pass


class _FakeContextClient:
    def __init__(self):
        self.calls = 0

    def GetDeviceEvents(self, _request):
        self.calls += 1
        if self.calls == 1:
            raise _FakeRpcError(grpc.StatusCode.UNAVAILABLE)
        if self.calls == 2:
            return _FakeStream(events=[_create_device_event()])
        raise _FakeRpcError(grpc.StatusCode.CANCELLED)


def _create_device_event():
    event = DeviceEvent()
    event.event.event_type = EventTypeEnum.EVENTTYPE_CREATE
    event.event.timestamp.timestamp = 1.0
    event.device_id.device_uuid.uuid = 'dev1'
    return event


def test_events_collector_retries_if_subscription_creation_fails(monkeypatch):
    monkeypatch.setattr(base_event_collector_module.time, 'sleep', lambda _seconds: None)

    collector = context_events_module.EventsCollector(
        _FakeContextClient(),
        activate_context_collector=False,
        activate_topology_collector=False,
        activate_device_collector=True,
        activate_link_collector=False,
        activate_service_collector=False,
        activate_slice_collector=False,
        activate_connection_collector=False,
    )

    collector.start()
    try:
        event = collector.get_event(block=True, timeout=1.0)
    finally:
        collector.stop()

    assert event.device_id.device_uuid.uuid == 'dev1'
+120 −26
Changes for src/device/service/drivers/gnmi_openconfig/handlers/Interface.py: 120 added lines, 26 removed lines.
Original line number Diff line number Diff line
@@ -20,24 +20,86 @@ from .YangHandler import YangHandler

LOGGER = logging.getLogger(__name__)

MIN_MTU = 68
EOS_TAGGED_L3_REPLACE_FIELD = '_eos_tagged_l3_replace'


def _normalize_interface_type(if_type: str, sif_index: int) -> str:
    if if_type == 'l3ipvlan' and sif_index is not None:
        return 'iana-if-type:ethernetCsmacd'
    if ':' not in if_type:
        return 'iana-if-type:{:s}'.format(if_type)
    return if_type

class InterfaceHandler(_Handler):
    def get_resource_key(self) -> str: return '/interface/subinterface'
    def get_path(self) -> str: return '/openconfig-interfaces:interfaces'

    @staticmethod
    def _create_subinterface(
        yang_sifs: libyang.DContainer, sif_index: int, enabled: bool = None,
        vlan_id: int = None, address_ip: str = None, address_prefix: int = None
    ) -> libyang.DContainer:
        yang_sif_path = 'subinterface[index="{:d}"]'.format(sif_index)
        yang_sif: libyang.DContainer = yang_sifs.create_path(yang_sif_path)
        yang_sif.create_path('config/index', sif_index)
        if enabled is not None:
            yang_sif.create_path('config/enabled', enabled)

        if vlan_id is not None:
            yang_subif_vlan : libyang.DContainer = yang_sif.create_path('openconfig-vlan:vlan')
            yang_subif_vlan.create_path('match/single-tagged/config/vlan-id', vlan_id)

        yang_ipv4 : libyang.DContainer = yang_sif.create_path('openconfig-if-ip:ipv4')
        if enabled is not None:
            yang_ipv4.create_path('config/enabled', enabled)

        if address_ip is not None and address_prefix is not None:
            yang_ipv4_addrs : libyang.DContainer = yang_ipv4.create_path('addresses')
            yang_ipv4_addr_path = 'address[ip="{:s}"]'.format(address_ip)
            yang_ipv4_addr : libyang.DContainer = yang_ipv4_addrs.create_path(yang_ipv4_addr_path)
            yang_ipv4_addr.create_path('config/ip',            address_ip)
            yang_ipv4_addr.create_path('config/prefix-length', address_prefix)

        return yang_sif

    def compose(
        self, resource_key : str, resource_value : Dict, yang_handler : YangHandler, delete : bool = False
    ) -> Tuple[str, str]:
        if_name   = get_str(resource_value, 'name'       )  # ethernet-1/1
        sif_index = get_int(resource_value, 'index', 0)  # 0
        sif_index = get_int(resource_value, 'index', None)  # 0
        vlan_id   = get_int(resource_value, 'vlan_id', None)
        eos_tagged_l3_replace = get_bool(resource_value, EOS_TAGGED_L3_REPLACE_FIELD, False)

        if delete:
            PATH_TMPL = '/interfaces/interface[name={:s}]/subinterfaces/subinterface[index={:d}]'
            str_path = PATH_TMPL.format(if_name, sif_index)
            if eos_tagged_l3_replace:
                root_node : libyang.DContainer = yang_handler.get_data_path(
                    '/openconfig-interfaces:interfaces'
                )
                str_path = '/interfaces/interface[name={:s}]'.format(if_name)
                yang_if = root_node.find_path('/'.join([
                    '',
                    'openconfig-interfaces:interfaces',
                    'interface[name="{:s}"]'.format(if_name),
                ]))
                if yang_if is not None:
                    yang_if.unlink()
                    yang_if.free()
                str_data = json.dumps({})
                return str_path, str_data

            if sif_index is None:
                return None, None

            root_node : libyang.DContainer = yang_handler.get_data_path(
                '/openconfig-interfaces:interfaces'
            )

            address_ip = get_str(resource_value, 'address_ip', None)
            if address_ip is None:
                PATH_TMPL = '/interfaces/interface[name={:s}]/subinterfaces/subinterface[index={:d}]'
                str_path = PATH_TMPL.format(if_name, sif_index)

                yang_sif = root_node.find_path('/'.join([
                    '', # add slash at the beginning
                    'openconfig-interfaces:interfaces',
@@ -45,15 +107,33 @@ class InterfaceHandler(_Handler):
                    'subinterfaces',
                    'subinterface[index="{:d}"]'.format(sif_index),
                ]))
            else:
                PATH_TMPL = (
                    '/interfaces/interface[name={:s}]/subinterfaces/subinterface[index={:d}]'
                    '/openconfig-if-ip:ipv4/addresses/address[ip={:s}]'
                )
                str_path = PATH_TMPL.format(if_name, sif_index, address_ip)

                yang_sif = root_node.find_path('/'.join([
                    '', # add slash at the beginning
                    'openconfig-interfaces:interfaces',
                    'interface[name="{:s}"]'.format(if_name),
                    'subinterfaces',
                    'subinterface[index="{:d}"]'.format(sif_index),
                    'openconfig-if-ip:ipv4',
                    'addresses',
                    'address[ip="{:s}"]'.format(address_ip)
                ]))

            if yang_sif is not None:
                yang_sif.unlink()
                yang_sif.free()

            str_data = json.dumps({})
            return str_path, str_data

        enabled        = get_bool(resource_value, 'enabled',  True) # True/False
        #if_type        = get_str (resource_value, 'type'         ) # 'l3ipvlan'
        vlan_id        = get_int (resource_value, 'vlan_id',      ) # 127
        enabled        = get_bool(resource_value, 'enabled'       ) # True/False
        if_type        = get_str (resource_value, 'type'          ) # 'l3ipvlan'
        address_ip     = get_str (resource_value, 'address_ip'    ) # 172.16.0.1
        address_prefix = get_int (resource_value, 'address_prefix') # 24
        mtu            = get_int (resource_value, 'mtu'           ) # 1500
@@ -62,28 +142,34 @@ class InterfaceHandler(_Handler):
        yang_if_path = 'interface[name="{:s}"]'.format(if_name)
        yang_if : libyang.DContainer = yang_ifs.create_path(yang_if_path)
        yang_if.create_path('config/name',    if_name   )
        if enabled is not None: yang_if.create_path('config/enabled', enabled)
        if mtu     is not None: yang_if.create_path('config/mtu',     mtu)
        if if_type is not None:
            yang_if.create_path('config/type', _normalize_interface_type(if_type, sif_index))
        if enabled is not None:
            yang_if.create_path('config/enabled', enabled)
        
        yang_sifs : libyang.DContainer = yang_if.create_path('subinterfaces')
        yang_sif_path = 'subinterface[index="{:d}"]'.format(sif_index)
        yang_sif : libyang.DContainer = yang_sifs.create_path(yang_sif_path)
        yang_sif.create_path('config/index', sif_index)
        if enabled is not None: yang_sif.create_path('config/enabled', enabled)
        if mtu is not None and mtu >= MIN_MTU:
            yang_if.create_path('config/mtu', mtu)

        if vlan_id is not None:
            yang_subif_vlan : libyang.DContainer = yang_sif.create_path('openconfig-vlan:vlan')
            yang_subif_vlan.create_path('match/single-tagged/config/vlan-id', vlan_id)

        yang_ipv4 : libyang.DContainer = yang_sif.create_path('openconfig-if-ip:ipv4')
        if enabled is not None: yang_ipv4.create_path('config/enabled', enabled)
        if sif_index is None:
            str_path = '/interfaces/interface[name={:s}]'.format(if_name)
            str_data = yang_if.print_mem('json')
            json_data = json.loads(str_data)
            json_data = json_data['openconfig-interfaces:interface'][0]
            str_data = json.dumps(json_data)
            return str_path, str_data

        if address_ip is not None and address_prefix is not None:
            yang_ipv4_addrs : libyang.DContainer = yang_ipv4.create_path('addresses')
            yang_ipv4_addr_path = 'address[ip="{:s}"]'.format(address_ip)
            yang_ipv4_addr : libyang.DContainer = yang_ipv4_addrs.create_path(yang_ipv4_addr_path)
            yang_ipv4_addr.create_path('config/ip',            address_ip)
            yang_ipv4_addr.create_path('config/prefix-length', address_prefix)
        yang_sifs : libyang.DContainer = yang_if.create_path('subinterfaces')
        if eos_tagged_l3_replace and sif_index == 0 and vlan_id is not None:
            self._create_subinterface(yang_sifs, 0, enabled=enabled)
            self._create_subinterface(
                yang_sifs, vlan_id, enabled=enabled, vlan_id=vlan_id,
                address_ip=address_ip, address_prefix=address_prefix
            )
        else:
            self._create_subinterface(
                yang_sifs, sif_index, enabled=enabled, vlan_id=vlan_id,
                address_ip=address_ip, address_prefix=address_prefix
            )

        str_path = '/interfaces/interface[name={:s}]'.format(if_name)
        str_data = yang_if.print_mem('json')
@@ -121,7 +207,6 @@ class InterfaceHandler(_Handler):
            _interface = {
                'name'         : interface_name,
                'type'         : interface_type,
                'mtu'          : interface_state['mtu'],
                'admin-status' : interface_state['admin-status'],
                'oper-status'  : interface_state['oper-status'],
                'management'   : interface_state['management'],
@@ -136,6 +221,9 @@ class InterfaceHandler(_Handler):
                _interface['hardware-port'] = interface_state['hardware-port']
            if 'transceiver' in interface_state:
                _interface['transceiver'] = interface_state['transceiver']
            if 'mtu' in interface_state:
                mtu = interface_state['mtu']
                if mtu > 0: _interface['mtu'] = mtu

            entry_interface_key = '/interface[{:s}]'.format(interface_name)
            entries.append((entry_interface_key, _interface))
@@ -164,6 +252,12 @@ class InterfaceHandler(_Handler):
                    _subinterface['name'] = subinterface_state['name']
                if 'enabled' in subinterface_state:
                    _subinterface['enabled'] = subinterface_state['enabled']
                if 'mtu' in subinterface_state:
                    mtu = subinterface_state['mtu']
                    if mtu > 0:
                        _subinterface['mtu'] = mtu
                        if 'mtu' not in _interface:
                            _interface['mtu'] = mtu

                if 'vlan' in subinterface:
                    vlan = subinterface['vlan']
Loading
Loading