From e89bb2ace97b5a127bf288a46cd73eeac6ee9e77 Mon Sep 17 00:00:00 2001 From: meh Date: Wed, 23 Sep 2026 11:49:52 +0700 Subject: [PATCH 1/2] Let the caller say when a server stops listen() claimed SIGTERM and SIGINT for the whole process, and nagoya allows one waiter per signal number, so a second server in the same process failed with EBUSY after it had bound its port. Its listener was dropped and a client connecting to that port waited forever. Any process that serves twice hits it, and every test harness that starts a server per test does: agency-proxy's WebSocket suite hung on exactly this. listen_until takes the stop condition as a future and registers no signal. listen() is unchanged and delegates to it with the process signals, so current callers keep their behaviour. The one stop future is now polled across loop iterations rather than rebuilt each time, so a caller's future is never dropped half-way. --- README.md | 12 ++++++++ src/libs/ws/server.rs | 64 ++++++++++++++++++++++++++++++++++++------- 2 files changed, 66 insertions(+), 10 deletions(-) diff --git a/README.md b/README.md index 7b86e06..ab42318 100644 --- a/README.md +++ b/README.md @@ -336,6 +336,18 @@ server.enable_mcp( server.listen() ``` +`listen()` blocks the calling thread and stops on SIGTERM or SIGINT, claiming +both for the process while it runs, so a second server in the same process +fails with `EBUSY`. A process that owns its own signals, or runs more than one +server (a test per server, say), stops it with a future instead: + +```rust +let (stop, stopped) = futures::channel::oneshot::channel::<()>(); +std::thread::spawn(move || server.listen_until(async move { let _ = stopped.await; })); +// ... later +let _ = stop.send(()); +``` + Behavior notes: - **Frame detection** — a frame is treated as JSON-RPC iff it carries a diff --git a/src/libs/ws/server.rs b/src/libs/ws/server.rs index cd3b534..5c0fed4 100644 --- a/src/libs/ws/server.rs +++ b/src/libs/ws/server.rs @@ -429,7 +429,36 @@ impl WebsocketServer { /// `TcpListener` has no `try_clone` and no `SO_REUSEPORT` bind, and a bound /// listener cannot be adopted by a second reactor. `shard_count` still reads the /// operator's intent, and the server says so when it cannot honour it. + /// + /// # Signals + /// + /// This stops on SIGTERM or SIGINT, and claims both for the process while it + /// runs: nagoya allows one waiter per signal number, so a second server in the + /// same process fails with `EBUSY`. A process that owns its signals, or runs + /// more than one server, calls [`Self::listen_until`] instead. pub fn listen(self) -> Result<()> { + self.listen_on(None::>) + } + + /// [`Self::listen`], stopping when `stop` resolves instead of on a signal. + /// + /// No signal is registered. Signals belong to a process, not to a library + /// server inside it, and taking them here is what made a second server in one + /// process fail: an application that serves twice, and every test harness + /// that starts a server per test. `stop` is polled on this server's own + /// thread, so it has to be a future that needs no particular runtime, such as + /// a oneshot receiver or [`crate::libs::signal::Shutdown::cancelled`]. + pub fn listen_until(self, stop: F) -> Result<()> + where + F: std::future::Future + 'static, + { + self.listen_on(Some(stop)) + } + + fn listen_on(self, stop: Option) -> Result<()> + where + F: std::future::Future + 'static, + { self.validate_protocol_mode()?; self.refuse_tls_config()?; debug!(ws_server = true, "Listening on {}", self.config.address); @@ -454,7 +483,7 @@ impl WebsocketServer { // `config.insecure` branch that chose between them are gone: this server // serves plain `ws://` and nothing else, with TLS terminated at the edge. block_on_with(&reactor, async move { - self.listen_impl(Arc::new(listener), &handle).await + self.listen_impl(Arc::new(listener), &handle, stop).await }) } @@ -493,11 +522,16 @@ impl WebsocketServer { /// grows one entry per connection ever accepted and panics past four billion. /// It fits a fixed population of tasks, which a server's connections are not. /// `FuturesUnordered` drops what finishes and wakes in O(woken) just the same. - async fn listen_impl( + async fn listen_impl( self, listener: Arc, handle: &Handle, - ) -> Result<()> { + stop: Option, + ) -> Result<()> + where + T: ConnectionListener + 'static, + F: std::future::Future + 'static, + { use futures::StreamExt; use futures::future::{Either, FutureExt, select}; use futures::stream::FuturesUnordered; @@ -520,17 +554,27 @@ impl WebsocketServer { ); } - // Registered on this reactor, and awaited on it below. A `Signal` only - // fires while its own reactor is polled, so creating it anywhere else + // The caller's stop, or the process signals when it gave none. Signals are + // registered on this reactor, and awaited on it below: a `Signal` only + // fires while its own reactor is polled, so creating one anywhere else // would be creating a wait that never ends. - let (mut sigterm, mut sigint) = crate::libs::signal::init_signals(handle)?; + let mut stop: futures::future::LocalBoxFuture<'static, ()> = match stop { + Some(stop) => stop.boxed_local(), + None => { + let (mut sigterm, mut sigint) = crate::libs::signal::init_signals(handle)?; + async move { crate::libs::signal::wait_for_signals(&mut sigterm, &mut sigint).await } + .boxed_local() + } + }; let mut connections = FuturesUnordered::new(); loop { - // Shutdown is the outermost left arm, so a pending signal is not stuck - // behind an accept or a live connection that is also ready. - let shutdown = crate::libs::signal::wait_for_signals(&mut sigterm, &mut sigint); + // Shutdown is the outermost left arm, so a pending stop is not stuck + // behind an accept or a live connection that is also ready. The one + // stop future is polled across iterations rather than rebuilt, so a + // caller's future is never dropped half-way. + let shutdown = stop.as_mut(); let accepted = listener.accept(); - futures::pin_mut!(shutdown, accepted); + futures::pin_mut!(accepted); let accepted = if connections.is_empty() { match select(shutdown, accepted).await { Either::Left(_) => break, From fc735cc01f7bd7d1d9070af9898634bebd94b81d Mon Sep 17 00:00:00 2001 From: meh Date: Wed, 23 Sep 2026 11:49:52 +0700 Subject: [PATCH 2/2] Release 3.3.0 Adds WebsocketServer::listen_until. --- Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Cargo.toml b/Cargo.toml index 77e25bd..d2b2ed8 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "endpoint-libs" -version = "3.2.1" +version = "3.3.0" edition = "2024" authors = ["Veon "] description = "Launch MCP services fast: describe endpoints once in RON, and endpoint-gen generates the Rust models, docs and MCP tool schemas that this crate serves over WebSocket RPC, with roles and typed errors built in."