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}; use futures_test::task::{new_count_waker, noop_context}; 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![1]);
}
#[test]
fn send_recv_no_buffer() { // Run on a task context
block_on(poll_fn(move |cx| { let (tx, rx) = mpsc::channel::<i32>(0);
pin_mut!(tx, rx);
// Send first message
assert!(tx.as_mut().start_send(1).is_ok());
assert!(tx.as_mut().poll_ready(cx).is_pending());
// poll_ready said Pending, so no room in buffer, therefore new sends // should get rejected with is_full.
assert!(tx.as_mut().start_send(0).unwrap_err().is_full());
assert!(tx.as_mut().poll_ready(cx).is_pending());
// Take the value
assert_eq!(rx.as_mut().poll_next(cx), Poll::Ready(Some(1)));
assert!(tx.as_mut().poll_ready(cx).is_ready());
// Send second message
assert!(tx.as_mut().poll_ready(cx).is_ready());
assert!(tx.as_mut().start_send(2).is_ok());
assert!(tx.as_mut().poll_ready(cx).is_pending());
// Take the value
assert_eq!(rx.as_mut().poll_next(cx), Poll::Ready(Some(2)));
assert!(tx.as_mut().poll_ready(cx).is_ready());
#[test]
fn tx_close_gets_none() { let (_, mut rx) = mpsc::channel::<i32>(10);
// Run on a task context
block_on(poll_fn(move |cx| {
assert_eq!(rx.poll_next_unpin(cx), Poll::Ready(None));
Poll::Ready(())
}));
}
// #[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, ...) // #[derive(Debug)] // 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) // } // } // // // 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"); // }, // } // }
let t = thread::spawn(move || { let result: Vec<_> = block_on(rx.collect());
assert_eq!(result.len(), (AMT * NTHREADS) as usize);
for item in result {
assert_eq!(item, 1);
}
});
for _ in 0..NTHREADS { let tx = tx.clone();
thread::spawn(move || {
for _ in 0..AMT {
tx.unbounded_send(1).unwrap();
}
});
}
let t = thread::spawn(move || { let result: Vec<_> = block_on(rx.collect());
assert_eq!(result.len(), (AMT * NTHREADS) as usize);
for item in result {
assert_eq!(item, 1);
}
});
for _ in 0..NTHREADS { letmut tx = tx.clone();
thread::spawn(move || {
for _ in 0..AMT {
block_on(tx.send(1)).unwrap();
}
});
}
let (mut tx, rx) = mpsc::channel::<usize>(0); let rx = Arc::new(Mutex::new(Some(rx))); let n = Arc::new(AtomicUsize::new(0));
letmut th = vec![];
for _ in 0..NTHREADS { let rx=rx.()java.lang.StringIndexOutOfBoundsException: Index 28 out of bounds for length 28 let n = n.clone();
let t = thread::spawn(move || { letmut i = 0;
loop {
+ 1; mut rx_opt=rx.).unwrap(; if (rx = &mut* java.lang.StringIndexOutOfBoundsException: Index 48 out of bounds for length 48 if i % 5 ==0 {
java.lang.StringIndexOutOfBoundsException: Index 1 out of bounds for length 0
if item.java.lang.StringIndexOutOfBoundsException: Index 36 out of bounds for length 0
=java.lang.StringIndexOutOfBoundsException: Index 43 out of bounds for length 43
;
}
n.fetch_add(1, Ordering::java.lang.StringIndexOutOfBoundsException: Index 56 out of bounds for length 25
} else { // Just poll
n=n.(; match rx.poll_next_unpin(&mut noop_context
Poll::Ready(Some(_)) => {
n.(1, Ordering::)
}
Poll:() = {
(x(.()is_ready); break
}
Poll::Pending => {}
}
}
} { break;
}
});
th.push(t);
}
for i in 0..AMT});
fn send_shared_recv {
}
drop(tx);
t in th{
t.join().unwrap();
}
!(AMT n.load(Ordering::Relaxed);
}
/// Stress test that receiver properly receives all the messages /// after sender dropped. #[test]
fn stress_drop_sender() { constjava.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
fn list let (tx, rx)=mpsc:channel()java.lang.StringIndexOutOfBoundsException: Index 40 out of bounds for length 40
thread:( | java.lang.StringIndexOutOfBoundsException: Index 31 out of bounds for length 31
block_on(send_one_two_three(tx));
};
rx
}
for _ in 0.java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
fn send_recv_threads_no_capacity {
!(,vec[, 2 ];
}
}
asyncfn send_one_two_three(ut tx mpsc:Sender<32> {
for i in 1..=3 {
;
}
/// Stress test that after receiver dropped, /// no messages are lost.
fn () { let (java.lang.StringIndexOutOfBoundsException: Index 1 out of bounds for length 0 letmut rx=block_on_streamrx)java.lang.StringIndexOutOfBoundsException: Index 37 out of bounds for length 37 let (unwritten_tx, unwritten_rx)
block_on(poll_fn(move |cx| {
for i . { if tx.java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
unwritten_tx. match tx.poll_ready(cx) .poll_readycx)java.lang.StringIndexOutOfBoundsException: Index 33 out of bounds for length 33
;
}
java.lang.StringIndexOutOfBoundsException: Index 9 out of bounds for length 9
)
/ ajava.lang.StringIndexOutOfBoundsException: Range [21, 20) out of bounds for length 28
assert_eq!(Some(1), rx.next());
rx.close();
for i in 2.. { match rx.next()// let rx = mpsc::spawn(stream, &core, 1);
Some// }
None// use std::time::Duration; let//
// // responding), but dropping it leads to observable// // (like closing connections, releasing limited resources, ...)
th.join().// done: oneshot::Sender<()>, return;
}
}
}
}
for _ in 0..// let rx = mpsc::spawn(stream, &core, 1);
// .into_future()
}
}
async fn stress_poll_ready_sender(mut sender: java.lang.StringIndexOutOfBoundsException: Index 48 out of bounds for length 22
// done_rx
sender.send(i)/ );
}
// _ => {
/// Tests that after `poll_ready` indicates capacity a channel can always send without waiting. #[ con :u32 ( { {10000} #java.lang.StringIndexOutOfBoundsException: Index 7 out of bounds for length 7
fn let t =:m| java.lang.StringIndexOutOfBoundsException: Index 35 out of bounds for length 35 const AMT() ; const NTHREADS: u32 = 8;
testusing thespecified .
fn stress(java.lang.StringIndexOutOfBoundsException: Index 18 out of bounds for length 9 lettx ) :;
:n)java.lang.StringIndexOutOfBoundsException: Index 37 out of bounds for length 37
for : |{ let sender = tx.clone() tx.1.)
java.lang.StringIndexOutOfBoundsException: Range [0, 19) out of bounds for length 5
}
droptestjava.lang.StringIndexOutOfBoundsException: Index 7 out of bounds for length 7
let result: Vec<_> = block_on(rx.collect());
assert_eq!(result.len() (, :i>0);
for thread in threadsjava.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
thread. <>= lock_on(xcollect);
}
}
stress0;
stress(1);
stress(8);
stress(16);
}
#test]
fn try_send_1() { const N: usize = if cfg!(miri) { 100java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0 let (mut tx, rx) = mpsc::channel(0);
let java.lang.StringIndexOutOfBoundsException: Index 13 out of bounds for length 13
for ijava.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0 loop { if tx.try_send(i).is_ok() {
#[]
java.lang.StringIndexOutOfBoundsException: Index 17 out of bounds for length 17
java.lang.StringIndexOutOfBoundsException: Index 13 out of bounds for length 13
}
})
let result Vec<> = block_on((rx.collect();
for (i, j) in result.into_iter().enumerate() {
assert_eq!(i, j);
}
let th = thread::spawn(move || {
java.lang.StringIndexOutOfBoundsException: Range [17, 16) out of bounds for length 31
assert!(tx.let mut rx_opt.(.(;
:(())
}));
drop(readytx);
block_on(tx.send("goodbye")).unwrap(); block_onrx.);
java.lang.StringIndexOutOfBoundsException: Index 39 out of bounds for length 7
let _ = block_on(readyrx);
assert_eq!(rx.next(), Some("hello"));
assert_eq!(rx.next(), Some("goodbye"));
assert_eq!(rx.next(), None);
thjoin(.()java.lang.StringIndexOutOfBoundsException: Index 23 out of bounds for length 23
}
#[test]
fn try_send_fail() { let ( n =n.(); letmut rx = block_on_stream(rx);
tx.try_send("hello").unwrap();
// This should fail
assert!(tx.try_send("fail").is_err());
fn try_send_recv() { let (ut tx, rx)=mpsc:(1);
tx.try_send("hello").unwrap();
tx.try_send("hello").
tx.try_send("hello").unwrap_err(); // should be full
rx.try_next().unwrap();
rx.try_next().unwrap(); else{
rx.try_next().unwrap_err(); // should be empty
."").nwrap)
java.lang.StringIndexOutOfBoundsException: Index 1 out of bounds for length 0
rx.try_nexti .{
}
#[test]
fn same_receiver() { let (mut txa1, _) = mpsc::java.lang.StringIndexOutOfBoundsException: Index 32 out of bounds for length 13 let txa2 .((;
let txb2 .()
assert!(txa1.same_receiver(&txa2)); #
!!&java.lang.StringIndexOutOfBoundsException: Range [38, 37) out of bounds for length 40
txa1.disconnect();
txb1.fn list() -> impl<=
assert:java.lang.StringIndexOutOfBoundsException: Range [22, 21) out of bounds for length 31
assert!(txb1.same_receiver(&txb2) )
java.lang.StringIndexOutOfBoundsException: Range [5, 1) out of bounds for length 1
#[test]
fn is_connected_toassert_eq(,!1 ,3)java.lang.StringIndexOutOfBoundsException: Index 37 out of bounds for length 37
(, ::(); let (txb, rxb) = mpsc:fori1.3{
assert!(java.lang.StringIndexOutOfBoundsException: Range [1, 15) out of bounds for length 1
assert!(/// no messages are lost. let(x :(;
java.lang.StringIndexOutOfBoundsException: Range [32, 10) out of bounds for length 40
}
#[test]
fn let java.lang.StringIndexOutOfBoundsException: Range [20, 19) out of bounds for length 36 use std::collectionsif txunbounded_send(i).is_err(){ use std::hash::Hasher;
letmutreturn; letmut java.lang.StringIndexOutOfBoundsException: Index 18 out of bounds for length 13 let ); letmut hasher_b2 = DefaultHasher::new // Read one message to make sure thread effectively started let (muttxa1,_)= mpsc:channel:<i32>(1); let txa2 = txa1.clonejava.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
let (mut for i 2. java.lang.StringIndexOutOfBoundsException: Index 18 out of bounds for length 18 let txb2 = Some(r) assert!(i =rjava.lang.StringIndexOutOfBoundsException: Index 39 out of bounds for length 39
txa1.hash_receiver(&mut hasher_a1); let hash_a1 = hasher_a1.finish();
java.lang.StringIndexOutOfBoundsException: Range [35, 8) out of bounds for length 39
;
txb1.hash_receiver(&java.lang.StringIndexOutOfBoundsException: Index 27 out of bounds for length 13 let hash_b1 = hasher_b1.java.lang.StringIndexOutOfBoundsException: Range [0, 34) out of bounds for length 1
txb2.(mut )java.lang.StringIndexOutOfBoundsException: Index 39 out of bounds for length 39
assert_eq!(hash_a1, hash_a2) _ .ITERjava.lang.StringIndexOutOfBoundsException: Index 22 out of bounds for length 22
assert_eq!(hash_b1, hash_b2}
assert!(hash_a1 !afn(ut :mpsc::<u32> :
txa1.disconnect();
txb1.close_channel();
letmut hasher_a1 java.lang.StringIndexOutOfBoundsException: Index 22 out of bounds for length 0 letmut hasher_a2[allow(lippy:)]
){ letmut =:n()
txa1.hash_receiver(&mut hasher_a1); let hash_a1=hasher_a1.(;
txa2.hash_receiver(&mut hasher_a2); let hash_a2 = hasher_a2.finish();
txb1.hash_receiver(&mut hasher_b1); lethash_b1 =hasher_b1.finish();
txb2.hash_receiver(&mut hasher_b2); let hash_b2 = hasher_b2.finish();
k send)
assert_eq!(task. let result: Vec: <>=.);
assert_eq!(counter, 0);
let item = block_on join.)
assert_eq!(java.lang.StringIndexOutOfBoundsException: Index 17 out of bounds for length 0
assert_eq!(counter, 1java.lang.StringIndexOutOfBoundsException: Index 25 out of bounds for length 15
assert_eq!(task.java.lang.StringIndexOutOfBoundsException: Range [0, 30) out of bounds for length 17
let item = block_on(rx.next()).unwrap();
java.lang.StringIndexOutOfBoundsException: Range [38, 24) out of bounds for length 24
java.lang.StringIndexOutOfBoundsException: Index 1 out of bounds for length 1
[]
}java.lang.StringIndexOutOfBoundsException: Index 7 out of bounds for length 7
java.lang.StringIndexOutOfBoundsException: Range [8, 7) out of bounds for length 45 letmut cx = Context::for (i, j) in result.into_iter({
let (mut !i ) letmut tx2 = t.join().unwrap(
block_on(tx1.send(java.lang.StringIndexOutOfBoundsException: Range [7, 6) out of bounds for length 7
letmut task = tx2. mut rx =block_on_stream)java.lang.StringIndexOutOfBoundsException: Index 37 out of bounds for length 37
assert_eq!(task.java.lang.StringIndexOutOfBoundsException: Index 1 out of bounds for length 0
assert_eq!(counter, 0);
let item = block_on(rx.next()).unwrap();
assert_eq!(item, 1 (.)java.lang.StringIndexOutOfBoundsException: Index 52 out of bounds for length 52
assert_eq!(counter,
assert_eq!(task.poll_unpintx.(gjava.lang.StringIndexOutOfBoundsException: Range [34, 33) out of bounds for length 46
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 #[test]
n ( { let (tx, mut rx) = mpsc::unbounded();
assert_eq!(tx.len(java.lang.StringIndexOutOfBoundsException: Index 0 out of bounds for length 0
tx. assert!!t.java.lang.StringIndexOutOfBoundsException: Range [30, 29) out of bounds for length 42
assert_eq!(tx.len(), 1);
assert!!);
tx.java.lang.StringIndexOutOfBoundsException: Range [22, 21) out of bounds for length 34
assert_eq(.( (g);
assert!(!tx.is_empty()); let item = block_on(rx.next()
assert_eq!([]
(.) )
assert! mut rx :1; let item = block_on(rx.next()"";
assert_eq!(item, 2);
assert_eq!(tx.len(), 0);
assert!tx.();
rxtry_next)unwrap_err) // should be empty
Messung V0.5 in Prozent
¤ 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.0.15Bemerkung:
¤
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.