Spec-Zone.ru › Tokio

Структура Notify

pub struct Notify { /* private fields */ }
Доступно только при включённой функции crate 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 может произойти следующее:

  1. Оба вызова try_recv возвращают None.
  2. Оба новых элемента добавляются в вектор.
  3. Метод notify_one вызывается дважды, но в Notify добавляется только одно разрешение.
  4. Оба вызова recv достигают future Notified. Один из них использует разрешение, а другой засыпает навсегда.

Добавив 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 Debug for Notify

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

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

impl Default for Notify

fn default() -> Notify

Возвращает «значение по умолчанию» для типа. Подробнее

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

Spec-Zone.ru

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