Структура Unbounded Receiver
pub struct UnboundedReceiver<T> { /* private fields */ }
sync.Получение значений из связанного UnboundedSender.
Экземпляры создаются функцией unbounded_channel.
Этот получатель можно преобразовать в Stream с помощью UnboundedReceiverStream.
Реализации
impl<T> UnboundedReceiver<T>
pub async fn recv(&mut self) -> Option<T>
Получает следующее значение для этого получателя.
Этот метод возвращает None, если канал закрыт и в его буфере не осталось сообщений. Это означает, что из этого Receiver больше нельзя будет получить значения. Канал закрывается, когда все отправители уничтожены или вызывается close.
Если в буфере канала нет сообщений, но канал ещё не закрыт, этот метод будет ждать, пока не будет отправлено сообщение или канал не закроется.
Безопасность отмены
Этот метод безопасен при отмене. Если recv используется как ветвь в tokio::select!, а другая ветвь завершается первой, гарантируется, что из этого канала не было получено ни одного сообщения.
Примеры
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::unbounded_channel();
tokio::spawn(async move {
tx.send("hello").unwrap();
});
assert_eq!(Some("hello"), rx.recv().await);
assert_eq!(None, rx.recv().await);Значения буферизуются:
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::unbounded_channel();
tx.send("hello").unwrap();
tx.send("world").unwrap();
assert_eq!(Some("hello"), rx.recv().await);
assert_eq!(Some("world"), rx.recv().await);pub async fn recv_many(&mut self, buffer: &mut Vec<T>, limit: usize) -> usize
Получает следующие значения для этого получателя и дополняет buffer.
Этот метод дополняет buffer не более чем фиксированным числом значений, заданным limit. Если limit равен нулю, функция немедленно возвращает 0. Возвращаемое значение — это количество значений, добавленных в buffer.
Для limit > 0, если в очереди канала нет сообщений, но канал ещё не закрыт, этот метод будет ждать, пока не будет отправлено сообщение или канал не закроется.
При ненулевых значениях limit этот метод никогда не вернёт 0, если только канал не закрыт и в его очереди не осталось сообщений. Это означает, что из этого Receiver больше нельзя будет получить значения. Канал закрывается, когда все отправители уничтожены или вызывается close.
Ёмкость buffer при необходимости увеличивается.
Безопасность отмены
Этот метод безопасен при отмене. Если recv_many используется как ветвь в tokio::select!, а другая ветвь завершается первой, гарантируется, что из этого канала не было получено ни одного сообщения.
Примеры
use tokio::sync::mpsc;
let mut buffer: Vec<&str> = Vec::with_capacity(2);
let limit = 2;
let (tx, mut rx) = mpsc::unbounded_channel();
let tx2 = tx.clone();
tx2.send("first").unwrap();
tx2.send("second").unwrap();
tx2.send("third").unwrap();
// Call `recv_many` to receive up to `limit` (2) values.
assert_eq!(2, rx.recv_many(&mut buffer, limit).await);
assert_eq!(vec!["first", "second"], buffer);
// If the buffer is full, the next call to `recv_many`
// reserves additional capacity.
assert_eq!(1, rx.recv_many(&mut buffer, limit).await);
tokio::spawn(async move {
tx.send("fourth").unwrap();
});
// 'tx' is dropped, but `recv_many`
// is guaranteed not to return 0 as the channel
// is not yet closed.
assert_eq!(1, rx.recv_many(&mut buffer, limit).await);
assert_eq!(vec!["first", "second", "third", "fourth"], buffer);
// Once the last sender is dropped, the channel is
// closed and `recv_many` returns 0, capacity unchanged.
drop(tx2);
assert_eq!(0, rx.recv_many(&mut buffer, limit).await);
assert_eq!(vec!["first", "second", "third", "fourth"], buffer);pub fn try_recv(&mut self) -> Result<T, TryRecvError>
Пытается получить следующее значение для этого получателя.
Этот метод возвращает ошибку Empty, если канал в данный момент пуст, но существуют активные отправители или разрешения.
Этот метод возвращает ошибку Disconnected, если канал в данный момент пуст и нет активных отправителей или разрешений.
В отличие от метода poll_recv, этот метод никогда не возвращает ошибку Empty без причины.
Примеры
use tokio::sync::mpsc;
use tokio::sync::mpsc::error::TryRecvError;
let (tx, mut rx) = mpsc::unbounded_channel();
tx.send("hello").unwrap();
assert_eq!(Ok("hello"), rx.try_recv());
assert_eq!(Err(TryRecvError::Empty), rx.try_recv());
tx.send("hello").unwrap();
// Drop the last sender, closing the channel.
drop(tx);
assert_eq!(Ok("hello"), rx.try_recv());
assert_eq!(Err(TryRecvError::Disconnected), rx.try_recv());pub fn blocking_recv(&mut self) -> Option<T>
Блокирующее получение для вызова вне асинхронного контекста.
Паника
Эта функция вызывает панику, если её вызвать в асинхронном контексте выполнения.
Примеры
use std::thread;
use tokio::sync::mpsc;
#[tokio::main]
async fn main() {
let (tx, mut rx) = mpsc::unbounded_channel::<u8>();
let sync_code = thread::spawn(move || {
assert_eq!(Some(10), rx.blocking_recv());
});
let _ = tx.send(10);
sync_code.join().unwrap();
}pub fn blocking_recv_many(&mut self, buffer: &mut Vec<T>, limit: usize) -> usize
Вариант Self::recv_many для блокирующих контекстов.
Применяются те же условия, что и для Self::blocking_recv.
pub fn close(&mut self)
Закрывает принимающую половину канала, не уничтожая её.
Это предотвращает отправку дальнейших сообщений по каналу, позволяя при этом получателю извлечь буферизованные сообщения.
Чтобы гарантировать, что ни одно сообщение не будет потеряно, после вызова close() необходимо вызывать recv() до тех пор, пока не будет возвращено None.
pub fn is_closed(&self) -> bool
Проверяет, закрыт ли канал.
Этот метод возвращает true, если канал был закрыт. Канал закрывается, когда все дескрипторы UnboundedSender уничтожены или когда вызывается UnboundedReceiver::close.
Примеры
use tokio::sync::mpsc;
let (_tx, mut rx) = mpsc::unbounded_channel::<()>();
assert!(!rx.is_closed());
rx.close();
assert!(rx.is_closed());pub fn is_empty(&self) -> bool
Проверяет, пуст ли канал.
Этот метод возвращает true, если в канале нет сообщений.
Примеры
use tokio::sync::mpsc;
let (tx, rx) = mpsc::unbounded_channel();
assert!(rx.is_empty());
tx.send(0).unwrap();
assert!(!rx.is_empty());
pub fn len(&self) -> usize
Возвращает количество сообщений в канале.
Примеры
use tokio::sync::mpsc;
let (tx, rx) = mpsc::unbounded_channel();
assert_eq!(0, rx.len());
tx.send(0).unwrap();
assert_eq!(1, rx.len());pub fn poll_recv(&mut self, cx: &mut Context<'_>) -> Poll<Option<T>>
Ожидает получения следующего сообщения из этого канала.
Этот метод возвращает:
-
Poll::Pending, если сообщения недоступны, но канал не закрыт, или произошёл ложный сбой. -
Poll::Ready(Some(message)), если доступно сообщение. -
Poll::Ready(None), если канал закрыт и получены все сообщения, отправленные до его закрытия.
Когда метод возвращает Poll::Pending, для Waker в переданном Context планируется пробуждение при отправке сообщения любому получателю или при закрытии канала. Обратите внимание: при нескольких вызовах poll_recv или poll_recv_many пробуждение планируется только для Waker из Context, переданного в последнем вызове.
Если этот метод возвращает Poll::Pending из-за ложного сбоя, Waker получит уведомление, когда ситуация, вызвавшая ложный сбой, будет устранена. Обратите внимание, что такое пробуждение не гарантирует успешного результата следующего вызова — он может завершиться ещё одним ложным сбоем.
pub fn poll_recv_many( &mut self, cx: &mut Context<'_>, buffer: &mut Vec<T>, limit: usize, ) -> Poll<usize>
Ожидает получения нескольких сообщений из этого канала, добавляя их в переданный буфер.
Этот метод возвращает:
-
Poll::Pending, если сообщения недоступны, но канал не закрыт, или произошёл ложный сбой. -
Poll::Ready(count), гдеcount— количество успешно полученных сообщений, сохранённых вbuffer. Это значение может быть меньше или равноlimit. -
Poll::Ready(0), еслиlimitустановлено в ноль или канал закрыт.
Когда метод возвращает Poll::Pending, для Waker в переданном Context планируется пробуждение при отправке сообщения любому получателю или при закрытии канала. Обратите внимание: при нескольких вызовах poll_recv или poll_recv_many пробуждение планируется только для Waker из Context, переданного в последнем вызове.
Обратите внимание, что этот метод не гарантирует получение ровно limit сообщений. Вместо этого, если доступно хотя бы одно сообщение, он возвращает столько сообщений, сколько возможно, но не больше указанного предела. Этот метод возвращает ноль только в том случае, если канал закрыт (или если limit равно нулю).
Примеры
use std::task::{Context, Poll};
use std::pin::Pin;
use tokio::sync::mpsc;
use futures::Future;
struct MyReceiverFuture<'a> {
receiver: mpsc::UnboundedReceiver<i32>,
buffer: &'a mut Vec<i32>,
limit: usize,
}
impl<'a> Future for MyReceiverFuture<'a> {
type Output = usize; // Number of messages received
fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
let MyReceiverFuture { receiver, buffer, limit } = &mut *self;
// Now `receiver` and `buffer` are mutable references, and `limit` is copied
match receiver.poll_recv_many(cx, *buffer, *limit) {
Poll::Pending => Poll::Pending,
Poll::Ready(count) => Poll::Ready(count),
}
}
}
let (tx, rx) = mpsc::unbounded_channel::<i32>();
let mut buffer = Vec::new();
let my_receiver_future = MyReceiverFuture {
receiver: rx,
buffer: &mut buffer,
limit: 3,
};
for i in 0..10 {
tx.send(i).expect("Unable to send integer");
}
let count = my_receiver_future.await;
assert_eq!(count, 3);
assert_eq!(buffer, vec![0,1,2])pub fn sender_strong_count(&self) -> usize
Возвращает количество дескрипторов UnboundedSender.
pub fn sender_weak_count(&self) -> usize
Возвращает количество дескрипторов WeakUnboundedSender.
Реализации трейтов
Автоматические реализации трейтов
impl<T> Freeze for UnboundedReceiver<T>
impl<T> RefUnwindSafe for UnboundedReceiver<T>
impl<T> Send for UnboundedReceiver<T>where T: Send,
impl<T> Sync for UnboundedReceiver<T>where T: Send,
impl<T> Unpin for UnboundedReceiver<T>
impl<T> UnsafeUnpin for UnboundedReceiver<T>
impl<T> UnwindSafe for UnboundedReceiver<T>
Универсальные реализации
impl<T> BorrowMut<T> for Twhere T: ?Sized,
fn borrow_mut(&mut self) -> &mut T
impl<T> Instrument for T
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
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.UnboundedReceiver.html