Структура Sender
pub struct Sender<T> { /* private fields */ }
sync.Отправляет значения связанному Receiver.
Экземпляры создаются функцией channel.
Чтобы преобразовать Sender в Sink или использовать его в функции опроса, можно воспользоваться утилитой PollSender.
Реализации
impl<T> Sender<T>
pub async fn send(&self, value: T) -> Result<(), SendError<T>>
Отправляет значение, ожидая появления свободной ёмкости.
Отправка считается успешной, если установлено, что другая сторона канала ещё не закрылась. Отправка считается неуспешной, если соответствующий получатель уже закрыт. Обратите внимание: возвращаемое значение Err означает, что данные никогда не будут получены, однако возвращаемое значение Ok не означает, что данные будут получены. Соответствующий получатель может закрыться сразу после того, как эта функция вернёт Ok.
Ошибки
Если принимающая часть канала закрыта — либо из-за вызова close, либо из-за удаления дескриптора Receiver, — функция возвращает ошибку. Ошибка содержит значение, переданное в send.
Безопасность при отмене
Если send используется как ветвь в tokio::select! и первой завершается другая ветвь, гарантируется, что сообщение не было отправлено. Однако в этом случае сообщение удаляется и теряется.
Чтобы избежать потери сообщений, используйте reserve для резервирования ёмкости, а затем отправьте сообщение с помощью возвращённого Permit.
В этом канале используется очередь, чтобы гарантировать, что вызовы send и reserve завершаются в порядке их поступления. Отмена вызова send приводит к потере места в очереди.
Примеры
В следующем примере каждый вызов send блокируется до тех пор, пока ранее отправленное значение не будет получено.
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(1);
tokio::spawn(async move {
for i in 0..10 {
if let Err(_) = tx.send(i).await {
println!("receiver dropped");
return;
}
}
});
while let Some(i) = rx.recv().await {
println!("got = {}", i);
}pub async fn closed(&self)
Завершается, когда получатель удалён.
Это позволяет отправителям узнать, что интерес к отправляемым значениям пропал, и немедленно прекратить работу.
Безопасность при отмене
Этот метод безопасен при отмене. После закрытия канала он остаётся закрытым навсегда, и все последующие вызовы closed возвращаются немедленно.
Примеры
use tokio::sync::mpsc;
let (tx1, rx) = mpsc::channel::<()>(1);
let tx2 = tx1.clone();
let tx3 = tx1.clone();
let tx4 = tx1.clone();
let tx5 = tx1.clone();
tokio::spawn(async move {
drop(rx);
});
futures::join!(
tx1.closed(),
tx2.closed(),
tx3.closed(),
tx4.closed(),
tx5.closed()
);
println!("Receiver dropped");pub fn try_send(&self, message: T) -> Result<(), TrySendError<T>>
Пытается немедленно отправить сообщение в этот Sender.
Этот метод отличается от send тем, что немедленно возвращает управление, если буфер канала заполнен или ни один получатель не ожидает получения данных. По сравнению с send эта функция имеет два случая ошибки вместо одного (один для отключения соединения, другой для заполненного буфера).
Ошибки
Если достигнута ёмкость канала, то есть канал содержит n буферизованных значений, где n — аргумент, переданный в channel, возвращается ошибка.
Если принимающая часть канала закрыта — либо из-за вызова close, либо из-за удаления дескриптора Receiver, — функция возвращает ошибку. Ошибка содержит значение, переданное в send.
Примеры
use tokio::sync::mpsc;
// Create a channel with buffer size 1
let (tx1, mut rx) = mpsc::channel(1);
let tx2 = tx1.clone();
tokio::spawn(async move {
tx1.send(1).await.unwrap();
tx1.send(2).await.unwrap();
// task waits until the receiver receives a value.
});
tokio::spawn(async move {
// This will return an error and send
// no message if the buffer is full
let _ = tx2.try_send(3);
});
let mut msg;
msg = rx.recv().await.unwrap();
println!("message {} received", msg);
msg = rx.recv().await.unwrap();
println!("message {} received", msg);
// Third message may have never been sent
match rx.recv().await {
Some(msg) => println!("message {} received", msg),
None => println!("the third message was never sent"),
}pub async fn send_timeout( &self, value: T, timeout: Duration, ) -> Result<(), SendTimeoutError<T>>
time.Отправляет значение, ожидая появления свободной ёмкости, но только в течение ограниченного времени.
Условия успеха и ошибки такие же, как у send; добавляется ещё одно условие неуспешной отправки: истёк указанный тайм-аут, а свободной ёмкости нет.
Ошибки
Если принимающая часть канала закрыта — либо из-за вызова close, либо из-за удаления Receiver, — функция возвращает ошибку. Ошибка содержит значение, переданное в send.
Паники
Эта функция вызывает панику, если её вызвать вне контекста среды выполнения Tokio, в которой включён механизм времени.
Примеры
В следующем примере каждый вызов send_timeout блокируется до тех пор, пока ранее отправленное значение не будет получено, если только не истечёт тайм-аут.
use tokio::sync::mpsc;
use tokio::time::{sleep, Duration};
let (tx, mut rx) = mpsc::channel(1);
tokio::spawn(async move {
for i in 0..10 {
if let Err(e) = tx.send_timeout(i, Duration::from_millis(100)).await {
println!("send error: #{:?}", e);
return;
}
}
});
while let Some(i) = rx.recv().await {
println!("got = {}", i);
sleep(Duration::from_millis(200)).await;
}pub fn blocking_send(&self, value: T) -> Result<(), SendError<T>>
Блокирующая отправка для вызова вне асинхронного контекста.
Этот метод предназначен для случаев, когда данные отправляются из синхронного кода в асинхронный, и работает, даже если получатель использует для получения сообщения не blocking_recv.
Паника
Эта функция вызывает панику, если её вызвать в асинхронном контексте выполнения.
Примеры
use std::thread;
use tokio::runtime::Runtime;
use tokio::sync::mpsc;
fn main() {
let (tx, mut rx) = mpsc::channel::<u8>(1);
let sync_code = thread::spawn(move || {
tx.blocking_send(10).unwrap();
});
Runtime::new().unwrap().block_on(async move {
assert_eq!(Some(10), rx.recv().await);
});
sync_code.join().unwrap()
}pub fn is_closed(&self) -> bool
Проверяет, закрыт ли канал. Это происходит, когда Receiver удаляется или вызывается метод Receiver::close.
let (tx, rx) = tokio::sync::mpsc::channel::<()>(42);
assert!(!tx.is_closed());
let tx2 = tx.clone();
assert!(!tx2.is_closed());
drop(rx);
assert!(tx.is_closed());
assert!(tx2.is_closed());pub async fn reserve(&self) -> Result<Permit<'_, T>, SendError<()>>
Ожидает появления свободной ёмкости в канале. Когда появляется возможность отправить одно сообщение, она резервируется для вызывающего кода.
Если канал заполнен, функция ожидает, пока количество непринятых сообщений не станет меньше ёмкости канала. Возможность отправить одно сообщение резервируется для вызывающего кода. Для отслеживания зарезервированной ёмкости возвращается Permit. Функция send у Permit использует зарезервированную ёмкость.
Удаление Permit без отправки сообщения возвращает ёмкость каналу.
Безопасность при отмене
Этот канал использует очередь, чтобы гарантировать, что вызовы send и reserve завершаются в том порядке, в котором они были запрошены. Отмена вызова reserve приводит к потере места в очереди.
Примеры
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(1);
// Reserve capacity
let permit = tx.reserve().await.unwrap();
// Trying to send directly on the `tx` will fail due to no
// available capacity.
assert!(tx.try_send(123).is_err());
// Sending on the permit succeeds
permit.send(456);
// The value sent on the permit is received
assert_eq!(rx.recv().await.unwrap(), 456);pub async fn reserve_many( &self, n: usize, ) -> Result<PermitIterator<'_, T>, SendError<()>>
Ожидает появления свободной ёмкости в канале. Когда появляется возможность отправить n сообщений, она резервируется для вызывающего кода.
Если канал заполнен или доступно меньше n разрешений, функция ожидает, пока количество непринятых сообщений не станет n меньше ёмкости канала. Затем для вызывающего кода резервируется возможность отправить n сообщение.
Для отслеживания зарезервированной ёмкости возвращается PermitIterator. Вы можете вызывать этот Iterator, пока он не будет исчерпан, чтобы получить Permit, а затем вызвать Permit::send. Эта функция похожа на try_reserve_many, но ожидает, пока свободные места не станут доступны.
Если канал закрыт, функция возвращает SendError.
Удаление PermitIterator до того, как он будет полностью использован, возвращает оставшиеся разрешения каналу.
Безопасность при отмене
Этот канал использует очередь, чтобы гарантировать, что вызовы send и reserve_many завершаются в том порядке, в котором они были запрошены. Отмена вызова reserve_many приводит к потере места в очереди.
Примеры
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(2);
// Reserve capacity
let mut permit = tx.reserve_many(2).await.unwrap();
// Trying to send directly on the `tx` will fail due to no
// available capacity.
assert!(tx.try_send(123).is_err());
// Sending with the permit iterator succeeds
permit.next().unwrap().send(456);
permit.next().unwrap().send(457);
// The iterator should now be exhausted
assert!(permit.next().is_none());
// The value sent on the permit is received
assert_eq!(rx.recv().await.unwrap(), 456);
assert_eq!(rx.recv().await.unwrap(), 457);pub async fn reserve_owned(self) -> Result<OwnedPermit<T>, SendError<()>>
Ожидает появления свободного места в канале, перемещая Sender и возвращая принадлежащее вызывающему коду разрешение. Как только появляется возможность отправить одно сообщение, это место резервируется для вызывающего кода.
Этот метод перемещает отправителя по значению и возвращает принадлежащее вызывающему коду разрешение, которое можно использовать для отправки сообщения в канал. В отличие от Sender::reserve, этот метод можно использовать в случаях, когда разрешение должно быть действительно на протяжении всего времени жизни 'static. Sender можно недорого клонировать (Sender::clone по сути является увеличением счётчика ссылок, сравнимым с Arc::clone), поэтому, если требуется несколько OwnedPermit или Sender нельзя переместить, его можно клонировать перед вызовом reserve_owned.
Если канал заполнен, функция ожидает, пока количество неполученных сообщений не станет меньше ёмкости канала. Место для отправки одного сообщения резервируется для вызывающего кода. Для отслеживания зарезервированного места возвращается OwnedPermit. Функция send у OwnedPermit использует зарезервированное место.
Если удалить OwnedPermit, не отправив сообщение, место освобождается и возвращается каналу.
Безопасность отмены
В этом канале используется очередь, чтобы вызовы send и reserve завершались в порядке их поступления. Отмена вызова reserve_owned приводит к потере места в очереди.
Примеры
Отправка сообщения с помощью OwnedPermit:
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(1);
// Reserve capacity, moving the sender.
let permit = tx.reserve_owned().await.unwrap();
// Send a message, consuming the permit and returning
// the moved sender.
let tx = permit.send(123);
// The value sent on the permit is received.
assert_eq!(rx.recv().await.unwrap(), 123);
// The sender can now be used again.
tx.send(456).await.unwrap();Если требуется несколько OwnedPermit или отправителя нельзя переместить по значению, его можно недорого клонировать перед вызовом reserve_owned:
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(1);
// Clone the sender and reserve capacity.
let permit = tx.clone().reserve_owned().await.unwrap();
// Trying to send directly on the `tx` will fail due to no
// available capacity.
assert!(tx.try_send(123).is_err());
// Sending on the permit succeeds.
permit.send(456);
// The value sent on the permit is received
assert_eq!(rx.recv().await.unwrap(), 456);pub fn try_reserve(&self) -> Result<Permit<'_, T>, TrySendError<()>>
Пытается занять место в канале, не ожидая, пока оно освободится.
Если канал заполнен, эта функция вернёт TrySendError. В противном случае, если свободное место есть, она вернёт Permit, которое позволит вам вызвать send для отправки в канал, гарантируя наличие свободного места. Эта функция похожа на reserve, но не ожидает появления свободного места.
Если удалить Permit, не отправив сообщение, место освобождается и возвращается каналу.
Примеры
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(1);
// Reserve capacity
let permit = tx.try_reserve().unwrap();
// Trying to send directly on the `tx` will fail due to no
// available capacity.
assert!(tx.try_send(123).is_err());
// Trying to reserve an additional slot on the `tx` will
// fail because there is no capacity.
assert!(tx.try_reserve().is_err());
// Sending on the permit succeeds
permit.send(456);
// The value sent on the permit is received
assert_eq!(rx.recv().await.unwrap(), 456);
pub fn try_reserve_many( &self, n: usize, ) -> Result<PermitIterator<'_, T>, TrySendError<()>>
Пытается занять n мест в канале, не ожидая, пока они освободятся.
Для отслеживания зарезервированного места возвращается PermitIterator. Вы можете вызывать этот Iterator, пока он не исчерпается, чтобы получить Permit, а затем вызвать Permit::send. Эта функция похожа на reserve_many, но не ожидает появления свободных мест.
Если в канале доступно меньше n разрешений, эта функция вернёт TrySendError::Full. Если канал закрыт, эта функция вернёт TrySendError::Closed.
Если удалить PermitIterator, не использовав его полностью, оставшиеся разрешения освобождаются и возвращаются каналу.
Примеры
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(2);
// Reserve capacity
let mut permit = tx.try_reserve_many(2).unwrap();
// Trying to send directly on the `tx` will fail due to no
// available capacity.
assert!(tx.try_send(123).is_err());
// Trying to reserve an additional slot on the `tx` will
// fail because there is no capacity.
assert!(tx.try_reserve().is_err());
// Sending with the permit iterator succeeds
permit.next().unwrap().send(456);
permit.next().unwrap().send(457);
// The iterator should now be exhausted
assert!(permit.next().is_none());
// The value sent on the permit is received
assert_eq!(rx.recv().await.unwrap(), 456);
assert_eq!(rx.recv().await.unwrap(), 457);
// Trying to call try_reserve_many with 0 will return an empty iterator
let mut permit = tx.try_reserve_many(0).unwrap();
assert!(permit.next().is_none());
// Trying to call try_reserve_many with a number greater than the channel
// capacity will return an error
let permit = tx.try_reserve_many(3);
assert!(permit.is_err());
// Trying to call try_reserve_many on a closed channel will return an error
drop(rx);
let permit = tx.try_reserve_many(1);
assert!(permit.is_err());
let permit = tx.try_reserve_many(0);
assert!(permit.is_err());pub fn try_reserve_owned(self) -> Result<OwnedPermit<T>, TrySendError<Self>>
Пытается получить место в канале, не ожидая его освобождения, и возвращает разрешение с владением.
Этот метод передаёт отправитель по значению и возвращает разрешение с владением, которое можно использовать для отправки сообщения в канал. В отличие от Sender::try_reserve, этот метод можно использовать в случаях, когда разрешение должно быть действительно в течение времени жизни 'static. Senders можно недорого клонировать (Sender::clone — это по сути увеличение счётчика ссылок, сопоставимое с Arc::clone), поэтому, если требуется несколько экземпляров OwnedPermit или Sender нельзя переместить, его можно клонировать перед вызовом try_reserve_owned.
Если канал заполнен, эта функция вернёт TrySendError. Поскольку отправитель передаётся по значению, в этом случае возвращённый TrySendError содержит отправитель, чтобы его можно было использовать повторно. Если же свободное место есть, этот метод вернёт OwnedPermit, которое затем можно использовать для вызова send в канале с гарантированно свободным местом. Эта функция похожа на reserve_owned, но не ожидает освобождения места.
Если удалить OwnedPermit, не отправив сообщение, доступная ёмкость возвращается каналу.
Примеры
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(1);
// Reserve capacity
let permit = tx.clone().try_reserve_owned().unwrap();
// Trying to send directly on the `tx` will fail due to no
// available capacity.
assert!(tx.try_send(123).is_err());
// Trying to reserve an additional slot on the `tx` will
// fail because there is no capacity.
assert!(tx.try_reserve().is_err());
// Sending on the permit succeeds
permit.send(456);
// The value sent on the permit is received
assert_eq!(rx.recv().await.unwrap(), 456);
pub fn same_channel(&self, other: &Self) -> bool
Возвращает true, если отправители принадлежат одному каналу.
Примеры
let (tx, rx) = tokio::sync::mpsc::channel::<()>(1);
let tx2 = tx.clone();
assert!(tx.same_channel(&tx2));
let (tx3, rx3) = tokio::sync::mpsc::channel::<()>(1);
assert!(!tx3.same_channel(&tx2));pub fn capacity(&self) -> usize
Возвращает текущую ёмкость канала.
Ёмкость уменьшается при отправке значения с помощью send или резервировании ёмкости с помощью reserve. Ёмкость увеличивается при получении значений через Receiver. Это значение отличается от max_capacity, которое всегда возвращает ёмкость буфера, изначально заданную при вызове channel
Примеры
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel::<()>(5);
assert_eq!(tx.capacity(), 5);
// Making a reservation drops the capacity by one.
let permit = tx.reserve().await.unwrap();
assert_eq!(tx.capacity(), 4);
// Sending and receiving a value increases the capacity by one.
permit.send(());
rx.recv().await.unwrap();
assert_eq!(tx.capacity(), 5);pub fn downgrade(&self) -> WeakSender<T>
Преобразует Sender в WeakSender, который не учитывается семантикой RAII: если удалить все экземпляры Sender канала и останутся только экземпляры WeakSender, канал будет закрыт.
pub fn max_capacity(&self) -> usize
Возвращает максимальную ёмкость буфера канала.
Максимальная ёмкость — это ёмкость буфера, изначально заданная при вызове channel. Она отличается от capacity, которое возвращает текущую доступную ёмкость буфера: по мере отправки и получения сообщений значение, возвращаемое capacity, будет увеличиваться или уменьшаться, тогда как значение, возвращаемое max_capacity, останется постоянным.
Примеры
use tokio::sync::mpsc;
let (tx, _rx) = mpsc::channel::<()>(5);
// both max capacity and capacity are the same at first
assert_eq!(tx.max_capacity(), 5);
assert_eq!(tx.capacity(), 5);
// Making a reservation doesn't change the max capacity.
let permit = tx.reserve().await.unwrap();
assert_eq!(tx.max_capacity(), 5);
// but drops the capacity by one
assert_eq!(tx.capacity(), 4);pub fn strong_count(&self) -> usize
Возвращает количество дескрипторов Sender.
pub fn weak_count(&self) -> usize
Возвращает количество дескрипторов WeakSender.
Реализации трейтов
impl<T> Clone for Sender<T>
fn clone_from(&mut self, source: &Self)
source. Подробнее
Автоматические реализации трейтов
impl<T> Freeze for Sender<T>
impl<T> RefUnwindSafe for Sender<T>
impl<T> Send for Sender<T>where T: Send,
impl<T> Sync for Sender<T>where T: Send,
impl<T> Unpin for Sender<T>
impl<T> UnsafeUnpin for Sender<T>
impl<T> UnwindSafe for Sender<T>
Общие реализации
impl<T> BorrowMut<T> for Twhere T: ?Sized,
fn borrow_mut(&mut self) -> &mut T
impl<T> CloneToUninit for Twhere T: Clone,
unsafe fn clone_to_uninit(&self, dest: *mut u8)
clone_to_uninit)
impl<T> Instrument for T
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
impl<T> ToOwned for Twhere T: Clone,
type Owned = T
fn to_owned(&self) -> T
fn clone_into(&self, target: &mut T)
impl<T> WithSubscriber for T
fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ
Subscriber к этому типу и возвращает оболочку WithDispatch. Подробнее
MIT License
Copyright © Tokio Contributors
https://docs.rs/tokio/1.53.1/tokio/sync/mpsc/struct.Sender.html