Commit dc0a6054 authored by Paris Stentoumis's avatar Paris Stentoumis
Browse files

feat: categorize an application profile based on requested resources. enhance...

feat: categorize an application profile based on requested resources. enhance application deployment to support nodeSelector, resource requests/limits and correct nodeport exposure
parent ebd25f26
Loading
Loading
Loading
Loading
+60 −5
Original line number Diff line number Diff line
@@ -69,6 +69,28 @@ class EdgeApplicationManager(EdgeCloudManagementInterface):
            request=None
        )

    
    def _calculate_tier(self, cpuRequest, memoryRequest, cpuLimit, memoryLimit):
        cpu_request = int(cpuRequest)  # max value = 1000
        memory_request = int(memoryRequest)
        cpu_limit = int(cpuLimit)
        memory_limit = int(memoryLimit)

        cpu_ratio = cpu_request / cpu_limit
        memory_ratio = memory_request / memory_limit

        cpu_score = min(cpu_request / 1000, 1.0)
        memory_score = min(memory_request / 1024, 1.0)

        score = (
            cpu_score * 0.4 + memory_score * 0.4 + cpu_ratio * 0.1 + memory_ratio * 0.1
        )
        if score >= 0.6:
            return "gold"
        else:
            return "silver"


    def _intent_to_dict(self, intent):
        ICM = Namespace("http://tio.models.tmforum.org/tio/v3.6.0/IntentCommonModel/")
        LOG = Namespace("http://tio.models.tmforum.org/tio/v3.6.0/LogicalOperators/")
@@ -96,6 +118,13 @@ class EdgeApplicationManager(EdgeCloudManagementInterface):
        ttl_string = intent["expression"]["expressionValue"]
        decoded_intent = self._intent_to_dict(intent=ttl_string)

        tier = self._calculate_tier(
            decoded_intent["cpuRequest"][:-1],
            decoded_intent["memoryRequest"][:-2],
            decoded_intent["cpuLimit"][:-1],
            decoded_intent["memoryLimit"][:-2],
        )

        body = {
            "appId": decoded_intent["appId"],
            "name": decoded_intent["app"],
@@ -113,7 +142,14 @@ class EdgeApplicationManager(EdgeCloudManagementInterface):
                    }
                }
            },
            "version": ""
            "version": "",
            "extra_fields": {
                "replicaCount": decoded_intent["replicaCount"],
                "cpuLimit": decoded_intent["cpuLimit"][:-1],
                "memoryLimit": decoded_intent["memoryLimit"][:-2],
                "nodePort": decoded_intent["nodePort"],
                "tier": tier,
            },
        }

        onboarding = self.onboard_app(body)
@@ -152,6 +188,7 @@ class EdgeApplicationManager(EdgeCloudManagementInterface):
                request=None,
            )
        
        
    def get_intents(self, fields, offset, limit):
        try:
            return self.connector_db.get_intents(fields, offset, limit)
@@ -206,6 +243,10 @@ class EdgeApplicationManager(EdgeCloudManagementInterface):
        req_resources = app_manifest.get("requiredResources")
        version = app_manifest.get("version")
        app_provider = app_manifest.get("appProvider")
        extra_fields = app_manifest.get("extra_fields")
        tier = app_manifest.get("tier")
        if tier is None and extra_fields is not None:
            tier = extra_fields.get("tier")
        for ni in network_interfaces:
            ports.append(ni.get("port"))
        insert_doc = ServiceFunctionRegistrationRequest(
@@ -217,6 +258,8 @@ class EdgeApplicationManager(EdgeCloudManagementInterface):
            required_resources=req_resources,
            app_provider=app_provider,
            version=version,
            tier=tier,
            extra_fields=extra_fields,
        )
        result = self.connector_db.insert_document_service_function(insert_doc.to_dict())
        if type(result) is str:
@@ -312,6 +355,10 @@ class EdgeApplicationManager(EdgeCloudManagementInterface):
        app = self.connector_db.get_documents_from_collection(
            "service_functions", input_type="_id", input_value=app_id
        )
    
        node_selector = {}
        if app[0].get("tier") != '':
            node_selector = {"node-tier": app[0].get("tier")}
        # success_response = []
        result = None
        response = None
@@ -321,13 +368,15 @@ class EdgeApplicationManager(EdgeCloudManagementInterface):
            sf = DeployServiceFunction(
                service_function_name=app[0].get("name"),
                service_function_instance_name=app[0].get("name"),
                node_ports=app[0].get("application_ports"),#fill this
                # service_function_instance_name=body.get("name"),
                # location=body.get('edgeCloudZoneId'),
            )
            result = deploy_service_function(
            (result, service_result) = deploy_service_function(
                service_function=sf,
                connector_db=self.connector_db,
                kubernetes_connector=self.k8s_connector,
                node_selector=node_selector if "tier" in node_selector else None
            )

        if type(result) is V1Deployment:
@@ -338,11 +387,17 @@ class EdgeApplicationManager(EdgeCloudManagementInterface):
            response["appInstanceId"] = result.metadata.uid
            response["appProvider"] = app[0].get("app_provider")
            response["status"] = "unknown"
            # interfaces = []
            if type(service_result) is V1Service:
                for port in service_result.spec.ports:
                    endpoint_info_list = []
                    endpoint_info = {}
                    access_point = {"port": port.node_port}
                    endpoint_info["accessPoints"] = access_point
                    endpoint_info_list.append(endpoint_info)
            # for port in deployment.get("ports"):
            #     access_point = {"port": port}
            #     
            #     interfaces.append({"interfaceId": "", "accessPoints": access_point})
            # response["componentEndpointInfo"] = interfaces
            response["componentEndpointInfo"] = endpoint_info_list
            response["kubernetesClusterRef"] = ""
            response["edgeCloudZoneId"] = app_zones[0].get("EdgeCloudZone").get("edgeCloudZoneId")

+30 −2
Original line number Diff line number Diff line
@@ -151,6 +151,7 @@ def deploy_service_function(
    service_function: DeployServiceFunction,
    connector_db: ConnectorDB,
    kubernetes_connector: KubernetesConnector,
    node_selector=None,
    paas_name=None,
):

@@ -177,10 +178,37 @@ def deploy_service_function(
    if service_function.location is not None:
        final_deploy_descriptor["location"] = service_function.location

    if "required_resources" in ser_function_[0]:
        final_deploy_descriptor["required_resources"] = ser_function_[0]["required_resources"]
    if "extra_fields" in ser_function_[0]:
        final_deploy_descriptor["extra_fields"] = ser_function_[0]["extra_fields"]
        replica_count = ser_function_[0]["extra_fields"].get("replicaCount")
        if replica_count is not None:
            final_deploy_descriptor["replica_count"] = replica_count

    containers = prepare_container(service_function, ser_function_)
    if isinstance(containers, tuple):
        return containers
    final_deploy_descriptor["containers"] = containers
    service_node_ports = service_function.service_node_ports
    if service_node_ports is None:
        extra_fields = ser_function_[0].get("extra_fields") or {}
        node_port = extra_fields.get("nodePort")
        if node_port is not None:
            service_node_ports = node_port if isinstance(node_port, list) else [node_port]
    if service_node_ports is not None:
        exposed_ports = containers[0].get("exposed_ports")
        if exposed_ports is None:
            return "Please expose ports before setting service_node_ports", 400
        if len(exposed_ports) != len(service_node_ports):
            return (
                "The number of service_node_ports must match the number of exposed ports",
                400,
            )
        final_deploy_descriptor["service_node_ports"] = service_node_ports

    if node_selector is not None:
        final_deploy_descriptor["nodeSelector"] = node_selector

    vol_result = prepare_volumes(service_function, ser_function_, final_deploy_descriptor)
    if vol_result is not None:
@@ -190,7 +218,7 @@ def deploy_service_function(
    if env_result is not None:
        return env_result

    response = kubernetes_connector.deploy_service_function(final_deploy_descriptor)
    (response, service_response) = kubernetes_connector.deploy_service_function(final_deploy_descriptor)
    deployed_service_function_db = {}
    deployed_service_function_db["service_function_name"] = ser_function_[0]["name"]
    if service_function.location is not None:
@@ -209,4 +237,4 @@ def deploy_service_function(
        connector_db.insert_document_deployed_service_function(
            document=deployed_service_function_db
        )
    return response
    return (response, service_response)
+4 −0
Original line number Diff line number Diff line
@@ -146,6 +146,10 @@ class ConnectorDB:
        insert_doc["app_provider"] = document["app_provider"]
        insert_doc["required_resources"] = document["required_resources"]
        insert_doc["version"] = document.get("version")
        if "tier" in document:
            insert_doc["tier"] = document.get("tier")
        if document.get("extra_fields") is not None:
            insert_doc["extra_fields"] = document.get("extra_fields")
        if document.get("application_ports") is not None:
            insert_doc["application_ports"] = document.get("application_ports")
        if document.get("autoscaling_policies") is not None:
+53 −18
Original line number Diff line number Diff line
@@ -6,8 +6,8 @@ import requests
import urllib3
from kubernetes import client
from kubernetes.client.rest import ApiException
from sunrise6g_opensdk.edgecloud.adapters.kubernetes.lib.models.node import Node

from sunrise6g_opensdk.edgecloud.adapters.kubernetes.lib.models.node import Node
from sunrise6g_opensdk.edgecloud.adapters.kubernetes.lib.utils import (
    auxiliary_functions,
)
@@ -321,7 +321,7 @@ class KubernetesConnector:
                self.namespace, body_deployment
            )
            # api_response_service = api_instance_apiregv1.create_api_service(body_service)
            self.v1.create_namespaced_service(self.namespace, body_service)
            api_response_service = self.v1.create_namespaced_service(self.namespace, body_service)
            if "autoscaling_policies" in descriptor_service_function:
                # V1 AUTOSCALER
                body_hpa = self.create_hpa(descriptor_service_function)
@@ -332,7 +332,7 @@ class KubernetesConnector:
                # body_hpa = create_hpa(descriptor_paas)
                # api_instance_v2beta1autoscale.create_namespaced_horizontal_pod_autoscaler("sunrise6g",body_hpa)

            return api_response_deployment
            return (api_response_deployment, api_response_service)
        except ApiException as e:
            # logging.error(traceback.format_exc())
            return (
@@ -366,12 +366,24 @@ class KubernetesConnector:
                    volume_mounts=volume_mounts if volume_mounts else None,
                    security_context=security_context,
                )
            else:
            elif "required_resources" in descriptor_service_function:
                resources = self._get_requested_resources(descriptor_service_function)
                con = client.V1Container(
                    name=descriptor_service_function["name"],
                    image=container["image"],
                    ports=ports,
                    image_pull_policy="Always",
                    resources=resources,
                    env=envs if envs else None,
                    volume_mounts=volume_mounts if volume_mounts else None,
                    security_context=security_context,
                )
            else:
                con = client.V1Container(
                    name=descriptor_service_function["name"],
                    image=container["image"],
                    ports=ports,
                    image_pull_policy="IfNotPresent",
                    env=envs if envs else None,
                    volume_mounts=volume_mounts if volume_mounts else None,
                    security_context=security_context,
@@ -384,7 +396,9 @@ class KubernetesConnector:
        spec = client.V1DeploymentSpec(
            selector=selector,
            template=template,
            replicas=descriptor_service_function["count-min"],
            replicas=descriptor_service_function.get(
                "replica_count", descriptor_service_function["count-min"]
            ),
        )

        body = client.V1Deployment(
@@ -475,6 +489,22 @@ class KubernetesConnector:
            request_dict[auto_scale_policy["metric"]] = auto_scale_policy["request"]
        return client.V1ResourceRequirements(limits=limits_dict, requests=request_dict)
    
    def _get_requested_resources(self, descriptor_service_function):
        request_dict = {}
        limits_dict = {}
        cpu_pool = descriptor_service_function["required_resources"]["applicationResources"]["cpuPool"]
        request_dict["cpu"] = str(cpu_pool["numCPU"])+"m"
        request_dict["memory"] = str(cpu_pool["memory"])+"Mi"
        extra_fields = descriptor_service_function.get("extra_fields") or {}
        if extra_fields.get("cpuLimit") is not None:
            limits_dict["cpu"] = str(extra_fields["cpuLimit"]) + "m"
        if extra_fields.get("memoryLimit") is not None:
            limits_dict["memory"] = str(extra_fields["memoryLimit"]) + "Mi"
        return client.V1ResourceRequirements(
            limits=limits_dict if limits_dict else None,
            requests=request_dict,
        )

    def _get_pod_spec(self, descriptor_service_function, containers, volumes):
        if "location" in descriptor_service_function:
            node_selector_dict = {"nodeName": descriptor_service_function["location"]}
@@ -485,6 +515,14 @@ class KubernetesConnector:
                restart_policy="Always",
                volumes=volumes if volumes else None,
            )
        elif "nodeSelector" in descriptor_service_function:
            return client.V1PodSpec(
                containers=containers,
                node_selector=descriptor_service_function["nodeSelector"],
                hostname=descriptor_service_function["name"],
                restart_policy="Always",
                volumes=volumes if volumes else None,
            )
        else:
            return client.V1PodSpec(
                containers=containers,
@@ -504,20 +542,17 @@ class KubernetesConnector:
            "exposed_ports" in descriptor_service_function["containers"][0]
        ):  # create NodePort svc object
            ports = []
            hepler = 0
            for port_id in descriptor_service_function["containers"][0]["exposed_ports"]:

                # if "grafana" in descriptor_service_function["name"]:
                #     ports_=client.V1ServicePort(port=port_id,
                #                                 node_port=31000,
                #                                 target_port=port_id, name=str(port_id))
                # else:
                #     ports_ = client.V1ServicePort(port=port_id,
                #                                   # node_port=descriptor_paas["containers"][0]["exposed_ports"][hepler],
                #                                   target_port=port_id, name=str(port_id))
                ports_ = client.V1ServicePort(port=port_id, target_port=port_id, name=str(port_id))
            service_node_ports = descriptor_service_function.get("service_node_ports")
            for index, port_id in enumerate(
                descriptor_service_function["containers"][0]["exposed_ports"]
            ):
                ports_ = client.V1ServicePort(
                    port=port_id,
                    node_port=service_node_ports[index] if service_node_ports else None,
                    target_port=port_id,
                    name=str(port_id),
                )
                ports.append(ports_)
                hepler = hepler + 1
            spec = client.V1ServiceSpec(selector=dict_label, ports=ports, type="NodePort")
            # body = client.V1Service(api_version="v1", kind="Service", metadata=metadata, spec=spec)
        else:  # create ClusterIP svc object