Commit ba25e87e authored by Javier Velázquez's avatar Javier Velázquez
Browse files

Merge branch 'feat/34-tid-new-monitoring-component' into 'develop'

Resolve "(TID) New Monitoring Component"

See merge request !34
parents 25bdf49e c076157e
Loading
Loading
Loading
Loading
Loading
+4 −1
Original line number Diff line number Diff line
@@ -27,6 +27,7 @@ from src.webui.gui import gui_bp
from src.config.config import create_config
from src.database.db import init_db as init_slice
from src.database.service_db import init_db as init_service
from src.database.telemetry_client_db import init_db as init_telemetry

# Paths that do not require authentication (Swagger UI and its static assets)
# /nsc       → Swagger UI
@@ -58,8 +59,10 @@ def create_app():
    """Create Flask application with configured API and namespaces."""
    init_slice()
    init_service()
    init_telemetry()
    app = Flask(__name__)
    app = create_config(app)
    # app.logger.setLevel(logging.INFO)
    CORS(app)

    # Configure logging to provide clear and informative log messages
@@ -104,4 +107,4 @@ def create_app():

if __name__ == "__main__":
    app = create_app()
    app.run(host="0.0.0.0", port=NSC_PORT, debug=True)
 No newline at end of file
    app.run(host="0.0.0.0", port=NSC_PORT, debug=True, threaded=True)
 No newline at end of file
+2 −2
Original line number Diff line number Diff line
@@ -14,7 +14,7 @@

# This file is an original contribution from Telefonica Innovación Digital S.L.

Flask
flask[async]
flask-cors
flask-restx
netmiko
@@ -25,4 +25,4 @@ coverage
pytest
libyang==3.3.0
sysrepo==1.7.6
aiohttp
 No newline at end of file
+266 −3
Original line number Diff line number Diff line
@@ -15,15 +15,20 @@
# This file is an original contribution from Telefonica Innovación Digital S.L.

from src.utils.send_response import send_response
import logging
import logging, json, time, requests, traceback, asyncio, aiohttp
from flask import current_app
from src.database.db import get_data, delete_data, get_all_data, delete_all_data
from src.database.service_db import delete_by_slice_id, get_data_by_slice_id
from src.database.telemetry_client_db import create_client, get_client, get_all_clients, delete_client, delete_all_clients, upsert_subscription, get_subscription, get_client_subscriptions, delete_subscription, delete_all_subscriptions
from src.realizer.tfs.helpers.tfs_connector import tfs_connector
from src.utils.safe_get import safe_get
from src.database.sysrepo_store import get_data_store, create_data_store, delete_data_store, update_data_store, normalize_libyang_data
from typing import Dict, Tuple
from src.realizer.restconf.connectors.tfs_connector import tfs_connector as tfs_restconf_connector
from src.planner.shortest_path import get_shortest_path

class Api:

    def __init__(self, slice_service):
        self.slice_service = slice_service

@@ -745,3 +750,261 @@ class Api:
            return send_response(False, code=404, message=str(e))
        except Exception as e:
            return send_response(False, code=500, message=str(e))

    # --- CLIENTS ---

    def get_clients(self, client_id=None):
        try:
            if client_id:
                return get_client(client_id), 200
            clients = get_all_clients()
            if not clients:
                raise ValueError("No clients found")
            return clients, 200

        except ValueError as e:
            return send_response(False, code=404, message=str(e))
        except Exception as e:
            return send_response(False, code=500, message=str(e))
        
    def add_client(self, client_id):
        try:
            create_client(client_id)
            logging.info(f"Client '{client_id}' created successfully")
            return send_response(
                True,
                code=201,
                message=f"Client '{client_id}' created successfully",
                data={"client_id": client_id}
            )

        except ValueError as e:
            return send_response(False, code=409, message=str(e))
        except Exception as e:
            return send_response(False, code=500, message=str(e))
    
    def delete_clients(self, client_id=None):
        try:
            if client_id:
                delete_client(client_id)
                logging.info(f"Client '{client_id}' removed successfully")
            else:
                delete_all_clients()
                logging.info("All clients removed successfully")
            return {}, 204
        except ValueError as e:
            return send_response(False, code=404, message=str(e))
        except Exception as e:
            return send_response(False, code=500, message=str(e))

    # --- SUBSCRIPTIONS ---

    def get_subscriptions(self, client_id, slice_id = None):
        try:
            if slice_id is not None:
                try:
                    subscription = get_subscription(client_id, slice_id)
                except ValueError:
                    subscription = None
                if not subscription:
                    raise ValueError(f"Client '{client_id}' has no subscription for slice '{slice_id}'")
                
                telemetry = self.get_telemetry(slice_id)

                return {
                    **subscription,
                    "telemetry": telemetry
                }, 200

            subscriptions = get_client_subscriptions(client_id)

            result = []

            for sub in subscriptions:
                slice_id = sub["slice_id"]

                telemetry = self.get_telemetry(slice_id)

                result.append({
                    **sub,
                    "telemetry": telemetry
                })

            return {
                "client_id": client_id,
                "subscriptions": result
            }, 200
            
        except ValueError as e:
            return send_response(False, code=404, message=str(e))
        except Exception as e:
            return send_response(False, code=500, message=str(e))

    def add_subscription(self, client_id, slice_id, frequency):
        try:
            try:
                subscription = get_subscription(client_id, slice_id)
            except ValueError:
                subscription = None
            if subscription:
                raise ValueError(f"Client '{client_id}' already has a subscription for slice '{slice_id}'")
            if not frequency:
                raise KeyError("Field 'frequency' is required")

            upsert_subscription(client_id, slice_id, frequency)
            logging.info(f"Subscription for slice '{slice_id}' and client '{client_id}' created successfully")
            return send_response(
                True,
                code=201,
                message="Subscription successfully created",
                data={
                    "sliceId": slice_id,
                    "frequency": frequency
                }
            )
        except KeyError as e:
            return send_response(False, code=400, message=str(e))
        except ValueError as e:
            return send_response(False, code=404, message=str(e))
        except Exception as e:
            return send_response(False, code=500, message=str(e))

    def update_subscription(self, client_id, slice_id, frequency):
        try:
            try:
                subscription = get_subscription(client_id, slice_id)
            except ValueError:
                subscription = None
            if not subscription:
                raise ValueError(f"Client '{client_id}' has no subscription for slice '{slice_id}'")
            if not frequency:
                raise KeyError("Field 'frequency' is required")

            upsert_subscription(client_id, slice_id, frequency)
            logging.info(f"Subscription for slice '{slice_id}' and client '{client_id}' modified successfully")
            return send_response(
                True,
                code=201,
                message="Subscription successfully modified",
                data={
                    "sliceId": slice_id,
                    "frequency": frequency
                }
            )
        except KeyError as e:
            return send_response(False, code=400, message=str(e))
        except ValueError as e:
            return send_response(False, code=404, message=str(e))
        except Exception as e:
            return send_response(False, code=500, message=str(e))

    def delete_subscriptions(self, client_id, slice_id = None):
        try:
            subscriptions = get_client_subscriptions(client_id)
            if slice_id:
                if slice_id not in subscriptions: 
                    raise ValueError(f"Client '{client_id}' has no subscription for slice '{slice_id}'")
                delete_subscription(client_id, slice_id)
                logging.info(f"Subscription for slice '{slice_id}' and client '{client_id}' removed successfully")
                return {}, 204
            delete_all_subscriptions(client_id)
            logging.info(f"All subscriptions for client '{client_id}' removed successfully")
            return {}, 204
        
        except ValueError as e:
            return send_response(False, code=404, message=str(e))
        except Exception as e:
            return send_response(False, code=500, message=str(e))

    # --- TELEMETRY ---

    def get_telemetry(self, slice_id = None):
        logging.debug(f"Getting telemetry for slice_id: {slice_id}")
        try:         
            if slice_id is not None:
                slice_sdps = get_data_store(f"/ietf-network-slice-service:network-slice-services/slice-service[id='{slice_id}']/sdps")
                logging.debug(f"SDPs for slice_id '{slice_id}': {slice_sdps}")
                if not slice_sdps:
                    raise ValueError("No SDPs found")
                if len(slice_sdps) > 2:
                    raise Exception(f"Monitoring for more than 2 SDPs is not supported. Found {len(slice_sdps)} SDPs.")
                xpath = f"/ietf-network-slice-service:network-slice-services/slice-service[id='{slice_id}']"
                existing_slice = get_data_store(xpath)
                if not existing_slice:
                    raise ValueError(f"There is no slice with id '{slice_id}' registered")
                template_id = safe_get(existing_slice, ["network-slice-services", "slice-service", slice_id, "slo-sle-template"])
                slo_sle_template = self.get_slo_sle_templates(template_id)[0]
                slo_sle_template = safe_get(slo_sle_template, ["network-slice-services", "slo-sle-templates", "slo-sle-template"])
                slo_sle_template = next(iter(slo_sle_template), None)
                if not slo_sle_template:
                    raise ValueError(f"SLO/SLE template '{template_id}' not found for slice '{slice_id}'")
                metrics = self.slice_service.monitoring(slice_id, slo_sle_template, slice_sdps)
                return metrics, 200
            
            telemetry_data = {}
            slices_data = self.get_slice_services()[0] 
            slice_service_list = slices_data["network-slice-services"]["slice-service"]
            if isinstance(slice_service_list, dict):
                slice_service_list = list(slice_service_list.values())
                
            for slice in slice_service_list:
                selected_template_id = slice.get("slo-sle-template")
                slo_sle_template = self.get_slo_sle_templates(selected_template_id)[0]
                slo_sle_template = safe_get(slo_sle_template, ["network-slice-services", "slo-sle-templates", "slo-sle-template"])
                slo_sle_template = next(iter(slo_sle_template), None)
                if not slo_sle_template:
                    raise ValueError(f"SLO/SLE template '{selected_template_id}' not found for slice '{slice['id']}'")
                slice_id = slice["id"] 
                slice_sdps = get_data_store(f"/ietf-network-slice-service:network-slice-services/slice-service[id='{slice_id}']/sdps") 
                if not slice_sdps:
                    raise ValueError("No SDPs found")
                if len(slice_sdps) > 2:
                    raise Exception(f"Monitoring for more than 2 SDPs is not supported. Found {len(slice_sdps)} SDPs.")
                telemetry_data[slice_id] = self.slice_service.monitoring(slice_id, slo_sle_template, slice_sdps)
            return telemetry_data, 200

        except ValueError as e:
            return send_response(False, code=404, message=str(e))
        except Exception as e:
            return send_response(False, code=500, message=str(e))

    def sync_stream(self, async_gen_func, *args, **kwargs):
        loop = asyncio.new_event_loop()
        asyncio.set_event_loop(loop)
        agen = async_gen_func(*args, **kwargs)

        try:
            while True: yield loop.run_until_complete(agen.__anext__())
        except StopAsyncIteration: pass
        finally: loop.close()

    async def stream_client_subscriptions(self, client_id):
        while True:
            try:
                data, code = self.get_subscriptions(client_id)
                if code == 200:
                    yield f"data: {json.dumps(data)}\n\n"
                    subs = data.get("subscriptions", [])
                    freq = max([s["frequency"] for s in subs]) if subs else 5
                    await asyncio.sleep(freq)
                else:
                    yield f"event: error\ndata: {json.dumps(data)}\n\n"
                    break
            except Exception as e:
                yield f"event: error\ndata: {json.dumps({'error': str(e)})}\n\n"
                break

    async def stream_slice_subscription(self, client_id, slice_id):
        while True:
            try:
                data, code = self.get_subscriptions(client_id, slice_id)
                if code == 200:
                    yield f"data: {json.dumps(data)}\n\n"
                    freq = data.get("frequency", 5)
                    await asyncio.sleep(freq)
                else:
                    yield f"event: error\ndata: {json.dumps(data)}\n\n"
                    break
            except Exception as e:
                yield f"event: error\ndata: {json.dumps({'error': str(e)})}\n\n"
                break
 No newline at end of file
+7 −2
Original line number Diff line number Diff line
@@ -42,7 +42,7 @@ E2E_OPTICAL_IP=127.0.0.1
# Realizer
# -------------------------
# If true, no config sent to controllers
DUMMY_MODE=false
DUMMY_MODE=true

# -------------------------
# Teraflow
@@ -66,12 +66,17 @@ TFS_E2E_IP=127.0.0.1
# -------------------------
# Restconf Controller
# -------------------------
RESTCONF_IP=127.0.0.1
RESTCONF_IP=192.168.27.189
# Options: TFS or IXIA
SDN_CONTROLLER_TYPE=TFS
# Options: FRR, CISCO
DATAPLANE_SUPPORT=CISCO

# -------------------------
# Monitoring
# -------------------------
SDN_SUBSCRIPTION_PERIOD=10

# -------------------------
# WebUI
# -------------------------
+5 −0
Original line number Diff line number Diff line
@@ -69,8 +69,13 @@ def create_config(app: Flask):
    app.config["SDN_CONTROLLER_TYPE"] = os.getenv("SDN_CONTROLLER_TYPE", "TFS")
    app.config["DATAPLANE_SUPPORT"] = os.getenv("DATAPLANE_SUPPORT", "FRR")

    # Monitoring
    app.config["SDN_SUBSCRIPTION_PERIOD"] = int(os.getenv("SDN_SUBSCRIPTION_PERIOD", "10"))
    app.config["TELEMETRY_CACHE"] = {}

    # API
    app.config["API_USERNAME"] = os.getenv("API_USERNAME", "admin")
    app.config["API_PASSWORD"] = os.getenv("API_PASSWORD", "admin")

    return app
Loading