Структура Notified
pub struct Notified<'a> { /* private fields */ }
sync.Future, возвращаемый методом Notify::notified().
Этот future является слитым: после завершения любые последующие вызовы poll немедленно вернут Poll::Ready.
Реализации
impl Notified<'_>
pub fn enable(self: Pin<&mut Self>) -> bool
Добавляет этот future в список future, готовых получать пробуждения от вызовов notify_one.
Опрос future также добавляет его в список, поэтому этот метод следует использовать только в том случае, если вы хотите добавить future в список до первого вызова poll. (Фактически этот метод эквивалентен вызову poll, за исключением того, что Waker не регистрируется.)
Это не влияет на уведомления, отправленные с помощью notify_waiters: они получаются, если возникают после создания Notified, независимо от того, вызывались ли enable или poll.
Метод возвращает true, если Notified готов. Это происходит в следующих случаях:
- Метод
notify_waitersбыл вызван между созданиемNotifiedи вызовом этого метода. - Это первый вызов
enableилиpollдля этого future, иNotifyудерживал разрешение, полученное при предыдущем вызовеnotify_one. В этом случае вызов использует разрешение. - Ранее future был активирован или опрошен, а затем помечен как готовый либо в результате использования разрешения из
Notify, либо вызовомnotify_oneилиnotify_waiters, который удалил его из списка future, готовых получать пробуждения.
Если этот метод возвращает true, любые последующие вызовы poll для того же future немедленно вернут Poll::Ready.
Примеры
Неограниченный канал «многие производители — многие потребители» (mpmc).
Вызов enable важен, поскольку в противном случае, если параллельно выполняются два вызова recv и два вызова send, может произойти следующее:
- Оба вызова
try_recvвозвращаютNone. - Оба новых элемента добавляются в вектор.
- Метод
notify_oneвызывается дважды, но вNotifyдобавляется только одно разрешение. - Оба вызова
recvдоходят до futureNotified. Один из них использует разрешение, а второй засыпает навсегда.
Добавив future Notified в список вызовом enable до try_recv, вызовы notify_one на третьем шаге удалили бы future из списка и пометили их как уведомлённые, вместо того чтобы добавлять разрешение в Notify. Это гарантирует, что оба future будут разбужены.
use tokio::sync::Notify;
use std::collections::VecDeque;
use std::sync::Mutex;
struct Channel<T> {
messages: Mutex<VecDeque<T>>,
notify_on_sent: Notify,
}
impl<T> Channel<T> {
pub fn send(&self, msg: T) {
let mut locked_queue = self.messages.lock().unwrap();
locked_queue.push_back(msg);
drop(locked_queue);
// Send a notification to one of the calls currently
// waiting in a call to `recv`.
self.notify_on_sent.notify_one();
}
pub fn try_recv(&self) -> Option<T> {
let mut locked_queue = self.messages.lock().unwrap();
locked_queue.pop_front()
}
pub async fn recv(&self) -> T {
let future = self.notify_on_sent.notified();
tokio::pin!(future);
loop {
// Make sure that no wakeup is lost if we get
// `None` from `try_recv`.
future.as_mut().enable();
if let Some(msg) = self.try_recv() {
return msg;
}
// Wait for a call to `notify_one`.
//
// This uses `.as_mut()` to avoid consuming the future,
// which lets us call `Pin::set` below.
future.as_mut().await;
// Reset the future in case another call to
// `try_recv` got the message before us.
future.set(self.notify_on_sent.notified());
}
}
}Реализации трейтов
impl<'a> Send for Notified<'a>
impl<'a> Sync for Notified<'a>
Автоматические реализации трейтов
impl<'a> !Freeze for Notified<'a>
impl<'a> !RefUnwindSafe for Notified<'a>
impl<'a> !Unpin for Notified<'a>
impl<'a> !UnsafeUnpin for Notified<'a>
impl<'a> !UnwindSafe for Notified<'a>
Универсальные реализации
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/futures/struct.Notified.html