Spec-Zone.ru › Tokio

Структура LocalSet

pub struct LocalSet { /* private fields */ }
Доступно только при включённой возможности crate rt.

Набор задач, выполняемых в одном потоке.

В некоторых случаях необходимо запустить один или несколько фьючерсов, которые не реализуют Send, а значит, их небезопасно передавать между потоками. В таких случаях для запуска одного или нескольких фьючерсов !Send в одном потоке можно использовать набор локальных задач.

Например, следующий код не скомпилируется:

ⓘ
use std::rc::Rc;

#[tokio::main]
async fn main() {
    // `Rc` does not implement `Send`, and thus may not be sent between
    // threads safely.
    let nonsend_data = Rc::new("my nonsend data...");

    let nonsend_data = nonsend_data.clone();
    // Because the `async` block here moves `nonsend_data`, the future is `!Send`.
    // Since `tokio::spawn` requires the spawned future to implement `Send`, this
    // will not compile.
    tokio::spawn(async move {
        println!("{}", nonsend_data);
        // ...
    }).await.unwrap();
}

Использование с run_until

Чтобы запускать фьючерсы !Send, мы можем использовать набор локальных задач, чтобы запланировать их выполнение в потоке, вызывающем Runtime::block_on. При работе внутри набора локальных задач можно использовать task::spawn_local, которая может запускать фьючерсы !Send. Например:

use std::rc::Rc;
use tokio::task;

let nonsend_data = Rc::new("my nonsend data...");

// Construct a local task set that can run `!Send` futures.
let local = task::LocalSet::new();

// Run the local task set.
local.run_until(async move {
    let nonsend_data = nonsend_data.clone();
    // `spawn_local` ensures that the future is spawned on the local
    // task set.
    task::spawn_local(async move {
        println!("{}", nonsend_data);
        // ...
    }).await.unwrap();
}).await;

Примечание: Метод run_until можно использовать только в #[tokio::main], #[tokio::test] или непосредственно внутри вызова Runtime::block_on. Его нельзя использовать внутри задачи, запущенной с помощью tokio::spawn.

Ожидание LocalSet

Кроме того, сам LocalSet реализует Future и завершается, когда завершаются все задачи, запущенные в LocalSet. Это можно использовать для запуска нескольких фьючерсов в LocalSet и ожидания завершения всего набора. Например:

use tokio::{task, time};
use std::rc::Rc;

let nonsend_data = Rc::new("world");
let local = task::LocalSet::new();

let nonsend_data2 = nonsend_data.clone();
local.spawn_local(async move {
    // ...
    println!("hello {}", nonsend_data2)
});

local.spawn_local(async move {
    time::sleep(time::Duration::from_millis(100)).await;
    println!("goodbye {}", nonsend_data)
});

// ...

local.await;

Примечание: Ожидать LocalSet можно только внутри #[tokio::main], #[tokio::test] или непосредственно внутри вызова Runtime::block_on. Это нельзя делать внутри задачи, запущенной с помощью tokio::spawn.

Использование внутри tokio::spawn

Два упомянутых выше метода нельзя использовать внутри tokio::spawn, поэтому для запуска фьючерсов !Send изнутри tokio::spawn нужно поступить иначе. Решение — создать LocalSet в другом месте и взаимодействовать с ним с помощью канала mpsc.

В следующем примере LocalSet размещается в новом потоке.

use tokio::runtime::Builder;
use tokio::sync::{mpsc, oneshot};
use tokio::task::LocalSet;

// This struct describes the task you want to spawn. Here we include
// some simple examples. The oneshot channel allows sending a response
// to the spawner.
#[derive(Debug)]
enum Task {
    PrintNumber(u32),
    AddOne(u32, oneshot::Sender<u32>),
}

#[derive(Clone)]
struct LocalSpawner {
   send: mpsc::UnboundedSender<Task>,
}

impl LocalSpawner {
    pub fn new() -> Self {
        let (send, mut recv) = mpsc::unbounded_channel();

        let rt = Builder::new_current_thread()
            .enable_all()
            .build()
            .unwrap();

        std::thread::spawn(move || {
            let local = LocalSet::new();

            local.spawn_local(async move {
                while let Some(new_task) = recv.recv().await {
                    tokio::task::spawn_local(run_task(new_task));
                }
                // If the while loop returns, then all the LocalSpawner
                // objects have been dropped.
            });

            // This will return once all senders are dropped and all
            // spawned tasks have returned.
            rt.block_on(local);
        });

        Self {
            send,
        }
    }

    pub fn spawn(&self, task: Task) {
        self.send.send(task).expect("Thread with LocalSet has shut down.");
    }
}

// This task may do !Send stuff. We use printing a number as an example,
// but it could be anything.
//
// The Task struct is an enum to support spawning many different kinds
// of operations.
async fn run_task(task: Task) {
    match task {
        Task::PrintNumber(n) => {
            println!("{}", n);
        },
        Task::AddOne(n, response) => {
            // We ignore failures to send the response.
            let _ = response.send(n + 1);
        },
    }
}

#[tokio::main]
async fn main() {
    let spawner = LocalSpawner::new();

    let (send, response) = oneshot::channel();
    spawner.spawn(Task::AddOne(10, send));
    let eleven = response.await.unwrap();
    assert_eq!(eleven, 11);
}

Реализации

impl LocalSet

pub fn new() -> LocalSet ⓘ

Возвращает новый набор локальных задач.

pub fn enter(&self) -> LocalEnterGuard

Входит в контекст этого LocalSet.

Метод spawn_local запускает задачи в LocalSet, в контексте которого вы находитесь.

pub fn spawn_local<F>(&self, future: F) -> JoinHandle<F::Output> ⓘ
where F: Future + 'static, F::Output: 'static,

Запускает задачу !Send в наборе локальных задач.

Гарантируется, что эта задача будет выполняться в текущем потоке.

В отличие от свободной функции spawn_local, этот метод можно использовать для запуска локальных задач, когда LocalSet не запущен. Переданный future начнёт выполняться при следующем запуске LocalSet, даже если вы не ожидаете возвращённый JoinHandle.

Примеры
use tokio::task;

let local = task::LocalSet::new();

// Spawn a future on the local set. This future will be run when
// we call `run_until` to drive the task set.
local.spawn_local(async {
    // ...
});

// Run the local task set.
local.run_until(async move {
    // ...
}).await;

// When `run` finishes, we can spawn _more_ futures, which will
// run in subsequent calls to `run_until`.
local.spawn_local(async {
    // ...
});

local.run_until(async move {
    // ...
}).await;

pub fn block_on<F>(&self, rt: &Runtime, future: F) -> F::Output
where F: Future,

Выполняет future до завершения в предоставленной среде выполнения, обслуживая в текущем потоке все локальные future, запущенные в этом наборе задач.

Этот метод выполняет переданный future в среде выполнения, блокируя выполнение до его завершения и возвращая полученный результат. Все задачи и таймеры, созданные future, будут выполняться в среде выполнения. Future также может вызвать spawn_local, чтобы spawn_local дополнительные локальные future в текущем потоке.

Этот метод нельзя вызывать из асинхронного контекста.

Панические завершения

Эта функция вызывает панику, если исполнитель достиг предельной нагрузки, если предоставленный future вызывает панику или если функция вызвана в контексте асинхронного выполнения.

Примечания

Поскольку эта функция внутри вызывает Runtime::block_on и обслуживает future из локального набора задач внутри этого вызова до block_on, локальные future не могут использовать блокировку на месте. Если из локальной задачи необходимо выполнить блокирующий вызов, вместо этого можно использовать API spawn_blocking.

Например, это вызовет панику:

ⓘ
use tokio::runtime::Runtime;
use tokio::task;

let rt  = Runtime::new().unwrap();
let local = task::LocalSet::new();
local.block_on(&rt, async {
    let join = task::spawn_local(async {
        let blocking_result = task::block_in_place(|| {
            // ...
        });
        // ...
    });
    join.await.unwrap();
})

Однако это не вызовет панику:

use tokio::runtime::Runtime;
use tokio::task;

let rt  = Runtime::new().unwrap();
let local = task::LocalSet::new();
local.block_on(&rt, async {
    let join = task::spawn_local(async {
        let blocking_result = task::spawn_blocking(|| {
            // ...
        }).await;
        // ...
    });
    join.await.unwrap();
})

pub async fn run_until<F>(&self, future: F) -> F::Output
where F: Future,

Выполняет future до завершения в локальном наборе и возвращает его результат.

Этот метод возвращает future, который выполняет переданный future с локальным набором, позволяя ему вызывать spawn_local для запуска дополнительных !Send future. Все локальные future, запущенные в локальном наборе, будут выполняться в фоновом режиме, пока не завершится future, переданный в run_until. Когда future, переданный в run_until, завершится, все незавершённые локальные future останутся в локальном наборе и будут выполняться при последующих вызовах run_until или при ожидании самого локального набора.

Безопасность при отмене

Этот метод безопасен при отмене, если future безопасен при отмене.

Примеры
use tokio::task;

task::LocalSet::new().run_until(async {
    task::spawn_local(async move {
        // ...
    }).await.unwrap();
    // ...
}).await;

pub fn id(&self) -> Id

Возвращает Id текущей среды выполнения LocalSet.

Примеры
use tokio::task;

let local_set = task::LocalSet::new();
println!("Local set id: {}", local_set.id());

impl LocalSet

pub fn unhandled_panic(&mut self, behavior: UnhandledPanic) -> &mut Self

Доступно только на tokio_unstable.

Настраивает поведение LocalSet при необработанной панике в запущенной задаче.

По умолчанию необработанная паника (то есть паника, не перехваченная с помощью std::panic::catch_unwind) не влияет на выполнение LocalSet. Значение ошибки паники передаётся в JoinHandle задачи, а все остальные запущенные задачи продолжают выполняться.

Параметр unhandled_panic позволяет настроить это поведение.

  • UnhandledPanic::Ignore — поведение по умолчанию. Паники в запущенных задачах не влияют на выполнение LocalSet.
  • UnhandledPanic::ShutdownRuntime заставит LocalSet немедленно завершить работу при панике в запущенной задаче, даже если её JoinHandle ещё не был удалён. Все остальные запущенные задачи немедленно завершатся, а последующие вызовы LocalSet::block_on и LocalSet::run_until вызовут панику.
Паники

Этот метод вызывает панику, если его вызвать после начала выполнения LocalSet.

Нестабильный API

Этот параметр в настоящее время является нестабильным, а его реализация не завершена. В будущем API может измениться или быть удалён. Подробнее см. tokio-rs/tokio#4516.

Примеры

В следующем примере показан LocalSet, настроенный на завершение работы при панике. Первая запущенная задача вызывает панику, в результате чего LocalSet завершает работу. Вторая запущенная задача не успевает выполниться. Вызов run_until вызовет панику, поскольку среда выполнения была принудительно остановлена.

ⓘ
use tokio::runtime::UnhandledPanic;

tokio::task::LocalSet::new()
    .unhandled_panic(UnhandledPanic::ShutdownRuntime)
    .run_until(async {
        tokio::task::spawn_local(async { panic!("boom"); });
        tokio::task::spawn_local(async {
            // This task never completes
        });

        // Do some work, but `run_until` will panic before it completes
    })
    .await;

Реализации трейтов

impl Debug for LocalSet

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

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

impl Default for LocalSet

fn default() -> LocalSet ⓘ

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

impl Drop for LocalSet

fn drop(&mut self)

Выполняет деструктор для этого типа. Подробнее

fn pin_drop(self: Pin<&mut Self>)

🔬Это экспериментальный API, доступный только в nightly-версии. (pin_ergonomics)
Выполняет деструктор для этого типа, но, в отличие от Drop::drop, требует, чтобы self был закреплён. Подробнее

impl Future for LocalSet

type Output = ()

Тип значения, возвращаемого при завершении.

fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>

Пытается получить итоговое значение из future; если значение ещё недоступно, регистрирует текущую задачу для пробуждения. Подробнее

Автоматические реализации трейтов

impl !Freeze for LocalSet

impl !RefUnwindSafe for LocalSet

impl !Send for LocalSet

impl !Sync for LocalSet

impl !UnwindSafe for LocalSet

impl Unpin for LocalSet

impl UnsafeUnpin for LocalSet

Общие реализации

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<F> IntoFuture for F
where F: Future,

type Output = <F as Future>::Output

Результат, который будет получен при завершении будущего значения.

type IntoFuture = F

В какой тип будущего значения мы преобразуем это значение?

fn into_future(self) -> <F as IntoFuture>::IntoFuture

Создаёт будущее значение из значения. Подробнее

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

Spec-Zone.ru

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