// tcp_io.h — управление TCP-соединением через uasync + ll_queue // // Один tcp_conn = одно TCP-соединение. Два пула: // entry_pool — struct ll_entry (без inline data) // data_pool — буферы (max(entry_data_size, write_chunk_size)) // // Чтение: read_cb → recv() в data_pool → entry в read_queue → callback внешнего кода. // high_water/low_water управляют backpressure (EPOLLIN on/off). // FIN (recv==0) → on_fin (сразу или через empty_callback после дренажа очереди). // // Запись: внешний код кладёт entry в write_queue через queue_data_put(). // write_queue_fetch_cb (deferred) → send(), при EAGAIN → write_buf + EPOLLOUT. // Всё отправлено → on_flushed. // // FIN/Close: tcp_conn_push_fin/close ставят сентинел в write_queue. // FIN → shutdown(SHUT_WR) → on_fin_sent // CLOSE → close сокета → on_closed // // Ошибка: любая ошибка сокета → tcp_conn_handle_error → close + null коллбэков → on_error. // Реализация on_error обязана вызвать tcp_conn_destroy. После возврата tc недействителен. #ifndef TCP_IO_H #define TCP_IO_H #ifdef __cplusplus extern "C" { #endif #include "u_async.h" #include "socket_compat.h" #include "ll_queue.h" #include "memory_pool.h" struct tcp_conn { socket_t sock; struct UASYNC* ua; void* socket_id; struct ll_queue* read_queue; // сокет → данные (блоки до entry_data_size) struct ll_queue* write_queue; // данные → сокет (блоки до write_chunk_size) int read_high_water; int read_low_water; int rcvbuf_size; uint8_t read_paused; uint8_t write_monitor; // 1 = EPOLLOUT активен uint8_t connected; uint8_t error; uint8_t fin_remote; // FIN получен от удалённой стороны (recv == 0) uint8_t fin_local; // FIN отправлен удалённой стороне (shutdown SHUT_WR) uint8_t closed; // сокет полностью закрыт (close) uint8_t destroyed; // 1 = tcp_conn_destroy вызван, tc ожидает отложенного free // Частичная отправка (из data_pool, не в очереди — досылается первой) uint8_t* write_buf; size_t write_len; size_t write_offset; size_t write_chunk_size; // Пулы памяти struct memory_pool* entry_pool; // sizeof(struct ll_entry) struct memory_pool* data_pool; // max(entry_data_size, write_chunk_size) size_t entry_data_size; struct queue_waiter_handle read_waiter; // Коллбэки void (*on_fin)(struct tcp_conn* tc, void* arg); // FIN получен от удалённой стороны (fin_remote=1) void (*on_fin_sent)(struct tcp_conn* tc, void* arg); // FIN отправлен удалённой стороне (fin_local=1) void (*on_error)(struct tcp_conn* tc, int err, void* arg); // фатальная ошибка (вызывается после close сокета; обязан вызвать tcp_conn_destroy) void (*on_flushed)(struct tcp_conn* tc, void* arg); // все данные записи отправлены (write_queue + write_buf пусты) void (*on_closed)(struct tcp_conn* tc, void* arg); // сокет закрыт через очередь (после close) void* arg; }; struct tcp_conn* tcp_conn_create( struct UASYNC* ua, socket_t sock, size_t entry_data_size, size_t write_chunk_size, int read_high_water, int read_low_water, int rcvbuf_size, void (*on_fin)(struct tcp_conn* tc, void* arg), void (*on_error)(struct tcp_conn* tc, int err, void* arg), void* arg); // Дренирует очереди, закрывает сокет, уничтожает пулы, освобождает tc. void tcp_conn_destroy(struct tcp_conn* tc); // Запись: queue_data_put(tc->write_queue, entry) напрямую. // Автозабор (deferred) сам шлёт когда сокет готов. // Порог: 32 записи, при заполнении queue_waiter_wait. // Поставить сигналы в очередь отправки. Все данные в очереди перед сигналом // будут отправлены до его обработки. Кодирование: dgram=NULL → close, len=0+dgram!=NULL → FIN. // tcp_conn_push_fin(tc) — FIN: после отправки предшествующих данных вызывает // shutdown(SHUT_WR) и on_fin_sent. // Повторный вызов игнорируется (fin_sent уже установлен). // tcp_conn_push_close(tc) — CLOSE: после отправки предшествующих данных закрывает // сокет (close), вызывает on_closed. // После close сокет удалён из uasync, tc->closed=1. // Повторный вызов игнорируется. int tcp_conn_push_fin(struct tcp_conn* tc); int tcp_conn_push_close(struct tcp_conn* tc); // Одноразовый коллбэк: вызывается когда write_queue + write_buf полностью опустели. // После вызова сбрасывается. Установить повторно можно в любой момент. void tcp_conn_set_flushed(struct tcp_conn* tc, void (*on_flushed)(struct tcp_conn* tc, void* arg)); // Принудительная остановка чтения: убирает EPOLLIN и отменяет read_waiter. void tcp_conn_pause_read(struct tcp_conn* tc); #ifdef __cplusplus } #endif #endif