Раздел 23 · Rust

Свой исполнитель

senior~80 мин

открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти

Свой исполнитель

Два прошлых урока показали половину машинерии: async fn это машина состояний с методом poll, а Pin не даёт её двигать. Но кто-то же должен poll звать. Этого «кого-то» мы всё время называли executor’ом и обещали собрать руками. Собираем. По дороге сделаем Waker через сырую vtable, минимальный block_on, многозадачный executor с очередью готовых задач и reactor поверх epoll. После этого урока в async не останется ни одной непрозрачной детали.

Три роли рантайма

Любой async-рантайм (tokio, smol) это связка трёх ролей. Executor крутит poll готовых задач. Reactor следит за внешним миром (сокеты, таймеры) и сообщает, когда появилась готовность. А связывает их Waker: футура отдаёт его reactor’у, reactor зовёт wake, и задача возвращается к executor’у в очередь. Начнём с Waker, потому что без него ни одна из сторон не работает.

Waker руками

В стандартной библиотеке Waker строится из RawWaker и RawWakerVTable: сырого указателя на данные и таблицы из четырёх функций. Это самый низ, на котором держится всё пробуждение. Соберём поверх него удобный трейт ArcWake: реализуй wake_by_ref, и мы превратим Arc<тебя> в настоящий Waker.

use std::sync::Arc;
use std::task::{RawWaker, RawWakerVTable, Waker};

pub trait ArcWake: Send + Sync + 'static {
    fn wake_by_ref(self: &Arc<Self>);
}

pub fn waker_from_arc<W: ArcWake>(w: Arc<W>) -> Waker {
    let raw = raw_waker(w);
    // SAFETY: raw_waker строит RawWaker с указателем на Arc<W> и vtable,
    // функции которой ниже соблюдают контракт для этого типа.
    unsafe { Waker::from_raw(raw) }
}

fn raw_waker<W: ArcWake>(w: Arc<W>) -> RawWaker {
    // into_raw: владение «уходит в указатель», refcount дальше ведём руками.
    let ptr = Arc::into_raw(w) as *const ();
    RawWaker::new(ptr, vtable::<W>())
}

fn vtable<W: ArcWake>() -> &'static RawWakerVTable {
    &RawWakerVTable::new(clone_raw::<W>, wake_raw::<W>, wake_by_ref_raw::<W>, drop_raw::<W>)
}

Вся соль в четырёх функциях vtable. Каждая получает сырой указатель, который мы когда-то получили из Arc::into_raw, и должна аккуратно поддержать счётчик ссылок:

unsafe fn clone_raw<W: ArcWake>(ptr: *const ()) -> RawWaker {
    // SAFETY: восстанавливаем Arc, поднимаем refcount клоном, а исходный
    // указатель оставляем валидным через forget.
    let arc = unsafe { Arc::from_raw(ptr as *const W) };
    let cloned = arc.clone();
    std::mem::forget(arc);
    raw_waker(cloned)
}

unsafe fn wake_raw<W: ArcWake>(ptr: *const ()) {
    // SAFETY: wake забирает владение Waker'ом, значит этот Arc надо уронить.
    let arc = unsafe { Arc::from_raw(ptr as *const W) };
    arc.wake_by_ref();
    // arc дропается здесь, освобождая ссылку Waker'а.
}

unsafe fn wake_by_ref_raw<W: ArcWake>(ptr: *const ()) {
    // SAFETY: wake_by_ref владение не забирает, поэтому Arc возвращаем через forget.
    let arc = unsafe { Arc::from_raw(ptr as *const W) };
    arc.wake_by_ref();
    std::mem::forget(arc);
}

unsafe fn drop_raw<W: ArcWake>(ptr: *const ()) {
    // SAFETY: drop Waker'а роняет ровно одну ссылку.
    let arc = unsafe { Arc::from_raw(ptr as *const W) };
    drop(arc);
}

Это полный unsafe, но логика простая: clone поднимает счётчик, wake и drop его опускают, wake_by_ref не трогает. Перепутаешь forget и drop хоть в одной, получишь либо утечку, либо двойное освобождение. Именно поэтому в реальном коде берут трейт std::task::Wake: он делает ровно это поверх Arc, но без единого unsafe. Мы написали руками один раз, чтобы увидеть механику; дальше пользуемся waker_from_arc. Если хочешь освежить, что такое сырые указатели и почему тут нужен unsafe, вернись к уроку про unsafe Rust.

Чтобы это не осталось словами, вот тот же будильник для потока, собранный без нашей vtable, на готовом трейте std::task::Wake:

use std::sync::Arc;
use std::task::{Wake, Waker};
use std::thread::{self, Thread};

struct ThreadWaker(Thread);

impl Wake for ThreadWaker {
    fn wake(self: Arc<Self>) {
        self.0.unpark(); // распарковать ждущий поток
    }
}

// Waker::from сам строит RawWaker и vtable за нас.
let waker = Waker::from(Arc::new(ThreadWaker(thread::current())));

Одна реализация wake вместо четырёх unsafe-функций, а Waker::from собирает ту же RawWaker с vtable под капотом. Дальше в уроке мы оставим свой waker_from_arc, чтобы держать механику на виду, но в боевом коде ты напишешь именно так.

block_on: прогнать одну футуру

Самый простой executor гоняет ровно одну футуру до конца. Его Waker будит сам поток: пока футура Pending, поток паркуется и не жжёт процессор, а wake его распарковывает.

use std::future::Future;
use std::pin::Pin;
use std::task::{Context, Poll};
use std::thread::{self, Thread};

pub fn block_on<F: Future>(future: F) -> F::Output {
    let mut future = future;
    // SAFETY: future живёт до конца функции и больше не двигается (работаем
    // только через &mut по этой Pin), поэтому закрепление на стеке корректно.
    let mut future = unsafe { Pin::new_unchecked(&mut future) };

    let parker = Arc::new(ThreadWaker { thread: thread::current() });
    let waker = waker_from_arc(parker);
    let mut cx = Context::from_waker(&waker);

    loop {
        match future.as_mut().poll(&mut cx) {
            Poll::Ready(value) => return value,
            Poll::Pending => thread::park(),
        }
    }
}

struct ThreadWaker { thread: Thread }

impl ArcWake for ThreadWaker {
    fn wake_by_ref(self: &Arc<Self>) {
        self.thread.unpark();
    }
}

Прочитай цикл: poll, если Ready, вернули результат; если Pending, паркуем поток. Кто-то (сама футура через yield_now или reactor) позовёт wake, поток распаркуется, и мы полим снова. Ноль холостого опроса, как и обещал принцип «не звони мне, я позвоню тебе».

Многозадачный executor

block_on гоняет одну футуру. Чтобы крутить много задач сразу, нужна очередь. Идея: задача держит свою футуру и умеет положить себя обратно в очередь готовых; её Waker именно это и делает.

use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
use std::sync::Mutex;

type BoxFuture = Pin<Box<dyn Future<Output = ()> + Send>>;

struct Task {
    future: Mutex<Option<BoxFuture>>,
    ready_queue: SyncSender<Arc<Task>>,
}

impl ArcWake for Task {
    fn wake_by_ref(self: &Arc<Self>) {
        // Пробуждение это «верни меня в очередь готовых». poll сделает run.
        let _ = self.ready_queue.send(self.clone());
    }
}

pub struct Executor { ready_queue: Receiver<Arc<Task>> }

#[derive(Clone)]
pub struct Spawner { ready_queue: SyncSender<Arc<Task>> }

pub fn new_executor_and_spawner() -> (Executor, Spawner) {
    let (tx, rx) = sync_channel(10_000);
    (Executor { ready_queue: rx }, Spawner { ready_queue: tx })
}

Футура лежит под Mutex, хотя executor однопоточный: так требуют типы, ведь Waker обязан быть Send + Sync. spawn заворачивает футуру в задачу и кладёт в очередь:

impl Spawner {
    pub fn spawn(&self, future: impl Future<Output = ()> + Send + 'static) {
        let task = Arc::new(Task {
            future: Mutex::new(Some(Box::pin(future))),
            ready_queue: self.ready_queue.clone(),
        });
        self.ready_queue.send(task).expect("очередь готовых задач закрыта");
    }
}

impl Executor {
    pub fn run(&self) {
        while let Ok(task) = self.ready_queue.recv() {
            let mut slot = task.future.lock().expect("future mutex отравлен");
            if let Some(mut future) = slot.take() {
                let waker = waker_from_arc(task.clone());
                let mut cx = Context::from_waker(&waker);
                match future.as_mut().poll(&mut cx) {
                    // Готова: футуру не возвращаем, slot остаётся None.
                    Poll::Ready(()) => {}
                    // Не готова: кладём обратно, задача вернётся в очередь по wake.
                    Poll::Pending => *slot = Some(future),
                }
            }
        }
    }
}

Весь цикл: вынь задачу, забери её футуру, полей один раз. Ready значит задача кончилась (футуру не возвращаем). Pending значит футура уже сохранила Waker где надо и вернётся в очередь, когда её разбудят, а пока кладём её обратно в задачу. Условие выхода из run: канал закрылся, то есть все Spawner сброшены и очередь пуста. Поэтому перед run исходный Spawner обычно дропают.

let (executor, spawner) = new_executor_and_spawner();
for i in 0..10 {
    spawner.spawn(async move {
        yield_now().await;
        println!("задача {i}");
    });
}
drop(spawner);     // иначе run никогда не выйдет
executor.run();

Останови взгляд на границах spawn: Future + Send + 'static. 'static тут не каприз, а прямое следствие устройства. Задача живёт в рантайме сколько угодно долго и не привязана к кадру стека того, кто её заспавнил; этот кадр давно исчезнет, а задача будет жить. Значит, одолжить ссылку с коротким временем жизни она не может: всё, что ей нужно, она обязана забрать во владение. Отсюда move в async move и 'static в сигнатуре. Send пока избыточен (executor однопоточный), но он понадобится сразу, как только задачи начнут мигрировать между потоками в work-stealing-рантайме, к которому мы придём в следующем уроке.

Reactor: встреча poll с готовностью I/O

Чего executor сам не умеет, так это узнать, что сетевой сокет стал читаемым. Крутить «а готово ли уже?» нельзя, это и есть busy-loop, который мы поклялись не делать. Эту работу берёт reactor поверх epoll. Кроссплатформенную обёртку над epoll/kqueue даёт крейт mio.

Схема: reactor держит mio::Poll и реестр Token -> Waker. Футура сокета при Pending кладёт свой Waker в реестр под своим токеном. Метод turn зовёт mio::Poll::poll, получает список готовых токенов и будит их Waker’ы.

use mio::{Events, Poll as MioPoll, Token};
use std::collections::HashMap;
use std::task::Waker;
use std::time::Duration;

pub struct Reactor {
    poll: Mutex<MioPoll>,
    wakers: Mutex<HashMap<Token, Waker>>,
    next_token: Mutex<usize>,
}

impl Reactor {
    /// Запомнить Waker для токена (зовёт футура сокета при Pending).
    fn store_waker(&self, token: Token, waker: Waker) {
        self.wakers.lock().unwrap().insert(token, waker);
    }

    /// Один оборот: ждём готовности до timeout и будим готовые токены.
    pub fn turn(&self, timeout: Option<Duration>) -> std::io::Result<usize> {
        let mut events = Events::with_capacity(64);
        self.poll.lock().unwrap().poll(&mut events, timeout)?;

        let mut woken = 0;
        let mut wakers = self.wakers.lock().unwrap();
        for event in events.iter() {
            if let Some(waker) = wakers.remove(&event.token()) {
                waker.wake();
                woken += 1;
            }
        }
        Ok(woken)
    }
}

Две детали в turn стоят отдельного слова. Первое: после события мы делаем wakers.remove(&token), то есть снимаем Waker, а не оставляем его. Это ровно поведение флага EPOLLONESHOT в сыром epoll: дескриптор после события снимается с наблюдения, и задача сама перевзводит интерес на следующем poll, когда снова упрётся в WouldBlock. Без этого epoll в уровневом режиме будил бы нас по кругу, пока в сокете есть хоть байт, и вместо сна мы получили бы busy-loop. Второе: HashMap<Token, Waker> держит ровно один Waker на токен. Жди один и тот же дескриптор сразу две задачи, новый клон затёр бы предыдущий, и одна из них зависла бы. В боевых рантаймах под каждый дескриптор хранят отдельные слоты на чтение и запись; мы держим по одному намеренно, чтобы не утонуть в деталях.

А вот как футура сокета встречается с готовностью. Асинхронный TCP-поток оборачивает неблокирующий mio::net::TcpStream. Его read пробует прочитать: если данные есть, отдаёт Ready; если ОС вернула WouldBlock, значит данных нет, и тогда футура кладёт Waker в reactor и засыпает.

impl Future for ReadFut<'_> {
    type Output = std::io::Result<usize>;

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
        let this = self.get_mut();
        match this.stream.inner.read(this.buf) {
            Ok(n) => Poll::Ready(Ok(n)),
            Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {
                // Готовности нет: оставляем Waker, reactor разбудит при данных.
                this.stream.reactor.store_waker(this.stream.token, cx.waker().clone());
                Poll::Pending
            }
            Err(e) => Poll::Ready(Err(e)),
        }
    }
}

Сложи всё вместе и увидишь полный цикл async I/O. Executor полит задачу. Задача читает сокет, данных нет, она кладёт Waker в reactor и возвращает Pending. Executor видит, что готовых задач не осталось, и вместо busy-loop зовёт reactor.turn(timeout), который блокируется в epoll на уровне ОС. Приходят данные, epoll просыпается, reactor находит Waker по токену и зовёт wake, задача возвращается в очередь, executor полит её снова, и на этот раз read отдаёт байты. Вот это и значит «poll встречается с эпохой готовности».

Свежий Waker на каждый poll

Посмотри ещё раз на poll у ReadFut: на ветке WouldBlock он каждый раз кладёт в reactor cx.waker().clone(), свежий клон из текущего Context. Это не лишняя перестраховка, а часть контракта Future. Правило жёсткое: на каждом poll футура обязана запомнить именно последний полученный Waker и забыть старый.

Соблазн закэшировать первый Waker и звонить потом только в него выходит боком в многопоточном рантайме. Tokio после work-stealing может перенести задачу на другой воркер-поток, и Context принесёт уже другой Waker, привязанный к новому потоку. Если футура хранит старый, пробуждение уйдёт в никуда: reactor позовёт wake у того, кто ждёт на потоке, где задачи давно нет. Симптом коварный: по всем логам задача обязана была проснуться, но спит, не потребляя при этом ни процента процессора. Это классический баг ручных футур, и ищется он мучительно. Лечится одной строкой: всегда обновляй сохранённый Waker через cx.waker().clone() и никогда не считай, что между опросами он тот же самый.

Полный async-сокет

ReadFut это половина картины. Чтобы цикл реально замкнулся, нужен сам сокет, который эти футуры выдаёт. Оборачиваем неблокирующий mio::net::TcpStream: при создании регистрируем интерес на чтение и запись в reactor, а методы read/write возвращают футуры, которые мы уже видели.

use mio::net::TcpStream;
use mio::Interest;
use std::net::SocketAddr;

pub struct AsyncTcpStream {
    inner: TcpStream,
    reactor: Arc<Reactor>,
    token: Token,
}

impl AsyncTcpStream {
    pub fn connect(reactor: Arc<Reactor>, addr: SocketAddr) -> std::io::Result<Self> {
        // mio connect неблокирующий: готовность на запись = соединение установлено.
        let stream = TcpStream::connect(addr)?;
        Self::from_mio(reactor, stream)
    }

    fn from_mio(reactor: Arc<Reactor>, mut stream: TcpStream) -> std::io::Result<Self> {
        let token = reactor.allocate_token();
        reactor.register(&mut stream, token, Interest::READABLE | Interest::WRITABLE)?;
        Ok(AsyncTcpStream { inner: stream, reactor, token })
    }

    pub fn read<'a>(&'a mut self, buf: &'a mut [u8]) -> ReadFut<'a> {
        ReadFut { stream: self, buf }
    }

    pub fn write<'a>(&'a mut self, buf: &'a [u8]) -> WriteFut<'a> {
        WriteFut { stream: self, buf }
    }
}

impl Drop for AsyncTcpStream {
    fn drop(&mut self) {
        // Снимаем регистрацию в reactor, чтобы он не будил мёртвый токен.
        self.reactor.deregister(&mut self.inner);
    }
}

WriteFut зеркалит ReadFut: пробует записать, на WouldBlock кладёт Waker и засыпает.

pub struct WriteFut<'a> {
    stream: &'a mut AsyncTcpStream,
    buf: &'a [u8],
}

impl Future for WriteFut<'_> {
    type Output = std::io::Result<usize>;

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
        let this = self.get_mut();
        match this.stream.inner.write(this.buf) {
            Ok(n) => Poll::Ready(Ok(n)),
            Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {
                this.stream.reactor.store_waker(this.stream.token, cx.waker().clone());
                Poll::Pending
            }
            Err(e) => Poll::Ready(Err(e)),
        }
    }
}

Листенер устроен так же, только интерес у него один (читаемость значит «есть новое соединение»), а его футура accept отдаёт готовый AsyncTcpStream:

use mio::net::TcpListener;

pub struct AsyncTcpListener {
    inner: TcpListener,
    reactor: Arc<Reactor>,
    token: Token,
}

impl AsyncTcpListener {
    pub fn bind(reactor: Arc<Reactor>, addr: SocketAddr) -> std::io::Result<Self> {
        let mut inner = TcpListener::bind(addr)?;
        let token = reactor.allocate_token();
        reactor.register(&mut inner, token, Interest::READABLE)?;
        Ok(AsyncTcpListener { inner, reactor, token })
    }

    pub fn accept(&self) -> Accept<'_> {
        Accept { listener: self }
    }
}

pub struct Accept<'a> {
    listener: &'a AsyncTcpListener,
}

impl Future for Accept<'_> {
    type Output = std::io::Result<AsyncTcpStream>;

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
        let listener = self.listener;
        match listener.inner.accept() {
            Ok((stream, _peer)) => {
                Poll::Ready(AsyncTcpStream::from_mio(listener.reactor.clone(), stream))
            }
            Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {
                listener.reactor.store_waker(listener.token, cx.waker().clone());
                Poll::Pending
            }
            Err(e) => Poll::Ready(Err(e)),
        }
    }
}

Связываем executor и reactor

block_on из начала урока паркует поток. Чтобы дождаться готовности I/O, между poll надо крутить не park, а reactor.turn. Тонкость одна: пробуждение может прийти между poll и turn, и его нельзя проспать. Поэтому Waker тут не паркует поток, а взводит атомарный флаг «меня будили», который цикл проверяет перед каждым poll.

use std::sync::atomic::{AtomicBool, Ordering};

pub fn block_on_with_reactor<F: Future>(reactor: &Reactor, future: F) -> F::Output {
    struct FlagWaker { woken: AtomicBool }
    impl ArcWake for FlagWaker {
        fn wake_by_ref(self: &Arc<Self>) {
            self.woken.store(true, Ordering::SeqCst); // не паркуем, а взводим флаг
        }
    }

    // Стартуем со взведённым флагом: первый poll обязателен.
    let flag = Arc::new(FlagWaker { woken: AtomicBool::new(true) });
    let waker = waker_from_arc(flag.clone());
    let mut cx = Context::from_waker(&waker);

    let mut future = future;
    // SAFETY: future живёт до конца функции и не двигается (только &mut через Pin).
    let mut future = unsafe { Pin::new_unchecked(&mut future) };

    loop {
        // Полим, только если нас будили: swap забирает флаг и сбрасывает его.
        if flag.woken.swap(false, Ordering::SeqCst) {
            if let Poll::Ready(value) = future.as_mut().poll(&mut cx) {
                return value;
            }
        }
        // Прогресса нет: ждём готовности I/O. Короткий таймаут страхует от
        // события, проскочившего между poll и turn (lost wakeup).
        let _ = reactor.turn(Some(Duration::from_millis(50)));
    }
}

Теперь весь блок собирается в одну работающую программу: эхо-сервер и клиент в одной задаче, целиком на нашем рантайме, без tokio.

let reactor = Reactor::new().expect("reactor");
let listener = AsyncTcpListener::bind(reactor.clone(), "127.0.0.1:0".parse().unwrap())
    .expect("bind");
let local = listener.local_addr().expect("addr");

let reactor_for_run = reactor.clone();
let echoed = block_on_with_reactor(&reactor, async move {
    let mut client = AsyncTcpStream::connect(reactor_for_run.clone(), local).expect("connect");
    let mut server = listener.accept().await.expect("accept");

    client.write(b"ping").await.expect("write");

    let mut buf = [0u8; 16];
    let n = server.read(&mut buf).await.expect("read");
    server.write(&buf[..n]).await.expect("echo"); // сервер вернул эхо обратно

    let mut back = [0u8; 16];
    let m = client.read(&mut back).await.expect("read echo");
    back[..m].to_vec()
});
assert_eq!(echoed, b"ping");

Каждый .await тут может вернуть Pending, тогда задача засыпает, block_on_with_reactor зовёт reactor.turn, ОС будит токен, и работа продолжается. Ни одного холостого опроса, и весь код, который это крутит, ты теперь видел целиком.

Как вернуть результат из задачи

Представь веб-сервер с единственным соединением к базе. Из двух мест сразу его трогать нельзя, поэтому соединение отдают одному фоновому работнику: тот держит его и гоняет запросы по очереди. Запросы летят от множества обработчиков, каждый кладёт своё задание в общую очередь. Сложное не в том, как задание отдать, а как вернуть ответ именно тому обработчику, который спросил, а не соседнему. Единого вызывающего, в которого можно было бы просто вернуться, больше нет.

Приём тот же, что дальше мы применим к задаче: дать каждому запросу одноразовый обратный канал. Обработчик создаёт пару из двух концов, отправляющего и принимающего; отправляющий конец он кладёт в задание, а сам ждёт на принимающем. Работник, закончив запрос, проталкивает результат в отправляющий конец, и тот приходит ровно на свой принимающий, и никуда больше. Такой одноразовый канал на одно значение и зовут oneshot-каналом. JoinHandle, который мы сейчас соберём, это его частный случай: задача играет роль отправителя, хендл это получатель.

Заметь: наш spawn берёт футуру с Output = (). Настоящий tokio::spawn возвращает JoinHandle<T>, по которому можно .await-ить результат. Соберём такую ручку поверх нашего executor’а: заведём разделяемую ячейку результата и Waker ожидающего.

struct Shared<T> {
    result: Mutex<Option<T>>,
    waker: Mutex<Option<Waker>>,
}

pub struct JoinHandle<T> {
    shared: Arc<Shared<T>>,
}

impl<T> Future for JoinHandle<T> {
    type Output = T;

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<T> {
        if let Some(value) = self.shared.result.lock().unwrap().take() {
            Poll::Ready(value) // задача уже положила результат
        } else {
            // Ещё нет: запоминаем, кого разбудить, когда результат появится.
            *self.shared.waker.lock().unwrap() = Some(cx.waker().clone());
            Poll::Pending
        }
    }
}

spawn_with_handle оборачивает исходную футуру: дожидается её результата, кладёт в ячейку и будит того, кто держит JoinHandle.

impl Spawner {
    pub fn spawn_with_handle<F>(&self, future: F) -> JoinHandle<F::Output>
    where
        F: Future + Send + 'static,
        F::Output: Send + 'static,
    {
        let shared = Arc::new(Shared {
            result: Mutex::new(None),
            waker: Mutex::new(None),
        });
        let shared_for_task = shared.clone();
        self.spawn(async move {
            let value = future.await;
            *shared_for_task.result.lock().unwrap() = Some(value);
            if let Some(waker) = shared_for_task.waker.lock().unwrap().take() {
                waker.wake(); // будим JoinHandle: результат готов
            }
        });
        JoinHandle { shared }
    }
}

Это ровно та же механика «не звони мне, я позвоню тебе», только теперь между задачей и тем, кто ждёт её результата. Каналы (oneshot и компания) из последнего урока блока обобщают этот приём.

Что унести из урока

Рантайм это три роли: executor крутит poll готовых задач, reactor ловит готовность I/O от ОС через epoll, а Waker их связывает. Waker на низком уровне это сырой указатель плюс vtable из четырёх функций с ручным счётчиком ссылок; в реальном коде его прячут за трейтом Wake. block_on гоняет одну футуру, паркуя поток между poll. Многозадачный executor это очередь готовых задач, где wake задачи просто кладёт её обратно в очередь. Reactor хранит Token -> Waker, и когда ОС сообщает о готовности, будит нужную задачу; так замыкается полный цикл неблокирующего I/O.

Дальше посмотрим, как из этой минимальной схемы вырастает настоящий многопоточный планировщик Tokio: work-stealing, локальные и глобальные очереди, таймерное колесо и кооперативная вытесняемость.

Домашка