From 367af4de8208cd853f4a74950001f7e5594f9f57 Mon Sep 17 00:00:00 2001 From: DuProcess <273172371+DuProcess@users.noreply.github.com> Date: Fri, 17 Jul 2026 18:26:45 -0400 Subject: [PATCH] Bound SDK integration test waits --- src/bin/dosh-client.rs | 10 ++++---- src/server.rs | 52 ++++++++++++++++++++++++++++++++++-------- 2 files changed, 48 insertions(+), 14 deletions(-) diff --git a/src/bin/dosh-client.rs b/src/bin/dosh-client.rs index eac71a1..4837842 100644 --- a/src/bin/dosh-client.rs +++ b/src/bin/dosh-client.rs @@ -4476,7 +4476,7 @@ fn latest_release_tag(repo: &str) -> Result> { } let effective = String::from_utf8_lossy(&output.stdout); let effective = effective.trim(); - Ok(release_tag_from_effective_url(web, &effective)) + Ok(release_tag_from_effective_url(web, effective)) } fn release_tag_download_url(repo: &str, tag: &str, artifact: &str) -> Option { @@ -5523,12 +5523,12 @@ fn resolve_forward_agent_endpoint(forward_agent: bool) -> Result #[cfg(unix)] { - return match std::env::var_os("SSH_AUTH_SOCK") { + match std::env::var_os("SSH_AUTH_SOCK") { Some(path) if !path.is_empty() => Ok(Some(PathBuf::from(path))), _ => Err(anyhow!( "agent forwarding requested but SSH_AUTH_SOCK is not set" )), - }; + } } #[cfg(windows)] @@ -6150,10 +6150,10 @@ fn local_username() -> String { local_username_from_env(|name| std::env::var(name).ok()) } -fn local_username_from_env(mut get: impl FnMut(&str) -> Option) -> String { +fn local_username_from_env(get: impl FnMut(&str) -> Option) -> String { ["USER", "USERNAME"] .into_iter() - .filter_map(|name| get(name)) + .filter_map(get) .find(|value| !value.is_empty()) .unwrap_or_else(|| "unknown".to_string()) } diff --git a/src/server.rs b/src/server.rs index e73f7b3..afe564b 100644 --- a/src/server.rs +++ b/src/server.rs @@ -578,20 +578,43 @@ mod tests { .user("sdk-user") .service("echo") .connect(); - let accept = async { + let handshake = async { + tokio::pin!(connect); + let mut connected = None; + let mut accepted = None; loop { - if let DoshServerEvent::Accepted(accepted) = server.recv().await.unwrap() { - break accepted; + tokio::select! { + result = &mut connect, if connected.is_none() => { + connected = Some(result?); + } + event = server.recv(), if accepted.is_none() => { + if let DoshServerEvent::Accepted(connection) = event? { + accepted = Some(connection); + } + } + } + if connected.is_some() && accepted.is_some() { + return Ok::<_, anyhow::Error>(( + connected.take().expect("connected result checked"), + accepted.take().expect("accepted result checked"), + )); } } }; - let (connected, accepted) = tokio::join!(connect, accept); - let mut client_transport = connected.unwrap().into_transport(); + let (connected, accepted) = tokio::time::timeout(Duration::from_secs(3), handshake) + .await + .expect("SDK client/server authentication timed out") + .unwrap(); + let mut client_transport = connected.into_transport(); let conn_id = accepted.conn_id; assert_eq!(accepted.services, vec!["echo".to_string()]); let stream_id = client_transport.open_service("echo").await.unwrap(); - match server.recv().await.unwrap() { + match tokio::time::timeout(Duration::from_secs(3), server.recv()) + .await + .expect("server timed out waiting for stream open") + .unwrap() + { DoshServerEvent::Session { event: SessionEvent::Stream(TransportEvent::Open(open)), .. @@ -602,7 +625,10 @@ mod tests { other => panic!("unexpected event {other:?}"), } assert!(matches!( - client_transport.recv().await.unwrap(), + tokio::time::timeout(Duration::from_secs(3), client_transport.recv()) + .await + .expect("client timed out waiting for stream-open acknowledgement") + .unwrap(), SessionEvent::Stream(TransportEvent::OpenOk { .. }) )); @@ -611,7 +637,11 @@ mod tests { .await .unwrap(); loop { - match server.recv().await.unwrap() { + match tokio::time::timeout(Duration::from_secs(3), server.recv()) + .await + .expect("server timed out waiting for stream data") + .unwrap() + { DoshServerEvent::Session { event: SessionEvent::Stream(TransportEvent::Data(data)), .. @@ -628,7 +658,11 @@ mod tests { } } loop { - match client_transport.recv().await.unwrap() { + match tokio::time::timeout(Duration::from_secs(3), client_transport.recv()) + .await + .expect("client timed out waiting for stream response") + .unwrap() + { SessionEvent::Stream(TransportEvent::Data(data)) => { assert_eq!(data.chunks, vec![b"pong".to_vec()]); break;