diff --git a/lib/tcp_io.c b/lib/tcp_io.c index bc9708a9..9dde20e0 100644 --- a/lib/tcp_io.c +++ b/lib/tcp_io.c @@ -20,7 +20,8 @@ static void error_cb(socket_t sock, void* arg); static void resume_read_cb(struct ll_queue* q, void* arg); static void fin_deferred_cb(struct ll_queue* q, void* arg); static void write_queue_fetch_cb(struct ll_queue* q, void* arg); -static void flush_write_buf(struct tcp_conn* tc); +static int flush_write_buf(struct tcp_conn* tc); +static void tcp_conn_handle_error(struct tcp_conn* tc, int err); static const uint8_t tcp_fin_sentinel; struct tcp_conn* tcp_conn_create( @@ -128,6 +129,45 @@ void tcp_conn_destroy(struct tcp_conn* tc) { u_free(tc); } +static void tcp_conn_handle_error(struct tcp_conn* tc, int err) +{ + tc->error = 1; + + queue_waiter_cancel(tc->read_queue, &tc->read_waiter); + queue_set_empty_callback(tc->read_queue, NULL, NULL); + queue_set_callback(tc->read_queue, NULL, NULL); + + if (tc->sock != SOCKET_INVALID) { + while (1) { + struct ll_entry* e = queue_entry_new_from_pool(tc->entry_pool); + if (!e) break; + uint8_t* buf = memory_pool_alloc(tc->data_pool); + if (!buf) { queue_entry_free(e); break; } + ssize_t n = recv(tc->sock, buf, tc->entry_data_size, 0); + if (n > 0) { e->dgram = buf; e->len = (uint16_t)n; queue_data_put(tc->read_queue, e); } + else { memory_pool_free(tc->data_pool, buf); queue_entry_free(e); break; } + } + } + + if (tc->socket_id) { + uasync_remove_socket_t(tc->ua, tc->sock); + tc->socket_id = NULL; + } + if (tc->sock != SOCKET_INVALID) { + socket_close_wrapper(tc->sock); + tc->sock = SOCKET_INVALID; + } + + queue_set_callback(tc->write_queue, NULL, NULL); + + if (tc->on_error) { + void (*on_err)(struct tcp_conn*, int, void*) = tc->on_error; + void* arg = tc->arg; + tc->on_error = NULL; + on_err(tc, err, arg); + } +} + // ==================================================================== // Чтение из сокета // ==================================================================== @@ -174,9 +214,8 @@ static void read_cb(socket_t sock, void* arg) { memory_pool_free(tc->data_pool, buf); queue_entry_free(e); if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINTR) return; - tc->error = 1; DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: recv error fd=%d errno=%d", (int)tc->sock, errno); - if (tc->on_error) tc->on_error(tc, errno, tc->arg); + tcp_conn_handle_error(tc, errno); } } @@ -198,7 +237,7 @@ static void fin_deferred_cb(struct ll_queue* q, void* arg) { // Запись в сокет // ==================================================================== -static void flush_write_buf(struct tcp_conn* tc) { +static int flush_write_buf(struct tcp_conn* tc) { while (tc->write_buf && tc->write_offset < tc->write_len) { ssize_t n = send(tc->sock, tc->write_buf + tc->write_offset, tc->write_len - tc->write_offset, MSG_NOSIGNAL); if (n > 0) { @@ -209,16 +248,16 @@ static void flush_write_buf(struct tcp_conn* tc) { tc->write_len = 0; tc->write_offset = 0; DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: flush_wbuf done fd=%d", (int)tc->sock); - return; + return 0; } continue; } - if (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINTR)) return; - tc->error = 1; + if (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINTR)) return 0; DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: send error fd=%d errno=%d", (int)tc->sock, errno); - if (tc->on_error) tc->on_error(tc, errno, tc->arg); - return; + tcp_conn_handle_error(tc, errno); + return 1; } + return 0; } static void write_queue_fetch_cb(struct ll_queue* q, void* arg) { @@ -227,7 +266,7 @@ static void write_queue_fetch_cb(struct ll_queue* q, void* arg) { if (tc->error || tc->closed) { DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch skip fd=%d error=%d closed=%d", (int)tc->sock, tc->error, tc->closed); return; } if (tc->write_buf) { - flush_write_buf(tc); + if (flush_write_buf(tc)) return; if (tc->write_buf) return; // остался остаток, ждём write_cb } @@ -278,9 +317,9 @@ static void write_queue_fetch_cb(struct ll_queue* q, void* arg) { tc->write_buf = memory_pool_alloc(tc->data_pool); if (!tc->write_buf) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: write_buf alloc failed fd=%d", (int)tc->sock); - tc->error = 1; memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); + tcp_conn_handle_error(tc, ENOMEM); return; } memcpy(tc->write_buf, e->dgram + n, rem); @@ -295,9 +334,9 @@ static void write_queue_fetch_cb(struct ll_queue* q, void* arg) { tc->write_buf = memory_pool_alloc(tc->data_pool); if (!tc->write_buf) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: write_buf alloc failed fd=%d", (int)tc->sock); - tc->error = 1; memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); + tcp_conn_handle_error(tc, ENOMEM); return; } memcpy(tc->write_buf, e->dgram, e->len); @@ -309,11 +348,10 @@ static void write_queue_fetch_cb(struct ll_queue* q, void* arg) { tc->write_monitor = 1; DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch POLLOUT ON fd=%d (EAGAIN, wbuf=%zu bytes)", (int)tc->sock, tc->write_len); } else { - tc->error = 1; DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: send error fd=%d errno=%d", (int)tc->sock, errno); - if (tc->on_error) tc->on_error(tc, errno, tc->arg); memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); + tcp_conn_handle_error(tc, errno); } } @@ -329,15 +367,13 @@ static void write_cb(socket_t sock, void* arg) { tc->connected = 1; DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: connect ok fd=%d", (int)tc->sock); } else { - tc->error = 1; - DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: write_cb connect fail fd=%d err=%d wmon=%d (POLLOUT still active)", (int)tc->sock, err, tc->write_monitor); DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: connect fail fd=%d err=%d", (int)tc->sock, err); - if (tc->on_error) tc->on_error(tc, err ? err : -1, tc->arg); + tcp_conn_handle_error(tc, err ? err : -1); return; } } - flush_write_buf(tc); + if (flush_write_buf(tc)) return; if (!tc->write_buf) { DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: write_cb flushed fd=%d wmon=%d (resume WQ cb)", (int)tc->sock, tc->write_monitor); queue_resume_callback(tc->write_queue); @@ -352,10 +388,8 @@ static void error_cb(socket_t sock, void* arg) { (void)sock; struct tcp_conn* tc = (struct tcp_conn*)arg; if (!tc || tc->sock == SOCKET_INVALID) return; - tc->error = 1; - DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: error_cb fd=%d connected=%d fin_remote=%d wmon=%d", (int)tc->sock, tc->connected, tc->fin_remote, tc->write_monitor); DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: async error fd=%d", (int)tc->sock); - if (tc->on_error) tc->on_error(tc, -1, tc->arg); + tcp_conn_handle_error(tc, -1); } // ==================================================================== diff --git a/lib/tcp_io.h b/lib/tcp_io.h index 19820684..c076f0d8 100644 --- a/lib/tcp_io.h +++ b/lib/tcp_io.h @@ -1,40 +1,23 @@ // tcp_io.h — управление TCP-соединением через uasync + ll_queue // -// Один tcp_conn = одно TCP-соединение. Две очереди: read_queue (сокет → данные) и -// write_queue (данные → сокет). Обе работают через автозабор ll_queue (deferred). +// Один tcp_conn = одно TCP-соединение. Два пула: +// entry_pool — struct ll_entry (без inline data) +// data_pool — буферы (max(entry_data_size, write_chunk_size)) // -// Пул-аллокация: -// entry_pool — только struct ll_entry (без inline data) -// data_pool — буферы данных чтения/записи (max(entry_data_size, write_chunk_size)) -// Данные и структуры аллоцируются раздельно — не копируются при recv/send. +// Чтение: read_cb → recv() в data_pool → entry в read_queue → callback внешнего кода. +// high_water/low_water управляют backpressure (EPOLLIN on/off). +// FIN (recv==0) → on_fin (сразу или через empty_callback после дренажа очереди). // -// Чтение: -// read_cb → memory_pool_alloc(read_pool) → recv() прямо в буфер → entry в read_queue -// read_queue автозабор(deferred) → внешний коллбэк (например, отправка в ETCP) -// high_water → пауза чтения (убираем EPOLLIN), low_water/waiter → возобновление -// FIN: если read_queue не пуст → empty_callback откладывает on_fin; иначе сразу +// Запись: внешний код кладёт entry в write_queue через queue_data_put(). +// write_queue_fetch_cb (deferred) → send(), при EAGAIN → write_buf + EPOLLOUT. +// Всё отправлено → on_flushed. // -// Запись: -// Внешний код: tcp_conn_push_write(data, len) → аллокация через write_pool → entry в write_queue -// write_queue автозабор(deferred): fetch_cb → send() в сокет -// EAGAIN → write_buf (из data_pool) + EPOLLOUT ON → write_cb досылает → resume автозабора -// Всё отправлено → EPOLLOUT OFF + on_flushed (если установлен) +// FIN/Close: tcp_conn_push_fin/close ставят сентинел в write_queue. +// FIN → shutdown(SHUT_WR) → on_fin_sent +// CLOSE → close сокета → on_closed // -// FIN / Close через очередь отправки: -// tcp_conn_push_fin(tc) — ставит FIN в write_queue (dgram=&sentinel, len=0) -// tcp_conn_push_close(tc) — ставит CLOSE в write_queue (dgram=NULL, len=0) -// Оба сигнала обрабатываются после всех предшествующих данных в очереди: -// FIN → shutdown(SHUT_WR), стоп чтения, on_fin_sent -// CLOSE → close сокета, on_closed -// После FIN входящие данные отбрасываются (read_cb → discard). -// После CLOSE сокет закрыт, tc->closed=1, дальнейшая отправка невозможна. -// -// Connect: -// tcp_conn_create регистрирует сокет с read_cb + write_cb в uasync -// После create вызывается connect() (неблокирующий, EINPROGRESS) -// write_cb детектит завершение connect через getsockopt(SO_ERROR) -// getpeername() сразу после create проверяет pre-connected сокеты (socketpair) -// До connect данные копятся в write_queue (автозабор не пытается send на unconnected сокет) +// Ошибка: любая ошибка сокета → tcp_conn_handle_error → close + null коллбэков → on_error. +// Реализация on_error обязана вызвать tcp_conn_destroy. После возврата tc недействителен. #ifndef TCP_IO_H #define TCP_IO_H @@ -63,7 +46,7 @@ struct tcp_conn { uint8_t fin_local; // FIN отправлен удалённой стороне (shutdown SHUT_WR) uint8_t closed; // сокет полностью закрыт (close) - // Частичная отправка (из write_pool, не в очереди — досылается первой) + // Частичная отправка (из data_pool, не в очереди — досылается первой) uint8_t* write_buf; size_t write_len; size_t write_offset; @@ -79,7 +62,7 @@ struct tcp_conn { // Коллбэки 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); // ошибка сокета + 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; @@ -96,17 +79,17 @@ struct tcp_conn* tcp_conn_create( void (*on_error)(struct tcp_conn* tc, int err, void* arg), void* arg); +// Дренирует очереди, закрывает сокет, уничтожает пулы, освобождает tc. void tcp_conn_destroy(struct tcp_conn* tc); -// Внешний код пишет данные в tc->write_queue напрямую (queue_data_put). -// Автозабор write_queue (deferred) сам отправляет когда сокет готов. -// Перед push проверять порог: queue_set_threshold в tcp_conn_create (32 entries). -// При заполнении — queue_waiter_wait на освобождение. +// Запись: 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. +// shutdown(SHUT_WR) и on_fin_sent. // Повторный вызов игнорируется (fin_sent уже установлен). // tcp_conn_push_close(tc) — CLOSE: после отправки предшествующих данных закрывает // сокет (close), вызывает on_closed.