Структура Notify
pub struct Notify { /* private fields */ }
sync.Уведомляет одну задачу о том, что ей нужно пробудиться.
Notify предоставляет базовый механизм уведомления одной задачи о событии. Сам Notify не содержит никаких данных. Вместо этого он используется для передачи сигнала другой задаче, чтобы та выполнила операцию.
Notify можно представить как Semaphore, изначально имеющий 0 разрешений. Метод notified().await ожидает, пока разрешение станет доступным, а notify_one() устанавливает разрешение, если в данный момент нет доступных разрешений.
Детали синхронизации Notify аналогичны thread::park и Thread::unpark из std. Значение Notify содержит одно разрешение. notified().await ожидает, пока разрешение станет доступным, использует его и продолжает выполнение. notify_one() устанавливает разрешение и пробуждает ожидающую задачу, если она есть.
Если notify_one() вызван до notified().await, следующий вызов notified().await немедленно завершится, использовав разрешение. Все последующие вызовы notified().await будут ожидать нового разрешения.
Если notify_one() вызван несколько раз до notified().await, сохраняется только одно разрешение. Следующий вызов notified().await немедленно завершится, а следующий за ним будет ожидать нового разрешения.
Примеры
Базовое использование.
use tokio::sync::Notify;
use std::sync::Arc;
let notify = Arc::new(Notify::new());
let notify2 = notify.clone();
let handle = tokio::spawn(async move {
notify2.notified().await;
println!("received notification");
});
println!("sending notification");
notify.notify_one();
// Wait for task to receive notification.
handle.await.unwrap();Неограниченный канал «многие отправители — один получатель» (mpsc).
При использовании этого канала пробуждения не могут быть потеряны, поскольку вызов notify_one() сохранит разрешение в Notify, которое будет использовано следующим вызовом notified().
use tokio::sync::Notify;
use std::collections::VecDeque;
use std::sync::Mutex;
struct Channel<T> {
values: Mutex<VecDeque<T>>,
notify: Notify,
}
impl<T> Channel<T> {
pub fn send(&self, value: T) {
self.values.lock().unwrap()
.push_back(value);
// Notify the consumer a value is available
self.notify.notify_one();
}
// This is a single-consumer channel, so several concurrent calls to
// `recv` are not allowed.
pub async fn recv(&self) -> T {
loop {
// Drain values
if let Some(value) = self.values.lock().unwrap().pop_front() {
return value;
}
// Wait for values to be available
self.notify.notified().await;
}
}
}Неограниченный канал «многие отправители — многие получатели» (mpmc).
Вызов enable важен, поскольку без него при наличии двух параллельных вызовов recv и двух вызовов send может произойти следующее:
- Оба вызова
try_recvвозвращаютNone. - Оба новых элемента добавляются в вектор.
- Метод
notify_oneвызывается дважды, но вNotifyдобавляется только одно разрешение. - Оба вызова
recvдостигают futureNotified. Один из них использует разрешение, а другой засыпает навсегда.
Добавив futures Notified в список вызовом enable до try_recv, вызовы notify_one на третьем шаге удалят futures из списка и отметят их как уведомлённые, вместо того чтобы добавлять разрешение в Notify. Это гарантирует пробуждение обоих futures.
Обратите внимание: такая ошибка может произойти только при наличии двух одновременных вызовов recv. Именно поэтому в приведённом выше примере с mpsc вызов enable не требуется.
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 Notify
pub fn new() -> Notify
Создает новый Notify, инициализированный без разрешения.
Примеры
use tokio::sync::Notify;
let notify = Notify::new();pub const fn const_new() -> Notify
Создает новый Notify, инициализированный без разрешения.
При использовании нестабильной возможности tracing, Notify, созданный с помощью const_new, не будет инструментирован. Поэтому он не будет отображаться в tokio-console. Вместо этого, если это необходимо, для создания инструментированного объекта следует использовать Notify::new.
Примеры
use tokio::sync::Notify;
static NOTIFY: Notify = Notify::const_new();pub fn notified(&self) -> Notified<'_> ⓘ
Ожидает уведомления.
Эквивалентно:
async fn notified(&self);
Каждое значение Notify содержит одно разрешение. Если разрешение доступно после предыдущего вызова notify_one(), то notified().await завершится немедленно, использовав это разрешение. В противном случае notified().await будет ожидать, пока следующий вызов notify_one() не предоставит разрешение.
Гарантируется, что будущее Notified не получит сигналов пробуждения от вызовов notify_one(), если оно еще не было опрошено. Подробнее см. документацию для Notified::enable().
Гарантируется, что будущее Notified получит сигналы пробуждения от notify_waiters() сразу после создания, даже если оно еще не было опрошено.
Безопасность отмены
Этот метод использует очередь для справедливого распределения уведомлений в порядке их запроса. Отмена вызова notified приводит к потере места в очереди.
Примеры
use tokio::sync::Notify;
use std::sync::Arc;
let notify = Arc::new(Notify::new());
let notify2 = notify.clone();
tokio::spawn(async move {
notify2.notified().await;
println!("received notification");
});
println!("sending notification");
notify.notify_one();pub fn notified_owned(self: Arc<Self>) -> OwnedNotified ⓘ
Ожидает уведомления, владея Future.
В отличие от Self::notified, который возвращает будущее, связанное со временем жизни Notify, notified_owned создает автономное будущее, владеющее состоянием уведомления, благодаря чему его можно безопасно перемещать между потоками.
Подробнее см. Self::notified.
Безопасность отмены
Этот метод использует очередь для справедливого распределения уведомлений в порядке их запроса. Отмена вызова notified_owned приводит к потере места в очереди.
Примеры
use std::sync::Arc;
use tokio::sync::Notify;
let notify = Arc::new(Notify::new());
for _ in 0..10 {
let notified = notify.clone().notified_owned();
tokio::spawn(async move {
notified.await;
println!("received notification");
});
}
println!("sending notification");
notify.notify_waiters();pub fn notify_one(&self)
Уведомляет первую ожидающую задачу.
Если задача ожидает в данный момент, она будет уведомлена. В противном случае разрешение сохраняется в этом значении Notify, и следующий вызов notified().await завершится немедленно, использовав разрешение, предоставленное этим вызовом notify_one().
Метод Notify может сохранять не более одного разрешения. Несколько последовательных вызовов notify_one приведут к сохранению одного разрешения. Следующий вызов notified().await завершится немедленно, а вызов после него будет ожидать.
Примеры
use tokio::sync::Notify;
use std::sync::Arc;
let notify = Arc::new(Notify::new());
let notify2 = notify.clone();
tokio::spawn(async move {
notify2.notified().await;
println!("received notification");
});
println!("sending notification");
notify.notify_one();pub fn notify_last(&self)
Уведомляет последнюю ожидающую задачу.
Эта функция работает аналогично notify_one. Единственное отличие состоит в том, что она пробуждает самого недавно добавленного ожидающего, а не ожидающего дольше всех.
Дополнительную информацию и примеры см. в документации к notify_one().
pub fn notify_waiters(&self)
Уведомляет все ожидающие задачи.
Если задача ожидает в данный момент, она будет уведомлена. В отличие от notify_one(), разрешение не сохраняется для использования при следующем вызове notified().await. Этот метод предназначен для уведомления всех уже зарегистрированных ожидающих задач. Регистрация для получения уведомления выполняется путем получения экземпляра будущего Notified посредством вызова notified().
Примеры
use tokio::sync::Notify;
use std::sync::Arc;
let notify = Arc::new(Notify::new());
let notify2 = notify.clone();
let notified1 = notify.notified();
let notified2 = notify.notified();
let handle = tokio::spawn(async move {
println!("sending notifications");
notify2.notify_waiters();
});
notified1.await;
notified2.await;
println!("received notifications");Реализации трейтов
impl RefUnwindSafe for Notify
impl UnwindSafe for Notify
Автоматические реализации трейтов
impl !Freeze for Notify
impl Send for Notify
impl Sync for Notify
impl Unpin for Notify
impl UnsafeUnpin for Notify
Общие реализации
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<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/struct.Notify.html