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(®istry.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 }