use crate::tokiort::TokioIo;
use crate::OSCNode;
use hyper::server::conn::http1;
use hyper::service::Service;
use hyper::{body::Incoming as IncomingBody, Request, Response};
use std::any::Any;
use std::future::Future;
use std::net::SocketAddr;
use std::pin::Pin;
use std::sync::Arc;
use std::thread::spawn;
use std::time::Duration;
use tokio::net::TcpListener;
use tokio::runtime::Runtime;
use zeroconf::prelude::*;
struct OscQueryStatic {
root: Arc<OSCNode>,
}
impl Service<Request<IncomingBody>> for OscQueryStatic {
type Response = Response<String>;
type Error = hyper::Error;
type Future = Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send>>;
fn call(&self, req: Request<IncomingBody>) -> Self::Future {
fn mk_response(s: String) -> Result<Response<String>, hyper::Error> {
println!("{}", s);
Ok(Response::builder()
.header("Content-Type", "application/json")
.body(s)
.unwrap())
}
println!("{:?} {:?}", req.uri(), req.method());
if let Ok(node) = self.root.get(req.uri().path().to_string()) {
if let Some(query) = req.uri().query() {
let res = match query {
"HOST_INFO" => mk_response(
serde_json::to_value(node)
.unwrap()
.get("HOST_INFO")
.unwrap()
.to_string(),
),
"VALUE" => mk_response(format!(
"{{\"VALUE\":{}}}",
serde_json::to_value(node).unwrap().get("VALUE").unwrap()
)),
"TYPE" => mk_response(
serde_json::to_value(node)
.unwrap()
.get("TYPE")
.unwrap()
.to_string(),
),
_ => Ok(Response::builder()
.status(204)
.body("not supported".to_string())
.unwrap()),
};
return Box::pin(async { res });
} else {
let res = mk_response(serde_json::to_string(node).unwrap());
return Box::pin(async { res });
}
}
let res = Ok(Response::builder()
.status(404)
.body("Not Found".to_string())
.unwrap());
Box::pin(async { res })
}
}
fn on_service_registered(
result: zeroconf::Result<zeroconf::ServiceRegistration>,
_: Option<Arc<dyn Any>>,
) {
let service = result.unwrap();
println!("Service registered: {:?}", service);
}
pub async fn run_oscquery_service(
root: OSCNode,
address: SocketAddr,
) -> tokio::io::Result<(tokio::task::JoinHandle<()>, tokio::task::JoinHandle<()>)> {
let arc_root = Arc::new(root);
println!("oscq_rs start tcp at {:?}", address);
let listener = TcpListener::bind(address).await?;
println!("oscq_rs started tcp at {:?}", address);
let handle = tokio::task::spawn(async move {
loop {
println!("oscq_rs wait for connection {:?}", address);
let (stream, con) = listener.accept().await.unwrap();
println!("oscq_rs serve connection {:?}", con);
let service = OscQueryStatic {
root: arc_root.clone(),
};
let io = TokioIo::new(stream);
tokio::task::spawn(async move {
println!("oscq_rs serve connection async {:?}", con);
if let Err(err) = http1::Builder::new()
.keep_alive(true)
.serve_connection(io, service)
.await
{
println!("Failed to serve connection: {:?}", err);
}
});
}
});
let handle1 = tokio::task::spawn(async move {
let mut service = zeroconf::MdnsService::new(
zeroconf::ServiceType::new("oscjson", "tcp").unwrap(),
address.port(),
);
service.set_name("oscq_rs");
service.set_registered_callback(Box::new(on_service_registered));
let event_loop = service.register().unwrap();
loop {
event_loop.poll(Duration::from_secs(10)).unwrap();
}
});
Ok((handle, handle1))
}
pub fn spawn_oscquery_service(root: OSCNode, address: SocketAddr) {
spawn(move || {
let rt = Runtime::new().unwrap();
rt.block_on(async move {
let (x, y) = run_oscquery_service(root, address).await.unwrap();
let res = tokio::join!(x, y);
res.0.unwrap();
res.1.unwrap();
});
loop {
panic!("oscQueryServer Stopped");
}
});
}
#[tokio::test]
async fn test_service() {
use crate::{OSCAccess, OSCUnit, OscHostInfo, OscQueryParameter};
use rosc::OscType;
use std::net::SocketAddr;
let info = OscHostInfo::new("OSCQuery Test".to_string(), "127.0.0.1".to_string(), 6668)
.with_ext_access()
.with_ext_unit()
.with_ext_description()
.with_ext_range();
let mut root = OSCNode::root(Some(Box::new(info)));
let par1 = OscQueryParameter::new("/group/test".to_string(), OscType::Float(1f32))
.with_description("My First Description".to_string())
.with_min_max(0f32, 10f32)
.with_access(OSCAccess::ReadWrite)
.with_unit(OSCUnit::Distance(crate::OSCDistance::Centimeter));
let par2 = OscQueryParameter::new("/group/test2".to_string(), OscType::Float(1f32))
.with_description("My First Description".to_string())
.with_min_max(0f32, 10f32)
.with_access(OSCAccess::ReadWrite)
.with_unit(OSCUnit::Distance(crate::OSCDistance::Meter));
let par3 = OscQueryParameter::new("/group/test/subtest".to_string(), OscType::Float(1f32))
.with_description("My First Description".to_string())
.with_min_max(0f32, 10f32)
.with_access(OSCAccess::ReadWrite)
.with_unit(OSCUnit::Distance(crate::OSCDistance::Meter));
root.add(par1).unwrap();
root.add(par2).unwrap();
root.add(par3).unwrap();
let addr: SocketAddr = ([127, 0, 0, 1], 3000).into();
let (x, y) = run_oscquery_service(root, addr).await.unwrap();
let addr_osc: SocketAddr = ([127, 0, 0, 1], 6669).into();
let sock = tokio::net::UdpSocket::bind(addr_osc).await.unwrap();
let mut buf = [0; 1024];
loop {
let (len, addr) = sock.recv_from(&mut buf).await.unwrap();
println!("{:?} bytes received from {:?}", len, addr);
let len = sock.send_to(&buf[..len], addr).await.unwrap();
println!("{:?} bytes sent", len);
}
let res = tokio::join!(x, y);
res.0.unwrap();
res.1.unwrap();
}