lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

socket_write_tests.rs (5130B)


      1 use super::*;
      2 use futures::{FutureExt, channel::mpsc, task::noop_waker_ref};
      3 
      4 fn message(text: &str) -> Message {
      5     Message::Text(text.to_owned())
      6 }
      7 
      8 fn channel_writer() -> (Arc<SocketWriter>, mpsc::Receiver<Message>) {
      9     let (tx, rx) = mpsc::channel(0);
     10     (
     11         SocketWriter::new(Box::new(tx.sink_map_err(TransportError::backend))),
     12         rx,
     13     )
     14 }
     15 
     16 #[test]
     17 fn concurrent_sdk_and_exact_messages_share_backpressure_and_order() {
     18     use futures::StreamExt;
     19     futures::executor::block_on(async {
     20         let (writer, mut rx) = channel_writer();
     21         let mut sdk = SharedSocketSink::new(Arc::clone(&writer));
     22         let mut first = Box::pin(sdk.send(message("AUTH")));
     23         assert!(first.as_mut().now_or_never().is_none());
     24         let mut raw = Box::pin(writer.send(message("EVENT")));
     25         assert!(raw.as_mut().now_or_never().is_none());
     26         assert_eq!(rx.next().await, Some(message("AUTH")));
     27         first.await.unwrap();
     28         assert!(raw.as_mut().now_or_never().is_none());
     29         assert_eq!(rx.next().await, Some(message("EVENT")));
     30         raw.await.unwrap();
     31         sdk.close().await.unwrap();
     32         sdk.close().await.unwrap();
     33         assert!(sdk.send(message("REQ")).await.is_err());
     34         assert!(writer.send(message("EVENT")).await.is_err());
     35         assert!(rx.next().await.is_none());
     36     });
     37 }
     38 
     39 #[test]
     40 fn cancelling_a_waiter_releases_its_lock_without_claiming_no_effect() {
     41     use futures::StreamExt;
     42     futures::executor::block_on(async {
     43         let (writer, mut rx) = channel_writer();
     44         let mut sdk = SharedSocketSink::new(Arc::clone(&writer));
     45         let mut pending = Box::pin(writer.send(message("first")));
     46         assert!(pending.as_mut().now_or_never().is_none());
     47         drop(pending);
     48         // It was already buffered: cancelling a send is not proof of no publication.
     49         assert_eq!(rx.next().await, Some(message("first")));
     50         let mut next = Box::pin(sdk.send(message("second")));
     51         assert!(next.as_mut().now_or_never().is_none());
     52         assert_eq!(rx.next().await, Some(message("second")));
     53         next.await.unwrap();
     54     });
     55 }
     56 
     57 #[test]
     58 fn pending_sdk_message_must_flush_before_next_send_or_close() {
     59     use futures::StreamExt;
     60     futures::executor::block_on(async {
     61         let (writer, mut rx) = channel_writer();
     62         let mut sdk = SharedSocketSink::new(writer);
     63         let mut context = Context::from_waker(noop_waker_ref());
     64         assert!(Pin::new(&mut sdk).poll_ready(&mut context).is_ready());
     65         Pin::new(&mut sdk).start_send(message("first")).unwrap();
     66         assert!(Pin::new(&mut sdk).start_send(message("unready")).is_err());
     67         assert!(Pin::new(&mut sdk).poll_close(&mut context).is_pending());
     68         assert_eq!(rx.next().await, Some(message("first")));
     69         sdk.close().await.unwrap();
     70         assert!(Pin::new(&mut sdk).start_send(message("closed")).is_err());
     71     });
     72 }
     73 
     74 #[test]
     75 fn sdk_drop_revokes_existing_raw_handles_and_failed_io_revokes_writer() {
     76     futures::executor::block_on(async {
     77         let (writer, rx) = channel_writer();
     78         drop(SharedSocketSink::new(Arc::clone(&writer)));
     79         assert!(writer.send(message("after drop")).await.is_err());
     80         drop(rx);
     81         let (writer, rx) = channel_writer();
     82         drop(rx);
     83         let mut sdk = SharedSocketSink::new(Arc::clone(&writer));
     84         assert!(sdk.send(message("failure")).await.is_err());
     85         assert!(!writer.open.load(Ordering::Acquire));
     86         assert!(writer.send(message("after error")).await.is_err());
     87     });
     88 }
     89 
     90 #[test]
     91 fn registry_is_configured_bounded_weak_and_reconnect_isolated() {
     92     let registry = WriterRegistry::new(["relay".to_owned()].into_iter());
     93     assert_eq!(format!("{registry:?}"), "WriterRegistry([redacted])");
     94     assert!(registry.get("relay").is_err());
     95     assert!(registry.get("unknown").is_err());
     96     let (first, _rx) = channel_writer();
     97     assert!(registry.install("unknown", &first).is_err());
     98     registry.install("relay", &first).unwrap();
     99     let retained = registry.get("relay").unwrap();
    100     assert!(Arc::ptr_eq(&retained, &first));
    101     let first_sdk = SharedSocketSink::new(Arc::clone(&first));
    102     let (second, _rx) = channel_writer();
    103     registry.install("relay", &second).unwrap();
    104     assert!(!retained.open.load(Ordering::Acquire));
    105     drop(first_sdk);
    106     assert!(!retained.open.load(Ordering::Acquire));
    107     assert!(Arc::ptr_eq(&registry.get("relay").unwrap(), &second));
    108     second.invalidate();
    109     assert!(registry.get("relay").is_err());
    110     drop(second);
    111     assert!(registry.get("relay").is_err());
    112     assert_eq!(registry.0.lock().unwrap().len(), 1);
    113 }
    114 
    115 #[test]
    116 fn poisoned_registry_fails_closed() {
    117     let registry = WriterRegistry::new(["relay".to_owned()].into_iter());
    118     let other = registry.clone();
    119     assert!(
    120         std::panic::catch_unwind(move || {
    121             let _guard = other.0.lock().unwrap();
    122             panic!("fixture registry poison");
    123         })
    124         .is_err()
    125     );
    126     let (writer, _rx) = channel_writer();
    127     assert!(registry.install("relay", &writer).is_err());
    128     assert!(registry.get("relay").is_err());
    129 }