Spec-Zone.ru › Tokio

Структура Receiver

pub struct Receiver<T> { /* private fields */ }
Доступно только при включённой функции crate sync.

Получает значения из связанного Sender.

Экземпляры создаются функцией channel.

Этот получатель можно преобразовать в Stream с помощью ReceiverStream.

Реализации

impl<T> Receiver<T>

pub async fn recv(&mut self) -> Option<T>

Получает следующее значение для этого получателя.

Этот метод возвращает None, если канал закрыт и в его буфере больше нет сообщений. Это означает, что из этого Receiver больше нельзя получить значения. Канал закрывается, когда все отправители уничтожены или вызывается close.

Если в буфере канала нет сообщений, но канал ещё не закрыт, этот метод будет ожидать, пока не будет отправлено сообщение или канал не закроется. Обратите внимание: если вызывается close, но всё ещё существуют Permits, выданные до закрытия, recv не считает канал закрытым, пока эти разрешения не будут освобождены.

Безопасность отмены

Этот метод безопасен при отмене. Если recv используется как ветвь в tokio::select! и первой завершается другая ветвь, гарантируется, что из этого канала не было получено ни одного сообщения.

Примеры
use tokio::sync::mpsc;

let (tx, mut rx) = mpsc::channel(100);

tokio::spawn(async move {
    tx.send("hello").await.unwrap();
});

assert_eq!(Some("hello"), rx.recv().await);
assert_eq!(None, rx.recv().await);

Значения буферизуются:

use tokio::sync::mpsc;

let (tx, mut rx) = mpsc::channel(100);

tx.send("hello").await.unwrap();
tx.send("world").await.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, если в очереди канала нет сообщений, но канал ещё не закрыт, этот метод будет ожидать, пока не будет отправлено сообщение или канал не закроется. Обратите внимание: если вызывается close, но всё ещё существуют Permits, выданные до закрытия, recv_many не считает канал закрытым, пока эти разрешения не будут освобождены.

При ненулевых значениях 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::channel(100);
let tx2 = tx.clone();
tx2.send("first").await.unwrap();
tx2.send("second").await.unwrap();
tx2.send("third").await.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, 1).await);

tokio::spawn(async move {
    tx.send("fourth").await.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, 1).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::channel(100);

tx.send("hello").await.unwrap();

assert_eq!(Ok("hello"), rx.try_recv());
assert_eq!(Err(TryRecvError::Empty), rx.try_recv());

tx.send("hello").await.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>

Блокирующий приём для вызова вне асинхронного контекста.

Этот метод возвращает None, если канал закрыт и в его буфере не осталось сообщений. Это означает, что из этого Receiver больше никогда нельзя будет получить значения. Канал закрывается, когда все отправители удалены или вызывается close.

Если в буфере канала нет сообщений, но канал ещё не закрыт, этот метод будет блокировать выполнение до отправки сообщения или закрытия канала.

Этот метод предназначен для случаев, когда данные отправляются из асинхронного кода в синхронный, и он будет работать, даже если отправитель не использует blocking_send для отправки сообщения.

Обратите внимание: если вызывается close, но остаются активные Permits, выданные до закрытия, канал не считается закрытым для blocking_recv, пока разрешения не будут освобождены.

Паника

Эта функция вызывает панику, если её вызвать в асинхронном контексте выполнения.

Примеры
use std::thread;
use tokio::runtime::Runtime;
use tokio::sync::mpsc;

fn main() {
    let (tx, mut rx) = mpsc::channel::<u8>(10);

    let sync_code = thread::spawn(move || {
        assert_eq!(Some(10), rx.blocking_recv());
    });

    Runtime::new()
        .unwrap()
        .block_on(async move {
            let _ = tx.send(10).await;
        });
    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)

Закрывает принимающую половину канала, не удаляя её.

Это предотвращает отправку дальнейших сообщений по каналу, но при этом позволяет получателю извлечь уже буферизованные сообщения. Все существующие значения Permit по-прежнему смогут отправлять сообщения.

Чтобы гарантировать, что сообщения не будут потеряны, после вызова close() необходимо вызывать recv() до тех пор, пока не будет возвращено None. Если существуют значения Permit или OwnedPermit, метод recv не вернёт None, пока они не будут освобождены.

Примеры
use tokio::sync::mpsc;

let (tx, mut rx) = mpsc::channel(20);

tokio::spawn(async move {
    let mut i = 0;
    while let Ok(permit) = tx.reserve().await {
        permit.send(i);
        i += 1;
    }
});

rx.close();

while let Some(msg) = rx.recv().await {
    println!("got {}", msg);
}

// Channel closed and no messages are lost.

pub fn is_closed(&self) -> bool

Проверяет, закрыт ли канал.

Этот метод возвращает true, если канал закрыт. Канал закрывается, когда все Sender удалены или вызывается Receiver::close.

Примеры
use tokio::sync::mpsc;

let (_tx, mut rx) = mpsc::channel::<()>(10);
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::channel(10);
assert!(rx.is_empty());

tx.send(0).await.unwrap();
assert!(!rx.is_empty());

pub fn len(&self) -> usize

Возвращает количество сообщений в канале.

Примеры
use tokio::sync::mpsc;

let (tx, rx) = mpsc::channel(10);
assert_eq!(0, rx.len());

tx.send(0).await.unwrap();
assert_eq!(1, rx.len());

pub fn capacity(&self) -> usize

Возвращает текущую ёмкость канала.

Ёмкость уменьшается, когда отправитель отправляет значение с помощью Sender::send или резервирует ёмкость с помощью Sender::reserve. Ёмкость увеличивается при получении значений. Это отличается от max_capacity, который всегда возвращает исходную ёмкость буфера, заданную при вызове channel.

Примеры
use tokio::sync::mpsc;

let (tx, mut rx) = mpsc::channel::<()>(5);

assert_eq!(rx.capacity(), 5);

// Making a reservation drops the capacity by one.
let permit = tx.reserve().await.unwrap();
assert_eq!(rx.capacity(), 4);
assert_eq!(rx.len(), 0);

// Sending and receiving a value increases the capacity by one.
permit.send(());
assert_eq!(rx.len(), 1);
rx.recv().await.unwrap();
assert_eq!(rx.capacity(), 5);

// Directly sending a message drops the capacity by one.
tx.send(()).await.unwrap();
assert_eq!(rx.capacity(), 4);
assert_eq!(rx.len(), 1);

// Receiving the message increases the capacity by one.
rx.recv().await.unwrap();
assert_eq!(rx.capacity(), 5);
assert_eq!(rx.len(), 0);

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!(rx.max_capacity(), 5);
assert_eq!(rx.capacity(), 5);

// Making a reservation doesn't change the max capacity.
let permit = tx.reserve().await.unwrap();
assert_eq!(rx.max_capacity(), 5);
// but drops the capacity by one
assert_eq!(rx.capacity(), 4);

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::Receiver<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::channel(32);
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).await.unwrap();
}

let count = my_receiver_future.await;
assert_eq!(count, 3);
assert_eq!(buffer, vec![0,1,2])

pub fn sender_strong_count(&self) -> usize

Возвращает количество дескрипторов Sender.

pub fn sender_weak_count(&self) -> usize

Возвращает количество дескрипторов WeakSender.

Реализации трейтов

impl<T> Debug for Receiver<T>

fn fmt(&self, fmt: &mut Formatter<'_>) -> Result

Форматирует значение с помощью указанного форматтера. Подробнее

impl<T> Unpin for Receiver<T>

Автоматические реализации трейтов

impl<T> Freeze for Receiver<T>

impl<T> RefUnwindSafe for Receiver<T>

impl<T> Send for Receiver<T>
where T: Send,

impl<T> Sync for Receiver<T>
where T: Send,

impl<T> UnsafeUnpin for Receiver<T>

impl<T> UnwindSafe for Receiver<T>

Обобщённые реализации

impl<T> Any for T
where T: 'static + ?Sized,

fn type_id(&self) -> TypeId

Получает TypeId для self. Подробнее

impl<T> Borrow<T> for T
where T: ?Sized,

fn borrow(&self) -> &T

Неизменяемо заимствует данные из принадлежащего ему значения. Подробнее

impl<T> BorrowMut<T> for T
where T: ?Sized,

fn borrow_mut(&mut self) -> &mut T

Изменяемо заимствует данные из принадлежащего ему значения. Подробнее

impl<T> From<T> for T

fn from(t: T) -> T

Возвращает аргумент без изменений.

impl<T> Instrument for T

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Инструментирует этот тип с помощью указанного Span, возвращая оболочку Instrumented. Подробнее

fn in_current_span(self) -> Instrumented<Self> ⓘ

Инструментирует этот тип с помощью текущего Span, возвращая оболочку Instrumented. Подробнее

impl<T, U> Into<U> for T
where U: From<T>,

fn into(self) -> U

Вызывает U::from(self).

То есть это преобразование выполняет действия, определяемые реализацией From<T> for U.

impl<T, U> TryFrom<U> for T
where U: Into<T>,

type Error = Infallible

Тип, возвращаемый в случае ошибки преобразования.

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Выполняет преобразование.

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

type Error = <U as TryFrom<T>>::Error

Тип, возвращаемый в случае ошибки преобразования.

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Выполняет преобразование.

impl<T> WithSubscriber for T

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Присоединяет предоставленный Subscriber к этому типу и возвращает обёртку WithDispatch. Подробнее

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.Receiver.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API