diff --git a/Cargo.lock b/Cargo.lock index fef60b0..0bc3f85 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2104,6 +2104,7 @@ dependencies = [ "bluer", "bollard", "can-iso-tp", + "can-isotp-interface", "chrono", "clap", "clover-hub-macros", diff --git a/clover-hub/Cargo.toml b/clover-hub/Cargo.toml index 629b4d2..eb250d7 100644 --- a/clover-hub/Cargo.toml +++ b/clover-hub/Cargo.toml @@ -107,3 +107,4 @@ rodio = "0.20.1" serde_bytes = { workspace = true } strum_macros = "0.28.0" embedded-can = "0.4.1" +can-isotp-interface = "0.1.0" 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 cb661a8..4f8d4ae 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 @@ -7,7 +7,10 @@ use anyhow::anyhow; use embedded_can::Id; use linux_socketcan_iso_tp::{ self, + flags, + IsoTpFlowControlOptions, IsoTpKernelOptions, + IsoTpSocketOptions, TokioSocketCanIsoTp, }; use regex::Regex; @@ -22,7 +25,10 @@ use tracing::{ use crate::server::modman::{ busses::proxies::group::can_2::{ - module_listener::can_bus_listener, + module_listener::{ + can_module_rx, + can_module_tx, + }, CAN2Bus, }, models::PortStatus, @@ -141,7 +147,14 @@ pub async fn setup_listener( match parse_id_str(id_tuple.0) { Ok(raw_rx_id) => match parse_id_str(id_tuple.1) { Ok(raw_tx_id) => { - let options = IsoTpKernelOptions::default(); + let rx_options = IsoTpKernelOptions::default(); + let tx_options = IsoTpKernelOptions { + socket: IsoTpSocketOptions { + flags: flags::CAN_ISOTP_LISTEN_MODE, + ..Default::default() + }, + ..Default::default() + }; match embedded_can::StandardId::new(raw_rx_id).ok_or(anyhow!( "We expect that an RX ID of {}, is valid. Check your config or there's a bug in manifest validation!", @@ -157,16 +170,36 @@ pub async fn setup_listener( &iface_name, Id::Standard(rx_id), Id::Standard(tx_id), - &options, + &rx_options, ) { - Ok(socket) => { - let listener_token = CancellationToken::new(); - - listener_registry.lock().await.insert(module_id.clone(), listener_token.clone()); - - tokio::task::spawn(async move { - can_bus_listener(ctx.clone(), listener_token.clone(), module_id, socket).await - }); + Ok(rx_socket) => { + match TokioSocketCanIsoTp::open( + &iface_name, + Id::Standard(rx_id), + Id::Standard(tx_id), + &tx_options + ) { + Ok(tx_socket) => { + let listener_token = CancellationToken::new(); + + listener_registry.lock().await.insert(module_id.clone(), listener_token.clone()); + + let rx_session = ctx.session.clone(); + let rx_token = listener_token.clone(); + let rx_id = module_id.clone(); + tokio::task::spawn(async move { + can_module_rx(rx_session, rx_token, rx_socket, rx_id).await; + }); + + let tx_session = ctx.session.clone(); + let tx_token = listener_token.clone(); + let tx_id = module_id.clone(); + tokio::task::spawn(async move { + can_module_tx(tx_session, tx_token, tx_socket, tx_id).await; + }); + }, + Err(err) => todo!(), + } }, Err(err) => { ret = Some(err.into()); 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 c9fc410..5ed32e4 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 @@ -1,52 +1,121 @@ -use std::sync::Arc; +use std::{ + sync::Arc, + time::Duration, +}; +use can_isotp_interface::{ + IsoTpAsyncEndpoint, + IsoTpEndpoint, + RecvControl, + RecvError, + RecvStatus, +}; use linux_socketcan_iso_tp::TokioSocketCanIsoTp; use tokio_util::sync::CancellationToken; -use tracing::instrument; +use tracing::{ + debug, + error, + instrument, + warn, +}; +use zenoh_ext::{ + AdvancedPublisherBuilderExt, + CacheConfig, +}; -use crate::server::modman::busses::proxies::group::can_2::CAN2Bus; - -// JK You thought this function actually did something lmaoooooo -#[instrument(skip(cancellation_token, raw_socket))] -pub async fn can_bus_listener( - ctx: Arc, - cancellation_token: CancellationToken, - module_id: String, - raw_socket: TokioSocketCanIsoTp, -) { - let socket = Arc::new(raw_socket); - - let rx_session = ctx.session.clone(); - let rx_token = cancellation_token.clone(); - let rx_socket = socket.clone(); - let rx_id = module_id.clone(); - tokio::task::spawn(async move { - can_module_rx(rx_session, rx_token, rx_socket, rx_id).await; - }); - - let tx_session = ctx.session.clone(); - let tx_token = cancellation_token.clone(); - let tx_socket = socket.clone(); - let tx_id = module_id.clone(); - tokio::task::spawn(async move { - can_module_tx(tx_session, tx_token, tx_socket, tx_id).await; - }); -} +use crate::server::modman::{ + busses::models::BusMessage, + MODULE_EVT_ID, +}; #[instrument(skip(session, cancellation_token, socket))] pub async fn can_module_rx( session: Arc, cancellation_token: CancellationToken, - socket: Arc, + mut socket: TokioSocketCanIsoTp, module_id: String, ) { + let key_expr = format!("{MODULE_EVT_ID}/modules/by-id/{}/recv", &module_id); + + match session + .declare_publisher(&key_expr) + .cache(CacheConfig::default().max_samples(1)) + .await + { + Ok(publisher) => {} + Err(err) => todo!(), + }; } #[instrument(skip(session, cancellation_token, socket))] pub async fn can_module_tx( session: Arc, cancellation_token: CancellationToken, - socket: Arc, + mut socket: TokioSocketCanIsoTp, module_id: String, ) { + let key_expr = format!("{MODULE_EVT_ID}/modules/by-id/{}/recv", &module_id); + + match session.declare_queryable(&key_expr).await { + Ok(queryable) => { + while !cancellation_token.is_cancelled() { + if let Ok(query) = queryable.recv_async().await { + match query.payload() { + Some(query_payload) => match query_payload.try_to_string() { + Ok(payload_str) => { + match serde_json_lenient::from_str::(&payload_str.to_owned()) { + Ok(message) => match rmp_serde::to_vec(&message) { + Ok(msg_bytes) => { + debug!( + "Dumping everything in the recv buffer so we don't get a memory leak." + ); + loop { + match socket + .recv_one(Duration::ZERO, |_meta, _payload| { + // just discard whatever showed up + Ok(RecvControl::Continue) + }) + .await + { + Ok(RecvStatus::DeliveredOne) => continue, + Ok(RecvStatus::TimedOut) => break, + Err(RecvError::BufferTooSmall { needed, got }) => { + error!( + needed, + got, "drain buffer undersized, this is a bug and should be reported!" + ); + break; + } + Err(RecvError::Backend(e)) => { + warn!(?e, "Isotp drain failed, socket likely dead!!"); + break; + } + } + } + + debug!("Sending CAN 2 message to module: {module_id}..."); + match socket + .send_to(0, &msg_bytes, Duration::from_millis(100)) + .await + { + Ok(_) => { + debug!("Successfully sent CAN 2 message to module: {module_id}!"); + } + Err(err) => todo!(), + } + } + Err(err) => todo!(), + }, + Err(err) => todo!(), + } + } + Err(err) => todo!(), + }, + None => {} + } + } + } + } + Err(err) => todo!(), + }; }