Bound SDK integration test waits
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
This commit is contained in:
+43
-9
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user