//! [`SinkCloseChannel`] use std::pin::Pin; use futures::Sink; use futures::channel::mpsc; //---------- principal API ---------- /// A [`Sink`] with a `close_channel` method like [`futures::channel::mpsc::Sender`]'s. pub trait SinkCloseChannel: Sink { /// Close the channel from the sending end, giving EOF at the receiver /// /// Future attempts to send will get a disconnected error. /// /// This closes *all* equivalent senders for the underlying data sink. /// For example, if `Self` is `Clone`, all clones are affected. /// /// If the Sink is for a channel, /// the receiver will see EOF after reading the messages that were successfully sent so far. fn close_channel(self: Pin<&mut Self>); } //---------- impl for futures::channel::mpsc ---------- impl SinkCloseChannel for mpsc::Sender { fn close_channel(self: Pin<&mut Self>) { let self_: &mut Self = Pin::into_inner(self); self_.close_channel(); } } #[cfg(test)] mod test { // @@ begin test lint list maintained by maint/add_warning @@ #![allow(clippy::bool_assert_comparison)] #![allow(clippy::clone_on_copy)] #![allow(clippy::dbg_macro)] #![allow(clippy::mixed_attributes_style)] #![allow(clippy::print_stderr)] #![allow(clippy::print_stdout)] #![allow(clippy::single_char_pattern)] #![allow(clippy::unwrap_used)] #![allow(clippy::unchecked_time_subtraction)] #![allow(clippy::useless_vec)] #![allow(clippy::needless_pass_by_value)] //! #![allow(clippy::arithmetic_side_effects)] // don't mind potential panicking ops in tests #![allow(clippy::useless_format)] // sorely use super::*; use futures::{SinkExt as _, StreamExt as _}; #[test] fn close_channel() { tor_rtmock::MockRuntime::test_with_various(|_rt| async move { let (mut tx, mut rx) = mpsc::channel::(20); tx.send(0).await.unwrap(); let mut tx2 = tx.clone(); tx2.send(1).await.unwrap(); tx2.close_channel(); let _: mpsc::SendError = tx.send(66).await.unwrap_err(); for i in 0..=1 { assert_eq!(rx.next().await.unwrap(), i); } assert_eq!(rx.next().await, None); }); } }