From 8f9238c1af14adff9852b83b8b99a0e5bf640176 Mon Sep 17 00:00:00 2001 From: andrejrakic Date: Fri, 11 Sep 2026 13:08:43 +0200 Subject: [PATCH] fix: Fix the Reconnection Counter Underflow issue --- .../sdk/src/stream/monitor_connection.rs | 1 - .../sdk/tests/stream_integration_tests.rs | 22 +++++++++++++++++++ .../sdk/tests/utils/mock_websocket_server.rs | 15 +++++++++++++ 3 files changed, 37 insertions(+), 1 deletion(-) diff --git a/rust/crates/sdk/src/stream/monitor_connection.rs b/rust/crates/sdk/src/stream/monitor_connection.rs index cb11db0..8cfb656 100644 --- a/rust/crates/sdk/src/stream/monitor_connection.rs +++ b/rust/crates/sdk/src/stream/monitor_connection.rs @@ -89,7 +89,6 @@ pub(crate) async fn run_stream( } else { info!("Connection closed"); } - stats.active_connections.fetch_sub(1, Ordering::SeqCst); } _ => { warn!("Received unhandled message."); diff --git a/rust/crates/sdk/tests/stream_integration_tests.rs b/rust/crates/sdk/tests/stream_integration_tests.rs index f8c02e4..70fc97a 100644 --- a/rust/crates/sdk/tests/stream_integration_tests.rs +++ b/rust/crates/sdk/tests/stream_integration_tests.rs @@ -316,6 +316,28 @@ async fn test_stream_ha_x_cll_origin_header() { stream.close().await.expect("Failed to close stream"); } +#[tokio::test] +async fn test_stream_ha_graceful_close_reconnect_preserves_active_connections() { + // Regression test for graceful close; must decrement active_connections only once, not twice. + let (mock_server, stream, _) = prepare_scenario().await; + + let initial_active = stream.get_stats().active_connections; + assert_eq!(initial_active, NUMBER_OF_CONNECTIONS); + + // Server-initiated graceful close: sends a Close frame to every connected client. + // The listener stays up so clients can reconnect immediately. + mock_server.close_connections().await; + + // Allow time for all connections to close and reconnect. + sleep(Duration::from_millis(500)).await; + + let stats = stream.get_stats(); + assert_eq!( + stats.active_connections, NUMBER_OF_CONNECTIONS, + "active_connections should be restored to its initial value after graceful close + reconnect" + ); +} + #[tokio::test] #[ignore] // Ignored because it takes a while to complete. To run it, use this command: cargo test -- --ignored async fn test_stream_ha_max_reconnection_attempts() { diff --git a/rust/crates/sdk/tests/utils/mock_websocket_server.rs b/rust/crates/sdk/tests/utils/mock_websocket_server.rs index 1d6ef58..1fed47c 100644 --- a/rust/crates/sdk/tests/utils/mock_websocket_server.rs +++ b/rust/crates/sdk/tests/utils/mock_websocket_server.rs @@ -13,6 +13,7 @@ use tokio_tungstenite::{ enum ServerCommand { Send(Vec), DropConnections, + CloseConnections, } #[derive(Clone)] @@ -96,6 +97,13 @@ impl MockWebSocketServer { println!("Dropping all client connections"); clients_command.lock().await.clear(); } + ServerCommand::CloseConnections => { + println!("Sending graceful Close frame to all client connections"); + let clients = clients_command.lock().await; + for client in clients.iter() { + let _ = client.send(Message::Close(None)).await; + } + } } } }); @@ -124,6 +132,13 @@ impl MockWebSocketServer { .await; } + pub async fn close_connections(&self) { + let _ = self + .command_sender + .send(ServerCommand::CloseConnections) + .await; + } + pub async fn shutdown(&self) { self.shutdown_notify.notify_waiters(); }