src/bin/h2server.rs: switch to async

Signed-off-by: Wolfgang Bumiller <w.bumiller@proxmox.com>
This commit is contained in:
Wolfgang Bumiller 2019-08-28 15:50:40 +02:00
parent 55d8a631fc
commit 74be6dc9b7

View File

@ -4,62 +4,47 @@ use futures::*;
// Simple H2 server to test H2 speed with h2client.rs // Simple H2 server to test H2 speed with h2client.rs
use tokio::net::TcpListener; use tokio::net::TcpListener;
use tokio::io::{AsyncRead, AsyncWrite};
use proxmox_backup::client::pipe_to_stream::*; use proxmox_backup::client::pipe_to_stream::PipeToSendStream;
pub fn main() -> Result<(), Error> {
start_h2_server()?;
Ok(())
}
pub fn start_h2_server() -> Result<(), Error> {
#[tokio::main]
async fn main() -> Result<(), Error> {
let listener = TcpListener::bind(&"127.0.0.1:8008".parse().unwrap()).unwrap(); let listener = TcpListener::bind(&"127.0.0.1:8008".parse().unwrap()).unwrap();
println!("listening on {:?}", listener.local_addr()); println!("listening on {:?}", listener.local_addr());
let server = listener.incoming().for_each(move |socket| { let mut incoming = listener.incoming();
while let Some(socket) = incoming.try_next().await? {
let connection = h2::server::handshake(socket) tokio::spawn(handle_connection(socket)
.map_err(Error::from) .map(|res| {
.and_then(|conn| { if let Err(err) = res {
println!("H2 connection bound"); eprintln!("Error: {}", err);
}
conn }));
.map_err(Error::from) }
.for_each(|(request, mut respond)| {
println!("GOT request: {:?}", request); Ok(())
}
let response = http::Response::builder().status(http::StatusCode::OK).body(()).unwrap();
async fn handle_connection<T: AsyncRead + AsyncWrite + Unpin>(socket: T) -> Result<(), Error> {
let send = respond.send_response(response, false).unwrap(); let mut conn = h2::server::handshake(socket).await?;
let data = vec![65u8; 1024*1024];
PipeToSendStream::new(bytes::Bytes::from(data), send) println!("H2 connection bound");
.and_then(|_| {
println!("DATA SENT"); while let Some((request, mut respond)) = conn.try_next().await? {
Ok(()) println!("GOT request: {:?}", request);
})
}) let response = http::Response::builder()
}) .status(http::StatusCode::OK)
.and_then(|_| { .body(())
println!("H2 connection CLOSE !"); .unwrap();
Ok(())
}) let send = respond.send_response(response, false).unwrap();
.then(|res| { let data = vec![65u8; 1024*1024];
if let Err(e) = res { PipeToSendStream::new(bytes::Bytes::from(data), send).await?;
println!(" -> err={:?}", e); println!("DATA SENT");
} }
Ok(())
});
tokio::spawn(Box::new(connection));
Ok(())
})
.map_err(|e| eprintln!("accept error: {}", e));
tokio::run(server);
Ok(()) Ok(())
} }