#![cfg(feature = "queue-messenger")]
use std::sync::Arc;
use std::time::Duration;
use bytes::Bytes;
use velo::messenger::Messenger;
use velo::transports::uds::UdsTransportBuilder;
use velo::queue::WorkQueueBackend;
use velo::queue::backends::messenger::{MessengerQueueBackend, MessengerQueueConfig};
use velo::queue::options::NextOptions;
async fn setup_two_messengers() -> (Arc<Messenger>, Arc<Messenger>) {
let socket_a =
std::env::temp_dir().join(format!("velo-queue-test-{}.sock", uuid::Uuid::new_v4()));
let socket_b =
std::env::temp_dir().join(format!("velo-queue-test-{}.sock", uuid::Uuid::new_v4()));
let transport_a = Arc::new(
UdsTransportBuilder::new()
.socket_path(&socket_a)
.build()
.expect("build transport A"),
);
let transport_b = Arc::new(
UdsTransportBuilder::new()
.socket_path(&socket_b)
.build()
.expect("build transport B"),
);
let m_a = Messenger::builder()
.add_transport(transport_a)
.build()
.await
.expect("build messenger A");
let m_b = Messenger::builder()
.add_transport(transport_b)
.build()
.await
.expect("build messenger B");
let peer_a = m_a.peer_info();
let peer_b = m_b.peer_info();
m_a.register_peer(peer_b).expect("register B on A");
m_b.register_peer(peer_a).expect("register A on B");
tokio::time::sleep(Duration::from_millis(200)).await;
(m_a, m_b)
}
#[tokio::test(flavor = "multi_thread")]
async fn test_messenger_local_send_recv() {
let (m_a, _m_b) = setup_two_messengers().await;
let backend = MessengerQueueBackend::new(
Arc::clone(&m_a),
m_a.instance_id(),
MessengerQueueConfig::default(),
);
let sender = backend.sender("local-test").await.unwrap();
let receiver = backend.receiver("local-test").await.unwrap();
let data = Bytes::from_static(b"hello local");
sender.send(data.clone()).await.unwrap();
let received = receiver.recv().await.unwrap();
assert_eq!(received, Some(data));
}
#[tokio::test(flavor = "multi_thread")]
async fn test_messenger_local_batch() {
let (m_a, _m_b) = setup_two_messengers().await;
let backend = MessengerQueueBackend::new(
Arc::clone(&m_a),
m_a.instance_id(),
MessengerQueueConfig::default(),
);
let sender = backend.sender("local-batch").await.unwrap();
let receiver = backend.receiver("local-batch").await.unwrap();
for i in 0u8..5 {
sender.send(Bytes::from(vec![i])).await.unwrap();
}
let opts = NextOptions::new()
.batch_size(3)
.timeout(Duration::from_millis(500));
let batch = receiver.recv_batch(&opts).await.unwrap();
assert_eq!(batch.len(), 3);
assert_eq!(batch[0], Bytes::from(vec![0u8]));
assert_eq!(batch[1], Bytes::from(vec![1u8]));
assert_eq!(batch[2], Bytes::from(vec![2u8]));
}
#[tokio::test(flavor = "multi_thread")]
async fn test_messenger_try_send_delivers() {
let (m_a, _m_b) = setup_two_messengers().await;
let backend = MessengerQueueBackend::new(
Arc::clone(&m_a),
m_a.instance_id(),
MessengerQueueConfig::default(),
);
let sender = backend.sender("try-send-test").await.unwrap();
let receiver = backend.receiver("try-send-test").await.unwrap();
sender.try_send(Bytes::from_static(b"try-sent")).unwrap();
let opts = NextOptions::new()
.batch_size(1)
.timeout(Duration::from_secs(2));
let batch = receiver.recv_batch(&opts).await.unwrap();
assert_eq!(batch.len(), 1);
assert_eq!(batch[0], Bytes::from_static(b"try-sent"));
}
#[tokio::test(flavor = "multi_thread")]
async fn test_messenger_remote_send_recv() {
let (m_a, m_b) = setup_two_messengers().await;
let backend_b = MessengerQueueBackend::new(
Arc::clone(&m_b),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let receiver = backend_b.receiver("remote-test").await.unwrap();
let backend_a = MessengerQueueBackend::new(
Arc::clone(&m_a),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let sender = backend_a.sender("remote-test").await.unwrap();
let data = Bytes::from_static(b"hello remote");
sender.send(data.clone()).await.unwrap();
let opts = NextOptions::new()
.batch_size(1)
.timeout(Duration::from_secs(5));
let batch = receiver.recv_batch(&opts).await.unwrap();
assert_eq!(batch.len(), 1);
assert_eq!(batch[0], data);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_messenger_remote_send_after_stale_handler_cache_refreshes() {
let (m_a, m_b) = setup_two_messengers().await;
let handlers = m_a.available_handlers(m_b.instance_id()).await.unwrap();
assert!(
!handlers.contains(&"velo.queue.rpc".to_string()),
"queue service should not be registered before the backend is created"
);
let backend_b = MessengerQueueBackend::new(
Arc::clone(&m_b),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let receiver = backend_b.receiver("stale-cache").await.unwrap();
let backend_a = MessengerQueueBackend::new(
Arc::clone(&m_a),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let sender = backend_a.sender("stale-cache").await.unwrap();
sender.send(Bytes::from_static(b"refresh")).await.unwrap();
assert_eq!(
receiver.recv().await.unwrap(),
Some(Bytes::from_static(b"refresh"))
);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_messenger_remote_recv_waits() {
let (m_a, m_b) = setup_two_messengers().await;
let backend_b = MessengerQueueBackend::new(
Arc::clone(&m_b),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let receiver = backend_b.receiver("remote-wait").await.unwrap();
let backend_a = MessengerQueueBackend::new(
Arc::clone(&m_a),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let sender = backend_a.sender("remote-wait").await.unwrap();
let sender_clone = Arc::clone(&sender);
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(200)).await;
sender_clone
.send(Bytes::from_static(b"delayed"))
.await
.unwrap();
});
let result = tokio::time::timeout(Duration::from_secs(5), receiver.recv())
.await
.expect("recv should not timeout waiting 5s for a 200ms delayed item");
assert_eq!(result.unwrap(), Some(Bytes::from_static(b"delayed")));
}
#[tokio::test(flavor = "multi_thread")]
async fn test_messenger_remote_batch() {
let (m_a, m_b) = setup_two_messengers().await;
let backend_b = MessengerQueueBackend::new(
Arc::clone(&m_b),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let receiver = backend_b.receiver("remote-batch").await.unwrap();
let backend_a = MessengerQueueBackend::new(
Arc::clone(&m_a),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let sender = backend_a.sender("remote-batch").await.unwrap();
for i in 0u8..3 {
sender.send(Bytes::from(vec![i])).await.unwrap();
}
let opts = NextOptions::new()
.batch_size(3)
.timeout(Duration::from_secs(5));
let batch = receiver.recv_batch(&opts).await.unwrap();
assert_eq!(batch.len(), 3);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_messenger_remote_batch_delayed_arrival() {
let (m_a, m_b) = setup_two_messengers().await;
let backend_b = MessengerQueueBackend::new(
Arc::clone(&m_b),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let receiver = backend_b.receiver("remote-delayed-batch").await.unwrap();
let backend_a = MessengerQueueBackend::new(
Arc::clone(&m_a),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let sender = backend_a.sender("remote-delayed-batch").await.unwrap();
let recv_handle = tokio::spawn({
let receiver = Arc::clone(&receiver);
async move {
let opts = NextOptions::new()
.batch_size(3)
.timeout(Duration::from_secs(5));
receiver.recv_batch(&opts).await
}
});
tokio::time::sleep(Duration::from_millis(500)).await;
for i in 0u8..3 {
sender.send(Bytes::from(vec![i])).await.unwrap();
}
let batch = recv_handle.await.unwrap().unwrap();
assert_eq!(
batch.len(),
3,
"batch should collect all 3 items despite initial empty polls"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_messenger_remote_batch_timeout_partial() {
let (m_a, m_b) = setup_two_messengers().await;
let backend_b = MessengerQueueBackend::new(
Arc::clone(&m_b),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let receiver = backend_b.receiver("remote-partial").await.unwrap();
let backend_a = MessengerQueueBackend::new(
Arc::clone(&m_a),
m_b.instance_id(),
MessengerQueueConfig::default(),
);
let sender = backend_a.sender("remote-partial").await.unwrap();
sender.send(Bytes::from_static(b"a")).await.unwrap();
sender.send(Bytes::from_static(b"b")).await.unwrap();
let opts = NextOptions::new()
.batch_size(5)
.timeout(Duration::from_secs(2));
let batch = receiver.recv_batch(&opts).await.unwrap();
assert_eq!(batch.len(), 2);
}