Commit 8621a7b3 authored by Lluis Gifre Renom's avatar Lluis Gifre Renom
Browse files

NBI component - SSE Telemetry:

- Updated to delegate establish/delete subscriptions to SIMAP Connector
- Implemented SSE-based resource
parent 1fab7278
Loading
Loading
Loading
Loading
+0 −169
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 json, logging
from random import choice
from sys import warnoptions
from typing import Dict, List, Optional, Set
from uuid import uuid4
from typing_extensions import TypedDict
from flask import jsonify, request
from flask_restful import Resource
from werkzeug.exceptions import BadRequest, NotFound, UnsupportedMediaType, InternalServerError
from common.proto.monitoring_pb2 import SSEMonitoringSubscriptionConfig
from common.tools.context_queries.Device import get_device
from common.tools.grpc.Tools import grpc_message_to_json_string
from common.proto.monitoring_pb2 import (
    SSEMonitoringSubscriptionConfig,
    SSEMonitoringSubscriptionResponse,
)
from common.tools.rest_conf.client.RestConfClient import RestConfClient
from context.client.ContextClient import ContextClient
from device.client.DeviceClient import DeviceClient
from nbi.service._tools.Authentication import HTTP_AUTH
from nbi.service.database.Engine import Engine
from nbi.service.sse_telemetry.database.Subscription import (
    SSESubsciprionDict,
    list_identifiers,
    set_subscription,
)
from .topology import (
    Controllers,
    SubscribedNotificationsSchema,
    decompose_subscription,
    get_controller_name,
)



class SubscriptionId(TypedDict):
    identifier: str
    uri: str


LOGGER = logging.getLogger(__name__)


class CreateSubscription(Resource):
    # @HTTP_AUTH.login_required
    def post(self):
        db = Engine.get_engine()
        if db is None:
            LOGGER.error('Database engine is not initialized')
            raise InternalServerError('Database engine is not initialized')
        if not request.is_json:
            LOGGER.error('JSON payload is required')
            raise UnsupportedMediaType('JSON payload is required')
        request_data: Optional[SubscribedNotificationsSchema] = request.json
        if request_data is None:
            LOGGER.error('JSON payload is required')
            raise UnsupportedMediaType('JSON payload is required')
        LOGGER.debug('Received subscription request data: {:s}'.format(str(request_data)))

        rest_conf_client = RestConfClient(
            '10.254.0.9', port=8080, scheme='http', username='admin', password='admin',
            logger=logging.getLogger('RestConfClient')
        )

        # break the request into its abstract components for telemetry subscription
        list_db_ids = list_identifiers(db)
        request_identifier = str(
            choice([x for x in range(1000, 10000) if x not in list_db_ids])
        )
        sub_subs = decompose_subscription(rest_conf_client, request_data)

        # subscribe to each component
        device_client = DeviceClient()
        context_client = ContextClient()
        for s in sub_subs:
            xpath_filter = s['ietf-subscribed-notifications:input'][
                'ietf-yang-push:datastore-xpath-filter'
            ]
            xpath_filter_prefix = xpath_filter.split('/ietf-network-topology:link')[0]
            xpath_network = rest_conf_client.get(xpath_filter_prefix)
            if not xpath_network:
                MSG = 'Resource({:s} => {:s}) not found in SIMAP Server'
                raise Exception(MSG.format(str(xpath_filter), str(xpath_filter_prefix)))
            networks = xpath_network.get('ietf-network:network', list())
            if len(networks) != 1:
                MSG = 'Resource({:s} => {:s}) wrong number of entries: {:s}'
                raise Exception(MSG.format(
                    str(xpath_filter), str(xpath_filter_prefix), str(xpath_network)
                ))
            network = networks[0]
            network_id = network['network-id']

            controller_name_map = {
                'e2e'      : 'TFS-E2E',
                'agg'      : 'TFS-AGG',
                'trans-pkt': 'TFS-IP',
                'trans-opt': 'NCE-T',
                'access'   : 'NCE-FAN',
            }
            controller_name = controller_name_map.get(network_id)
            if controller_name is None:
                LOGGER.warning(
                    'Controllerless device detected, skipping subscription for: {:s}'.format(xpath_filter)
                )
                continue

            #SERVICE_ID = ''
            #device_controller = get_controller_name(xpath, SERVICE_ID, context_client)
            #if device_controller == Controllers.CONTROLLERLESS:
            #    LOGGER.warning(
            #        'Controllerless device detected, skipping subscription for: {:s}'.format(xpath)
            #    )
            #    continue

            sampling_interval = s['ietf-subscribed-notifications:input'][
                'ietf-yang-push:periodic'
            ]['ietf-yang-push:period']

            s_req = SSEMonitoringSubscriptionConfig()
            #s_req.device_id.device_uuid.uuid = device_controller.value
            s_req.device_id.device_uuid.uuid = controller_name
            s_req.config_type = SSEMonitoringSubscriptionConfig.Subscribe
            s_req.uri = xpath_filter
            s_req.sampling_interval = str(sampling_interval)
            r: SSEMonitoringSubscriptionResponse = device_client.SSETelemetrySubscribe(s_req)
            s = SSESubsciprionDict(
                uuid=str(uuid4()),
                identifier=r.identifier,
                uri=r.uri,
                xpath=xpath_filter,
                sampling_interval=sampling_interval,
                main_subscription=False,
                main_subscription_id=request_identifier,
            )
            _ = set_subscription(db, s)

        # save the main subscription to the database
        r_uri = f'/restconf/data/subscriptions/{request_identifier}'
        s = SSESubsciprionDict(
            uuid=str(uuid4()),
            identifier=request_identifier,
            uri=r_uri,
            xpath=request_data['ietf-subscribed-notifications:input'][
                'ietf-yang-push:datastore-xpath-filter'
            ],
            sampling_interval=sampling_interval,
            main_subscription=True,
            main_subscription_id=None,
        )
        _ = set_subscription(db, s)

        # Return the subscription ID
        sub_id = SubscriptionId(identifier=request_identifier, uri=r_uri)
        return jsonify(sub_id)
+96 −77
Original line number Diff line number Diff line
@@ -14,28 +14,30 @@


import logging
from typing import Optional
#from typing import Optional
from flask import jsonify, request
from flask_restful import Resource
from werkzeug.exceptions import NotFound, InternalServerError, UnsupportedMediaType
from common.proto.monitoring_pb2 import (
    SSEMonitoringSubscriptionConfig,
    SSEMonitoringSubscriptionResponse,
)
from device.client.DeviceClient import DeviceClient
from context.client.ContextClient import ContextClient
from nbi.service._tools.Authentication import HTTP_AUTH
from nbi.service.database.Engine import Engine
from nbi.service.sse_telemetry.database.Subscription import (
    get_main_subscription,
    get_sub_subscription,
    delete_subscription,
)
from nbi.service.sse_telemetry.topology import (
    Controllers,
    UnsubscribedNotificationsSchema,
    get_controller_name,
)
from werkzeug.exceptions import BadRequest, UnsupportedMediaType #, NotFound, InternalServerError
from common.proto.simap_connector_pb2 import SubscriptionId
from simap_connector.client.SimapConnectorClient import SimapConnectorClient
#from common.proto.monitoring_pb2 import (
#    SSEMonitoringSubscriptionConfig,
#    SSEMonitoringSubscriptionResponse,
#)
#from device.client.DeviceClient import DeviceClient
#from context.client.ContextClient import ContextClient
#from nbi.service._tools.Authentication import HTTP_AUTH
#from nbi.service.database.Engine import Engine
#from nbi.service.sse_telemetry.database.Subscription import (
#    get_main_subscription,
#    get_sub_subscription,
#    delete_subscription,
#)
#from nbi.service.sse_telemetry.topology import (
#    Controllers,
#    UnsubscribedNotificationsSchema,
#    get_controller_name,
#)


LOGGER = logging.getLogger(__name__)
@@ -44,69 +46,86 @@ LOGGER = logging.getLogger(__name__)
class DeleteSubscription(Resource):
    # @HTTP_AUTH.login_required
    def post(self):
        db = Engine.get_engine()
        if db is None:
            LOGGER.error('Database engine is not initialized')
            raise InternalServerError('Database engine is not initialized')
#        db = Engine.get_engine()
#        if db is None:
#            LOGGER.error('Database engine is not initialized')
#            raise InternalServerError('Database engine is not initialized')

        if not request.is_json:
            LOGGER.error('JSON payload is required')
            raise UnsupportedMediaType('JSON payload is required')
        request_data: Optional[UnsubscribedNotificationsSchema] = request.json
        if request_data is None:
            LOGGER.error('JSON payload is required')
#            LOGGER.error('JSON payload is required')
            raise UnsupportedMediaType('JSON payload is required')
        main_subscription_id = request_data['delete-subscription']['identifier']
        LOGGER.debug(
            'Received delete subscription request for ID: {:s}'.format(main_subscription_id)
        )

        # Get the main subscription
        main_subscription = get_main_subscription(db, main_subscription_id)
        if main_subscription is None:
            LOGGER.error('Subscription not found: {:s}'.format(main_subscription_id))
            raise NotFound('Subscription not found')

        # Get all sub-subscriptions associated with this main subscription
        sub_subscriptions = get_sub_subscription(db, main_subscription_id)

        device_client = DeviceClient()
        context_client = ContextClient()

        # Unsubscribe from each sub-subscription
        for sub_sub in sub_subscriptions:
            # Create unsubscribe request
            SERVICE_ID = ''
            device_controller = get_controller_name(sub_sub['xpath'], SERVICE_ID, context_client)
            if device_controller == Controllers.CONTROLLERLESS:
                LOGGER.warning(
                    'Controllerless device detected, skipping subscription for: {:s}'.format(
                        sub_sub['xpath']
                    )
                )
                continue
            unsub_req = SSEMonitoringSubscriptionConfig()
            unsub_req.device_id.device_uuid.uuid = device_controller.value
            unsub_req.config_type = SSEMonitoringSubscriptionConfig.Unsubscribe
            unsub_req.uri = sub_sub['xpath']
            unsub_req.identifier = sub_sub['identifier']
        request_data = request.json
        LOGGER.debug('[post] Unsubscription request: {:s}'.format(str(request_data)))
#        if request_data is None:
#            LOGGER.error('JSON payload is required')
#            raise UnsupportedMediaType('JSON payload is required')

            # Send unsubscribe request to device
            device_client.SSETelemetrySubscribe(unsub_req)
        if 'ietf-subscribed-notifications:input' not in request_data:
            raise BadRequest('Missing field(ietf-subscribed-notifications:input)')
        input_data = request_data['ietf-subscribed-notifications:input']

            delete_subscription(db, sub_sub['identifier'], False)
        subscription_id = SubscriptionId()

            LOGGER.info('Unsubscribed from {:s} successfully'.format(sub_sub.get('uri', '')))
        if 'id' not in input_data:
            raise BadRequest('Missing field(ietf-subscribed-notifications:input/id)')
        subscription_id.subscription_id = input_data['id']

        # Delete the main subscription from database
        delete_subscription(db, main_subscription_id, True)

        LOGGER.info('Successfully deleted main subscription: {:s}'.format(main_subscription_id))

        #if SERVICE_ID == 'simap1':
        #    SERVICE_ID = 'simap2'
        #elif SERVICE_ID == 'simap2':
        #    SERVICE_ID = 'simap1'
        #else:
        #    LOGGER.warning('Unknown service ID, not switching: {:s}'.format(SERVICE_ID))
        simap_connector_client = SimapConnectorClient()
        simap_connector_client.DeleteSubscription(subscription_id)

#        main_subscription_id = request_data['delete-subscription']['identifier']
#        LOGGER.debug(
#            'Received delete subscription request for ID: {:s}'.format(main_subscription_id)
#        )
#
#        # Get the main subscription
#        main_subscription = get_main_subscription(db, main_subscription_id)
#        if main_subscription is None:
#            LOGGER.error('Subscription not found: {:s}'.format(main_subscription_id))
#            raise NotFound('Subscription not found')
#
#        # Get all sub-subscriptions associated with this main subscription
#        sub_subscriptions = get_sub_subscription(db, main_subscription_id)
#
#        device_client = DeviceClient()
#        context_client = ContextClient()
#
#        # Unsubscribe from each sub-subscription
#        for sub_sub in sub_subscriptions:
#            # Create unsubscribe request
#            SERVICE_ID = ''
#            device_controller = get_controller_name(sub_sub['xpath'], SERVICE_ID, context_client)
#            if device_controller == Controllers.CONTROLLERLESS:
#                LOGGER.warning(
#                    'Controllerless device detected, skipping subscription for: {:s}'.format(
#                        sub_sub['xpath']
#                    )
#                )
#                continue
#            unsub_req = SSEMonitoringSubscriptionConfig()
#            unsub_req.device_id.device_uuid.uuid = device_controller.value
#            unsub_req.config_type = SSEMonitoringSubscriptionConfig.Unsubscribe
#            unsub_req.uri = sub_sub['xpath']
#            unsub_req.identifier = sub_sub['identifier']
#
#            # Send unsubscribe request to device
#            device_client.SSETelemetrySubscribe(unsub_req)
#
#            delete_subscription(db, sub_sub['identifier'], False)
#
#            LOGGER.info('Unsubscribed from {:s} successfully'.format(sub_sub.get('uri', '')))
#
#        # Delete the main subscription from database
#        delete_subscription(db, main_subscription_id, True)
#
#        LOGGER.info('Successfully deleted main subscription: {:s}'.format(main_subscription_id))
#
#        #if SERVICE_ID == 'simap1':
#        #    SERVICE_ID = 'simap2'
#        #elif SERVICE_ID == 'simap2':
#        #    SERVICE_ID = 'simap1'
#        #else:
#        #    LOGGER.warning('Unknown service ID, not switching: {:s}'.format(SERVICE_ID))
#
        return jsonify({})
+198 −0
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 logging #, json
#from random import choice
#from typing import Dict, List, Optional, Set
#from uuid import uuid4
#from typing_extensions import TypedDict
from flask import jsonify, request, url_for
from flask_restful import Resource
from werkzeug.exceptions import BadRequest, UnsupportedMediaType #, NotFound, InternalServerError
from common.proto.simap_connector_pb2 import Subscription #, SubscriptionId
from simap_connector.client.SimapConnectorClient import SimapConnectorClient
#from common.proto.monitoring_pb2 import SSEMonitoringSubscriptionConfig
#from common.tools.context_queries.Device import get_device
#from common.tools.grpc.Tools import grpc_message_to_json_string
#from common.proto.monitoring_pb2 import (
#    SSEMonitoringSubscriptionConfig,
#    SSEMonitoringSubscriptionResponse,
#)
#from common.tools.rest_conf.client.RestConfClient import RestConfClient
#from context.client.ContextClient import ContextClient
#from device.client.DeviceClient import DeviceClient
#from nbi.service._tools.Authentication import HTTP_AUTH
#from nbi.service.database.Engine import Engine
#from nbi.service.sse_telemetry.database.Subscription import (
#    SSESubsciprionDict,
#    list_identifiers,
#    set_subscription,
#)
#from .topology import (
#    Controllers,
#    SubscribedNotificationsSchema,
#    decompose_subscription,
#    get_controller_name,
#)



#class SubscriptionId(TypedDict):
#    identifier: str
#    uri: str


LOGGER = logging.getLogger(__name__)


class CreateSubscription(Resource):
    # @HTTP_AUTH.login_required
    def post(self):
        if not request.is_json:
            raise UnsupportedMediaType('JSON payload is required')

        request_data = request.json
        LOGGER.debug('[post] Subscription request: {:s}'.format(str(request_data)))

        if 'ietf-subscribed-notifications:input' not in request_data:
            raise BadRequest('Missing field(ietf-subscribed-notifications:input)')
        input_data = request_data['ietf-subscribed-notifications:input']

        subscription = Subscription()

        if 'datastore' not in input_data:
            raise BadRequest('Missing field(ietf-subscribed-notifications:input/datastore)')
        subscription.datastore = input_data['datastore']

        if 'ietf-yang-push:datastore-xpath-filter' not in input_data:
            raise BadRequest('Missing field(ietf-subscribed-notifications:input/ietf-yang-push:datastore-xpath-filter)')
        subscription.xpath_filter = input_data['ietf-yang-push:datastore-xpath-filter']

        if 'ietf-yang-push:periodic' not in input_data:
            raise BadRequest('Missing field(ietf-subscribed-notifications:input/ietf-yang-push:periodic)')
        periodic = input_data['ietf-yang-push:periodic']

        if 'ietf-yang-push:period' not in periodic:
            raise BadRequest('Missing field(ietf-subscribed-notifications:input/ietf-yang-push:periodic/ietf-yang-push:period)')
        subscription.period = float(periodic['ietf-yang-push:period'])

        simap_connector_client = SimapConnectorClient()
        subscription_id = simap_connector_client.EstablishSubscription(subscription)
        subscription_id = subscription_id.subscription_id

        subscription_uri = url_for('sse.stream', subscription_id=subscription_id)
        sub_id = {'identifier': subscription_id, 'uri': subscription_uri}
        return jsonify(sub_id)


#        db = Engine.get_engine()
#        if db is None:
#            LOGGER.error('Database engine is not initialized')
#            raise InternalServerError('Database engine is not initialized')
#        rest_conf_client = RestConfClient(
#            '10.254.0.9', port=8080, scheme='http', username='admin', password='admin',
#            logger=logging.getLogger('RestConfClient')
#        )
#
#        # break the request into its abstract components for telemetry subscription
#        list_db_ids = list_identifiers(db)
#        request_identifier = str(
#            choice([x for x in range(1000, 10000) if x not in list_db_ids])
#        )
#        sub_subs = decompose_subscription(rest_conf_client, request_data)
#
#        # subscribe to each component
#        device_client = DeviceClient()
#        context_client = ContextClient()
#        for s in sub_subs:
#            xpath_filter = s['ietf-subscribed-notifications:input'][
#                'ietf-yang-push:datastore-xpath-filter'
#            ]
#            xpath_filter_prefix = xpath_filter.split('/ietf-network-topology:link')[0]
#            xpath_network = rest_conf_client.get(xpath_filter_prefix)
#            if not xpath_network:
#                MSG = 'Resource({:s} => {:s}) not found in SIMAP Server'
#                raise Exception(MSG.format(str(xpath_filter), str(xpath_filter_prefix)))
#            networks = xpath_network.get('ietf-network:network', list())
#            if len(networks) != 1:
#                MSG = 'Resource({:s} => {:s}) wrong number of entries: {:s}'
#                raise Exception(MSG.format(
#                    str(xpath_filter), str(xpath_filter_prefix), str(xpath_network)
#                ))
#            network = networks[0]
#            network_id = network['network-id']
#
#            controller_name_map = {
#                'e2e'      : 'TFS-E2E',
#                'agg'      : 'TFS-AGG',
#                'trans-pkt': 'TFS-IP',
#                'trans-opt': 'NCE-T',
#                'access'   : 'NCE-FAN',
#            }
#            controller_name = controller_name_map.get(network_id)
#            if controller_name is None:
#                LOGGER.warning(
#                    'Controllerless device detected, skipping subscription for: {:s}'.format(xpath_filter)
#                )
#                continue
#
#            #SERVICE_ID = ''
#            #device_controller = get_controller_name(xpath, SERVICE_ID, context_client)
#            #if device_controller == Controllers.CONTROLLERLESS:
#            #    LOGGER.warning(
#            #        'Controllerless device detected, skipping subscription for: {:s}'.format(xpath)
#            #    )
#            #    continue
#
#            sampling_interval = s['ietf-subscribed-notifications:input'][
#                'ietf-yang-push:periodic'
#            ]['ietf-yang-push:period']
#
#            s_req = SSEMonitoringSubscriptionConfig()
#            #s_req.device_id.device_uuid.uuid = device_controller.value
#            s_req.device_id.device_uuid.uuid = controller_name
#            s_req.config_type = SSEMonitoringSubscriptionConfig.Subscribe
#            s_req.uri = xpath_filter
#            s_req.sampling_interval = str(sampling_interval)
#            r: SSEMonitoringSubscriptionResponse = device_client.SSETelemetrySubscribe(s_req)
#            s = SSESubsciprionDict(
#                uuid=str(uuid4()),
#                identifier=r.identifier,
#                uri=r.uri,
#                xpath=xpath_filter,
#                sampling_interval=sampling_interval,
#                main_subscription=False,
#                main_subscription_id=request_identifier,
#            )
#            _ = set_subscription(db, s)
#
#        # save the main subscription to the database
#        r_uri = f'/restconf/data/subscriptions/{request_identifier}'
#        s = SSESubsciprionDict(
#            uuid=str(uuid4()),
#            identifier=request_identifier,
#            uri=r_uri,
#            xpath=request_data['ietf-subscribed-notifications:input'][
#                'ietf-yang-push:datastore-xpath-filter'
#            ],
#            sampling_interval=sampling_interval,
#            main_subscription=True,
#            main_subscription_id=None,
#        )
#        _ = set_subscription(db, s)

#        # Return the subscription ID
#        sub_id = SubscriptionId(identifier=request_identifier, uri=r_uri)
#        return jsonify(sub_id)
+154 −0

File added.

Preview size limit exceeded, changes collapsed.

+14 −4
Original line number Diff line number Diff line
@@ -12,18 +12,23 @@
# See the License for the specific language governing permissions and
# limitations under the License.

# RFC 8299 - YANG Data Model for L3VPN Service Delivery
# Ref: https://datatracker.ietf.org/doc/rfc8299

# RFC 8639 - Subscription to YANG Notifications
# Ref: https://datatracker.ietf.org/doc/html/rfc8639

# RFC 8641 - Subscription to YANG Notifications for Datastore Updates
# Ref: https://datatracker.ietf.org/doc/html/rfc8641


from nbi.service.NbiApplication import NbiApplication
from .CreateSubscription import CreateSubscription
from .EstablishSubscription import EstablishSubscription
from .DeleteSubscription import DeleteSubscription
from .StreamSubscription import StreamSubscription


def register_telemetry_subscription(nbi_app: NbiApplication):
    nbi_app.add_rest_api_resource(
        CreateSubscription,
        EstablishSubscription,
        '/restconf/operations/subscriptions:establish-subscription',
        '/restconf/operations/subscriptions:establish-subscription/',
    )
@@ -32,3 +37,8 @@ def register_telemetry_subscription(nbi_app: NbiApplication):
        '/restconf/operations/subscriptions:delete-subscription',
        '/restconf/operations/subscriptions:delete-subscription/',
    )
    nbi_app.add_rest_api_resource(
        StreamSubscription,
        '/restconf/stream/<int:subscription_id>',
        '/restconf/stream/<int:subscription_id>/',
    )