use bytes::{buf::Reader, Bytes};
use std::sync::{Arc, Condvar, Mutex};
+use tokio::sync::Notify;
use warp::http;
pub struct HttpListener {
pub incoming: std::sync::mpsc::Receiver<HttpRequest>,
- pub warp_shutdown: tokio::sync::mpsc::Sender<()>,
+ pub warp_shutdown: Arc<Notify>,
}
pub struct HttpRequest {
use std::sync::LazyLock;
#[cfg(feature = "http")]
use std::sync::{Arc, Condvar, Mutex};
+use tokio::sync::Notify;
use chrono::{offset::Local, DateTime};
#[cfg(not(target_arch = "wasm32"))]
let (tx, rx) = std::sync::mpsc::sync_channel(1024);
// warp shutdown channel
- let (warp_shutdown_tx, mut warp_shutdown_rx) = tokio::sync::mpsc::channel(1);
+ let warp_shutdown = Arc::new(Notify::new());
let runtime = tokio::runtime::Handle::current();
let _guard = runtime.enter();
},
);
+ let warp_shutdown_clone = warp_shutdown.clone();
runtime.spawn(async move {
match ssl_server {
Some((key, cert)) => {
.key(key)
.cert(cert)
.bind_with_graceful_shutdown(addr, async move {
- warp_shutdown_rx.recv().await;
+ warp_shutdown_clone.notified().await;
});
tokio::task::spawn(server);
None => {
let (_addr, server) =
warp::serve(serve).bind_with_graceful_shutdown(addr, async move {
- warp_shutdown_rx.recv().await;
+ warp_shutdown_clone.notified().await;
});
tokio::task::spawn(server);
let http_listener = HttpListener {
incoming: rx,
- warp_shutdown: warp_shutdown_tx,
+ warp_shutdown: warp_shutdown,
};
let http_listener: TypedArenaPtr<HttpListener> =
arena_alloc!(http_listener, &mut self.machine_st.arena);
(HeapCellValueTag::Cons, cons_ptr) => {
match_untyped_arena_ptr!(cons_ptr,
(ArenaHeaderTag::HttpListener, http_listener) => {
- let _ = futures::executor::block_on(http_listener.warp_shutdown.send(()));
+ http_listener.warp_shutdown.notify_one();
}
_ => {
unreachable!();