use std::fmt::Debug;
use std::fs;
use std::path::PathBuf;
use std::str::FromStr;
use std::sync::{Arc as StdArc, atomic::AtomicUsize};
use log::LevelFilter;
use serde::{Serialize, de::DeserializeOwned};
use tokio::sync::{Notify, mpsc};
use tokio::task;
use tokio::time::{Duration, sleep};
use crate::containerisation::client_broker::ClientBroker;
use crate::containerisation::hyperion_container::HyperionContainer;
use crate::containerisation::traits::{
ContainerIdentidy, HeartbeatConfigProvider, HyperionContainerDirectiveMessage,
HyperionHeartbeatMessage, Initialisable, LogLevel, Run,
};
use crate::heartbeat::handler::{HeartbeatMissedHandler, HeartbeatTimeoutHandler};
use crate::logging::logging_service::initialise_logger;
use crate::network::network_topology::NetworkTopology;
use crate::network::server::Server;
use crate::utilities::load_config;
pub async fn create<A, C, T>(
config_path_str: &str,
network_topology_path_str: &str,
container_state: StdArc<AtomicUsize>,
container_state_notify: StdArc<Notify>,
main_rx: mpsc::Receiver<T>,
timeout_handler: Option<Box<dyn HeartbeatTimeoutHandler>>,
missed_handler: Option<Box<dyn HeartbeatMissedHandler>>,
) -> HyperionContainer<T>
where
A: Initialisable<ConfigType = C> + Run<Message = T> + Send + 'static + Sync + Debug,
C: Debug
+ Send
+ 'static
+ DeserializeOwned
+ Sync
+ LogLevel
+ ContainerIdentidy
+ HeartbeatConfigProvider,
T: HyperionContainerDirectiveMessage
+ HyperionHeartbeatMessage
+ Debug
+ Send
+ 'static
+ DeserializeOwned
+ Sync
+ Clone
+ Serialize,
{
let config_path: PathBuf = fs::canonicalize(config_path_str)
.unwrap_or_else(|e| panic!("Could not canonicalize '{config_path_str}': {e}"));
let component_config: StdArc<C> = load_config::load_config::<C>(&config_path)
.unwrap_or_else(|e| panic!("Failed to load component config from '{config_path:?}': {e}"));
let network_topology_path: PathBuf = fs::canonicalize(network_topology_path_str)
.unwrap_or_else(|e| panic!("Could not canonicalize '{network_topology_path_str}': {e}"));
let network_topology: StdArc<NetworkTopology> =
load_config::load_config::<NetworkTopology>(&network_topology_path).unwrap_or_else(|e| {
panic!("Failed to load network topology from '{network_topology_path:?}': {e}")
});
let log_level: LevelFilter = LevelFilter::from_str(component_config.log_level())
.unwrap_or_else(|e| {
println!("Log level was not parsed correctly: {e:?}\nDefaulting to 'Trace' log level.");
LevelFilter::Trace
});
initialise_logger(log_level).unwrap_or_else(|e| panic!("Failed to initialise logger: {e:?}"));
for (key, value) in component_config.container_identity().iter() {
log::debug!("{key}: {value}");
}
log::info!(
"Building Hyperion Container for {}...",
component_config
.container_identity()
.get("name")
.unwrap_or(&"Unknown".to_string())
);
let component_archetype = A::initialise(
container_state.clone(),
container_state_notify.clone(),
component_config.clone(),
);
let (server_tx, server_rx) = mpsc::channel::<T>(32);
let arc_server: StdArc<Server<T>> = Server::new(
network_topology.server_address.clone(),
server_tx,
container_state.clone(),
container_state_notify.clone(),
);
task::spawn(async move {
if let Err(e) = Server::run(arc_server).await {
log::error!("Server encountered an error: {e:?}");
}
});
sleep(Duration::from_secs(2)).await;
let client_broker: ClientBroker<T> = ClientBroker::init(
network_topology,
container_state.clone(),
container_state_notify.clone(),
);
sleep(Duration::from_secs(2)).await;
let container_name = component_config
.container_identity()
.get("name")
.cloned()
.unwrap_or_else(|| "Unknown".to_string());
let heartbeat_config = component_config.heartbeat_config();
HyperionContainer::<T>::create(
component_archetype,
container_state,
container_state_notify,
client_broker,
main_rx,
server_rx,
heartbeat_config,
container_name,
timeout_handler,
missed_handler,
)
}