Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
import os
import shutil
from nebula.utils import DockerUtils, APIUtils
from datetime import datetime, timezone
import docker
from nebula.controller.federation.federation_controller import FederationController
from nebula.controller.federation.scenario_builder import ScenarioBuilder
Expand All @@ -13,9 +14,10 @@
from nebula.config.config import Config
from nebula.core.utils.certificate import generate_ca_certificate
from nebula.core.utils.locker import Locker
from nebula.controller.federation.resource_manager import ResourceManager, ReleaseDevicesEvent, RAMOverusedEvent

class NebulaFederationDocker():
def __init__(self):
def __init__(self, timestamp: datetime):
self.scenario_name = ""
self.participants_alive = 0
self.round_per_participant = {}
Expand All @@ -31,6 +33,7 @@ def __init__(self):
self.participants_alive_lock = Locker("participants_alive_lock", async_lock=True)
self.config_dir = ""
self.log_dir = ""
self.timestamp = timestamp

async def get_additionals_to_be_deployed(self, config) -> list:
async with self.federation_deployment_lock:
Expand Down Expand Up @@ -85,12 +88,12 @@ def nfp(self):
###############################
"""

async def run_scenario(self, federation_id: str, scenario_data: Dict, user: str):
async def run_scenario(self, federation_id: str, scenario_data: Dict, user: str, rol: str):
#TODO maintain files on memory, not read them again
federation = await self._add_nebula_federation_to_pool(federation_id, user)
scenario_info = {}
if federation:
scenario_builder = ScenarioBuilder(federation_id, user=user)
scenario_builder = ScenarioBuilder(federation_id, user=user, rol=rol)
await self._initialize_scenario(scenario_builder, scenario_data, federation)
generate_ca_certificate(dir_path=self.cert_dir)
await self._load_configuration_and_start_nodes(scenario_builder, federation)
Expand Down Expand Up @@ -200,6 +203,8 @@ async def node_done(self, federation_id: str, node_done_request: NodeDoneRequest
self.logger.info(f"Node-Done received from node on federation ID: ({federation_id})")

if await nebula_federation.is_experiment_finish():
asyncio.create_task(self._release_devices(federation_id))

payload = node_done_request.model_dump()
self.logger.info(f"All nodes have finished on federation ID: ({federation_id}), reporting to hub..")
await self._remove_nebula_federation_from_pool(federation_id)
Expand Down Expand Up @@ -248,7 +253,7 @@ async def _add_nebula_federation_to_pool(self, federation_id: str, user: str):
fed = None
async with self._federations_dict_lock:
if not federation_id in self.nfp:
fed = NebulaFederationDocker()
fed = NebulaFederationDocker(datetime.now(timezone.utc))
self.nfp[federation_id] = fed
self.logger.info(f"SUCCESS: new ID: ({federation_id}) added to the pool")
else:
Expand All @@ -264,6 +269,13 @@ async def _remove_nebula_federation_from_pool(self, federation_id: str) -> Nebul
else:
self.logger.info(f"ERROR: trying to remove ({federation_id}) from federations pool..")
return None

async def _get_most_recent_federation(self):
async with self._federations_dict_lock:
if not self.nfp:
return None
federation_id = max(self.nfp, key=lambda k: self.nfp[k].timestamp)
return federation_id

async def _check_active_federation(self, federation_id: str) -> bool:
async with self._federations_dict_lock:
Expand Down Expand Up @@ -360,7 +372,7 @@ async def _initialize_scenario(self, sb: ScenarioBuilder, scenario_data, federat
self.logger.info(f"ERROR while creating files: {e}")

try:
participant_config = sb.build_scenario_config_for_node(index, node)
participant_config = await sb.build_scenario_config_for_node(index, node)
#self.logger.info(f"dictionary: {participant_config}")
except Exception as e:
self.logger.info(f"ERROR while building configuration for node: {e}")
Expand Down Expand Up @@ -615,3 +627,21 @@ def _start_node(self, scenario_name, node, network_name, base_network_name, base
json.dump(metadata, f, indent=2)

return success

""" ###############################
# RESOURCE MANAGEMENT #
###############################
"""

async def initialize_resources_functionalities(self):
await ResourceManager.get_instance().subscribe_resource_event(RAMOverusedEvent, self._ram_overused_event_callback)

async def _ram_overused_event_callback(self, roe: RAMOverusedEvent):
federation_id = await self._get_most_recent_federation()
if federation_id:
await self.stop_scenario(federation_id)

async def _release_devices(self, federation_id: str):
rde = ReleaseDevicesEvent(federation_id)
await ResourceManager.get_instance().publish_recource_event(rde)

Original file line number Diff line number Diff line change
Expand Up @@ -4,19 +4,19 @@
import os
import shutil
from nebula.utils import APIUtils
import docker
from datetime import datetime, timezone
from nebula.controller.federation.federation_controller import FederationController
from nebula.controller.federation.scenario_builder import ScenarioBuilder
from nebula.controller.federation.utils_requests import factory_requests
from nebula.controller.federation.utils_requests import RemoveScenarioRequest, NodeUpdateRequest, NodeDoneRequest
from typing import Dict
from fastapi import Request
from nebula.config.config import Config
from nebula.core.utils.certificate import generate_ca_certificate
from nebula.core.utils.locker import Locker
from nebula.controller.federation.resource_manager import ResourceManager, ReleaseDevicesEvent, RAMOverusedEvent

class NebulaFederationProcesses():
def __init__(self):
def __init__(self, timestamp: datetime):
self.scenario_name = ""
self.participants_alive = 0
self.round_per_participant = {}
Expand All @@ -32,6 +32,7 @@ def __init__(self):
self.participants_alive_lock = Locker("participants_alive_lock", async_lock=True)
self.config_dir = ""
self.log_dir = ""
self.timestamp = timestamp

async def get_additionals_to_be_deployed(self, config) -> list:
async with self.federation_deployment_lock:
Expand Down Expand Up @@ -86,12 +87,12 @@ def nfp(self):
###############################
"""

async def run_scenario(self, federation_id: str, scenario_data: Dict, user: str):
async def run_scenario(self, federation_id: str, scenario_data: Dict, user: str, rol: str):
#TODO maintain files on memory, not read them again
federation = await self._add_nebula_federation_to_pool(federation_id, user)
scenario_info = {}
if federation:
scenario_builder = ScenarioBuilder(federation_id, user=user)
scenario_builder = ScenarioBuilder(federation_id, user=user, rol=rol)
await self._initialize_scenario(scenario_builder, scenario_data, federation)
generate_ca_certificate(dir_path=self.cert_dir)
await self._load_configuration_and_start_nodes(scenario_builder, federation)
Expand Down Expand Up @@ -186,6 +187,8 @@ async def node_done(self, federation_id: str, node_done_request: NodeDoneRequest
self.logger.info(f"Node-Done received from node on federation ID: ({federation_id})")

if await nebula_federation.is_experiment_finish():
asyncio.create_task(self._release_devices(federation_id))

payload = node_done_request.model_dump()
self.logger.info(f"All nodes have finished on federation ID: ({federation_id}), reporting to hub..")
await self._remove_nebula_federation_from_pool(federation_id)
Expand Down Expand Up @@ -234,7 +237,7 @@ async def _add_nebula_federation_to_pool(self, federation_id: str, user: str):
fed = None
async with self._federations_dict_lock:
if not federation_id in self.nfp:
fed = NebulaFederationProcesses()
fed = NebulaFederationProcesses(datetime.now(timezone.utc))
self.nfp[federation_id] = fed
self.logger.info(f"SUCCESS: new ID: ({federation_id}) added to the pool")
else:
Expand All @@ -251,6 +254,13 @@ async def _remove_nebula_federation_from_pool(self, federation_id: str) -> Nebul
self.logger.info(f"ERROR: trying to remove ({federation_id}) from federations pool..")
return None

async def _get_most_recent_federation(self):
async with self._federations_dict_lock:
if not self.nfp:
return None
federation_id = max(self.nfp, key=lambda k: self.nfp[k].timestamp)
return federation_id

async def _check_active_federation(self, federation_id: str) -> bool:
async with self._federations_dict_lock:
if federation_id in self.nfp:
Expand Down Expand Up @@ -335,7 +345,7 @@ async def _initialize_scenario(self, sb: ScenarioBuilder, scenario_data, federat
self.logger.info(f"ERROR while creating files: {e}")

try:
participant_config = sb.build_scenario_config_for_node(index, node)
participant_config = await sb.build_scenario_config_for_node(index, node)
#self.logger.info(f"dictionary: {participant_config}")
except Exception as e:
self.logger.info(f"ERROR while building configuration for node: {e}")
Expand Down Expand Up @@ -560,3 +570,20 @@ def _write_commands_on_file(self, commands: str, federation: NebulaFederationPro
os.chmod(f"{federation.config_dir}/current_scenario_commands.sh", 0o755)
except Exception as e:
raise Exception(f"Error starting nodes as processes: {e}")

""" ###############################
# RESOURCE MANAGEMENT #
###############################
"""

async def initialize_resources_functionalities(self):
await ResourceManager.get_instance().subscribe_resource_event(RAMOverusedEvent, self._ram_overused_event_callback)

async def _ram_overused_event_callback(self, roe: RAMOverusedEvent):
federation_id = await self._get_most_recent_federation()
if federation_id:
await self.stop_scenario(federation_id)

async def _release_devices(self, federation_id: str):
rde = ReleaseDevicesEvent(federation_id)
await ResourceManager.get_instance().publish_recource_event(rde)
8 changes: 7 additions & 1 deletion nebula/controller/federation/federation_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
from nebula.controller.federation.federation_controller import FederationController
from nebula.controller.federation.factory_federation_controller import federation_controller_factory
from nebula.controller.federation.utils_requests import RemoveScenarioRequest, RunScenarioRequest, StopScenarioRequest, NodeUpdateRequest, NodeDoneRequest, Routes
from nebula.controller.federation.resource_manager import ResourceManager

fed_controllers: Dict[str, FederationController] = {}

Expand All @@ -30,9 +31,14 @@ async def lifespan(app: FastAPI):
controller_host = os.environ.get("NEBULA_CONTROLLER_HOST")
hub_url = f"http://{controller_host}:{hub_port}"

# Initialize resource manager to assign devices availables to federations
#TODO get maxRAM from environ
ResourceManager.get_instance(logger=logger, verbose=False).init()

#["docker", "processes", "physical"]
for exp_type in ["docker", "process"]:
fed_controllers[exp_type] = federation_controller_factory(exp_type, hub_url, logger)
await fed_controllers[exp_type].initialize_resources_functionalities()
logger.info(f"{exp_type} Federation controller created.")

yield
Expand All @@ -59,7 +65,7 @@ async def run_scenario(run_scenario_request: RunScenarioRequest):
logger.info(f"[API]: run experiment request for deployment type: {experiment_type}")
controller = fed_controllers.get(experiment_type, None)
if controller:
return await controller.run_scenario(run_scenario_request.federation_id, run_scenario_request.scenario_data, run_scenario_request.user)
return await controller.run_scenario(run_scenario_request.federation_id, run_scenario_request.scenario_data, run_scenario_request.user, run_scenario_request.rol)
else:
return {"message": "Experiment type not allowed"}

Expand Down
6 changes: 5 additions & 1 deletion nebula/controller/federation/federation_controller.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ def logger(self):
return self._logger

@abstractmethod
async def run_scenario(self, federation_id: str, scenario_data: Dict, user: str):
async def run_scenario(self, federation_id: str, scenario_data: Dict, user: str, rol: str):
pass

@abstractmethod
Expand All @@ -36,4 +36,8 @@ async def node_done(self, federation_id: str, node_done_request: NodeDoneRequest

abstractmethod
async def remove_scenario(self, federation_id: str, remove_scenario_request: RemoveScenarioRequest):
pass

abstractmethod
async def initialize_resources_functionalities(self):
pass
Loading