diff --git a/Cargo.lock b/Cargo.lock index d9dd10c..b00153e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2097,6 +2097,7 @@ dependencies = [ "clap", "clover-hub-macros", "decorum", + "embedded-can", "futures", "git2", "i2c", @@ -4664,6 +4665,7 @@ dependencies = [ "libc", "socket2 0.5.10", "socketcan", + "tokio", ] [[package]] diff --git a/clover-hub/Cargo.toml b/clover-hub/Cargo.toml index 9ec9733..629b4d2 100644 --- a/clover-hub/Cargo.toml +++ b/clover-hub/Cargo.toml @@ -92,7 +92,7 @@ bevy = "0.15.0" # Busses can-iso-tp = { workspace = true, optional = true } -linux-socketcan-iso-tp = { version = "0.1.3", optional = true } +linux-socketcan-iso-tp = { version = "0.1.3", optional = true, features = ["tokio"] } bluer = { version = "0.17.3", optional = true } spidev = { version = "0.6.0", optional = true } i2c = { version = "0.1.0", optional = true } @@ -106,3 +106,4 @@ rmp-serde = "1.3.0" rodio = "0.20.1" serde_bytes = { workspace = true } strum_macros = "0.28.0" +embedded-can = "0.4.1" 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 new file mode 100644 index 0000000..97febf9 --- /dev/null +++ b/clover-hub/src/server/modman/busses/proxies/group/can_2/bus_manager.rs @@ -0,0 +1,88 @@ +use std::sync::Arc; + +use can_iso_tp::{ + self, + IsoTpNode, +}; +use embedded_can::Id; +use linux_socketcan_iso_tp::{ + self, + IsoTpKernelOptions, + TokioSocketCanIsoTp, +}; +use tokio::sync::{ + broadcast::{ + channel as broadcast_channel, + Receiver as BroadcastReceiver, + Sender as BroadcastSender, + }, + oneshot::{ + channel as oneshot_channel, + Receiver as OneShotReceiver, + Sender as OneShotSender, + }, +}; +use tokio_util::sync::CancellationToken; +use tracing::instrument; + +use crate::server::modman::busses::{ + models::BusMessage, + proxies::group::can_2::CAN2Bus, +}; + +#[instrument(skip(ctx, cancellation_token))] +pub async fn can_bus_manager( + ctx: Arc, + cancellation_token: CancellationToken, + iface_details: (String, u32), +) { + let (iface_name, iface_index) = iface_details; + + // Modules shouldn't be sending events over CAN before we introduce ourselves, but we have a buffer just in case. + let (to_bus_tx, to_bus) = broadcast_channel(16); + let (from_bus, from_bus_rx) = broadcast_channel(16); + + // I'd rather return the JoinHandle from the function, but due to using libc, we should avoid doing that. + let (status_channel, status_channel_rx) = oneshot_channel(); + + tokio::task::spawn(async move { + can_bus_listener( + to_bus, + from_bus, + cancellation_token.clone(), + status_channel, + // We recreate the iface_details tuple due to an error created by the instrument macro. + (iface_name.clone(), iface_index), + ) + .await + }); + + match status_channel_rx.await { + Ok(_) => while !ctx.cancellation_token.is_cancelled() {}, + Err(err) => todo!(), + } +} + +#[instrument(skip(to_bus, from_bus, cancellation_token))] +pub async fn can_bus_listener( + to_bus: BroadcastReceiver, + from_bus: BroadcastSender, + cancellation_token: CancellationToken, + status_channel: OneShotSender>, + iface_details: (String, u32), +) { + let options = IsoTpKernelOptions::default(); + let (iface_name, _iface_index) = iface_details; + + match TokioSocketCanIsoTp::open( + &iface_name, + Id::Standard(embedded_can::StandardId::new(0x101).expect("0x101 is a valid standard CAN ID")), + Id::Standard(embedded_can::StandardId::new(0x201).expect("0x201 is a valid standard CAN ID")), + &options, + ) { + Ok(socket) => { + + } + Err(err) => todo!(), + } +} diff --git a/clover-hub/src/server/modman/busses/proxies/group/can_2/lookout.rs b/clover-hub/src/server/modman/busses/proxies/group/can_2/lookout.rs new file mode 100644 index 0000000..fbcb18d --- /dev/null +++ b/clover-hub/src/server/modman/busses/proxies/group/can_2/lookout.rs @@ -0,0 +1,158 @@ +use std::sync::Arc; +use std::time::Duration; +use std::{ + collections::HashMap, + thread::sleep as std_sleep, +}; + +use crate::server::modman::busses::proxies::group::can_2::CAN2Bus; + +use nix::net::if_::if_nameindex; +use tokio::sync::mpsc::{ + unbounded_channel, + UnboundedReceiver, + UnboundedSender, +}; +use tokio_util::sync::CancellationToken; +use tracing::{ + debug, + error, + info, + instrument, +}; + +#[derive(Debug, Clone)] +pub enum CanLookoutEvent { + IFaceCreate((String, u32)), + IFaceDestroy(String), +} + +#[instrument(skip(ctx))] +pub async fn can_lookout_thread(ctx: Arc) { + let (lookout_tx, lookout_rx) = unbounded_channel::(); + + let registrar_ctx = ctx.clone(); + let lookout_ctx = ctx.clone(); + + // Something something tokio task pool is limited. + std::thread::spawn(|| can_interface_lookout(lookout_ctx, lookout_tx)); + tokio::task::spawn(async move { can_bus_registrar(registrar_ctx, lookout_rx).await }).await; +} + +/// Detects network interfaces to try and bind. +#[instrument(skip(ctx, channel))] +pub fn can_interface_lookout(ctx: Arc, channel: UnboundedSender) { + let mut known_ifaces: HashMap = HashMap::new(); + let mut retries = 0; + + while !ctx.cancellation_token.is_cancelled() { + match if_nameindex() { + Ok(detected_ifaces) => { + let mut ifaces_to_save = Vec::new(); + + // Detect new interfaces, and mark known ones to be preserved. + for detected_iface in detected_ifaces.iter() { + let detected_iface_name = detected_iface.name().to_string_lossy().to_string(); + + match known_ifaces.get(&detected_iface_name) { + Some(_iface) => { + ifaces_to_save.push(detected_iface_name); + } + None => { + debug!("Detected new network interface: {}!", &detected_iface_name); + match channel.send(CanLookoutEvent::IFaceCreate(( + detected_iface_name.clone(), + detected_iface.index(), + ))) { + Ok(_) => {} + Err(err) => { + error!("Failed to inform async registry thread that the '{}' interface was discovered, due to:\n{err}", detected_iface_name.clone()); + } + } + + known_ifaces.insert(detected_iface_name.clone(), detected_iface.index()); + ifaces_to_save.push(detected_iface_name); + } + } + } + + // Copy the current status into a new hashmap to avoid borrow errors. + let mut known_ifaces_snapshot = HashMap::new(); + + for known_iface in known_ifaces.iter() { + known_ifaces_snapshot.insert(known_iface.0.clone(), known_iface.1.clone()); + } + + // Remove all interfaces that shouldn't be saved. + for known_iface in known_ifaces_snapshot.iter() { + let mut should_save = false; + + for iface_to_save in ifaces_to_save.iter() { + if iface_to_save == known_iface.0 { + should_save = true; + } + } + + if !should_save { + debug!( + "Interface: {}, dissapeared, removing it...", + known_iface.0.clone() + ); + known_ifaces.remove(known_iface.0); + match channel.send(CanLookoutEvent::IFaceDestroy(known_iface.0.clone())) { + Ok(_) => {} + Err(err) => { + error!("Failed to inform async registry thread that the '{}' interface was destroyed, due to:\n{err}", known_iface.0.clone()); + } + } + } + } + } + Err(err) => { + if retries == 5 { + error!("Continously failed to get network interfaces even after inital check, did something happen to the network manager?"); + debug!("{err}"); + ctx.cancellation_token.cancel(); + break; + } else { + std_sleep(Duration::from_millis(500)); + retries += 1; + } + } + } + } +} + +#[instrument(skip(ctx, channel))] +pub async fn can_bus_registrar(ctx: Arc, mut channel: UnboundedReceiver) { + let mut bus_registry: HashMap = HashMap::new(); + // Prevent the mutex from being locked for the entire time we run the bus proxy. + let config = ctx.store.config.lock().await; + let can2_config = config.modman.group_busses.can_2.clone(); + drop(config); + + while !ctx.cancellation_token.is_cancelled() { + if let Some(lookout_event) = channel.recv().await { + match lookout_event { + CanLookoutEvent::IFaceCreate(iface_details) => { + let (iface_name, iface_index) = iface_details; + + for permitted_iface in can2_config.permitted_interfaces.clone() { + if permitted_iface == iface_name { + match bus_registry.get(&iface_name) { + Some(_) => {} + None => { + info!( + "Found configured interface: {}, starting up a CAN bus listener...", + iface_name.clone() + ); + } + } + } + } + } + CanLookoutEvent::IFaceDestroy(iface_name) => todo!(), + } + } + } +} diff --git a/clover-hub/src/server/modman/busses/proxies/group/can_2/mod.rs b/clover-hub/src/server/modman/busses/proxies/group/can_2/mod.rs index 662b3ec..a2718b3 100644 --- a/clover-hub/src/server/modman/busses/proxies/group/can_2/mod.rs +++ b/clover-hub/src/server/modman/busses/proxies/group/can_2/mod.rs @@ -1,25 +1,36 @@ //! # CAN 2.0 Bus Proxy //! +//! The primary proxy bus for Clover. //! -//! -use std::sync::Arc; +pub mod bus_manager; +pub mod lookout; + +use std::{ + sync::Arc, + time::Duration, +}; use crate::server::modman::{ - busses::models::{ - Bus, - BusTypes, - }, + busses::proxies::group::can_2::lookout::can_lookout_thread, models::store::ModManStore, }; -use can_iso_tp; -use linux_socketcan_iso_tp; +use anyhow::anyhow; use nix::net::if_::if_nameindex; +use serde::{ + Deserialize, + Serialize, +}; +use tokio::time::sleep; +use tokio_util::sync::CancellationToken; +use tracing::instrument; #[derive(Debug, Clone)] pub struct CAN2Bus { + pub session: Arc, pub store: Arc, + pub cancellation_token: CancellationToken, } #[derive(Debug, Clone)] @@ -28,15 +39,32 @@ pub struct InterfaceToBind { pub path: String, } -impl Bus for CAN2Bus { - async fn subscribe_to_bus( - mut self, - session: Arc, - ) -> Result, anyhow::Error> { - todo!() +#[instrument(skip(ctx))] +pub async fn spawn_lookout_thread( + ctx: Arc, +) -> Result, anyhow::Error> { + let mut error = None; + let mut retries = 0; + + while retries != 5 { + match if_nameindex() { + Ok(_) => return Ok(tokio::task::spawn(async { can_lookout_thread(ctx).await })), + Err(err) => { + error = Some(err.into()); + retries += 1; + sleep(Duration::from_millis(500)).await; + } + } } - fn get_type() -> BusTypes { - BusTypes::CAN2 + match error { + Some(err) => Err(err), + None => Err(anyhow!("")), } } + +#[derive(Serialize, Deserialize, Debug, Clone, Default)] +pub struct CAN2Config { + /// Interfaces that we should bind to. + pub permitted_interfaces: Vec, +} diff --git a/clover-hub/src/server/modman/busses/proxies/group/mod.rs b/clover-hub/src/server/modman/busses/proxies/group/mod.rs index 9c4d848..2c1ebc1 100644 --- a/clover-hub/src/server/modman/busses/proxies/group/mod.rs +++ b/clover-hub/src/server/modman/busses/proxies/group/mod.rs @@ -17,6 +17,14 @@ //! Module manager threads handle the actual zenoh endpoints for modules (like how drivers expose device paths in linux). //! +use serde::{ + Deserialize, + Serialize, +}; + +#[cfg(feature = "can_2")] +use crate::server::modman::busses::proxies::group::can_2::CAN2Config; + #[cfg(feature = "can_2")] pub mod can_2; #[cfg(feature = "can_fd")] @@ -25,3 +33,9 @@ pub mod can_fd; pub mod i2c; #[cfg(feature = "spi")] pub mod spi; + +#[derive(Serialize, Deserialize, Debug, Clone, Default)] +pub struct GroupBusConfigs { + #[cfg(feature = "can_2")] + pub can_2: CAN2Config, +} diff --git a/clover-hub/src/server/modman/models/config.rs b/clover-hub/src/server/modman/models/config.rs index 3a33c6e..168897e 100644 --- a/clover-hub/src/server/modman/models/config.rs +++ b/clover-hub/src/server/modman/models/config.rs @@ -5,21 +5,23 @@ use serde::{ Serialize, }; -use crate::server::modman::models::{ - components::{ - CloverComponent, - CloverComponentMeta, +use crate::server::modman::{ + busses::proxies::group::GroupBusConfigs, + models::{ + components::{ + CloverComponent, + CloverComponentMeta, + }, + gestures::GestureStates, + modules::Module, }, - gestures::GestureStates, - modules::Module, }; #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ModManConfig { /// All ports available for modman to use to connect to modules. pub uart_ports: Vec, - /// All (CAN 2) network interfaces available for modman to use to connect to modules. - pub can_2_interfaces: Vec, + pub group_busses: GroupBusConfigs, /// Whether to restart paused gestures automatically on startup. pub restart_gestures: bool, pub gesture_states: HashMap, @@ -140,7 +142,7 @@ impl Default for ModManConfig { static_components, static_modules, uart_ports: Default::default(), - can_2_interfaces: Default::default(), + group_busses: Default::default(), restart_gestures: Default::default(), gesture_states: Default::default(), gestures_bg_by_default: Default::default(), diff --git a/clover-hub/src/server/warehouse/repos/mod.rs b/clover-hub/src/server/warehouse/repos/mod.rs index 4e96cb5..75669d4 100644 --- a/clover-hub/src/server/warehouse/repos/mod.rs +++ b/clover-hub/src/server/warehouse/repos/mod.rs @@ -49,6 +49,7 @@ use tokio_stream::{ wrappers::ReadDirStream, StreamExt, }; +use tracing::instrument; use super::models::WarehouseStore; @@ -81,6 +82,7 @@ impl From for Error { /// /// // TODO: Make these a constant so all built-in strings get updated at once! +#[instrument] pub fn builtin_rfqdn(is_core: bool) -> String { if is_core { String::from("com.reboot-codes.clover.CORE") @@ -90,6 +92,7 @@ pub fn builtin_rfqdn(is_core: bool) -> String { } /// Replace `@here`, `@base`, and `@builtin` manifest value directives. +#[instrument] pub fn replace_simple_directives(value: String, resolution_ctx: ResolutionCtx) -> String { debug!( "replace_simple_directives (provided): {} + {:#?}", @@ -167,6 +170,7 @@ pub fn replace_simple_directives(value: String, resolution_ctx: ResolutionCtx) - String::from(val) } +#[instrument] pub async fn resolve_list_entry( raw_list: HashMap>, resolution_ctx: ResolutionCtx, @@ -174,7 +178,7 @@ pub async fn resolve_list_entry( ) -> Result, SimpleError> where K: ManifestCompilationFrom, - T: for<'a> Deserialize<'a>, + T: for<'a> Deserialize<'a> + std::fmt::Debug, { let mut err = None; let mut entries = HashMap::new(); @@ -327,6 +331,7 @@ where /// Source of filesystem structures like `/opt/clover/repos/com/reboot-codes/clover/@repo`. /// /// Using `@repo` for the actual repository keeps everything unique and organized and allows for nested repo bases (e.g. an unstable repo for testing out the latest apps). For this reason `@repo` is a banned directory name in Clover-compatible remote repositories. +#[instrument(skip(store))] pub async fn update_repo_dir_structure( repo_dir_path: OsPath, store: Arc, @@ -364,6 +369,7 @@ pub async fn update_repo_dir_structure( /// Used to resolve repo manifest entry **values** that may have directives (`@import`, `@base`, `@here`, `@builtin`) in them. /// Hands off to [replace_simple_directives] if it isn't an import. +#[instrument] pub async fn resolve_entry_value( value: String, resolution_ctx: ResolutionCtx, @@ -521,6 +527,7 @@ pub async fn resolve_entry_value( /// Downloads repository updates from their origin remote using git. /// Git implicitly supports both HTTP(S) and SSH, so users have options when getting updates. +#[instrument(skip(store))] pub async fn download_repo_updates( store: Arc, repo_dir_path: OsPath, diff --git a/core/modules/two-phase-tail/module.clover.jsonc b/core/modules/two-phase-tail/module.clover.jsonc index 12157e0..b91d16a 100644 --- a/core/modules/two-phase-tail/module.clover.jsonc +++ b/core/modules/two-phase-tail/module.clover.jsonc @@ -12,7 +12,7 @@ // Data type, in this case, each message must contain 2 floats. "input": "vec2d", // Data type, this phase has a position reporting function and will send a 2D vector back regularly - "output": "vec2d", + "output": "vec2d" }, "phase-2": { "type": "movement", @@ -21,7 +21,7 @@ }, "tip-light": { "type": "indicator", - "input": "rgba", + "input": "rgba" } } }