diff --git a/clover-hub/src/server/modman/busses/proxies/group/can_2/bus_manager.rs b/clover-hub/src/server/modman/busses/proxies/group/can_2/bus_manager.rs index 7ce770c..c7d397e 100644 --- a/clover-hub/src/server/modman/busses/proxies/group/can_2/bus_manager.rs +++ b/clover-hub/src/server/modman/busses/proxies/group/can_2/bus_manager.rs @@ -1,6 +1,7 @@ use std::{ collections::HashMap, sync::Arc, + time::Duration, }; use anyhow::anyhow; @@ -20,6 +21,7 @@ use tracing::{ error, info, instrument, + warn, }; use crate::server::modman::{ @@ -35,7 +37,15 @@ use crate::server::modman::{ #[instrument] pub fn parse_id_str(id_str: &str) -> Result { - todo!() + let digits = id_str + .strip_prefix("0x") + .or_else(|| id_str.strip_prefix("0X")) + .unwrap_or(id_str); + + match u16::from_str_radix(digits, 16) { + Ok(val) => Ok(val), + Err(err) => Err(err.into()), + } } pub fn match_to_str(match_struct: regex::Match<'_>) -> &str { @@ -52,6 +62,8 @@ pub async fn can_bus_manager( let listener_registry: Arc>> = Arc::new(Mutex::new(HashMap::new())); + debug!("Now listening for requests for bus: {iface_name}..."); + while !cancellation_token.is_cancelled() { let port_statuses = can_2_port_status_mutex.lock().await; let mut port_statuses_snapshot = Vec::new(); @@ -59,10 +71,9 @@ pub async fn can_bus_manager( // Matches against strings like: `"can0/0xFFF:0xFFF"`. // This regex does not validate if the specified IDs are within the CAN range though. // That's done later by `embedded_can`. - let port_specifier_re = Regex::new( - r"^(?(?:\\w|[0-9])+)\\/(?0[xX][0-9a-fA-F]{3}):(?0[xX][0-9a-fA-F]{3})$", - ) - .unwrap(); + let port_specifier_re = + Regex::new(r"^(?\w+)\/(?0[xX][0-9a-fA-F]{3}):(?0[xX][0-9a-fA-F]{3})$") + .unwrap(); // We need to make sure that we're not leaving that mutex locked for too long. for port_status in port_statuses.iter() { @@ -77,32 +88,18 @@ pub async fn can_bus_manager( match port_specifier_re.captures(&port_path) { Some(re_captures) => { // Known good value since the regex is static, and the haystack matched. - let requested_iface = re_captures.name("iface").unwrap(); + let requested_iface = match_to_str(re_captures.get(1).unwrap()); + let rx_id = match_to_str(re_captures.get(2).unwrap()); + let tx_id = match_to_str(re_captures.get(3).unwrap()); - if match_to_str(requested_iface) == &iface_name { + if requested_iface == &iface_name { match port_status { PortStatus::Requested(module_id) => { setup_listener( ctx.clone(), - port_path.clone(), - module_id, - ( - match_to_str(re_captures.name("rx_id").unwrap()), - match_to_str(re_captures.name("tx_id").unwrap()), - ), - listener_registry.clone(), - ) - .await; - } - PortStatus::Unavailable(module_id) => { - setup_listener( - ctx.clone(), - port_path.clone(), + iface_name.clone(), module_id, - ( - match_to_str(re_captures.name("tx_id").unwrap()), - match_to_str(re_captures.name("tx_id").unwrap()), - ), + (rx_id, tx_id), listener_registry.clone(), ) .await; @@ -125,6 +122,9 @@ pub async fn can_bus_manager( None => {} } } + + // We don't wanna obliterate the CPU. + tokio::time::sleep(Duration::from_millis(100)).await; } for (module_id, listener_token) in listener_registry.lock().await.iter() { @@ -141,6 +141,8 @@ pub async fn setup_listener( id_tuple: (&str, &str), listener_registry: Arc>>, ) { + let mut bound_port = false; + match parse_id_str(id_tuple.0) { Ok(raw_rx_id) => match parse_id_str(id_tuple.1) { Ok(raw_tx_id) => { @@ -177,6 +179,8 @@ pub async fn setup_listener( &tx_options ) { Ok(tx_socket) => { + info!("Bound port: {}, for module: {module_id}!", format!("{iface_name}/{}:{}", id_tuple.0, id_tuple.1)); + let listener_token = CancellationToken::new(); listener_registry.lock().await.insert(module_id.clone(), listener_token.clone()); @@ -194,6 +198,8 @@ pub async fn setup_listener( tokio::task::spawn(async move { can_module_tx(tx_session, tx_token, tx_socket, tx_id).await; }); + + bound_port = true; }, Err(err) => { error!("Error while binding socketcan port: {} for TX, due to:\n{err}", format!("{iface_name}/{}:{}", id_tuple.0, id_tuple.1)); @@ -223,4 +229,30 @@ pub async fn setup_listener( error!("Error while parsing the rx id: {}:\n{err}", id_tuple.0); } } + + if bound_port { + debug!( + "Letting the rest of ModMan know that: {}, was bound!", + format!("{iface_name}/{}:{}", id_tuple.0, id_tuple.1) + ); + + let mut port_statuses = ctx.store.port_statuses.can_2.lock().await; + + port_statuses.insert( + format!("{iface_name}/{}:{}", id_tuple.0, id_tuple.1), + PortStatus::Bound(module_id), + ); + } else { + warn!( + "Letting the rest of ModMan know that: {}, was not bound.", + format!("{iface_name}/{}:{}", id_tuple.0, id_tuple.1) + ); + + let mut port_statuses = ctx.store.port_statuses.can_2.lock().await; + + port_statuses.insert( + format!("{iface_name}/{}:{}", id_tuple.0, id_tuple.1), + PortStatus::Unavailable(module_id), + ); + } } diff --git a/clover-hub/src/server/modman/busses/proxies/group/can_2/interface_lookout.rs b/clover-hub/src/server/modman/busses/proxies/group/can_2/interface_lookout.rs index 47fa085..673d6f5 100644 --- a/clover-hub/src/server/modman/busses/proxies/group/can_2/interface_lookout.rs +++ b/clover-hub/src/server/modman/busses/proxies/group/can_2/interface_lookout.rs @@ -109,6 +109,9 @@ pub fn can_interface_lookout(ctx: Arc, channel: UnboundedSender { if retries == 5 { diff --git a/clover-hub/src/server/modman/busses/proxies/group/can_2/module_listener.rs b/clover-hub/src/server/modman/busses/proxies/group/can_2/module_listener.rs index 3b52948..7760245 100644 --- a/clover-hub/src/server/modman/busses/proxies/group/can_2/module_listener.rs +++ b/clover-hub/src/server/modman/busses/proxies/group/can_2/module_listener.rs @@ -55,25 +55,29 @@ pub async fn can_module_rx( }) .await { - Ok(_) => match rmp_serde::from_slice::(&payload) { - Ok(decoded_payload) => match serde_json::to_string(&decoded_payload) { - Ok(json_payload) => match publisher.put(json_payload).await { - Ok(_) => { - debug!("Successfully proxied CAN 2 message from module: {module_id}!"); - } + Ok(_) => { + if payload.len() > 0 { + match rmp_serde::from_slice::(&payload) { + Ok(decoded_payload) => match serde_json::to_string(&decoded_payload) { + Ok(json_payload) => match publisher.put(json_payload).await { + Ok(_) => { + debug!("Successfully proxied CAN 2 message from module: {module_id}!"); + } + Err(err) => { + error!("Failed to publish message to Zenoh, at this stage, we've either lost connectivity, or there's a massive problem. (Might be a bug) Due to:\n{err}"); + } + }, + Err(err) => { + error!("Failed to produce a JSON payload from the decoded message, this is a bug and should be reported! Due to:\n{err}"); + } + }, Err(err) => { - error!("Failed to publish message to Zenoh, at this stage, we've either lost connectivity, or there's a massive problem. (Might be a bug) Due to:\n{err}"); + // TODO: Do we want to tell the module that it fucked up? + error!("Invalid message from module: {module_id}, this is a bug (or bad connection) and should (probably) be reported to the module maintainer! Happened due to:\n{err}"); } - }, - Err(err) => { - error!("Failed to produce a JSON payload from the decoded message, this is a bug and should be reported! Due to:\n{err}"); } - }, - Err(err) => { - // TODO: Do we want to tell the module that it fucked up? - error!("Invalid message from module: {module_id}, this is a bug (or bad connection) and should (probably) be reported to the module maintainer! Happened due to:\n{err}"); } - }, + } Err(err) => match err { RecvError::BufferTooSmall { needed, got } => { error!( @@ -95,6 +99,8 @@ pub async fn can_module_rx( // I trust rustc, but also we need to save memory!!! drop(payload); } + + debug!("Shutting down CAN 2 RX thread."); } Err(err) => { error!("Unable to create a zenoh broadcaster at: {key_expr}, there's probably an error in your configuration. Due to:\n{err}"); @@ -207,6 +213,8 @@ pub async fn can_module_tx( } } } + + debug!("Shutting down CAN 2 TX thread."); } Err(err) => { error!("Unable to create a zenoh queryable at: {key_expr}, there's probably an error in your configuration. Due to:\n{err}"); diff --git a/clover-hub/src/server/modman/connections.rs b/clover-hub/src/server/modman/connections.rs index 30ee40b..15feb8b 100644 --- a/clover-hub/src/server/modman/connections.rs +++ b/clover-hub/src/server/modman/connections.rs @@ -8,7 +8,8 @@ use serde::{ pub enum ModuleConnection { /// Device ID Simulated(String), - App(AppConnection), + // RFQDN of the App to bind to. Module ID will be passed to the app automatically. + App(String), /// Bus path and Device ID #[cfg(feature = "can_fd")] CANFD(CANFDConnection), @@ -32,12 +33,6 @@ pub enum ModuleConnection { UART(UARTConnection), } -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct AppConnection { - pub app_id: String, - pub device_id: String, -} - #[derive(Debug, Clone, Serialize, Deserialize)] pub struct CANFDConnection { pub bus_id: String, @@ -48,6 +43,7 @@ pub struct CANFDConnection { pub struct CAN2Connection { pub bus_id: String, pub device_id: String, + pub reply_id: String, } #[derive(Debug, Clone, Serialize, Deserialize)] diff --git a/clover-hub/src/server/modman/mod.rs b/clover-hub/src/server/modman/mod.rs index 84720cb..0f9e6bb 100644 --- a/clover-hub/src/server/modman/mod.rs +++ b/clover-hub/src/server/modman/mod.rs @@ -91,6 +91,7 @@ pub async fn modman_main( start_busses(bus_store, bus_session, bus_token).await.await; }); + let init_session = session.clone(); let init_store = Arc::new(store.clone()); let init_results = cancellation_tokens .0 @@ -149,8 +150,13 @@ pub async fn modman_main( if total_modules > 0 { // Initialize modules that were registered already via configuration and persistence. for (id, module) in modules_to_init.iter() { - let (initialized, _components_initialized) = - init_module(&init_store, id.clone(), module.clone()).await; + let (initialized, _components_initialized) = init_module( + &init_store, + id.clone(), + module.clone(), + init_session.clone(), + ) + .await; if initialized { modules_initialized += 1; diff --git a/clover-hub/src/server/modman/modules/connections/can_2.rs b/clover-hub/src/server/modman/modules/connections/can_2.rs new file mode 100644 index 0000000..8f9b95f --- /dev/null +++ b/clover-hub/src/server/modman/modules/connections/can_2.rs @@ -0,0 +1,38 @@ +use std::sync::Arc; + +use tracing::{ + debug, + instrument, +}; + +use crate::server::modman::{ + connections::CAN2Connection, + models::{ + modules::Module, + store::ModManStore, + PortStatus, + }, +}; + +#[instrument(skip(store, session))] +pub async fn setup_can_2_connection( + store: &ModManStore, + module: &Module, + id: &String, + connection: CAN2Connection, + session: Arc, +) -> Result<(), anyhow::Error> { + // can0/0x201:0x101 + let requested_port = format!( + "{}/{}:{}", + connection.bus_id, connection.reply_id, connection.device_id + ); + + debug!("Requesting CAN 2 port: {requested_port}, for module: {id}..."); + + let mut can_2_ports = store.port_statuses.can_2.lock().await; + can_2_ports.insert(requested_port, PortStatus::Requested(id.clone())); + drop(can_2_ports); + + Ok(()) +} diff --git a/clover-hub/src/server/modman/modules/connections/mod.rs b/clover-hub/src/server/modman/modules/connections/mod.rs new file mode 100644 index 0000000..ce3394a --- /dev/null +++ b/clover-hub/src/server/modman/modules/connections/mod.rs @@ -0,0 +1,2 @@ +#[cfg(feature = "can_2")] +pub mod can_2; diff --git a/clover-hub/src/server/modman/modules.rs b/clover-hub/src/server/modman/modules/mod.rs similarity index 79% rename from clover-hub/src/server/modman/modules.rs rename to clover-hub/src/server/modman/modules/mod.rs index 2143389..a40e695 100644 --- a/clover-hub/src/server/modman/modules.rs +++ b/clover-hub/src/server/modman/modules/mod.rs @@ -40,6 +40,8 @@ //! Similar to Level 3, however, the asymmetric keys are changed constantly to ensure perfect forward secrecy. //! +pub mod connections; + use super::{ components::models::CloverComponentTrait, models::{ @@ -47,6 +49,7 @@ use super::{ store::ModManStore, }, }; +use anyhow::anyhow; use std::sync::Arc; use tracing::{ debug, @@ -125,8 +128,13 @@ pub async fn init_component( } } -#[instrument(skip(store))] -pub async fn init_module(store: &ModManStore, id: String, module: Module) -> (bool, usize) { +#[instrument(skip(store, session))] +pub async fn init_module( + store: &ModManStore, + id: String, + module: Module, + session: Arc, +) -> (bool, usize) { let mut initialized_module = module.initialized; let mut initialized_module_components = 0; @@ -147,44 +155,87 @@ pub async fn init_module(store: &ModManStore, id: String, module: Module) -> (bo } else { let mut critical_failiure = None; - for (component_id, is_critical) in module.components.iter() { - match init_component( - &store.clone(), - id.clone(), - component_id.clone(), - is_critical, - ) - .await - { - Ok(_) => { - initialized_module_components += 1; - } - Err(e) => { - if is_critical.to_owned() { - critical_failiure = Some((component_id.clone(), e)); - } else { - error!( - "Module: {}, Failed to initialize component \"{}\", due to: {}", - id.clone(), - component_id.clone(), - e - ); + match module.connection.clone() { + crate::server::modman::connections::ModuleConnection::Simulated(_) => { + // TODO: Check for simulation bindings, and load them + // Otherwise fail. + } + crate::server::modman::connections::ModuleConnection::App(app_connection) => { + // TODO: App Module binding. + critical_failiure = Some(anyhow!(format!( + "Module: {id}, failed to bind to app: {app_connection}, due to: Unimplemented.", + ))); + } + #[cfg(feature = "can_2")] + crate::server::modman::connections::ModuleConnection::CAN2(can2_connection) => { + use crate::server::modman::modules::connections::can_2::setup_can_2_connection; + + match setup_can_2_connection( + &store, + &module, + &id, + can2_connection.clone(), + session.clone(), + ) + .await + { + Ok(_) => {} + Err(err) => { + critical_failiure = Some(anyhow!(format!( + "Module: {id}, failed to bind CAN 2 bus proxy: {}, due to:\n{err}", + format!( + "{}/{}:{}", + can2_connection.bus_id, can2_connection.reply_id, can2_connection.device_id + ) + ))); } } } } - debug!("Module: {id}, Done initializing components!"); + match critical_failiure { + None => { + for (component_id, is_critical) in module.components.iter() { + match init_component( + &store.clone(), + id.clone(), + component_id.clone(), + is_critical, + ) + .await + { + Ok(_) => { + initialized_module_components += 1; + } + Err(e) => { + if is_critical.to_owned() { + critical_failiure = Some(anyhow!(format!( + "Module: {id}, failed to initialize critical component: {}, due to: {}", + component_id.clone(), + e + ))); + } else { + error!( + "Module: {}, Failed to initialize component \"{}\", due to: {}", + id.clone(), + component_id.clone(), + e + ); + } + } + } + } + + info!("Module: {id}, Done initializing components!"); + } + Some(_) => { + warn!("There was a critical failiure when binding the module to the Bus, we're skipping component initialization!"); + } + } match critical_failiure { Some(failiure) => { - let (component_id, e) = failiure; - error!( - "Module: {}, failed to initialize critical component: {}, due to: {}\nSkipping rest of Module init...", - id.clone(), - component_id, - e - ); + error!("{failiure}\nSkipping rest of Module init..."); } Option::None => { if initialized_module_components != module.components.len() { @@ -196,20 +247,20 @@ pub async fn init_module(store: &ModManStore, id: String, module: Module) -> (bo module.components.len() ); initialized_module = true; + debug!("Module: {id}, Marked as initialized."); } else { error!("Module: {}, failed to initialize!", id.clone()); } } else { debug!("Finished initializing module."); initialized_module = true; + debug!("Module: {id}, Marked as initialized."); } } } } } - debug!("Module: {id}, Marked as initialized."); - // Update the store with new state of the module. if initialized_module { debug!("Module: {id}, Waiting on in-memory store lock..."); diff --git a/flake.nix b/flake.nix index 5a2138f..ae1b444 100644 --- a/flake.nix +++ b/flake.nix @@ -196,6 +196,7 @@ cloverHubCMD cloverHubZenoh flutter + can-utils ]) ++ libraries; };