Loading src/sunrise6g_opensdk/edgecloud/adapters/kubernetes/client.py +80 −0 Original line number Diff line number Diff line Loading @@ -44,6 +44,7 @@ class EdgeApplicationManager(EdgeCloudManagementInterface): storage_uri = kwargs.get("EMP_STORAGE_URI") username = kwargs.get("KUBERNETES_USERNAME") namespace = kwargs.get("K8S_NAMESPACE") self.service_info_url = kwargs.get("SERVICE_INFO_URL") if base_url is not None and base_url != "": self.k8s_connector = KubernetesConnector( ip=self.kubernetes_host, Loading Loading @@ -113,6 +114,85 @@ class EdgeApplicationManager(EdgeCloudManagementInterface): decoded_intent[prop] = inner_blank_node.value return decoded_intent def save_service_level(self, intent, service_level_group): ttl_string = intent["expression"]["expressionValue"] decoded_intent = self._intent_to_dict(intent=ttl_string) service_level = intent["name"].split("_")[3] decoded_intent["service_level_group"] = service_level_group if service_level == "L1": decoded_intent["service_level_id"] = 1 if service_level == "L2": decoded_intent["service_level_id"] = 2 if service_level == "L3": decoded_intent["service_level_id"] = 3 self.connector_db.insert_document_service_level( decoded_intent, service_level_group ) return self.connector_db.count_service_levels_by_group(service_level_group) def _build_required_applications(self, data): result = {} for key, value in data.items(): if not key.startswith("required_applications_"): continue _, _, app_name, field = key.split("_", 3) if app_name not in result: result[app_name] = {} if field == "supported_service_level_ids": value = [int(x.strip()) for x in value.split(",") if x.strip()] result[app_name][field] = value return result def _build_combined_service_levels_payload(self, service_level_group): service_levels = self.connector_db.get_service_levels_by_group( service_level_group ) combined_service_levels = {} combined_service_levels["vertical_id"] = service_levels[0]["vertical_id"] combined_service_levels["timestamp"] = service_levels[0]["timestamp"] combined_service_levels["quality_vs_deployment_cost_weight"] = float( service_levels[0]["quality_vs_deployment_cost_weight"].to_decimal() ) combined_service_levels["required_applications"] = ( self._build_required_applications(service_levels[0]) ) service_level_catalog = [] for _service_level in service_levels: entry = {} entry["service_level_id"] = _service_level["service_level_id"] entry["service_type"] = _service_level["service_type"] entry["name"] = _service_level["name"] entry["service_area_m2"] = _service_level["service_area_m2"] entry["target_precision_m"] = float(_service_level["target_precision_m"].to_decimal()) entry["confidence_level"] = float(_service_level["confidence_level"].to_decimal()) entry["periodicity_per_min"] = _service_level["periodicity_per_min"] entry["duration_minutes"] = _service_level["duration_minutes"] entry["target_latency_s"] = float(_service_level["target_latency_s"].to_decimal()) entry["max_user_density_per_m2"] = _service_level["max_user_density_per_m2"] entry["shareable"] = _service_level["shareable"] service_level_catalog.append(entry) combined_service_levels["service_level_catalog"] = service_level_catalog return combined_service_levels def submit_combined_service_levels(self, service_level_group): combined_service_levels = self._build_combined_service_levels_payload( service_level_group ) return requests.post( self.service_info_url, json=combined_service_levels, headers={ "Accept": "application/json", "Content-Type": "application/json", }, timeout=(30, 120), ) def onboard_intent(self, intent): ttl_string = intent["expression"]["expressionValue"] Loading src/sunrise6g_opensdk/edgecloud/adapters/kubernetes/lib/utils/connector_db.py +41 −0 Original line number Diff line number Diff line Loading @@ -325,6 +325,47 @@ class ConnectorDB: except Exception as ce_: raise Exception("An exception occurred :", ce_) def count_service_levels_by_group(self, service_level_group): collection = "service_levels" myclient = pymongo.MongoClient(self._storage_url) mydbmongo = myclient[self.mydb_mongo] mycol = mydbmongo[collection] myquery = {"service_level_group": service_level_group} count = mycol.count_documents(myquery) return count def get_service_levels_by_group(self, service_level_group): collection = "service_levels" myclient = pymongo.MongoClient(self._storage_url) mydbmongo = myclient[self.mydb_mongo] mycol = mydbmongo[collection] myquery = {"service_level_group": service_level_group} service_levels_cur = mycol.find(myquery) service_levels = list(service_levels_cur) if len(service_levels) == 0: return 404 return service_levels def insert_document_service_level(self, service_level=None, service_level_group=None): collection = "service_levels" myclient = pymongo.MongoClient(self._storage_url) mydbmongo = myclient[self.mydb_mongo] mycol = mydbmongo[collection] myquery = {"service_level_id": service_level["service_level_id"], "service_level_group": service_level_group } mydoc = mycol.find_one(myquery) # keeps the last record (contains registrationStatus) if mydoc is not None: return 409 try: mycol.insert_one(service_level) return 200 except Exception as ce_: raise Exception("An exception occurred :", ce_) def get_documents_from_collection( self, collection_input, input_type=None, input_value=None Loading Loading
src/sunrise6g_opensdk/edgecloud/adapters/kubernetes/client.py +80 −0 Original line number Diff line number Diff line Loading @@ -44,6 +44,7 @@ class EdgeApplicationManager(EdgeCloudManagementInterface): storage_uri = kwargs.get("EMP_STORAGE_URI") username = kwargs.get("KUBERNETES_USERNAME") namespace = kwargs.get("K8S_NAMESPACE") self.service_info_url = kwargs.get("SERVICE_INFO_URL") if base_url is not None and base_url != "": self.k8s_connector = KubernetesConnector( ip=self.kubernetes_host, Loading Loading @@ -113,6 +114,85 @@ class EdgeApplicationManager(EdgeCloudManagementInterface): decoded_intent[prop] = inner_blank_node.value return decoded_intent def save_service_level(self, intent, service_level_group): ttl_string = intent["expression"]["expressionValue"] decoded_intent = self._intent_to_dict(intent=ttl_string) service_level = intent["name"].split("_")[3] decoded_intent["service_level_group"] = service_level_group if service_level == "L1": decoded_intent["service_level_id"] = 1 if service_level == "L2": decoded_intent["service_level_id"] = 2 if service_level == "L3": decoded_intent["service_level_id"] = 3 self.connector_db.insert_document_service_level( decoded_intent, service_level_group ) return self.connector_db.count_service_levels_by_group(service_level_group) def _build_required_applications(self, data): result = {} for key, value in data.items(): if not key.startswith("required_applications_"): continue _, _, app_name, field = key.split("_", 3) if app_name not in result: result[app_name] = {} if field == "supported_service_level_ids": value = [int(x.strip()) for x in value.split(",") if x.strip()] result[app_name][field] = value return result def _build_combined_service_levels_payload(self, service_level_group): service_levels = self.connector_db.get_service_levels_by_group( service_level_group ) combined_service_levels = {} combined_service_levels["vertical_id"] = service_levels[0]["vertical_id"] combined_service_levels["timestamp"] = service_levels[0]["timestamp"] combined_service_levels["quality_vs_deployment_cost_weight"] = float( service_levels[0]["quality_vs_deployment_cost_weight"].to_decimal() ) combined_service_levels["required_applications"] = ( self._build_required_applications(service_levels[0]) ) service_level_catalog = [] for _service_level in service_levels: entry = {} entry["service_level_id"] = _service_level["service_level_id"] entry["service_type"] = _service_level["service_type"] entry["name"] = _service_level["name"] entry["service_area_m2"] = _service_level["service_area_m2"] entry["target_precision_m"] = float(_service_level["target_precision_m"].to_decimal()) entry["confidence_level"] = float(_service_level["confidence_level"].to_decimal()) entry["periodicity_per_min"] = _service_level["periodicity_per_min"] entry["duration_minutes"] = _service_level["duration_minutes"] entry["target_latency_s"] = float(_service_level["target_latency_s"].to_decimal()) entry["max_user_density_per_m2"] = _service_level["max_user_density_per_m2"] entry["shareable"] = _service_level["shareable"] service_level_catalog.append(entry) combined_service_levels["service_level_catalog"] = service_level_catalog return combined_service_levels def submit_combined_service_levels(self, service_level_group): combined_service_levels = self._build_combined_service_levels_payload( service_level_group ) return requests.post( self.service_info_url, json=combined_service_levels, headers={ "Accept": "application/json", "Content-Type": "application/json", }, timeout=(30, 120), ) def onboard_intent(self, intent): ttl_string = intent["expression"]["expressionValue"] Loading
src/sunrise6g_opensdk/edgecloud/adapters/kubernetes/lib/utils/connector_db.py +41 −0 Original line number Diff line number Diff line Loading @@ -325,6 +325,47 @@ class ConnectorDB: except Exception as ce_: raise Exception("An exception occurred :", ce_) def count_service_levels_by_group(self, service_level_group): collection = "service_levels" myclient = pymongo.MongoClient(self._storage_url) mydbmongo = myclient[self.mydb_mongo] mycol = mydbmongo[collection] myquery = {"service_level_group": service_level_group} count = mycol.count_documents(myquery) return count def get_service_levels_by_group(self, service_level_group): collection = "service_levels" myclient = pymongo.MongoClient(self._storage_url) mydbmongo = myclient[self.mydb_mongo] mycol = mydbmongo[collection] myquery = {"service_level_group": service_level_group} service_levels_cur = mycol.find(myquery) service_levels = list(service_levels_cur) if len(service_levels) == 0: return 404 return service_levels def insert_document_service_level(self, service_level=None, service_level_group=None): collection = "service_levels" myclient = pymongo.MongoClient(self._storage_url) mydbmongo = myclient[self.mydb_mongo] mycol = mydbmongo[collection] myquery = {"service_level_id": service_level["service_level_id"], "service_level_group": service_level_group } mydoc = mycol.find_one(myquery) # keeps the last record (contains registrationStatus) if mydoc is not None: return 409 try: mycol.insert_one(service_level) return 200 except Exception as ce_: raise Exception("An exception occurred :", ce_) def get_documents_from_collection( self, collection_input, input_type=None, input_value=None Loading