Multiple fixes for http libraries
* use reqwest for http_open (still uses Hyper underneath) * use Hyper 1.0.0-rc3 for server * Modify all internal handling of server
This commit is contained in:
57
src/http.rs
57
src/http.rs
@@ -1,25 +1,54 @@
|
||||
use std::sync::Arc;
|
||||
use std::convert::Infallible;
|
||||
|
||||
use hyper::{Response, Request, Body};
|
||||
use tokio::sync::Mutex;
|
||||
use tokio::sync::mpsc::{channel, Receiver, Sender};
|
||||
use std::sync::{Arc, Mutex, Condvar};
|
||||
use std::future::Future;
|
||||
use std::pin::Pin;
|
||||
use http_body_util::Full;
|
||||
use bytes::Bytes;
|
||||
use hyper::service::Service;
|
||||
use hyper::{body::Incoming as IncomingBody, Request, Response};
|
||||
|
||||
pub struct HttpListener {
|
||||
pub incoming: Receiver<HttpRequest>
|
||||
pub incoming: std::sync::mpsc::Receiver<HttpRequest>
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct HttpRequest {
|
||||
pub request: Request<Body>,
|
||||
pub request: Request<IncomingBody>,
|
||||
pub response: HttpResponse,
|
||||
}
|
||||
|
||||
pub type HttpResponse = Sender<Response<Body>>;
|
||||
pub type HttpResponse = Arc<(Mutex<bool>, Mutex<Option<Response<Full<Bytes>>>>, Condvar)>;
|
||||
|
||||
pub async fn serve_req(req: Request<Body>, tx: Arc<Mutex<Sender<HttpRequest>>>) -> Result<Response<Body>, Infallible> {
|
||||
let (response_tx, mut rx) = channel(1);
|
||||
let http_request = HttpRequest { request: req, response: response_tx };
|
||||
tx.lock().await.send(http_request).await.unwrap();
|
||||
Ok(rx.recv().await.unwrap())
|
||||
pub struct HttpService {
|
||||
pub tx: std::sync::mpsc::SyncSender<HttpRequest>,
|
||||
}
|
||||
|
||||
impl Service<Request<IncomingBody>> for HttpService {
|
||||
type Response = Response<Full<Bytes>>;
|
||||
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();
|
||||
|
||||
// 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)
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -121,37 +121,44 @@ http_loop(HttpListener, Handlers) :-
|
||||
send_response(ResponseHandle, http_response(StatusCode0, text(ResponseText), ResponseHeaders0)) :-
|
||||
default(StatusCode0, 200, StatusCode),
|
||||
maplist(map_header_kv_2, ResponseHeaders, ResponseHeaders0),
|
||||
'$http_answer'(ResponseHandle, StatusCode, ResponseHeaders, ResponseStream),
|
||||
call_cleanup(
|
||||
format(ResponseStream, "~s", [ResponseText]),
|
||||
close(ResponseStream)
|
||||
'$http_answer'(ResponseHandle, StatusCode, ResponseHeaders, ResponseStream0),
|
||||
open(stream(ResponseStream0), write, ResponseStream, [type(text)]),
|
||||
catch(
|
||||
call_cleanup(format(ResponseStream, "~s", [ResponseText]),close(ResponseStream)),
|
||||
error(existence_error(stream, _), _),
|
||||
true
|
||||
).
|
||||
|
||||
send_response(ResponseHandle, http_response(StatusCode0, bytes(ResponseBytes), ResponseHeaders0)) :-
|
||||
default(StatusCode0, 200, StatusCode),
|
||||
maplist(map_header_kv_2, ResponseHeaders, ResponseHeaders0),
|
||||
'$http_answer'(ResponseHandle, StatusCode, ResponseHeaders, ResponseStream),
|
||||
call_cleanup(
|
||||
format(ResponseStream, "~s", [ResponseBytes]),
|
||||
close(ResponseStream)
|
||||
catch(
|
||||
call_cleanup(format(ResponseStream, "~s", [ResponseBytes]),close(ResponseStream)),
|
||||
error(existence_error(stream, _), _),
|
||||
true
|
||||
).
|
||||
|
||||
send_response(ResponseHandle, http_response(StatusCode0, file(Filename), ResponseHeaders0)) :-
|
||||
default(StatusCode0, 200, StatusCode),
|
||||
maplist(map_header_kv_2, ResponseHeaders, ResponseHeaders0),
|
||||
'$http_answer'(ResponseHandle, StatusCode, ResponseHeaders, ResponseStream),
|
||||
call_cleanup(
|
||||
setup_call_cleanup(
|
||||
open(Filename, read, FileStream, [type(binary)]),
|
||||
(
|
||||
get_n_chars(FileStream, _, FileCs),
|
||||
format(ResponseStream, "~s", [FileCs])
|
||||
catch(
|
||||
call_cleanup(
|
||||
setup_call_cleanup(
|
||||
open(Filename, read, FileStream, [type(binary)]),
|
||||
(
|
||||
get_n_chars(FileStream, _, FileCs),
|
||||
format(ResponseStream, "~s", [FileCs])
|
||||
),
|
||||
close(FileStream)
|
||||
),
|
||||
close(FileStream)
|
||||
close(ResponseStream)
|
||||
),
|
||||
close(ResponseStream)
|
||||
error(existence_error(stream, _), _),
|
||||
true
|
||||
).
|
||||
|
||||
|
||||
|
||||
default(Var, Default, Out) :-
|
||||
(var(Var) -> Out = Default
|
||||
|
||||
@@ -9,6 +9,7 @@ use crate::machine::machine_errors::*;
|
||||
use crate::machine::machine_indices::*;
|
||||
use crate::machine::machine_state::*;
|
||||
use crate::types::*;
|
||||
use crate::http::HttpResponse;
|
||||
|
||||
pub use modular_bitfield::prelude::*;
|
||||
|
||||
@@ -26,7 +27,6 @@ use std::ops::{Deref, DerefMut};
|
||||
use std::ptr;
|
||||
|
||||
use native_tls::TlsStream;
|
||||
use hyper::body::{Bytes, Sender};
|
||||
|
||||
#[derive(Debug, BitfieldSpecifier, Clone, Copy, PartialEq, Eq, Hash)]
|
||||
#[bits = 1]
|
||||
@@ -276,7 +276,10 @@ impl Read for HttpReadStream {
|
||||
}
|
||||
|
||||
pub struct HttpWriteStream {
|
||||
body_writer: Sender,
|
||||
status_code: u16,
|
||||
headers: hyper::HeaderMap,
|
||||
response: TypedArenaPtr<HttpResponse>,
|
||||
buffer: Vec<u8>,
|
||||
}
|
||||
|
||||
impl Debug for HttpWriteStream {
|
||||
@@ -288,17 +291,28 @@ impl Debug for HttpWriteStream {
|
||||
impl Write for HttpWriteStream {
|
||||
#[inline]
|
||||
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
|
||||
let bytes = Bytes::copy_from_slice(buf);
|
||||
let len = bytes.len();
|
||||
match self.body_writer.try_send_data(bytes) {
|
||||
Ok(()) => Ok(len),
|
||||
Err(_) => Err(std::io::Error::from(ErrorKind::Interrupted))
|
||||
}
|
||||
self.buffer.extend_from_slice(buf);
|
||||
Ok(buf.len())
|
||||
}
|
||||
|
||||
#[inline]
|
||||
fn flush(&mut self) -> std::io::Result<()> {
|
||||
Ok(())
|
||||
let (ready, response, cvar) = &**self.response;
|
||||
|
||||
let mut ready = ready.lock().unwrap();
|
||||
{
|
||||
let mut response = response.lock().unwrap();
|
||||
|
||||
let bytes = bytes::Bytes::copy_from_slice(&self.buffer);
|
||||
let mut response_ = hyper::Response::builder()
|
||||
.status(self.status_code);
|
||||
*response_.headers_mut().unwrap() = self.headers.clone();
|
||||
*response = Some(response_.body(http_body_util::Full::new(bytes)).unwrap());
|
||||
}
|
||||
*ready = true;
|
||||
cvar.notify_one();
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1084,15 +1098,20 @@ impl Stream {
|
||||
|
||||
#[inline]
|
||||
pub(crate) fn from_http_sender(
|
||||
body_writer: Sender,
|
||||
arena: &mut Arena,
|
||||
response: TypedArenaPtr<HttpResponse>,
|
||||
status_code: u16,
|
||||
headers: hyper::HeaderMap,
|
||||
arena: &mut Arena,
|
||||
) -> Self {
|
||||
Stream::HttpWrite(arena_alloc!(
|
||||
StreamLayout::new(CharReader::new(HttpWriteStream {
|
||||
body_writer
|
||||
})),
|
||||
arena
|
||||
))
|
||||
Stream::HttpWrite(arena_alloc!(
|
||||
StreamLayout::new(CharReader::new(HttpWriteStream {
|
||||
response,
|
||||
status_code,
|
||||
headers,
|
||||
buffer: Vec::new(),
|
||||
})),
|
||||
arena
|
||||
))
|
||||
}
|
||||
|
||||
#[inline]
|
||||
@@ -1139,10 +1158,10 @@ impl Stream {
|
||||
|
||||
Ok(())
|
||||
}
|
||||
Stream::HttpWrite(ref mut http_stream) => {
|
||||
Stream::HttpWrite(ref mut http_stream) => {
|
||||
unsafe {
|
||||
http_stream.set_tag(ArenaHeaderTag::Dropped);
|
||||
std::ptr::drop_in_place(&mut http_stream.inner_mut().body_writer as *mut _);
|
||||
std::ptr::drop_in_place(&mut http_stream.inner_mut().buffer as *mut _);
|
||||
}
|
||||
|
||||
Ok(())
|
||||
|
||||
@@ -9,7 +9,7 @@ use crate::forms::*;
|
||||
use crate::ffi::*;
|
||||
use crate::heap_iter::*;
|
||||
use crate::heap_print::*;
|
||||
use crate::http::{self, HttpListener, HttpResponse};
|
||||
use crate::http::{HttpService, HttpListener, HttpResponse};
|
||||
use crate::instructions::*;
|
||||
use crate::machine;
|
||||
use crate::machine::{Machine, VERIFY_ATTR_INTERRUPT_LOC, get_structure_index};
|
||||
@@ -39,7 +39,7 @@ use ref_thread_local::{RefThreadLocal, ref_thread_local};
|
||||
use std::cell::Cell;
|
||||
use std::cmp::Ordering;
|
||||
use std::collections::BTreeSet;
|
||||
use std::convert::{TryFrom, Infallible};
|
||||
use std::convert::{TryFrom};
|
||||
use std::env;
|
||||
use std::ffi::CString;
|
||||
use std::fs;
|
||||
@@ -52,7 +52,6 @@ use std::num::NonZeroU32;
|
||||
use std::ops::Sub;
|
||||
use std::process;
|
||||
use std::str::FromStr;
|
||||
use std::sync::Arc;
|
||||
|
||||
use chrono::{offset::Local, DateTime};
|
||||
use cpu_time::ProcessTime;
|
||||
@@ -80,13 +79,12 @@ use base64;
|
||||
use roxmltree;
|
||||
use select;
|
||||
|
||||
use hyper::{Body, Server, Client, HeaderMap, Method, Request, Response, Uri};
|
||||
use hyper::header::{HeaderName, HeaderValue};
|
||||
use hyper::body::Buf;
|
||||
use hyper::service::{make_service_fn, service_fn};
|
||||
use hyper_tls::HttpsConnector;
|
||||
use tokio::sync::Mutex;
|
||||
use tokio::sync::mpsc::channel;
|
||||
use hyper::server::conn::http1;
|
||||
use hyper::header::{HeaderValue, HeaderName};
|
||||
use hyper::{HeaderMap, Method};
|
||||
use http_body_util::BodyExt;
|
||||
use bytes::Buf;
|
||||
use reqwest::Url;
|
||||
|
||||
ref_thread_local! {
|
||||
pub(crate) static managed RANDOM_STATE: RandState<'static> = RandState::new();
|
||||
@@ -3125,7 +3123,7 @@ impl Machine {
|
||||
}
|
||||
|
||||
bytes.push(c as u8);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
bytes = string.as_str().bytes().collect();
|
||||
}
|
||||
@@ -4184,63 +4182,65 @@ impl Machine {
|
||||
};
|
||||
if let Some(address_sink) = self.machine_st.value_to_str_like(address_sink) {
|
||||
let address_string = address_sink.as_str(); //to_string();
|
||||
let address: Uri = address_string.parse().unwrap();
|
||||
let address: Url = address_string.parse().unwrap();
|
||||
|
||||
let stream = self.runtime.block_on(async {
|
||||
let https = HttpsConnector::new();
|
||||
let client = Client::builder()
|
||||
.build::<_, hyper::Body>(https);
|
||||
let client = reqwest::blocking::Client::builder()
|
||||
.build()
|
||||
.unwrap();
|
||||
|
||||
// request
|
||||
let mut req = Request::builder()
|
||||
.method(method)
|
||||
.uri(address)
|
||||
.body(Body::from(bytes))
|
||||
.unwrap();
|
||||
// request headers
|
||||
*req.headers_mut() = headers;
|
||||
// do it!
|
||||
let resp = client.request(req).await.unwrap();
|
||||
// status code
|
||||
let status = resp.status().as_u16();
|
||||
self.machine_st.unify_fixnum(Fixnum::build_with(status as i64), address_status);
|
||||
// headers
|
||||
let headers: Vec<HeapCellValue> = resp.headers().iter().map(|(header_name, header_value)| {
|
||||
let h = self.machine_st.heap.len();
|
||||
// request
|
||||
let mut req = reqwest::blocking::Request::new(method, address);
|
||||
|
||||
let header_term = functor!(
|
||||
self.machine_st.atom_tbl.build_with(header_name.as_str()),
|
||||
[cell(string_as_cstr_cell!(self.machine_st.atom_tbl.build_with(header_value.to_str().unwrap())))]
|
||||
);
|
||||
*req.headers_mut() = headers;
|
||||
if bytes.len() > 0 {
|
||||
*req.body_mut() = Some(reqwest::blocking::Body::from(bytes));
|
||||
}
|
||||
|
||||
self.machine_st.heap.extend(header_term.into_iter());
|
||||
str_loc_as_cell!(h)
|
||||
}).collect();
|
||||
// do it!
|
||||
match client.execute(req) {
|
||||
Ok(resp) => {
|
||||
// status code
|
||||
let status = resp.status().as_u16();
|
||||
self.machine_st.unify_fixnum(Fixnum::build_with(status as i64), address_status);
|
||||
// headers
|
||||
let headers: Vec<HeapCellValue> = resp.headers().iter().map(|(header_name, header_value)| {
|
||||
let h = self.machine_st.heap.len();
|
||||
|
||||
let headers_list = iter_to_heap_list(&mut self.machine_st.heap, headers.into_iter());
|
||||
unify!(self.machine_st, heap_loc_as_cell!(headers_list), self.machine_st.registers[6]);
|
||||
// body
|
||||
let buf = hyper::body::aggregate(resp).await.unwrap();
|
||||
let reader = buf.reader();
|
||||
let header_term = functor!(
|
||||
self.machine_st.atom_tbl.build_with(header_name.as_str()),
|
||||
[cell(string_as_cstr_cell!(self.machine_st.atom_tbl.build_with(header_value.to_str().unwrap())))]
|
||||
);
|
||||
|
||||
let mut stream = Stream::from_http_stream(
|
||||
self.machine_st.atom_tbl.build_with(&address_string),
|
||||
Box::new(reader),
|
||||
&mut self.machine_st.arena
|
||||
);
|
||||
*stream.options_mut() = StreamOptions::default();
|
||||
if let Some(alias) = stream.options().get_alias() {
|
||||
self.indices.stream_aliases.insert(alias, stream);
|
||||
}
|
||||
self.machine_st.heap.extend(header_term.into_iter());
|
||||
str_loc_as_cell!(h)
|
||||
}).collect();
|
||||
|
||||
self.indices.streams.insert(stream);
|
||||
let headers_list = iter_to_heap_list(&mut self.machine_st.heap, headers.into_iter());
|
||||
unify!(self.machine_st, heap_loc_as_cell!(headers_list), self.machine_st.registers[6]);
|
||||
// body
|
||||
let reader = resp.bytes().unwrap().reader();
|
||||
|
||||
stream_as_cell!(stream)
|
||||
});
|
||||
let mut stream = Stream::from_http_stream(
|
||||
self.machine_st.atom_tbl.build_with(&address_string),
|
||||
Box::new(reader),
|
||||
&mut self.machine_st.arena
|
||||
);
|
||||
*stream.options_mut() = StreamOptions::default();
|
||||
if let Some(alias) = stream.options().get_alias() {
|
||||
self.indices.stream_aliases.insert(alias, stream);
|
||||
}
|
||||
|
||||
let stream_addr = self.deref_register(2);
|
||||
self.machine_st.bind(stream_addr.as_var().unwrap(), stream);
|
||||
self.indices.streams.insert(stream);
|
||||
|
||||
let stream = stream_as_cell!(stream);
|
||||
|
||||
let stream_addr = self.deref_register(2);
|
||||
self.machine_st.bind(stream_addr.as_var().unwrap(), stream);
|
||||
},
|
||||
Err(_) => {
|
||||
self.machine_st.fail = true;
|
||||
}
|
||||
}
|
||||
} else {
|
||||
let err = self.machine_st.domain_error(DomainErrorType::SourceSink, address_sink);
|
||||
let stub = functor_stub(atom!("http_open"), 3);
|
||||
@@ -4264,26 +4264,31 @@ impl Machine {
|
||||
}
|
||||
};
|
||||
|
||||
let (tx, rx) = channel(1);
|
||||
let tx = Arc::new(Mutex::new(tx));
|
||||
let (tx, rx) = std::sync::mpsc::sync_channel(1024);
|
||||
|
||||
let _guard = self.runtime.enter();
|
||||
let server = match Server::try_bind(&addr) {
|
||||
Ok(server) => server,
|
||||
let listener = match self.runtime.block_on(async { tokio::net::TcpListener::bind(addr).await }) {
|
||||
Ok(listener) => listener,
|
||||
Err(_) => {
|
||||
return Err(self.machine_st.open_permission_error(address_sink, atom!("http_listen"), 2));
|
||||
}
|
||||
};
|
||||
|
||||
self.runtime.spawn(async move {
|
||||
let make_svc = make_service_fn(move |_conn| {
|
||||
loop {
|
||||
let tx = tx.clone();
|
||||
async move { Ok::<_, Infallible>(service_fn(move |req| http::serve_req(req, tx.clone()))) }
|
||||
});
|
||||
let server = server.serve(make_svc);
|
||||
let (stream, _) = listener.accept().await.unwrap();
|
||||
|
||||
if let Err(_) = server.await {
|
||||
eprintln!("server error");
|
||||
tokio::task::spawn(async move {
|
||||
if let Err(err) = http1::Builder::new()
|
||||
.serve_connection(stream, HttpService {
|
||||
tx
|
||||
})
|
||||
.await
|
||||
{
|
||||
eprintln!("Error serving connection: {:?}", err);
|
||||
}
|
||||
});
|
||||
}
|
||||
});
|
||||
let http_listener = HttpListener { incoming: rx };
|
||||
@@ -4306,8 +4311,8 @@ impl Machine {
|
||||
(HeapCellValueTag::Cons, cons_ptr) => {
|
||||
match_untyped_arena_ptr!(cons_ptr,
|
||||
(ArenaHeaderTag::HttpListener, http_listener) => {
|
||||
match http_listener.incoming.blocking_recv() {
|
||||
Some(request) => {
|
||||
match http_listener.incoming.recv() {
|
||||
Ok(request) => {
|
||||
let method_atom = match *request.request.method() {
|
||||
Method::GET => atom!("get"),
|
||||
Method::POST => atom!("post"),
|
||||
@@ -4338,7 +4343,7 @@ impl Machine {
|
||||
let query_cell = string_as_cstr_cell!(query_atom);
|
||||
|
||||
let hyper_req = request.request;
|
||||
let buf = self.runtime.block_on(async {hyper::body::aggregate(hyper_req).await.unwrap()});
|
||||
let buf = self.runtime.block_on(async {hyper_req.collect().await.unwrap().aggregate()});
|
||||
let reader = buf.reader();
|
||||
|
||||
let mut stream = Stream::from_http_stream(
|
||||
@@ -4360,7 +4365,7 @@ impl Machine {
|
||||
self.machine_st.bind(stream_addr.as_var().unwrap(), stream);
|
||||
self.machine_st.bind(handle_addr.as_var().unwrap(), typed_arena_ptr_as_cell!(handle));
|
||||
}
|
||||
None => {
|
||||
Err(_) => {
|
||||
self.machine_st.fail = true;
|
||||
}
|
||||
}
|
||||
@@ -4418,15 +4423,10 @@ impl Machine {
|
||||
(HeapCellValueTag::Cons, cons_ptr) => {
|
||||
match_untyped_arena_ptr!(cons_ptr,
|
||||
(ArenaHeaderTag::HttpResponse, http_response) => {
|
||||
let mut response = Response::builder()
|
||||
.status(status_code);
|
||||
*response.headers_mut().unwrap() = headers;
|
||||
let (sender, body) = Body::channel();
|
||||
let response = response.body(body).unwrap();
|
||||
http_response.blocking_send(response).unwrap();
|
||||
|
||||
let mut stream = Stream::from_http_sender(
|
||||
sender,
|
||||
http_response,
|
||||
status_code,
|
||||
headers,
|
||||
&mut self.machine_st.arena
|
||||
);
|
||||
*stream.options_mut() = StreamOptions::default();
|
||||
|
||||
Reference in New Issue
Block a user