Структура Receiver
pub struct Receiver<T> { /* private fields */ }
sync.Получает значения от связанного Sender.
Экземпляры создаются функцией channel.
Чтобы преобразовать этот receiver в Stream, можно использовать обёртку WatchStream.
Реализации
impl<T> Receiver<T>
pub fn borrow(&self) -> Ref<'_, T>
Возвращает ссылку на последнее отправленное значение.
Этот метод не помечает возвращённое значение как просмотренное, поэтому последующие вызовы changed могут немедленно вернуть результат, даже если вы уже просмотрели значение вызовом borrow.
Активные заимствования удерживают блокировку для чтения внутреннего значения. Это означает, что длительные заимствования могут привести к блокировке стороны-отправителя. Рекомендуется делать заимствование как можно короче. Кроме того, если вы работаете в среде, допускающей !Send futures, необходимо убедиться, что возвращённый тип Ref не удерживается активным через точку .await, иначе это может привести к взаимной блокировке.
Политика приоритетов блокировки зависит от базовой реализации блокировки, и этот тип не гарантирует применение какой-либо определённой политики. В частности, ожидающий получения блокировки в send отправитель может как блокировать, так и не блокировать параллельные вызовы borrow, например:
// Task 1 (on thread A) | // Task 2 (on thread B)
let _ref1 = rx.borrow(); |
| // will block
| let _ = tx.send(());
// may deadlock |
let _ref2 = rx.borrow(); |Подробнее о том, когда использовать этот метод вместо borrow_and_update, см. здесь.
Примеры
use tokio::sync::watch;
let (_, rx) = watch::channel("hello");
assert_eq!(*rx.borrow(), "hello");pub fn borrow_and_update(&mut self) -> Ref<'_, T>
Возвращает ссылку на последнее отправленное значение и помечает его как просмотренное.
Этот метод помечает текущее значение как просмотренное. Последующие вызовы changed не будут немедленно возвращать результат, пока Sender снова не изменит общее значение.
Активные заимствования удерживают блокировку для чтения внутреннего значения. Это означает, что длительные заимствования могут привести к блокировке стороны-отправителя. Рекомендуется делать заимствование как можно короче. Кроме того, если вы работаете в среде, допускающей !Send futures, необходимо убедиться, что возвращённый тип Ref не удерживается активным через точку .await, иначе это может привести к взаимной блокировке.
Политика приоритетов блокировки зависит от базовой реализации блокировки, и этот тип не гарантирует применение какой-либо определённой политики. В частности, ожидающий получения блокировки в send отправитель может как блокировать, так и не блокировать параллельные вызовы borrow, например:
// Task 1 (on thread A) | // Task 2 (on thread B)
let _ref1 = rx1.borrow_and_update(); |
| // will block
| let _ = tx.send(());
// may deadlock |
let _ref2 = rx2.borrow_and_update(); |Подробнее о том, когда использовать этот метод вместо borrow, см. здесь.
pub fn has_changed(&self) -> Result<bool, RecvError>
Проверяет, содержит ли этот канал сообщение, которое receiver ещё не просматривал. Текущее значение не будет помечено как просмотренное.
Несмотря на то что этот метод называется has_changed, он не сравнивает сообщения на равенство, поэтому вызов вернёт true, даже если текущее сообщение равно предыдущему.
Ошибки
Возвращает RecvError тогда и только тогда, когда канал закрыт.
Примеры
Базовое использование
use tokio::sync::watch;
let (tx, mut rx) = watch::channel("hello");
tx.send("goodbye").unwrap();
assert!(rx.has_changed().unwrap());
assert_eq!(*rx.borrow_and_update(), "goodbye");
// The value has been marked as seen
assert!(!rx.has_changed().unwrap());Пример закрытого канала
use tokio::sync::watch;
let (tx, rx) = watch::channel("hello");
tx.send("goodbye").unwrap();
drop(tx);
// The channel is closed
assert!(rx.has_changed().is_err());pub fn mark_changed(&mut self)
Помечает состояние как изменённое.
После вызова этого метода has_changed() возвращает true, а changed() возвращает результат немедленно, независимо от того, было ли отправлено новое значение.
Это полезно, чтобы инициировать первоначальное уведомление об изменении после подписки и синхронизировать новые receiver.
pub fn mark_unchanged(&mut self)
Помечает состояние как неизменённое.
Receiver будет считать текущее значение просмотренным.
Это полезно, если вас не интересует текущее значение, доступное в receiver.
pub async fn changed(&mut self) -> Result<(), RecvError>
Ожидает уведомления об изменении, а затем помечает текущее значение как замеченное.
Если текущее значение в канале ещё не было помечено как замеченное на момент вызова этого метода, метод помечает его как замеченное и немедленно возвращает управление. Если последнее значение уже было помечено как замеченное, метод приостанавливается до тех пор, пока новое сообщение не будет отправлено через Sender, подключённый к этому Receiver, или пока все Sender не будут удалены.
Дополнительные сведения см. в разделе Уведомления об изменениях документации модуля.
Ошибки
Возвращает RecvError, если канал закрыт И текущее значение помечено как замеченное.
Безопасность отмены
Этот метод безопасен при отмене. Если использовать его в качестве ветви в tokio::select!, а другая ветвь завершится первой, гарантируется, что этот вызов changed не пометил ни одного значения как замеченное.
Примеры
use tokio::sync::watch;
let (tx, mut rx) = watch::channel("hello");
tokio::spawn(async move {
tx.send("goodbye").unwrap();
});
assert!(rx.changed().await.is_ok());
assert_eq!(*rx.borrow_and_update(), "goodbye");
// The `tx` handle has been dropped
assert!(rx.changed().await.is_err());pub async fn wait_for( &mut self, f: impl FnMut(&T) -> bool, ) -> Result<Ref<'_, T>, RecvError>
Ожидает значение, удовлетворяющее заданному условию.
Этот метод вызывает переданное замыкание при каждой отправке чего-либо в канал. Как только замыкание возвращает true, метод возвращает ссылку на значение, переданное замыканию.
Перед тем как wait_for начнёт ожидать изменений, он вызовет замыкание для текущего значения. Если замыкание возвращает true для текущего значения, wait_for немедленно вернёт ссылку на это значение. Это произойдёт, даже если текущее значение уже считается замеченным.
Канал watch хранит только последнее значение, поэтому, если несколько сообщений отправляются быстрее, чем wait_for успевает вызывать замыкание, некоторые обновления могут быть пропущены. При каждом вызове замыкания ему передаётся последнее значение.
Когда эта функция возвращает результат, значение, переданное замыканию при возврате им true, считается замеченным.
Если канал закрыт, wait_for вернёт RecvError. После этого отправить в канал новые сообщения уже невозможно. При возврате ошибки гарантируется, что замыкание было вызвано для последнего значения и вернуло для него false. (Если бы замыкание вернуло true, вместо ошибки было бы возвращено последнее значение.)
Как и метод borrow, возвращённое заимствование удерживает блокировку для чтения внутреннего значения. Это означает, что длительное заимствование может заблокировать сторону-производителя. Рекомендуется как можно быстрее завершать заимствование. Дополнительные сведения см. в документации к borrow.
Безопасность отмены
Этот метод безопасен при отмене. Если использовать его в качестве ветви в tokio::select!, а другая ветвь завершится первой, гарантируется, что последнее замеченное значение val (если оно есть) удовлетворяет условию f(val) == false.
Паники
Паника возникает тогда и только тогда, когда замыкание f. В этом случае ни один ресурс, принадлежащий этому Receiver или совместно используемый им, не будет отравлен.
Примеры
use tokio::sync::watch;
use tokio::time::{sleep, Duration};
#[tokio::main(flavor = "current_thread", start_paused = true)]
async fn main() {
let (tx, mut rx) = watch::channel("hello");
tokio::spawn(async move {
sleep(Duration::from_secs(1)).await;
tx.send("goodbye").unwrap();
});
assert!(rx.wait_for(|val| *val == "goodbye").await.is_ok());
assert_eq!(*rx.borrow(), "goodbye");
}pub fn same_channel(&self, other: &Self) -> bool
Возвращает true, если получатели принадлежат одному каналу.
Примеры
let (tx, rx) = tokio::sync::watch::channel(true);
let rx2 = rx.clone();
assert!(rx.same_channel(&rx2));
let (tx3, rx3) = tokio::sync::watch::channel(true);
assert!(!rx3.same_channel(&rx2));Реализации трейтов
impl<T> Clone for Receiver<T>
fn clone_from(&mut self, source: &Self)
source. Подробнее
Автоматические реализации трейтов
impl<T> !RefUnwindSafe for Receiver<T>
impl<T> !UnwindSafe for Receiver<T>
impl<T> Freeze for Receiver<T>
impl<T> Send for Receiver<T>
impl<T> Sync for Receiver<T>
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> 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/watch/struct.Receiver.html