1use rosc::{encoder, OscTime}; 2use rosc::{OscBundle, OscMessage, OscPacket, OscType}; 3 4use std::net::UdpSocket; 5 6use serde::Deserialize; 7use std::sync::Arc; 8use std::thread::sleep; 9use std::time::{Duration}; 10use tokio::sync::{mpsc, Mutex}; 11 12use crate::loggerbridge::Logger; 13pub struct OscMsg { 14 pub msg_buf: Vec<u8>, 15 pub timestamp: f64, 16} 17 18pub struct AsyncInputTransmit { 19 pub inner: Mutex<mpsc::Sender<Vec<OscMsg>>>, 20} 21 22const UNIX_OFFSET: u64 = 2_208_988_800; // 70 years in seconds 23const TWO_POW_32: f64 = (u32::MAX as f64) + 1.0; // Number of bits in a `u32` 24const NANOS_PER_SECOND: f64 = 1.0e9; 25const SECONDS_PER_NANO: f64 = 1.0 / NANOS_PER_SECOND; 26 27pub fn init( 28 logger: Logger, 29 async_input_receiver: mpsc::Receiver<Vec<OscMsg>>, 30 mut async_output_receiver: mpsc::Receiver<Vec<OscMsg>>, 31 async_output_transmitter: mpsc::Sender<Vec<OscMsg>>, 32) { 33 tauri::async_runtime::spawn(async move { 34 async_process_model(async_input_receiver, async_output_transmitter).await 35 }); 36 let message_queue: Arc<Mutex<Vec<OscMsg>>> = Arc::new(Mutex::new(Vec::new())); 37 /* ........................................................... 38 Listen For incoming messages and add to queue 39 ............................................................*/ 40 let message_queue_clone = Arc::clone(&message_queue); 41 tauri::async_runtime::spawn(async move { 42 loop { 43 if let Some(package) = async_output_receiver.recv().await { 44 let mut message_queue = message_queue_clone.lock().await; 45 let messages = package; 46 for message in messages { 47 (*message_queue).push(message); 48 } 49 } 50 } 51 }); 52 53 let message_queue_clone = Arc::clone(&message_queue); 54 tauri::async_runtime::spawn(async move { 55 /* ........................................................... 56 Open OSC Ports 57 ............................................................*/ 58 let sock = UdpSocket::bind("127.0.0.1:57122").unwrap(); 59 let to_addr = String::from("127.0.0.1:57120"); 60 sock.set_nonblocking(true).unwrap(); 61 sock.connect(to_addr) 62 .expect("could not connect to OSC address"); 63 64 /* ........................................................... 65 Process queued messages 66 ............................................................*/ 67 68 loop { 69 let mut message_queue = message_queue_clone.lock().await; 70 71 message_queue.retain(|message| { 72 let result = sock.send(&message.msg_buf); 73 if result.is_err() { 74 logger.log( 75 format!( 76 "OSC Message failed to send, the server might no longer be available" 77 ), 78 "error".to_string(), 79 ); 80 } 81 return false; 82 }); 83 84 sleep(Duration::from_millis(1)); 85 } 86 }); 87} 88 89pub async fn async_process_model( 90 mut input_reciever: mpsc::Receiver<Vec<OscMsg>>, 91 output_transmitter: mpsc::Sender<Vec<OscMsg>>, 92) -> Result<(), Box<dyn std::error::Error + Send + Sync>> { 93 while let Some(input) = input_reciever.recv().await { 94 let output = input; 95 output_transmitter.send(output).await?; 96 } 97 Ok(()) 98} 99 100#[derive(Deserialize)] 101pub struct Param { 102 name: String, 103 value: String, 104 valueisnumber: bool, 105} 106#[derive(Deserialize)] 107pub struct MessageFromJS { 108 params: Vec<Param>, 109 timestamp: f64, 110 target: String, 111}
Called from JS
113#[tauri::command] 114pub async fn sendosc( 115 messagesfromjs: Vec<MessageFromJS>, 116 state: tauri::State<'_, AsyncInputTransmit>, 117) -> Result<(), String> { 118 let async_proc_input_tx = state.inner.lock().await; 119 let mut messages_to_process: Vec<OscMsg> = Vec::new(); 120 for m in messagesfromjs { 121 let mut args = Vec::new(); 122 for p in m.params { 123 args.push(OscType::String(p.name)); 124 if p.valueisnumber { 125 args.push(OscType::Float(p.value.parse().unwrap())); 126 } else { 127 args.push(OscType::String(p.value)); 128 } 129 } 130 // let start = SystemTime::now().duration_since(UNIX_EPOCH).unwrap(); 131 132 let time_delay = Duration::from_secs_f64(m.timestamp); 133 let duration_since_epoch = 134 time_delay + Duration::new(UNIX_OFFSET, 0); 135 136 let seconds = u32::try_from(duration_since_epoch.as_secs()) 137 .map_err(|_| "bit conversion failed for osc message timetag")?; 138 139 let nanos = duration_since_epoch.subsec_nanos() as f64; 140 let fractional = (nanos * SECONDS_PER_NANO * TWO_POW_32).round() as u32; 141 142 let timetag = OscTime::from((seconds, fractional)); 143 144 let packet = OscPacket::Message(OscMessage { 145 addr: m.target, 146 args, 147 }); 148 149 let bundle = OscBundle { 150 content: vec![packet], 151 timetag, 152 }; 153 154 let msg_buf = encoder::encode(&OscPacket::Bundle(bundle)).unwrap(); 155 156 let message_to_process = OscMsg { 157 msg_buf, 158 timestamp: m.timestamp, 159 }; 160 messages_to_process.push(message_to_process); 161 } 162 163 async_proc_input_tx 164 .send(messages_to_process) 165 .await 166 .map_err(|e| e.to_string()) 167}