//! Tests copied from `std::sync::mpsc`. //! //! This is a copy of tests for the `std::sync::mpsc` channels from the standard library, but //! modified to work with `crossbeam-channel` instead. //! //! Minor tweaks were needed to make the tests compile: //! //! - Replace `box` syntax with `Box::new`. //! - Replace all uses of `Select` with `select!`. //! - Change the imports. //! - Join all spawned threads. //! - Removed assertion from oneshot_multi_thread_send_close_stress tests. //! //! Source: //! - https://github.com/rust-lang/rust/tree/master/src/libstd/sync/mpsc //! //! Copyright & License: //! - Copyright 2013-2014 The Rust Project Developers //! - Apache License, Version 2.0 or MIT license, at your option //! - https://github.com/rust-lang/rust/blob/master/COPYRIGHT //! - https://www.rust-lang.org/en-US/legal.html
use std::sync::mpsc::{RecvError, RecvTimeoutError, TryRecvError}; use std::sync::mpsc::{SendError, TrySendError}; use std::thread::JoinHandle; use std::time::Duration;
pubfn channel<T>() -> (Sender<T>, Receiver<T>) { let (s, r) = cc::unbounded(); let s = Sender { inner: s }; let r = Receiver { inner: r };
(s, r)
}
pubfn sync_channel<T>(bound: usize) -> (SyncSender<T>, Receiver<T>) { let (s, r) = cc::bounded(bound); let s = SyncSender { inner: s }; let r = Receiver { inner: r };
(s, r)
}
let (tx, rx) = channel::<i32>(); let t = thread::spawn(move || { for _ in0..COUNT {
tx.send(1).unwrap();
}
}); for _ in0..COUNT {
assert_eq!(rx.recv().unwrap(), 1);
}
t.join().ok().unwrap();
}
#[test] fn stress_shared() { let amt: u32 = if cfg!(miri) { 100 } else { 10_000 }; let nthreads: u32 = if cfg!(miri) { 4 } else { 8 }; let (tx, rx) = channel::<i32>();
let t = thread::spawn(move || { for _ in0..amt * nthreads {
assert_eq!(rx.recv().unwrap(), 1);
}
assert!(rx.try_recv().is_err());
});
letmut ts = Vec::with_capacity(nthreads as usize); for _ in0..nthreads { let tx = tx.clone(); let t = thread::spawn(move || { for _ in0..amt {
tx.send(1).unwrap();
}
});
ts.push(t);
}
drop(tx);
t.join().ok().unwrap(); for t in ts {
t.join().unwrap();
}
}
#[test] fn send_from_outside_runtime() { let (tx1, rx1) = channel::<()>(); let (tx2, rx2) = channel::<i32>(); let t1 = thread::spawn(move || {
tx1.send(()).unwrap(); for _ in0..40 {
assert_eq!(rx2.recv().unwrap(), 1);
}
});
rx1.recv().unwrap(); let t2 = thread::spawn(move || { for _ in0..40 {
tx2.send(1).unwrap();
}
});
t1.join().ok().unwrap();
t2.join().ok().unwrap();
}
#[test] fn recv_from_outside_runtime() { let (tx, rx) = channel::<i32>(); let t = thread::spawn(move || { for _ in0..40 {
assert_eq!(rx.recv().unwrap(), 1);
}
}); for _ in0..40 {
tx.send(1).unwrap();
}
t.join().ok().unwrap();
}
#[test] fn no_runtime() { let (tx1, rx1) = channel::<i32>(); let (tx2, rx2) = channel::<i32>(); let t1 = thread::spawn(move || {
assert_eq!(rx1.recv().unwrap(), 1);
tx2.send(2).unwrap();
}); let t2 = thread::spawn(move || {
tx1.send(1).unwrap();
assert_eq!(rx2.recv().unwrap(), 2);
});
t1.join().ok().unwrap();
t2.join().ok().unwrap();
}
#[test] fn oneshot_single_thread_close_port_first() { // Simple test of closing without sending let (_tx, rx) = channel::<i32>();
drop(rx);
}
#[test] fn oneshot_single_thread_close_chan_first() { // Simple test of closing without sending let (tx, _rx) = channel::<i32>();
drop(tx);
}
#[test] fn oneshot_single_thread_send_port_close() { // Testing that the sender cleans up the payload if receiver is closed let (tx, rx) = channel::<Box<i32>>();
drop(rx);
assert!(tx.send(Box::new(0)).is_err());
}
#[test] fn oneshot_multi_task_recv_then_close() { let (tx, rx) = channel::<Box<i32>>(); let t = thread::spawn(move || {
drop(tx);
});
thread::spawn(move || {
assert_eq!(rx.recv(), Err(RecvError));
})
.join()
.unwrap();
t.join().unwrap();
}
#[test] fn oneshot_multi_thread_close_stress() { let stress_factor = stress_factor(); letmut ts = Vec::with_capacity(stress_factor); for _ in0..stress_factor { let (tx, rx) = channel::<i32>(); let t = thread::spawn(move || {
drop(rx);
});
ts.push(t);
drop(tx);
} for t in ts {
t.join().unwrap();
}
}
#[test] fn oneshot_multi_thread_send_close_stress() { let stress_factor = stress_factor(); letmut ts = Vec::with_capacity(2 * stress_factor); for _ in0..stress_factor { let (tx, rx) = channel::<i32>(); let t = thread::spawn(move || {
drop(rx);
});
ts.push(t);
thread::spawn(move || { let _ = tx.send(1);
})
.join()
.unwrap();
} for t in ts {
t.join().unwrap();
}
}
#[test] fn oneshot_multi_thread_recv_close_stress() { let stress_factor = stress_factor(); letmut ts = Vec::with_capacity(2 * stress_factor); for _ in0..stress_factor { let (tx, rx) = channel::<i32>(); let t = thread::spawn(move || {
thread::spawn(move || {
assert_eq!(rx.recv(), Err(RecvError));
})
.join()
.unwrap();
});
ts.push(t); let t2 = thread::spawn(move || { let t = thread::spawn(move || {
drop(tx);
});
t.join().unwrap();
});
ts.push(t2);
} for t in ts {
t.join().unwrap();
}
}
#[test] fn oneshot_multi_thread_send_recv_stress() { let stress_factor = stress_factor(); letmut ts = Vec::with_capacity(stress_factor); for _ in0..stress_factor { let (tx, rx) = channel::<Box<isize>>(); let t = thread::spawn(move || {
tx.send(Box::new(10)).unwrap();
});
ts.push(t);
assert!(*rx.recv().unwrap() == 10);
} for t in ts {
t.join().unwrap();
}
}
#[test] fn stream_send_recv_stress() { let stress_factor = stress_factor(); letmut ts = Vec::with_capacity(2 * stress_factor); for _ in0..stress_factor { let (tx, rx) = channel();
#[test] fn stress_recv_timeout_shared() { let (tx, rx) = channel(); let stress = stress_factor() + 100;
letmut ts = Vec::with_capacity(stress); for i in0..stress { let tx = tx.clone(); let t = thread::spawn(move || {
thread::sleep(Duration::from_millis(i as u64 * 10));
tx.send(1usize).unwrap();
});
ts.push(t);
}
// Regression test that we don't run out of stack in scheduler context let (tx, rx) = channel(); for _ in0..N {
tx.send(()).unwrap();
} for _ in0..N {
rx.recv().unwrap();
}
}
#[test] fn shared_recv_timeout() { let (tx, rx) = channel(); let total = 5; letmut ts = Vec::with_capacity(total); for _ in0..total { let tx = tx.clone(); let t = thread::spawn(move || {
tx.send(()).unwrap();
});
ts.push(t);
}
for _ in0..total {
rx.recv().unwrap();
}
assert_eq!(
rx.recv_timeout(Duration::from_millis(1)),
Err(RecvTimeoutError::Timeout)
);
tx.send(()).unwrap();
assert_eq!(rx.recv_timeout(Duration::from_millis(1)), Ok(())); for t in ts {
t.join().unwrap();
}
}
#[test] fn shared_chan_stress() { let (tx, rx) = channel(); let total = stress_factor() + 100; letmut ts = Vec::with_capacity(total); for _ in0..total { let tx = tx.clone(); let t = thread::spawn(move || {
tx.send(()).unwrap();
});
ts.push(t);
}
for _ in0..total {
rx.recv().unwrap();
} for t in ts {
t.join().unwrap();
}
}
#[test] fn test_nested_recv_iter() { let (tx, rx) = channel::<i32>(); let (total_tx, total_rx) = channel::<i32>();
let t = thread::spawn(move || { letmut acc = 0; for x in rx.iter() {
acc += x;
}
total_tx.send(acc).unwrap();
});
#[test] fn test_recv_try_iter() { let (request_tx, request_rx) = channel(); let (response_tx, response_rx) = channel();
// Request `x`s until we have `6`. let t = thread::spawn(move || { letmut count = 0; loop { for x in response_rx.try_iter() {
count += x; if count == 6 { return count;
}
}
request_tx.send(()).unwrap();
}
});
for _ in request_rx.iter() { if response_tx.send(2).is_err() { break;
}
}
assert_eq!(t.join().unwrap(), 6);
}
#[test] fn test_recv_into_iter_owned() { letmut iter = { let (tx, rx) = channel::<i32>();
tx.send(1).unwrap();
tx.send(2).unwrap();
// This bug used to end up in a livelock inside of the Receiver destructor // because the internal state of the Shared packet was corrupted #[test] fn destroy_upgraded_shared_port_when_sender_still_active() { let (tx, rx) = channel(); let (tx2, rx2) = channel(); let t = thread::spawn(move || {
rx.recv().unwrap(); // wait on a oneshot
drop(rx); // destroy a shared
tx2.send(()).unwrap();
}); // make sure the other thread has gone to sleep for _ in0..5000 {
thread::yield_now();
}
// upgrade to a shared chan and send a message let tx2 = tx.clone();
drop(tx);
tx2.send(()).unwrap();
// wait for the child thread to exit before we exit
rx2.recv().unwrap();
t.join().unwrap();
}
#[test] fn issue_32114() { let (tx, _) = channel(); let _ = tx.send(123);
assert_eq!(tx.send(123), Err(SendError(123)));
}
}
assert_eq!(recv_count, AMT * NTHREADS);
assert!(rx.try_recv().is_err());
dtx.send(()).unwrap();
});
letmut ts = Vec::with_capacity(NTHREADS as usize); for _ in0..NTHREADS { let tx = tx.clone(); let t = thread::spawn(move || { for _ in0..AMT {
tx.send(1).unwrap();
}
});
ts.push(t);
}
drop(tx);
drx.recv().unwrap(); for t in ts {
t.join().unwrap();
}
t.join().unwrap();
}
let t = thread::spawn(move || { for _ in0..AMT * NTHREADS {
assert_eq!(rx.recv().unwrap(), 1);
}
assert!(rx.try_recv().is_err());
dtx.send(()).unwrap();
});
letmut ts = Vec::with_capacity(NTHREADS as usize); for _ in0..NTHREADS { let tx = tx.clone(); let t = thread::spawn(move || { for _ in0..AMT {
tx.send(1).unwrap();
}
});
ts.push(t);
}
drop(tx);
drx.recv().unwrap(); for t in ts {
t.join().unwrap();
}
t.join().unwrap();
}
#[test] fn oneshot_single_thread_close_port_first() { // Simple test of closing without sending let (_tx, rx) = sync_channel::<i32>(0);
drop(rx);
}
#[test] fn oneshot_single_thread_close_chan_first() { // Simple test of closing without sending let (tx, _rx) = sync_channel::<i32>(0);
drop(tx);
}
#[test] fn oneshot_single_thread_send_port_close() { // Testing that the sender cleans up the payload if receiver is closed let (tx, rx) = sync_channel::<Box<i32>>(0);
drop(rx);
assert!(tx.send(Box::new(0)).is_err());
}
#[test] fn oneshot_multi_task_recv_then_close() { let (tx, rx) = sync_channel::<Box<i32>>(0); let t = thread::spawn(move || {
drop(tx);
});
thread::spawn(move || {
assert_eq!(rx.recv(), Err(RecvError));
})
.join()
.unwrap();
t.join().unwrap();
}
#[test] fn oneshot_multi_thread_close_stress() { let stress_factor = stress_factor(); letmut ts = Vec::with_capacity(stress_factor); for _ in0..stress_factor { let (tx, rx) = sync_channel::<i32>(0); let t = thread::spawn(move || {
drop(rx);
});
ts.push(t);
drop(tx);
} for t in ts {
t.join().unwrap();
}
}
#[test] fn oneshot_multi_thread_send_close_stress() { let stress_factor = stress_factor(); letmut ts = Vec::with_capacity(stress_factor); for _ in0..stress_factor { let (tx, rx) = sync_channel::<i32>(0); let t = thread::spawn(move || {
drop(rx);
});
ts.push(t);
thread::spawn(move || { let _ = tx.send(1);
})
.join()
.unwrap();
} for t in ts {
t.join().unwrap();
}
}
#[test] fn oneshot_multi_thread_recv_close_stress() { let stress_factor = stress_factor(); letmut ts = Vec::with_capacity(2 * stress_factor); for _ in0..stress_factor { let (tx, rx) = sync_channel::<i32>(0); let t = thread::spawn(move || {
thread::spawn(move || {
assert_eq!(rx.recv(), Err(RecvError));
})
.join()
.unwrap();
});
ts.push(t); let t2 = thread::spawn(move || {
thread::spawn(move || {
drop(tx);
});
});
ts.push(t2);
} for t in ts {
t.join().unwrap();
}
}
#[test] fn oneshot_multi_thread_send_recv_stress() { let stress_factor = stress_factor(); letmut ts = Vec::with_capacity(stress_factor); for _ in0..stress_factor { let (tx, rx) = sync_channel::<Box<i32>>(0); let t = thread::spawn(move || {
tx.send(Box::new(10)).unwrap();
});
ts.push(t);
assert!(*rx.recv().unwrap() == 10);
} for t in ts {
t.join().unwrap();
}
}
#[test] fn stream_send_recv_stress() { let stress_factor = stress_factor(); letmut ts = Vec::with_capacity(2 * stress_factor); for _ in0..stress_factor { let (tx, rx) = sync_channel::<Box<i32>>(0);
// Regression test that we don't run out of stack in scheduler context let (tx, rx) = sync_channel(N); for _ in0..N {
tx.send(()).unwrap();
} for _ in0..N {
rx.recv().unwrap();
}
}
#[test] fn shared_chan_stress() { let (tx, rx) = sync_channel(0); let total = stress_factor() + 100; letmut ts = Vec::with_capacity(total); for _ in0..total { let tx = tx.clone(); let t = thread::spawn(move || {
tx.send(()).unwrap();
});
ts.push(t);
}
for _ in0..total {
rx.recv().unwrap();
} for t in ts {
t.join().unwrap();
}
}
#[test] fn test_nested_recv_iter() { let (tx, rx) = sync_channel::<i32>(0); let (total_tx, total_rx) = sync_channel::<i32>(0);
let t = thread::spawn(move || { letmut acc = 0; for x in rx.iter() {
acc += x;
}
total_tx.send(acc).unwrap();
});
// This bug used to end up in a livelock inside of the Receiver destructor // because the internal state of the Shared packet was corrupted #[test] fn destroy_upgraded_shared_port_when_sender_still_active() { let (tx, rx) = sync_channel::<()>(0); let (tx2, rx2) = sync_channel::<()>(0); let t = thread::spawn(move || {
rx.recv().unwrap(); // wait on a oneshot
drop(rx); // destroy a shared
tx2.send(()).unwrap();
}); // make sure the other thread has gone to sleep for _ in0..5000 {
thread::yield_now();
}
// upgrade to a shared chan and send a message let tx2 = tx.clone();
drop(tx);
tx2.send(()).unwrap();
// wait for the child thread to exit before we exit
rx2.recv().unwrap();
t.join().unwrap();
}
#[test] fn send1() { let (tx, rx) = sync_channel::<i32>(0); let t = thread::spawn(move || {
rx.recv().unwrap();
});
assert_eq!(tx.send(1), Ok(()));
t.join().unwrap();
}
#[test] fn send2() { let (tx, rx) = sync_channel::<i32>(0); let t = thread::spawn(move || {
drop(rx);
});
assert!(tx.send(1).is_err());
t.join().unwrap();
}
#[test] fn send3() { let (tx, rx) = sync_channel::<i32>(1);
assert_eq!(tx.send(1), Ok(())); let t = thread::spawn(move || {
drop(rx);
});
assert!(tx.send(1).is_err());
t.join().unwrap();
}
#[test] fn send4() { let (tx, rx) = sync_channel::<i32>(0); let tx2 = tx.clone(); let (done, donerx) = channel(); let done2 = done.clone(); let t = thread::spawn(move || {
assert!(tx.send(1).is_err());
done.send(()).unwrap();
}); let t2 = thread::spawn(move || {
assert!(tx2.send(2).is_err());
done2.send(()).unwrap();
});
drop(rx);
donerx.recv().unwrap();
donerx.recv().unwrap();
t.join().unwrap();
t2.join().unwrap();
}
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.