From fac268226805a9133c23178e46330f1f963eb044 Mon Sep 17 00:00:00 2001 From: ssmithx Date: Fri, 24 Jul 2026 12:26:06 +0000 Subject: [PATCH] fix(server): stop federation peer congestion from freezing the web UI MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The main (5678) and FIPS peer/federation (5679) listeners shared one 1024-permit connection semaphore, and accept_loop awaited that permit before spawning the connection handler — inside the accept loop itself. When federation peer connections piled up in CLOSE-WAIT without ever completing, they exhausted the shared pool and froze accept() entirely, taking the web UI down with them (production outage, 2026-07-24: no login possible, auth.isSetup/server.echo hung indefinitely even hitting the backend directly). Two changes: - Acquire the permit inside the spawned task (with a bounded 30s wait) instead of in the accept loop, so a saturated pool can no longer block accept() from continuing to accept new connections. - Give the peer listener its own semaphore, separate from the main listener's, so peer/federation congestion can never starve local web UI connections again. Co-Authored-By: Claude Sonnet 5 --- core/archipelago/src/server.rs | 52 ++++++++++++++++++++++++++++++---- 1 file changed, 47 insertions(+), 5 deletions(-) diff --git a/core/archipelago/src/server.rs b/core/archipelago/src/server.rs index 67d58a17..3ba6f038 100644 --- a/core/archipelago/src/server.rs +++ b/core/archipelago/src/server.rs @@ -976,13 +976,21 @@ impl Server { main_addr: SocketAddr, shutdown: impl std::future::Future, ) -> Result<()> { - let active_connections = Arc::new(tokio::sync::Semaphore::new(1024)); + // Separate pools per listener. Federation/peer connections (fips0, + // often over Tor, from nodes we don't control the behavior of) used + // to share one pool with the local web UI listener — when peer + // connections piled up in CLOSE-WAIT without ever completing, they + // starved the shared pool and took the web UI down with them + // (production outage, 2026-07-24). Peer congestion must never be + // able to block a local login. + let main_connections = Arc::new(tokio::sync::Semaphore::new(1024)); + let peer_connections = Arc::new(tokio::sync::Semaphore::new(256)); let (tx, rx_main) = tokio::sync::watch::channel(false); let main_task = tokio::spawn(accept_loop( self.api_handler.clone(), TcpListener::bind(main_addr).await?, - active_connections.clone(), + main_connections.clone(), false, // main listener: no path filter rx_main, main_addr, @@ -992,7 +1000,7 @@ impl Server { // restart when fips0 comes up after onboarding. let peer_task = tokio::spawn(peer_late_bind_loop( self.api_handler.clone(), - active_connections.clone(), + peer_connections.clone(), tx.subscribe(), )); @@ -1003,7 +1011,9 @@ impl Server { // Wait up to 5s for in-flight requests. let drain_start = std::time::Instant::now(); let drain_timeout = std::time::Duration::from_secs(5); - while active_connections.available_permits() < 1024 { + while main_connections.available_permits() < 1024 + || peer_connections.available_permits() < 256 + { if drain_start.elapsed() > drain_timeout { warn!("Drain timeout reached, forcing shutdown"); break; @@ -1096,6 +1106,11 @@ pub fn is_peer_allowed_path(path: &str) -> bool { || path.starts_with("/content/") } +/// How long a freshly-accepted connection will wait for a connection-pool +/// permit before it's dropped. Bounds worst-case fd/task growth if the pool +/// is ever genuinely saturated; under normal load this never triggers. +const PERMIT_ACQUIRE_TIMEOUT: Duration = Duration::from_secs(30); + async fn accept_loop( handler: Arc, listener: TcpListener, @@ -1115,8 +1130,35 @@ async fn accept_loop( } }; let handler = handler.clone(); - let permit = active_connections.clone().acquire_owned().await; + let active_connections = active_connections.clone(); + // Acquire the permit *inside* the spawned task, not here. + // This loop must never block on anything but accept()/shutdown: + // the main (5678) and FIPS peer (5679) listeners share one + // semaphore, and a single slow/hung connection holding the + // last permit used to freeze this whole loop — including for + // the OTHER listener — since accept() couldn't be called + // again until a permit freed up. That took down the entire + // web UI in production (2026-07-24) when federation peer + // connections piled up. Now a saturated pool just delays + // (and, past PERMIT_ACQUIRE_TIMEOUT, drops) individual + // connections instead of wedging the acceptor itself. tokio::spawn(async move { + let permit = match tokio::time::timeout( + PERMIT_ACQUIRE_TIMEOUT, + active_connections.acquire_owned(), + ) + .await + { + Ok(Ok(permit)) => permit, + Ok(Err(_)) => return, // semaphore closed during shutdown + Err(_) => { + warn!( + "{} connection from {} dropped — connection pool saturated for {}s", + local_addr, peer_addr, PERMIT_ACQUIRE_TIMEOUT.as_secs() + ); + return; + } + }; let _permit = permit; let service = service_fn(move |mut req: hyper::Request| { let handler = handler.clone();