Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 0 additions & 1 deletion rust/crates/sdk/src/stream/monitor_connection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.");
Expand Down
22 changes: 22 additions & 0 deletions rust/crates/sdk/tests/stream_integration_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down
15 changes: 15 additions & 0 deletions rust/crates/sdk/tests/utils/mock_websocket_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ use tokio_tungstenite::{
enum ServerCommand {
Send(Vec<u8>),
DropConnections,
CloseConnections,
}

#[derive(Clone)]
Expand Down Expand Up @@ -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;
}
}
}
}
});
Expand Down Expand Up @@ -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();
}
Expand Down
Loading