use futures::StreamExt;
use std::time::Duration;
#[cfg(feature = "js")]
use wasm_bindgen_test::wasm_bindgen_test;
use crate::{droppable_loop_channel, loop_channel};
use remoc::{
exec,
exec::time::sleep,
rch::{
base::SendErrorKind,
watch::{self, ChangedError, ReceiverStream, SendError, TransferStrategy, WatchExt},
},
};
async fn simple(transfer_strategy: TransferStrategy) {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let start_value = 2;
let end_value = 124;
println!("Sending remote mpsc channel receiver");
let (mut tx, rx) = watch::channel(start_value).with_transfer_strategy(transfer_strategy);
a_tx.send(rx).await.unwrap();
println!("Receiving remote mpsc channel receiver");
let mut rx = b_rx.recv().await.unwrap().unwrap();
{
let value = rx.borrow().unwrap();
println!("Initial value: {value:?}");
}
let recv_task = exec::spawn(async move {
let mut value = *rx.borrow().unwrap();
assert_eq!(value, start_value);
while rx.changed().await.is_ok() {
value = *rx.borrow_and_update().unwrap();
println!("Received value change: {value}");
}
value = *rx.borrow_and_update().unwrap();
assert_eq!(value, end_value);
});
for value in start_value..=end_value {
println!("Sending {value}");
tx.send(value).unwrap();
assert_eq!(*tx.borrow(), value);
if value % 10 == 0 {
sleep(Duration::from_millis(20)).await;
}
}
tx.check().unwrap();
drop(tx);
println!("Waiting for receive task");
recv_task.await.unwrap();
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn simple_single_transfer_strategy() {
simple(TransferStrategy::Single).await;
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn simple_global_buffered_transfer_strategy() {
simple(TransferStrategy::GlobalBuffered).await;
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn simple_channel_buffered_transfer_strategy() {
simple(TransferStrategy::ChannelBuffered).await;
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn simple_stream() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let start_value = 2;
let end_value = 124;
println!("Sending remote mpsc channel receiver");
let (mut tx, rx) = watch::channel(start_value);
a_tx.send(rx).await.unwrap();
println!("Receiving remote mpsc channel receiver");
let rx = b_rx.recv().await.unwrap().unwrap();
let mut rx = ReceiverStream::from(rx);
let recv_task = exec::spawn(async move {
let mut value = 0;
while let Some(rxed_value) = rx.next().await {
value = rxed_value.unwrap();
println!("Received value change: {value}");
}
assert_eq!(value, end_value);
});
let mut prev_value = start_value;
for value in start_value..=end_value {
println!("Sending {value}");
let last_value = tx.send_replace(value);
assert_eq!(last_value, prev_value);
assert_eq!(*tx.borrow(), value);
prev_value = value;
if value % 10 == 0 {
sleep(Duration::from_millis(20)).await;
println!("Modifying");
tx.send_modify(|v| *v -= 1);
prev_value -= 1;
}
}
tx.check().unwrap();
drop(tx);
println!("Waiting for receive task");
recv_task.await.unwrap();
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn forward() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let start_value = 2;
let end_value = 124;
let (tx, local_rx) = tokio::sync::watch::channel(start_value);
println!("Forwarding remote mpsc channel receiver");
let (forward, rx) = watch::forward(local_rx);
a_tx.send(rx).await.unwrap();
println!("Receiving remote mpsc channel receiver");
let mut rx = b_rx.recv().await.unwrap().unwrap();
{
let value = rx.borrow().unwrap();
println!("Initial value: {value:?}");
}
let recv_task = exec::spawn(async move {
let mut value = *rx.borrow().unwrap();
assert_eq!(value, start_value);
while rx.changed().await.is_ok() {
value = *rx.borrow_and_update().unwrap();
println!("Received value change: {value}");
}
value = *rx.borrow_and_update().unwrap();
assert_eq!(value, end_value);
});
for value in start_value..=end_value {
println!("Sending {value}");
tx.send(value).unwrap();
assert_eq!(*tx.borrow(), value);
if value % 10 == 0 {
sleep(Duration::from_millis(20)).await;
}
}
drop(tx);
println!("Waiting for receive task");
recv_task.await.unwrap();
println!("Waiting for forward task");
forward.await.unwrap();
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn forwarded() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let start_value = 2;
let end_value = 32;
let (tx, local_rx) = tokio::sync::watch::channel(start_value);
println!("Forwarding local watch receiver via Receiver::forwarded");
let rx = watch::Receiver::forwarded(local_rx);
a_tx.send(rx).await.unwrap();
let mut rx = b_rx.recv().await.unwrap().unwrap();
let recv_task = exec::spawn(async move {
let mut value = *rx.borrow().unwrap();
assert_eq!(value, start_value);
while rx.changed().await.is_ok() {
value = *rx.borrow_and_update().unwrap();
println!("Received value change: {value}");
}
value = *rx.borrow_and_update().unwrap();
assert_eq!(value, end_value);
});
for value in start_value..=end_value {
tx.send(value).unwrap();
if value % 5 == 0 {
sleep(Duration::from_millis(10)).await;
}
}
drop(tx);
println!("Waiting for receive task");
recv_task.await.unwrap();
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn modify_stream() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let start_value = 2;
let end_value = 124;
println!("Sending remote mpsc channel receiver");
let (mut tx, rx) = watch::channel(start_value);
a_tx.send(rx).await.unwrap();
println!("Receiving remote mpsc channel receiver");
let rx = b_rx.recv().await.unwrap().unwrap();
let mut rx = ReceiverStream::from(rx);
let recv_task = exec::spawn(async move {
let mut value = 0;
while let Some(rxed_value) = rx.next().await {
value = rxed_value.unwrap();
println!("Received value change: {value}");
}
assert_eq!(value, end_value);
});
for value in (start_value + 1)..=end_value {
println!("Modifying {value}");
tx.send_modify(|v| *v += 1);
if value % 10 == 0 {
sleep(Duration::from_millis(20)).await;
}
}
tx.check().unwrap();
drop(tx);
println!("Waiting for receive task");
recv_task.await.unwrap();
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn send_if_modified() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let start_value = 2;
let end_value = 124;
println!("Sending remote watch channel receiver");
let (mut tx, rx) = watch::channel(start_value);
a_tx.send(rx).await.unwrap();
println!("Receiving remote watch channel receiver");
let rx = b_rx.recv().await.unwrap().unwrap();
let mut rx = ReceiverStream::from(rx);
let recv_task = exec::spawn(async move {
let mut last_value = start_value;
while let Some(rxed_value) = rx.next().await {
let value = rxed_value.unwrap();
println!("Received value change: {value}");
assert_eq!(value % 2, 0);
last_value = value;
}
assert_eq!(last_value, end_value);
});
for value in (start_value + 1)..=end_value {
let notified = tx.send_if_modified(|v| {
*v = value;
value % 2 == 0
});
assert_eq!(notified, value % 2 == 0);
assert_eq!(*tx.borrow(), value);
if value % 10 == 0 {
sleep(Duration::from_millis(20)).await;
}
}
tx.check().unwrap();
drop(tx);
println!("Waiting for receive task");
recv_task.await.unwrap();
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn send_if_different() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let start_value = 2;
let end_value = 124;
println!("Sending remote watch channel receiver");
let (mut tx, rx) = watch::channel(start_value);
a_tx.send(rx).await.unwrap();
println!("Receiving remote watch channel receiver");
let rx = b_rx.recv().await.unwrap().unwrap();
let mut rx = ReceiverStream::from(rx);
let recv_task = exec::spawn(async move {
let mut last_value = start_value - 1;
while let Some(rxed_value) = rx.next().await {
let value = rxed_value.unwrap();
println!("Received value change: {value}");
assert_ne!(value, last_value);
last_value = value;
}
assert_eq!(last_value, end_value);
});
for value in (start_value + 1)..=end_value {
let current = *tx.borrow();
let notified = tx.send_if_different(current);
assert!(!notified);
assert_eq!(*tx.borrow(), value - 1);
let notified = tx.send_if_different(value);
assert!(notified);
assert_eq!(*tx.borrow(), value);
if value % 10 == 0 {
sleep(Duration::from_millis(20)).await;
}
}
tx.check().unwrap();
drop(tx);
println!("Waiting for receive task");
recv_task.await.unwrap();
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn close() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Sender<i16>>().await;
println!("Sending remote mpsc channel sender");
let (tx, rx) = watch::channel(123);
a_tx.send(tx).await.unwrap();
println!("Receiving remote mpsc channel sender");
let mut tx = b_rx.recv().await.unwrap().unwrap();
println!("Cloning receiver");
let rx2 = rx.clone();
assert!(!tx.is_closed());
println!("Dropping first receiver");
drop(rx);
assert!(!tx.is_closed());
println!("Dropping second receiver");
drop(rx2);
println!("Waiting for close notification");
tx.closed().await;
assert!(tx.is_closed());
tx.check().unwrap();
println!("Attempting to send");
match tx.send(15) {
Ok(()) => panic!("send succeeded after close"),
Err(err) if err.is_closed() => (),
Err(err) => panic!("wrong error after close: {err}"),
}
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn conn_failure() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx), conn) = droppable_loop_channel::<watch::Sender<i16>>().await;
println!("Sending remote mpsc channel sender");
let (tx, rx) = watch::channel(123);
a_tx.send(tx).await.unwrap();
println!("Receiving remote mpsc channel sender");
let mut tx = b_rx.recv().await.unwrap().unwrap();
println!("Cloning receiver");
let _rx2 = rx.clone();
assert!(!tx.is_closed());
println!("Dropping connection");
drop(conn);
println!("Waiting for close notification");
tx.closed().await;
assert!(tx.is_closed());
tx.check().unwrap();
println!("Attempting to send");
match tx.send(15) {
Ok(()) => panic!("send succeeded after close"),
Err(err) if err.is_closed() => (),
Err(err) => panic!("wrong error after close: {err}"),
}
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn max_item_size_exceeded() {
crate::init();
if !remoc::exec::are_threads_available().await {
println!("test requires threads");
return;
}
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<Vec<u8>>>().await;
println!("Sending remote mpsc channel receiver");
let (mut tx, rx) = watch::channel(Vec::new());
a_tx.send(rx).await.unwrap();
println!("Receiving remote mpsc channel receiver");
let mut rx = b_rx.recv().await.unwrap().unwrap();
assert_eq!(tx.max_item_size(), rx.max_item_size());
let max_item_size = tx.max_item_size();
println!("Maximum send and recv item size is {max_item_size}");
{
let value = rx.borrow().unwrap();
println!("Initial value: {value:?}");
}
let recv_task = exec::spawn(async move {
loop {
let res = rx.changed().await;
println!("RX changed result: {res:?}");
if res.is_err() {
break res;
}
let value = rx.borrow_and_update().unwrap().clone();
println!("Received value change: {} elements", value.len());
}
});
let elems = max_item_size / 10;
println!("Sending {elems} elements");
let value = vec![100; elems];
tx.send(value.clone()).unwrap();
assert_eq!(*tx.borrow(), value);
sleep(Duration::from_millis(100)).await;
let elems = max_item_size * 10;
println!("Sending {elems} elements");
let value = vec![100; elems];
tx.send(value.clone()).unwrap();
assert_eq!(*tx.borrow(), value);
println!("Waiting for receive task");
assert!(matches!(recv_task.await.unwrap(), Err(ChangedError::Closed)));
println!("Sending one more element to obtain error");
let res = tx.send(vec![1; 1]);
println!("Result: {res:?}");
assert!(matches!(res, Err(SendError::RemoteSend(SendErrorKind::MaxItemSizeExceeded))));
assert!(matches!(tx.error(), Some(SendError::RemoteSend(SendErrorKind::MaxItemSizeExceeded))));
tx.clear_error();
assert!(tx.error().is_none());
tx.check().unwrap();
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn max_item_size_exceeded_check() {
crate::init();
if !remoc::exec::are_threads_available().await {
println!("test requires threads");
return;
}
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<Vec<u8>>>().await;
println!("Sending remote mpsc channel receiver");
let (mut tx, rx) = watch::channel(Vec::new());
a_tx.send(rx).await.unwrap();
println!("Receiving remote mpsc channel receiver");
let mut rx = b_rx.recv().await.unwrap().unwrap();
assert_eq!(tx.max_item_size(), rx.max_item_size());
let max_item_size = tx.max_item_size();
println!("Maximum send and recv item size is {max_item_size}");
{
let value = rx.borrow().unwrap();
println!("Initial value: {value:?}");
}
let recv_task = exec::spawn(async move {
loop {
let res = rx.changed().await;
println!("RX changed result: {res:?}");
if res.is_err() {
break res;
}
let value = rx.borrow_and_update().unwrap().clone();
println!("Received value change: {} elements", value.len());
}
});
let elems = max_item_size / 10;
println!("Sending {elems} elements");
let value = vec![100; elems];
tx.send(value.clone()).unwrap();
assert_eq!(*tx.borrow(), value);
sleep(Duration::from_millis(100)).await;
let elems = max_item_size * 10;
println!("Sending {elems} elements");
let value = vec![100; elems];
tx.send(value.clone()).unwrap();
assert_eq!(*tx.borrow(), value);
println!("Wait for sender close");
tx.closed().await;
let res = tx.check();
println!("Sender check result: {res:?}");
assert!(matches!(res, Err(SendError::RemoteSend(SendErrorKind::MaxItemSizeExceeded))));
println!("Waiting for receive task");
assert!(matches!(recv_task.await.unwrap(), Err(ChangedError::Closed)));
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn has_changed() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let (tx, rx) = watch::channel(10);
a_tx.send(rx).await.unwrap();
let mut rx = b_rx.recv().await.unwrap().unwrap();
let value = *rx.borrow_and_update().unwrap();
assert_eq!(value, 10);
assert!(!rx.has_changed().unwrap());
tx.send(20).unwrap();
let mut n = 1000;
while !rx.has_changed().unwrap() {
assert!(n > 0, "timed out waiting for has_changed");
sleep(Duration::from_millis(10)).await;
n -= 1;
}
let value = *rx.borrow_and_update().unwrap();
assert_eq!(value, 20);
assert!(!rx.has_changed().unwrap());
drop(tx);
let mut n = 1000;
loop {
match rx.has_changed() {
Err(_) => break,
Ok(_) => {
assert!(n > 0, "timed out waiting for closed error");
sleep(Duration::from_millis(10)).await;
n -= 1;
}
}
}
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn mark_changed_and_unchanged() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let (tx, rx) = watch::channel(10);
a_tx.send(rx).await.unwrap();
let mut rx = b_rx.recv().await.unwrap().unwrap();
let _ = rx.borrow_and_update().unwrap();
assert!(!rx.has_changed().unwrap());
rx.mark_changed();
assert!(rx.has_changed().unwrap());
let _ = rx.borrow_and_update().unwrap();
assert!(!rx.has_changed().unwrap());
tx.send(30).unwrap();
let mut n = 1000;
while !rx.has_changed().unwrap() {
assert!(n >= 0, "timed out waiting for has_changed");
sleep(Duration::from_millis(10)).await;
n -= 1;
}
rx.mark_unchanged();
assert!(!rx.has_changed().unwrap());
let value = *rx.borrow().unwrap();
assert_eq!(value, 30);
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn wait_for_immediate() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let (_tx, rx) = watch::channel(42);
a_tx.send(rx).await.unwrap();
let mut rx = b_rx.recv().await.unwrap().unwrap();
let value = rx.wait_for(|v| *v == 42).await.unwrap();
assert_eq!(*value, 42);
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn wait_for_future_value() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let (tx, rx) = watch::channel(0);
a_tx.send(rx).await.unwrap();
let mut rx = b_rx.recv().await.unwrap().unwrap();
let send_task = exec::spawn(async move {
sleep(Duration::from_millis(50)).await;
tx.send(5).unwrap();
sleep(Duration::from_millis(50)).await;
tx.send(10).unwrap();
sleep(Duration::from_millis(50)).await;
tx.send(100).unwrap();
tx
});
let value = rx.wait_for(|v| *v >= 100).await.unwrap();
assert!(*value >= 100);
let _tx = send_task.await.unwrap();
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn wait_for_closed() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
let (tx, rx) = watch::channel(0);
a_tx.send(rx).await.unwrap();
let mut rx = b_rx.recv().await.unwrap().unwrap();
let _drop_task = exec::spawn(async move {
sleep(Duration::from_millis(50)).await;
drop(tx);
});
let result = rx.wait_for(|v| *v == 999).await;
assert!(result.is_err(), "wait_for should fail when sender is dropped");
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn conn_failure_receiver_changed() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx), conn) = droppable_loop_channel::<watch::Receiver<i16>>().await;
println!("Sending remote watch receiver");
let (tx, rx) = watch::channel(123);
a_tx.send(rx).await.unwrap();
println!("Receiving remote watch receiver");
let mut rx = b_rx.recv().await.unwrap().unwrap();
let initial = *rx.borrow_and_update().unwrap();
assert_eq!(initial, 123);
tx.send(456).unwrap();
let changed_res = rx.changed().await;
println!("changed() after normal send returned: {changed_res:?}");
assert!(changed_res.is_ok());
assert_eq!(*rx.borrow_and_update().unwrap(), 456);
println!("Dropping connection");
drop(conn);
let _keep_tx_alive = tx;
println!("Waiting for changed() on receiver after connection drop");
let res = rx.changed().await;
println!("changed() result after connection drop: {res:?}");
let err = match res {
Err(ChangedError::Recv(err)) => err,
other => panic!("expected ChangedError::RecvError, got: {other:?}"),
};
assert!(err.is_final(), "RecvError reported by changed() must be final");
let borrow_res = rx.borrow().map(|v| *v);
println!("borrow() after connection drop: {borrow_res:?}");
assert!(borrow_res.is_err());
let res2 = rx.changed().await;
println!("Second changed() result: {res2:?}");
assert!(matches!(res2, Err(ChangedError::Recv(ref e)) if e.is_final()));
assert!(matches!(rx.has_changed(), Err(ChangedError::Recv(ref e)) if e.is_final()));
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn sender_drop_receiver_changed() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i16>>().await;
println!("Sending remote watch receiver");
let (tx, rx) = watch::channel(1);
a_tx.send(rx).await.unwrap();
println!("Receiving remote watch receiver");
let mut rx = b_rx.recv().await.unwrap().unwrap();
assert_eq!(*rx.borrow_and_update().unwrap(), 1);
tx.send(2).unwrap();
let res = rx.changed().await;
println!("changed() after normal send returned: {res:?}");
assert!(res.is_ok());
assert_eq!(*rx.borrow_and_update().unwrap(), 2);
println!("Dropping remote sender (clean close)");
drop(tx);
println!("Waiting for changed() after sender drop");
let res = rx.changed().await;
println!("changed() result after sender drop: {res:?}");
assert!(matches!(res, Err(ChangedError::Closed)));
let res2 = rx.changed().await;
println!("Second changed() result: {res2:?}");
assert!(matches!(res2, Err(ChangedError::Closed)));
assert!(matches!(rx.has_changed(), Err(ChangedError::Closed)));
assert_eq!(*rx.borrow().unwrap(), 2);
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn rate_limit_default_and_set() {
crate::init();
let (mut tx, _rx): (watch::Sender<i16>, watch::Receiver<i16>) = watch::channel(0i16);
assert_eq!(tx.rate_limit(), Duration::ZERO);
tx.set_rate_limit(Duration::from_millis(50));
assert_eq!(tx.rate_limit(), Duration::from_millis(50));
tx.set_rate_limit(Duration::from_millis(75));
assert_eq!(tx.rate_limit(), Duration::from_millis(75));
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn rate_limit_coalesces_burst() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i32>>().await;
let rate_limit = Duration::from_millis(100);
let end_value = 1000;
println!("Sending remote watch channel receiver");
let (mut tx, rx) = watch::channel(0);
tx.set_rate_limit(rate_limit);
assert_eq!(tx.rate_limit(), rate_limit);
a_tx.send(rx).await.unwrap();
println!("Receiving remote watch channel receiver");
let mut rx = b_rx.recv().await.unwrap().unwrap();
let recv_task = exec::spawn(async move {
let mut count = 0u32;
let mut last = *rx.borrow_and_update().unwrap();
while rx.changed().await.is_ok() {
last = *rx.borrow_and_update().unwrap();
count += 1;
println!("Received update #{count}: {last}");
}
(count, last)
});
println!("Sending burst of {end_value} updates");
for value in 1..=end_value {
tx.send(value).unwrap();
}
sleep(rate_limit * 5).await;
tx.check().unwrap();
drop(tx);
println!("Waiting for receive task");
let (count, last) = recv_task.await.unwrap();
println!("Total received updates: {count}, last value: {last}");
assert_eq!(last, end_value);
assert!(count < 10, "expected coalescing, but got {count} updates for {end_value} sends");
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn rate_limit_throttles_rate() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i32>>().await;
let rate_limit = Duration::from_millis(100);
let send_interval = Duration::from_millis(10);
let end_value = 100;
println!("Sending remote watch channel receiver");
let (mut tx, rx) = watch::channel(0);
tx.set_rate_limit(rate_limit);
a_tx.send(rx).await.unwrap();
println!("Receiving remote watch channel receiver");
let mut rx = b_rx.recv().await.unwrap().unwrap();
let recv_task = exec::spawn(async move {
let mut count = 0u32;
let mut last = *rx.borrow_and_update().unwrap();
while rx.changed().await.is_ok() {
last = *rx.borrow_and_update().unwrap();
count += 1;
}
(count, last)
});
println!("Sending {end_value} updates spaced {send_interval:?} apart");
for value in 1..=end_value {
tx.send(value).unwrap();
sleep(send_interval).await;
}
sleep(rate_limit * 5).await;
tx.check().unwrap();
drop(tx);
println!("Waiting for receive task");
let (count, last) = recv_task.await.unwrap();
println!("Total received updates: {count}, last value: {last}");
assert_eq!(last, end_value);
assert!(count > 1, "expected multiple throttled updates, got {count}");
assert!(count < (end_value / 2) as u32, "expected throttling, but got {count} updates for {end_value} sends");
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn rate_limit_survives_sender_transport() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Sender<i32>>().await;
let rate_limit = Duration::from_millis(100);
let end_value = 1000;
let (mut tx, rx) = watch::channel(0);
tx.set_rate_limit(rate_limit);
println!("Sending remote watch channel sender");
a_tx.send(tx).await.unwrap();
println!("Receiving remote watch channel sender");
let mut remote_tx = b_rx.recv().await.unwrap().unwrap();
assert_eq!(remote_tx.rate_limit(), rate_limit);
let mut rx = rx;
let recv_task = exec::spawn(async move {
let mut count = 0u32;
let mut last = *rx.borrow_and_update().unwrap();
while rx.changed().await.is_ok() {
last = *rx.borrow_and_update().unwrap();
count += 1;
}
(count, last)
});
println!("Sending burst of {end_value} updates from transported sender");
for value in 1..=end_value {
remote_tx.send(value).unwrap();
}
sleep(rate_limit * 5).await;
remote_tx.check().unwrap();
drop(remote_tx);
println!("Waiting for receive task");
let (count, last) = recv_task.await.unwrap();
println!("Total received updates: {count}, last value: {last}");
assert_eq!(last, end_value);
assert!(count < 10, "expected coalescing, but got {count} updates for {end_value} sends");
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn receiver_rate_limit_default_and_set() {
crate::init();
let (_tx, mut rx): (watch::Sender<i16>, watch::Receiver<i16>) = watch::channel(0i16);
assert_eq!(rx.rate_limit(), Duration::ZERO);
rx.set_rate_limit(Duration::from_millis(50));
assert_eq!(rx.rate_limit(), Duration::from_millis(50));
rx.set_rate_limit(Duration::from_millis(75));
assert_eq!(rx.rate_limit(), Duration::from_millis(75));
rx.set_rate_limit(Duration::ZERO);
assert_eq!(rx.rate_limit(), Duration::ZERO);
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn receiver_rate_limit_coalesces_burst() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i32>>().await;
let rate_limit = Duration::from_millis(100);
let end_value = 1000;
println!("Sending remote watch channel receiver");
let (mut tx, rx) = watch::channel(0);
assert_eq!(tx.rate_limit(), Duration::ZERO);
a_tx.send(rx).await.unwrap();
println!("Receiving remote watch channel receiver");
let mut rx = b_rx.recv().await.unwrap().unwrap();
assert_eq!(rx.rate_limit(), Duration::ZERO);
rx.set_rate_limit(rate_limit);
assert_eq!(rx.rate_limit(), rate_limit);
sleep(rate_limit).await;
let recv_task = exec::spawn(async move {
let mut count = 0u32;
let mut last = *rx.borrow_and_update().unwrap();
while rx.changed().await.is_ok() {
last = *rx.borrow_and_update().unwrap();
count += 1;
println!("Received update #{count}: {last}");
}
(count, last)
});
println!("Sending burst of {end_value} updates");
for value in 1..=end_value {
tx.send(value).unwrap();
}
sleep(rate_limit * 5).await;
tx.check().unwrap();
drop(tx);
println!("Waiting for receive task");
let (count, last) = recv_task.await.unwrap();
println!("Total received updates: {count}, last value: {last}");
assert_eq!(last, end_value);
assert!(count < 10, "expected coalescing, but got {count} updates for {end_value} sends");
}
#[cfg_attr(not(feature = "js"), tokio::test)]
#[cfg_attr(feature = "js", wasm_bindgen_test)]
async fn receiver_rate_limit_combines_with_sender_via_max() {
crate::init();
let ((mut a_tx, _), (_, mut b_rx)) = loop_channel::<watch::Receiver<i32>>().await;
let sender_rate_limit = Duration::from_millis(10);
let receiver_rate_limit = Duration::from_millis(200);
let send_interval = Duration::from_millis(10);
let end_value = 100;
println!("Sending remote watch channel receiver");
let (mut tx, rx) = watch::channel(0);
tx.set_rate_limit(sender_rate_limit);
a_tx.send(rx).await.unwrap();
println!("Receiving remote watch channel receiver");
let mut rx = b_rx.recv().await.unwrap().unwrap();
rx.set_rate_limit(receiver_rate_limit);
sleep(receiver_rate_limit).await;
let recv_task = exec::spawn(async move {
let mut count = 0u32;
let mut last = *rx.borrow_and_update().unwrap();
while rx.changed().await.is_ok() {
last = *rx.borrow_and_update().unwrap();
count += 1;
}
(count, last)
});
println!("Sending {end_value} updates spaced {send_interval:?} apart");
for value in 1..=end_value {
tx.send(value).unwrap();
sleep(send_interval).await;
}
sleep(receiver_rate_limit * 5).await;
tx.check().unwrap();
drop(tx);
println!("Waiting for receive task");
let (count, last) = recv_task.await.unwrap();
println!("Total received updates: {count}, last value: {last}");
assert_eq!(last, end_value);
assert!(count > 1, "expected multiple throttled updates, got {count}");
assert!(count < 20, "expected receiver rate limit (max) to dominate, but got {count} updates");
}