initial commit

This commit is contained in:
Franz Heinzmann (Frando) 2021-02-19 14:40:58 +01:00
commit 091e6319cd
9 changed files with 700 additions and 0 deletions

158
src/main.rs Normal file
View file

@ -0,0 +1,158 @@
use async_std::stream::{StreamExt};
use log::*;
use rosc::{OscMessage, OscPacket, OscType};
use studiox::switcher::{SwitcherError, SwitcherEvent, SwitcherHandle};
// use studiox::udp::UdpStream;
use studiox::osc::OscStream;
use studiox::spawn::FaustHandle;
pub enum Event {
FaustRx(Result<OscPacket, SwitcherError>),
Switcher(OscMessage),
CommandRx(Result<OscPacket, SwitcherError>),
}
#[async_std::main]
async fn main() -> anyhow::Result<()> {
env_logger::init();
// let child_task: JoinHandle<io::Result<()>> = task::spawn(async {
// let handle = FaustHandle::spawn().await;
// match handle {
// Ok(handle) => {
// info!("FAUST mixer launched: {:?}", handle)
// }
// Err(e) => error!("Launching FAUST mixer failed: {:?}", e),
// }
// eprintln!("spawn done");
// Ok(())
// // eprintln!("HANDLE: {:?}", handle);
// });
let handle = FaustHandle::spawn().await;
match handle {
Ok(handle) => {
info!("FAUST mixer launched: {:?}", handle)
}
Err(e) => error!("Launching FAUST mixer failed: {:?}", e),
}
let port_faust_tx = 5510;
let port_faust_rx = 5511;
let port_commands = 5590;
let addr_faust_tx = format!("localhost:{}", port_faust_tx);
let addr_faust_rx = format!("localhost:{}", port_faust_rx);
let addr_commands = format!("localhost:{}", port_commands);
let mut switcher = SwitcherHandle::run();
let switcher_rx = switcher.take_receiver().unwrap().map(Event::Switcher);
let faust_rx = OscStream::bind(addr_faust_rx).await?.map(Event::FaustRx);
let command_osc_rx = OscStream::bind(addr_commands).await?;
let command_osc_tx = command_osc_rx.sender();
let command_osc_rx = command_osc_rx.map(Event::CommandRx);
let mut events = faust_rx.merge(switcher_rx).merge(command_osc_rx);
info!("start main loop");
while let Some(event) = events.next().await {
match event {
Event::Switcher(message) => {
debug!("switcher tx: {:?}", message);
let packet = OscPacket::Message(message);
command_osc_tx.send_to(packet, &addr_faust_tx).await?;
}
Event::FaustRx(packet) => {
let packet = packet?;
debug!("faust rx: {:?}", packet);
}
Event::CommandRx(packet) => {
let packet = packet?;
debug!("command rx: {:?}", packet);
let event = parse_switcher_command(packet);
if let Some(event) = event {
switcher.send(event).await?;
}
}
}
}
switcher.join().await?;
// child_task.await?;
Ok(())
}
fn parse_switcher_command(packet: OscPacket) -> Option<SwitcherEvent> {
let event = match packet {
OscPacket::Message(message) => match (&message.addr as &str, &message.args[..]) {
("/switcher", [OscType::Int(channel_id), OscType::Int(on_value)]) => match on_value {
1 => Some(SwitcherEvent::Enable(*channel_id as usize)),
0 => Some(SwitcherEvent::Disable(*channel_id as usize)),
_ => {
error!(
"Invalid value for {}: {} (valid is 0 and 1)",
message.addr, on_value
);
None
}
},
_ => {
error!(
"Invalid address or arguments: {} {:?}",
message.addr, message.args
);
None
}
},
_ => {
error!("Bundles are not supported");
None
}
};
event
}
// match msg {
// Some(Incoming::Osc(packet, addr)) => match rosc::decoder::decode(&packet) {
// Ok(packet) => {
// trace!("RECV OSC from {}: {:?}", &addr, &packet);
// if let OscPacket::Message(message) = &packet {
// if message.addr.starts_with("/player") {
// async fn osc_stream(
// addr: &str,
// ) -> Result<
// Map<UdpStream, dyn FnMut((Vec<u8>, SocketAddr)) -> Result<OscMessage, SwitcherError> + 'static>,
// SwitcherError,
// >
// // where
// // S: Stream<Item = Result<OscPacket, rosc::OscError>>,
// {
// // Result<impl Stream<Result<OscPacket, OscError>> {
// let socket = UdpSocket::bind(addr).await?;
// let socket = Arc::new(socket);
// let stream = UdpStream::new(socket);
// let stream = stream.map(|packet| {
// packet
// .map_err(|e| SwitcherError::from(e))
// .map(|buf| rosc::decoder::decode(&buf.0[..]))
// .map_err(|e| SwitcherError::from(e))
// });
// Ok(stream)
// // Ok(stream)
// }
// let addr = message.addr;
// let args = message.args;
// let parts: Vec<&str> = message.addr.split("/").skip(1).collect();