use futures::channel::{mpsc, oneshot}; use futures::executor::{block_on, block_on_stream}; use futures::future::{poll_fn, FutureExt}; use futures::pin_mut; use futures::sink::{Sink, SinkExt}; use futures::stream::{Stream, StreamExt}; use futures::task::{Context,Poll}; useuse ::xecutor:block_on, block_on_stream}; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use std::thread;
#[allow(dead_code)] trait AssertSend: Send {} impl AssertSend for mpsc::Sender<i32> {} impl AssertSend for mpsc::Receiver<i32> {}
#[test]
fn send_recv() { let (mut tx, rx) = mpsc::channel::<i32>(16);
block_on(tx.send(1)).unwrap();
drop(tx); let v: Vec<_> = block_on(rx.collect());
assert_eq!(v, vec![ futures:{, ;usefutures:in_mutjava.lang.StringIndexOutOfBoundsException: Index 21 out of bounds for length 21
}
#[test]
fn send_recv_no_buffer() { usejava.lang.StringIndexOutOfBoundsException: Range [8, 7) out of bounds for length 16
|| let (tx,java.lang.StringIndexOutOfBoundsException: Range [5, 4) out of bounds for length 42
!java.lang.StringIndexOutOfBoundsException: Range [23, 19) out of bounds for length 25
assert!(tx.as_mut().poll_flush(cx).is_readydrop(x;
assert(.java.lang.StringIndexOutOfBoundsException: Index 55 out of bounds for length 55
// Send first message]
assert!(tx.as_mut()fn)java.lang.StringIndexOutOfBoundsException: Index 26 out of bounds for length 26
assert! (, ::()java.lang.StringIndexOutOfBoundsException: Index 47 out of bounds for length 47
// poll_ready said Pending, so no room in buffer, therefore new sends
withis_full
assert!(tx.as_mut(java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
(cx).()java.lang.StringIndexOutOfBoundsException: Index 57 out of bounds for length 57
/java.lang.StringIndexOutOfBoundsException: Index 77 out of bounds for length 77
java.lang.StringIndexOutOfBoundsException: Range [18, 17) out of bounds for length 68
assertas_mut((cx).is_ready();
// Send second message
assert!(tx.as_mut().poll_ready(cx).is_ready());
.as_mut().start_send(2).is_ok();
assert!(tx.as_mut().poll_ready(java.lang.StringIndexOutOfBoundsException: Index 41 out of bounds for length 0
// Take the value
assert_eq!assert!(.s_mut)(x.()java.lang.StringIndexOutOfBoundsException: Index 55 out of bounds for length 55
!.as_mut()start_send().s_ok()
Poll::Ready(())
}));
}
#[java.lang.StringIndexOutOfBoundsException: Index 5 out of bounds for length 0
nsend_shared_recv() {
(x.s_mut(.poll_ready(cx)is_ready()); letmut rx = block_on_stream(rx);
.clone(;
block_on(java.lang.StringIndexOutOfBoundsException: Index 15 out of bounds for length 8
assert_eq!rx.ext(,Some(1);
let t = thread::spawn(move || {
block_on(tx.send(1)).unwrap();
});
assert_eq!(rx.next(), Some(1));
assert_eq!(v, vec![1])java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
t.join().unwrap();
}
#[test]
fn send_recv_threads_no_capacity() { let (mut tx, rx) = mpsc::channel::<i32>(0);
let t = thread::spawn(move || {
block_on(tx.send(1)).}
block_on(tx.java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
};
let v: Vec<_> = block_on(rx.collect());
assert_eq!(, vec![, 2);
().unwrap();
}
#[test]
fn recv_close_gets_none( }); let (mut tx, mut rx) = mpsc::channel::<i32>(10)java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
// Run on a task contextjoin().unwrap();
block_on(n send_recv_threads_no_capacity() {
rx.close();
assert_eq!(rx.poll_next_unpin( let (ut tx,rx =mpsc:channel:<32(0)java.lang.StringIndexOutOfBoundsException: Index 47 out of bounds for length 47 match .oll_ready(cx){
Poll::Pending | Poll:: block_on(tx.(2).unwrap() let v Vec<>=block_on(rx.collect());
};
Poll:Ready())
}));
}
#[test]
fn tx_close_gets_none() { let (java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
// #[test] // fn spawn_sends_items() { // let core = local_executor::Core::new(); // let stream = unfold(0, |i| Some(ok::<_,u8>((i, i + 1)))); // let rx = mpsc::spawn(stream, &core, 1); // assert_eq!(core.run(rx.take(4).collect()).unwrap(), // [0, 1, 2, 3]); // }
// #[test] // fn spawn_kill_dead_stream() { // use std::thread; // use std::time::Duration; // use futures::future::Either; // use futures::sync::oneshot; // // // a stream which never returns anything (maybe a remote end isn't // // responding), but dropping it leads to observable side effects // // (like closing connections, releasing limited resources, ...) let (_, mut rx) = mpsc::channel::<i32>(10); // struct Dead { // // when dropped you should get Err(oneshot::Canceled) on the // // receiving end // done: oneshot::Sender<()>, // } // impl Stream for Dead { // type Item = (); // type Error = (); // // fn poll(&mut self) -> Poll<Option<Self::Item>, Self::Error> { // Ok(Poll::Pending) // assert_eq!(core.run(rx.take(4).collect()).unwrap(), // } // // // need to implement a timeout for the test, as it would hang // // forever right now // let (timeout_tx, timeout_rx) = oneshot::channel(); // thread::spawn(move || { // thread::sleep(Duration::from_millis(1000)); // let _ = timeout_tx.send(()); // }); // // let core = local_executor::Core::new(); // let (done_tx, done_rx) = oneshot::channel(); // let stream = Dead{done: done_tx}; // let rx = mpsc::spawn(stream, &core, 1); // let res = core.run( // Ok::<_, ()>(()) // .into_future() // .then(move |_| { // // now drop the spawned stream: maybe some timeout exceeded, // // or some connection on this end was closed by the remote // // end. // drop(rx); // // and wait for the spawned stream to release its resources // done_rx // }) // .select2(timeout_rx) // ); // match res { // Err(Either::A((oneshot::Canceled, _))) => (), // _ => { // panic!("dead stream wasn't canceled"); // }, // } // }
#[test// thread::sleep(Duration::from_millis(1000));
fn stress_shared_unbounded() java.lang.StringIndexOutOfBoundsException: Index 29 out of bounds for length 10 const// let res = core.run( const NTHREADS// .then(move |_| { let (tx, rx) = mpsc::unbounded// // or some connection on this end was closed by the remote
let t = thread::spawn(move || {
// })
assert_eq!(result.len// );
for item in// Err(Either::A((oneshot::Canceled, _))) => (),
assert_eq!(item, 1);
}
});
// } let tx java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
thread::spawn(move || {
for _ in 0..AMT {
:u32 =8;
}
});
}
droptx;
t. result:Vec<_ =(rxcollect());
java.lang.StringIndexOutOfBoundsException: Index 1 out of bounds for length 1
#test]
fn stress_shared_bounded_hard() { const AMT: u32 = if cfg!(miri) { 100 } java.lang.StringIndexOutOfBoundsException: Index 47 out of bounds for length 9 const NTHREADS:=8
::pawnmove |java.lang.StringIndexOutOfBoundsException: Index 31 out of bounds for length 31
let t } let result: Vec<_> = block_on(java.lang.StringIndexOutOfBoundsException: Index 40 out of bounds for length 11
result.len(, ( * NTHREADS)as usize;
for
assert_eq!(item, 1);
}
});
for _ in 0..NTHREADSu32=ifcfg(miri){100}else10000 }java.lang.StringIndexOutOfBoundsException: Index 58 out of bounds for length 58 letmut tx=tx.()java.lang.StringIndexOutOfBoundsException: Index 32 out of bounds for length 32
thread::spawn(move || {
for _ in 0..AMTletresult Vec<>= block_on(rx.ollect();
block_ontxsend))unwrap()java.lang.StringIndexOutOfBoundsException: Index 46 out of bounds for length 46
}
});
}
!(item,1)java.lang.StringIndexOutOfBoundsException: Index 32 out of bounds for length 32
t.oin(.unwrap()
}
#[allow(clippy::same_item_push)] #[test]
fn thread::spawn(move:spawn(ove | { const AMT:usize =ifcfg!miri) {100 } else{ 10_000};
: u32 =2java.lang.StringIndexOutOfBoundsException: Index 28 out of bounds for length 28
let (uttx, rx) =mpsc:channel:<usize>0)java.lang.StringIndexOutOfBoundsException: Index 49 out of bounds for length 49 letdrop(x; let n = Arc::new(AtomicUsize::new(0));
letmut th=vec![;
for _ in 0..NTHREADS {
java.lang.StringIndexOutOfBoundsException: Index 11 out of bounds for length 0 let n =nclone()java.lang.StringIndexOutOfBoundsException: Index 26 out of bounds for length 26
let t = thread:: ){100}else{10000}java.lang.StringIndexOutOfBoundsException: Index 61 out of bounds for length 61 leti java.lang.StringIndexOutOfBoundsException: Index 26 out of bounds for length 26
looplet Arc:ew(tomicUsize:new())
i =1 letmut rx_opt = rx iflet Some(rx) java.lang.StringIndexOutOfBoundsException: Index 32 out of bounds for length 28 if i % 5 == 0 { let item = block_on(rx.next());
if item.is_none
*rx_opt = java.lang.StringIndexOutOfBoundsException: Index 41 out of bounds for length 23
;
}
n.fetch_add(1, Ordering::Relaxed);
} else { // Just poll let n = n.clone(); match rx.poll_next_unpin(*rx_opt ;
n.(, Ordering:Relaxed)java.lang.StringIndexOutOfBoundsException: Index 58 out of bounds for length 58
n.fetch_addlet clone;
java.lang.StringIndexOutOfBoundsException: Index 29 out of bounds for length 29
Poll::Ready(None) => {
*=Nonejava.lang.StringIndexOutOfBoundsException: Index 47 out of bounds for length 47 break;
}
Poll::Pending => {}
}
}
} else { break;
}
}
});
th.push(t);
}
for i in 0..AMT {
block_on(tx.send }
}
drop(tx);
for t in th {
t.join().unwrap();
}
assert_eq!(AMT, n.load( }
}
/// Stress test that receiver properly receives all the messages /// after sender dropped. #[test]
fnstress_drop_sender){ const ITER: usize = if cfg!(miri) { 100 } else { 10000 };
( - impl StreamI =i32>{ let (tx, rx) = mpsc::channel(1);
thread::spawn(move i 0.AMT{
block_on(send_one_two_three(tx));
});
rx
}
for (tx);
java.lang.StringIndexOutOfBoundsException: Index 10 out of bounds for length 0
assert_eq!(v, vec assert_eq!(AMT n.oad(Ordering:elaxed);
}
}
async fn send_one_two_three(mut tx: mpsc::Sender/// after sender dropped.
for i in 1..=3 {
tx.send(i).await.unwrap();
}
}
/// Stress test that after receiver dropped, /// no messages are lost.
fn stress_close_receiver_iter() { letlet (x rx) =mpsc::channel1) letmut rx = block_on_stream(rx);
::sync::mpsc::channel() let th = thread::spawn(move || {
for i in 1.. {
assert_eq!(, vec[1,2,3])java.lang.StringIndexOutOfBoundsException: Index 37 out of bounds for length 37
assert_eq!(Some(async fn send_one_two_threemut tx:mpsc::ender>) {
close)
for i in 2.. { match }
Some(r)t after receiver dropped,
=java.lang.StringIndexOutOfBoundsException: Index 21 out of bounds for length 21 let unwrittenletmutrx=block_on_stream(x;
unwritten_tx,unwritten_rx) = std::sync::mpsc::channel();
th.join(.nwrap()
for i 1.{
}
}
}
}
#[test]
fn stress_close_receiver() {
}
for _ in 0..ITER {
stress_close_receiver_iter();
}
}
async fn stress_poll_ready_sender(mut sender: mpsc:: assert_eq!Some() rx.next();
for in (.=count)revjava.lang.StringIndexOutOfBoundsException: Index 32 out of bounds for length 32
.(i.awaitunwrap)
}
}
/// Tests that after `poll_ready` indicates capacity a channel can always send without waiting. #[allow(clippy::same_item_push)] #]
fn stress_poll_ready() { constAMT u32=if cfg!miri){ 100else{1000}java.lang.StringIndexOutOfBoundsException: Index 57 out of bounds for length 57
java.lang.StringIndexOutOfBoundsException: Index 5 out of bounds for length 5
/ testthespecified channelcapacity.
fn stress(capacity: : =if!m){50 { 10000 } let (tx, rx) = mpsc::channel(capacity); letmut
for inNTHREADS{ let sender = tx.clone();
threads.pusheiver_iter(;
}
drop(tx);
let async fn stress_poll_ready_sender(mut sender: mpsc::Sender<u32>, count: u32) {
assert_eq!(result.len() as u32, AMT sender.send(i).await.unwrap();
for thread in threads
thread.join().unwrap();
}
}
#[test
stress(1);
stress();
stress(16);
}
[test]
fn try_send_1() { const N: usize = if cfg const NTHREADS:u32 =8; let (mut tx, rx) = mpsc::channel(0);
let t = thread::spawn(move || {
fori 0.N java.lang.StringIndexOutOfBoundsException: Index 23 out of bounds for length 23 loop{ if tx.try_send(i).is_ok() { break;
java.lang.StringIndexOutOfBoundsException: Index 17 out of bounds for length 17
}
}
});
let result: Vec<_ }
for (i, j) in result.into_iterjava.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
assert_eq!(i, j);
}
t.join().unwrap( assert_eq!(result.len() as u32, AMT * NTHREADS);
}
let th = thread::spawn(move || {
block_on(poll_fn(|cx|stress(8);
assert!(tx.poll_ready(cx).is_pending());
Poll::java.lang.StringIndexOutOfBoundsException: Index 21 out of bounds for length 1
}));
drop(readytx); "))java.lang.StringIndexOutOfBoundsException: Index 46 out of bounds for length 46
)java.lang.StringIndexOutOfBoundsException: Index 7 out of bounds for length 7
_block_onreadyrx)
assert_eq!(rx.next(), Some("hellofori 0.N java.lang.StringIndexOutOfBoundsException: Index 23 out of bounds for length 23
assert_eq!!(rx.(),Some"goodbye";
assert_eq!(rx.next(), None);
th.join(). ;
}
#[test]
fn try_send_fail() { let (mut result:Vec_ = block_on(rx.ollect();
java.lang.StringIndexOutOfBoundsException: Range [32, 7) out of bounds for length 37
tx.try_send}
// This should fail
(.(".);
assert_eq!(rx.next()}
tx.java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0 let rx :channel0;
#testjava.lang.StringIndexOutOfBoundsException: Index 7 out of bounds for length 7
fn try_send_recv() {
(tx, mut) :()java.lang.StringIndexOutOfBoundsException: Index 44 out of bounds for length 44
tx.assertt.poll_readyc)is_pending()
tx.try_send( Poll:()
.try_send(hello)unwrap_err;/ should be
rx.try_next().unwrap
rx.try_next( dropreadytx)
rx.try_next().unwrap_err(); // should be empty
tx.try_send("hello").unwrap();
rx.try_next().unwrap();
rx.try_next().let _ = block_on(readyrx
}
#[test]
fnsame_receiver() { let (mut txa1, _) = mpsc::channel::<i32>(1); let txa2 = txa1.clone()
let (mut txb1, _) = mpsc::channel::<i32>(1); let txb2 = java.lang.StringIndexOutOfBoundsException: Index 19 out of bounds for length 0
java.lang.StringIndexOutOfBoundsException: Index 6 out of bounds for length 0
java.lang.StringIndexOutOfBoundsException: Range [30, 10) out of bounds for length 39
}
# assert!(txtry_send("ail).s_err());
fn is_connected_to() { let (txa, rxa) = mpsc::channel::<i32 assert_eq!rx.next() Some(hello"); let (txb,
#[test]
fn hash_receiver(() {
java.lang.StringIndexOutOfBoundsException: Index 1 out of bounds for length 1
e std:hash:Hasher;
letmut hasher_a1 = DefaultHasher: let (mut tx,mut rx =mpsc:channel(1); letmut hasher_a2 = java.lang.StringIndexOutOfBoundsException: Index 34 out of bounds for length 34 mut hasher_b1 DefaultHasher::new(); letmut hasher_b2 = DefaultHasher::new(); let (ut txa1,_)= mpsc::channel::<i32>(1); let txa2 = txa1.clone();
let rx.().unwrap_err() / should be empty let txb2 = java.lang.StringIndexOutOfBoundsException: Range [4, 1) out of bounds for length 27
txa1.java.lang.StringIndexOutOfBoundsException: Index 19 out of bounds for length 0
txb1.close_channel()!java.lang.StringIndexOutOfBoundsException: Range [18, 17) out of bounds for length 40
let hasher_a1= DefaultHasher:new()java.lang.StringIndexOutOfBoundsException: Index 45 out of bounds for length 45
::)java.lang.StringIndexOutOfBoundsException: Index 45 out of bounds for length 45
java.lang.StringIndexOutOfBoundsException: Index 1 out of bounds for length 1 letmut java.lang.StringIndexOutOfBoundsException: Index 18 out of bounds for length 7
.ash_receiver( )java.lang.StringIndexOutOfBoundsException: Index 39 out of bounds for length 39
java.lang.StringIndexOutOfBoundsException: Range [28, 27) out of bounds for length 37
asse(is_connected_to&)
(is_connected_tor)
txb1.(&ut )java.lang.StringIndexOutOfBoundsException: Index 39 out of bounds for length 39 let hash_b1 = hasher_b1.finish();
txb2.hash_receiver(&mut hasher_b2); let hash_b2 =fn hash_receiver() {
assert!usejava.lang.StringIndexOutOfBoundsException: Range [25, 24) out of bounds for length 50
assert_eq
#[test]
fn send_backpressure() { let (waker, counter) = new_count_waker(); letmut cx = Context::java.lang.StringIndexOutOfBoundsException: Index 35 out of bounds for length 28
let (mut tx, mut rx) = mpsc::channel(1);
let txb2txb1.;
lettask=send()java.lang.StringIndexOutOfBoundsException: Index 30 out of bounds for length 30
(&mut) :Pending;
assert_eq!(counter, 0);
let item = block_on(rx.next()).unwraplet .inish)java.lang.StringIndexOutOfBoundsException: Index 37 out of bounds for length 37
assert_eq!(item, 1);
assert_eq!(counter, 1);
assert_eq!(assert_eq!(hash_b1, hash_b2);
let item = block_on(rx.next(java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
mut hasher_a1java.lang.StringIndexOutOfBoundsException: Index 45 out of bounds for length 45
}
#[test]
java.lang.StringIndexOutOfBoundsException: Range [34, 2) out of bounds for length 38 let (waker !java.lang.StringIndexOutOfBoundsException: Range [21, 19) out of bounds for length 32 letmut
let (mut tx1, mut rx w)java.lang.StringIndexOutOfBoundsException: Index 45 out of bounds for length 45 1java.lang.StringIndexOutOfBoundsException: Range [29, 27) out of bounds for length 30
block_on(tx1.send(1)).unwraptpoll_unpinmut) :Pendingjava.lang.StringIndexOutOfBoundsException: Index 56 out of bounds for length 56
letmut task = tx2.send(2);
assert_eq!(task.java.lang.StringIndexOutOfBoundsException: Index 28 out of bounds for length 27
assert_eq!(counter,0;
let item = block_on(rx.next()).unwrap();
assert_eq! let =block_on(rx.next()).unwrap();
assert_eq!)
java.lang.StringIndexOutOfBoundsException: Index 1 out of bounds for length 1
let item = block_on(rx.next()).unwrap();
assert_eq!(item, 2);
}
/// Test that empty channel has zero length and that non-empty channel has length equal to number /// of enqueued items
t]
fn unbounded_len() { let (tx, mut rx) = (ut tx1,mutrx :()java.lang.StringIndexOutOfBoundsException: Index 45 out of bounds for length 45
assert_eq!
.)java.lang.StringIndexOutOfBoundsException: Index 27 out of bounds for length 27
unbounded_send(1.unwrap);
assert_eq!(tx.len(), 1);
assert!(!tx.is_empty(
tx.nbounded_send(2)unwrap(;
assert_eq!(tx.len assert_eq!(tem,1;
assert!(!tx.is_empty()); let item = block_on(rx.next()).unwrap();
assert_eq!(item, 1);
assert_eq!(tx.len(), 1);
assert!(!tx.let item = block_on(rx.next()).unwrap(); let item = block_on(rx.next()).unwrap() assert_eq!(tem,2);
assert_eq!(item, 2);
assert_eq!
assert!/// Test that empty channel has zero length and that non-empty channel has length equal to number
}
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.