From 0cd0199c602fb824eea5825f8796c315d9eba9c0 Mon Sep 17 00:00:00 2001 From: Reboot-Codes Date: Sat, 20 Jun 2026 11:06:56 -0700 Subject: [PATCH] Update connection logic, update UART bus to continous scanning. - Proxied connections should be global to the module, and not per-component. - Every bus proxy should also continously scan for requested bindings due to the dynamic nature of App Modules. --- clover-hub/src/server/modman/busses/mod.rs | 38 +- clover-hub/src/server/modman/busses/models.rs | 11 +- .../src/server/modman/busses/proxies/app.rs | 9 +- .../modman/busses/proxies/bt_classic.rs | 9 +- .../src/server/modman/busses/proxies/bt_le.rs | 9 +- .../src/server/modman/busses/proxies/can_2.rs | 9 +- .../server/modman/busses/proxies/can_fd.rs | 9 +- .../src/server/modman/busses/proxies/i2c.rs | 13 +- .../src/server/modman/busses/proxies/mod.rs | 4 +- .../src/server/modman/busses/proxies/spi.rs | 9 +- .../src/server/modman/busses/proxies/uart.rs | 391 +++++++++++------- .../server/modman/components/audio/impls.rs | 4 +- .../server/modman/components/audio/models.rs | 7 +- .../src/server/modman/components/models.rs | 40 +- .../modman/components/movement/models.rs | 11 +- .../modman/components/sensors/models.rs | 7 +- .../modman/components/video/cameras/models.rs | 7 +- .../components/video/displays/models.rs | 7 +- clover-hub/src/server/modman/connections.rs | 69 ++++ clover-hub/src/server/modman/ipc/displays.rs | 2 +- clover-hub/src/server/modman/mod.rs | 1 + clover-hub/src/server/modman/models.rs | 48 ++- clover-hub/src/server/modman/modules.rs | 2 + clover-hub/src/server/renderer/mod.rs | 4 +- .../system_ui/systems/displays/mod.rs | 28 +- clover-hub/src/server/warehouse/mod.rs | 2 +- .../src/server/warehouse/repos/models.rs | 4 +- listeners.md | 38 ++ 28 files changed, 485 insertions(+), 307 deletions(-) create mode 100644 clover-hub/src/server/modman/connections.rs create mode 100644 listeners.md diff --git a/clover-hub/src/server/modman/busses/mod.rs b/clover-hub/src/server/modman/busses/mod.rs index b8f3d9f..4c6fb27 100644 --- a/clover-hub/src/server/modman/busses/mod.rs +++ b/clover-hub/src/server/modman/busses/mod.rs @@ -8,6 +8,7 @@ pub mod proxies; use super::models::ModManStore; use log::info; +use models::Bus; use std::sync::Arc; pub async fn start_busses( @@ -17,14 +18,6 @@ pub async fn start_busses( info!("Starting ModMan Proxy Busses..."); let mut handles = vec![]; - // TODO: ASAP: Move to creating a sub-user for each module! - let proxy_session = session.clone(); - handles.push(tokio::task::spawn(async move { - /*while let Ok(raw_msg) = rx.recv().await { - debug!("Got proxy message from bus: {:?raw_msg}"); - }*/ - })); - // TODO: Add config options for each bus! handles.push(tokio::task::spawn(async move { info!("Starting App Bus..."); @@ -62,20 +55,21 @@ pub async fn start_busses( #[cfg(feature = "uart")] { - handles.push(tokio::task::spawn(async move { - info!("Starting UART Bus..."); - /*match (UARTBus { - store: store.clone(), - }) - .subscribe_to_bus(tx.clone(), uart_channel.0) - .await - { - Ok(handles) => { - futures::future::join_all(handles).await; - } - Err(_) => todo!(), - }*/ - })); + use proxies::uart::UARTBus; + + let uart_session = session.clone(); + info!("Starting UART Bus..."); + match (UARTBus { + store: store.clone(), + }) + .subscribe_to_bus(uart_session) + .await + { + Ok(handle) => { + handles.push(handle); + } + Err(_) => todo!(), + } } futures::future::join_all(handles) diff --git a/clover-hub/src/server/modman/busses/models.rs b/clover-hub/src/server/modman/busses/models.rs index e33e8a4..e71efda 100644 --- a/clover-hub/src/server/modman/busses/models.rs +++ b/clover-hub/src/server/modman/busses/models.rs @@ -1,3 +1,5 @@ +use std::sync::Arc; + use serde::{ Deserialize, Serialize, @@ -25,10 +27,9 @@ pub enum BusTypes { pub trait Bus { /// Send a message and expect a reply to that message. /// Listen to the Bus (does NOT contain IDs.) - async fn subscribe_to_bus( - &mut self, - //from_bus: tokio::sync::broadcast::Sender, - //to_bus: tokio::sync::broadcast::Sender, - ) -> Result>, anyhow::Error>; + fn subscribe_to_bus( + self, + session: Arc, + ) -> impl std::future::Future, anyhow::Error>> + Send; fn get_type() -> BusTypes; } diff --git a/clover-hub/src/server/modman/busses/proxies/app.rs b/clover-hub/src/server/modman/busses/proxies/app.rs index 6b76fc2..b2d2e7d 100644 --- a/clover-hub/src/server/modman/busses/proxies/app.rs +++ b/clover-hub/src/server/modman/busses/proxies/app.rs @@ -1,3 +1,5 @@ +use std::sync::Arc; + use crate::server::modman::busses::models::{ Bus, BusTypes, @@ -7,10 +9,9 @@ pub struct AppBus {} impl Bus for AppBus { async fn subscribe_to_bus( - &mut self, - //from_bus: tokio::sync::broadcast::Sender, - //to_bus: tokio::sync::broadcast::Sender, - ) -> Result>, anyhow::Error> { + mut self, + session: Arc, + ) -> Result, anyhow::Error> { todo!() } diff --git a/clover-hub/src/server/modman/busses/proxies/bt_classic.rs b/clover-hub/src/server/modman/busses/proxies/bt_classic.rs index 463e36c..00b891f 100644 --- a/clover-hub/src/server/modman/busses/proxies/bt_classic.rs +++ b/clover-hub/src/server/modman/busses/proxies/bt_classic.rs @@ -1,3 +1,5 @@ +use std::sync::Arc; + use crate::server::modman::busses::models::{ Bus, BusTypes, @@ -9,10 +11,9 @@ pub struct BluetoothBus {} impl Bus for BluetoothBus { async fn subscribe_to_bus( - &mut self, - from_bus: tokio::sync::broadcast::Sender, - to_bus: tokio::sync::broadcast::Sender, - ) -> Result>, anyhow::Error> { + mut self, + session: Arc, + ) -> Result, anyhow::Error> { todo!() } diff --git a/clover-hub/src/server/modman/busses/proxies/bt_le.rs b/clover-hub/src/server/modman/busses/proxies/bt_le.rs index d51b8ff..838b63c 100644 --- a/clover-hub/src/server/modman/busses/proxies/bt_le.rs +++ b/clover-hub/src/server/modman/busses/proxies/bt_le.rs @@ -1,3 +1,5 @@ +use std::sync::Arc; + use crate::server::modman::busses::models::{ Bus, BusTypes, @@ -9,10 +11,9 @@ pub struct BluetoothLEBus {} impl Bus for BluetoothLEBus { async fn subscribe_to_bus( - &mut self, - from_bus: tokio::sync::broadcast::Sender, - to_bus: tokio::sync::broadcast::Sender, - ) -> Result>, anyhow::Error> { + mut self, + session: Arc, + ) -> Result, anyhow::Error> { todo!() } diff --git a/clover-hub/src/server/modman/busses/proxies/can_2.rs b/clover-hub/src/server/modman/busses/proxies/can_2.rs index ab17fdb..494acc6 100644 --- a/clover-hub/src/server/modman/busses/proxies/can_2.rs +++ b/clover-hub/src/server/modman/busses/proxies/can_2.rs @@ -1,3 +1,5 @@ +use std::sync::Arc; + use crate::server::modman::busses::models::{ Bus, BusTypes, @@ -10,10 +12,9 @@ pub struct CAN2Bus {} impl Bus for CAN2Bus { async fn subscribe_to_bus( - &mut self, - from_bus: tokio::sync::broadcast::Sender, - to_bus: tokio::sync::broadcast::Sender, - ) -> Result>, anyhow::Error> { + mut self, + session: Arc, + ) -> Result, anyhow::Error> { todo!() } diff --git a/clover-hub/src/server/modman/busses/proxies/can_fd.rs b/clover-hub/src/server/modman/busses/proxies/can_fd.rs index 7109366..22bd318 100644 --- a/clover-hub/src/server/modman/busses/proxies/can_fd.rs +++ b/clover-hub/src/server/modman/busses/proxies/can_fd.rs @@ -1,3 +1,5 @@ +use std::sync::Arc; + use crate::server::modman::busses::models::{ Bus, BusTypes, @@ -10,10 +12,9 @@ pub struct CANFDBus {} impl Bus for CANFDBus { async fn subscribe_to_bus( - &mut self, - from_bus: tokio::sync::broadcast::Sender, - to_bus: tokio::sync::broadcast::Sender, - ) -> Result>, anyhow::Error> { + mut self, + session: Arc, + ) -> Result, anyhow::Error> { todo!() } diff --git a/clover-hub/src/server/modman/busses/proxies/i2c.rs b/clover-hub/src/server/modman/busses/proxies/i2c.rs index 5e5fd43..f4b08b7 100644 --- a/clover-hub/src/server/modman/busses/proxies/i2c.rs +++ b/clover-hub/src/server/modman/busses/proxies/i2c.rs @@ -1,23 +1,20 @@ +use std::sync::Arc; + use crate::server::modman::busses::models::{ Bus, BusTypes, }; use i2c; use i2cdev; -use serde::{ - Deserialize, - Serialize, -}; #[derive(Debug, Clone)] pub struct I2CBus {} impl Bus for I2CBus { async fn subscribe_to_bus( - &mut self, - from_bus: tokio::sync::broadcast::Sender, - to_bus: tokio::sync::broadcast::Sender, - ) -> Result>, anyhow::Error> { + mut self, + session: Arc, + ) -> Result, anyhow::Error> { todo!() } diff --git a/clover-hub/src/server/modman/busses/proxies/mod.rs b/clover-hub/src/server/modman/busses/proxies/mod.rs index 3ef5bbc..b844567 100644 --- a/clover-hub/src/server/modman/busses/proxies/mod.rs +++ b/clover-hub/src/server/modman/busses/proxies/mod.rs @@ -4,7 +4,7 @@ //! pub mod app; -/*#[cfg(feature = "bt_classic")] +#[cfg(feature = "bt_classic")] pub mod bt_classic; #[cfg(feature = "bt_le")] pub mod bt_le; @@ -17,4 +17,4 @@ pub mod i2c; #[cfg(feature = "spi")] pub mod spi; #[cfg(feature = "uart")] -pub mod uart;*/ +pub mod uart; diff --git a/clover-hub/src/server/modman/busses/proxies/spi.rs b/clover-hub/src/server/modman/busses/proxies/spi.rs index c2d343b..4526927 100644 --- a/clover-hub/src/server/modman/busses/proxies/spi.rs +++ b/clover-hub/src/server/modman/busses/proxies/spi.rs @@ -1,3 +1,5 @@ +use std::sync::Arc; + use crate::server::modman::busses::models::{ Bus, BusTypes, @@ -9,10 +11,9 @@ pub struct SPIBus {} impl Bus for SPIBus { async fn subscribe_to_bus( - &mut self, - from_bus: tokio::sync::broadcast::Sender, - to_bus: tokio::sync::broadcast::Sender, - ) -> Result>, anyhow::Error> { + mut self, + session: Arc, + ) -> Result, anyhow::Error> { todo!() } diff --git a/clover-hub/src/server/modman/busses/proxies/uart.rs b/clover-hub/src/server/modman/busses/proxies/uart.rs index a354ff3..7601031 100644 --- a/clover-hub/src/server/modman/busses/proxies/uart.rs +++ b/clover-hub/src/server/modman/busses/proxies/uart.rs @@ -1,3 +1,8 @@ +//! # UART Proxy Bus +//! +//! The UART proxy bus is designed bind one serial port per module, then expose I/O over Zenoh. +//! + use std::sync::Arc; use crate::server::modman::{ @@ -5,197 +10,289 @@ use crate::server::modman::{ Bus, BusTypes, }, + connections::ModuleConnection, models::{ ModManStore, PortStatus, }, + MODULE_EVT_ID, }; -use log::{ - debug, - warn, +use anyhow::anyhow; +use serialport::{ + self, + SerialPortInfo, }; -use serialport; -use tokio::io::{ - split, - AsyncReadExt, - AsyncWriteExt, +use tokio::{ + io::{ + split, + AsyncReadExt, + AsyncWriteExt, + ReadHalf, + WriteHalf, + }, + task::JoinHandle, }; use tokio_serial::{ self, SerialStream, }; +use tracing::{ + debug, + error, + instrument, + warn, +}; #[derive(Debug, Clone)] pub struct UARTBus { pub store: Arc, } +#[derive(Debug, Clone)] +pub struct PortToBind { + module_id: String, + path: String, +} + impl Bus for UARTBus { + #[instrument(name = "uart_bus", skip(self, session))] async fn subscribe_to_bus( - &mut self, - from_bus: tokio::sync::broadcast::Sender, - to_bus: tokio::sync::broadcast::Sender, - ) -> Result>, anyhow::Error> { - let warn_config = self.store.config.lock().await; - if !(warn_config.modman.uart_ports.len() > 0) { - warn!("No UART ports configured to proxy messages to! You may see message delivery errors from `nexus::user` due to a closed channel (the UARTBus proxy); they can be safely ignored."); - } - drop(warn_config); - + self, + session: Arc, + ) -> Result, anyhow::Error> { match serialport::available_ports() { Ok(port_info) => { - let mut handles = vec![]; - - for port in port_info { - debug!("Found port: {:#?}!", port.clone()); - - let mut attempt_bind = None; - let config = self.store.config.lock().await; - let mut port_statuses = self.store.port_statuses.uart.lock().await; - - for allowed_port in config.modman.uart_ports.clone() { - if allowed_port.0 == port.port_name.clone() { - let mut bound = false; - let port_path = allowed_port.0.clone(); - - for port_status in port_statuses.iter() { - match port_status.1 { - PortStatus::Available => {} - PortStatus::Requested(component_id) => { - if port_status.0 == &port_path { - attempt_bind = Some((component_id.clone(), allowed_port)); + return Ok(tokio::task::spawn(async move { + loop { + let mut handles = vec![]; + + for port in port_info.clone() { + let mut attempt_bind: Option = None; + let config = self.store.config.lock().await; + let mut port_statuses = self.store.port_statuses.uart.lock().await; + let allowed_ports = config.modman.uart_ports.clone(); + drop(config); + + for allowed_port in allowed_ports { + if allowed_port == port.port_name.clone() { + let mut bound = false; + let port_path = allowed_port.clone(); + + for port_status in port_statuses.iter() { + match port_status.1 { + PortStatus::Available => {} + PortStatus::Requested(module_id) => { + debug!("Port: {allowed_port}, requested for Module ID: {module_id}"); + + if port_status.0 == &port_path { + attempt_bind = Some(PortToBind { + module_id: module_id.clone(), + path: allowed_port, + }); + } + + break; + } + PortStatus::Bound(_port_status) => { + bound = true; + } + PortStatus::Unavailable(module_id) => { + if port_status.0 == &port_path { + attempt_bind = Some(PortToBind { + module_id: module_id.clone(), + path: allowed_port, + }); + } + break; + } } - break; } - PortStatus::Bound(_) => { - bound = true; - } - PortStatus::Unavailable(component_id) => { - if port_status.0 == &port_path { - attempt_bind = Some((component_id.clone(), allowed_port)); + + match attempt_bind.clone() { + Some(_bind_info) => {} + Option::None => { + if !bound { + port_statuses.insert(port_path.clone(), PortStatus::Available); + } } - break; } + + break; } } - match attempt_bind { - Some(_) => {} - Option::None => { - if !bound { - port_statuses.insert(port_path.clone(), PortStatus::Available); + std::mem::drop(port_statuses); + + match attempt_bind.clone() { + Some(bind_info) => { + match bind_uart_port(bind_info, port, self.store.clone(), session.clone()).await { + Ok(mut sub_handles) => { + handles.append(&mut sub_handles); + } + Err(err) => { + error!("Failed to bind port due to:\n{err}"); + } } } + None => {} } - - break; } + + futures::future::join_all(handles).await; } + })); + } + Err(err) => Err(anyhow!("Failed to get available ports due to:\n{err}")), + } + } - std::mem::drop(port_statuses); - - match attempt_bind { - Some(bind_info) => { - // TODO: Create new Nexus user for this component!! - - let (component_id, port_info) = bind_info; - debug!( - "Component: {}, Attempting to bind to: {}...", - component_id.clone(), - port.port_name.clone() - ); - match SerialStream::open(&tokio_serial::new(port_info.0.clone(), port_info.1)) { - Ok(bound_port) => { - let (mut port_read, mut port_write) = split(bound_port); - let from_port = from_bus.clone(); - let mut to_port_rx = to_bus.subscribe(); - self - .store - .port_statuses - .uart - .lock() - .await - .insert(port_info.0.clone(), PortStatus::Bound(component_id.clone())); - debug!( - "Component: {}, Port: {}, bound!", - component_id.clone(), - port.port_name.clone() - ); - - handles.push(tokio::task::spawn(async move { - let mut sub_handles = vec![]; - - let tx_port_ctx = (component_id.clone(), port.port_name.clone()); - sub_handles.push(tokio::task::spawn(async move { - while let Ok(msg) = to_port_rx.recv().await { - debug!( - "Component: {}, Port: {}, sending message from Nexus...", - tx_port_ctx.0.clone(), - tx_port_ctx.1.clone() - ); - - // TODO: Encrypt - match rmp_serde::to_vec(&msg) { - Ok(msg_vec) => { - port_write.write(msg_vec.as_slice()).await; - } - Err(_) => { - // TODO: - } - } - } - })); - - let rx_port_ctx = (component_id.clone(), port.port_name.clone()); - sub_handles.push(tokio::task::spawn(async move { - let mut buffer = Vec::new(); - while let Ok(msg_size) = port_read.read(buffer.as_mut_slice()).await { - debug!( - "Component: {}, Port: {}, got message of length: {}, parsing...", - rx_port_ctx.0.clone(), - rx_port_ctx.1.clone(), - msg_size - ); - - // TODO: Decryption - match rmp_serde::from_slice::(buffer.as_slice()) { - Ok(msg) => { - debug!( - "Component: {}, Port: {}, parsed message: {:#?}", - rx_port_ctx.0.clone(), - rx_port_ctx.1.clone(), - msg.clone() - ); - from_port.send(msg); - } - Err(_) => { - // TODO: - } - } + fn get_type() -> BusTypes { + BusTypes::UART + } +} - buffer = vec![]; - } - })); +/// Reusable code to bind a serial port for a function, and then run Zenoh endpoints for that binding. +/// Should be called on startup for static modules, and then dynamically for app modules. +#[instrument(skip(store, session))] +pub async fn bind_uart_port( + bind_info: PortToBind, + port: SerialPortInfo, + store: Arc, + session: Arc, +) -> Result>, anyhow::Error> { + debug!( + "Module: {}, Attempting to bind to: {}...", + bind_info.module_id.clone(), + port.port_name.clone() + ); + + let modules_mutex = store.modules.lock().await; + let module_config; + match modules_mutex.get(&bind_info.module_id) { + Some(remote_module_config) => { + module_config = remote_module_config.clone(); + } + None => { + error!("Unable to find Module ID: {} in the modules store. Won't bind UART port: {} for a non-existant module!", bind_info.module_id, bind_info.path); + return Err(anyhow!("Module not in store.")); + } + } + drop(modules_mutex); + + match module_config.connection { + ModuleConnection::UART(connection_config) => { + match SerialStream::open(&tokio_serial::new( + bind_info.path.clone(), + connection_config.baud, + )) { + Ok(bound_port) => { + let port_session = session.clone(); + let (port_read, port_write) = split(bound_port); + + store.port_statuses.uart.lock().await.insert( + bind_info.path.clone(), + PortStatus::Bound(bind_info.module_id.clone()), + ); + + debug!( + "Module: {}, Port: {}, bound!", + bind_info.module_id, port.port_name + ); + + let mut sub_handles = vec![]; + + let tx_port_ctx = (bind_info.module_id.clone(), port.port_name.clone()); + let tx_bind_info = bind_info.clone(); + let tx_session = port_session.clone(); + sub_handles.push(tokio::task::spawn(async move { + uart_tx_thread(tx_bind_info, tx_session, tx_port_ctx, port_write).await; + })); + + let rx_port_ctx = (bind_info.module_id.clone(), port.port_name.clone()); + let rx_session = port_session.clone(); + sub_handles.push(tokio::task::spawn(async move { + uart_rx_thread(rx_port_ctx, port_read, rx_session).await; + })); + + Ok(sub_handles) + } + Err(err) => { + return Err(err.into()); + } + } + } + _ => { + error!("UART port: {}, was requested for Module ID: {}, but the module connection configuration is not for UART. This is unnaceptable and a bug, please report!", bind_info.path, bind_info.module_id); + Err(anyhow!("Module Connection is not UART.")) + } + } +} + +#[instrument(skip(port_session, port_write))] +pub async fn uart_tx_thread( + tx_bind_info: PortToBind, + port_session: Arc, + tx_port_ctx: (String, String), + mut port_write: WriteHalf, +) { + let key_expr = format!( + "{MODULE_EVT_ID}/modules/by-id/{}/send", + tx_bind_info.module_id.clone() + ); + let (module_id, port_name) = tx_port_ctx; + + match port_session.declare_queryable(&key_expr).await { + Ok(queryable) => { + match queryable.recv_async().await { + Ok(query) => { + match query.payload() { + Some(payload) => { + match payload.try_to_string() { + Ok(payload_str) => { + debug!("Sending message: {payload_str}..."); - futures::future::join_all(sub_handles).await; - })); + // TODO: Encrypt + match rmp_serde::to_vec(&payload_str) { + Ok(msg_vec) => match port_write.write(msg_vec.as_slice()).await { + Ok(size) => match query.reply(key_expr, format!("{size}")).await { + Ok(_) => {} + Err(err) => todo!(), + }, + Err(err) => todo!(), + }, + Err(err) => todo!(), + } } - Err(_) => todo!(), + Err(err) => todo!(), } } - None => {} + None => todo!(), } - - std::mem::drop(config); } - - Ok(handles) + Err(err) => todo!(), } - Err(e) => Err(e.into()), } + Err(err) => todo!(), } +} - fn get_type() -> BusTypes { - BusTypes::UART +#[instrument(skip(port_read, port_session))] +pub async fn uart_rx_thread( + rx_port_ctx: (String, String), + mut port_read: ReadHalf, + port_session: Arc, +) { + let mut buffer = Vec::new(); + + loop { + match port_read.read(buffer.as_mut_slice()).await { + Ok(0) => { + warn!("Got EOF, was the module disconnected?"); + // TODO: Send tokio oneshot to reconnect to serial port. + break; + } + Ok(bytes) => {} + Err(err) => todo!(), + } } } diff --git a/clover-hub/src/server/modman/components/audio/impls.rs b/clover-hub/src/server/modman/components/audio/impls.rs index 404825a..86f9cb3 100644 --- a/clover-hub/src/server/modman/components/audio/impls.rs +++ b/clover-hub/src/server/modman/components/audio/impls.rs @@ -99,7 +99,7 @@ impl CloverComponentTrait for AudioInputComponent { } } } - super::models::ConnectionType::ModManProxy(proxied_connection) => todo!(), + super::models::ConnectionType::ModManProxy => todo!(), super::models::ConnectionType::Stream(streaming_connection) => todo!(), } @@ -195,7 +195,7 @@ impl CloverComponentTrait for AudioOutputComponent { } } // TODO: ModMan Proxied Audio Output device (like a GROVE speaker). - super::models::ConnectionType::ModManProxy(proxied_connection) => todo!(), + super::models::ConnectionType::ModManProxy => todo!(), super::models::ConnectionType::Stream(streaming_connection) => todo!(), } diff --git a/clover-hub/src/server/modman/components/audio/models.rs b/clover-hub/src/server/modman/components/audio/models.rs index bd01515..f6201e4 100644 --- a/clover-hub/src/server/modman/components/audio/models.rs +++ b/clover-hub/src/server/modman/components/audio/models.rs @@ -1,10 +1,7 @@ use std::sync::Arc; use crate::server::modman::{ - components::models::{ - ProxiedConnection, - StreamingConnection, - }, + components::models::StreamingConnection, models::GestureConfig, }; use serde::{ @@ -31,7 +28,7 @@ pub enum ConnectionType { Direct(DirectConnection), #[serde(rename = "modman-proxy")] #[strum(serialize = "modman-proxy")] - ModManProxy(ProxiedConnection), + ModManProxy, #[serde(rename = "stream")] #[strum(serialize = "stream")] Stream(StreamingConnection), diff --git a/clover-hub/src/server/modman/components/models.rs b/clover-hub/src/server/modman/components/models.rs index 1fee2ce..bb044e0 100644 --- a/clover-hub/src/server/modman/components/models.rs +++ b/clover-hub/src/server/modman/components/models.rs @@ -8,9 +8,15 @@ use std::sync::Arc; /// All components must implement this trait, ensures standardization between component types, etc. pub trait CloverComponentTrait: Sized { /// Should initalize the component in the store, and ensure that 2-way communication is setup. - async fn init(&mut self, store: Arc) -> Result<(), anyhow::Error>; + fn init( + &mut self, + store: Arc, + ) -> impl std::future::Future> + Send; /// Tells the component that it will not be used in the *near* future, and may even power it down. - async fn deinit(&mut self, store: Arc) -> Result<(), anyhow::Error>; + fn deinit( + &mut self, + store: Arc, + ) -> impl std::future::Future> + Send; } /// Known and supported streaming protocols for Video and Audio @@ -33,33 +39,3 @@ pub struct StreamingConnection { pub protocol: StreamProtocol, pub path: Option, } - -/// The Bus Proxy this component is connected through. -#[derive(Debug, Clone, Serialize, Deserialize)] -pub enum ProxiedConnection { - /// Device ID - Simulated(String), - /// App ID and Device ID - App(String, String), - /// Bus path and Device ID - #[cfg(feature = "can_fd")] - CANFD(String, String), - /// Bus path and Device ID - #[cfg(feature = "can_2")] - CAN2(String, String), - /// Device ID - #[cfg(feature = "bt_classic")] - BT(String), - /// Device ID - #[cfg(feature = "bt_le")] - BTLE(String), - /// Bus path and Device ID - #[cfg(feature = "spi")] - SPI(String, String), - /// Bus path and Device ID - #[cfg(feature = "i2c")] - I2C(String, String), - /// Bus path - #[cfg(feature = "uart")] - UART(String), -} diff --git a/clover-hub/src/server/modman/components/movement/models.rs b/clover-hub/src/server/modman/components/movement/models.rs index 60c1ad1..87d1253 100644 --- a/clover-hub/src/server/modman/components/movement/models.rs +++ b/clover-hub/src/server/modman/components/movement/models.rs @@ -1,9 +1,6 @@ -use crate::server::modman::{ - components::models::ProxiedConnection, - models::{ - GestureConfig, - GestureParameters, - }, +use crate::server::modman::models::{ + GestureConfig, + GestureParameters, }; use serde::{ Deserialize, @@ -15,7 +12,7 @@ use strum::VariantNames; pub enum ConnectionType { #[serde(rename = "modman-proxy")] #[strum(serialize = "modman-proxy")] - ModManProxy(ProxiedConnection), + ModManProxy, } #[derive(Debug, Clone, Serialize, Deserialize)] diff --git a/clover-hub/src/server/modman/components/sensors/models.rs b/clover-hub/src/server/modman/components/sensors/models.rs index 38e2205..504e603 100644 --- a/clover-hub/src/server/modman/components/sensors/models.rs +++ b/clover-hub/src/server/modman/components/sensors/models.rs @@ -1,7 +1,4 @@ -use crate::server::modman::{ - components::models::ProxiedConnection, - models::GestureConfig, -}; +use crate::server::modman::models::GestureConfig; use serde::{ Deserialize, Serialize, @@ -23,5 +20,5 @@ pub struct OutputSensorComponent { pub enum ConnectionType { #[serde(rename = "modman-proxy")] #[strum(serialize = "modman-proxy")] - ModManProxy(ProxiedConnection), + ModManProxy, } diff --git a/clover-hub/src/server/modman/components/video/cameras/models.rs b/clover-hub/src/server/modman/components/video/cameras/models.rs index a5a3905..3ba1e4a 100644 --- a/clover-hub/src/server/modman/components/video/cameras/models.rs +++ b/clover-hub/src/server/modman/components/video/cameras/models.rs @@ -1,8 +1,5 @@ use crate::server::modman::components::{ - models::{ - ProxiedConnection, - StreamingConnection, - }, + models::StreamingConnection, video::VideoResolution, }; use serde::{ @@ -19,7 +16,7 @@ pub enum ConnectionType { Video4Linux(String), #[serde(rename = "modman-proxy")] #[strum(serialize = "modman-proxy")] - ModManProxy(ProxiedConnection), + ModManProxy, #[serde(rename = "stream")] #[strum(serialize = "stream")] Stream(StreamingConnection), diff --git a/clover-hub/src/server/modman/components/video/displays/models.rs b/clover-hub/src/server/modman/components/video/displays/models.rs index ffc4ff7..715a2a3 100644 --- a/clover-hub/src/server/modman/components/video/displays/models.rs +++ b/clover-hub/src/server/modman/components/video/displays/models.rs @@ -1,9 +1,6 @@ use crate::server::modman::{ components::{ - models::{ - ProxiedConnection, - StreamingConnection, - }, + models::StreamingConnection, video::VideoResolution, }, models::GestureConfig, @@ -68,7 +65,7 @@ pub enum ConnectionType { Direct(DirectConnection), #[serde(rename = "modman-proxy")] #[strum(serialize = "modman-proxy")] - ModManProxy(ProxiedConnection), + ModManProxy, #[serde(rename = "stream")] #[strum(serialize = "stream")] Stream(StreamingConnection), diff --git a/clover-hub/src/server/modman/connections.rs b/clover-hub/src/server/modman/connections.rs new file mode 100644 index 0000000..30ee40b --- /dev/null +++ b/clover-hub/src/server/modman/connections.rs @@ -0,0 +1,69 @@ +use serde::{ + Deserialize, + Serialize, +}; + +/// The Bus Proxy this module is connected through. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub enum ModuleConnection { + /// Device ID + Simulated(String), + App(AppConnection), + /// Bus path and Device ID + #[cfg(feature = "can_fd")] + CANFD(CANFDConnection), + /// Bus path and Device ID + #[cfg(feature = "can_2")] + CAN2(CAN2Connection), + /// Device ID + #[cfg(feature = "bt_classic")] + BT(String), + /// Device ID + #[cfg(feature = "bt_le")] + BTLE(String), + /// Bus path and Device ID + #[cfg(feature = "spi")] + SPI(SPIConnection), + /// Bus path and Device ID + #[cfg(feature = "i2c")] + I2C(I2CConnection), + /// Bus path + #[cfg(feature = "uart")] + 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, + pub device_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CAN2Connection { + pub bus_id: String, + pub device_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SPIConnection { + pub bus_id: String, + pub device_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct I2CConnection { + pub bus_id: String, + pub device_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct UARTConnection { + pub port: String, + pub baud: u32, +} diff --git a/clover-hub/src/server/modman/ipc/displays.rs b/clover-hub/src/server/modman/ipc/displays.rs index cff8819..ec118bd 100644 --- a/clover-hub/src/server/modman/ipc/displays.rs +++ b/clover-hub/src/server/modman/ipc/displays.rs @@ -24,7 +24,7 @@ pub async fn display_queryable( session: Arc, cancellation_token: CancellationToken, ) { - let key_expr = format!("{MODULE_EVT_ID}/displays/get"); + let key_expr = format!("{MODULE_EVT_ID}/components/by-type/video/displays/all"); let queryable = session.declare_queryable(&key_expr).await.unwrap(); diff --git a/clover-hub/src/server/modman/mod.rs b/clover-hub/src/server/modman/mod.rs index e8921ac..7c9b413 100644 --- a/clover-hub/src/server/modman/mod.rs +++ b/clover-hub/src/server/modman/mod.rs @@ -8,6 +8,7 @@ pub mod busses; pub mod components; +pub mod connections; pub mod gestures; pub mod ipc; pub mod models; diff --git a/clover-hub/src/server/modman/models.rs b/clover-hub/src/server/modman/models.rs index fa1390e..c89f14a 100644 --- a/clover-hub/src/server/modman/models.rs +++ b/clover-hub/src/server/modman/models.rs @@ -4,23 +4,26 @@ //! use crate::server::{ - modman::components::{ - audio::models::{ - AudioInputComponent, - AudioOutputComponent, - }, - movement::models::MovementComponent, - sensors::models::{ - InputSensorComponent, - OutputSensorComponent, - }, - video::{ - cameras::models::CameraComponent, - displays::models::{ - PhysicalDisplayComponent, - VirtualDisplayComponent, + modman::{ + components::{ + audio::models::{ + AudioInputComponent, + AudioOutputComponent, + }, + movement::models::MovementComponent, + sensors::models::{ + InputSensorComponent, + OutputSensorComponent, + }, + video::{ + cameras::models::CameraComponent, + displays::models::{ + PhysicalDisplayComponent, + VirtualDisplayComponent, + }, }, }, + connections::ModuleConnection, }, warehouse::config::models::Config, }; @@ -52,6 +55,8 @@ pub struct Module { pub components: Vec<(String, bool)>, /// Either `com.reboot-codes.clover.hub` or the RFQDN of the app that manages this module. pub registered_by: String, + /// How is this module connected to modman? + pub connection: ModuleConnection, } impl Module { @@ -195,13 +200,13 @@ pub enum PortStatus { /// Available but unused. #[serde(rename = "available")] Available, - /// Requested by $COMPONENT_ID, but the UART bus isn't initalized yet + /// Requested by $MODULE_ID, but the UART bus isn't initalized yet #[serde(rename = "requested")] Requested(String), - /// Currently being used by $COMPONENT_ID + /// Currently being used by $MODULE_ID #[serde(rename = "bound")] Bound(String), - /// Unavailable, but still requested by $COMPONENT_ID + /// Unavailable, but still requested by $MODULE_ID #[serde(rename = "unavailable")] Unavailable(String), } @@ -267,7 +272,9 @@ pub struct GestureOverride { #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ModManConfig { - pub uart_ports: Vec<(String, u32)>, + /// All ports available for modman to use to connect to modules. + pub uart_ports: Vec, + /// Whether to restart paused gestures automatically on startup. pub restart_gestures: bool, pub gesture_states: HashMap, pub gestures_bg_by_default: bool, @@ -312,6 +319,9 @@ impl Default for ModManConfig { (external_display_id.clone(), true), ], registered_by: "com.reboot-codes.clover.modman.default".to_string(), + connection: ModuleConnection::Simulated( + "com.reboot-codes.clover.debug-display:0".to_string(), + ), }, ); diff --git a/clover-hub/src/server/modman/modules.rs b/clover-hub/src/server/modman/modules.rs index f52ba8f..9519739 100644 --- a/clover-hub/src/server/modman/modules.rs +++ b/clover-hub/src/server/modman/modules.rs @@ -222,6 +222,7 @@ pub async fn init_module(store: &ModManStore, id: String, module: Module) -> (bo initialized: true, components: module.components.clone(), registered_by: module.registered_by.clone(), + connection: module.connection.clone(), }, ); debug!("Module: {id}, In-memory store updated!"); @@ -398,6 +399,7 @@ pub async fn deinit_module(store: &ModManStore, id: String, module: Module) -> ( initialized: false, components: module.components.clone(), registered_by: module.registered_by.clone(), + connection: module.connection.clone(), }, ); diff --git a/clover-hub/src/server/renderer/mod.rs b/clover-hub/src/server/renderer/mod.rs index 3ed7a66..4b47447 100644 --- a/clover-hub/src/server/renderer/mod.rs +++ b/clover-hub/src/server/renderer/mod.rs @@ -174,7 +174,9 @@ pub async fn renderer_main( let displays_payload = loop { match init_session - .get(format!("{MODMAN_EVT_ID}/displays/get")) + .get(format!( + "{MODMAN_EVT_ID}/components/by-type/video/displays/all" + )) .await { Ok(reply_fifo) => match reply_fifo.recv_async().await { diff --git a/clover-hub/src/server/renderer/system_ui/systems/displays/mod.rs b/clover-hub/src/server/renderer/system_ui/systems/displays/mod.rs index 1a306af..d5921a8 100644 --- a/clover-hub/src/server/renderer/system_ui/systems/displays/mod.rs +++ b/clover-hub/src/server/renderer/system_ui/systems/displays/mod.rs @@ -1,11 +1,8 @@ +use crate::server::modman::components::video::displays::models::VirtualDisplayComponent; use crate::server::renderer::system_ui::{ AnyDisplayComponent, SystemUIIPC, }; -use crate::server::{ - modman::components::video::displays::models::VirtualDisplayComponent, - renderer::system_ui::systems::view_management::Composition, -}; use bevy::prelude::*; use bevy::render::camera::RenderTarget; #[cfg(feature = "compositor")] @@ -96,13 +93,13 @@ pub fn display_registrar( match physical_display_component.virtual_display { Some(vdisplay_id) => { for ( - queried_vdisplay_entity, - _queried_vdisplay_component, + queried_vdisplay_entity, + _queried_vdisplay_component, queried_vdisplay_id ) in &vdisplay_query { if vdisplay_id == queried_vdisplay_id.id { use bevy::window::WindowTheme; - + let window = commands.get_entity(queried_vdisplay_entity).unwrap().insert((Window { title: format!("Clover SystemUI: Virtual Display: {}", display_id.clone()), mode: windowed, @@ -110,7 +107,7 @@ pub fn display_registrar( position, ..Default::default() }, DisplayWindow { id: vdisplay_id.to_string() })).id(); - + commands.spawn(( Camera3d::default(), Camera { @@ -120,7 +117,7 @@ pub fn display_registrar( Transform::from_xyz(6.0, 0.0, 0.0).looking_at(Vec3::ZERO, Vec3::Y), DisplayCamera { id: vdisplay_id.to_string() } )); - + break; } } @@ -152,19 +149,22 @@ pub fn display_registrar( } } } - crate::server::modman::components::video::displays::models::ConnectionType::ModManProxy( - proxied_connection, - ) => {} + crate::server::modman::components::video::displays::models::ConnectionType::ModManProxy => todo!(), crate::server::modman::components::video::displays::models::ConnectionType::Stream( stream_config, - ) => {} + ) => todo!() } } AnyDisplayComponent::Virtual(virtual_display_component) => { debug!("Spawing VDisplay: {}", display_id.clone()); // TODO: Spawn a composition for this display. - commands.spawn((virtual_display_component, VirtualDisplayID { id: display_id.clone() })); + commands.spawn(( + virtual_display_component, + VirtualDisplayID { + id: display_id.clone(), + }, + )); debug!("Done spawning VDisplay: {}!", display_id.clone()); } diff --git a/clover-hub/src/server/warehouse/mod.rs b/clover-hub/src/server/warehouse/mod.rs index 6e60c8b..8b79baa 100644 --- a/clover-hub/src/server/warehouse/mod.rs +++ b/clover-hub/src/server/warehouse/mod.rs @@ -63,7 +63,7 @@ pub enum Error { /// 1. Ensures that the data directory exists, /// 2. Loads the core [configuration file](config) (paired management devices, permanently attached hardware, core Modules to initalize, etc), /// 3. and preps [Repository storage](repos). -#[instrument] +#[instrument(skip(store))] pub async fn setup_warehouse(data_dir: String, store: Arc) -> Result<(), Error> { let mut err = None; let mut data_dir_path = OsPath::new().join(data_dir.clone()); diff --git a/clover-hub/src/server/warehouse/repos/models.rs b/clover-hub/src/server/warehouse/repos/models.rs index d909da2..25c68ca 100644 --- a/clover-hub/src/server/warehouse/repos/models.rs +++ b/clover-hub/src/server/warehouse/repos/models.rs @@ -333,11 +333,11 @@ pub struct StaticGestureSpec { // TODO: Specify trait bounds (resolve async_fn_in_trait). pub trait ManifestCompilationFrom { /// Perform the compilation on the RAW manifest value type to get the COMPILED manifest value with its dependencies and directives resolved. Put the *parsed* (use [Deserialize]), *`Raw`* value specification in the `spec` parameter. - async fn compile( + fn compile( spec: T, resolution_ctx: ResolutionCtx, repo_dir_path: OsPath, - ) -> Result + ) -> impl std::future::Future> where Self: Sized, T: for<'a> Deserialize<'a>; diff --git a/listeners.md b/listeners.md new file mode 100644 index 0000000..4cb6831 --- /dev/null +++ b/listeners.md @@ -0,0 +1,38 @@ +# Zenoh Endpoints + +- com/reboot-codes/clover/ + - hub + - telemetry + - B(C1000):log-stream ! OTelTrace + - warehouse + - B(C1):status + - modman + - B(C1):status + - gestures + - Q:update Vec Result<(), GestureError> + - components + - @routes + - by-type + - video + - displays + - Q:all ! Vec + - @endpoints + - Q:config + - B(C1):state ComponentState + - modules + - @routes + - by-id + - $MODULE_ID + - by-rfqdn + - $MODULE_RFQDN + - $INIT_ORDER + - @endpoints + - Q:config + - Q:send BusMessage Result<(), BusError> + - B(C100?):recv BusMessage + - renderer + - B(C1):status + - inference_engine + - B(C1):status + - appdaemon + - B(C1):status -- 2.51.2