jevstrudel.git / src-tauri / src / oscbridge.rs
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}