Updated hyper to 1.0.0-rc.4

This commit is contained in:
Fayeed Pawaskar
2023-08-08 16:20:49 +05:30
parent cf63b588bc
commit 72ceceb7ae
4 changed files with 101 additions and 38 deletions

View File

@@ -27,28 +27,28 @@ impl Service<Request<IncomingBody>> for HttpService {
type Error = hyper::Error;
type Future = Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send>>;
fn call(&mut self, req: Request<IncomingBody>) -> Self::Future {
// new connection!
// we send the Request info to Prolog
let response = Arc::new((Mutex::new(false), Mutex::new(None), Condvar::new()));
let http_request = HttpRequest { request: req, response: Arc::clone(&response) };
self.tx.send(http_request).unwrap();
fn call(self: &HttpService, req: Request<IncomingBody>) -> Self::Future {
// new connection!
// we send the Request info to Prolog
let response = Arc::new((Mutex::new(false), Mutex::new(None), Condvar::new()));
let http_request = HttpRequest { request: req, response: Arc::clone(&response) };
self.tx.send(http_request).unwrap();
// we wait for the Response info from Prolog
{
let (ready, _response, cvar) = &*response;
let mut ready = ready.lock().unwrap();
while !*ready {
ready = cvar.wait(ready).unwrap();
}
}
{
let (_, response, _) = &*response;
let response = response.lock().unwrap().take();
let res = response.expect("Data race error in HTTP Server");
Box::pin(async move {
Ok(res)
})
}
// we wait for the Response info from Prolog
{
let (ready, _response, cvar) = &*response;
let mut ready = ready.lock().unwrap();
while !*ready {
ready = cvar.wait(ready).unwrap();
}
}
{
let (_, response, _) = &*response;
let response = response.lock().unwrap().take();
let res = response.expect("Data race error in HTTP Server");
Box::pin(async move {
Ok(res)
})
}
}
}

View File

@@ -88,6 +88,7 @@ use hyper::{HeaderMap, Method};
use http_body_util::BodyExt;
use bytes::Buf;
use reqwest::Url;
use hyper_util::rt::TokioIo;
pub(crate) fn get_key() -> KeyEvent {
let key;
@@ -4349,14 +4350,16 @@ impl Machine {
let (stream, _) = listener.accept().await.unwrap();
tokio::task::spawn(async move {
if let Err(err) = http1::Builder::new()
.serve_connection(stream, HttpService {
tx
})
.await
{
eprintln!("Error serving connection: {:?}", err);
}
let io = TokioIo::new(stream);
if let Err(err) = http1::Builder::new()
.serve_connection(io, HttpService {
tx
})
.await
{
eprintln!("Error serving connection: {:?}", err);
}
});
}
});