Skip to content

Сетевая подсистема

Это ядро производительности Sharp. Все TCP-соединения (клиентские и серверные) обслуживаются одним и тем же паттерном.

Проблема: TCP deadlock

Наивный подход — две таски (reader + writer) или select!-based unified loop:

// ПЛОХО: deadlock при высокой нагрузке
loop {
    tokio::select! {
        pkt = framed.next() => { /* handle */ }
        pkt = outbox.recv() => { framed.feed(pkt).await; }  // ← блокирует!
    }
}

При высоких объёмах (например, 10000 × 16KB) framed.feed().await блокируется на TCP backpressure — send buffer полон, ОС не принимает данные. Пока задача висит на feed(), она не может читать входящие данные. Обе стороны соединения одновременно заблокированы на записи → классический TCP deadlock.

Две отдельные таски (reader + writer) решают deadlock, но создают overhead: два waker'а, двойные epoll_ctl сисколлы, overhead на синхронизацию между тасками.

Решение: poll_fn четырёхфазный цикл

run_connection() в sharp/src/duplex.rs — одна таска, один waker, четыре неблокирующих фазы:

pub async fn run_connection(
    stream: TcpStream,
    outbox: Arc<NotifyChannel<Packet>>,
    mut on_recv: impl FnMut(Packet),
)
poll_fn(|cx| {
    // 1. READ: вычитать всё доступное из TCP
    loop {
        match framed.poll_next_unpin(cx) {
            Ready(Some(Ok(pkt))) => on_recv(pkt),    // синхронный callback
            Ready(Some(Err(e))) => return error,
            Ready(None) => return EOF,
            Pending => break,                         // буфер пуст
        }
    }

    // 2. NOTIFY: забрать пакеты из outbox
    loop {
        match notified.poll(cx) {
            Ready(()) => {
                write_buf.extend(outbox.drain());     // батч drain
                notified.set(outbox.notified());      // новый future
            }
            Pending => break,
        }
    }

    // 3. WRITE: запихнуть в codec
    while let Some(pkt) = write_buf.pop_front() {
        match framed.poll_ready_unpin(cx) {
            Ready(Ok(())) => framed.start_send_unpin(pkt),  // в буфер кодека
            Pending => { write_buf.push_front(pkt); break } // backpressure
            Ready(Err(e)) => return error,
        }
    }

    // 4. FLUSH: сбросить буфер кодека в TCP
    if needs_flush {
        match framed.poll_flush_unpin(cx) {
            Ready(Ok(())) => needs_flush = false,
            Pending => {},                            // продолжим позже
            Ready(Err(e)) => return error,
        }
    }

    Poll::Pending  // один waker на все источники
})

Почему не дедлочит: каждая фаза неблокирующая. Если TCP write buffer полон (WRITE возвращает Pending), пакет возвращается в write_buf, цикл завершается. При следующем wakeup (от входящих данных или от освобождения write buffer) READ фаза снова отработает первой.

Почему быстрее двух тасок (~30% прирост):

  • Один waker вместо двух → меньше epoll_ctl сисколлов
  • Нет inter-task синхронизации (всё в одной таске)
  • Батчинг: NOTIFY дренит весь outbox за один lock → несколько пакетов за один wakeup
  • Нет overhead select! macro (создание futures + pin на каждый poll)

NotifyChannel

pub struct NotifyChannel<T> {
    queue: Mutex<VecDeque<T>>,
    notify: Notify,
}

Легковесный канал для outbox. Используется всеми компонентами:

Метод Тип Описание
push(item) sync Кладёт в очередь + notify_one()
drain() sync Забирает всё из очереди за один lock
recv() async Ждёт items, loop для spurious wakeups
notified() sync → Future Возвращает Notified future для poll_fn

Важный нюанс: recv() содержит внутренний цикл. Два быстрых push() могут колапснуться в один permit (второй permit теряется). Цикл гарантирует, что recv() не вернёт пустой результат.

В контексте poll_fn используется notified() + drain() вместо recv().

Использование по компонентам

Компонент Использует run_connection()? Свой poll_fn?
Publisher Да Нет
Subscriber Да Нет
Server Нет Да (тот же паттерн, но с серверной логикой в callback)

Сервер реализует тот же четырёхфазный цикл в handle_connection(), но с более сложной обработкой входящих пакетов (routing, sliding window, subscriber management).