Skip to content
2 changes: 1 addition & 1 deletion nebula/config/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,7 @@ def __default_config(self):

def __set_default_logging(self, mode="w"):
experiment_name = self.participant["scenario_args"]["name"]
self.log_dir = os.path.join(self.participant["tracking_args"]["log_dir"], experiment_name)
self.log_dir =self.participant["tracking_args"]["log_dir"]
if not os.path.exists(self.log_dir):
os.makedirs(self.log_dir)
self.log_filename = f"{self.log_dir}/participant_{self.participant['device_args']['idx']}"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,8 @@ def __init__(self):
self.federation_round: int = 0
self.federation_deployment_lock = Locker("federation_deployment_lock", async_lock=True)
self.participants_alive_lock = Locker("participants_alive_lock", async_lock=True)
self.config_dir = ""
self.log_dir = ""

async def get_additionals_to_be_deployed(self, config) -> list:
async with self.federation_deployment_lock:
Expand Down Expand Up @@ -89,7 +91,7 @@ async def run_scenario(self, federation_id: str, scenario_data: Dict, user: str)
federation = await self._add_nebula_federation_to_pool(federation_id, user)
scenario_info = {}
if federation:
scenario_builder = ScenarioBuilder(federation_id)
scenario_builder = ScenarioBuilder(federation_id, user=user)
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 @@ -301,12 +303,14 @@ async def _initialize_scenario(self, sb: ScenarioBuilder, scenario_data, federat
# Initialize Scenario builder using scenario_data from user
self.logger.info("🔧 Initializing Scenario Builder using scenario data")
sb.set_scenario_data(scenario_data)
scenario_name = sb.get_scenario_name()
scenario_name = sb.get_scenario_name(user_to=True)

self.root_path = os.environ.get("NEBULA_ROOT_HOST")
self.host_platform = os.environ.get("NEBULA_HOST_PLATFORM")
self.config_dir = os.path.join(os.environ.get("NEBULA_CONFIG_DIR"), scenario_name)
self.log_dir = os.environ.get("NEBULA_LOGS_DIR")
# self.config_dir = os.path.join(os.environ.get("NEBULA_CONFIG_DIR"), scenario_name)
# self.log_dir = os.path.join(os.environ.get("NEBULA_LOGS_DIR"), scenario_name)
federation.config_dir = os.path.join(os.environ.get("NEBULA_CONFIG_DIR"), scenario_name)
federation.log_dir = os.path.join(os.environ.get("NEBULA_LOGS_DIR"), scenario_name)
self.cert_dir = os.environ.get("NEBULA_CERTS_DIR")
self.advanced_analytics = os.environ.get("NEBULA_ADVANCED_ANALYTICS", "False") == "True"
#self.config = Config(entity="FederationController")
Expand All @@ -317,17 +321,17 @@ async def _initialize_scenario(self, sb: ScenarioBuilder, scenario_data, federat
self.url = f"{os.environ.get('NEBULA_CONTROLLER_HOST')}:{os.environ.get('NEBULA_FEDERATION_CONTROLLER_PORT')}"

# Create Scenario management dirs
os.makedirs(self.config_dir, exist_ok=True)
os.makedirs(os.path.join(self.log_dir, scenario_name), exist_ok=True)
os.makedirs(federation.config_dir, exist_ok=True)
os.makedirs(federation.log_dir, exist_ok=True)
os.makedirs(self.cert_dir, exist_ok=True)

# Give permissions to the directories
os.chmod(self.config_dir, 0o777)
os.chmod(os.path.join(self.log_dir, scenario_name), 0o777)
os.chmod(federation.config_dir, 0o777)
os.chmod(federation.log_dir, 0o777)
os.chmod(self.cert_dir, 0o777)

# Save the scenario configuration
scenario_file = os.path.join(self.config_dir, "scenario.json")
scenario_file = os.path.join(federation.config_dir, "scenario.json")
with open(scenario_file, "w") as f:
json.dump(scenario_data, f, sort_keys=False, indent=2)

Expand All @@ -337,13 +341,13 @@ async def _initialize_scenario(self, sb: ScenarioBuilder, scenario_data, federat
settings = {
"scenario_name": scenario_name,
"root_path": self.root_path,
"config_dir": self.config_dir,
"log_dir": self.log_dir,
"config_dir": federation.config_dir,
"log_dir": federation.log_dir,
"cert_dir": self.cert_dir,
"env": None,
}

settings_file = os.path.join(self.config_dir, "settings.json")
settings_file = os.path.join(federation.config_dir, "settings.json")
with open(settings_file, "w") as f:
json.dump(settings, f, sort_keys=False, indent=2)

Expand All @@ -359,7 +363,7 @@ async def _initialize_scenario(self, sb: ScenarioBuilder, scenario_data, federat
self.logger.info(f"Creating .json file for participant: {index}, Configuration: {node}")
node_config = node
try:
participant_file = os.path.join(self.config_dir, f"participant_{node_config['id']}.json")
participant_file = os.path.join(federation.config_dir, f"participant_{node_config['id']}.json")
self.logger.info(f"Filename: {participant_file}")
os.makedirs(os.path.dirname(participant_file), exist_ok=True)
except Exception as e:
Expand All @@ -383,7 +387,7 @@ async def _initialize_scenario(self, sb: ScenarioBuilder, scenario_data, federat
async def _load_configuration_and_start_nodes(self, sb: ScenarioBuilder, federation: NebulaFederationDocker):
self.logger.info("🔧 Loading Scenario configuration...")
# Get participants configurations
participant_files = glob.glob(f"{self.config_dir}/participant_*.json")
participant_files = glob.glob(f"{federation.config_dir}/participant_*.json")
participant_files.sort()
if len(participant_files) == 0:
raise ValueError("No participant files found in config folder")
Expand All @@ -408,19 +412,19 @@ async def _load_configuration_and_start_nodes(self, sb: ScenarioBuilder, federat
self.logger.info("🔧 Building preload configuration for initial nodes...")
for i in range(n_nodes):
try:
with open(f"{self.config_dir}/participant_" + str(i) + ".json") as f:
with open(f"{federation.config_dir}/participant_" + str(i) + ".json") as f:
participant_config = json.load(f)
except Exception as e:
self.logger.info(f"ERROR: open/load participant .json")

self.logger.info(f"Building preload conf for participant {i}")
try:
sb.build_preload_initial_node_configuration(i, participant_config, self.log_dir, self.config_dir, self.cert_dir, self.advanced_analytics)
sb.build_preload_initial_node_configuration(i, participant_config, federation.log_dir, federation.config_dir, self.cert_dir, self.advanced_analytics)
except Exception as e:
self.logger.info(f"ERROR: cannot build preload configuration")

try:
with open(f"{self.config_dir}/participant_" + str(i) + ".json", "w") as f:
with open(f"{federation.config_dir}/participant_" + str(i) + ".json", "w") as f:
json.dump(participant_config, f, sort_keys=False, indent=2)
except Exception as e:
self.logger.info(f"ERROR: cannot dump preload configuration into participant .json file")
Expand All @@ -441,7 +445,7 @@ async def _load_configuration_and_start_nodes(self, sb: ScenarioBuilder, federat
federation.config.set_participants_config(participant_files)

# Add role to the topology (visualization purposes)
sb.visualize_topology(config_participants, path=f"{self.config_dir}/topology.png", plot=False)
sb.visualize_topology(config_participants, path=f"{federation.config_dir}/topology.png", plot=False)

# Additional participants
self.logger.info("🔧 Building preload configuration for additional nodes...")
Expand All @@ -451,7 +455,7 @@ async def _load_configuration_and_start_nodes(self, sb: ScenarioBuilder, federat
last_participant_index = len(participant_files)

for i, _ in enumerate(additional_participants):
additional_participant_file = f"{self.config_dir}/participant_{last_participant_index + i}.json"
additional_participant_file = f"{federation.config_dir}/participant_{last_participant_index + i}.json"
shutil.copy(last_participant_file, additional_participant_file)

with open(additional_participant_file) as f:
Expand All @@ -475,7 +479,7 @@ async def _load_configuration_and_start_nodes(self, sb: ScenarioBuilder, federat
self.logger.info("✅ Loading Scenario configuration done")

# Build dataset
dataset = sb.configure_dataset(self.config_dir)
dataset = sb.configure_dataset(federation.config_dir)
self.logger.info(f"🔧 Splitting {sb.get_dataset_name()} dataset...")
dataset.initialize_dataset()
self.logger.info(f"✅ Splitting {sb.get_dataset_name()} dataset... Done")
Expand All @@ -502,7 +506,7 @@ def _get_participant_container_name(self, scenario_name, idx: int) -> str:

def _start_initial_nodes(self, sb: ScenarioBuilder, federation: NebulaFederationDocker):
self.logger.info("Starting nodes using Docker Compose...")
federation.network_name = self._get_network_name(f"{sb.get_scenario_name()}-net-scenario")
federation.network_name = self._get_network_name(f"{sb.get_scenario_name(user_to=True)}-net-scenario")
federation.base_network_name = self._get_network_name("net-base")

# Create the Docker network
Expand All @@ -520,7 +524,7 @@ def _start_initial_nodes(self, sb: ScenarioBuilder, federation: NebulaFederation
# deploy initial nodes
self.logger.info(f"Deployment starting for participant {idx}")
federation.round_per_participant[idx] = 0
deployed_successfully = self._start_node(sb.get_scenario_name(), node, federation.network_name, federation.base_network_name, federation.base, federation.last_index_deployed, federation)
deployed_successfully = self._start_node(sb.get_scenario_name(user_to=True), node, federation.network_name, federation.base_network_name, federation.base, federation.last_index_deployed, federation)
if deployed_successfully:
federation.last_index_deployed += 1
federation.participants_alive += 1
Expand Down Expand Up @@ -568,8 +572,8 @@ def _start_node(self, scenario_name, node, network_name, base_network_name, base
),
base_network_name: client.api.create_endpoint_config(),
})
node["tracking_args"]["log_dir"] = "/nebula/app/logs"
node["tracking_args"]["config_dir"] = f"/nebula/app/config/{scenario_name}"
node["tracking_args"]["log_dir"] = federation.log_dir
node["tracking_args"]["config_dir"] = federation.config_dir
node["scenario_args"]["controller"] = self.url
node["scenario_args"]["deployment"] = "docker"
node["security_args"]["certfile"] = f"/nebula/app/certs/participant_{node['device_args']['idx']}_cert.pem"
Expand All @@ -583,7 +587,7 @@ def _start_node(self, scenario_name, node, network_name, base_network_name, base
except docker.errors.NotFound:
pass # No conflict, safe to proceed
# Write the config file in config directory
with open(f"{self.config_dir}/participant_{node['device_args']['idx']}.json", "w") as f:
with open(f"{federation.config_dir}/participant_{node['device_args']['idx']}.json", "w") as f:
json.dump(node, f, indent=4)
try:
container_id = client.api.create_container(
Expand All @@ -610,14 +614,14 @@ def _start_node(self, scenario_name, node, network_name, base_network_name, base

# Write scenario-level metadata for cleanup
scenario_metadata = {"containers": container_names, "network": network_name}
with open(os.path.join(self.config_dir, "scenario.metadata"), "a") as f:
with open(os.path.join(federation.config_dir, "scenario.metadata"), "a") as f:
if i == 2:
json.dump(scenario_metadata, f, indent=2)
else:
with open(os.path.join(self.config_dir, "scenario.metadata"), "r") as f:
with open(os.path.join(federation.config_dir, "scenario.metadata"), "r") as f:
metadata = json.load(f)
metadata["containers"].extend(container_names)
with open(os.path.join(self.config_dir, "scenario.metadata"), "w") as f:
with open(os.path.join(federation.config_dir, "scenario.metadata"), "w") as f:
json.dump(metadata, f, indent=2)

return success
Loading