Replace Hyper with Warp for HTTP server

- Use Warp
- Optimize clones
- HTTPS server
- Content-Length limit
- HTTP Basic Auth
- Stop server with Ctrl-C
This commit is contained in:
Adrián Arroyo Calle
2023-09-05 20:40:58 +02:00
parent 841aaa78b4
commit 660860bccf
9 changed files with 682 additions and 339 deletions

View File

@@ -14,7 +14,7 @@ use crate::forms::*;
use crate::heap_iter::*;
use crate::heap_print::*;
#[cfg(feature = "http")]
use crate::http::{HttpListener, HttpResponse, HttpService};
use crate::http::{HttpRequestData, HttpListener, HttpResponse, HttpRequest};
use crate::instructions::*;
use crate::machine;
use crate::machine::code_walker::*;
@@ -44,14 +44,14 @@ pub(crate) use ref_thread_local::RefThreadLocal;
use std::cell::Cell;
use std::cmp::Ordering;
use std::collections::BTreeSet;
use std::collections::{BTreeSet};
use std::convert::TryFrom;
use std::env;
#[cfg(feature = "ffi")]
use std::ffi::CString;
use std::fs;
use std::hash::{BuildHasher, BuildHasherDefault};
use std::io::{ErrorKind, Read, Write};
use std::io::{ErrorKind, Read, BufRead, Write};
use std::iter::{once, FromIterator};
use std::mem;
use std::net::{SocketAddr, TcpListener, TcpStream, ToSocketAddrs};
@@ -59,6 +59,7 @@ use std::num::NonZeroU32;
use std::ops::Sub;
use std::process;
use std::str::FromStr;
use std::sync::{Mutex, Arc, Condvar};
use chrono::{offset::Local, DateTime};
#[cfg(not(target_arch = "wasm32"))]
@@ -91,16 +92,15 @@ use base64;
use roxmltree;
use select;
use bytes::Buf;
use http_body_util::BodyExt;
#[cfg(feature = "http")]
use hyper::header::{HeaderName, HeaderValue};
use warp::hyper::header::{HeaderValue, HeaderName};
#[cfg(feature = "http")]
use hyper::server::conn::http1;
use warp::hyper::{HeaderMap, Method};
#[cfg(feature = "http")]
use hyper::{HeaderMap, Method};
use warp::{Buf, Filter};
#[cfg(feature = "http")]
use reqwest::Url;
use futures::future;
#[cfg(feature = "repl")]
pub(crate) fn get_key() -> KeyEvent {
@@ -4412,54 +4412,119 @@ impl Machine {
#[inline(always)]
pub(crate) fn http_listen(&mut self) -> CallResult {
let address_sink = self.deref_register(1);
if let Some(address_str) = self.machine_st.value_to_str_like(address_sink) {
let address_string = address_str.as_str();
let addr: SocketAddr = match address_string
.to_socket_addrs()
.ok()
.and_then(|mut s| s.next())
{
Some(addr) => addr,
let tls_key = self.deref_register(3);
let tls_cert = self.deref_register(4);
let content_length_limit = self.deref_register(5);
const CONTENT_LENGTH_LIMIT_DEFAULT: u64 = 32768;
let content_length_limit = match Number::try_from(content_length_limit) {
Ok(Number::Fixnum(n)) => if n.get_num() >= 0 {
n.get_num() as u64
} else {
CONTENT_LENGTH_LIMIT_DEFAULT
},
Ok(Number::Integer(n)) => {
let n: Result<u64, _> = (&*n).try_into();
match n {
Ok(u) => u,
Err(_) => CONTENT_LENGTH_LIMIT_DEFAULT,
}
}
_ => CONTENT_LENGTH_LIMIT_DEFAULT,
};
let ssl_server: Option<(String,String)> = {
match self.machine_st.value_to_str_like(tls_key) {
Some(key) => {
match self.machine_st.value_to_str_like(tls_cert) {
Some(cert) => {
let key_str = key.as_str();
let cert_str = cert.as_str();
if key_str.is_empty() || cert_str.is_empty() {
None
} else {
Some((key_str.to_string(), cert_str.to_string()))
}
}
None => None
}
}
None => None
}
};
if let Some(address_str) = self.machine_st.value_to_str_like(address_sink) {
let address_string = address_str.as_str();
let addr: SocketAddr = match address_string.to_socket_addrs().ok().and_then(|mut s| s.next()) {
Some(addr) => addr,
_ => {
self.machine_st.fail = true;
return Ok(());
}
};
let (tx, rx) = std::sync::mpsc::sync_channel(1024);
let (tx, rx) = std::sync::mpsc::sync_channel(1024);
let _guard = self.runtime.enter();
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,
));
}
};
fn get_reader(body: impl Buf + Send + 'static) -> Box<dyn BufRead + Send> {
Box::new(body.reader())
}
self.runtime.spawn(async move {
loop {
let tx = tx.clone();
let (stream, _) = listener.accept().await.unwrap();
let serve = warp::body::aggregate()
.and(warp::header::optional::<u64>(warp::http::header::CONTENT_LENGTH.as_str()))
.and(warp::method())
.and(warp::header::headers_cloned())
.and(warp::path::full())
.and(warp::query::raw().or_else(|_| future::ready(Ok::<(String,), warp::Rejection>(("".to_string(),)))))
.map(move |body, content_length, method, headers: warp::http::HeaderMap, path: warp::filters::path::FullPath, query| {
if let Some(content_length) = content_length {
if content_length > content_length_limit {
return warp::http::Response::builder()
.status(413)
.body(warp::hyper::Body::empty())
.unwrap();
}
}
let http_request_data = HttpRequestData {
method,
headers,
path: path.as_str().to_string(),
query,
body: get_reader(body),
};
let response = Arc::new((Mutex::new(false), Mutex::new(None), Condvar::new()));
let http_request = HttpRequest { request_data: http_request_data, response: Arc::clone(&response) };
// we send the request to http_accept
tx.send(http_request).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 http_listener = HttpListener { incoming: rx };
let http_listener = arena_alloc!(http_listener, &mut self.machine_st.arena);
// 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();
response.expect("Data race error in HTTP server")
}
});
self.runtime.spawn(async move {
match ssl_server {
Some((key, cert)) => {
warp::serve(serve).tls().key(key).cert(cert).run(addr).await
}
None => {
warp::serve(serve).run(addr).await
}
}
});
let http_listener = HttpListener { incoming: rx };
let http_listener = arena_alloc!(http_listener, &mut self.machine_st.arena);
let addr = self.deref_register(2);
self.machine_st.bind(
addr.as_var().unwrap(),
@@ -4472,74 +4537,94 @@ impl Machine {
#[cfg(feature = "http")]
#[inline(always)]
pub(crate) fn http_accept(&mut self) -> CallResult {
let culprit = self.deref_register(1);
let method = self.deref_register(2);
let path = self.deref_register(3);
let query = self.deref_register(5);
let stream_addr = self.deref_register(6);
let handle_addr = self.deref_register(7);
read_heap_cell!(culprit,
(HeapCellValueTag::Cons, cons_ptr) => {
match_untyped_arena_ptr!(cons_ptr,
(ArenaHeaderTag::HttpListener, http_listener) => {
match http_listener.incoming.recv() {
Ok(request) => {
let method_atom = match *request.request.method() {
Method::GET => atom!("get"),
Method::POST => atom!("post"),
Method::PUT => atom!("put"),
Method::DELETE => atom!("delete"),
Method::PATCH => atom!("patch"),
Method::HEAD => atom!("head"),
_ => unreachable!(),
};
let path_atom = AtomTable::build_with(&self.machine_st.atom_tbl, request.request.uri().path());
let path_cell = atom_as_cstr_cell!(path_atom);
let headers: Vec<HeapCellValue> = request.request.headers().iter().map(|(header_name, header_value)| {
let h = self.machine_st.heap.len();
let culprit = self.deref_register(1);
let method = self.deref_register(2);
let path = self.deref_register(3);
let query = self.deref_register(5);
let stream_addr = self.deref_register(6);
let handle_addr = self.deref_register(7);
read_heap_cell!(culprit,
(HeapCellValueTag::Cons, cons_ptr) => {
match_untyped_arena_ptr!(cons_ptr,
(ArenaHeaderTag::HttpListener, http_listener) => {
loop {
match http_listener.incoming.recv_timeout(std::time::Duration::from_millis(200)) {
Ok(request) => {
let method_atom = match request.request_data.method {
Method::GET => atom!("get"),
Method::POST => atom!("post"),
Method::PUT => atom!("put"),
Method::DELETE => atom!("delete"),
Method::PATCH => atom!("patch"),
Method::HEAD => atom!("head"),
Method::OPTIONS => atom!("options"),
Method::TRACE => atom!("trace"),
Method::CONNECT => atom!("connect"),
_ => atom!("unsupported_extension"),
};
let path_atom = AtomTable::build_with(&self.machine_st.atom_tbl, &request.request_data.path);
let path_cell = atom_as_cstr_cell!(path_atom);
let headers: Vec<HeapCellValue> = request.request_data.headers.iter().map(|(header_name, header_value)| {
let h = self.machine_st.heap.len();
let header_term = functor!(AtomTable::build_with(&self.machine_st.atom_tbl, header_name.as_str()), [cell(string_as_cstr_cell!(AtomTable::build_with(&self.machine_st.atom_tbl, header_value.to_str().unwrap())))]);
let header_term = functor!(
AtomTable::build_with(&self.machine_st.atom_tbl, header_name.as_str()),
[cell(string_as_cstr_cell!(AtomTable::build_with(&self.machine_st.atom_tbl, header_value.to_str().unwrap())))]
);
self.machine_st.heap.extend(header_term.into_iter());
str_loc_as_cell!(h)
}).collect();
self.machine_st.heap.extend(header_term.into_iter());
str_loc_as_cell!(h)
}).collect();
let headers_list = iter_to_heap_list(&mut self.machine_st.heap, headers.into_iter());
let headers_list = iter_to_heap_list(&mut self.machine_st.heap, headers.into_iter());
let query_str = request.request_data.query;
let query_atom = AtomTable::build_with(&self.machine_st.atom_tbl, &query_str);
let query_cell = string_as_cstr_cell!(query_atom);
let query_str = request.request.uri().query().unwrap_or("");
let query_atom = AtomTable::build_with(&self.machine_st.atom_tbl, query_str);
let query_cell = string_as_cstr_cell!(query_atom);
let mut stream = Stream::from_http_stream(
path_atom,
request.request_data.body,
&mut self.machine_st.arena
);
*stream.options_mut() = StreamOptions::default();
stream.options_mut().set_stream_type(StreamType::Binary);
self.indices.streams.insert(stream);
let stream = stream_as_cell!(stream);
let hyper_req = request.request;
let buf = self.runtime.block_on(async {hyper_req.collect().await.unwrap().aggregate()});
let reader = buf.reader();
let handle = arena_alloc!(request.response, &mut self.machine_st.arena);
let mut stream = Stream::from_http_stream(
path_atom,
Box::new(reader),
&mut self.machine_st.arena
);
*stream.options_mut() = StreamOptions::default();
stream.options_mut().set_stream_type(StreamType::Binary);
self.indices.streams.insert(stream);
let stream = stream_as_cell!(stream);
self.machine_st.bind(method.as_var().unwrap(), atom_as_cell!(method_atom));
self.machine_st.bind(path.as_var().unwrap(), path_cell);
unify!(self.machine_st, heap_loc_as_cell!(headers_list), self.machine_st.registers[4]);
self.machine_st.bind(query.as_var().unwrap(), query_cell);
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));
break
}
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
let interrupted = machine::INTERRUPT.load(std::sync::atomic::Ordering::Relaxed);
let handle = arena_alloc!(request.response, &mut self.machine_st.arena);
match machine::INTERRUPT.compare_exchange(
interrupted,
false,
std::sync::atomic::Ordering::Relaxed,
std::sync::atomic::Ordering::Relaxed,
) {
Ok(interruption) => {
if interruption {
self.machine_st.throw_interrupt_exception();
self.machine_st.backtrack();
let old_runtime = std::mem::replace(&mut self.runtime, tokio::runtime::Runtime::new().unwrap());
old_runtime.shutdown_background();
break
}
}
Err(_) => unreachable!(),
}
self.machine_st.bind(method.as_var().unwrap(), atom_as_cell!(method_atom));
self.machine_st.bind(path.as_var().unwrap(), path_cell);
unify!(self.machine_st, heap_loc_as_cell!(headers_list), self.machine_st.registers[4]);
self.machine_st.bind(query.as_var().unwrap(), query_cell);
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));
}
Err(_) => {
self.machine_st.fail = true;
}
}
}
Err(_) => {
self.machine_st.fail = true;
}
}
}
_ => {
unreachable!();