From 80986904e16e95649117586653c0c48da8aa1e52 Mon Sep 17 00:00:00 2001 From: RoboClaw Date: Wed, 30 Sep 2026 09:14:20 -0700 Subject: [PATCH] improve(tests): replace Rust session fixture delays with signals (#161614) Keep real socket event ordering and activity publication while replacing timing guesses with response barriers and fixture readiness. Co-authored-by: steipete <58493+steipete@users.noreply.github.com> --- .../openclaw-gateway-client/tests/session.rs | 37 ++++++++++++++++--- 1 file changed, 32 insertions(+), 5 deletions(-) diff --git a/crates/openclaw-gateway-client/tests/session.rs b/crates/openclaw-gateway-client/tests/session.rs index 907523cd83c3..5de3fba1fb1f 100644 --- a/crates/openclaw-gateway-client/tests/session.rs +++ b/crates/openclaw-gateway-client/tests/session.rs @@ -662,6 +662,7 @@ async fn raw_event_subscription_closes_with_the_session() { async fn default_buffer_retains_256_small_events() { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); + let (assertions_done_tx, assertions_done_rx) = tokio::sync::oneshot::channel(); let server = tokio::spawn(async move { let (tcp, _) = listener.accept().await.unwrap(); let mut socket = accept_async(tcp).await.unwrap(); @@ -688,7 +689,8 @@ async fn default_buffer_retains_256_small_events() { ) .await; } - tokio::time::sleep(Duration::from_millis(100)).await; + acknowledge_buffered_events(&mut socket).await; + let _ = assertions_done_rx.await; }); let session = GatewayClient::connect( @@ -697,12 +699,13 @@ async fn default_buffer_retains_256_small_events() { ) .await .unwrap(); - tokio::time::sleep(Duration::from_millis(25)).await; + session.request("test.buffered", Value::Null).await.unwrap(); for seq in 0..256 { let event = session.next_event().await.unwrap(); assert_eq!(event.event, "node.small"); assert_eq!(event.seq, Some(seq)); } + assertions_done_tx.send(()).unwrap(); server.await.unwrap(); } @@ -710,6 +713,7 @@ async fn default_buffer_retains_256_small_events() { async fn oversized_retained_event_lags_without_closing_the_session() { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); + let (assertions_done_tx, assertions_done_rx) = tokio::sync::oneshot::channel(); let server = tokio::spawn(async move { let (tcp, _) = listener.accept().await.unwrap(); let mut socket = accept_async(tcp).await.unwrap(); @@ -742,7 +746,8 @@ async fn oversized_retained_event_lags_without_closing_the_session() { json!({"type":"event", "event":"node.after-large", "seq":2}), ) .await; - tokio::time::sleep(Duration::from_millis(50)).await; + acknowledge_buffered_events(&mut socket).await; + let _ = assertions_done_rx.await; }); let config = GatewayClientConfig::new(format!("ws://{address}")) @@ -754,6 +759,7 @@ async fn oversized_retained_event_lags_without_closing_the_session() { }) .await .unwrap(); + session.request("test.buffered", Value::Null).await.unwrap(); assert!(matches!( session.next_event().await, Err(ClientError::EventLagged(1)) @@ -762,6 +768,8 @@ async fn oversized_retained_event_lags_without_closing_the_session() { session.next_event().await.unwrap().event, "node.after-large" ); + assert!(!session.is_closed()); + assertions_done_tx.send(()).unwrap(); server.await.unwrap(); } @@ -910,6 +918,7 @@ async fn drains_a_queued_event_before_reporting_disconnect() { async fn surfaces_websocket_ping_as_transport_activity() { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); + let (subscribed_tx, subscribed_rx) = tokio::sync::oneshot::channel(); let server = tokio::spawn(async move { let (tcp, _) = listener.accept().await.unwrap(); let mut socket = accept_async(tcp).await.unwrap(); @@ -929,7 +938,7 @@ async fn surfaces_websocket_ping_as_transport_activity() { }), ) .await; - tokio::time::sleep(Duration::from_millis(20)).await; + subscribed_rx.await.unwrap(); socket .send(Message::Ping(vec![1, 2, 3].into())) .await @@ -945,6 +954,7 @@ async fn surfaces_websocket_ping_as_transport_activity() { .await .unwrap(); let mut activity = session.subscribe_transport_activity(); + subscribed_tx.send(()).unwrap(); tokio::time::timeout(Duration::from_secs(1), activity.changed()) .await .expect("ping activity timeout") @@ -1022,6 +1032,7 @@ async fn websocket_ping_remains_available_when_rpc_capacity_is_full() { async fn malformed_idle_text_is_activity_and_does_not_close_the_session() { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); + let (subscribed_tx, subscribed_rx) = tokio::sync::oneshot::channel(); let server = tokio::spawn(async move { let (tcp, _) = listener.accept().await.unwrap(); let mut socket = accept_async(tcp).await.unwrap(); @@ -1042,7 +1053,7 @@ async fn malformed_idle_text_is_activity_and_does_not_close_the_session() { }), ) .await; - tokio::time::sleep(Duration::from_millis(20)).await; + subscribed_rx.await.unwrap(); socket .send(Message::Text("not gateway json".into())) .await @@ -1065,6 +1076,7 @@ async fn malformed_idle_text_is_activity_and_does_not_close_the_session() { .await .unwrap(); let mut activity = session.subscribe_transport_activity(); + subscribed_tx.send(()).unwrap(); tokio::time::timeout(Duration::from_secs(1), activity.changed()) .await .expect("malformed idle frame activity timeout") @@ -1274,6 +1286,21 @@ async fn pinned_trust_rejects_plaintext_before_connecting() { assert!(matches!(result, Err(ClientError::Tls(_)))); } +async fn acknowledge_buffered_events(socket: &mut tokio_tungstenite::WebSocketStream) +where + S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin, +{ + // The response follows every event on the same socket, so the client must + // process those events before completing the request and starting to drain. + let barrier = receive_json(socket).await; + assert_eq!(barrier["method"], "test.buffered"); + send_json( + socket, + json!({"type":"res", "id":barrier["id"], "ok":true, "payload":null}), + ) + .await; +} + async fn send_json(socket: &mut tokio_tungstenite::WebSocketStream, value: Value) where S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin,