Replace futures::executor::block_on with tokio::block_in_place

This commit is contained in:
revue_2_presse
2025-02-25 20:19:14 +01:00
parent 4fc4152eac
commit 5d468d3e19
4 changed files with 110 additions and 51 deletions

View File

@@ -91,6 +91,8 @@ use roxmltree;
use futures::future; use futures::future;
#[cfg(feature = "http")] #[cfg(feature = "http")]
use reqwest::Url; use reqwest::Url;
use tokio::runtime::Handle;
use tokio::task;
#[cfg(feature = "http")] #[cfg(feature = "http")]
use warp::hyper::header::{HeaderName, HeaderValue}; use warp::hyper::header::{HeaderName, HeaderValue};
#[cfg(feature = "http")] #[cfg(feature = "http")]
@@ -4378,65 +4380,70 @@ impl Machine {
} }
// do it! // do it!
match futures::executor::block_on(req.send()) { task::block_in_place(move || {
Ok(resp) => { match Handle::current().block_on(req.send()) {
// status code Ok(resp) => {
let status = resp.status().as_u16(); // status code
self.machine_st let status = resp.status().as_u16();
.unify_fixnum(Fixnum::build_with(status as i64), address_status); self.machine_st
// headers .unify_fixnum(Fixnum::build_with(status as i64), address_status);
let headers: Vec<HeapCellValue> = resp // headers
.headers() let headers: Vec<HeapCellValue> = resp
.iter() .headers()
.map(|(header_name, header_value)| { .iter()
let h = self.machine_st.heap.len(); .map(|(header_name, header_value)| {
let h = self.machine_st.heap.len();
let header_term = functor!( let header_term = functor!(
AtomTable::build_with( AtomTable::build_with(
&self.machine_st.atom_tbl, &self.machine_st.atom_tbl,
header_name.as_str() header_name.as_str()
), ),
[cell(string_as_cstr_cell!(AtomTable::build_with( [cell(string_as_cstr_cell!(AtomTable::build_with(
&self.machine_st.atom_tbl, &self.machine_st.atom_tbl,
header_value.to_str().unwrap() header_value.to_str().unwrap()
)))] )))]
); );
self.machine_st.heap.extend(header_term); self.machine_st.heap.extend(header_term);
str_loc_as_cell!(h) str_loc_as_cell!(h)
}) })
.collect(); .collect();
let headers_list = let headers_list =
iter_to_heap_list(&mut self.machine_st.heap, headers.into_iter()); 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 mut stream = Stream::from_http_stream( unify!(
AtomTable::build_with(&self.machine_st.atom_tbl, &address_string), self.machine_st,
reader, heap_loc_as_cell!(headers_list),
&mut self.machine_st.arena, self.machine_st.registers[6]
); );
*stream.options_mut() = StreamOptions::default();
self.indices // body
.add_stream(stream, atom!("http_open"), 3) let reader = futures::executor::block_on(resp.bytes()).unwrap().reader();
.map_err(|stub_gen| stub_gen(&mut self.machine_st))?;
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.indices
self.machine_st.bind(stream_addr.as_var().unwrap(), stream); .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;
}
} }
Err(_) => { });
self.machine_st.fail = true;
}
}
} else { } else {
let err = self let err = self
.machine_st .machine_st

View File

@@ -0,0 +1,23 @@
:- module(http_open_hanging, [submit_request/0]).
:- use_module(library(charsio)).
:- use_module(library(http/http_open)).
send_request :-
Options = [
method('get'),
status_code(StatusCode),
request_headers([]),
headers(_)
],
http_open("https://scryer.pl", _Stream, Options),
write_term('received response with status code':StatusCode, []), nl.
main :-
send_request,
send_request,
send_request,
send_request,
send_request.
:- initialization(main).

View File

@@ -1,3 +1,5 @@
use scryer_prolog::MachineBuilder;
pub(crate) trait Expectable { pub(crate) trait Expectable {
#[track_caller] #[track_caller]
fn assert_eq(self, other: &[u8]); fn assert_eq(self, other: &[u8]);
@@ -31,3 +33,17 @@ pub(crate) fn load_module_test<T: Expectable>(file: &str, expected: T) {
let mut wam = MachineBuilder::default().build(); let mut wam = MachineBuilder::default().build();
expected.assert_eq(wam.test_load_file(file).as_slice()); 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())
});
}

View File

@@ -1,4 +1,6 @@
use crate::helper::load_module_test; use crate::helper::load_module_test;
#[cfg(not(target_arch = "wasm32"))]
use crate::helper::load_module_test_with_tokio_runtime;
use serial_test::serial; use serial_test::serial;
// issue #831 // issue #831
@@ -43,3 +45,14 @@ fn load_context_unreachable() {
fn issue2725_dcg_without_module() { fn issue2725_dcg_without_module() {
load_module_test("tests-pl/issue2725.pl", ""); 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"
);
}