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

feat: tf-sdk now retrieves and saves available nodes during startup

parent aff7f212
Loading
Loading
Loading
Loading
+15 −0
Original line number Diff line number Diff line
@@ -50,6 +50,21 @@ class EdgeApplicationManager(EdgeCloudManagementInterface):
            )
        if storage_uri is not None:
            self.connector_db = ConnectorDB(storage_uri)
            self.connector_db.wait_until_ready()
            self.get_and_save_nodes()


    def get_and_save_nodes(self):
        nodes = self.k8s_connector.get_nodes()
        self.connector_db.insert_document_k8s_nodes(nodes=nodes)
        return build_custom_http_response(
            status_code=200,
            content={"Found " + str(len(nodes)) + " nodes"}, headers={"Content-Type": "application/json"},
            encoding="utf-8",
            url=None,
            request=None
        )


    def onboard_app(self, app_manifest: AppManifest) -> Response:
        print(f"Submitting application: {app_manifest}")
+22 −0
Original line number Diff line number Diff line
from enum import StrEnum

from pydantic import BaseModel


class NodeTier(StrEnum):
    BRONZE = "bronze"
    SILVER = "silver"
    GOLD = "gold"
    NO_TIER = "no-tier"


class Node(BaseModel):
    id: str
    edge_cloud_zone_id: str
    tier: NodeTier
    available_cpu: int | None = None
    available_memory: int | None = None
    available_gpu: int | None = 0
    total_cpu: int
    total_memory: int
    total_gpu: int | None = 0
+19 −0
Original line number Diff line number Diff line
@@ -10,6 +10,11 @@ class ConnectorDB:
        self._storage_url = host
        self.mydb_mongo = "pi-edge"


    def wait_until_ready(self):
        myclient = pymongo.MongoClient(self._storage_url)
        myclient.admin.command("ping")

    # def insert_document_k8s_platform(self, document=None, _id=None):
    #     collection = "kubernetes_platforms"
    #     myclient = pymongo.MongoClient(self._storage_url)
@@ -251,6 +256,20 @@ class ConnectorDB:
        except Exception as ce_:
            raise Exception("An exception occurred :", ce_)

    def insert_document_k8s_nodes(self, nodes=None):
        collection = "nodes"
        myclient = pymongo.MongoClient(self._storage_url)
        mydbmongo = myclient[self.mydb_mongo]
        mycol = mydbmongo[collection]

        for node in nodes or []:
            try:
                node_json = node.model_dump()
                node_json["_id"] = node.id
                mycol.update_one({"_id": node.id}, {"$set": node_json}, upsert=True)
            except Exception as ce_:
                raise Exception("An exception occurred :", ce_)

    def get_documents_from_collection(
        self, collection_input, input_type=None, input_value=None
    ) -> List[dict]:
+19 −0
Original line number Diff line number Diff line
@@ -6,6 +6,7 @@ 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.utils import (
    auxiliary_functions,
@@ -182,6 +183,24 @@ class KubernetesConnector:

        return pop_output
        
    def get_nodes(self):
        namespace = self.api_instance_corev1api.read_namespace("kube-system")
        cluster_id = namespace.metadata.uid
        k8s_nodes = self.api_instance_corev1api.list_node()
        nodes = []
        for node in k8s_nodes.items:
            _node = {
                "edge_cloud_zone_id": cluster_id,
                "id": node.metadata.uid,
                "tier": node.metadata.labels.get("node-tier", "no-tier"),
                "total_cpu": int(node.status.allocatable["cpu"]),
                "total_memory": int(node.status.allocatable["memory"][:-2])
            }
            nodes.append(Node(**_node))

        return nodes


    def get_PoPs(self):

        try: