From: Danil Platonov Date: Sat, 30 May 2026 22:47:50 +0000 (-0700) Subject: allow to stop and start http server at will with SIGINT, stop the server on uncaught... X-Git-Url: https://git.sagredo.dev/?a=commitdiff_plain;h=f8f7b3b7e2c8e38e4fbad96376eeffe6271d2222;p=scryer-prolog.git allow to stop and start http server at will with SIGINT, stop the server on uncaught throws, display timing info for each request, use setup_call_cleanup in more places to avoid potential leaks, when handler for a path doesn't have a matching case, throw a detailed error instead of silently continuing --- diff --git a/build/instructions_template.rs b/build/instructions_template.rs index 7b8be1f8..2dad5291 100644 --- a/build/instructions_template.rs +++ b/build/instructions_template.rs @@ -599,6 +599,8 @@ enum SystemClauseType { HttpOpen, #[strum_discriminants(strum(props(Arity = "5", Name = "$http_listen")))] HttpListen, + #[strum_discriminants(strum(props(Arity = "1", Name = "$http_listen_stop")))] + HttpListenStop, #[strum_discriminants(strum(props(Arity = "7", Name = "$http_accept")))] HttpAccept, #[strum_discriminants(strum(props(Arity = "4", Name = "$http_answer")))] diff --git a/src/http.rs b/src/http.rs index bd6ec33b..939dbe7d 100644 --- a/src/http.rs +++ b/src/http.rs @@ -5,6 +5,7 @@ use warp::http; pub struct HttpListener { pub incoming: std::sync::mpsc::Receiver, + pub warp_shutdown: tokio::sync::mpsc::Sender<()>, } pub struct HttpRequest { diff --git a/src/instructions.rs b/src/instructions.rs index f029e463..2fc94c63 100644 --- a/src/instructions.rs +++ b/src/instructions.rs @@ -803,6 +803,7 @@ impl Instruction { | &Instruction::CallDeterministicLengthRundown | &Instruction::CallHttpOpen | &Instruction::CallHttpListen + | &Instruction::CallHttpListenStop | &Instruction::CallHttpAccept | &Instruction::CallHttpAnswer | &Instruction::CallLoadForeignLib @@ -1062,6 +1063,7 @@ impl Instruction { | &Instruction::ExecuteDeterministicLengthRundown | &Instruction::ExecuteHttpOpen | &Instruction::ExecuteHttpListen + | &Instruction::ExecuteHttpListenStop | &Instruction::ExecuteHttpAccept | &Instruction::ExecuteHttpAnswer | &Instruction::ExecuteLoadForeignLib diff --git a/src/lib/http/http_server.pl b/src/lib/http/http_server.pl index 242e0241..cc4acb59 100644 --- a/src/lib/http/http_server.pl +++ b/src/lib/http/http_server.pl @@ -112,12 +112,29 @@ module_qualification(M, H0, H) :- H0 =.. [Method, Path, Goal], H =.. [Method, Path, M:Goal]. +http_listen__(Addr, HttpListener, TLSKey, TLSCert, ContentLengthLimit) :- + '$http_listen'(Addr, HttpListener, TLSKey, TLSCert, ContentLengthLimit). + +http_listen_stop_(HttpListener) :- + '$http_listen_stop'(HttpListener). + +http_accept_(HttpListener, RequestMethod, RequestPath, RequestHeaders, RequestQuery, RequestStream, ResponseHandle) :- + '$http_accept'(HttpListener, RequestMethod, RequestPath, RequestHeaders, RequestQuery, RequestStream, ResponseHandle). + +http_answer_(ResponseHandle, Code, Headers, ResponseStream) :- + '$http_answer'(ResponseHandle, Code, Headers, ResponseStream). + http_listen_(Port, Handlers, Options) :- parse_options(Options, TLSKey, TLSCert, ContentLengthLimit), phrase(format_("0.0.0.0:~d", [Port]), Addr), - '$http_listen'(Addr, HttpListener, TLSKey, TLSCert, ContentLengthLimit),!, - format("Listening at ~s\n", [Addr]), - http_loop(HttpListener, Handlers). + setup_call_cleanup( + ( + http_listen__(Addr, HttpListener, TLSKey, TLSCert, ContentLengthLimit), + format("Listening at http://~s\n", [Addr]) + ), + http_loop(HttpListener, Handlers), + http_listen_stop_(HttpListener) + ). parse_options(Options, TLSKey, TLSCert, ContentLengthLimit) :- member_option_default(tls_key, Options, "", TLSKey), @@ -131,37 +148,56 @@ member_option_default(Key, List, _Default, Value) :- member_option_default(Key, List, Default, Default) :- X =.. [Key, _], \+ member(X, List). - http_loop(HttpListener, Handlers) :- - '$http_accept'(HttpListener, RequestMethod, RequestPath, RequestHeaders, RequestQuery, RequestStream, ResponseHandle), - current_time(Time), - phrase(format_time("%Y-%m-%d (%H:%M:%S)", Time), TimeString), - format("~s ~w ~s\n", [TimeString, RequestMethod, RequestPath]), - maplist(map_header_kv, RequestHeaders, RequestHeadersKV), - phrase(parse_queries(RequestQueries), RequestQuery), - ( - match_handler(Handlers, RequestMethod, RequestPath, Handler) -> - ( - HttpRequest = http_request(RequestHeadersKV, stream(RequestStream), RequestQueries), - HttpResponse = http_response(_, _, _), - (call(Handler, HttpRequest, HttpResponse) -> - send_response(ResponseHandle, HttpResponse) - ; ( - '$http_answer'(ResponseHandle, 500, [], ResponseStream), - call_cleanup(format(ResponseStream, "Internal Server Error", []), close(ResponseStream))) - ) - ) - ; ( - '$http_answer'(ResponseHandle, 404, [], ResponseStream), - call_cleanup(format(ResponseStream, "Not Found", []), close(ResponseStream))) - ), + time(( + http_accept_(HttpListener, RequestMethod, RequestPath, RequestHeaders, RequestQuery, RequestStream, ResponseHandle), + current_time(Time), + phrase(format_time("%Y-%m-%d (%H:%M:%S)", Time), TimeString), + format("~s ~w ~s", [TimeString, RequestMethod, RequestPath]), + maplist(map_header_kv, RequestHeaders, RequestHeadersKV), + phrase(parse_queries(RequestQueries), RequestQuery), + ( + match_handler(Handlers, RequestMethod, RequestPath, Handler) -> + ( + HttpRequest = http_request(RequestHeadersKV, stream(RequestStream), RequestQueries), + HttpResponse = http_response(_, _, _), + catch( + (call(Handler, HttpRequest, HttpResponse) -> + send_response(ResponseHandle, HttpResponse) + ; + setup_call_cleanup( + http_answer_(ResponseHandle, 500, [], ResponseStream), + format(ResponseStream, "Internal Server Error", []), + close(ResponseStream) + ), + throw(handler_not_available(Handler, RequestMethod, RequestPath, RequestQuery, RequestHeaders)) + ), + HandlerError, + ( + setup_call_cleanup( + http_answer_(ResponseHandle, 500, [], ResponseStream), + format(ResponseStream, "Internal Server Error", []), + close(ResponseStream) + ), + throw(HandlerError) + ) + ) + ) + ; + setup_call_cleanup( + http_answer_(ResponseHandle, 404, [], ResponseStream), + format(ResponseStream, "Not Found", []), + close(ResponseStream) + ) + ) + )), http_loop(HttpListener, Handlers). send_response(ResponseHandle, http_response(StatusCode0, text(ResponseText), ResponseHeaders0)) :- default(StatusCode0, 200, StatusCode), maplist(map_header_kv_2, ResponseHeaders, ResponseHeaders0), - '$http_answer'(ResponseHandle, StatusCode, ResponseHeaders, ResponseStream0), + http_answer_(ResponseHandle, StatusCode, ResponseHeaders, ResponseStream0), open(stream(ResponseStream0), write, ResponseStream, [type(text)]), catch( call_cleanup(format(ResponseStream, "~s", [ResponseText]),close(ResponseStream)), @@ -172,9 +208,9 @@ send_response(ResponseHandle, http_response(StatusCode0, text(ResponseText), Res send_response(ResponseHandle, http_response(StatusCode0, bytes(ResponseBytes), ResponseHeaders0)) :- default(StatusCode0, 200, StatusCode), maplist(map_header_kv_2, ResponseHeaders, ResponseHeaders0), - '$http_answer'(ResponseHandle, StatusCode, ResponseHeaders, ResponseStream), + http_answer_(ResponseHandle, StatusCode, ResponseHeaders, ResponseStream), catch( - call_cleanup(format(ResponseStream, "~s", [ResponseBytes]),close(ResponseStream)), + call_cleanup(format(ResponseStream, "~s", [ResponseBytes]),close(ResponseStream)), error(existence_error(stream, _), _), true ). @@ -182,7 +218,7 @@ send_response(ResponseHandle, http_response(StatusCode0, bytes(ResponseBytes), R send_response(ResponseHandle, http_response(StatusCode0, file(Filename), ResponseHeaders0)) :- default(StatusCode0, 200, StatusCode), maplist(map_header_kv_2, ResponseHeaders, ResponseHeaders0), - '$http_answer'(ResponseHandle, StatusCode, ResponseHeaders, ResponseStream), + http_answer_(ResponseHandle, StatusCode, ResponseHeaders, ResponseStream), catch( call_cleanup( setup_call_cleanup( diff --git a/src/machine/dispatch.rs b/src/machine/dispatch.rs index 90a65cb1..e356bc38 100644 --- a/src/machine/dispatch.rs +++ b/src/machine/dispatch.rs @@ -4690,6 +4690,16 @@ impl Machine { try_or_throw!(self.machine_st, self.http_listen(), continue); step_or_fail!(self.machine_st, self.machine_st.p = self.machine_st.cp); } + &Instruction::CallHttpListenStop => { + #[cfg(feature = "http")] + try_or_throw!(self.machine_st, self.http_listen_stop(), continue); + step_or_fail!(self.machine_st, self.machine_st.p += 1); + } + &Instruction::ExecuteHttpListenStop => { + #[cfg(feature = "http")] + try_or_throw!(self.machine_st, self.http_listen_stop(), continue); + step_or_fail!(self.machine_st, self.machine_st.p = self.machine_st.cp); + } &Instruction::CallHttpAccept => { #[cfg(feature = "http")] try_or_throw!(self.machine_st, self.http_accept(), continue); diff --git a/src/machine/system_calls.rs b/src/machine/system_calls.rs index df15aeb5..becb48f2 100644 --- a/src/machine/system_calls.rs +++ b/src/machine/system_calls.rs @@ -4593,6 +4593,10 @@ impl Machine { 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_tx_timeout = warp_shutdown_tx.clone(); + let runtime = tokio::runtime::Handle::current(); let _guard = runtime.enter(); @@ -4657,13 +4661,23 @@ impl Machine { runtime.spawn(async move { match ssl_server { Some((key, cert)) => { - warp::serve(serve).tls().key(key).cert(cert).run(addr).await + let (_addr, server) = warp::serve(serve).tls().key(key).cert(cert).bind_with_graceful_shutdown(addr, async move { + warp_shutdown_rx.recv().await; + }); + + tokio::task::spawn(server); + } + None => { + let (_addr, server) = warp::serve(serve).bind_with_graceful_shutdown(addr, async move { + warp_shutdown_rx.recv().await; + }); + + tokio::task::spawn(server); } - None => warp::serve(serve).run(addr).await, } }); - let http_listener = HttpListener { incoming: rx }; + let http_listener = HttpListener { incoming: rx, warp_shutdown: warp_shutdown_tx }; let http_listener: TypedArenaPtr = arena_alloc!(http_listener, &mut self.machine_st.arena); @@ -4676,6 +4690,30 @@ impl Machine { Ok(()) } + #[cfg(feature = "http")] + #[inline(always)] + pub(crate) fn http_listen_stop(&mut self) -> CallResult { + let culprit = self.deref_register(1); + + read_heap_cell!(culprit, + (HeapCellValueTag::Cons, cons_ptr) => { + match_untyped_arena_ptr!(cons_ptr, + (ArenaHeaderTag::HttpListener, http_listener) => { + let _ = futures::executor::block_on(http_listener.warp_shutdown.send(())); + } + _ => { + unreachable!(); + } + ); + } + _ => { + unreachable!(); + } + ); + + Ok(()) + } + #[cfg(feature = "http")] #[inline(always)] pub(crate) fn http_accept(&mut self) -> CallResult {