Структура Receiver
pub struct Receiver<T> { /* private fields */ }
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> !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> 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<F> IntoFuture for Fwhere F: Future,
type Output = <F as Future>::Output
type IntoFuture = F
fn into_future(self) -> <F as IntoFuture>::IntoFuture
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/oneshot/struct.Receiver.html