Spec-Zone.ru › Tokio

Структура Receiver

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

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

Пара из Sender и Receiver создается функцией channel.

У этого канала нет метода recv, поскольку сам получатель реализует трейт Future. Чтобы получить Result<T, error::RecvError>, .await объект Receiver напрямую.

Метод poll трейта Future может без причины возвращать Poll::Pending, даже если сообщение уже отправлено. Если происходит такой ложный сбой, вызывающий код будет разбужен после его устранения, чтобы он мог повторить попытку получения сообщения. Обратите внимание, что такое пробуждение не гарантирует успех следующего вызова — он может завершиться еще одним ложным сбоем. (Ложный сбой не означает, что сообщение потеряно. Просто его получение задерживается.)

Безопасность при отмене

Ожидание &mut Receiver<T> безопасно при отмене. Если оно используется как ветвь в tokio::select! и другая ветвь завершается первой, гарантируется, что через этот канал не было получено ни одного сообщения.

Примеры

use tokio::sync::oneshot;

let (tx, rx) = oneshot::channel();

tokio::spawn(async move {
    if let Err(_) = tx.send(3) {
        println!("the receiver dropped");
    }
});

match rx.await {
    Ok(v) => println!("got = {:?}", v),
    Err(_) => println!("the sender dropped"),
}

Если отправитель удаляется, не отправив сообщение, прием завершится ошибкой error::RecvError:

use tokio::sync::oneshot;

let (tx, rx) = oneshot::channel::<u32>();

tokio::spawn(async move {
    drop(tx);
});

match rx.await {
    Ok(_) => panic!("This doesn't happen"),
    Err(_) => println!("the sender dropped"),
}

Чтобы использовать Receiver в цикле tokio::select!, добавьте &mut перед каналом.

use tokio::sync::oneshot;
use tokio::time::{interval, sleep, Duration};

let (send, mut recv) = oneshot::channel();
let mut interval = interval(Duration::from_millis(100));

tokio::spawn(async move {
    sleep(Duration::from_secs(1)).await;
    send.send("shut down").unwrap();
});

loop {
    tokio::select! {
        _ = interval.tick() => println!("Another 100ms"),
        msg = &mut recv => {
            println!("Got message: {}", msg.unwrap());
            break;
        }
    }
}

Реализации

impl<T> Receiver<T>

pub fn close(&mut self)

Не позволяет связанному дескриптору Sender отправить значение.

Любая операция send, выполняемая после вызова close, гарантированно завершится неудачей. После вызова close следует вызвать try_recv, чтобы получить значение, если оно было отправлено до завершения вызова close.

Эта функция полезна для корректного завершения работы и гарантирует, что значение не будет отправлено в канал и так и не получено.

close ничего не делает, если сообщение уже получено или канал уже закрыт.

Примеры

Запретить отправку значения

use tokio::sync::oneshot;
use tokio::sync::oneshot::error::TryRecvError;

let (tx, mut rx) = oneshot::channel();

assert!(!tx.is_closed());

rx.close();

assert!(tx.is_closed());
assert!(tx.send("never received").is_err());

match rx.try_recv() {
    Err(TryRecvError::Closed) => {}
    _ => unreachable!(),
}

Получить значение, отправленное до вызова close

use tokio::sync::oneshot;

let (tx, mut rx) = oneshot::channel();

assert!(tx.send("will receive").is_ok());

rx.close();

let msg = rx.try_recv().unwrap();
assert_eq!(msg, "will receive");

pub fn is_terminated(&self) -> bool

Проверяет, завершена ли работа этого получателя.

Эта функция возвращает true, если этот получатель уже выдал результат Poll::Ready. В этом случае получателя больше не следует опрашивать.

Примеры

Отправка значения и его опрос.

use tokio::sync::oneshot;

use std::task::Poll;

let (tx, mut rx) = oneshot::channel();

// A receiver is not terminated when it is initialized.
assert!(!rx.is_terminated());

// A receiver is not terminated it is polled and is still pending.
let poll = futures::poll!(&mut rx);
assert_eq!(poll, Poll::Pending);
assert!(!rx.is_terminated());

// A receiver is not terminated if a value has been sent, but not yet read.
tx.send(0).unwrap();
assert!(!rx.is_terminated());

// A receiver *is* terminated after it has been polled and yielded a value.
assert_eq!((&mut rx).await, Ok(0));
assert!(rx.is_terminated());

Удаление отправителя.

use tokio::sync::oneshot;

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

// A receiver is not immediately terminated when the sender is dropped.
drop(tx);
assert!(!rx.is_terminated());

// A receiver *is* terminated after it has been polled and yielded an error.
let _ = (&mut rx).await.unwrap_err();
assert!(rx.is_terminated());

pub fn is_empty(&self) -> bool

Проверяет, пуст ли канал.

Этот метод возвращает true, если в канале нет сообщений.

Опрос пустого получателя не обязательно безопасен: он мог уже выдать значение. Вместо этого используйте is_terminated(), чтобы проверить, можно ли безопасно опрашивать получателя.

Примеры

Отправка значения.

use tokio::sync::oneshot;

let (tx, mut rx) = oneshot::channel();
assert!(rx.is_empty());

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

let _ = (&mut rx).await;
assert!(rx.is_empty());

Удаление отправителя.

use tokio::sync::oneshot;

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

// A channel is empty if the sender is dropped.
drop(tx);
assert!(rx.is_empty());

// A closed channel still yields an error, however.
(&mut rx).await.expect_err("should yield an error");
assert!(rx.is_empty());

Завершённые каналы пусты.

ⓘ
use tokio::sync::oneshot;

#[tokio::main]
async fn main() {
    let (tx, mut rx) = oneshot::channel();
    tx.send(0).unwrap();
    let _ = (&mut rx).await;

    // NB: an empty channel is not necessarily safe to poll!
    assert!(rx.is_empty());
    let _ = (&mut rx).await;
}

pub fn try_recv(&mut self) -> Result<T, TryRecvError>

Пытается получить значение.

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

Эту функцию удобно вызывать вне контекста асинхронной задачи.

Обратите внимание: в отличие от метода poll, метод try_recv не может завершиться неудачей без причины. Любое событие отправки или закрытия, произошедшее до этого вызова try_recv, будет корректно возвращено вызывающему коду.

Возвращаемое значение
  • Ok(T), если в канале ожидает значение.
  • Err(TryRecvError::Empty), если значение ещё не было отправлено.
  • Err(TryRecvError::Closed), если отправитель был удалён, не отправив значение, или если сообщение уже было получено.
Примеры

try_recv до отправки значения, а затем после неё.

use tokio::sync::oneshot;
use tokio::sync::oneshot::error::TryRecvError;

let (tx, mut rx) = oneshot::channel();

match rx.try_recv() {
    // The channel is currently empty
    Err(TryRecvError::Empty) => {}
    _ => unreachable!(),
}

// Send a value
tx.send("hello").unwrap();

match rx.try_recv() {
     Ok(value) => assert_eq!(value, "hello"),
     _ => unreachable!(),
}

try_recv, если отправитель был удалён до отправки значения

use tokio::sync::oneshot;
use tokio::sync::oneshot::error::TryRecvError;

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

drop(tx);

match rx.try_recv() {
    // The channel will never receive a value.
    Err(TryRecvError::Closed) => {}
    _ => unreachable!(),
}

pub fn blocking_recv(self) -> Result<T, RecvError>

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

Паника

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

Примеры
use std::thread;
use tokio::sync::oneshot;

#[tokio::main]
async fn main() {
    let (tx, rx) = oneshot::channel::<u8>();

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

    let _ = tx.send(10);
    sync_code.join().unwrap();
}

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

impl<T: Debug> Debug for Receiver<T>

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

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

impl<T> Drop for Receiver<T>

fn drop(&mut self)

Выполняет деструктор для этого типа. Подробнее

fn pin_drop(self: Pin<&mut Self>)

🔬Это экспериментальный API, доступный только в nightly-сборке. (pin_ergonomics)
Выполняет деструктор для этого типа, но, в отличие от Drop::drop, требует, чтобы self был закреплён. Подробнее

impl<T> Future for Receiver<T>

type Output = Result<T, RecvError>

Тип значения, возвращаемого при завершении.

fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>

Пытается разрешить будущее значение до окончательного результата; если значение ещё недоступно, регистрирует текущую задачу для пробуждения. Подробнее

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

impl<T> !RefUnwindSafe for Receiver<T>

impl<T> !UnwindSafe for Receiver<T>

impl<T> Freeze for Receiver<T>

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

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

impl<T> Unpin for Receiver<T>

impl<T> UnsafeUnpin 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<F> IntoFuture for F
where F: Future,

type Output = <F as Future>::Output

Результат, который будет получен при завершении будущего значения.

type IntoFuture = F

В какой тип будущего значения мы преобразуем это значение?

fn into_future(self) -> <F as IntoFuture>::IntoFuture

Создаёт будущее значение из значения. Подробнее

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/oneshot/struct.Receiver.html

Spec-Zone.ru

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