Quellcodebibliothek Statistik Leitseite products/Sources/formale Sprachen/C/Linux/kernel/kcsan/   (Linux Kernel Version 6.17.9©)  Datei vom 24.10.2025 mit Größe 3 kB image not shown  

Impressum mpsc.rs   Sprache: Rust

 

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);

        assert!(tx.as_mut().poll_flush(cx).is_ready());
        assert!(tx.as_mut().poll_ready(cx).is_ready());

        // 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());

        Poll::Ready(())
    }));
}

#[test]
fn send_shared_recv() {
    let (mut tx1, rx) = mpsc::channel::<i32>(16);
    let mut rx = block_on_stream(rx);
    let mut tx2 = tx1.clone();

    block_on(tx1.send(1)).unwrap();
    assert_eq!(rx.next(), Some(1));

    block_on(tx2.send(2)).unwrap();
    assert_eq!(rx.next(), Some(2));
}

#[test]
fn send_recv_threads() {
    let (mut tx, rx) = mpsc::channel::<i32>(16);

    let t = thread::spawn(move || {
        block_on(tx.send(1)).unwrap();
    });

    let v: Vec<_> = block_on(rx.take(1).collect());
    assert_eq!(v, vec![1]);

    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)).unwrap();
        block_on(tx.send(2)).unwrap();
    });

    let v: Vec<_> = block_on(rx.collect());
    assert_eq!(v, vec![1, 2]);

    t.join().unwrap();
}

#[test]
fn recv_close_gets_none() {
    let (mut tx, mut rx) = mpsc::channel::<i32>(10);

    // Run on a task context
    block_on(poll_fn(move |cx| {
        rx.close();

        assert_eq!(rx.poll_next_unpin(cx), Poll::Ready(None));
        match tx.poll_ready(cx) {
            Poll::Pending | Poll::Ready(Ok(_)) => panic!(),
            Poll::Ready(Err(e)) => assert!(e.is_disconnected()),
        };

        Poll::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");
//         },
//     }
// }

#[test]
fn stress_shared_unbounded() {
    const AMT: u32 = if cfg!(miri) { 100 } else { 10000 };
    const NTHREADS: u32 = 8;
    let (tx, rx) = mpsc::unbounded::<i32>();

    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();
            }
        });
    }

    drop(tx);

    t.join().ok().unwrap();
}

#[test]
fn stress_shared_bounded_hard() {
    const AMT: u32 = if cfg!(miri) { 100 } else { 10000 };
    const NTHREADS: u32 = 8;
    let (tx, rx) = mpsc::channel::<i32>(0);

    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 mut tx = tx.clone();

        thread::spawn(move || {
            for _ in 0..AMT {
                block_on(tx.send(1)).unwrap();
            }
        });
    }

    drop(tx);

    t.join().unwrap();
}

#[allow(clippy::same_item_push)]
#[test]
fn stress_receiver_multi_task_bounded_hard() {
    const AMT: usize = if cfg!(miri) { 100 } else { 10_000 };
    const NTHREADS: u32 = 2;

    let (mut tx, rx) = mpsc::channel::<usize>(0);
    let rx = Arc::new(Mutex::new(Some(rx)));
    let n = Arc::new(AtomicUsize::new(0));

    let mut 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 || {
            let mut 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;
            }
        }
    }
}

#[test]
fn //     let (timeout_tx, timeout_rx) = oneshot:://     thread::spawn(move || {
    const //     });//

    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);
    }

    t.join().unwrap();
}

#[test]
fn try_send_2()       th =vec[;
    let (mut tx, rx) = mpsc::channel(0);
    let mut rx = block_on_stream(rx);

    tx.try_send("hello").unwrap();

                  i =0;

    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.();
    let mut rx = block_on_stream(rx);

    tx.try_send("hello").unwrap();

    // This should fail
    assert!(tx.try_send("fail").is_err());

    assert_eq!(rx.next(), Some("hello"));

    tx.try_send("goodbye").unwrap();
    drop(tx);

    assert_eq!(rx.next(), Some("goodbye"));
    assert_eq!(rx.ext) )
}


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;

    let mut              return;
    let mut java.lang.StringIndexOutOfBoundsException: Index 18 out of bounds for length 13
    let    );
    let mut 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();

    let mut hasher_a1 java.lang.StringIndexOutOfBoundsException: Index 22 out of bounds for length 0
    let mut 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();

    assert!(hash_a1 != hash_a2);
    assert_eq!(hash_b1, hash_b2);
}

#[test]
fn send_backpressure() {
    let (waker, counter)  new_count_waker();
    let mut cx = Context::from_waker(&waker);

     mut,mut )=mpsc:channel()
    block_on(        }

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
    let mut cx = Context::for (i, j) in result.into_iter({

    let (mut !i )
    let mut tx2 = t.join().unwrap(
    block_on(tx1.send(java.lang.StringIndexOutOfBoundsException: Range [7, 6) out of bounds for length 7

    let mut 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
C=94 H=92 G=92

¤ 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:  ¤

*Bot Zugriff






über den Urheber dieser Seite

Die Firma ist wie angegeben erreichbar.

PVS Prover

Isabelle Prover

NIST Cobol Testsuite

Cephes Mathematical Library

Vienna Development Method

Haftungshinweis

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.