use std::sync::Arc;
use std::time::Duration;
use futures::StreamExt;
use velo::messenger::Messenger;
use velo::streaming::velo_transport::VeloFrameTransport;
use velo::streaming::{AnchorManagerBuilder, FrameTransport, StreamAnchorHandle, StreamFrame};
use velo::transports::tcp::TcpTransportBuilder;
use velo_ext::WorkerId;
const FRAMES: u32 = 200;
fn new_tcp_transport() -> Arc<velo::transports::tcp::TcpTransport> {
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
Arc::new(
TcpTransportBuilder::new()
.from_listener(listener)
.unwrap()
.build()
.unwrap(),
)
}
async fn make_two_messengers() -> (Arc<Messenger>, Arc<Messenger>) {
let t1 = new_tcp_transport();
let t2 = new_tcp_transport();
let m1 = Messenger::builder()
.add_transport(t1)
.build()
.await
.expect("create messenger 1");
let m2 = Messenger::builder()
.add_transport(t2)
.build()
.await
.expect("create messenger 2");
let p1 = m1.peer_info();
let p2 = m2.peer_info();
m2.register_peer(p1).expect("register m1 on m2");
m1.register_peer(p2).expect("register m2 on m1");
tokio::time::sleep(Duration::from_millis(200)).await;
(m1, m2)
}
#[tokio::test(flavor = "multi_thread")]
async fn test_remote_stream_preserves_send_order() {
let (messenger_a, messenger_b) = make_two_messengers().await;
let worker_id_a = messenger_a.instance_id().worker_id();
let worker_id_b = messenger_b.instance_id().worker_id();
let vft_a =
Arc::new(VeloFrameTransport::new(Arc::clone(&messenger_a), None).expect("VFT worker A"));
let am_a: Arc<velo::streaming::AnchorManager> = Arc::new(
AnchorManagerBuilder::default()
.worker_id(worker_id_a)
.transport(Arc::clone(&vft_a) as Arc<dyn FrameTransport>)
.build()
.expect("AM worker A"),
);
am_a.register_handlers(Arc::clone(&messenger_a))
.expect("register_handlers on worker A");
let mut anchor_stream = am_a.create_anchor::<u32>();
let handle = anchor_stream.handle();
let handle_raw: u128 = handle.as_u128();
let hi = (handle_raw >> 64) as u64;
let lo = handle_raw as u64;
let handle_transferred = StreamAnchorHandle::pack(WorkerId::from_u64(hi), lo);
let vft_b =
Arc::new(VeloFrameTransport::new(Arc::clone(&messenger_b), None).expect("VFT worker B"));
let am_b: Arc<velo::streaming::AnchorManager> = Arc::new(
AnchorManagerBuilder::default()
.worker_id(worker_id_b)
.transport(Arc::clone(&vft_b) as Arc<dyn FrameTransport>)
.build()
.expect("AM worker B"),
);
am_b.register_handlers(Arc::clone(&messenger_b))
.expect("register_handlers on worker B");
let sender = am_b
.attach_stream_anchor::<u32>(handle_transferred)
.await
.expect("remote attach must succeed");
let send_task = tokio::spawn(async move {
for i in 0u32..FRAMES {
sender.send(i).await.expect("send item");
}
tokio::time::sleep(Duration::from_millis(50)).await;
sender.finalize().expect("finalize");
});
let collect = async {
let mut items = Vec::with_capacity(FRAMES as usize);
while let Some(frame) = anchor_stream.next().await {
match frame.expect("no stream error") {
StreamFrame::Item(v) => {
items.push(v);
tokio::task::yield_now().await;
}
StreamFrame::Finalized => break,
other => panic!("unexpected frame: {:?}", other),
}
}
items
};
let items = tokio::time::timeout(Duration::from_secs(15), collect)
.await
.expect("timed out waiting for items");
send_task.await.expect("send task panicked");
assert_eq!(
items.len(),
FRAMES as usize,
"expected exactly {} items, got {}",
FRAMES,
items.len()
);
let expected: Vec<u32> = (0..FRAMES).collect();
if items != expected {
let first_oos = items
.iter()
.zip(expected.iter())
.position(|(got, want)| got != want)
.unwrap_or(items.len());
panic!(
"frames out of order at index {}: got {:?}..., want {:?}...",
first_oos,
&items[first_oos..(first_oos + 5).min(items.len())],
&expected[first_oos..(first_oos + 5).min(expected.len())],
);
}
}