use futures::future;
#[cfg(feature = "http")]
use reqwest::Url;
+use tokio::runtime::Handle;
+use tokio::task;
#[cfg(feature = "http")]
use warp::hyper::header::{HeaderName, HeaderValue};
#[cfg(feature = "http")]
}
// do it!
- match futures::executor::block_on(req.send()) {
- 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 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()
- )))]
- );
+ task::block_in_place(move || {
+ match Handle::current().block_on(req.send()) {
+ 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 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);
- str_loc_as_cell!(h)
- })
- .collect();
+ self.machine_st.heap.extend(header_term);
+ str_loc_as_cell!(h)
+ })
+ .collect();
- 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 = futures::executor::block_on(resp.bytes()).unwrap().reader();
+ let headers_list =
+ iter_to_heap_list(&mut self.machine_st.heap, headers.into_iter());
- let mut stream = Stream::from_http_stream(
- AtomTable::build_with(&self.machine_st.atom_tbl, &address_string),
- reader,
- &mut self.machine_st.arena,
- );
- *stream.options_mut() = StreamOptions::default();
+ unify!(
+ self.machine_st,
+ heap_loc_as_cell!(headers_list),
+ self.machine_st.registers[6]
+ );
- self.indices
- .add_stream(stream, atom!("http_open"), 3)
- .map_err(|stub_gen| stub_gen(&mut self.machine_st))?;
+ // body
+ let reader = futures::executor::block_on(resp.bytes()).unwrap().reader();
- let stream = stream_as_cell!(stream);
+ let mut stream = Stream::from_http_stream(
+ AtomTable::build_with(&self.machine_st.atom_tbl, &address_string),
+ reader,
+ &mut self.machine_st.arena,
+ );
+ *stream.options_mut() = StreamOptions::default();
- let stream_addr = self.deref_register(2);
- self.machine_st.bind(stream_addr.as_var().unwrap(), stream);
- }
- Err(_) => {
- self.machine_st.fail = true;
+ self.indices
+ .add_stream(stream, atom!("http_open"), 3)
+ .map_err(|stub_gen| stub_gen(&mut self.machine_st))
+ .unwrap();
+
+ 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
+use scryer_prolog::MachineBuilder;
+
pub(crate) trait Expectable {
#[track_caller]
fn assert_eq(self, other: &[u8]);
let mut wam = MachineBuilder::default().build();
expected.assert_eq(wam.test_load_file(file).as_slice());
}
+
+/// Same as `load_module_test` with tokio runtime
+#[cfg(not(target_arch = "wasm32"))]
+pub(crate) fn load_module_test_with_tokio_runtime<T: Expectable>(file: &str, expected: T) {
+ let runtime = tokio::runtime::Builder::new_multi_thread()
+ .enable_all()
+ .build()
+ .unwrap();
+
+ runtime.block_on(async move {
+ let mut wam = MachineBuilder::default().build();
+ expected.assert_eq(wam.test_load_file(file).as_slice())
+ });
+}
use crate::helper::load_module_test;
+#[cfg(not(target_arch = "wasm32"))]
+use crate::helper::load_module_test_with_tokio_runtime;
use serial_test::serial;
// issue #831
fn issue2725_dcg_without_module() {
load_module_test("tests-pl/issue2725.pl", "");
}
+
+#[test]
+#[cfg(feature = "http")]
+#[cfg(not(target_arch = "wasm32"))]
+#[cfg_attr(miri, ignore = "it takes too long to run")]
+fn http_open_hanging() {
+ load_module_test_with_tokio_runtime(
+ "tests-pl/issue-http_open-hanging.pl",
+ "received response with status code:200\nreceived response with status code:200\nreceived response with status code:200\nreceived response with status code:200\nreceived response with status code:200\n"
+ );
+}