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>
This commit is contained in:
RoboClaw 2026-09-30 09:14:20 -07:00 • committed by GitHub
parent 63b43fc09e
commit 80986904e1
No known key found for this signature in database
GPG key ID: B5690EEEBB952194

View file

@ -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<S>(socket: &mut tokio_tungstenite::WebSocketStream<S>)
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<S>(socket: &mut tokio_tungstenite::WebSocketStream<S>, value: Value)
where
S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin,