Сетевая подсистема¶
Это ядро производительности 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¶
Легковесный канал для 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).