8.7 KiB
tcp_io — Управление TCP-соединением на uasync + ll_queue
1. Назначение
Асинхронный TCP-коннектор, построенный на event loop uasync и lock-free очередях ll_queue.
Одна структура struct tcp_conn = одно TCP-соединение.
Используется в SOCKS/HTTP-прокси, STCP, TCP-прокси и тестах.
Ключевые свойства:
- Чтение, запись, ошибки и connect — через
uasync(epoll/kqueue/poll). - Два независимых
ll_queue:read_queue(входящие данные) иwrite_queue(исходящие). - Два memory pool:
entry_pool(struct ll_entry) иdata_pool(буферы), всё выделение на горячем пути через пулы. - Backpressure на чтение: high_water/low_water +
queue_waiter_wait, EPOLLIN вкл/выкл по уровню заполнения очереди. - Backpressure на запись: порог 32 записи в write_queue +
queue_waiter_wait, EPOLLOUT вкл/выкл при частичной отправке. - Graceful shutdown: FIN (shutdown SHUT_WR) и CLOSE (close сокета) — оба как сентинелы в write_queue, все предшествующие данные гарантированно отправлены.
2. Как пользоваться
Создание соединения
struct tcp_conn* tc = tcp_conn_create(
ua, sock,
1500, // entry_data_size — размер буфера для одного recv
8192, // write_chunk_size — размер буфера для одной send-операции
32, // read_high_water — порог приостановки EPOLLIN
8, // read_low_water — порог возобновления EPOLLIN
0, // rcvbuf_size — размер буфера сокета (0 = не менять)
on_fin_cb, on_error_cb, arg);
После создания нужно установить callback на read_queue и включить отложенную обработку:
queue_set_callback(tc->read_queue, on_read_cb, my_conn);
queue_set_waiter_defer(tc->read_queue, 1);
Чтение данных (в on_read_cb)
static void on_read_cb(struct ll_queue* q, void* arg) {
struct my_conn* c = (struct my_conn*)arg;
struct ll_entry* e = queue_data_get(q);
if (!e) { queue_resume_callback(q); return; }
// обработать e->dgram (e->len байт)
queue_entry_free(e);
queue_resume_callback(q); // обязательно!
}
Запись данных
Пользователь напрямую кладёт entry в write_queue, отправка авто (deferred callback сам шлёт когда сокет готов):
struct ll_entry* e = queue_entry_new_from_pool(tc->entry_pool);
uint8_t* buf = memory_pool_alloc(tc->data_pool);
memcpy(buf, data, len);
e->dgram = buf; e->len = (uint16_t)len;
queue_data_put(tc->write_queue, e);
При заполнении (порог 32) нужно ждать через queue_waiter_wait.
Graceful shutdown
tcp_conn_push_fin(tc); // shutdown(SHUT_WR) после отправки всех данных → on_fin_sent
tcp_conn_push_close(tc); // close сокета после отправки всех данных → on_closed
Можно вызывать из callback-ов (on_fin, on_read_cb и т.д.). Повторные вызовы игнорируются (fin_local=1/closed=1 — возвращают -1).
Уничтожение
void tcp_conn_destroy(struct tcp_conn* tc);
Идемпотентен (destroyed флаг). Дренирует обе очереди, закрывает сокет, освобождает пулы через отложенный uasync_call_soon.
Обработка ошибок
on_error(tc, err, arg) вызывается после закрытия сокета. Реализация on_error обязана вызвать tcp_conn_destroy. После возврата из on_error tc недействителен.
static void on_error_cb(struct tcp_conn* tc, int err, void* arg) {
struct my_conn* c = (struct my_conn*)arg;
// почистить свои ресурсы
tcp_conn_destroy(tc); // tc после этого использовать нельзя
socket_close_wrapper(other_sock);
u_free(c);
}
on_flushed (одноразовый)
tcp_conn_set_flushed(tc, my_flushed_cb);
// коллбэк вызовется один раз когда write_queue + write_buf полностью опустеют, затем сбросится
Приостановка чтения
tcp_conn_pause_read(tc); // убирает EPOLLIN и отменяет read_waiter
Ограничения и многопоточность
- Однопоточный. Всё работает в одном
uasyncevent loop. Нельзя вызывать из другого потока. - FIN-сентинел в очереди. После
push_finданные продолжат отправляться, сам FIN выполнится только когда все предшествующие данные отправлены. - Нельзя вызывать close сокета напрямую. Только через
tcp_conn_push_close— иначе нарушится порядок обработки очереди. - Можно вызывать push_fin/push_close из callback-ов — защита от повторного вызова встроена (fin_local/closed флаги).
3. API
Структуры
| Поле | Назначение |
|---|---|
struct tcp_conn |
Всё состояние TCP-соединения: сокет, очереди, пулы, коллбэки, флаги |
Жизненный цикл
| Функция | Описание |
|---|---|
tcp_conn_create(ua, sock, entry_data_size, write_chunk_size, read_high_water, read_low_water, rcvbuf_size, on_fin, on_error, arg) |
Создать соединение: аллоцирует tc, два memory pool, два ll_queue, регистрирует сокет в uasync с EPOLLIN+EPOLLOUT. Вовращает NULL при ошибке |
tcp_conn_destroy(tc) |
Дренирует read_queue и write_queue, закрывает сокет, освобождает пулы через uasync_call_soon. Идемпотентен |
Коллбэки (поля struct tcp_conn)
| Поле | Когда вызывается |
|---|---|
on_fin(tc, arg) |
FIN получен от удалённой стороны (recv==0): сразу если read_queue пуста, иначе через queue_set_empty_callback после дренажа |
on_fin_sent(tc, arg) |
FIN отправлен (shutdown(SHUT_WR)) — после отправки всех данных перед FIN-сентинелом |
on_error(tc, err, arg) |
Фатальная ошибка сокета. Вызывается после close сокета. Обязан вызвать tcp_conn_destroy. После возврата tc недействителен |
on_flushed(tc, arg) |
write_queue + write_buf полностью опустели. Одноразовый — после вызова сбрасывается в NULL. Установка: tcp_conn_set_flushed() |
on_closed(tc, arg) |
Сокет закрыт через CLOSE-сентинел (tcp_conn_push_close) после отправки всех предшествующих данных |
Graceful shutdown
| Функция | Описание |
|---|---|
tcp_conn_push_fin(tc) |
Помещает FIN-сентинел (dgram=&tcp_fin_sentinel, len=0) в write_queue. После отправки предшествующих данных: shutdown(SHUT_WR), fin_local=1, вызывает on_fin_sent. Повторный вызов → -1 |
tcp_conn_push_close(tc) |
Помещает CLOSE-сентинел (dgram=NULL, len=0) в write_queue. После отправки предшествующих данных: close сокета, удаление из uasync, closed=1, вызывает on_closed. Повторный вызов → -1 |
Прочее
| Функция | Описание |
|---|---|
tcp_conn_set_flushed(tc, cb) |
Установить одноразовый коллбэк на полное опустошение write_queue |
tcp_conn_pause_read(tc) |
Принудительно выключить EPOLLIN и отменить ожидание низкого порога |