diff --git a/.direnv/flake-profile b/.direnv/flake-profile index c0bc4eb..69f4eef 120000 --- a/.direnv/flake-profile +++ b/.direnv/flake-profile @@ -1 +1 @@ -flake-profile-59-link \ No newline at end of file +flake-profile-60-link \ No newline at end of file diff --git a/Cargo.lock b/Cargo.lock index 5b0032d..92639da 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1637,6 +1637,15 @@ dependencies = [ "serde", ] +[[package]] +name = "bincode" +version = "1.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b1f45e9417d87227c7a56d22e471c6206462cba514c7590c09aff4cf6d1ddcad" +dependencies = [ + "serde", +] + [[package]] name = "bindgen" version = "0.70.1" @@ -2260,6 +2269,7 @@ dependencies = [ "url", "uuid", "zenoh", + "zenoh-ext", ] [[package]] @@ -5422,6 +5432,12 @@ dependencies = [ "spin 0.9.8", ] +[[package]] +name = "leb128" +version = "0.2.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c83bff1d572d6b9aeef67ddfc8448e4a3737909cb28e81f97c791b9018703e52" + [[package]] name = "leb128fmt" version = "0.1.0" @@ -8764,6 +8780,15 @@ dependencies = [ "serde", ] +[[package]] +name = "serde_spanned" +version = "1.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6662b5879511e06e8999a8a235d848113e942c9124f211511b16466ee2995f26" +dependencies = [ + "serde_core", +] + [[package]] name = "serde_urlencoded" version = "0.7.1" @@ -9673,7 +9698,7 @@ dependencies = [ "cfg-expr", "heck 0.5.0", "pkg-config", - "toml", + "toml 0.8.22", "version-compare", ] @@ -10096,11 +10121,26 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "05ae329d1f08c4d17a59bed7ff5b5a769d062e64a62d34a3261b219e62cd5aae" dependencies = [ "serde", - "serde_spanned", - "toml_datetime", + "serde_spanned 0.6.8", + "toml_datetime 0.6.9", "toml_edit", ] +[[package]] +name = "toml" +version = "0.9.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ae2a4cf385da23d1d53bc15cdfa5c2109e93d8d362393c801e87da2f72f0e201" +dependencies = [ + "indexmap 2.9.0", + "serde_core", + "serde_spanned 1.1.1", + "toml_datetime 0.7.5+spec-1.1.0", + "toml_parser", + "toml_writer", + "winnow 0.7.10", +] + [[package]] name = "toml_datetime" version = "0.6.9" @@ -10110,6 +10150,15 @@ dependencies = [ "serde", ] +[[package]] +name = "toml_datetime" +version = "0.7.5+spec-1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "92e1cfed4a3038bc5a127e35a2d360f145e1f4b971b551a2ba5fd7aedf7e1347" +dependencies = [ + "serde_core", +] + [[package]] name = "toml_edit" version = "0.22.26" @@ -10118,11 +10167,26 @@ checksum = "310068873db2c5b3e7659d2cc35d21855dbafa50d1ce336397c666e3cb08137e" dependencies = [ "indexmap 2.9.0", "serde", - "serde_spanned", - "toml_datetime", - "winnow", + "serde_spanned 0.6.8", + "toml_datetime 0.6.9", + "winnow 0.7.10", +] + +[[package]] +name = "toml_parser" +version = "1.1.2+spec-1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a2abe9b86193656635d2411dc43050282ca48aa31c2451210f4202550afb7526" +dependencies = [ + "winnow 1.0.3", ] +[[package]] +name = "toml_writer" +version = "1.1.1+spec-1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "756daf9b1013ebe47a8776667b466417e2d4c5679d441c26230efd9ef78692db" + [[package]] name = "tower-service" version = "0.3.3" @@ -12208,6 +12272,12 @@ dependencies = [ "memchr", ] +[[package]] +name = "winnow" +version = "1.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0592e1c9d151f854e6fd382574c3a0855250e1d9b2f99d9281c6e6391af352f1" + [[package]] name = "wit-bindgen" version = "0.51.0" @@ -12487,7 +12557,7 @@ dependencies = [ "uds_windows", "uuid", "windows-sys 0.61.2", - "winnow", + "winnow 0.7.10", "zbus_macros", "zbus_names", "zvariant", @@ -12515,7 +12585,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ffd8af6d5b78619bab301ff3c560a5bd22426150253db278f164d6cf3b72c50f" dependencies = [ "serde", - "winnow", + "winnow 0.7.10", "zvariant", ] @@ -12626,6 +12696,7 @@ dependencies = [ "serde_json", "serde_with", "serde_yaml", + "toml 0.9.6", "tracing", "uhlc", "validated_struct", @@ -12663,6 +12734,26 @@ dependencies = [ "zenoh-result", ] +[[package]] +name = "zenoh-ext" +version = "1.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4e4de171bff459816b1ca9e6657c4ae9783d2c1fc569e1721a380521425e897d" +dependencies = [ + "async-trait", + "bincode", + "flume", + "futures", + "leb128", + "serde", + "tokio", + "tracing", + "uhlc", + "zenoh", + "zenoh-macros", + "zenoh-util", +] + [[package]] name = "zenoh-keyexpr" version = "1.9.0" @@ -13160,7 +13251,7 @@ dependencies = [ "endi", "enumflags2", "serde", - "winnow", + "winnow 0.7.10", "zvariant_derive", "zvariant_utils", ] @@ -13188,5 +13279,5 @@ dependencies = [ "quote", "serde", "syn 2.0.114", - "winnow", + "winnow 0.7.10", ] diff --git a/Cargo.toml b/Cargo.toml index 7e294f4..0a032a4 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -48,3 +48,4 @@ tracing-log = "0.2.0" opentelemetry_sdk = "0.32.0" opentelemetry-appender-tracing = "0.32.0" zenoh = "1.9.0" +zenoh-ext = { version = "1.9.0", features = ["unstable"] } diff --git a/clover-hub/Cargo.toml b/clover-hub/Cargo.toml index 16e5c97..7db1147 100644 --- a/clover-hub/Cargo.toml +++ b/clover-hub/Cargo.toml @@ -65,6 +65,7 @@ anyhow = "1.0.97" # HTTP/WS nexus = { path = "../../nexus" } zenoh = { workspace = true } +zenoh-ext = { workspace = true } url = "2.5.0" image = { version = "0.25.5", features = ["serde"] } diff --git a/clover-hub/src/server/appd/mod.rs b/clover-hub/src/server/appd/mod.rs index 62b3a24..4bc0cb9 100644 --- a/clover-hub/src/server/appd/mod.rs +++ b/clover-hub/src/server/appd/mod.rs @@ -7,10 +7,7 @@ pub mod docker; pub mod ipc; pub mod models; -use crate::{ - server::appd::ipc::handle_ipc, - utils::one_off_message, -}; +use crate::server::appd::ipc::handle_ipc; use self::docker::init_app; use bollard::{ @@ -19,11 +16,6 @@ use bollard::{ }; use docker::remove_app; use models::AppDStore; -use nexus::{ - arbiter::models::ApiKeyWithoutUID, - server::models::UserConfig, - user::NexusUser, -}; use std::sync::Arc; use tokio_util::sync::CancellationToken; use tracing::{ @@ -33,32 +25,18 @@ use tracing::{ instrument, warn, }; +use zenoh_ext::{ + AdvancedPublisherBuilderExt, + CacheConfig, +}; // TODO: Create application manifest schema/models pub const MODULE_EVT_ID: &str = "com/reboot-codes/clover/hub/appdaemon"; -pub async fn gen_user() -> UserConfig { - UserConfig { - user_type: "com.reboot-codes.com.clover.appd".to_string(), - pretty_name: "Clover: AppD".to_string(), - api_keys: vec![ApiKeyWithoutUID { - allowed_events_to: vec![ - "^nexus://com.reboot-codes.clover.appd(\\.(.*))*(\\/.*)*$".to_string() - ], - allowed_events_from: vec![ - "^nexus://com.reboot-codes.clover.appd(\\.(.*))*(\\/.*)*$".to_string() - ], - echo: false, - proxy: false, - }], - } -} - -#[instrument(skip(store, user, cancellation_tokens))] +#[instrument(skip(store, cancellation_tokens))] pub async fn appd_main( store: AppDStore, - user: NexusUser, cancellation_tokens: (CancellationToken, CancellationToken), ) { info!("Starting AppDaemon..."); @@ -73,11 +51,23 @@ pub async fn appd_main( let mut zenoh_config = zenoh::Config::default(); zenoh_config.insert_json5("connect/endpoints", "tcp/localhost:6699"); + zenoh_config + .insert_json5( + "timestamping/enabled", + r#"{ router: true, peer: true, client: true }"#, + ) + .unwrap(); debug!("Connecting to Zenoh..."); let session = Arc::new(zenoh::open(zenoh_config).await.unwrap()); debug!("Connected to Zenoh!"); + let status_publisher = session + .declare_publisher(format!("{MODULE_EVT_ID}/status")) + .cache(CacheConfig::default().max_samples(1)) + .await + .unwrap(); + let ipc_token = cancellation_tokens.0.clone(); let ipc_session = session.clone(); let ipc_handle = tokio::task::spawn(handle_ipc(ipc_token, ipc_session)); @@ -123,23 +113,21 @@ pub async fn appd_main( apps_initialized, init_apps.len() ); - one_off_message( - init_session.clone(), - &format!("{MODULE_EVT_ID}/status"), - "ready:incomplete", - ) - .await; + status_publisher + .put("ready:incomplete") + .await + .unwrap_or_else(|e| error!("Failed to publish status due to:\n{e}")); } else { if apps_initialized != 0 { info!("Initialized all {} apps!", apps_initialized); } - one_off_message( - init_session.clone(), - &format!("{MODULE_EVT_ID}/status"), - "ready", - ) - .await; + status_publisher + .put("ready") + .await + .unwrap_or_else(|e| error!("Failed to publish status due to:\n{e}")); } + + info!("AppDaemon Ready!"); }) .await; diff --git a/clover-hub/src/server/inference_engine/mod.rs b/clover-hub/src/server/inference_engine/mod.rs index 5744028..07d8264 100644 --- a/clover-hub/src/server/inference_engine/mod.rs +++ b/clover-hub/src/server/inference_engine/mod.rs @@ -22,33 +22,17 @@ use tracing::{ instrument, span, }; - -use crate::{ - server::inference_engine::ipc::handle_ipc, - utils::one_off_message, +use zenoh_ext::{ + AdvancedPublisherBuilderExt, + CacheConfig, }; +use crate::server::inference_engine::ipc::handle_ipc; + use super::warehouse::config::models::Config; pub const MODULE_EVT_ID: &str = "com/reboot-codes/clover/hub/inference_engine"; -pub async fn gen_user() -> UserConfig { - UserConfig { - user_type: "com.reboot-codes.com.clover.inference-engine".to_string(), - pretty_name: "Clover: Inference Engine".to_string(), - api_keys: vec![ApiKeyWithoutUID { - allowed_events_to: vec![ - "^nexus://com.reboot-codes.clover.inference-engine(\\.(.*))*(\\/.*)*$".to_string(), - ], - allowed_events_from: vec![ - "^nexus://com.reboot-codes.clover.inference-engine(\\.(.*))*(\\/.*)*$".to_string(), - ], - echo: false, - proxy: false, - }], - } -} - #[derive(Debug, Clone)] pub struct InferenceEngineStore { config: Arc>, @@ -65,10 +49,9 @@ impl InferenceEngineStore { } } -#[instrument(skip(inference_engine_store, user, cancellation_tokens))] +#[instrument(skip(inference_engine_store, cancellation_tokens))] pub async fn inference_engine_main( inference_engine_store: InferenceEngineStore, - user: NexusUser, cancellation_tokens: (CancellationToken, CancellationToken), ) { info!("Starting Inference Engine..."); @@ -76,11 +59,23 @@ pub async fn inference_engine_main( let mut zenoh_config = zenoh::Config::default(); zenoh_config.insert_json5("connect/endpoints", "tcp/localhost:6699"); + zenoh_config + .insert_json5( + "timestamping/enabled", + r#"{ router: true, peer: true, client: true }"#, + ) + .unwrap(); debug!("Connecting to Zenoh..."); let session = Arc::new(zenoh::open(zenoh_config).await.unwrap()); debug!("Connected to Zenoh!"); + let status_publisher = session + .declare_publisher(format!("{MODULE_EVT_ID}/status")) + .cache(CacheConfig::default().max_samples(1)) + .await + .unwrap(); + let ipc_token = cancellation_tokens.0.clone(); let ipc_session = session.clone(); let ipc_handle = tokio::task::spawn(handle_ipc(ipc_token, ipc_session)); @@ -88,7 +83,12 @@ pub async fn inference_engine_main( cancellation_tokens .0 .run_until_cancelled(async move { - one_off_message(session.clone(), &format!("{MODULE_EVT_ID}/status"), "ready").await; + status_publisher + .put("ready") + .await + .unwrap_or_else(|e| error!("Failed to publish status due to:\n{e}")); + + info!("InferenceEngine Ready!"); }) .await; diff --git a/clover-hub/src/server/mod.rs b/clover-hub/src/server/mod.rs index 099a543..e117fd9 100644 --- a/clover-hub/src/server/mod.rs +++ b/clover-hub/src/server/mod.rs @@ -132,42 +132,6 @@ pub async fn server_main( ); debug!("Created master Nexus user!"); - debug!("Creating Warehouse Nexus user..."); - let warehouse_user_config = Arc::new( - nexus_store - .add_user( - warehouse::gen_user().await, - Some(master_user_config.id.clone()), - None, - ) - .await - .unwrap(), - ); - debug!("Creating Warehouse NexusUser object..."); - let (warehouse_user, from_warehouse, to_warehouse) = nexus_store - .connect_user(&warehouse_user_config.api_keys[0].clone()) - .await - .unwrap(); - debug!("Created Warehouse Nexus user!"); - - debug!("Creating Renderer Nexus user..."); - let renderer_user_config = Arc::new( - nexus_store - .add_user( - renderer::gen_user().await, - Some(master_user_config.id.clone()), - None, - ) - .await - .unwrap(), - ); - debug!("Creating Renderer NexusUser object..."); - let (renderer_user, from_renderer, to_renderer) = nexus_store - .connect_user(&renderer_user_config.api_keys[0].clone()) - .await - .unwrap(); - debug!("Created Renderer Nexus user!"); - debug!("Creating Modman Nexus user..."); let modman_user_config = Arc::new( nexus_store @@ -186,42 +150,6 @@ pub async fn server_main( .unwrap(); debug!("Created Modman Nexus user!"); - debug!("Creating Inference Engine Nexus user..."); - let inference_engine_user_config = Arc::new( - nexus_store - .add_user( - inference_engine::gen_user().await, - Some(master_user_config.id.clone()), - None, - ) - .await - .unwrap(), - ); - debug!("Creating InferenceEngine NexusUser object..."); - let (inference_engine_user, from_inference_engine, to_inference_engine) = nexus_store - .connect_user(&inference_engine_user_config.api_keys[0].clone()) - .await - .unwrap(); - debug!("Created Inference Engine Nexus user!"); - - debug!("Creating AppDaemon Nexus user..."); - let appd_user_config = Arc::new( - nexus_store - .add_user( - appd::gen_user().await, - Some(master_user_config.id.clone()), - None, - ) - .await - .unwrap(), - ); - debug!("Creating AppDaemon NexusUser object..."); - let (appd_user, from_appd, to_appd) = nexus_store - .connect_user(&appd_user_config.api_keys[0].clone()) - .await - .unwrap(); - debug!("Created AppDaemon Nexus user!"); - // Create oneshot channel for ready signal let (nexus_ready_tx, nexus_ready_rx) = tokio::sync::oneshot::channel(); @@ -237,24 +165,8 @@ pub async fn server_main( nexus_listener( *nexus_port.to_owned(), nexus_store_clone, - vec![ - (warehouse_user_config.clone(), 0, to_warehouse), - (renderer_user_config.clone(), 0, to_renderer), - (modman_user_config.clone(), 0, to_modman), - (inference_engine_user_config.clone(), 0, to_inference_engine), - (appd_user_config.clone(), 0, to_appd), - ], - vec![ - (warehouse_user_config.clone(), 0, from_warehouse), - (renderer_user_config.clone(), 0, from_renderer), - (modman_user_config.clone(), 0, from_modman), - ( - inference_engine_user_config.clone(), - 0, - from_inference_engine, - ), - (appd_user_config.clone(), 0, from_appd), - ], + vec![(modman_user_config.clone(), 0, to_modman)], + vec![(modman_user_config.clone(), 0, from_modman)], nexus_tokens_clone, Some(nexus_ready_tx), ) @@ -284,12 +196,7 @@ pub async fn server_main( let warehouse_tokens = (CancellationToken::new(), CancellationToken::new()); let warehouse_tokens_clone = warehouse_tokens.clone(); let warehouse_handle = tokio::task::spawn(async move { - warehouse_main( - warehouse_store_clone.clone(), - warehouse_user, - warehouse_tokens_clone, - ) - .await; + warehouse_main(warehouse_store_clone.clone(), warehouse_tokens_clone).await; }); services.push(InitializedService { name: "Warehouse", @@ -315,7 +222,7 @@ pub async fn server_main( let renderer_tokens = (CancellationToken::new(), CancellationToken::new()); let renderer_tokens_clone = renderer_tokens.clone(); let renderer_handle = tokio::task::spawn(async move { - renderer_main(renderer_store, renderer_user, renderer_tokens_clone).await; + renderer_main(renderer_store, renderer_tokens_clone).await; }); services.push(InitializedService { name: "Renderer", @@ -328,12 +235,7 @@ pub async fn server_main( let inference_engine_tokens = (CancellationToken::new(), CancellationToken::new()); let inference_engine_tokens_clone = inference_engine_tokens.clone(); let inference_engine_handle = tokio::task::spawn(async move { - inference_engine_main( - inference_engine_store, - inference_engine_user, - inference_engine_tokens_clone, - ) - .await; + inference_engine_main(inference_engine_store, inference_engine_tokens_clone).await; }); services.push(InitializedService { name: "Inference Engine", @@ -346,7 +248,7 @@ pub async fn server_main( let appd_tokens = (CancellationToken::new(), CancellationToken::new()); let appd_tokens_clone = appd_tokens.clone(); let appd_handle = tokio::task::spawn(async move { - appd_main(appd_store, appd_user, appd_tokens_clone).await; + appd_main(appd_store, appd_tokens_clone).await; }); services.push(InitializedService { name: "AppDaemon", diff --git a/clover-hub/src/server/modman/mod.rs b/clover-hub/src/server/modman/mod.rs index f9dbb87..2a545ed 100644 --- a/clover-hub/src/server/modman/mod.rs +++ b/clover-hub/src/server/modman/mod.rs @@ -24,7 +24,10 @@ use nexus::{ server::models::UserConfig, user::NexusUser, }; -use std::sync::Arc; +use std::{ + collections::HashMap, + sync::Arc, +}; use tokio_util::sync::CancellationToken; use tracing::{ debug, @@ -33,10 +36,14 @@ use tracing::{ instrument, warn, }; +use zenoh_ext::{ + AdvancedPublisherBuilderExt, + CacheConfig, +}; -use crate::{ - server::modman::ipc::handle_ipc, - utils::one_off_message, +use crate::server::modman::{ + ipc::handle_ipc, + models::Module, }; pub const MODULE_EVT_ID: &str = "com/reboot-codes/clover/hub/modman"; @@ -71,11 +78,23 @@ pub async fn modman_main( let mut zenoh_config = zenoh::Config::default(); zenoh_config.insert_json5("connect/endpoints", "tcp/localhost:6699"); + zenoh_config + .insert_json5( + "timestamping/enabled", + r#"{ router: true, peer: true, client: true }"#, + ) + .unwrap(); debug!("Connecting to Zenoh..."); let session = Arc::new(zenoh::open(zenoh_config).await.unwrap()); debug!("Connected to Zenoh!"); + let status_publisher = session + .declare_publisher(format!("{MODULE_EVT_ID}/status")) + .cache(CacheConfig::default().max_samples(1)) + .await + .unwrap(); + let ipc_token = cancellation_tokens.0.clone(); let ipc_session = session.clone(); let ipc_store = store.clone(); @@ -93,15 +112,23 @@ pub async fn modman_main( } }); - let init_session = session.clone(); let init_store = Arc::new(store.clone()); - cancellation_tokens + let init_results = cancellation_tokens .0 .run_until_cancelled(async move { let config = init_store.config.lock().await; let static_modules = &config.modman.static_modules; let static_components = &config.modman.static_components; - let mut modules = init_store.modules.lock().await; + let mut modules_to_init: HashMap = { + let mut hashmap = HashMap::new(); + let modules = init_store.modules.lock().await; + + for (module_id, module_config) in modules.iter() { + hashmap.insert(module_id.clone(), module_config.clone()); + } + + hashmap + }; let mut modules_initalized: usize = 0; debug!("Checking for statically defined modules to init..."); @@ -111,7 +138,7 @@ pub async fn modman_main( static_modules.len() ); for (module_id, module) in static_modules { - modules.insert(module_id.clone(), module.clone()); + modules_to_init.insert(module_id.clone(), module.clone()); } } else { debug!("No statically defined modules to put into store, skipping!"); @@ -135,12 +162,14 @@ pub async fn modman_main( debug!("No statically defined components to put into store, skipping!"); } + let total_modules = modules_to_init.len(); + drop(config); info!("Initalizing modules..."); - if modules.len() > 0 { + if total_modules > 0 { // Initialize modules that were registered already via configuration and persistence. - for (id, module) in modules.iter() { + for (id, module) in modules_to_init.iter() { let (initialized, _components_initialized) = init_module(&init_store, id.clone(), module.clone()).await; @@ -152,65 +181,68 @@ pub async fn modman_main( info!("No static modules to initialize."); } - if modules_initalized != modules.len() { - warn!( - "Initalized {} out of {} module(s)!", - modules_initalized, - modules.len() - ); - one_off_message( - init_session.clone(), - &format!("{MODULE_EVT_ID}/status"), - "ready:incomplete", - ) - .await; - } else { - if modules_initalized != 0 { - info!("Initalized all {} module(s)", modules_initalized); - } - one_off_message( - init_session.clone(), - &format!("{MODULE_EVT_ID}/status"), - "ready", - ) - .await; - } + (modules_initalized, total_modules) }) .await; + if let Some((modules_initalized, total_modules)) = init_results { + if modules_initalized != total_modules { + warn!("Initalized {modules_initalized} out of {total_modules} module(s)!"); + status_publisher + .put("ready:incomplete") + .await + .unwrap_or_else(|e| error!("Failed to publish status due to:\n{e}")); + } else { + if modules_initalized != 0 { + info!("Initalized all {modules_initalized} module(s)"); + } + status_publisher + .put("ready") + .await + .unwrap_or_else(|e| error!("Failed to publish status due to:\n{e}")); + } + + info!("ModMan Ready!"); + } + let mod_clean_token = cancellation_tokens.0.clone(); tokio::select! { _ = mod_clean_token.cancelled() => { bus_handle.abort(); ipc_handle.abort(); + drop(status_publisher); info!("Cleaning up modules..."); // TODO: Add override cancellation token to force stop! - tokio::select! { - modules = store.modules.lock() => { - debug!("done waiting for lock"); - if modules.len() > 0 { - let mut modules_deinitalized: usize = 0; - - for (id, module) in modules.iter() { - if module.initialized { - let (de_initialized, _components_deinitialized) = deinit_module(&store, id.clone(), module.clone()).await; - - if de_initialized { modules_deinitalized += 1; } - } - } - - if modules_deinitalized != modules.len() { - warn!("Deinitalized {} out of {} module(s)!", modules_deinitalized, modules.len()); - } else { - info!("Deinitalized all {} module(s)", modules_deinitalized); - } - } else { - debug!("No modules to deinit."); - } + let modules_snapshot: Vec<(String, Module)> = { + let modules = store.modules.lock().await; + debug!("done waiting for lock"); + modules + .iter() + .filter(|(_, m)| m.initialized) + .map(|(id, m)| (id.clone(), m.clone())) + .collect() + // lock drops here + }; + + let total = modules_snapshot.len(); + let mut modules_deinitalized: usize = 0; + + if total > 0 { + for (id, module) in modules_snapshot { + let (de_initialized, _) = deinit_module(&store, id, module).await; + if de_initialized { modules_deinitalized += 1; } } + + if modules_deinitalized != total { + warn!("Deinitalized {} out of {} module(s)!", modules_deinitalized, total); + } else { + info!("Deinitalized all {} module(s)", modules_deinitalized); + } + } else { + debug!("No modules to deinit."); } std::mem::drop(store); diff --git a/clover-hub/src/server/modman/modules.rs b/clover-hub/src/server/modman/modules.rs index 12b9ce2..f52ba8f 100644 --- a/clover-hub/src/server/modman/modules.rs +++ b/clover-hub/src/server/modman/modules.rs @@ -374,10 +374,12 @@ pub async fn deinit_module(store: &ModManStore, id: String, module: Module) -> ( deinitialized_module_components, module.components.len() ); - initialized_module = true; + initialized_module = false; } else { error!("Module: {}, failed to deinitialize!", id.clone()); } + } else { + initialized_module = false; } } } diff --git a/clover-hub/src/server/renderer/mod.rs b/clover-hub/src/server/renderer/mod.rs index b51bb6b..3ed7a66 100644 --- a/clover-hub/src/server/renderer/mod.rs +++ b/clover-hub/src/server/renderer/mod.rs @@ -16,13 +16,7 @@ use crate::server::{ modman::MODULE_EVT_ID as MODMAN_EVT_ID, renderer::ipc::handle_ipc, }; -use crate::utils::one_off_message; use crate::utils::RecvSync; -use nexus::{ - arbiter::models::ApiKeyWithoutUID, - server::models::UserConfig, - user::NexusUser, -}; use queues::*; use std::sync::Arc; use std::sync::Mutex as StdMutex; @@ -35,29 +29,18 @@ use tracing::{ info, instrument, }; +use zenoh_ext::{ + AdvancedPublisherBuilderExt, + AdvancedSubscriberBuilderExt, + CacheConfig, + HistoryConfig, + RecoveryConfig, +}; use super::warehouse::config::models::Config; pub const MODULE_EVT_ID: &str = "com/reboot-codes/clover/hub/renderer"; -pub async fn gen_user() -> UserConfig { - UserConfig { - user_type: "com.reboot-codes.com.clover.renderer".to_string(), - pretty_name: "Clover: Renderer".to_string(), - api_keys: vec![ApiKeyWithoutUID { - allowed_events_to: vec![ - "^nexus://com.reboot-codes.clover.renderer(\\.(.*))*(\\/.*)*$".to_string(), - ], - allowed_events_from: vec![ - "^nexus://com.reboot-codes.clover.renderer(\\.(.*))*(\\/.*)*$".to_string(), - "^nexus://com.reboot-codes.clover.modman(\\.(.*))*(\\/.*)*$".to_string(), - ], - echo: false, - proxy: false, - }], - } -} - #[derive(Debug, Clone)] pub struct RendererStore { pub config: Arc>, @@ -74,10 +57,9 @@ impl RendererStore { } } -#[instrument(skip(store, user, cancellation_tokens))] +#[instrument(skip(store, cancellation_tokens))] pub async fn renderer_main( store: RendererStore, - user: NexusUser, cancellation_tokens: (CancellationToken, CancellationToken), ) { info!("Starting Renderer..."); @@ -85,11 +67,23 @@ pub async fn renderer_main( let mut zenoh_config = zenoh::Config::default(); zenoh_config.insert_json5("connect/endpoints", "tcp/localhost:6699"); + zenoh_config + .insert_json5( + "timestamping/enabled", + r#"{ router: true, peer: true, client: true }"#, + ) + .unwrap(); debug!("Connecting to Zenoh..."); let session = Arc::new(zenoh::open(zenoh_config).await.unwrap()); debug!("Connected to Zenoh!"); + let status_publisher = session + .declare_publisher(format!("{MODULE_EVT_ID}/status")) + .cache(CacheConfig::default().max_samples(1)) + .await + .unwrap(); + let ipc_token = cancellation_tokens.0.clone(); let ipc_session = session.clone(); let ipc_handle = tokio::task::spawn(handle_ipc(ipc_token, ipc_session)); @@ -112,62 +106,127 @@ pub async fn renderer_main( cancellation_tokens .0 .run_until_cancelled(async move { - info!("Requesting displays to setup in SystemUI from ModMan..."); + debug!("Waiting for ModMan to be ready..."); - match init_session - .get(format!("{MODMAN_EVT_ID}/displays/get")) + let modman_status_subscriber = init_session + .declare_subscriber(format!("{MODMAN_EVT_ID}/status")) + .history(HistoryConfig::default().detect_late_publishers()) + .recovery(RecoveryConfig::default().heartbeat()) .await - { - Ok(reply_fifo) => match reply_fifo.recv_async().await { - Ok(reply) => match reply.result() { - Ok(sample) => { - let payload = sample - .payload() - .try_to_string() - .unwrap_or_else(|e| e.to_string().into()); - - debug!("Got displays: {:?}", payload); - - one_off_message( - init_session.clone(), - &format!("{MODULE_EVT_ID}/status"), - "ready", - ) - .await; + .unwrap(); + + let modman_ready = loop { + match tokio::time::timeout( + std::time::Duration::from_millis(500), + modman_status_subscriber.recv_async(), + ) + .await + { + Ok(Ok(sample)) => { + let status = sample + .payload() + .try_to_string() + .unwrap_or_else(|e| e.to_string().into()); + debug!("ModMan Status: {status}"); + break status.to_string(); + } + Ok(Err(e)) => { + error!("Subscriber channel error: {e}"); + break "error".to_string(); + } + Err(_) => { + debug!("Timed out waiting for ModMan status, querying cache directly..."); + match init_session.get(format!("{MODMAN_EVT_ID}/status")).await { + Ok(replies) => { + if let Ok(reply) = replies.recv_async().await { + if let Ok(sample) = reply.into_result() { + let status = sample + .payload() + .try_to_string() + .unwrap_or_else(|e| e.to_string().into()); + debug!("ModMan Status (from cache query): {status}"); + break status.to_string(); + } + } + debug!("No cached ModMan status yet, retrying..."); + } + Err(e) => { + error!("Failed to query ModMan status: {e}"); + break "error".to_string(); + } } + } + } + }; + + drop(modman_status_subscriber); + + if modman_ready == "error" { + error!("Failed to get ModMan ready status!"); + status_publisher + .put("ready:incomplete") + .await + .unwrap_or_else(|e| error!("Failed to publish status due to:\n{e}")); + return; + } + + info!("Requesting displays to setup in SystemUI from ModMan..."); + + let displays_payload = loop { + match init_session + .get(format!("{MODMAN_EVT_ID}/displays/get")) + .await + { + Ok(reply_fifo) => match reply_fifo.recv_async().await { + Ok(reply) => match reply.result() { + Ok(sample) => { + break Ok( + sample + .payload() + .try_to_string() + .unwrap_or_else(|e| e.to_string().into()) + .to_string(), + ); + } + Err(e) => { + break Err( + e.payload() + .try_to_string() + .unwrap_or_else(|e| e.to_string().into()) + .to_string(), + ); + } + }, Err(e) => { - let payload = e - .payload() - .try_to_string() - .unwrap_or_else(|e| e.to_string().into()); - error!(">> Received (ERROR: '{payload}')"); - - one_off_message( - init_session.clone(), - &format!("{MODULE_EVT_ID}/status"), - "ready:incomplete", - ) - .await; + debug!("displays/get queryable not ready yet ({e}), retrying in 100ms..."); + tokio::time::sleep(std::time::Duration::from_millis(100)).await; } }, - Err(_) => { - one_off_message( - init_session.clone(), - &format!("{MODULE_EVT_ID}/status"), - "ready:incomplete", - ) - .await; + Err(e) => { + error!("Unable to setup Querier:\n{e}"); + break Err(e.to_string()); } - }, - Err(e) => { - one_off_message( - init_session.clone(), - &format!("{MODULE_EVT_ID}/status"), - "ready:incomplete", - ) - .await; + } + }; + + match displays_payload { + Ok(payload) => { + debug!("Got displays: {:?}", payload); + status_publisher + .put("ready") + .await + .unwrap_or_else(|e| error!("Failed to publish status due to:\n{e}")); + } + Err(payload) => { + error!("Got error payload:\n{payload}"); + status_publisher + .put("ready:incomplete") + .await + .unwrap_or_else(|e| error!("Failed to publish status due to:\n{e}")); } } + + info!("Renderer Ready!"); }) .await; diff --git a/clover-hub/src/server/renderer/system_ui/mod.rs b/clover-hub/src/server/renderer/system_ui/mod.rs index 51052ab..78bbb5d 100644 --- a/clover-hub/src/server/renderer/system_ui/mod.rs +++ b/clover-hub/src/server/renderer/system_ui/mod.rs @@ -120,7 +120,7 @@ pub fn system_ui_main(custom_bevy_ipc: SystemUIIPC, disable_winit: Option) .set(WindowPlugin { primary_window: Some(Window { title: "Clover".into(), - resolution: (200.0, 500.0).into(), + resolution: (500.0, 200.0).into(), present_mode: PresentMode::AutoVsync, // Tells Wasm to resize the window according to the available canvas fit_canvas_to_parent: true, diff --git a/clover-hub/src/server/warehouse/mod.rs b/clover-hub/src/server/warehouse/mod.rs index c95b245..55c3259 100644 --- a/clover-hub/src/server/warehouse/mod.rs +++ b/clover-hub/src/server/warehouse/mod.rs @@ -40,9 +40,12 @@ use tracing::{ info, instrument, }; +use zenoh_ext::{ + AdvancedPublisherBuilderExt, + CacheConfig, +}; use crate::server::warehouse::ipc::handle_ipc; -use crate::utils::one_off_message; /// The primary startup Error enum. /// Warehouse will return this enum in [`setup_warehouse`]. @@ -243,10 +246,9 @@ pub async fn gen_user() -> UserConfig { pub const MODULE_EVT_ID: &str = "com/reboot-codes/clover/hub/warehouse"; /// Main service function for Warehouse. Maintains an ongoing connection to Zenoh, and will manage filesystem operations as needed. -#[instrument(skip(store, user, cancellation_tokens))] +#[instrument(skip(store, cancellation_tokens))] pub async fn warehouse_main( store: Arc, - user: NexusUser, cancellation_tokens: (CancellationToken, CancellationToken), ) { info!("Starting Warehouse..."); @@ -254,11 +256,23 @@ pub async fn warehouse_main( let mut zenoh_config = zenoh::Config::default(); zenoh_config.insert_json5("connect/endpoints", "tcp/localhost:6699"); + zenoh_config + .insert_json5( + "timestamping/enabled", + r#"{ router: true, peer: true, client: true }"#, + ) + .unwrap(); debug!("Connecting to Zenoh..."); let session = Arc::new(zenoh::open(zenoh_config).await.unwrap()); debug!("Connected to Zenoh!"); + let status_publisher = session + .declare_publisher(format!("{MODULE_EVT_ID}/status")) + .cache(CacheConfig::default().max_samples(1)) + .await + .unwrap(); + let ipc_token = cancellation_tokens.0.clone(); let ipc_session = session.clone(); let ipc_handle = tokio::task::spawn(handle_ipc(ipc_token, ipc_session)); @@ -271,7 +285,6 @@ pub async fn warehouse_main( .await; let init_store = Arc::new(store.clone()); - let init_user = Arc::new(user.clone()); let init_tokens = cancellation_tokens.clone(); cancellation_tokens .0 @@ -286,12 +299,12 @@ pub async fn warehouse_main( } } - one_off_message( - session.clone(), - "com/reboot-codes/clover/server/warehouse/status", - "ready", - ) - .await; + status_publisher + .put("ready") + .await + .unwrap_or_else(|e| error!("Failed to publish status due to:\n{e}")); + + info!("Warehouse Ready!"); }) .await; diff --git a/clover-hub/src/utils.rs b/clover-hub/src/utils.rs index f4dea1a..b0e17bd 100644 --- a/clover-hub/src/utils.rs +++ b/clover-hub/src/utils.rs @@ -24,6 +24,7 @@ use tokio::{ }; use tokio_util::sync::CancellationToken; use tracing::instrument; +use zenoh_ext::AdvancedPublisher; pub struct RecvSync(pub std::sync::mpsc::Receiver); @@ -113,21 +114,3 @@ where }, } } - -/// One time publisher to a topic over a Zenoh session. For any longer term publishing, setup a publisher manually! -#[instrument(skip(session))] -pub async fn one_off_message(session: Arc, topic: &str, message: &str) { - let publisher = session.declare_publisher(topic).await.unwrap(); - - match publisher.put(message).await { - Ok(_) => { - debug!("Sucessfully sent message to topic: {}", topic); - } - Err(_) => { - error!( - "Failed to send message to zenoh topic: {}, this may be a sign of a misconfiguration!!", - topic - ); - } - } -} diff --git a/flake.nix b/flake.nix index 9b0bb2c..b1c9d2a 100644 --- a/flake.nix +++ b/flake.nix @@ -166,7 +166,8 @@ exec zenohd -l tcp/0.0.0.0:6699 --adminspace-permissions rw \ --cfg='adminspace/enabled:true' \ --cfg='adminspace/permissions/read:true' \ - --cfg='adminspace/permissions/write:true' + --cfg='adminspace/permissions/write:true' \ + --cfg='timestamping/enabled/client:true' ''; in craneLib.devShell {