Merge pull request #1832 from aarroyoc/http-fixes

Multiple fixes for http libraries
This commit is contained in:
Mark Thom
2023-07-06 16:30:56 -06:00
committed by GitHub
6 changed files with 545 additions and 190 deletions

View File

@@ -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();
@@ -3142,7 +3140,7 @@ impl Machine {
}
bytes.push(c as u8);
}
}
} else {
bytes = string.as_str().bytes().collect();
}
@@ -4201,63 +4199,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);
@@ -4281,26 +4281,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 };
@@ -4323,8 +4328,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"),
@@ -4355,7 +4360,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(
@@ -4377,7 +4382,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;
}
}
@@ -4435,15 +4440,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();