use std::{sync::Arc, time::Duration};
use bytes::Bytes;
use pocketscion::util::topologies::{IA132, IA212, UnderlayType, minimal::two_path_topology};
use scion_stack::stack::ScionStackBuilder;
use snap_tokens::v0::dummy_snap_token;
use test_log::test;
use tokio::{
sync::Barrier,
time::{sleep, timeout},
};
use tracing::info;
#[test(tokio::test)]
#[ntest::timeout(5_000)]
async fn should_failover_on_link_error() {
let ps_handle = two_path_topology(UnderlayType::Snap).await;
let sender_stack = ScionStackBuilder::new()
.with_endhost_api(ps_handle.endhost_api(IA132).unwrap())
.with_auth_token(dummy_snap_token())
.build()
.await
.unwrap();
let receiver_stack = ScionStackBuilder::new()
.with_endhost_api(ps_handle.endhost_api(IA212).unwrap())
.with_auth_token(dummy_snap_token())
.build()
.await
.unwrap();
let sender_socket = Arc::new(sender_stack.bind(None).await.unwrap());
let sender_addr = sender_socket.local_addr();
let receiver_socket = receiver_stack.bind(None).await.unwrap();
let receiver_addr = receiver_socket.local_addr();
let test_data = Bytes::from("Hello, World!");
let mut recv_buffer = [0u8; 1024];
let failover_send_barrier = Arc::new(Barrier::new(2));
let sender_socket_clone = sender_socket.clone();
let sender_recv_task = tokio::spawn(async move {
let mut recv_buffer = [0u8; 1024];
let (..) = sender_socket_clone
.recv_from(&mut recv_buffer)
.await
.unwrap();
panic!("Sender should not receive udp packets");
});
let sender_task = tokio::spawn({
let failover_send_barrier = failover_send_barrier.clone();
let test_data = test_data.clone();
async move {
sender_socket
.send_to(test_data.as_ref(), receiver_addr)
.await
.unwrap_or_else(|_| {
panic!("error sending from {sender_addr:?} to {receiver_addr:?}")
});
failover_send_barrier.wait().await;
sender_socket
.send_to(test_data.as_ref(), receiver_addr)
.await
.unwrap_or_else(|_| {
panic!("error sending from {sender_addr:?} to {receiver_addr:?}")
});
loop {
sleep(Duration::from_millis(100)).await;
sender_socket
.send_to(test_data.as_ref(), receiver_addr)
.await
.unwrap_or_else(|_| {
panic!("error sending from {sender_addr:?} to {receiver_addr:?}")
});
}
}
});
let (_, source, path) = receiver_socket
.recv_from_with_path(&mut recv_buffer)
.await
.unwrap();
assert_eq!(
source, sender_addr,
"receiver should receive packets from the sender"
);
let egress = path
.first_egress_interface()
.expect("path should have first hop egress interface");
ps_handle
.runtime
.set_link_state(egress.isd_asn, egress.id, false)
.unwrap();
failover_send_barrier.wait().await;
let mut recv_buffer = [0u8; 1024];
let (_size, _addr, new_path) = timeout(
Duration::from_millis(500),
receiver_socket.recv_from_with_path(&mut recv_buffer),
)
.await
.expect("should not time out waiting for packet after failover")
.expect("should receive packet after failover");
info!(old_path = ?path, new_path = ?new_path, "path changed?");
assert_ne!(path, new_path, "should use a different path after failover");
sender_task.abort();
sender_recv_task.abort();
}