Spec-Zone.ru › Tokio

Макрос select

macro_rules! select {
    {
        $(
            biased;
        )?
        $(
            $bind:pat = $fut:expr $(, if $cond:expr)? => $handler:expr,
        )*
        $(
            else => $els:expr $(,)?
        )?
    } => { ... };
}
Доступен только при включённой функции crate macros.

Ожидает завершения нескольких параллельных ветвей, возвращая управление, когда завершается первая ветвь, и отменяя остальные.

Макрос select! должен использоваться внутри асинхронных функций, замыканий и блоков.

Макрос select! принимает одну или несколько ветвей следующего вида:

<pattern> = <async expression> (, if <precondition>)? => <handler>,

Кроме того, макрос select! может содержать одну необязательную ветвь else, которая выполняется, если ни одна из других ветвей не соответствует своим шаблонам:

else => <expression>

Макрос объединяет все выражения <async expression> и выполняет их параллельно в рамках текущей задачи. Как только первое выражение завершается значением, соответствующим его <pattern>, макрос select! возвращает результат вычисления выражения <handler> завершившейся ветви.

Кроме того, каждая ветвь может содержать необязательное предварительное условие if. Если предварительное условие возвращает false, ветвь отключается. Указанное <async expression> всё равно вычисляется, но полученное будущее никогда не опрашивается. Это удобно при использовании select! в цикле.

Полный жизненный цикл выражения select! выглядит следующим образом:

  1. Вычислить все указанные выражения <precondition>. Если предварительное условие возвращает false, отключить ветвь до конца текущего вызова select!. Повторный вход в select! из-за цикла сбрасывает состояние «отключена».
  2. Объединить <async expression> всех ветвей, включая отключённые. Если ветвь отключена, <async expression> всё равно вычисляется, но полученное будущее не опрашивается.
  3. Если отключены все ветви, перейти к шагу 6.
  4. Ожидать параллельно результаты всех оставшихся <async expression>.
  5. Когда <async expression> возвращает значение, попытаться сопоставить его с указанным <pattern>. Если шаблон совпадает, вычислить <handler> и выполнить возврат. Если шаблон не совпадает, отключить текущую ветвь до конца текущего вызова select!. Продолжить с шага 3.
  6. Вычислить выражение else. Если выражение else не указано, вызвать панику.

Особенности выполнения

Поскольку все асинхронные выражения выполняются в текущей задаче, они могут выполняться параллельно, но не одновременно. Это означает, что все выражения выполняются в одном потоке, и если одна ветвь блокирует поток, остальные выражения не смогут продолжить выполнение. Если требуется параллельное выполнение, запустите каждое асинхронное выражение с помощью tokio::spawn и передайте дескриптор присоединения в select!.

Справедливость

По умолчанию select! случайным образом выбирает ветвь, которую проверит первой. Это обеспечивает некоторую степень справедливости при вызове select! в цикле с ветвями, которые всегда готовы к выполнению.

Это поведение можно изменить, добавив biased; в начало использования макроса. Подробности см. в примерах. В этом случае select будет опрашивать будущие в порядке их появления сверху вниз. Это может быть нужно по нескольким причинам:

  • Генерация случайных чисел в tokio::select! требует ненулевых затрат процессорного времени.
  • Будущие могут взаимодействовать таким образом, что известный порядок опроса имеет значение.

Однако у этого режима есть важный нюанс. Вы сами должны следить за справедливостью порядка опроса будущих. Например, если вы выбираете между потоком и будущим завершения работы, а поток содержит огромное количество сообщений с нулевым или почти нулевым интервалом между ними, следует поместить будущее завершения работы раньше в списке select!, чтобы оно всегда опрашивалось и не игнорировалось из-за того, что поток постоянно готов.

Паники

Макрос select! вызывает панику, если отключены все ветви и ветвь else не указана. Ветвь отключается, если указанное предварительное условие if возвращает false или если шаблон не соответствует результату <async expression>.

Безопасность отмены

При использовании select! в цикле для получения сообщений из нескольких источников убедитесь, что вызов получения безопасен при отмене, чтобы избежать потери сообщений. В этом разделе рассматриваются различные распространённые методы и описывается, безопасны ли они при отмене. Списки в этом разделе не являются исчерпывающими.

Безопасность отмены описывает, что происходит, когда будущее отбрасывается до завершения. Безопасность отмены зависит от поведения будущего, переданного в select!. Такое будущее может быть получено из асинхронного метода, асинхронного выражения или другой операции, создающей будущее.

Следующие методы безопасны при отмене:

  • tokio::sync::mpsc::Receiver::recv
  • tokio::sync::mpsc::UnboundedReceiver::recv
  • tokio::sync::broadcast::Receiver::recv
  • tokio::sync::watch::Receiver::changed
  • tokio::net::TcpListener::accept
  • tokio::net::UnixListener::accept
  • tokio::signal::unix::Signal::recv
  • tokio::io::AsyncReadExt::read для любого AsyncRead
  • tokio::io::AsyncReadExt::read_buf для любого AsyncRead
  • tokio::io::AsyncWriteExt::write для любого AsyncWrite
  • tokio::io::AsyncWriteExt::write_buf для любого AsyncWrite
  • tokio_stream::StreamExt::next для любого Stream
  • futures::stream::StreamExt::next для любого Stream

Следующие методы небезопасны при отмене и могут привести к потере данных:

  • tokio::io::AsyncReadExt::read_exact
  • tokio::io::AsyncReadExt::read_to_end
  • tokio::io::AsyncReadExt::read_to_string
  • tokio::io::AsyncWriteExt::write_all

Следующие методы небезопасны при отмене, поскольку для обеспечения справедливости они используют очередь, а при отмене вы теряете своё место в очереди:

  • tokio::sync::Mutex::lock
  • tokio::sync::RwLock::read
  • tokio::sync::RwLock::write
  • tokio::sync::Semaphore::acquire
  • tokio::sync::Notify::notified

Чтобы определить, безопасны ли при отмене ваши собственные методы, найдите места использования .await. Это связано с тем, что отмена асинхронного метода всегда происходит в точке .await. Если функция работает корректно, даже когда её выполнение начинается заново во время ожидания в точке .await, значит, она безопасна при отмене.

Безопасность отмены можно определить следующим образом: если у вас есть ещё не завершившееся будущее, его отбрасывание и создание заново не должны иметь никаких побочных эффектов. Это определение обусловлено ситуацией, когда select! используется в цикле. Без этой гарантии вы потеряете достигнутый прогресс, когда завершится другая ветвь и вы заново запустите select!, перейдя на следующую итерацию цикла.

Помните, что отмена операции, небезопасной при отмене, не обязательно является ошибкой. Например, если вы отменяете задачу из-за завершения работы приложения, вероятно, вас не беспокоит потеря частично прочитанных данных.

Примеры

Простой select с двумя ветвями.

async fn do_stuff_async() {
    // async work
}

async fn more_async_work() {
    // more here
}

tokio::select! {
    _ = do_stuff_async() => {
        println!("do_stuff_async() completed first")
    }
    _ = more_async_work() => {
        println!("more_async_work() completed first")
    }
};

Простой выбор из потоков.

use tokio_stream::{self as stream, StreamExt};

let mut stream1 = stream::iter(vec![1, 2, 3]);
let mut stream2 = stream::iter(vec![4, 5, 6]);

let next = tokio::select! {
    v = stream1.next() => v.unwrap(),
    v = stream2.next() => v.unwrap(),
};

assert!(next == 1 || next == 4);

Собрать содержимое двух потоков. В этом примере мы полагаемся на сопоставление с шаблоном и на то, что stream::iter является «слитым», то есть после завершения потока все вызовы next() возвращают None.

use tokio_stream::{self as stream, StreamExt};

let mut stream1 = stream::iter(vec![1, 2, 3]);
let mut stream2 = stream::iter(vec![4, 5, 6]);

let mut values = vec![];

loop {
    tokio::select! {
        Some(v) = stream1.next() => values.push(v),
        Some(v) = stream2.next() => values.push(v),
        else => break,
    }
}

values.sort();
assert_eq!(&[1, 2, 3, 4, 5, 6], &values[..]);

Одно и то же будущее можно использовать в нескольких выражениях select!, передав ссылку на него. Для этого будущее должно реализовывать Unpin. Будущее можно сделать Unpin, используя Box::pin или закрепив его в стеке.

В этом примере поток обрабатывается не более 1 секунды.

use tokio_stream::{self as stream, StreamExt};
use tokio::time::{self, Duration};

let mut stream = stream::iter(vec![1, 2, 3]);
let sleep = time::sleep(Duration::from_secs(1));
tokio::pin!(sleep);

loop {
    tokio::select! {
        maybe_v = stream.next() => {
            if let Some(v) = maybe_v {
                println!("got = {}", v);
            } else {
                break;
            }
        }
        _ = &mut sleep => {
            println!("timeout");
            break;
        }
    }
}

Объединение двух значений с помощью select!.

use tokio::sync::oneshot;

let (tx1, mut rx1) = oneshot::channel();
let (tx2, mut rx2) = oneshot::channel();

tokio::spawn(async move {
    tx1.send("first").unwrap();
});

tokio::spawn(async move {
    tx2.send("second").unwrap();
});

let mut a = None;
let mut b = None;

while a.is_none() || b.is_none() {
    tokio::select! {
        v1 = (&mut rx1), if a.is_none() => a = Some(v1.unwrap()),
        v2 = (&mut rx2), if b.is_none() => b = Some(v2.unwrap()),
    }
}

let res = (a.unwrap(), b.unwrap());

assert_eq!(res.0, "first");
assert_eq!(res.1, "second");

Использование режима biased; для управления порядком опроса.

let mut count = 0u8;

loop {
    tokio::select! {
        // If you run this example without `biased;`, the polling order is
        // pseudo-random, and the assertions on the value of count will
        // (probably) fail.
        biased;

        _ = async {}, if count < 1 => {
            count += 1;
            assert_eq!(count, 1);
        }
        _ = async {}, if count < 2 => {
            count += 1;
            assert_eq!(count, 2);
        }
        _ = async {}, if count < 3 => {
            count += 1;
            assert_eq!(count, 3);
        }
        _ = async {}, if count < 4 => {
            count += 1;
            assert_eq!(count, 4);
        }

        else => {
            break;
        }
    };
}

Избегайте гонок при использовании предварительных условий if

Поскольку предварительные условия if используются для отключения ветвей select!, необходимо соблюдать осторожность, чтобы не пропустить значения.

Например, ниже показано неправильное использование sleep с if. Цель — многократно выполнять асинхронную задачу в течение не более 50 миллисекунд. Однако существует риск пропустить завершение sleep.

ⓘ
use tokio::time::{self, Duration};

async fn some_async_work() {
    // do work
}

let sleep = time::sleep(Duration::from_millis(50));
tokio::pin!(sleep);

while !sleep.is_elapsed() {
    tokio::select! {
        _ = &mut sleep, if !sleep.is_elapsed() => {
            println!("operation timed out");
        }
        _ = some_async_work() => {
            println!("operation completed");
        }
    }
}

panic!("This example shows how not to do it!");

В приведённом выше примере sleep.is_elapsed() может вернуть true, даже если sleep.poll() никогда не возвращал Ready. Это создаёт потенциальное состояние гонки: срок действия sleep может истечь между проверкой while !sleep.is_elapsed() и вызовом select!, в результате чего вызов some_async_work() выполнится без прерывания, хотя время ожидания уже истекло.

Один из способов переписать приведённый выше пример без гонки:

use tokio::time::{self, Duration};

async fn some_async_work() {
    // do work
}

let sleep = time::sleep(Duration::from_millis(50));
tokio::pin!(sleep);

loop {
    tokio::select! {
        _ = &mut sleep => {
            println!("operation timed out");
            break;
        }
        _ = some_async_work() => {
            println!("operation completed");
        }
    }
}

Альтернативы в экосистеме

Макрос select! — мощный инструмент для управления несколькими асинхронными ветвями, позволяющий выполнять задачи параллельно в одном потоке. Однако его использование может создавать трудности, особенно в отношении безопасности отмены, что способно приводить к незаметным ошибкам, которые сложно отлаживать. Во многих случаях предпочтительнее альтернативы из экосистемы: они помогают избежать этих проблем благодаря более ясному синтаксису, предсказуемому управлению потоком выполнения и отсутствию необходимости вручную обрабатывать такие вопросы, как семантика слияния или безопасность отмены.

Объединение потоков

Если loop { select! { ... } } используется для опроса нескольких задач, объединение потоков представляет собой лаконичную альтернативу: оно изначально обеспечивает безопасную при отмене обработку и устраняет риск потери данных. Такие библиотеки, как tokio_stream, futures::stream и futures_concurrency, предоставляют средства для объединения потоков и последовательной обработки их результатов.

Пример с select!

struct File;
struct Channel;
struct Socket;

impl Socket {
    async fn read_packet(&mut self) -> Vec<u8> {
        vec![]
    }
}

async fn read_send(_file: &mut File, _channel: &mut Channel) {
    // do work that is not cancel safe
}

// open our IO types
let mut file = File;
let mut channel = Channel;
let mut socket = Socket;

loop {
    tokio::select! {
        _ = read_send(&mut file, &mut channel) => { /* ... */ },
        _data = socket.read_packet() => { /* ... */ }
        _ = futures::future::ready(()) => break
    }
}

Переход на merge

С помощью merge можно объединить несколько асинхронных задач в единый поток, избавившись от необходимости управлять задачами вручную и снизив риск нежелательного поведения, например потери данных.

use std::pin::pin;

use futures::stream::unfold;
use tokio_stream::StreamExt;

struct File;
struct Channel;
struct Socket;

impl Socket {
    async fn read_packet(&mut self) -> Vec<u8> {
        vec![]
    }
}

async fn read_send(_file: &mut File, _channel: &mut Channel) {
    // do work that is not cancel safe
}

enum Message {
    Stop,
    Sent,
    Data(Vec<u8>),
}

// open our IO types
let file = File;
let channel = Channel;
let socket = Socket;

let a = unfold((file, channel), |(mut file, mut channel)| async {
    read_send(&mut file, &mut channel).await;
    Some((Message::Sent, (file, channel)))
});
let b = unfold(socket, |mut socket| async {
    let data = socket.read_packet().await;
    Some((Message::Data(data), socket))
});
let c = tokio_stream::iter([Message::Stop]);

let mut s = pin!(a.merge(b).merge(c));
while let Some(msg) = s.next().await {
    match msg {
        Message::Data(_data) => { /* ... */ }
        Message::Sent => continue,
        Message::Stop => break,
    }
}

Гонка будущих

Если нужно дождаться завершения первой из нескольких асинхронных задач, утилиты из экосистемы, такие как futures, futures-lite или futures-concurrency, предоставляют лаконичный синтаксис для организации гонки будущих:

  • futures_concurrency::future::Race
  • futures::select
  • futures::stream::select_all (для потоков)
  • futures_lite::future::or
  • futures_lite::future::race
use futures_concurrency::future::Race;

let task_a = async { Ok("ok") };
let task_b = async { Err("error") };
let result = (task_a, task_b).race().await;

match result {
    Ok(output) => println!("First task completed with: {output}"),
    Err(err) => eprintln!("Error occurred: {err}"),
}

MIT License
Copyright © Tokio Contributors
https://docs.rs/tokio/1.53.1/tokio/macro.select.html

Spec-Zone.ru

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